diff --git a/docs/source/reference/cli.md b/docs/source/reference/cli.md index f46cb17dd5..c63b377189 100644 --- a/docs/source/reference/cli.md +++ b/docs/source/reference/cli.md @@ -51,9 +51,11 @@ supports that capability. Docker validation uses two containers concurrently during the fresh-container replay. Both use the same immutable image and are removed even if a check fails. -The seven implemented runtime checks cover startup, rewards, observation schemas, -state continuity, seed control, determinism and emitted trajectory records. Other -applicable Level 2 checks appear as `SKIP`; they make +The 11 implemented runtime checks cover startup, rewards, observation schemas, +state continuity, seed control, determinism, emitted trajectory records, tool and +task discovery, rubric introspection and reward attribution. Five checks remain +unimplemented: network policy, host containment, resource bounds, episode isolation +and oracle containment. Applicable unimplemented checks appear as `SKIP`; they make the result `WARN`, which exits zero and does not mean Level 2 is complete. `FAIL` exits 1, unsupported package formats exit 2, and internal or policy errors exit 3. `--level semantic` diff --git a/rfcs/008-environment-auto-validation.md b/rfcs/008-environment-auto-validation.md index 930f01311c..c840c20381 100644 --- a/rfcs/008-environment-auto-validation.md +++ b/rfcs/008-environment-auto-validation.md @@ -660,6 +660,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 @@ -668,6 +673,27 @@ 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. +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 +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 0000000000..87a63fa4ba --- /dev/null +++ b/src/openenv/validation/graders/runtime/discovery.py @@ -0,0 +1,342 @@ +"""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 +from .repeatability import _difference + + +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 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): + 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 ( + "declared_task_count" in manifest.capabilities.model_fields_set + ) + + 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 + 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] + 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 if has_counts else "task counts were not declared" + + +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 + 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( + f"step {index}: attribution identity differs from the wire" + ) + nodes = _rubric_nodes(record["rubric"]) + if _difference(definition(nodes), definition(baseline)): + problems.append(f"step {index}: rubric configuration changed") + data = json.loads(step.response_json)["data"] + reward = data["reward"] + root = nodes["root"] + 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) + ): + 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" + ) + 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/src/openenv/validation/manifest.py b/src/openenv/validation/manifest.py index a4623da3b6..d5509d6d3c 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 19dd007e4d..4dfe36e673 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, @@ -101,6 +107,10 @@ def default_grader_registry(policy: SeverityPolicy) -> GraderRegistry: registry.register(SeedControlGrader()) registry.register(EpisodeDeterminismGrader()) registry.register(TrajectoryRecordGrader()) + registry.register(ToolDeclarationAccuracyGrader()) + registry.register(TaskDeclarationAccuracyGrader()) + registry.register(RubricIntrospectableGrader()) + registry.register(RewardAttributionGrader()) return registry @@ -120,9 +130,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 @@ -196,6 +206,8 @@ def _runtime(subject, graders, *, skip_build, provider): manifest.resources.episode_timeout_s, REPLAY_BUDGET_SECONDS ), validation_token=spec.env_vars["OPENENV_VALIDATION_TOKEN"], + collect_tools=_applicable("runtime.tool_declaration_accuracy", manifest), + collect_tasks=_applicable("runtime.task_declaration_accuracy", manifest), ) # Protocol collection must finish before its dependent contract checks run. if evidence.failure_reason: diff --git a/src/openenv/validation/runtime/artifacts.py b/src/openenv/validation/runtime/artifacts.py index 71a01c57ab..d0a69d3c54 100644 --- a/src/openenv/validation/runtime/artifacts.py +++ b/src/openenv/validation/runtime/artifacts.py @@ -135,6 +135,21 @@ def write_runtime_bundle( ) files["collector-trace.json"] = trace files["collector-evidence.json"] = collector_metadata + if evidence: + discovery = {} + omitted = [] + for name in ("tools", "tasks"): + raw = getattr(evidence, name + "_json") + discovery[name + "_available"] = raw is not None + discovery[name + "_error"] = getattr(evidence, name + "_error") + try: + discovery[name] = json.loads(raw) if raw is not None else None + except (ValueError, RecursionError): + discovery[name] = "[malformed evidence omitted]" + omitted.append(name) + discovery["omitted_fields"] = omitted + discovery["redacted"] = bool(omitted) or _redact(discovery) != discovery + files["discovery.json"] = discovery if evidence and (evidence.replays or evidence.replay_failure_reason): samples = [] for replay in evidence.replays: @@ -183,7 +198,12 @@ def write_runtime_bundle( "failure_reason": evidence.replay_failure_reason, "samples": samples, } - for name in ("runtime-plan.json", "session-telemetry.json", "replays.json"): + for name in ( + "runtime-plan.json", + "session-telemetry.json", + "replays.json", + "discovery.json", + ): if name not in files: (directory / name).unlink(missing_ok=True) digests = [] diff --git a/src/openenv/validation/runtime/collector.py b/src/openenv/validation/runtime/collector.py index 21b264fbd8..35d67084c5 100644 --- a/src/openenv/validation/runtime/collector.py +++ b/src/openenv/validation/runtime/collector.py @@ -11,6 +11,7 @@ from ...core.env_server.types import WSErrorCode from .contracts import RuntimeEvidence, RuntimePlan, WireExchange +from .discovery import collect_task_evidence from .transport import abort_socket, http_deadline MAX_MESSAGE_BYTES = 1024 * 1024 @@ -72,6 +73,8 @@ def collect_runtime_evidence( episode_timeout_s: float, request_timeout_s: float | None = None, validation_token: str | None = None, + collect_tools: bool = False, + collect_tasks: bool = False, ) -> RuntimeEvidence: """ Preserve schema and reset/step/state responses without model coercion. @@ -92,6 +95,10 @@ def collect_runtime_evidence( remaining declared 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. + 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. @@ -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 server_code = None def remaining() -> float: @@ -209,10 +217,10 @@ def receive_response(request): raise return connection.recv(timeout=remaining()) - def telemetry_request(operation, data): + def telemetry_request(operation, data, max_bytes=MAX_TRACE_BYTES): request = json.dumps({"type": operation, "data": data}) raw = receive_response(request) - 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): @@ -268,6 +276,61 @@ 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 collect_tasks: + phase = "tasks" + try: + discovered, tasks_error = collect_task_evidence( + base_url, + 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, server_code phase = operation @@ -372,6 +435,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, reset_observation_schema_json=reset_schema_json, ) except KeyboardInterrupt: @@ -384,6 +451,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: @@ -396,4 +467,8 @@ def exchange(operation: str, data: dict | None = None) -> dict: failure_reason=f"{phase} failed ({server_code or 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 bf4cf65bd0..6aa4edf788 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. reset_observation_schema_json (`str`, *optional*): Explicit `/schema` reset_observation value; absent means use the step schema. """ @@ -281,6 +289,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 reset_observation_schema_json: str | None = None diff --git a/src/openenv/validation/runtime/discovery.py b/src/openenv/validation/runtime/discovery.py new file mode 100644 index 0000000000..0b792e07bf --- /dev/null +++ b/src/openenv/validation/runtime/discovery.py @@ -0,0 +1,128 @@ +"""Bounded task metadata sampling through the production task API.""" + +import json +import time +from urllib.parse import quote + +import httpx + +from .transport import http_deadline + +MAX_TASK_SPLITS = 64 +MAX_RESPONSE_BYTES = 1024 * 1024 +MAX_DISCOVERY_BYTES = 8 * 1024 * 1024 + + +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. 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 = 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 + + try: + prefix = base_url.rstrip("/") + 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 + # 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") + .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) + + 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") + 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( + { + "environments": environments, + "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, + httpx.InvalidURL, + OSError, + ValueError, + TypeError, + KeyError, + RecursionError, + ) as error: + return None, f"task discovery failed ({type(error).__name__})" diff --git a/src/openenv/validation/runtime/replay.py b/src/openenv/validation/runtime/replay.py index 6cea193d59..88000b54ae 100644 --- a/src/openenv/validation/runtime/replay.py +++ b/src/openenv/validation/runtime/replay.py @@ -22,6 +22,8 @@ def _evidence_bytes(evidence): evidence.observation_schema_json or "", evidence.reset_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/fixtures/validation/runtime/served_probe/app.py b/tests/fixtures/validation/runtime/served_probe/app.py index fddd8dcca6..e90d4dfb4f 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 @@ -145,6 +186,15 @@ 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 + 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) @@ -169,6 +219,14 @@ 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", + "changed_rubric_config", "slow_step", }: raise ValueError(f"Unknown fixture mode: {mode}") diff --git a/tests/fixtures/validation/runtime/served_probe/openenv.yaml b/tests/fixtures/validation/runtime/served_probe/openenv.yaml index 8dcb0f534d..58314397a4 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 0000000000..a5b456b377 --- /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 1177226b9a..42addd9b61 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,51 @@ 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"] + 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 "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) 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 @@ -316,11 +345,25 @@ 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"), + ("changed_rubric_config", "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" @@ -330,6 +373,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", @@ -340,7 +384,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" @@ -469,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 27c924ff94..7d4a67320d 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,15 @@ 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"), + ("changed_rubric_config", "runtime.reward_attribution"), + ("empty_tools", None), + ("namespace_mismatch", None), ], ) def test_installed_server_collector_and_graders_over_loopback( @@ -177,10 +188,23 @@ def test_installed_server_collector_and_graders_over_loopback( / "process" / mode ) - provider = ProcessProvider(artifacts / "subject", mode) + target = FIXTURE + 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()) + 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", "good" if mode == "namespace_mismatch" else mode + ) try: report = run_validation( - FIXTURE, + target, max_level=Level.RUNTIME, provider=provider, artifacts_dir=artifacts / "report", @@ -200,6 +224,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 +236,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 + 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 new file mode 100644 index 0000000000..c095668ae2 --- /dev/null +++ b/tests/test_validation/test_discovery_collection.py @@ -0,0 +1,267 @@ +"""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 == "/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"): + 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, collect_tasks=True + ) + 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", + "/list_environments", + "/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, + collect_tasks=True, + 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, + collect_tasks=True, + 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_manifest.py b/tests/test_validation/test_manifest.py index 7f7dd36709..fb091e3fb2 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_runner.py b/tests/test_validation/test_runner.py index a484152534..9fde1daa15 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_artifacts.py b/tests/test_validation/test_runtime_artifacts.py index 84ac372e85..4f93f14677 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 @@ -386,6 +386,95 @@ def test_rewriting_bundle_removes_stale_optional_evidence(tmp_path): 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() + + def test_redaction_preserves_token_metadata_and_filters_secret_keys(): from openenv.validation.runtime.artifacts import _redact diff --git a/tests/test_validation/test_runtime_discovery.py b/tests/test_validation/test_runtime_discovery.py new file mode 100644 index 0000000000..6b17e4ab5b --- /dev/null +++ b/tests/test_validation/test_runtime_discovery.py @@ -0,0 +1,521 @@ +"""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 + + +@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 = [] + 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("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( + 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("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": + 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.model_fields_set.discard("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 fb9fc16aa6..77b37e4a15 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,28 @@ 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 + ) + document["validation"]["capabilities"].pop("declared_task_count", None) + 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 +128,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 ): @@ -351,7 +391,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) @@ -368,7 +408,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" @@ -383,7 +423,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] @@ -629,6 +677,184 @@ 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]["collect_tasks"] is False + 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]["collect_tasks"] is True + for name in ( + "tool_declaration_accuracy", + "task_declaration_accuracy", + "rubric_introspectable", + "reward_attribution", + ): + 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 +): + 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 + 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( + 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]["collect_tasks"] is True + + def grader(check_id, depends_on=(), *, status=CheckStatus.PASS, requires=frozenset()): return SimpleNamespace( check_id=check_id, diff --git a/tests/test_validation/test_runtime_replay.py b/tests/test_validation/test_runtime_replay.py index 24d4031ea3..45be9035e2 100644 --- a/tests/test_validation/test_runtime_replay.py +++ b/tests/test_validation/test_runtime_replay.py @@ -278,6 +278,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, ): 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 0000000000..0ebca06ae3 --- /dev/null +++ b/tests/test_validation/test_task_discovery_collection.py @@ -0,0 +1,278 @@ +import json +import threading +import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + +import httpx +import pytest +from openenv.validation.runtime import discovery + + +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["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) + return discovery.collect_task_evidence( + "http://127.0.0.1:8000", 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) == { + "environments": ["name?secret"], + "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( + "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", + [ + "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("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("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) + + 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 ef5dc249ae..4861358725 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`, `slow_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 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 seed, and in an independently inspected fresh container. 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 20 required cases: three provider lifecycle tests, -16 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 cd58914133..4445840ea5 100644 --- a/tests/validation_runtime/acceptance.json +++ b/tests/validation_runtime/acceptance.json @@ -32,7 +32,17 @@ "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[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]" ], "docker": [ "tests.test_validation.integration.test_docker_lifecycle::test_docker_lifecycle_effective_limits_and_owned_cleanup", @@ -54,7 +64,15 @@ "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[changed_rubric_config-runtime.reward_attribution]", + "tests.test_validation.integration.test_runtime_cli::test_cli_runtime_contract_findings[empty_tools-None]" ] } }