From 78c6bfe6cadea75101b3ecbcdfdefd363f29bc7b Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Thu, 24 Sep 2026 12:26:31 +0200 Subject: [PATCH 1/8] feat: validate discovery and rubric attribution --- rfcs/008-environment-auto-validation.md | 17 + .../validation/graders/runtime/discovery.py | 329 +++++++++++++++ src/openenv/validation/runner.py | 39 +- src/openenv/validation/runtime/artifacts.py | 24 +- src/openenv/validation/runtime/collector.py | 80 +++- src/openenv/validation/runtime/contracts.py | 12 + src/openenv/validation/runtime/discovery.py | 95 +++++ .../validation/runtime/served_probe/app.py | 54 ++- .../runtime/served_probe/openenv.yaml | 5 +- .../integration/test_discovery_protocol.py | 82 ++++ .../integration/test_runtime_cli.py | 73 +++- .../integration/test_runtime_process.py | 26 +- .../test_discovery_collection.py | 264 ++++++++++++ .../test_validation/test_runtime_artifacts.py | 93 ++++- .../test_validation/test_runtime_discovery.py | 385 ++++++++++++++++++ .../test_validation/test_runtime_execution.py | 175 +++++++- .../test_task_discovery_collection.py | 127 ++++++ tests/validation_runtime/README.md | 21 +- tests/validation_runtime/acceptance.json | 19 +- 19 files changed, 1876 insertions(+), 44 deletions(-) create mode 100644 src/openenv/validation/graders/runtime/discovery.py create mode 100644 src/openenv/validation/runtime/discovery.py create mode 100644 tests/test_validation/integration/test_discovery_protocol.py create mode 100644 tests/test_validation/test_discovery_collection.py create mode 100644 tests/test_validation/test_runtime_discovery.py create mode 100644 tests/test_validation/test_task_discovery_collection.py diff --git a/rfcs/008-environment-auto-validation.md b/rfcs/008-environment-auto-validation.md index f49732d08..2fc615a20 100644 --- a/rfcs/008-environment-auto-validation.md +++ b/rfcs/008-environment-auto-validation.md @@ -629,6 +629,23 @@ limited to 100 steps/202 operations and 8 MiB, and marks truncation explicitly. Missing configuration, attribution or records cannot be inferred from other successful operations. This transport introduces no new passing grader by itself. +#### Declaration and rubric evidence (PR6) + +The validator discovers MCP tools over the authenticated replay connection before +reset and compares their names with the declaration, including an explicitly empty +set. Task discovery queries split names and authoritative per-split counts, then +reads at most two individual tasks per split. Preview length is never treated as +the dataset size. Requests share the collection deadline and bounded response +budgets; unavailable or malformed discovery cannot become an empty passing set. + +Rubric checks require the declared tree and explicit public configuration, then +compare every observed step reward with its fresh root score and named child +attribution. Built-in weighted, sequential and gated aggregation is checked using +its published configuration. Unknown custom aggregation remains incomplete. +Optional task and rubric checks are selected only when declared. `discovery.json` +retains raw discovery results, availability and sanitized errors alongside session +telemetry; ordinary discovery responses never enter the episode trajectory. + Applicability predicates must distinguish empty declared sets from absent capabilities. Missing subject features, missing provider support and checks whose implementation has not shipped are distinct outcomes. Independent graders must diff --git a/src/openenv/validation/graders/runtime/discovery.py b/src/openenv/validation/graders/runtime/discovery.py new file mode 100644 index 000000000..8c76540bd --- /dev/null +++ b/src/openenv/validation/graders/runtime/discovery.py @@ -0,0 +1,329 @@ +"""Declaration and rubric checks over bounded, independently collected evidence.""" + +import json +import math +import time + +from ...report import CheckResult +from ...types import CheckStatus +from .basic import _RuntimeGrader + + +def _finite(value): + try: + return type(value) in (int, float) and math.isfinite(value) + except OverflowError: + return False + + +def _json_float(value): + number = float(value) + if not math.isfinite(number): + raise ValueError("runtime JSON numbers must be finite") + return number + + +class _EvidenceGrader(_RuntimeGrader): + """Discovery can be measured without an episode step; missing evidence cannot pass.""" + + def run(self, subject): + started = time.monotonic() + problems, incomplete = [], None + evidence = subject.runtime_evidence + if not self.applies_to(subject.manifest): + incomplete = "declared capability does not apply" + elif evidence is None: + incomplete = "runtime evidence is unavailable" + elif getattr(evidence, self.error_field): + problems = [f"{self.label} collection failed"] + elif getattr(evidence, self.field) is None: + incomplete = f"missing prerequisite: {self.label} evidence" + else: + try: + payload = json.loads( + getattr(evidence, self.field), + parse_float=_json_float, + parse_constant=_json_float, + ) + if not isinstance(payload, dict): + raise ValueError("evidence must be an object") + problems, incomplete = self.check(subject, payload) + except (ValueError, TypeError, KeyError, RecursionError, OverflowError): + problems = [f"malformed {self.label} evidence"] + status = ( + CheckStatus.FAIL + if problems + else CheckStatus.SKIP + if incomplete + else CheckStatus.PASS + ) + return CheckResult( + check_id=self.check_id, + status=status, + evidence=problems[:20] + or [incomplete or f"observed {self.label} contract holds"], + duration_s=time.monotonic() - started, + ) + + +class ToolDeclarationAccuracyGrader(_EvidenceGrader): + check_id = "runtime.tool_declaration_accuracy" + field, error_field, label = "tools_json", "tools_error", "tool discovery" + + def check(self, subject, payload): + tools = payload["tools"] + if not isinstance(tools, list): + raise ValueError("tools must be a list") + names = [tool["name"] for tool in tools] + if any(not isinstance(name, str) or not name for name in names): + raise ValueError("invalid tool name") + cursor = payload.get("nextCursor") + if cursor is not None and not isinstance(cursor, str): + raise ValueError("invalid tool pagination cursor") + declared = subject.manifest.capabilities.declared_tools + problems = [] + if len(set(names)) != len(names): + problems.append("discovery contains duplicate tool names") + if len(set(declared)) != len(declared): + problems.append("declaration contains duplicate tool names") + missing = set(declared) - set(names) + extra = set(names) - set(declared) + if missing and cursor is None: + problems.append(f"{len(missing)} declared tools are missing from discovery") + if extra: + problems.append(f"{len(extra)} discovered tools are undeclared") + return problems, ( + "tool inventory is incomplete: pagination was not collected" + if cursor is not None + else None + ) + + +class TaskDeclarationAccuracyGrader(_EvidenceGrader): + check_id = "runtime.task_declaration_accuracy" + field, error_field, label = "tasks_json", "tasks_error", "task discovery" + + def applies_to(self, manifest): + return manifest.capabilities.task_api or bool( + manifest.capabilities.declared_task_count + ) + + def check(self, subject, payload): + splits, counts, previews = ( + payload["splits"], + payload["counts"], + payload["previews"], + ) + if ( + not isinstance(splits, list) + or not isinstance(counts, dict) + or not isinstance(previews, dict) + ): + raise ValueError("invalid task discovery") + names = [split["name"] for split in splits] + if any(not isinstance(name, str) or not name for name in names) or len( + set(names) + ) != len(names): + raise ValueError("invalid split names") + if set(names) != counts.keys() or set(names) != previews.keys(): + raise ValueError("incomplete split discovery") + problems = [] + declared = subject.manifest.capabilities.declared_task_count + if set(names) != declared.keys(): + problems.append("discovered splits differ from declared task counts") + for name in names: + count, preview = counts[name], previews[name] + if type(count) is not int or count < 0: + raise ValueError("invalid task count") + if not isinstance(preview, list): + raise ValueError("invalid task preview") + if len(preview) > count or (count > 0 and not preview): + problems.append( + "task preview is inconsistent with the advertised count" + ) + if name in declared and count != declared[name]: + problems.append("advertised task count differs from the declaration") + return problems, None + + +def _rubric_nodes(payload): + if not isinstance(payload, list) or not payload: + raise ValueError("missing rubric tree") + nodes = {} + for node in payload: + name = node["name"] + if not isinstance(name, str) or not all(name.split(".")) or name in nodes: + raise ValueError("invalid rubric name") + if not isinstance(node["class_name"], str) or not node["class_name"]: + raise ValueError("missing rubric class") + children = node["children"] + if not isinstance(children, list) or any( + not isinstance(child, str) for child in children + ): + raise ValueError("invalid rubric children") + if len(set(children)) != len(children): + raise ValueError("duplicate rubric children") + if type(node["config_available"]) is not bool or not isinstance( + node["config"], dict + ): + raise ValueError("invalid rubric configuration") + if type(node["evaluated"]) is not bool: + raise ValueError("invalid evaluation marker") + if node["evaluated"] and not _finite(node["score"]): + raise ValueError("evaluated score is not finite") + if not node["evaluated"] and node["score"] is not None: + raise ValueError("unevaluated score must be absent") + aggregation = node["aggregation"] + if aggregation not in {"weighted_sum", "sequential", "gate", "leaf", "unknown"}: + raise ValueError("invalid aggregation") + if aggregation == "leaf" and children: + raise ValueError("leaf has children") + if aggregation == "weighted_sum": + weights = node["config"]["weights"] + if ( + not isinstance(weights, list) + or len(weights) != len(children) + or not all(_finite(weight) for weight in weights) + ): + raise ValueError("invalid rubric weights") + if not math.isclose(sum(weights), 1, rel_tol=0, abs_tol=1e-6): + raise ValueError("rubric weights do not sum to one") + if aggregation == "gate" and ( + len(children) != 1 or not _finite(node["config"]["threshold"]) + ): + raise ValueError("invalid rubric gate") + nodes[name] = node + if "root" not in nodes: + raise ValueError("missing rubric root") + linked = set() + for name, node in nodes.items(): + for child in node["children"]: + if ( + child not in nodes + or child in linked + or child.rsplit(".", 1)[0] != name + or child == name + ): + raise ValueError("invalid rubric edge") + linked.add(child) + if linked != nodes.keys() - {"root"}: + raise ValueError("disconnected rubric tree") + return nodes + + +class RubricIntrospectableGrader(_EvidenceGrader): + check_id = "runtime.rubric_introspectable" + requires_capabilities = frozenset({"rubric_tree"}) + field, error_field, label = "telemetry_json", "telemetry_error", "rubric telemetry" + + def applies_to(self, manifest): + return manifest.capabilities.rubric_tree + + def check(self, subject, payload): + if type(payload["schema_version"]) is not int or payload["schema_version"] != 1: + raise ValueError("unsupported telemetry version") + if payload.get("rubric_error") is not None: + return ["subject reported unavailable rubric introspection"], None + nodes = _rubric_nodes(payload["rubric"]) + return ( + ["rubric configuration is unavailable"] + if any(not node["config_available"] for node in nodes.values()) + else [] + ), None + + +class RewardAttributionGrader(RubricIntrospectableGrader): + check_id = "runtime.reward_attribution" + depends_on = ("runtime.startup", "runtime.rubric_introspectable") + + def check(self, subject, payload): + problems, _ = super().check(subject, payload) + if subject.runtime_evidence.failure_reason: + problems.append("wire replay is incomplete") + if problems: + return problems, None + baseline = _rubric_nodes(payload["rubric"]) + steps = [ + row for row in subject.runtime_evidence.exchanges if row.operation == "step" + ] + records = payload["attribution"] + if not isinstance(records, list): + raise ValueError("invalid attribution") + if not steps: + return problems, "missing prerequisite: an observed step reward" + if len(records) != len(steps): + return problems + ["attribution does not cover every observed step"], None + + def definition(tree): + return { + name: { + key: value + for key, value in node.items() + if key not in {"score", "evaluated"} + } + for name, node in tree.items() + } + + unsupported = False + for index, (record, step) in enumerate(zip(records, steps, strict=True)): + if type(record["step_index"]) is not int or record["step_index"] != index: + problems.append( + f"step {index}: attribution identity differs from the wire" + ) + nodes = _rubric_nodes(record["rubric"]) + if definition(nodes) != definition(baseline): + problems.append(f"step {index}: rubric configuration changed") + reward = json.loads(step.response_json)["data"]["reward"] + root = nodes["root"] + if ( + not root["evaluated"] + or not _finite(reward) + or not math.isclose(root["score"], reward, rel_tol=1e-9, abs_tol=1e-9) + ): + problems.append( + f"step {index}: fresh root score differs from emitted reward" + ) + for node in nodes.values(): + if not node["evaluated"]: + continue + children = [nodes[name] for name in node["children"]] + aggregation = node["aggregation"] + if aggregation == "unknown": + unsupported = True + continue + if aggregation == "leaf": + continue + expected = 1.0 if aggregation == "sequential" else 0.0 + visited = 0 + for child in children: + if not child["evaluated"]: + problems.append( + f"step {index}: required child score was not evaluated" + ) + break + visited += 1 + score = child["score"] + if aggregation == "weighted_sum": + expected += score * node["config"]["weights"][visited - 1] + elif aggregation == "gate": + expected = 0.0 if score < node["config"]["threshold"] else score + else: + expected = score + if score == 0: + break + if aggregation == "sequential" and any( + child["evaluated"] for child in children[visited:] + ): + problems.append( + f"step {index}: sequential evaluated children after its gate" + ) + if not _finite(expected) or not math.isclose( + node["score"], expected, rel_tol=1e-9, abs_tol=1e-9 + ): + problems.append( + f"step {index}: child attribution does not match the parent score" + ) + return ( + problems, + "unsupported custom rubric aggregation semantics" if unsupported else None, + ) diff --git a/src/openenv/validation/runner.py b/src/openenv/validation/runner.py index 180e77d85..9fef32f29 100644 --- a/src/openenv/validation/runner.py +++ b/src/openenv/validation/runner.py @@ -15,6 +15,12 @@ RewardWellFormedGrader, StateContractGrader, ) +from .graders.runtime.discovery import ( + RewardAttributionGrader, + RubricIntrospectableGrader, + TaskDeclarationAccuracyGrader, + ToolDeclarationAccuracyGrader, +) from .graders.runtime.repeatability import ( EpisodeDeterminismGrader, SeedControlGrader, @@ -170,6 +176,10 @@ def _runtime(subject, *, skip_build, provider): manifest.resources.episode_timeout_s, REPLAY_BUDGET_SECONDS ), validation_token=spec.env_vars["OPENENV_VALIDATION_TOKEN"], + collect_tools=True, + task_env_name=manifest.name + if _applicable("runtime.task_declaration_accuracy", manifest) + else None, ) # A health endpoint without a functioning protocol isn't a startup success. if not evidence.exchanges and evidence.failure_reason: @@ -194,15 +204,20 @@ def _runtime(subject, *, skip_build, provider): subject = replace( subject, image_ref=image_ref, running=running, runtime_evidence=evidence ) + runtime_graders = [ + RewardWellFormedGrader(), + ObservationSchemaGrader(), + StateContractGrader(), + SeedControlGrader(), + EpisodeDeterminismGrader(), + TrajectoryRecordGrader(), + ToolDeclarationAccuracyGrader(), + TaskDeclarationAccuracyGrader(), + RubricIntrospectableGrader(), + RewardAttributionGrader(), + ] checks = execute_graders( - [ - RewardWellFormedGrader(), - ObservationSchemaGrader(), - StateContractGrader(), - SeedControlGrader(), - EpisodeDeterminismGrader(), - TrajectoryRecordGrader(), - ], + [grader for grader in runtime_graders if grader.applies_to(manifest)], subject, provider_capabilities=provider.capabilities, prior=[result], @@ -389,6 +404,10 @@ def run_validation( "runtime.seed_control", "runtime.episode_determinism", "runtime.trajectory_record", + "runtime.tool_declaration_accuracy", + "runtime.task_declaration_accuracy", + "runtime.rubric_introspectable", + "runtime.reward_attribution", }: reason = "unmet dependency: runtime.startup" results.append(_outcome(entry.check_id, CheckStatus.SKIP, reason)) @@ -421,6 +440,10 @@ def run_validation( "runtime.seed_control", "runtime.episode_determinism", "runtime.trajectory_record", + "runtime.tool_declaration_accuracy", + "runtime.task_declaration_accuracy", + "runtime.rubric_introspectable", + "runtime.reward_attribution", } else r for r in results diff --git a/src/openenv/validation/runtime/artifacts.py b/src/openenv/validation/runtime/artifacts.py index 4bf43e395..cdfe55826 100644 --- a/src/openenv/validation/runtime/artifacts.py +++ b/src/openenv/validation/runtime/artifacts.py @@ -10,7 +10,7 @@ _SECRET_KEY = re.compile(r"(?i)(password|secret|token|authorization|api[_-]?key)") _TOKEN = re.compile( - r"(?:hf_[A-Za-z0-9]{8,}|(?:sk|ghp|github_pat)[-_][A-Za-z0-9_-]{8,}|(?i:bearer)\s+\S+)" + r"(? RuntimeEvidence: """ Preserve schema and reset/step/state responses without model coercion. @@ -92,6 +95,10 @@ def collect_runtime_evidence( Per-operation deadline, capped by the remaining episode budget. validation_token (`str`, *optional*): Run-scoped telemetry authorization; never retained in evidence. + collect_tools (`bool`, *optional*, defaults to `False`): + Discover tools on the measured WebSocket before reset. + task_env_name (`str`, *optional*): + Sample this environment's task metadata through its HTTP task API. Returns: [`~openenv.validation.runtime.contracts.RuntimeEvidence`]: raw evidence. @@ -103,6 +110,7 @@ def collect_runtime_evidence( trace_bytes = 0 telemetry_json = None telemetry_error = None + tools_json = tools_error = tasks_json = tasks_error = None def remaining() -> float: value = min(request_timeout_s, deadline - time.monotonic()) @@ -169,11 +177,11 @@ def remaining() -> float: complete = False try: - def telemetry_request(operation, data): + def telemetry_request(operation, data, max_bytes=MAX_TRACE_BYTES): request = json.dumps({"type": operation, "data": data}) _bounded_call(connection, lambda: connection.send(request), remaining()) raw = connection.recv(timeout=remaining()) - if not isinstance(raw, str) or len(raw.encode()) > MAX_TRACE_BYTES: + if not isinstance(raw, str) or len(raw.encode()) > max_bytes: raise ValueError("invalid telemetry response") response = json.loads(raw) if not isinstance(response, dict): @@ -215,6 +223,62 @@ def contains_credential(value): return any(contains_credential(child) for child in value) return False + if collect_tools: + phase = "tools/list" + try: + response = telemetry_request( + "mcp", + { + "jsonrpc": "2.0", + "id": 1, + "method": "tools/list", + "params": {}, + }, + MAX_MESSAGE_BYTES, + ) + if contains_credential(response): + raise ValueError("discovery contains validation credentials") + rpc = response.get("data") + if ( + response.get("type") != "mcp" + or not isinstance(rpc, dict) + or rpc.get("jsonrpc") != "2.0" + or type(rpc.get("id")) is not int + or rpc["id"] != 1 + or rpc.get("error") is not None + or not isinstance(rpc.get("result"), dict) + or not isinstance(rpc["result"].get("tools"), list) + ): + raise ValueError("invalid tool discovery response") + discovered = json.dumps( + rpc["result"], allow_nan=False, ensure_ascii=False + ) + if len(discovered.encode()) > MAX_MESSAGE_BYTES: + raise ValueError("discovery exceeds size bound") + tools_json = discovered + except Exception as exc: + tools_error = f"tool discovery failed ({type(exc).__name__})" + + if task_env_name is not None: + phase = "tasks" + try: + discovered, tasks_error = collect_task_evidence( + base_url, + task_env_name, + deadline=deadline, + request_timeout_s=request_timeout_s, + ) + if discovered is not None: + if contains_credential(discovered) or contains_credential( + json.loads(discovered) + ): + raise ValueError( + "discovery contains validation credentials" + ) + tasks_json = discovered + except Exception as exc: + tasks_error = f"task discovery failed ({type(exc).__name__})" + def exchange(operation: str, data: dict | None = None) -> dict: nonlocal phase, trace_bytes phase = operation @@ -316,6 +380,10 @@ def exchange(operation: str, data: dict | None = None) -> dict: observation_schema_json=schema_json, telemetry_json=telemetry_json, telemetry_error=telemetry_error, + tools_json=tools_json, + tools_error=tools_error, + tasks_json=tasks_json, + tasks_error=tasks_error, ) except KeyboardInterrupt: raise RuntimeCollectionInterrupted( @@ -326,6 +394,10 @@ def exchange(operation: str, data: dict | None = None) -> dict: failure_reason=f"{phase} failed (KeyboardInterrupt)", telemetry_json=telemetry_json, telemetry_error=telemetry_error, + tools_json=tools_json, + tools_error=tools_error, + tasks_json=tasks_json, + tasks_error=tasks_error, ) ) from None except Exception as exc: @@ -337,4 +409,8 @@ def exchange(operation: str, data: dict | None = None) -> dict: failure_reason=f"{phase} failed ({type(exc).__name__})", telemetry_json=telemetry_json, telemetry_error=telemetry_error, + tools_json=tools_json, + tools_error=tools_error, + tasks_json=tasks_json, + tasks_error=tasks_error, ) diff --git a/src/openenv/validation/runtime/contracts.py b/src/openenv/validation/runtime/contracts.py index 6d116f998..89d05d642 100644 --- a/src/openenv/validation/runtime/contracts.py +++ b/src/openenv/validation/runtime/contracts.py @@ -269,6 +269,14 @@ class RuntimeEvidence: Subject-emitted session snapshot, independent of the wire transcript. telemetry_error (`str`, *optional*): Bounded telemetry failure without invalidating completed wire evidence. + tools_json (`str`, *optional*): + Raw tools/list result from the measured WebSocket session. + tools_error (`str`, *optional*): + Sanitized tool-discovery failure, distinct from a successful empty list. + tasks_json (`str`, *optional*): + Task split descriptors, counts and bounded item samples. + tasks_error (`str`, *optional*): + Sanitized task-discovery failure, independent of tool discovery. """ exchanges: tuple[WireExchange, ...] = () @@ -279,6 +287,10 @@ class RuntimeEvidence: telemetry_error: str | None = None replays: tuple["ReplayEvidence", ...] = () replay_failure_reason: str | None = None + tools_json: str | None = None + tools_error: str | None = None + tasks_json: str | None = None + tasks_error: str | None = None @dataclass(frozen=True) diff --git a/src/openenv/validation/runtime/discovery.py b/src/openenv/validation/runtime/discovery.py new file mode 100644 index 000000000..138d9dfd8 --- /dev/null +++ b/src/openenv/validation/runtime/discovery.py @@ -0,0 +1,95 @@ +"""Bounded task metadata sampling through the production task API.""" + +import json +import time +from urllib.parse import quote + +import httpx + +MAX_TASK_SPLITS = 64 +MAX_RESPONSE_BYTES = 1024 * 1024 +MAX_DISCOVERY_BYTES = 8 * 1024 * 1024 + + +def collect_task_evidence(base_url, env_name, *, deadline, request_timeout_s=5.0): + """Return raw counts and at most two task specs per split, never a full listing. + + ``deadline`` is the collector's absolute monotonic episode deadline. Errors + are sanitized; a failed/unsupported endpoint never becomes an empty success. + Task metadata routes intentionally use independent environment instances. + """ + total_bytes = 0 + + def remaining(): + budget = min(request_timeout_s, deadline - time.monotonic()) + if budget <= 0: + raise TimeoutError("task discovery deadline exceeded") + return budget + + try: + prefix = base_url.rstrip("/") + "/" + quote(env_name, safe="") + with httpx.Client(trust_env=False, follow_redirects=False) as client: + + def request(method, route, payload=None): + nonlocal total_bytes + with client.stream( + method, + prefix + route, + json=payload, + timeout=remaining(), + headers={"Accept-Encoding": "identity"}, + ) as response: + response.raise_for_status() + if ( + response.headers.get("Content-Encoding", "identity") + .strip() + .lower() + != "identity" + ): + raise ValueError("compressed discovery response") + body = bytearray() + for chunk in response.iter_bytes(): + remaining() + total_bytes += len(chunk) + if ( + len(body) + len(chunk) > MAX_RESPONSE_BYTES + or total_bytes > MAX_DISCOVERY_BYTES + ): + raise ValueError("task discovery exceeds byte budget") + body.extend(chunk) + return json.loads(body) + + splits = request("GET", "/splits") + if not isinstance(splits, list) or len(splits) > MAX_TASK_SPLITS: + raise ValueError("invalid or excessive splits") + counts, previews = {}, {} + for split in splits: + name = split["name"] + if not isinstance(name, str) or not name or name in counts: + raise ValueError("invalid or duplicate split name") + count = request("POST", "/num_tasks", {"split": name})["num_tasks"] + if type(count) is not int or count < 0: + raise ValueError("invalid task count") + counts[name] = count + previews[name] = [ + request("POST", "/task", {"split": name, "index": index})["task"] + for index in range(min(2, count)) + ] + discovered = json.dumps( + {"splits": splits, "counts": counts, "previews": previews}, + allow_nan=False, + ensure_ascii=False, + separators=(",", ":"), + ) + if len(discovered.encode()) > MAX_DISCOVERY_BYTES: + raise ValueError("task evidence exceeds byte budget") + return discovered, None + except ( + httpx.HTTPError, + OSError, + ValueError, + TypeError, + KeyError, + RecursionError, + ) as error: + return None, f"task discovery failed ({type(error).__name__})" diff --git a/tests/fixtures/validation/runtime/served_probe/app.py b/tests/fixtures/validation/runtime/served_probe/app.py index e89b71500..18c45e816 100644 --- a/tests/fixtures/validation/runtime/served_probe/app.py +++ b/tests/fixtures/validation/runtime/served_probe/app.py @@ -6,6 +6,7 @@ from typing import Any import uvicorn +from fastmcp import FastMCP from openenv.core.env_server.http_server import create_app from openenv.core.env_server.interfaces import Environment from openenv.core.env_server.types import Action, Observation, State @@ -22,11 +23,15 @@ class ProbeObservation(Observation): class CounterRubric(Rubric): + def __init__(self, public_config=True): + super().__init__() + self.public_config = public_config + def forward(self, action, observation): return float(observation.counter >= 2) def validation_config(self): - return {"threshold": 2} + return {"threshold": 2} if self.public_config else None class ControlledJudge(Rubric): @@ -53,7 +58,43 @@ def __init__(self, mode="good"): self.ordinal = 0 self._state = State(episode_id="uninitialized", step_count=0) self.counter = 0 - self.rubric = WeightedSum([CounterRubric(), CounterRubric()], [0.5, 0.5]) + public_config = mode != "missing_rubric_config" + self.rubric = WeightedSum( + [CounterRubric(public_config), CounterRubric(public_config)], [0.5, 0.5] + ) + self.mcp_server = FastMCP("validation_probe") + if mode != "empty_tools": + self.mcp_server.tool(self.increment) + if mode != "missing_tool": + self.mcp_server.tool(self.read_counter) + if mode == "extra_tool": + self.mcp_server.tool(name="unexpected")(self.read_counter) + if mode == "tool_discovery_error": + self.mcp_server = None + + def increment(self, amount: int = 1) -> dict: + """Advance the probe using its ordinary step implementation.""" + return self.step(ProbeAction(increment=amount)).model_dump(mode="json") + + def read_counter(self) -> int: + """Read the session's current counter without changing it.""" + return self.counter + + def list_splits(self): + return ["train", "test"] + + def num_tasks(self, split): + counts = {"train": 4, "test": 2} + return counts[split] + int(self.mode == "bad_task_count" and split == "train") + + def get_task(self, split, index): + if not 0 <= index < {"train": 4, "test": 2}[split]: + raise IndexError(index) + return {"id": f"{split}-{index}", "index": index, "split": split} + + def list_tasks(self, split): + # Deliberately bounded: listing length is never the authoritative count. + return [self.get_task(split, 0)] def reset(self, seed=None, episode_id=None, **kwargs): self.counter = 0 @@ -143,6 +184,8 @@ async def fault_send(message: dict[str, Any]): data["data"]["trajectory"]["records"][0]["response"]["data"][ "observation" ]["counter"] = 999 + elif self.mode == "bad_attribution": + data["data"]["attribution"][0]["rubric"][1]["score"] = 0.25 message = {**message, "text": json.dumps(data)} await send(message) @@ -167,6 +210,13 @@ def make_app(mode="good"): "trace_mismatch", "judged_stable", "judged_noisy", + "missing_tool", + "extra_tool", + "empty_tools", + "tool_discovery_error", + "bad_task_count", + "missing_rubric_config", + "bad_attribution", }: raise ValueError(f"Unknown fixture mode: {mode}") environment = IgnoredSeedEnvironment if mode == "ignored_seed" else ProbeEnvironment diff --git a/tests/fixtures/validation/runtime/served_probe/openenv.yaml b/tests/fixtures/validation/runtime/served_probe/openenv.yaml index 8dcb0f534..58314397a 100644 --- a/tests/fixtures/validation/runtime/served_probe/openenv.yaml +++ b/tests/fixtures/validation/runtime/served_probe/openenv.yaml @@ -18,7 +18,10 @@ validation: capabilities: verifier: kind: reward_channel - declared_tools: [] + declared_tools: [increment, read_counter] + rubric_tree: true + task_api: true + declared_task_count: {train: 4, test: 2} types: tags: [demo] execution: diff --git a/tests/test_validation/integration/test_discovery_protocol.py b/tests/test_validation/integration/test_discovery_protocol.py new file mode 100644 index 000000000..a5b456b37 --- /dev/null +++ b/tests/test_validation/integration/test_discovery_protocol.py @@ -0,0 +1,82 @@ +"""Real registered MCP tools and bounded task metadata use existing APIs.""" + +from starlette.testclient import TestClient +from test_served_probe import fixture + + +def test_real_tools_share_the_replayed_instance_and_tasks_allow_bounded_listing( + monkeypatch, +): + token = "fixture-tool-test-token-" + "0" * 32 + monkeypatch.setenv("OPENENV_VALIDATION_TOKEN", token) + with TestClient(fixture.make_app()) as client: + counts = {} + for split in ("train", "test"): + counts[split] = client.post( + "/validation_probe/num_tasks", json={"split": split} + ).json()["num_tasks"] + preview = client.post( + "/validation_probe/tasks", json={"split": split} + ).json()["tasks"] + assert len(preview) == 1 < counts[split] + assert counts == {"train": 4, "test": 2} + with client.websocket_connect("/ws") as ws: + ws.send_json( + { + "type": "validation_open", + "data": {"schema_version": 1, "token": token}, + } + ) + capability = ws.receive_json()["data"]["capability"] + + def rpc(method, params): + ws.send_json( + { + "type": "mcp", + "data": { + "jsonrpc": "2.0", + "id": 1, + "method": method, + "params": params, + }, + } + ) + return ws.receive_json()["data"] + + tools = rpc("tools/list", {})["result"]["tools"] + assert {tool["name"] for tool in tools} == {"increment", "read_counter"} + assert ( + rpc("tools/call", {"name": "read_counter", "arguments": {}})["result"] + == 0 + ) + result = rpc( + "tools/call", {"name": "increment", "arguments": {"amount": 2}} + ) + assert result["result"]["counter"] == 2 + assert ( + rpc("tools/call", {"name": "read_counter", "arguments": {}})["result"] + == 2 + ) + ws.send_json( + {"type": "reset", "data": {"seed": 1, "episode_id": "after-tools"}} + ) + assert ws.receive_json()["data"]["observation"]["counter"] == 0 + ws.send_json({"type": "state"}) + assert ws.receive_json()["data"]["step_count"] == 0 + ws.send_json( + { + "type": "validation_read", + "data": {"schema_version": 1, "capability": capability}, + } + ) + telemetry = ws.receive_json()["data"] + assert telemetry["seed"] == { + "requested": True, + "value": 1, + "accepted": True, + } + assert [row["operation"] for row in telemetry["trajectory"]["records"]] == [ + "reset", + "state", + ] + assert telemetry["trajectory"]["complete"] is True diff --git a/tests/test_validation/integration/test_runtime_cli.py b/tests/test_validation/integration/test_runtime_cli.py index bae6f2f50..ec89f19dd 100644 --- a/tests/test_validation/integration/test_runtime_cli.py +++ b/tests/test_validation/integration/test_runtime_cli.py @@ -26,20 +26,18 @@ "runtime.seed_control", "runtime.episode_determinism", "runtime.trajectory_record", + "runtime.tool_declaration_accuracy", + "runtime.task_declaration_accuracy", + "runtime.rubric_introspectable", + "runtime.reward_attribution", } PENDING = { - "runtime.tool_declaration_accuracy", "runtime.network_policy", "runtime.host_containment", "runtime.resource_bounds", "runtime.episode_isolation", "runtime.oracle_containment", } -NOT_APPLICABLE = { - "runtime.rubric_introspectable", - "runtime.reward_attribution", - "runtime.task_declaration_accuracy", -} def _copy_asset(source, destination): @@ -272,20 +270,46 @@ def _invoke_cli( assert not marker.exists(), "--skip-build invoked Docker" checks = {row["check_id"]: row for row in report["results"]} assert len(checks) == len(report["results"]), "Report contains duplicate check IDs" - assert { - key for key in checks if key.startswith("runtime.") - } == IMPLEMENTED | PENDING - assert not NOT_APPLICABLE & checks.keys() + applicable = IMPLEMENTED | PENDING + capabilities = report["manifest"]["capabilities"] + if not capabilities["rubric_tree"]: + applicable -= {"runtime.rubric_introspectable", "runtime.reward_attribution"} + if not capabilities["task_api"] and not capabilities["declared_task_count"]: + applicable -= {"runtime.task_declaration_accuracy"} + assert {key for key in checks if key.startswith("runtime.")} == applicable assert all(checks[key]["status"] == "skip" for key in PENDING) assert all(checks[key]["evidence"] for key in PENDING) - assert ( - set(artifacts["coverage.json"]["requested_runtime_checks"]) - == IMPLEMENTED | PENDING - ) + assert set(artifacts["coverage.json"]["requested_runtime_checks"]) == applicable assert artifacts["cleanup.json"]["completed"] is True return result, report, checks, artifacts +def _assert_discovery(discovery, mode): + assert discovery["omitted_fields"] == [] + if mode == "tool_discovery_error": + assert discovery["tools_available"] is False + assert discovery["tools_error"] + assert discovery["tools"] is None + else: + assert discovery["tools_available"] is True + assert discovery["tools_error"] is None + expected = {"increment", "read_counter"} + if mode == "empty_tools": + expected = set() + elif mode == "missing_tool": + expected.remove("read_counter") + elif mode == "extra_tool": + expected.add("unexpected") + assert {tool["name"] for tool in discovery["tools"]["tools"]} == expected + assert discovery["tasks_available"] is True + assert discovery["tasks_error"] is None + assert discovery["tasks"]["counts"] == { + "train": 5 if mode == "bad_task_count" else 4, + "test": 2, + } + assert all(len(preview) == 2 for preview in discovery["tasks"]["previews"].values()) + + def _assert_fresh_replays(artifacts, identical_samples=3): replay = artifacts["replays.json"] assert replay["failure_reason"] is None @@ -315,11 +339,24 @@ def _assert_fresh_replays(artifacts, identical_samples=3): ("nondeterministic", "runtime.episode_determinism"), ("missing_record", "runtime.trajectory_record"), ("trace_mismatch", "runtime.trajectory_record"), + ("missing_tool", "runtime.tool_declaration_accuracy"), + ("extra_tool", "runtime.tool_declaration_accuracy"), + ("tool_discovery_error", "runtime.tool_declaration_accuracy"), + ("bad_task_count", "runtime.task_declaration_accuracy"), + ("missing_rubric_config", "runtime.rubric_introspectable"), + ("bad_attribution", "runtime.reward_attribution"), + ("empty_tools", None), ], ) def test_cli_runtime_contract_findings(cli_context, tmp_path, mode, failed_check): with (cli_context / "Dockerfile").open("a") as stream: stream.write(f"\nENV VALIDATION_FAULT={mode}\n") + if mode == "empty_tools": + manifest = cli_context / "openenv.yaml" + data = yaml.safe_load(manifest.read_text()) + data["validation"]["capabilities"]["declared_tools"] = [] + manifest.unlink() + manifest.write_text(yaml.safe_dump(data, sort_keys=False)) result, report, checks, artifacts = _invoke_cli(cli_context, tmp_path, mode) assert report["levels_run"] == [1, 2] assert checks["static.manifest"]["status"] == "pass" @@ -329,6 +366,7 @@ def test_cli_runtime_contract_findings(cli_context, tmp_path, mode, failed_check assert artifacts["run-manifest.json"]["provider"]["container_id"] assert artifacts["run-manifest.json"]["source_digest"] == report["source_digest"] _assert_fresh_replays(artifacts) + _assert_discovery(artifacts["discovery.json"], mode) trace = artifacts["collector-trace.json"] assert [row["operation"] for row in trace] == [ "reset", @@ -339,7 +377,11 @@ def test_cli_runtime_contract_findings(cli_context, tmp_path, mode, failed_check "state", ] if failed_check: - warning_only = failed_check == "runtime.trajectory_record" + warning_only = failed_check in { + "runtime.trajectory_record", + "runtime.rubric_introspectable", + "runtime.reward_attribution", + } assert result.returncode == (0 if warning_only else 1) assert report["verdict"] == ("warn" if warning_only else "fail") assert checks[failed_check]["status"] == "fail" @@ -475,6 +517,7 @@ def test_cli_controlled_judge_requires_twenty_samples( assert determinism["status"] == expected assert determinism["measured"]["completed_replays"] == 20 _assert_fresh_replays(artifacts, identical_samples=20) + _assert_discovery(artifacts["discovery.json"], mode) assert determinism["measured"]["variance_units"] == "reward_squared" variances = determinism["measured"]["reward_population_variance"] assert len(variances) == 2 diff --git a/tests/test_validation/integration/test_runtime_process.py b/tests/test_validation/integration/test_runtime_process.py index 8afe549ea..0117b4495 100644 --- a/tests/test_validation/integration/test_runtime_process.py +++ b/tests/test_validation/integration/test_runtime_process.py @@ -19,9 +19,11 @@ import httpx import pytest +import yaml from openenv.validation.providers import StartupError, UnsupportedCapability from openenv.validation.runner import run_validation from openenv.validation.types import CheckStatus, Level, ProviderCapability +from test_runtime_cli import _assert_discovery FIXTURE = Path(__file__).parents[2] / "fixtures/validation/runtime/served_probe" SERVER = """ @@ -167,6 +169,13 @@ def start(self, spec): ("ignored_seed", "runtime.seed_control"), ("missing_record", "runtime.trajectory_record"), ("trace_mismatch", "runtime.trajectory_record"), + ("missing_tool", "runtime.tool_declaration_accuracy"), + ("extra_tool", "runtime.tool_declaration_accuracy"), + ("tool_discovery_error", "runtime.tool_declaration_accuracy"), + ("bad_task_count", "runtime.task_declaration_accuracy"), + ("missing_rubric_config", "runtime.rubric_introspectable"), + ("bad_attribution", "runtime.reward_attribution"), + ("empty_tools", None), ], ) def test_installed_server_collector_and_graders_over_loopback( @@ -177,10 +186,18 @@ def test_installed_server_collector_and_graders_over_loopback( / "process" / mode ) + target = FIXTURE + if mode == "empty_tools": + target = tmp_path / "empty-tools" + shutil.copytree(FIXTURE, target, ignore=shutil.ignore_patterns("__pycache__")) + manifest_path = target / "openenv.yaml" + manifest = yaml.safe_load(manifest_path.read_text()) + manifest["validation"]["capabilities"]["declared_tools"] = [] + manifest_path.write_text(yaml.safe_dump(manifest, sort_keys=False)) provider = ProcessProvider(artifacts / "subject", mode) try: report = run_validation( - FIXTURE, + target, max_level=Level.RUNTIME, provider=provider, artifacts_dir=artifacts / "report", @@ -200,6 +217,10 @@ def test_installed_server_collector_and_graders_over_loopback( "state_contract", "seed_control", "trajectory_record", + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", ): assert checks[f"runtime.{check}"].status is CheckStatus.PASS if mode != "startup_failure": @@ -208,6 +229,9 @@ def test_installed_server_collector_and_graders_over_loopback( manifest = json.loads((artifacts / "report/run-manifest.json").read_text()) assert manifest["provider"]["isolation"] == "process" assert manifest["provider"]["container_build_exercised"] is False + _assert_discovery( + json.loads((artifacts / "report/discovery.json").read_text()), mode + ) telemetry_path = artifacts / "report/session-telemetry.json" telemetry = json.loads(telemetry_path.read_text()) assert telemetry["seed"]["accepted"] is (mode != "ignored_seed") diff --git a/tests/test_validation/test_discovery_collection.py b/tests/test_validation/test_discovery_collection.py new file mode 100644 index 000000000..926974453 --- /dev/null +++ b/tests/test_validation/test_discovery_collection.py @@ -0,0 +1,264 @@ +"""Opt-in discovery uses the measured socket and never retains its credentials.""" + +import json +from dataclasses import asdict + +import httpx +import pytest +from openenv.validation.runtime import collector +from openenv.validation.runtime.contracts import RuntimePlan +from test_runtime_transport import EpisodeConnection + +TOKEN = "test-validation-token-" * 2 +CAPABILITY = "test-session-capability-" * 2 + + +class DiscoveryConnection(EpisodeConnection): + def __init__(self, response=None): + super().__init__() + self.requests = [] + self.response = response + + def send(self, payload): + super().send(payload) + self.requests.append(json.loads(payload)) + + def recv(self, timeout): + if self.operation == "validation_open": + return json.dumps( + { + "type": "validation_open", + "data": {"schema_version": 1, "capability": CAPABILITY}, + } + ) + if self.operation == "validation_read": + return json.dumps( + {"type": "validation", "data": {"schema_version": 1, "rubric": []}} + ) + if self.operation == "mcp": + if isinstance(self.response, Exception): + raise self.response + if isinstance(self.response, str): + return self.response + return json.dumps(self.response or tool_response()) + return super().recv(timeout) + + +def tool_response(tools=None): + return { + "type": "mcp", + "data": { + "jsonrpc": "2.0", + "id": 1, + "result": {"tools": [] if tools is None else tools}, + "error": None, + }, + } + + +def collect(monkeypatch, connection, *, task_fault=None, **kwargs): + original_client = httpx.Client + urls = [] + + def respond(request): + urls.append(request.url.path) + if request.url.path == "/schema": + return httpx.Response(200, json={"observation": {"type": "object"}}) + if task_fault == "unsupported": + return httpx.Response(501, text="private server details") + if request.url.path.endswith("/splits"): + return httpx.Response(200, json=[{"name": "train"}]) + if request.url.path.endswith("/num_tasks"): + return httpx.Response(200, json={"num_tasks": 100}) + assert request.url.path.endswith("/task") + return httpx.Response(200, json={"task": {"id": task_fault or "public-task"}}) + + monkeypatch.setattr( + collector.httpx, + "Client", + lambda **options: original_client( + transport=httpx.MockTransport(respond), **options + ), + ) + monkeypatch.setattr(collector, "connect", lambda *args, **options: connection) + plan = RuntimePlan.model_validate( + { + "plan_schema_version": "1", + "reset": {"episode_id": "bounded", "seed": 7}, + "actions": [{"increment": 1}], + } + ) + evidence = collector.collect_runtime_evidence( + "http://127.0.0.1:8000", + plan, + episode_timeout_s=2, + validation_token=TOKEN, + **kwargs, + ) + return evidence, urls + + +def test_default_collection_does_not_discover_tools_or_tasks(monkeypatch): + connection = DiscoveryConnection() + evidence, urls = collect(monkeypatch, connection) + assert "mcp" not in [request["type"] for request in connection.requests] + assert urls == ["/schema"] + assert ( + evidence.tools_json + is evidence.tools_error + is evidence.tasks_json + is evidence.tasks_error + is None + ) + + +def test_discovery_stays_on_measured_socket_outside_trajectory(monkeypatch): + connection = DiscoveryConnection( + tool_response([{"name": "echo", "inputSchema": {"type": "object"}}]) + ) + evidence, urls = collect( + monkeypatch, connection, collect_tools=True, task_env_name="probe" + ) + assert [request["type"] for request in connection.requests] == [ + "validation_open", + "mcp", + "reset", + "state", + "step", + "state", + "validation_read", + "close", + ] + assert connection.requests[1] == { + "type": "mcp", + "data": {"jsonrpc": "2.0", "id": 1, "method": "tools/list", "params": {}}, + } + assert [row.operation for row in evidence.exchanges] == [ + "reset", + "state", + "step", + "state", + ] + assert ( + evidence.failure_reason is evidence.tools_error is evidence.tasks_error is None + ) + assert json.loads(evidence.tools_json) == { + "tools": [{"name": "echo", "inputSchema": {"type": "object"}}] + } + assert json.loads(evidence.tasks_json)["counts"] == {"train": 100} + assert urls == [ + "/schema", + "/probe/splits", + "/probe/num_tasks", + "/probe/task", + "/probe/task", + ] + assert TOKEN not in json.dumps(asdict(evidence)) and CAPABILITY not in json.dumps( + asdict(evidence) + ) + + +@pytest.mark.parametrize( + "fault", + [ + "unsupported", + "wrong_id", + "boolean_id", + "missing_tools", + "bad_json", + "oversize", + "transport", + ], +) +def test_failed_tool_discovery_is_not_empty_success_and_preserves_episode( + monkeypatch, fault +): + response = tool_response() + if fault == "unsupported": + response["data"].update( + result=None, error={"code": -32601, "message": "private error"} + ) + elif fault == "wrong_id": + response["data"]["id"] = 2 + elif fault == "boolean_id": + response["data"]["id"] = True + elif fault == "missing_tools": + response["data"]["result"] = {} + elif fault == "bad_json": + response = "{private invalid json" + elif fault == "oversize": + response = "x" * (collector.MAX_MESSAGE_BYTES + 1) + else: + response = ValueError("private transport error") + evidence, _ = collect( + monkeypatch, DiscoveryConnection(response), collect_tools=True + ) + assert evidence.tools_json is None and evidence.tools_error + assert "private" not in evidence.tools_error + assert evidence.failure_reason is None + assert len(evidence.exchanges) == 4 + + +@pytest.mark.parametrize("secret", [TOKEN, CAPABILITY]) +@pytest.mark.parametrize("escaped", [False, True]) +def test_discovery_credential_echo_is_discarded_even_with_unicode_escapes( + monkeypatch, secret, escaped +): + response = json.dumps(tool_response([{"name": "echo", "description": secret}])) + if escaped: + response = response.replace( + secret, "".join(f"\\u{ord(char):04x}" for char in secret) + ) + evidence, _ = collect( + monkeypatch, + DiscoveryConnection(response), + collect_tools=True, + task_env_name="probe", + task_fault=secret, + ) + assert evidence.tools_json is evidence.tasks_json is None + assert evidence.tools_error and evidence.tasks_error + assert evidence.failure_reason is None + assert secret not in json.dumps(asdict(evidence)) + + +def test_task_failure_does_not_poison_successful_tools_or_episode(monkeypatch): + evidence, _ = collect( + monkeypatch, + DiscoveryConnection(), + collect_tools=True, + task_env_name="probe", + task_fault="unsupported", + ) + assert json.loads(evidence.tools_json) == {"tools": []} + assert evidence.tools_error is None + assert evidence.tasks_json is None and evidence.tasks_error + assert evidence.failure_reason is None and len(evidence.exchanges) == 4 + + +def test_tool_pagination_remains_visible_to_grader(monkeypatch): + response = tool_response([{"name": "echo"}]) + response["data"]["result"]["nextCursor"] = "next-page" + evidence, _ = collect( + monkeypatch, DiscoveryConnection(response), collect_tools=True + ) + assert json.loads(evidence.tools_json)["nextCursor"] == "next-page" + assert evidence.tools_error is evidence.failure_reason is None + assert len(evidence.exchanges) == 4 + + +def test_socket_failure_keeps_discovery_error_and_fails_replay(monkeypatch): + connection = DiscoveryConnection(ConnectionError("private socket error")) + send = connection.send + + def broken_send(payload): + if json.loads(payload)["type"] == "reset": + raise ConnectionError("closed socket") + send(payload) + + connection.send = broken_send + evidence, _ = collect(monkeypatch, connection, collect_tools=True) + assert evidence.tools_error == "tool discovery failed (ConnectionError)" + assert evidence.failure_phase == "reset" + assert evidence.failure_reason == "reset failed (ConnectionError)" + assert evidence.exchanges == () diff --git a/tests/test_validation/test_runtime_artifacts.py b/tests/test_validation/test_runtime_artifacts.py index ddcb124ba..6a037e740 100644 --- a/tests/test_validation/test_runtime_artifacts.py +++ b/tests/test_validation/test_runtime_artifacts.py @@ -11,14 +11,14 @@ StateContractGrader, ) from openenv.validation.manifest import NormalizedManifest -from openenv.validation.report import ValidationReportV2 +from openenv.validation.report import CheckResult, ValidationReportV2 from openenv.validation.runtime.artifacts import write_runtime_bundle from openenv.validation.runtime.contracts import ( ReplayEvidence, RuntimeEvidence, WireExchange, ) -from openenv.validation.types import Lane, Level, SignatureKind, Verdict +from openenv.validation.types import CheckStatus, Lane, Level, SignatureKind, Verdict from support.runtime import evidence, exchange @@ -365,3 +365,92 @@ def test_rewriting_bundle_removes_stale_optional_evidence(tmp_path): assert not (tmp_path / "replays.json").exists() assert not (tmp_path / "session-telemetry.json").exists() assert "replays.json" not in (tmp_path / "SHA256SUMS").read_text() + + +def test_discovery_artifact_retains_true_counts_bounded_previews_and_digest(tmp_path): + tools = {"tools": []} + tasks = { + "splits": [{"name": "train"}], + "counts": {"train": 100}, + "previews": {"train": ["task-0", "task-1"]}, + } + original = replace( + measured(), tools_json=json.dumps(tools), tasks_json=json.dumps(tasks) + ) + write_runtime_bundle(tmp_path, report(), evidence=original) + path = tmp_path / "discovery.json" + artifact = json.loads(path.read_text()) + assert artifact["tools"] == tools and artifact["tasks"] == tasks + assert artifact["tools_available"] is artifact["tasks_available"] is True + assert artifact["tools_error"] is artifact["tasks_error"] is None + assert artifact["redacted"] is False + assert ( + f"{hashlib.sha256(path.read_bytes()).hexdigest()} discovery.json" + in (tmp_path / "SHA256SUMS").read_text() + ) + + +def test_discovery_artifact_redacts_tool_task_and_error_credentials(tmp_path): + original = replace( + measured(), + tools_json='{"tools":[{"api_key":"private-tool"}]}', + tasks_json='{"previews":{"train":[{"password":"private-task"}]}}', + tools_error="Authorization: Bearer private-error", + ) + write_runtime_bundle(tmp_path, report(), evidence=original) + text = (tmp_path / "discovery.json").read_text() + assert "private-" not in text + assert json.loads(text)["redacted"] is True + + +def test_token_redaction_preserves_task_check_identifiers(tmp_path): + validation_report = report() + validation_report.results = [ + CheckResult( + check_id="runtime.task_declaration_accuracy", + status=CheckStatus.PASS, + duration_s=0, + ) + ] + original = replace( + measured(), + tools_json='{"tools":[{"name":"task_declaration_accuracy","description":"sk_thisisafaketoken123"}]}', + ) + write_runtime_bundle(tmp_path, validation_report, evidence=original) + assert json.loads( + (tmp_path / "report.json").read_text() + ) == validation_report.model_dump(mode="json") + artifact = json.loads((tmp_path / "discovery.json").read_text()) + assert artifact["tools"]["tools"][0] == { + "name": "task_declaration_accuracy", + "description": "[REDACTED]", + } + + +def test_malformed_discovery_omissions_preserve_collection_failures(tmp_path): + original = replace( + measured(), + tools_json="invalid-private-tool", + tasks_json="invalid-private-task", + tools_error="tool discovery failed (ValueError)", + tasks_error="task discovery failed (ValueError)", + ) + write_runtime_bundle(tmp_path, report(), evidence=original) + text = (tmp_path / "discovery.json").read_text() + assert "invalid-private" not in text + artifact = json.loads(text) + assert artifact["redacted"] is True + assert artifact["omitted_fields"] == ["tools", "tasks"] + assert artifact["tools_error"] == original.tools_error + assert artifact["tasks_error"] == original.tasks_error + + +def test_discovery_distinguishes_missing_from_null_and_removes_stale_artifact(tmp_path): + original = replace(measured(), tasks_json="null") + write_runtime_bundle(tmp_path, report(), evidence=original) + artifact = json.loads((tmp_path / "discovery.json").read_text()) + assert artifact["tools"] is artifact["tasks"] is None + assert artifact["tools_available"] is False and artifact["tasks_available"] is True + write_runtime_bundle(tmp_path, report()) + assert not (tmp_path / "discovery.json").exists() + assert "discovery.json" not in (tmp_path / "SHA256SUMS").read_text() diff --git a/tests/test_validation/test_runtime_discovery.py b/tests/test_validation/test_runtime_discovery.py new file mode 100644 index 000000000..6e4ed1d58 --- /dev/null +++ b/tests/test_validation/test_runtime_discovery.py @@ -0,0 +1,385 @@ +"""Pure declaration/attribution checks; process tests cover their wire inputs.""" + +import copy +import json +from dataclasses import replace +from types import SimpleNamespace + +import pytest +from conftest import load_fixture_manifest +from openenv.validation.graders import Subject +from openenv.validation.graders.runtime.discovery import ( + RewardAttributionGrader, + RubricIntrospectableGrader, + TaskDeclarationAccuracyGrader, + ToolDeclarationAccuracyGrader, +) +from openenv.validation.manifest import NormalizedManifest +from openenv.validation.types import CheckStatus +from support.runtime import exchange + + +def node( + name="root", + score=0.75, + *, + children=None, + aggregation="leaf", + config=None, + evaluated=True, +): + return { + "name": name, + "class_name": "fixture.Rubric", + "children": children or [], + "aggregation": aggregation, + "config": config or {}, + "config_available": True, + "score": score if evaluated else None, + "evaluated": evaluated, + } + + +def telemetry(nodes=None): + nodes = nodes or [ + node( + children=["root.a", "root.b"], + aggregation="weighted_sum", + config={"weights": [0.25, 0.75]}, + ), + node("root.a", 0), + node("root.b", 1), + ] + return { + "schema_version": 1, + "rubric": copy.deepcopy(nodes), + "attribution": [{"step_index": 0, "rubric": copy.deepcopy(nodes)}], + } + + +def subject(tmp_path, *, snapshot=None, reward=0.75, **fields): + manifest = load_fixture_manifest("served_min_pass") + manifest["capabilities"].update( + rubric_tree=True, + task_api=True, + declared_tools=["echo"], + declared_task_count={"train": 100}, + ) + evidence = SimpleNamespace( + tools_json=json.dumps({"tools": [{"name": "echo"}]}), + tools_error=None, + tasks_json=json.dumps( + { + "splits": [{"name": "train"}], + "counts": {"train": 100}, + "previews": {"train": [{"id": 0}, {"id": 1}]}, + } + ), + tasks_error=None, + telemetry_json=json.dumps(telemetry() if snapshot is None else snapshot), + telemetry_error=None, + failure_reason=None, + exchanges=( + exchange( + "step", + {"type": "step", "data": {}}, + { + "type": "observation", + "data": {"reward": reward, "done": False, "observation": {}}, + }, + ), + ), + ) + for name, value in fields.items(): + setattr(evidence, name, value) + return Subject( + tmp_path, + NormalizedManifest.model_validate(manifest), + None, + None, + tmp_path, + evidence, + ) + + +@pytest.mark.parametrize( + "grader", + [ + ToolDeclarationAccuracyGrader, + TaskDeclarationAccuracyGrader, + RubricIntrospectableGrader, + RewardAttributionGrader, + ], +) +def test_matching_discovery_and_fresh_attribution_pass(tmp_path, grader): + assert grader().run(subject(tmp_path)).status is CheckStatus.PASS + + +def test_genuinely_empty_tools_pass_but_unsupported_discovery_does_not(tmp_path): + measured = subject(tmp_path, tools_json='{"tools": []}') + measured.manifest.capabilities.declared_tools = [] + assert ToolDeclarationAccuracyGrader().run(measured).status is CheckStatus.PASS + measured.runtime_evidence.tools_error = "unsupported" + assert ToolDeclarationAccuracyGrader().run(measured).status is CheckStatus.FAIL + + +@pytest.mark.parametrize("tools", [[], [{"name": "echo"}]]) +def test_matching_or_partial_first_tool_page_is_explicitly_incomplete(tmp_path, tools): + result = ToolDeclarationAccuracyGrader().run( + subject(tmp_path, tools_json=json.dumps({"tools": tools, "nextCursor": "next"})) + ) + assert result.status is CheckStatus.SKIP + assert "pagination" in " ".join(result.evidence) + + +@pytest.mark.parametrize("cursor", [False, 1, [], {}]) +def test_malformed_tool_cursor_fails(tmp_path, cursor): + result = ToolDeclarationAccuracyGrader().run( + subject(tmp_path, tools_json=json.dumps({"tools": [], "nextCursor": cursor})) + ) + assert result.status is CheckStatus.FAIL + + +def test_undeclared_tool_on_first_page_still_fails(tmp_path): + result = ToolDeclarationAccuracyGrader().run( + subject( + tmp_path, + tools_json=json.dumps({"tools": [{"name": "extra"}], "nextCursor": "next"}), + ) + ) + assert result.status is CheckStatus.FAIL + assert "undeclared" in " ".join(result.evidence) + + +def test_task_preview_preserves_any_json_task_spec(tmp_path): + measured = subject(tmp_path) + tasks = json.loads(measured.runtime_evidence.tasks_json) + tasks["previews"]["train"] = ["task-0", ["task-1", None]] + measured.runtime_evidence.tasks_json = json.dumps(tasks) + assert TaskDeclarationAccuracyGrader().run(measured).status is CheckStatus.PASS + + +def test_rubric_checks_do_not_apply_plan_size_limits_to_trajectory(tmp_path): + payload = telemetry() + payload["trajectory"] = { + "records": [{"observation": [0] * 200} for _ in range(100)] + } + measured = subject(tmp_path, snapshot=payload) + for grader in (RubricIntrospectableGrader, RewardAttributionGrader): + assert grader().run(measured).status is CheckStatus.PASS + + +@pytest.mark.parametrize( + "tools,reason", + [ + ([], "missing"), + ([{"name": "echo"}, {"name": "extra"}], "undeclared"), + ([{"name": "echo"}, {"name": "echo"}], "duplicate"), + ([{"name": True}], "malformed"), + ], +) +def test_tool_mismatches_have_distinct_findings(tmp_path, tools, reason): + result = ToolDeclarationAccuracyGrader().run( + subject(tmp_path, tools_json=json.dumps({"tools": tools})) + ) + assert result.status is CheckStatus.FAIL + assert reason in " ".join(result.evidence) + + +@pytest.mark.parametrize("count", [True, "100", 100.0, -1, 99]) +def test_invalid_or_wrong_task_count_cannot_be_coerced(tmp_path, count): + measured = subject(tmp_path) + tasks = json.loads(measured.runtime_evidence.tasks_json) + tasks["counts"]["train"] = count + measured.runtime_evidence.tasks_json = json.dumps(tasks) + assert TaskDeclarationAccuracyGrader().run(measured).status is CheckStatus.FAIL + + +@pytest.mark.parametrize( + "change", ["extra_split", "missing_count", "oversized_preview", "empty_preview"] +) +def test_inconsistent_task_discovery_fails(tmp_path, change): + measured = subject(tmp_path) + tasks = json.loads(measured.runtime_evidence.tasks_json) + if change == "extra_split": + tasks["splits"].append({"name": "test"}) + tasks["counts"]["test"] = 0 + tasks["previews"]["test"] = [] + elif change == "missing_count": + tasks["counts"] = {} + elif change == "oversized_preview": + tasks["counts"]["train"] = 1 + measured.manifest.capabilities.declared_task_count["train"] = 1 + else: + tasks["previews"]["train"] = [] + measured.runtime_evidence.tasks_json = json.dumps(tasks) + assert TaskDeclarationAccuracyGrader().run(measured).status is CheckStatus.FAIL + + +@pytest.mark.parametrize( + "grader,field,error", + [ + (ToolDeclarationAccuracyGrader, "tools_json", "tools_error"), + (TaskDeclarationAccuracyGrader, "tasks_json", "tasks_error"), + (RubricIntrospectableGrader, "telemetry_json", "telemetry_error"), + (RewardAttributionGrader, "telemetry_json", "telemetry_error"), + ], +) +def test_missing_prerequisite_and_collection_failure_are_distinct( + tmp_path, grader, field, error +): + assert grader().run(subject(tmp_path, **{field: None})).status is CheckStatus.SKIP + result = grader().run( + subject(tmp_path, **{field: None, error: "private-exception-text"}) + ) + assert result.status is CheckStatus.FAIL + assert "private-exception-text" not in result.model_dump_json() + assert ( + grader().run(replace(subject(tmp_path), runtime_evidence=None)).status + is CheckStatus.SKIP + ) + + +@pytest.mark.parametrize( + "change", + [ + "no_config", + "no_root", + "broken_edge", + "duplicate_node", + "nonfinite_score", + "stale_score", + "bad_version", + "rubric_error", + "empty_child_segment", + ], +) +def test_missing_or_malformed_rubric_introspection_fails(tmp_path, change): + payload = telemetry() + if change == "no_config": + payload["rubric"][0]["config_available"] = False + elif change == "no_root": + payload["rubric"] = [] + elif change == "broken_edge": + payload["rubric"][0]["children"] = ["root.missing", "root.b"] + elif change == "duplicate_node": + payload["rubric"].append(payload["rubric"][0]) + elif change == "nonfinite_score": + payload["rubric"][0]["score"] = float("nan") + elif change == "stale_score": + payload["rubric"][0]["evaluated"] = False + elif change == "bad_version": + payload["schema_version"] = True + elif change == "empty_child_segment": + payload["rubric"][0]["children"][0] = "root." + payload["rubric"][1]["name"] = "root." + else: + payload["rubric_error"] = "private-error" + result = RubricIntrospectableGrader().run(subject(tmp_path, snapshot=payload)) + assert result.status is CheckStatus.FAIL + assert "private-error" not in result.model_dump_json() + + +@pytest.mark.parametrize( + "change", + [ + "root_reward", + "child_total", + "missing_child", + "stale_root", + "missing_step", + "wrong_step", + "config_changed", + ], +) +def test_incorrect_or_stale_attribution_fails(tmp_path, change): + payload = telemetry() + record = payload["attribution"][0] + if change == "root_reward": + record["rubric"][0]["score"] = 0.5 + elif change == "child_total": + record["rubric"][1]["score"] = 1 + elif change == "missing_child": + record["rubric"][1].update(evaluated=False, score=None) + elif change == "stale_root": + record["rubric"][0].update(evaluated=False, score=None) + elif change == "missing_step": + payload["attribution"] = [] + elif change == "wrong_step": + record["step_index"] = True + else: + record["rubric"][0]["config"]["weights"] = [0.5, 0.5] + assert ( + RewardAttributionGrader().run(subject(tmp_path, snapshot=payload)).status + is CheckStatus.FAIL + ) + + +@pytest.mark.parametrize("aggregation", ["leaf", "gate", "sequential", "weighted_sum"]) +def test_stock_aggregation_rules_and_sequential_short_circuit(tmp_path, aggregation): + if aggregation == "leaf": + nodes, reward = [node()], 0.75 + elif aggregation == "gate": + nodes, reward = ( + [ + node( + score=0, + children=["root.a"], + aggregation="gate", + config={"threshold": 1}, + ), + node("root.a", 0.75), + ], + 0, + ) + elif aggregation == "sequential": + nodes, reward = ( + [ + node(score=0, children=["root.a", "root.b"], aggregation="sequential"), + node("root.a", 0), + node("root.b", evaluated=False), + ], + 0, + ) + else: + nodes, reward = telemetry()["rubric"], 0.75 + assert ( + RewardAttributionGrader() + .run(subject(tmp_path, snapshot=telemetry(nodes), reward=reward)) + .status + is CheckStatus.PASS + ) + + +def test_unknown_aggregation_is_explicitly_incomplete_but_wrong_root_still_fails( + tmp_path, +): + nodes = [node(children=["root.a"], aggregation="unknown"), node("root.a")] + result = RewardAttributionGrader().run(subject(tmp_path, snapshot=telemetry(nodes))) + assert result.status is CheckStatus.SKIP + assert "unsupported custom" in " ".join(result.evidence) + assert ( + RewardAttributionGrader() + .run(subject(tmp_path, snapshot=telemetry(nodes), reward=0)) + .status + is CheckStatus.FAIL + ) + + +def test_attribution_cannot_pass_from_a_truncated_wire_prefix(tmp_path): + measured = subject(tmp_path, failure_reason="step failed (TimeoutError)") + assert RewardAttributionGrader().run(measured).status is CheckStatus.FAIL + + +def test_undeclared_task_and_rubric_capabilities_do_not_run(tmp_path): + measured = subject(tmp_path) + measured.manifest.capabilities.task_api = False + measured.manifest.capabilities.declared_task_count = {} + measured.manifest.capabilities.rubric_tree = False + for grader in ( + TaskDeclarationAccuracyGrader, + RubricIntrospectableGrader, + RewardAttributionGrader, + ): + assert not grader().applies_to(measured.manifest) + assert grader().run(measured).status is CheckStatus.SKIP diff --git a/tests/test_validation/test_runtime_execution.py b/tests/test_validation/test_runtime_execution.py index f6f792b15..8302628e6 100644 --- a/tests/test_validation/test_runtime_execution.py +++ b/tests/test_validation/test_runtime_execution.py @@ -7,6 +7,7 @@ from types import SimpleNamespace import pytest +import yaml from openenv.validation.policy import load_policy, PolicyError from openenv.validation.providers import ProviderError, StartupError from openenv.validation.report import CheckResult @@ -27,9 +28,27 @@ def package(tmp_path): root = tmp_path / "subject" shutil.copytree(FIXTURE, root, ignore=shutil.ignore_patterns("__pycache__")) + # Lifecycle tests use the minimal declaration; discovery tests opt in below. + path = root / "openenv.yaml" + document = yaml.safe_load(path.read_text()) + document["validation"]["capabilities"].update( + declared_tools=[], rubric_tree=False, task_api=False, declared_task_count={} + ) + path.write_text(yaml.safe_dump(document)) return root +@pytest.fixture +def discovery_package(package): + path = package / "openenv.yaml" + document = yaml.safe_load(path.read_text()) + document["validation"]["capabilities"].update( + task_api=True, declared_task_count={"train": 100}, rubric_tree=True + ) + path.write_text(yaml.safe_dump(document)) + return package + + @pytest.fixture def baseline_only(monkeypatch): """Isolate original-subject failure handling from separately tested replays.""" @@ -108,6 +127,26 @@ def measured_episode(seed=42): ) +def discovery_episode(): + from test_runtime_discovery import node, telemetry + + result = measured_episode() + snapshot = json.loads(result.telemetry_json) + snapshot.update(telemetry([node(score=1.0)])) + return replace( + result, + tools_json='{"tools":[]}', + tasks_json=json.dumps( + { + "splits": [{"name": "train"}], + "counts": {"train": 100}, + "previews": {"train": ["task-0", "task-1"]}, + } + ), + telemetry_json=json.dumps(snapshot), + ) + + def test_runtime_collects_primary_and_session_replays_then_cleans_up( package, monkeypatch, tmp_path ): @@ -263,7 +302,7 @@ def failed_build(*args): @pytest.mark.parametrize("mutation", ["content", "symlink", "unreadable"]) def test_source_change_withdraws_dependent_runtime_results( - package, monkeypatch, tmp_path, mutation, baseline_only + package, monkeypatch, tmp_path, mutation, baseline_only, discovery_package ): provider = FakeRuntimeProvider() original_digest = source_digest(package) @@ -280,7 +319,7 @@ def collect(*args, **kwargs): monkeypatch.setattr( "openenv.validation.runner.source_digest", unreadable_source ) - return measured_episode() + return discovery_episode() monkeypatch.setattr("openenv.validation.runner.collect_runtime_evidence", collect) bundle = tmp_path / "bundle" @@ -295,7 +334,15 @@ def collect(*args, **kwargs): else "package source could not be verified after validation" ) assert results["runtime.startup"].evidence == [reason] - for name in ("reward_well_formed", "observation_schema", "state_contract"): + for name in ( + "reward_well_formed", + "observation_schema", + "state_contract", + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", + ): result = results[f"runtime.{name}"] assert result.status is CheckStatus.SKIP assert "runtime.startup" in result.evidence[0] @@ -479,6 +526,128 @@ def test_semantic_ceiling_does_not_claim_semantic_execution(package): ) +def test_runner_checks_empty_tools_and_omits_undeclared_optional_capabilities( + package, monkeypatch, baseline_only +): + observed = [] + + def collect(*args, **kwargs): + observed.append(kwargs) + return replace(measured_episode(), tools_json='{"tools":[]}') + + monkeypatch.setattr("openenv.validation.runner.collect_runtime_evidence", collect) + result = run_validation( + package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() + ) + checks = {row.check_id: row for row in result.results} + assert checks["runtime.tool_declaration_accuracy"].status is CheckStatus.PASS + assert observed[0]["collect_tools"] is True + assert observed[0]["task_env_name"] is None + for name in ( + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", + ): + assert "runtime." + name not in checks + + +def test_runner_collects_and_grades_claimed_discovery_and_rubric_evidence( + discovery_package, monkeypatch, baseline_only +): + observed = [] + + def collect(*args, **kwargs): + observed.append(kwargs) + return discovery_episode() + + monkeypatch.setattr("openenv.validation.runner.collect_runtime_evidence", collect) + result = run_validation( + discovery_package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() + ) + checks = {row.check_id: row for row in result.results} + assert observed[0]["task_env_name"] == "validation_probe" + for name in ( + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", + ): + assert checks["runtime." + name].status is CheckStatus.PASS + + +@pytest.mark.parametrize("failed_collection", [False, True]) +def test_claims_without_evidence_are_incomplete_or_fail_not_pass( + discovery_package, monkeypatch, baseline_only, failed_collection +): + reason = "discovery failed (ValueError)" if failed_collection else None + measured = replace( + measured_episode(), + telemetry_json=None, + telemetry_error=reason, + tools_error=reason, + tasks_error=reason, + ) + monkeypatch.setattr( + "openenv.validation.runner.collect_runtime_evidence", lambda *a, **k: measured + ) + result = run_validation( + discovery_package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() + ) + checks = {row.check_id: row for row in result.results} + for name in ( + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + ): + assert checks["runtime." + name].status is ( + CheckStatus.FAIL if failed_collection else CheckStatus.SKIP + ) + assert checks["runtime.reward_attribution"].status is CheckStatus.SKIP + + +def test_all_claimed_discovery_checks_name_the_missing_startup_dependency( + discovery_package, +): + result = run_validation(discovery_package, max_level=Level.RUNTIME, skip_build=True) + checks = {row.check_id: row for row in result.results} + for name in ( + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", + ): + check = checks["runtime." + name] + assert check.status is CheckStatus.SKIP + assert check.evidence == ["unmet dependency: runtime.startup"] + + +def test_task_count_claim_is_checked_even_without_task_api_flag( + package, monkeypatch, baseline_only +): + path = package / "openenv.yaml" + document = yaml.safe_load(path.read_text()) + document["validation"]["capabilities"]["declared_task_count"] = {"train": 2} + path.write_text(yaml.safe_dump(document)) + observed = [] + + def collect(*args, **kwargs): + observed.append(kwargs) + return measured_episode() + + monkeypatch.setattr("openenv.validation.runner.collect_runtime_evidence", collect) + result = run_validation( + package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() + ) + task_check = next( + row + for row in result.results + if row.check_id == "runtime.task_declaration_accuracy" + ) + assert task_check.status is CheckStatus.SKIP + assert "missing prerequisite" in task_check.evidence[0] + assert observed[0]["task_env_name"] == "validation_probe" + + def grader(check_id, depends_on=(), *, status=CheckStatus.PASS, requires=frozenset()): return SimpleNamespace( check_id=check_id, diff --git a/tests/test_validation/test_task_discovery_collection.py b/tests/test_validation/test_task_discovery_collection.py new file mode 100644 index 000000000..5e51127ff --- /dev/null +++ b/tests/test_validation/test_task_discovery_collection.py @@ -0,0 +1,127 @@ +import json +import time + +import httpx +import pytest +from openenv.validation.runtime import discovery + + +def collect(monkeypatch, handler, **kwargs): + original = httpx.Client + + def client(**options): + assert options == {"trust_env": False, "follow_redirects": False} + return original(transport=httpx.MockTransport(handler), **options) + + monkeypatch.setattr(discovery.httpx, "Client", client) + return discovery.collect_task_evidence( + "http://127.0.0.1:8000", "name?secret", deadline=time.monotonic() + 1, **kwargs + ) + + +def test_task_sampling_uses_true_counts_and_two_bounded_specs(monkeypatch): + seen = [] + + def handler(request): + seen.append(request) + assert request.url.query == b"" + assert request.url.raw_path.startswith(b"/name%3Fsecret/") + if request.url.path.endswith("/splits"): + return httpx.Response(200, json=[{"name": "train"}, {"name": "empty"}]) + data = json.loads(request.content) + if request.url.path.endswith("/num_tasks"): + return httpx.Response( + 200, json={"num_tasks": 10**9 if data["split"] == "train" else 0} + ) + assert request.url.path.endswith("/task") and data["index"] < 2 + return httpx.Response(200, json={"task": [data["split"], data["index"]]}) + + payload, error = collect(monkeypatch, handler) + assert error is None + assert json.loads(payload) == { + "splits": [{"name": "train"}, {"name": "empty"}], + "counts": {"train": 10**9, "empty": 0}, + "previews": {"train": [["train", 0], ["train", 1]], "empty": []}, + } + assert len(seen) == 5 + + +@pytest.mark.parametrize( + "fault", + [ + "unsupported", + "redirect", + "too_large", + "duplicate", + "boolean_count", + "too_many_splits", + "compressed", + "missing_task", + ], +) +def test_bad_discovery_is_explicit_failure_without_private_values(monkeypatch, fault): + def handler(request): + if request.url.path.endswith("/splits"): + if fault == "unsupported": + return httpx.Response(501, text="private-error") + if fault == "redirect": + return httpx.Response( + 302, headers={"Location": "http://secret.invalid"} + ) + if fault == "too_large": + return httpx.Response( + 200, content=b"x" * (discovery.MAX_RESPONSE_BYTES + 1) + ) + if fault == "duplicate": + return httpx.Response(200, json=[{"name": "train"}] * 2) + if fault == "too_many_splits": + return httpx.Response(200, json=[{}] * (discovery.MAX_TASK_SPLITS + 1)) + if fault == "compressed": + return httpx.Response( + 200, headers={"Content-Encoding": "br"}, content=b"invalid" + ) + return httpx.Response(200, json=[{"name": "train"}]) + if request.url.path.endswith("/num_tasks"): + return httpx.Response( + 200, json={"num_tasks": True if fault == "boolean_count" else 1} + ) + return httpx.Response( + 200, + json={"wrong": "private-error"} + if fault == "missing_task" + else {"task": "valid"}, + ) + + payload, error = collect(monkeypatch, handler) + assert payload is None + assert error.startswith("task discovery failed (") + assert "private" not in error and "secret" not in error + + +def test_expired_deadline_makes_no_network_request(monkeypatch): + def handler(request): + pytest.fail("expired discovery budget must not contact the subject") + + payload, error = collect(monkeypatch, handler, request_timeout_s=0) + assert payload is None and error == "task discovery failed (TimeoutError)" + + +@pytest.mark.parametrize("split,task", [("train", "😀" * 40), ("s" * 200, None)]) +def test_retained_task_evidence_respects_byte_bound(monkeypatch, split, task): + monkeypatch.setattr(discovery, "MAX_DISCOVERY_BYTES", 400) + + def handler(request): + if request.url.path.endswith("/splits"): + return httpx.Response(200, json=[{"name": split}]) + if request.url.path.endswith("/num_tasks"): + return httpx.Response(200, json={"num_tasks": 1 if task else 0}) + return httpx.Response(200, json={"task": task}) + + payload, error = collect(monkeypatch, handler) + if task: + assert error is None + assert "😀" in payload and len(payload.encode()) <= 400 + else: + # Split names occur in three sections; received bytes alone do not bound + # the normalized evidence document. + assert payload is None and error == "task discovery failed (ValueError)" diff --git a/tests/validation_runtime/README.md b/tests/validation_runtime/README.md index d81b4ea37..86a300dc4 100644 --- a/tests/validation_runtime/README.md +++ b/tests/validation_runtime/README.md @@ -33,10 +33,19 @@ with networking disabled. A cold cache is supported. The tests run outside the checkout with `PYTHONPATH` removed, exercising installed package data and the production OpenEnv `/ws` endpoint. Each launch uses a fresh subject. -The image supports controlled `VALIDATION_FAULT` modes: `good`, `bad_reward`, -`bad_observation`, `missing_done`, `bad_state`, `hung_step`, and `startup_failure`. All fault -switches and wire corruption remain inside test assets. They share one fixture -and one public runtime plan, so a defect changes one property at a time. +The image includes controlled faults for malformed observations/rewards/state, +startup and timeout failures, ignored seeds, session drift, missing/mismatched +records, tool/task declarations, rubric configuration and child attribution. +`judged_stable` and `judged_noisy` exercise the bounded judge procedure. All fault +switches and wire corruption remain inside test assets with one public runtime plan. + +The fixture registers two real MCP tools and declares task counts of four train +and two test tasks. Its task listing returns only one preview item; the collector +uses `num_tasks` and samples at most two items per split. Empty tool declarations +are checked through a real empty FastMCP registry. `discovery.json` preserves raw +results and distinguishes failed discovery from an empty success. The protocol +suite contains 35 required cases, including real tool calls and discovery/rubric +faults through the installed-wheel process provider. Repeatability checks compare the original episode against a fresh session and an independently inspected fresh container, then exercise a different seed. Controlled @@ -46,8 +55,8 @@ telemetry, container identity and cleanup outcome. A process-only provider expli skips fresh-container determinism. Subject-emitted records are compared with the independently collected wire trace. -The Docker suite contains 19 required cases: three provider lifecycle tests, -15 CLI fault/control cases, and one real `echo_env` canary. The hung-step case +The Docker suite contains 26 required cases: three provider lifecycle tests, +22 CLI fault/control cases, and one real `echo_env` canary. The hung-step case checks the episode deadline; the interruption case sends SIGINT only after a container log confirms the second step has begun. Both must retain the completed reset/state/step/state prefix and remove their own containers. Each CLI case uses diff --git a/tests/validation_runtime/acceptance.json b/tests/validation_runtime/acceptance.json index 2a2eb4c53..8c788e090 100644 --- a/tests/validation_runtime/acceptance.json +++ b/tests/validation_runtime/acceptance.json @@ -28,7 +28,15 @@ "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[startup_failure-runtime.startup]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[ignored_seed-runtime.seed_control]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[missing_record-runtime.trajectory_record]", - "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[trace_mismatch-runtime.trajectory_record]" + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[trace_mismatch-runtime.trajectory_record]", + "tests.test_validation.integration.test_discovery_protocol::test_real_tools_share_the_replayed_instance_and_tasks_allow_bounded_listing", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[missing_tool-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[extra_tool-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[tool_discovery_error-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_task_count-runtime.task_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[missing_rubric_config-runtime.rubric_introspectable]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_attribution-runtime.reward_attribution]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[empty_tools-None]" ], "docker": [ "tests.test_validation.integration.test_docker_lifecycle::test_docker_lifecycle_effective_limits_and_owned_cleanup", @@ -49,7 +57,14 @@ "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[missing_record-runtime.trajectory_record]", "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[trace_mismatch-runtime.trajectory_record]", "tests.test_validation.integration.test_runtime_cli::test_cli_controlled_judge_requires_twenty_samples[judged_stable-pass]", - "tests.test_validation.integration.test_runtime_cli::test_cli_controlled_judge_requires_twenty_samples[judged_noisy-fail]" + "tests.test_validation.integration.test_runtime_cli::test_cli_controlled_judge_requires_twenty_samples[judged_noisy-fail]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[missing_tool-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[extra_tool-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[tool_discovery_error-runtime.tool_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[bad_task_count-runtime.task_declaration_accuracy]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[missing_rubric_config-runtime.rubric_introspectable]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[bad_attribution-runtime.reward_attribution]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[empty_tools-None]" ] } } From ae2c8d4bd7c447aea197edd9637a2e56d7f97d53 Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Thu, 24 Sep 2026 12:36:29 +0200 Subject: [PATCH 2/8] fix: include discovery in replay byte budget --- src/openenv/validation/runtime/replay.py | 2 ++ tests/test_validation/test_runtime_replay.py | 17 +++++++++++++++++ 2 files changed, 19 insertions(+) diff --git a/src/openenv/validation/runtime/replay.py b/src/openenv/validation/runtime/replay.py index fca6700cb..6bc876374 100644 --- a/src/openenv/validation/runtime/replay.py +++ b/src/openenv/validation/runtime/replay.py @@ -21,6 +21,8 @@ def _evidence_bytes(evidence): for value in ( evidence.observation_schema_json or "", evidence.telemetry_json or "", + evidence.tools_json or "", + evidence.tasks_json or "", *( value for row in evidence.exchanges diff --git a/tests/test_validation/test_runtime_replay.py b/tests/test_validation/test_runtime_replay.py index 47300b4d2..7a2a8f0d6 100644 --- a/tests/test_validation/test_runtime_replay.py +++ b/tests/test_validation/test_runtime_replay.py @@ -277,6 +277,23 @@ def test_retained_byte_budget_counts_primary_and_all_raw_fields( assert not container.evidence.exchanges +@pytest.mark.parametrize("field", ["tools_json", "tasks_json"]) +@pytest.mark.parametrize("budget,collections", [(3, 0), (5, 1)]) +def test_discovery_payload_counts_toward_retained_replay_budget( + monkeypatch, field, budget, collections +): + values = list(inputs(monkeypatch)) + values[5] = RuntimeEvidence(**{field: '"é"'}) # Four UTF-8 bytes. + values[-1].return_value = RuntimeEvidence(observation_schema_json="{}") + monkeypatch.setattr(replay, "MAX_REPLAY_BYTES", budget) + result = run(values) + assert "total retained replay evidence exceeds" in result.replay_failure_reason + assert not result.replays + assert replay._evidence_bytes(result) <= budget + assert values[-1].call_count == collections + values[0].start.assert_not_called() + + def test_oversized_primary_evidence_is_explicitly_incomplete_and_not_retained( monkeypatch, ): From 8645817f9e36bf78248eca824160e7202d684749 Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Fri, 25 Sep 2026 09:29:02 +0200 Subject: [PATCH 3/8] fix: honor optional discovery declarations --- rfcs/008-environment-auto-validation.md | 4 ++ .../validation/graders/runtime/discovery.py | 14 ++-- src/openenv/validation/manifest.py | 17 ++++- src/openenv/validation/runner.py | 12 ++-- src/openenv/validation/runtime/collector.py | 9 ++- src/openenv/validation/runtime/discovery.py | 24 ++++++- .../integration/test_runtime_cli.py | 7 +- .../integration/test_runtime_process.py | 22 +++++-- .../test_discovery_collection.py | 9 ++- tests/test_validation/test_manifest.py | 30 ++++++++- .../test_validation/test_runtime_discovery.py | 64 +++++++++++++++++++ .../test_validation/test_runtime_execution.py | 62 ++++++++++++++++-- .../test_task_discovery_collection.py | 53 ++++++++++++++- tests/validation_runtime/README.md | 2 +- tests/validation_runtime/acceptance.json | 3 +- 15 files changed, 291 insertions(+), 41 deletions(-) diff --git a/rfcs/008-environment-auto-validation.md b/rfcs/008-environment-auto-validation.md index 3496ceeb6..43e2878c0 100644 --- a/rfcs/008-environment-auto-validation.md +++ b/rfcs/008-environment-auto-validation.md @@ -638,6 +638,10 @@ set. Task discovery queries split names and authoritative per-split counts, then reads at most two individual tasks per split. Preview length is never treated as the dataset size. Requests share the collection deadline and bounded response budgets; unavailable or malformed discovery cannot become an empty passing set. +The task namespace comes from the server's single `/list_environments` entry, +independently of the package name. Omitted tool/count declarations stay unknown; +explicit `[]`/`{}` declarations assert zero tools/splits. A declared task API +without declared counts leaves count accuracy incomplete. Rubric checks require the declared tree and explicit public configuration, then compare every observed step reward with its fresh root score and named child diff --git a/src/openenv/validation/graders/runtime/discovery.py b/src/openenv/validation/graders/runtime/discovery.py index 8c76540bd..a129ed76f 100644 --- a/src/openenv/validation/graders/runtime/discovery.py +++ b/src/openenv/validation/graders/runtime/discovery.py @@ -70,6 +70,9 @@ class ToolDeclarationAccuracyGrader(_EvidenceGrader): check_id = "runtime.tool_declaration_accuracy" field, error_field, label = "tools_json", "tools_error", "tool discovery" + def applies_to(self, manifest): + return "declared_tools" in manifest.capabilities.model_fields_set + def check(self, subject, payload): tools = payload["tools"] if not isinstance(tools, list): @@ -104,8 +107,8 @@ class TaskDeclarationAccuracyGrader(_EvidenceGrader): field, error_field, label = "tasks_json", "tasks_error", "task discovery" def applies_to(self, manifest): - return manifest.capabilities.task_api or bool( - manifest.capabilities.declared_task_count + return manifest.capabilities.task_api or ( + "declared_task_count" in manifest.capabilities.model_fields_set ) def check(self, subject, payload): @@ -129,7 +132,10 @@ def check(self, subject, payload): raise ValueError("incomplete split discovery") problems = [] declared = subject.manifest.capabilities.declared_task_count - if set(names) != declared.keys(): + has_counts = ( + "declared_task_count" in subject.manifest.capabilities.model_fields_set + ) + if has_counts and set(names) != declared.keys(): problems.append("discovered splits differ from declared task counts") for name in names: count, preview = counts[name], previews[name] @@ -143,7 +149,7 @@ def check(self, subject, payload): ) if name in declared and count != declared[name]: problems.append("advertised task count differs from the declaration") - return problems, None + return problems, None if has_counts else "task counts were not declared" def _rubric_nodes(payload): diff --git a/src/openenv/validation/manifest.py b/src/openenv/validation/manifest.py index de5a14333..3412834b7 100644 --- a/src/openenv/validation/manifest.py +++ b/src/openenv/validation/manifest.py @@ -9,7 +9,14 @@ from pathlib import PurePosixPath from typing import Annotated, Any, Literal -from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator +from pydantic import ( + BaseModel, + ConfigDict, + Field, + field_validator, + model_serializer, + model_validator, +) from .types import SignatureKind @@ -241,6 +248,14 @@ class CapabilitiesSpec(BaseModel): declared_tools: list[str] = Field(default_factory=list) declared_task_count: dict[str, int] = Field(default_factory=dict) + @model_serializer(mode="wrap") + def _preserve_declaration_presence(self, serialize): + data = serialize(self) + for field in ("declared_tools", "declared_task_count"): + if field not in self.model_fields_set: + data.pop(field, None) + return data + @model_validator(mode="after") def _injected_state_requires_set_state(self) -> "CapabilitiesSpec": if ( diff --git a/src/openenv/validation/runner.py b/src/openenv/validation/runner.py index 6874c89b1..8de526673 100644 --- a/src/openenv/validation/runner.py +++ b/src/openenv/validation/runner.py @@ -105,9 +105,9 @@ def _applicable(check_id, manifest): if check_id in {"runtime.rubric_introspectable", "runtime.reward_attribution"}: return manifest.capabilities.rubric_tree if check_id == "runtime.task_declaration_accuracy": - return manifest.capabilities.task_api or bool( - manifest.capabilities.declared_task_count - ) + return TaskDeclarationAccuracyGrader().applies_to(manifest) + if check_id == "runtime.tool_declaration_accuracy": + return ToolDeclarationAccuracyGrader().applies_to(manifest) return True @@ -177,10 +177,8 @@ def _runtime(subject, *, skip_build, provider): manifest.resources.episode_timeout_s, REPLAY_BUDGET_SECONDS ), validation_token=spec.env_vars["OPENENV_VALIDATION_TOKEN"], - collect_tools=True, - task_env_name=manifest.name - if _applicable("runtime.task_declaration_accuracy", manifest) - else None, + collect_tools=_applicable("runtime.tool_declaration_accuracy", manifest), + collect_tasks=_applicable("runtime.task_declaration_accuracy", manifest), ) # A health endpoint without a functioning protocol isn't a startup success. if not evidence.exchanges and evidence.failure_reason: diff --git a/src/openenv/validation/runtime/collector.py b/src/openenv/validation/runtime/collector.py index a50cbbfa6..5375d617a 100644 --- a/src/openenv/validation/runtime/collector.py +++ b/src/openenv/validation/runtime/collector.py @@ -75,7 +75,7 @@ def collect_runtime_evidence( request_timeout_s: float = 5.0, validation_token: str | None = None, collect_tools: bool = False, - task_env_name: str | None = None, + collect_tasks: bool = False, ) -> RuntimeEvidence: """ Preserve schema and reset/step/state responses without model coercion. @@ -97,8 +97,8 @@ def collect_runtime_evidence( Run-scoped telemetry authorization; never retained in evidence. collect_tools (`bool`, *optional*, defaults to `False`): Discover tools on the measured WebSocket before reset. - task_env_name (`str`, *optional*): - Sample this environment's task metadata through its HTTP task API. + collect_tasks (`bool`, *optional*, defaults to `False`): + Discover the server's environment namespace and sample its task metadata. Returns: [`~openenv.validation.runtime.contracts.RuntimeEvidence`]: raw evidence. @@ -265,12 +265,11 @@ def contains_credential(value): except Exception as exc: tools_error = f"tool discovery failed ({type(exc).__name__})" - if task_env_name is not None: + if collect_tasks: phase = "tasks" try: discovered, tasks_error = collect_task_evidence( base_url, - task_env_name, deadline=deadline, request_timeout_s=request_timeout_s, ) diff --git a/src/openenv/validation/runtime/discovery.py b/src/openenv/validation/runtime/discovery.py index 138d9dfd8..810294e31 100644 --- a/src/openenv/validation/runtime/discovery.py +++ b/src/openenv/validation/runtime/discovery.py @@ -11,7 +11,7 @@ MAX_DISCOVERY_BYTES = 8 * 1024 * 1024 -def collect_task_evidence(base_url, env_name, *, deadline, request_timeout_s=5.0): +def collect_task_evidence(base_url, *, deadline, request_timeout_s=5.0): """Return raw counts and at most two task specs per split, never a full listing. ``deadline`` is the collector's absolute monotonic episode deadline. Errors @@ -27,7 +27,7 @@ def remaining(): return budget try: - prefix = base_url.rstrip("/") + "/" + quote(env_name, safe="") + prefix = base_url.rstrip("/") with httpx.Client(trust_env=False, follow_redirects=False) as client: def request(method, route, payload=None): @@ -59,6 +59,18 @@ def request(method, route, payload=None): body.extend(chunk) return json.loads(body) + environments = request("GET", "/list_environments") + if ( + not isinstance(environments, list) + or len(environments) != 1 + or not isinstance(environments[0], str) + or not environments[0] + or environments[0] in {".", ".."} + or "/" in environments[0] + or "\\" in environments[0] + ): + raise ValueError("expected one valid task environment namespace") + prefix += "/" + quote(environments[0], safe="") splits = request("GET", "/splits") if not isinstance(splits, list) or len(splits) > MAX_TASK_SPLITS: raise ValueError("invalid or excessive splits") @@ -76,7 +88,12 @@ def request(method, route, payload=None): for index in range(min(2, count)) ] discovered = json.dumps( - {"splits": splits, "counts": counts, "previews": previews}, + { + "environments": environments, + "splits": splits, + "counts": counts, + "previews": previews, + }, allow_nan=False, ensure_ascii=False, separators=(",", ":"), @@ -86,6 +103,7 @@ def request(method, route, payload=None): return discovered, None except ( httpx.HTTPError, + httpx.InvalidURL, OSError, ValueError, TypeError, diff --git a/tests/test_validation/integration/test_runtime_cli.py b/tests/test_validation/integration/test_runtime_cli.py index ec89f19dd..b26af734b 100644 --- a/tests/test_validation/integration/test_runtime_cli.py +++ b/tests/test_validation/integration/test_runtime_cli.py @@ -272,9 +272,14 @@ def _invoke_cli( assert len(checks) == len(report["results"]), "Report contains duplicate check IDs" applicable = IMPLEMENTED | PENDING capabilities = report["manifest"]["capabilities"] + declarations = yaml.safe_load((context / "openenv.yaml").read_text())["validation"][ + "capabilities" + ] + if "declared_tools" not in declarations: + applicable -= {"runtime.tool_declaration_accuracy"} if not capabilities["rubric_tree"]: applicable -= {"runtime.rubric_introspectable", "runtime.reward_attribution"} - if not capabilities["task_api"] and not capabilities["declared_task_count"]: + if not capabilities["task_api"] and "declared_task_count" not in declarations: applicable -= {"runtime.task_declaration_accuracy"} assert {key for key in checks if key.startswith("runtime.")} == applicable assert all(checks[key]["status"] == "skip" for key in PENDING) diff --git a/tests/test_validation/integration/test_runtime_process.py b/tests/test_validation/integration/test_runtime_process.py index 9518e67dc..f5af66784 100644 --- a/tests/test_validation/integration/test_runtime_process.py +++ b/tests/test_validation/integration/test_runtime_process.py @@ -176,6 +176,7 @@ def start(self, spec): ("missing_rubric_config", "runtime.rubric_introspectable"), ("bad_attribution", "runtime.reward_attribution"), ("empty_tools", None), + ("namespace_mismatch", None), ], ) def test_installed_server_collector_and_graders_over_loopback( @@ -187,14 +188,19 @@ def test_installed_server_collector_and_graders_over_loopback( / mode ) target = FIXTURE - if mode == "empty_tools": - target = tmp_path / "empty-tools" + if mode in {"empty_tools", "namespace_mismatch"}: + target = tmp_path / mode shutil.copytree(FIXTURE, target, ignore=shutil.ignore_patterns("__pycache__")) manifest_path = target / "openenv.yaml" manifest = yaml.safe_load(manifest_path.read_text()) - manifest["validation"]["capabilities"]["declared_tools"] = [] + if mode == "empty_tools": + manifest["validation"]["capabilities"]["declared_tools"] = [] + else: + manifest["name"] = "package_metadata_name" manifest_path.write_text(yaml.safe_dump(manifest, sort_keys=False)) - provider = ProcessProvider(artifacts / "subject", mode) + provider = ProcessProvider( + artifacts / "subject", "good" if mode == "namespace_mismatch" else mode + ) try: report = run_validation( target, @@ -229,9 +235,11 @@ def test_installed_server_collector_and_graders_over_loopback( manifest = json.loads((artifacts / "report/run-manifest.json").read_text()) assert manifest["provider"]["isolation"] == "process" assert manifest["provider"]["container_build_exercised"] is False - _assert_discovery( - json.loads((artifacts / "report/discovery.json").read_text()), mode - ) + discovery = json.loads((artifacts / "report/discovery.json").read_text()) + _assert_discovery(discovery, mode) + assert discovery["tasks"]["environments"] == ["validation_probe"] + if mode == "namespace_mismatch": + assert report.manifest.name == "package_metadata_name" telemetry_path = artifacts / "report/session-telemetry.json" telemetry = json.loads(telemetry_path.read_text()) assert telemetry["seed"]["accepted"] is (mode != "ignored_seed") diff --git a/tests/test_validation/test_discovery_collection.py b/tests/test_validation/test_discovery_collection.py index 926974453..c095668ae 100644 --- a/tests/test_validation/test_discovery_collection.py +++ b/tests/test_validation/test_discovery_collection.py @@ -66,6 +66,8 @@ def respond(request): return httpx.Response(200, json={"observation": {"type": "object"}}) if task_fault == "unsupported": return httpx.Response(501, text="private server details") + if request.url.path == "/list_environments": + return httpx.Response(200, json=["probe"]) if request.url.path.endswith("/splits"): return httpx.Response(200, json=[{"name": "train"}]) if request.url.path.endswith("/num_tasks"): @@ -117,7 +119,7 @@ def test_discovery_stays_on_measured_socket_outside_trajectory(monkeypatch): tool_response([{"name": "echo", "inputSchema": {"type": "object"}}]) ) evidence, urls = collect( - monkeypatch, connection, collect_tools=True, task_env_name="probe" + monkeypatch, connection, collect_tools=True, collect_tasks=True ) assert [request["type"] for request in connection.requests] == [ "validation_open", @@ -148,6 +150,7 @@ def test_discovery_stays_on_measured_socket_outside_trajectory(monkeypatch): assert json.loads(evidence.tasks_json)["counts"] == {"train": 100} assert urls == [ "/schema", + "/list_environments", "/probe/splits", "/probe/num_tasks", "/probe/task", @@ -213,7 +216,7 @@ def test_discovery_credential_echo_is_discarded_even_with_unicode_escapes( monkeypatch, DiscoveryConnection(response), collect_tools=True, - task_env_name="probe", + collect_tasks=True, task_fault=secret, ) assert evidence.tools_json is evidence.tasks_json is None @@ -227,7 +230,7 @@ def test_task_failure_does_not_poison_successful_tools_or_episode(monkeypatch): monkeypatch, DiscoveryConnection(), collect_tools=True, - task_env_name="probe", + collect_tasks=True, task_fault="unsupported", ) assert json.loads(evidence.tools_json) == {"tools": []} diff --git a/tests/test_validation/test_manifest.py b/tests/test_validation/test_manifest.py index 7f7dd3670..fb091e3fb 100644 --- a/tests/test_validation/test_manifest.py +++ b/tests/test_validation/test_manifest.py @@ -1,10 +1,16 @@ +import json + import pytest from conftest import ( INVALID_MANIFEST_FIXTURES, load_fixture_manifest, VALID_MANIFEST_FIXTURES, ) -from openenv.validation.manifest import NormalizedManifest, VerifierBinding +from openenv.validation.manifest import ( + CapabilitiesSpec, + NormalizedManifest, + VerifierBinding, +) from pydantic import ValidationError @@ -17,6 +23,28 @@ def test_valid_fixture_manifests_round_trip(name): ) +@pytest.mark.parametrize("tools", [False, True]) +@pytest.mark.parametrize("counts", [False, True]) +def test_declaration_omission_survives_dict_and_json_round_trips(tools, counts): + declarations = {} + if tools: + declarations["declared_tools"] = [] + if counts: + declarations["declared_task_count"] = {} + capabilities = CapabilitiesSpec(verifier={"kind": "reward_channel"}, **declarations) + for payload in ( + capabilities.model_dump(), + json.loads(capabilities.model_dump_json()), + ): + assert ("declared_tools" in payload) is tools + assert ("declared_task_count" in payload) is counts + assert payload["task_api"] is False # Other defaults remain serialized. + restored = CapabilitiesSpec.model_validate(payload) + assert ("declared_tools" in restored.model_fields_set) is tools + assert ("declared_task_count" in restored.model_fields_set) is counts + assert restored.model_dump() == payload + + @pytest.mark.parametrize("name", INVALID_MANIFEST_FIXTURES) def test_invalid_fixture_manifests_are_rejected(name): data = load_fixture_manifest(name) diff --git a/tests/test_validation/test_runtime_discovery.py b/tests/test_validation/test_runtime_discovery.py index 6e4ed1d58..855111702 100644 --- a/tests/test_validation/test_runtime_discovery.py +++ b/tests/test_validation/test_runtime_discovery.py @@ -123,6 +123,69 @@ def test_genuinely_empty_tools_pass_but_unsupported_discovery_does_not(tmp_path) assert ToolDeclarationAccuracyGrader().run(measured).status is CheckStatus.FAIL +@pytest.mark.parametrize("declared", [False, True]) +@pytest.mark.parametrize("tools", [[], [{"name": "echo"}]]) +def test_omitted_tools_are_unknown_but_explicit_empty_tools_are_checked( + tmp_path, declared, tools +): + measured = subject(tmp_path, tools_json=json.dumps({"tools": tools})) + capabilities = measured.manifest.capabilities + data = capabilities.model_dump(exclude={"declared_tools"}) + if declared: + data["declared_tools"] = [] + measured.manifest.capabilities = type(capabilities).model_validate(data) + grader = ToolDeclarationAccuracyGrader() + assert grader.applies_to(measured.manifest) is declared + assert grader.run(measured).status is ( + CheckStatus.SKIP + if not declared + else CheckStatus.FAIL + if tools + else CheckStatus.PASS + ) + + +@pytest.mark.parametrize("task_api", [False, True]) +@pytest.mark.parametrize("declared", [False, True]) +@pytest.mark.parametrize("empty", [False, True]) +def test_omitted_task_counts_are_unknown_but_explicit_empty_counts_are_checked( + tmp_path, task_api, declared, empty +): + measured = subject(tmp_path) + capabilities = measured.manifest.capabilities + data = capabilities.model_dump(exclude={"declared_task_count"}) + data["task_api"] = task_api + if declared: + data["declared_task_count"] = {} + measured.manifest.capabilities = type(capabilities).model_validate(data) + if empty: + measured.runtime_evidence.tasks_json = json.dumps( + {"splits": [], "counts": {}, "previews": {}} + ) + grader = TaskDeclarationAccuracyGrader() + assert grader.applies_to(measured.manifest) is (task_api or declared) + result = grader.run(measured) + assert result.status is ( + CheckStatus.SKIP + if not declared + else CheckStatus.PASS + if empty + else CheckStatus.FAIL + ) + if task_api and not declared: + assert result.evidence == ["task counts were not declared"] + + +def test_omitted_counts_do_not_hide_an_invalid_declared_task_api(tmp_path): + measured = subject(tmp_path) + measured.manifest.capabilities.model_fields_set.discard("declared_task_count") + measured.manifest.capabilities.declared_task_count.clear() + payload = json.loads(measured.runtime_evidence.tasks_json) + payload["counts"]["train"] = True + measured.runtime_evidence.tasks_json = json.dumps(payload) + assert TaskDeclarationAccuracyGrader().run(measured).status is CheckStatus.FAIL + + @pytest.mark.parametrize("tools", [[], [{"name": "echo"}]]) def test_matching_or_partial_first_tool_page_is_explicitly_incomplete(tmp_path, tools): result = ToolDeclarationAccuracyGrader().run( @@ -375,6 +438,7 @@ def test_undeclared_task_and_rubric_capabilities_do_not_run(tmp_path): measured = subject(tmp_path) measured.manifest.capabilities.task_api = False measured.manifest.capabilities.declared_task_count = {} + measured.manifest.capabilities.model_fields_set.discard("declared_task_count") measured.manifest.capabilities.rubric_tree = False for grader in ( TaskDeclarationAccuracyGrader, diff --git a/tests/test_validation/test_runtime_execution.py b/tests/test_validation/test_runtime_execution.py index cf601939f..fff8a4800 100644 --- a/tests/test_validation/test_runtime_execution.py +++ b/tests/test_validation/test_runtime_execution.py @@ -32,8 +32,9 @@ def package(tmp_path): path = root / "openenv.yaml" document = yaml.safe_load(path.read_text()) document["validation"]["capabilities"].update( - declared_tools=[], rubric_tree=False, task_api=False, declared_task_count={} + declared_tools=[], rubric_tree=False, task_api=False ) + document["validation"]["capabilities"].pop("declared_task_count", None) path.write_text(yaml.safe_dump(document)) return root @@ -597,7 +598,7 @@ def collect(*args, **kwargs): checks = {row.check_id: row for row in result.results} assert checks["runtime.tool_declaration_accuracy"].status is CheckStatus.PASS assert observed[0]["collect_tools"] is True - assert observed[0]["task_env_name"] is None + assert observed[0]["collect_tasks"] is False for name in ( "task_declaration_accuracy", "rubric_introspectable", @@ -620,7 +621,7 @@ def collect(*args, **kwargs): discovery_package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() ) checks = {row.check_id: row for row in result.results} - assert observed[0]["task_env_name"] == "validation_probe" + assert observed[0]["collect_tasks"] is True for name in ( "tool_declaration_accuracy", "task_declaration_accuracy", @@ -630,6 +631,59 @@ def collect(*args, **kwargs): assert checks["runtime." + name].status is CheckStatus.PASS +@pytest.mark.parametrize("declared", [False, True]) +@pytest.mark.parametrize("task_api", [False, True]) +def test_parsed_declaration_presence_controls_collection_and_findings( + package, monkeypatch, baseline_only, declared, task_api +): + path = package / "openenv.yaml" + document = yaml.safe_load(path.read_text()) + capabilities = document["validation"]["capabilities"] + capabilities["task_api"] = task_api + for field, empty in (("declared_tools", []), ("declared_task_count", {})): + if declared: + capabilities[field] = empty + else: + capabilities.pop(field, None) + path.write_text(yaml.safe_dump(document)) + observed = [] + + def collect(*args, **kwargs): + observed.append(kwargs) + # Omitted counts are unknown even when a working API advertises tasks. + return replace( + discovery_episode(), + tools_json='{"tools":[]}', + tasks_json='{"splits":[],"counts":{},"previews":{}}' + if declared + else discovery_episode().tasks_json, + ) + + monkeypatch.setattr("openenv.validation.runner.collect_runtime_evidence", collect) + result = run_validation( + package, max_level=Level.RUNTIME, provider=FakeRuntimeProvider() + ) + parsed = result.manifest.capabilities.model_fields_set + assert ("declared_tools" in parsed) is declared + assert ("declared_task_count" in parsed) is declared + restored = type(result).model_validate_json(result.model_dump_json()) + restored_fields = restored.manifest.capabilities.model_fields_set + assert ("declared_tools" in restored_fields) is declared + assert ("declared_task_count" in restored_fields) is declared + assert observed[0]["collect_tools"] is declared + assert observed[0]["collect_tasks"] is (declared or task_api) + checks = {row.check_id: row for row in result.results} + assert ("runtime.tool_declaration_accuracy" in checks) is declared + assert ("runtime.task_declaration_accuracy" in checks) is (declared or task_api) + if declared: + assert checks["runtime.tool_declaration_accuracy"].status is CheckStatus.PASS + assert checks["runtime.task_declaration_accuracy"].status is CheckStatus.PASS + elif task_api: + check = checks["runtime.task_declaration_accuracy"] + assert check.status is CheckStatus.SKIP + assert check.evidence == ["task counts were not declared"] + + @pytest.mark.parametrize("failed_collection", [False, True]) def test_claims_without_evidence_are_incomplete_or_fail_not_pass( discovery_package, monkeypatch, baseline_only, failed_collection @@ -700,7 +754,7 @@ def collect(*args, **kwargs): ) assert task_check.status is CheckStatus.SKIP assert "missing prerequisite" in task_check.evidence[0] - assert observed[0]["task_env_name"] == "validation_probe" + assert observed[0]["collect_tasks"] is True def grader(check_id, depends_on=(), *, status=CheckStatus.PASS, requires=frozenset()): diff --git a/tests/test_validation/test_task_discovery_collection.py b/tests/test_validation/test_task_discovery_collection.py index 5e51127ff..d2a653db1 100644 --- a/tests/test_validation/test_task_discovery_collection.py +++ b/tests/test_validation/test_task_discovery_collection.py @@ -6,16 +6,25 @@ from openenv.validation.runtime import discovery -def collect(monkeypatch, handler, **kwargs): +def collect(monkeypatch, handler, *, advertised=None, **kwargs): original = httpx.Client + def respond(request): + if request.url.path == "/list_environments": + return ( + advertised + if advertised is not None + else httpx.Response(200, json=["name?secret"]) + ) + return handler(request) + def client(**options): assert options == {"trust_env": False, "follow_redirects": False} - return original(transport=httpx.MockTransport(handler), **options) + return original(transport=httpx.MockTransport(respond), **options) monkeypatch.setattr(discovery.httpx, "Client", client) return discovery.collect_task_evidence( - "http://127.0.0.1:8000", "name?secret", deadline=time.monotonic() + 1, **kwargs + "http://127.0.0.1:8000", deadline=time.monotonic() + 1, **kwargs ) @@ -39,6 +48,7 @@ def handler(request): payload, error = collect(monkeypatch, handler) assert error is None assert json.loads(payload) == { + "environments": ["name?secret"], "splits": [{"name": "train"}, {"name": "empty"}], "counts": {"train": 10**9, "empty": 0}, "previews": {"train": [["train", 0], ["train", 1]], "empty": []}, @@ -46,6 +56,43 @@ def handler(request): assert len(seen) == 5 +@pytest.mark.parametrize( + "names", [[], {}, "probe", [None], [""], ["a", "b"], ["a", "a"], [".."], ["a/b"]] +) +def test_invalid_environment_inventory_is_not_guessed(monkeypatch, names): + def handler(request): + pytest.fail("invalid environment inventory must not request task routes") + + payload, error = collect( + monkeypatch, handler, advertised=httpx.Response(200, json=names) + ) + assert payload is None and error == "task discovery failed (ValueError)" + + +@pytest.mark.parametrize( + "fault", ["unsupported", "compressed", "oversized", "long_namespace"] +) +def test_environment_inventory_uses_discovery_response_bounds(monkeypatch, fault): + def handler(request): + pytest.fail("failed environment discovery must not request task routes") + + if fault == "unsupported": + response = httpx.Response(501, text="private-error") + elif fault == "compressed": + response = httpx.Response( + 200, headers={"Content-Encoding": "br"}, stream=httpx.ByteStream(b"invalid") + ) + elif fault == "long_namespace": + response = httpx.Response(200, json=["x" * 65537]) + else: + response = httpx.Response( + 200, content=b"x" * (discovery.MAX_RESPONSE_BYTES + 1) + ) + payload, error = collect(monkeypatch, handler, advertised=response) + assert payload is None and error.startswith("task discovery failed (") + assert "private" not in error + + @pytest.mark.parametrize( "fault", [ diff --git a/tests/validation_runtime/README.md b/tests/validation_runtime/README.md index 86a300dc4..5f2746a77 100644 --- a/tests/validation_runtime/README.md +++ b/tests/validation_runtime/README.md @@ -44,7 +44,7 @@ and two test tasks. Its task listing returns only one preview item; the collecto uses `num_tasks` and samples at most two items per split. Empty tool declarations are checked through a real empty FastMCP registry. `discovery.json` preserves raw results and distinguishes failed discovery from an empty success. The protocol -suite contains 35 required cases, including real tool calls and discovery/rubric +suite contains 36 required cases, including real tool calls and discovery/rubric faults through the installed-wheel process provider. Repeatability checks compare the original episode against a fresh session and an diff --git a/tests/validation_runtime/acceptance.json b/tests/validation_runtime/acceptance.json index 8c788e090..381046b7a 100644 --- a/tests/validation_runtime/acceptance.json +++ b/tests/validation_runtime/acceptance.json @@ -36,7 +36,8 @@ "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_task_count-runtime.task_declaration_accuracy]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[missing_rubric_config-runtime.rubric_introspectable]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_attribution-runtime.reward_attribution]", - "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[empty_tools-None]" + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[empty_tools-None]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[namespace_mismatch-None]" ], "docker": [ "tests.test_validation.integration.test_docker_lifecycle::test_docker_lifecycle_effective_limits_and_owned_cleanup", From c6ae9052dbfd7a332c5956ed4dc3cd4962d0af9a Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Thu, 1 Oct 2026 12:06:23 +0200 Subject: [PATCH 4/8] fix: honor task discovery deadline --- src/openenv/validation/runtime/discovery.py | 11 ++++--- .../test_task_discovery_collection.py | 30 +++++++++++++++++++ 2 files changed, 37 insertions(+), 4 deletions(-) diff --git a/src/openenv/validation/runtime/discovery.py b/src/openenv/validation/runtime/discovery.py index 810294e31..f9c2bb9cd 100644 --- a/src/openenv/validation/runtime/discovery.py +++ b/src/openenv/validation/runtime/discovery.py @@ -11,17 +11,20 @@ MAX_DISCOVERY_BYTES = 8 * 1024 * 1024 -def collect_task_evidence(base_url, *, deadline, request_timeout_s=5.0): +def collect_task_evidence(base_url, *, deadline, request_timeout_s=None): """Return raw counts and at most two task specs per split, never a full listing. - ``deadline`` is the collector's absolute monotonic episode deadline. Errors - are sanitized; a failed/unsupported endpoint never becomes an empty success. + ``deadline`` is the collector's absolute monotonic episode deadline. Optional + ``request_timeout_s`` may further cap each request. Errors are sanitized; + a failed/unsupported endpoint never becomes an empty success. Task metadata routes intentionally use independent environment instances. """ total_bytes = 0 def remaining(): - budget = min(request_timeout_s, deadline - time.monotonic()) + budget = deadline - time.monotonic() + if request_timeout_s is not None: + budget = min(request_timeout_s, budget) if budget <= 0: raise TimeoutError("task discovery deadline exceeded") return budget diff --git a/tests/test_validation/test_task_discovery_collection.py b/tests/test_validation/test_task_discovery_collection.py index d2a653db1..2acd195ec 100644 --- a/tests/test_validation/test_task_discovery_collection.py +++ b/tests/test_validation/test_task_discovery_collection.py @@ -153,6 +153,36 @@ def handler(request): assert payload is None and error == "task discovery failed (TimeoutError)" +@pytest.mark.parametrize("request_timeout_s,expected", [(None, 12.0), (3.0, 3.0)]) +def test_request_timeout_uses_episode_deadline_unless_explicitly_capped( + monkeypatch, request_timeout_s, expected +): + original = httpx.Client + timeouts = [] + + def respond(request): + timeouts.append(request.extensions["timeout"]) + return httpx.Response( + 200, + json=["probe"] if request.url.path == "/list_environments" else [], + ) + + monkeypatch.setattr(discovery.time, "monotonic", lambda: 100.0) + monkeypatch.setattr( + discovery.httpx, + "Client", + lambda **options: original(transport=httpx.MockTransport(respond), **options), + ) + payload, error = discovery.collect_task_evidence( + "http://127.0.0.1:8000", + deadline=112.0, + request_timeout_s=request_timeout_s, + ) + assert error is None and json.loads(payload)["splits"] == [] + assert len(timeouts) == 2 + assert all(set(timeout.values()) == {expected} for timeout in timeouts) + + @pytest.mark.parametrize("split,task", [("train", "😀" * 40), ("s" * 200, None)]) def test_retained_task_evidence_respects_byte_bound(monkeypatch, split, task): monkeypatch.setattr(discovery, "MAX_DISCOVERY_BYTES", 400) From 4296c67e83987c0c5904070fde4378d235b275a3 Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Thu, 1 Oct 2026 12:13:03 +0200 Subject: [PATCH 5/8] fix: bound task HTTP requests --- src/openenv/validation/runtime/discovery.py | 28 +++++-- .../test_task_discovery_collection.py | 76 ++++++++++++++++++- 2 files changed, 95 insertions(+), 9 deletions(-) diff --git a/src/openenv/validation/runtime/discovery.py b/src/openenv/validation/runtime/discovery.py index f9c2bb9cd..0b792e07b 100644 --- a/src/openenv/validation/runtime/discovery.py +++ b/src/openenv/validation/runtime/discovery.py @@ -6,6 +6,8 @@ import httpx +from .transport import http_deadline + MAX_TASK_SPLITS = 64 MAX_RESPONSE_BYTES = 1024 * 1024 MAX_DISCOVERY_BYTES = 8 * 1024 * 1024 @@ -31,17 +33,27 @@ def remaining(): try: prefix = base_url.rstrip("/") - with httpx.Client(trust_env=False, follow_redirects=False) as client: + with httpx.Client( + trust_env=False, + follow_redirects=False, + limits=httpx.Limits(max_keepalive_connections=0), + ) as client: def request(method, route, payload=None): nonlocal total_bytes - with client.stream( - method, - prefix + route, - json=payload, - timeout=remaining(), - headers={"Accept-Encoding": "identity"}, - ) as response: + # Each request opens a transport for its deadline watchdog; + # pooled connections do not emit the connection trace event. + with ( + http_deadline(remaining()) as extensions, + client.stream( + method, + prefix + route, + json=payload, + timeout=remaining(), + headers={"Accept-Encoding": "identity"}, + extensions=extensions, + ) as response, + ): response.raise_for_status() if ( response.headers.get("Content-Encoding", "identity") diff --git a/tests/test_validation/test_task_discovery_collection.py b/tests/test_validation/test_task_discovery_collection.py index 2acd195ec..0ebca06ae 100644 --- a/tests/test_validation/test_task_discovery_collection.py +++ b/tests/test_validation/test_task_discovery_collection.py @@ -1,5 +1,7 @@ import json +import threading import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer import httpx import pytest @@ -19,7 +21,11 @@ def respond(request): return handler(request) def client(**options): - assert options == {"trust_env": False, "follow_redirects": False} + assert options["limits"].max_keepalive_connections == 0 + assert {key: value for key, value in options.items() if key != "limits"} == { + "trust_env": False, + "follow_redirects": False, + } return original(transport=httpx.MockTransport(respond), **options) monkeypatch.setattr(discovery.httpx, "Client", client) @@ -183,6 +189,74 @@ def respond(request): assert all(set(timeout.values()) == {expected} for timeout in timeouts) +@pytest.mark.parametrize("phase", ["headers", "chunk_header"]) +@pytest.mark.parametrize("request_timeout_s", [None, 0.2]) +def test_dripping_task_response_cannot_extend_the_http_deadline( + phase, request_timeout_s +): + stopped = threading.Event() + paths = [] + + class Handler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def do_GET(self): + paths.append(self.path) + if self.path == "/list_environments": + body = b'["probe"]' + self.send_response(200) + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + return + assert self.path == "/probe/splits" + self.close_connection = True + prefix = ( + b"HTTP/1.1 200 OK\r\nX-Slow: " + if phase == "headers" + else b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n2;" + ) + suffix = ( + b"\r\nContent-Length: 2\r\n\r\n[]" + if phase == "headers" + else b"\r\n[]\r\n0\r\n\r\n" + ) + try: + self.connection.sendall(prefix) + # Each byte arrives before the inactivity timeout. The full + # header or chunk extension still exceeds the request budget. + for _ in range(100): + if stopped.wait(0.01): + return + self.connection.sendall(b"x") + self.connection.sendall(suffix) + except OSError: + pass # The deadline closes the peer's transport mid-response. + + def log_message(self, *args): + pass + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + started = time.monotonic() + payload, error = discovery.collect_task_evidence( + f"http://127.0.0.1:{server.server_port}", + deadline=started + (0.2 if request_timeout_s is None else 5), + request_timeout_s=request_timeout_s, + ) + elapsed = time.monotonic() - started + finally: + stopped.set() + server.shutdown() + server.server_close() + thread.join(timeout=1) + assert paths == ["/list_environments", "/probe/splits"] + assert payload is None and error == "task discovery failed (TimeoutError)" + assert elapsed < 0.75 + + @pytest.mark.parametrize("split,task", [("train", "😀" * 40), ("s" * 200, None)]) def test_retained_task_evidence_respects_byte_bound(monkeypatch, split, task): monkeypatch.setattr(discovery, "MAX_DISCOVERY_BYTES", 400) From 5b3d4b0393ec12e584fff9f4643ea10645d4792d Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Fri, 2 Oct 2026 13:33:18 +0200 Subject: [PATCH 6/8] fix: handle unscored attribution --- rfcs/008-environment-auto-validation.md | 5 ++ .../validation/graders/runtime/discovery.py | 16 ++++-- .../test_validation/test_runtime_discovery.py | 49 +++++++++++++++++++ 3 files changed, 65 insertions(+), 5 deletions(-) diff --git a/rfcs/008-environment-auto-validation.md b/rfcs/008-environment-auto-validation.md index 43e2878c0..bbaecd3aa 100644 --- a/rfcs/008-environment-auto-validation.md +++ b/rfcs/008-environment-auto-validation.md @@ -622,6 +622,11 @@ evaluation flags. Unevaluated gated children never reuse an earlier score. Stock container aggregation is named explicitly; custom rubrics may supply `validation_config()` to expose public JSON configuration. Arbitrary attributes and `state_dict()` are never serialized as configuration. +Attribution compares each numeric step reward with its freshly evaluated root +score. Null non-terminal rewards have no emitted score to compare; their records +must still have the correct identity, unchanged configuration, and valid evaluated +aggregation. An entirely unscored episode is incomplete rather than passing +attribution. Terminal null rewards remain invalid. The subject server also emits a bounded record of the reset/step/state request and response envelopes it executed. The validator captures the wire separately diff --git a/src/openenv/validation/graders/runtime/discovery.py b/src/openenv/validation/graders/runtime/discovery.py index a129ed76f..77f05f66f 100644 --- a/src/openenv/validation/graders/runtime/discovery.py +++ b/src/openenv/validation/graders/runtime/discovery.py @@ -271,6 +271,7 @@ def definition(tree): } unsupported = False + scored_steps = 0 for index, (record, step) in enumerate(zip(records, steps, strict=True)): if type(record["step_index"]) is not int or record["step_index"] != index: problems.append( @@ -279,9 +280,12 @@ def definition(tree): nodes = _rubric_nodes(record["rubric"]) if definition(nodes) != definition(baseline): problems.append(f"step {index}: rubric configuration changed") - reward = json.loads(step.response_json)["data"]["reward"] + data = json.loads(step.response_json)["data"] + reward = data["reward"] root = nodes["root"] - if ( + unscored = reward is None and data.get("done") is False + scored_steps += not unscored + if not unscored and ( not root["evaluated"] or not _finite(reward) or not math.isclose(root["score"], reward, rel_tol=1e-9, abs_tol=1e-9) @@ -329,7 +333,9 @@ def definition(tree): problems.append( f"step {index}: child attribution does not match the parent score" ) - return ( - problems, - "unsupported custom rubric aggregation semantics" if unsupported else None, + incomplete = ( + "unsupported custom rubric aggregation semantics" if unsupported else None ) + if not scored_steps: + incomplete = "missing prerequisite: an observed numeric step reward" + return problems, incomplete diff --git a/tests/test_validation/test_runtime_discovery.py b/tests/test_validation/test_runtime_discovery.py index 855111702..2a9dc9f3a 100644 --- a/tests/test_validation/test_runtime_discovery.py +++ b/tests/test_validation/test_runtime_discovery.py @@ -115,6 +115,55 @@ def test_matching_discovery_and_fresh_attribution_pass(tmp_path, grader): assert grader().run(subject(tmp_path)).status is CheckStatus.PASS +@pytest.mark.parametrize("scored", [True, False]) +@pytest.mark.parametrize("evaluated", [True, False]) +def test_unscored_steps_do_not_require_an_emitted_root_score( + tmp_path, scored, evaluated +): + payload = telemetry() + if not evaluated: + for item in payload["attribution"][0]["rubric"]: + item.update(evaluated=False, score=None) + measured = subject(tmp_path, snapshot=payload, reward=None) + if scored: + payload["attribution"].append( + {"step_index": 1, "rubric": copy.deepcopy(telemetry()["rubric"])} + ) + measured.runtime_evidence.telemetry_json = json.dumps(payload) + measured.runtime_evidence.exchanges += subject( + tmp_path + ).runtime_evidence.exchanges + result = RewardAttributionGrader().run(measured) + assert result.status is (CheckStatus.PASS if scored else CheckStatus.SKIP) + + +@pytest.mark.parametrize("mutation", ["identity", "configuration", "aggregation"]) +def test_unscored_steps_still_validate_supplied_attribution(tmp_path, mutation): + payload = telemetry() + record = payload["attribution"][0] + if mutation == "identity": + record["step_index"] = 1 + elif mutation == "configuration": + record["rubric"][0]["config"]["weights"] = [0.5, 0.5] + else: + record["rubric"][0]["score"] = 0.2 + result = RewardAttributionGrader().run( + subject(tmp_path, snapshot=payload, reward=None) + ) + assert result.status is CheckStatus.FAIL + + +def test_null_terminal_reward_is_not_an_unscored_step(tmp_path): + measured = subject(tmp_path, reward=None) + step = measured.runtime_evidence.exchanges[0] + response = json.loads(step.response_json) + response["data"]["done"] = True + measured.runtime_evidence.exchanges = ( + replace(step, response_json=json.dumps(response)), + ) + assert RewardAttributionGrader().run(measured).status is CheckStatus.FAIL + + def test_genuinely_empty_tools_pass_but_unsupported_discovery_does_not(tmp_path): measured = subject(tmp_path, tools_json='{"tools": []}') measured.manifest.capabilities.declared_tools = [] From 3722100493be0c097df68a0ae3e984498a45288a Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Mon, 5 Oct 2026 13:47:10 +0200 Subject: [PATCH 7/8] fix: detect rubric config type changes --- .../validation/graders/runtime/discovery.py | 3 ++- .../validation/runtime/served_probe/app.py | 8 +++++++ .../integration/test_runtime_cli.py | 1 + .../integration/test_runtime_process.py | 1 + .../test_validation/test_runtime_discovery.py | 23 +++++++++++++++++++ tests/validation_runtime/README.md | 6 ++--- tests/validation_runtime/acceptance.json | 2 ++ 7 files changed, 40 insertions(+), 4 deletions(-) diff --git a/src/openenv/validation/graders/runtime/discovery.py b/src/openenv/validation/graders/runtime/discovery.py index 77f05f66f..87a63fa4b 100644 --- a/src/openenv/validation/graders/runtime/discovery.py +++ b/src/openenv/validation/graders/runtime/discovery.py @@ -7,6 +7,7 @@ from ...report import CheckResult from ...types import CheckStatus from .basic import _RuntimeGrader +from .repeatability import _difference def _finite(value): @@ -278,7 +279,7 @@ def definition(tree): f"step {index}: attribution identity differs from the wire" ) nodes = _rubric_nodes(record["rubric"]) - if definition(nodes) != definition(baseline): + if _difference(definition(nodes), definition(baseline)): problems.append(f"step {index}: rubric configuration changed") data = json.loads(step.response_json)["data"] reward = data["reward"] diff --git a/tests/fixtures/validation/runtime/served_probe/app.py b/tests/fixtures/validation/runtime/served_probe/app.py index 23e865bf4..e90d4dfb4 100644 --- a/tests/fixtures/validation/runtime/served_probe/app.py +++ b/tests/fixtures/validation/runtime/served_probe/app.py @@ -188,6 +188,13 @@ async def fault_send(message: dict[str, Any]): ]["counter"] = 999 elif self.mode == "bad_attribution": data["data"]["attribution"][0]["rubric"][1]["score"] = 0.25 + elif self.mode == "changed_rubric_config": + snapshot = data["data"] + for tree in [snapshot["rubric"]] + [ + record["rubric"] for record in snapshot["attribution"] + ]: + tree[1]["config"]["option"] = True + snapshot["attribution"][0]["rubric"][1]["config"]["option"] = 1 message = {**message, "text": json.dumps(data)} await send(message) @@ -219,6 +226,7 @@ def make_app(mode="good"): "bad_task_count", "missing_rubric_config", "bad_attribution", + "changed_rubric_config", "slow_step", }: raise ValueError(f"Unknown fixture mode: {mode}") diff --git a/tests/test_validation/integration/test_runtime_cli.py b/tests/test_validation/integration/test_runtime_cli.py index 76aeab7e4..42addd9b6 100644 --- a/tests/test_validation/integration/test_runtime_cli.py +++ b/tests/test_validation/integration/test_runtime_cli.py @@ -351,6 +351,7 @@ def _assert_fresh_replays(artifacts, identical_samples=3): ("bad_task_count", "runtime.task_declaration_accuracy"), ("missing_rubric_config", "runtime.rubric_introspectable"), ("bad_attribution", "runtime.reward_attribution"), + ("changed_rubric_config", "runtime.reward_attribution"), ("empty_tools", None), ], ) diff --git a/tests/test_validation/integration/test_runtime_process.py b/tests/test_validation/integration/test_runtime_process.py index f5af66784..7d4a67320 100644 --- a/tests/test_validation/integration/test_runtime_process.py +++ b/tests/test_validation/integration/test_runtime_process.py @@ -175,6 +175,7 @@ def start(self, spec): ("bad_task_count", "runtime.task_declaration_accuracy"), ("missing_rubric_config", "runtime.rubric_introspectable"), ("bad_attribution", "runtime.reward_attribution"), + ("changed_rubric_config", "runtime.reward_attribution"), ("empty_tools", None), ("namespace_mismatch", None), ], diff --git a/tests/test_validation/test_runtime_discovery.py b/tests/test_validation/test_runtime_discovery.py index 2a9dc9f3a..6b17e4ab5 100644 --- a/tests/test_validation/test_runtime_discovery.py +++ b/tests/test_validation/test_runtime_discovery.py @@ -427,6 +427,29 @@ def test_incorrect_or_stale_attribution_fails(tmp_path, change): ) +@pytest.mark.parametrize("before,after", [(True, 1), (False, 0), (1, 1.0)]) +def test_attribution_detects_rubric_configuration_type_changes(tmp_path, before, after): + payload = telemetry() + payload["rubric"][1]["config"] = {"nested": [{"value": after}]} + payload["attribution"][0]["rubric"][1]["config"] = {"nested": [{"value": before}]} + result = RewardAttributionGrader().run(subject(tmp_path, snapshot=payload)) + assert result.status is CheckStatus.FAIL + assert result.evidence == ["step 0: rubric configuration changed"] + + +def test_attribution_ignores_rubric_configuration_key_order(tmp_path): + payload = telemetry() + payload["rubric"][1]["config"] = {"enabled": True, "threshold": 1} + payload["attribution"][0]["rubric"][1]["config"] = { + "threshold": 1, + "enabled": True, + } + assert ( + RewardAttributionGrader().run(subject(tmp_path, snapshot=payload)).status + is CheckStatus.PASS + ) + + @pytest.mark.parametrize("aggregation", ["leaf", "gate", "sequential", "weighted_sum"]) def test_stock_aggregation_rules_and_sequential_short_circuit(tmp_path, aggregation): if aggregation == "leaf": diff --git a/tests/validation_runtime/README.md b/tests/validation_runtime/README.md index d27f7e04d..562f7e95b 100644 --- a/tests/validation_runtime/README.md +++ b/tests/validation_runtime/README.md @@ -44,7 +44,7 @@ and two test tasks. Its task listing returns only one preview item; the collecto uses `num_tasks` and samples at most two items per split. Empty tool declarations are checked through a real empty FastMCP registry. `discovery.json` preserves raw results and distinguishes failed discovery from an empty success. The protocol -suite contains 40 required cases, including real tool calls and discovery/rubric +suite contains 41 required cases, including real tool calls and discovery/rubric faults through the installed-wheel process provider. Repeatability checks replay the original plan in a fresh session, with a different @@ -55,8 +55,8 @@ telemetry, container identity and cleanup outcome. A process-only provider expli skips fresh-container determinism. Subject-emitted records are compared with the independently collected wire trace. -The Docker suite contains 27 required cases: three provider lifecycle tests, -23 CLI fault/control cases, and one real `echo_env` canary. The slow-step case +The Docker suite contains 28 required cases: three provider lifecycle tests, +24 CLI fault/control cases, and one real `echo_env` canary. The slow-step case completes a tool call taking more than five seconds within the declared episode budget. The hung-step case checks the episode deadline; the interruption case sends SIGINT only after a diff --git a/tests/validation_runtime/acceptance.json b/tests/validation_runtime/acceptance.json index 035382c80..4445840ea 100644 --- a/tests/validation_runtime/acceptance.json +++ b/tests/validation_runtime/acceptance.json @@ -40,6 +40,7 @@ "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_task_count-runtime.task_declaration_accuracy]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[missing_rubric_config-runtime.rubric_introspectable]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[bad_attribution-runtime.reward_attribution]", + "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[changed_rubric_config-runtime.reward_attribution]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[empty_tools-None]", "tests.test_validation.integration.test_runtime_process::test_installed_server_collector_and_graders_over_loopback[namespace_mismatch-None]" ], @@ -70,6 +71,7 @@ "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[bad_task_count-runtime.task_declaration_accuracy]", "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[missing_rubric_config-runtime.rubric_introspectable]", "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[bad_attribution-runtime.reward_attribution]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[changed_rubric_config-runtime.reward_attribution]", "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[empty_tools-None]" ] } From 2a9d1b2d7eef6bdf05290e2fa420f7d402ed540e Mon Sep 17 00:00:00 2001 From: burtenshaw Date: Mon, 5 Oct 2026 14:04:13 +0200 Subject: [PATCH 8/8] test: align stacked grader metadata --- tests/test_validation/test_runner.py | 2 +- tests/test_validation/test_runtime_execution.py | 5 ++++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/tests/test_validation/test_runner.py b/tests/test_validation/test_runner.py index a48415253..9fde1daa1 100644 --- a/tests/test_validation/test_runner.py +++ b/tests/test_validation/test_runner.py @@ -29,7 +29,7 @@ class Probe: level = Level.RUNTIME requires_provider = frozenset() requires_capabilities = ( - frozenset({"task_api"}) if mode == "capability" else frozenset() + frozenset({"set_state"}) if mode == "capability" else frozenset() ) def __init__(self, check_id, dependency): diff --git a/tests/test_validation/test_runtime_execution.py b/tests/test_validation/test_runtime_execution.py index bd5f9a764..77b37e4a1 100644 --- a/tests/test_validation/test_runtime_execution.py +++ b/tests/test_validation/test_runtime_execution.py @@ -822,7 +822,10 @@ def test_all_claimed_discovery_checks_name_the_missing_startup_dependency( ): check = checks["runtime." + name] assert check.status is CheckStatus.SKIP - assert check.evidence == ["unmet dependency: runtime.startup"] + dependencies = "runtime.startup" + if name == "reward_attribution": + dependencies += ", runtime.rubric_introspectable" + assert check.evidence == [f"unmet dependency: {dependencies}"] def test_task_count_claim_is_checked_even_without_task_api_flag(