Heal part 2: split the cybergym benchmark under the module caps, regenerate the manifest

The drift landed four cybergym modules over the 1600-line cap, which the
ratchet manifest generator structurally refuses to except. Behavior-
preserving splits along real seams: the executor sheds its docker/
container machinery (cybergym_docker.py) and run-settle accounting
lifecycle (cybergym_lifecycle.py) as mixin layers assembled in the
original class; the adapter sheds its stateless protocol layer
(cybergym_protocol.py: constants, validators, provenance, argv builders);
the sidecar sheds observation parsing and connectivity evaluation
(cybergym_observations.py); the executor test suite sheds its wire-
protocol cluster (test_cybergym_executor_wire.py). Every name any other
file imports stays importable from its original module.

loop.py returns to the upstream byte-debt record via wording trims on
this branch's own additions (no pinned phrase touched). The manifest
regenerates cleanly with band rationales for every drift arrival; the
blocking CI ratchet lane (-m size_ratchet) is green against the drifted
base - it was red even on the pristine target tip.

All cybergym suites pass (protocol, server, executor, wire); ruff -F
clean; full canonical pytest re-run on this tree.
This commit is contained in:
Anton Razzhigaev 2026-08-29 12:23:29 +00:00
parent 4ace26d339
commit 515dcfa3f6
11 changed files with 5022 additions and 4726 deletions

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,336 @@
"""Sanitised observation parsing for the CyberGym sidecar contracts.
This module owns the seam that reads facts out of sanitised runner
observations: Docker ``inspect``-shaped payload readers and the
connectivity-probe evaluation, together with the attestation schema
identifier and the sha256-digest pattern shared with the argv/spec
validators in ``cybergym_sidecar``. It stays pure (stdlib only) so the
rules remain testable on CI hosts without Docker, and it must never
import from ``cybergym_sidecar`` (that module imports from here).
"""
from __future__ import annotations
import re
from typing import Any, Mapping, Sequence
SCHEMA_VERSION = "ouroboros.benchmark.cybergym.sidecar_attestation.v1"
_DIGEST = re.compile(r"^sha256:[0-9a-f]{64}$")
_CONNECTIVITY_EXPECTATIONS = {
"agent_to_server": True,
"verifier_to_private": True,
"agent_to_public": True,
"agent_to_verifier": False,
"agent_socket_visible": False,
}
def _probe_value(value: Any) -> bool | None:
if isinstance(value, bool):
return value
if isinstance(value, Mapping):
for key in ("reachable", "visible", "ok", "success", "passed"):
if isinstance(value.get(key), bool):
return value[key]
return None
_PROTECTED_ROUTE_KEY = "agent_to_server_protected"
# A wrong HTTP method (405) reaches Starlette routing before FastAPI's private
# API-key dependency, so it is not evidence that the route is protected. The
# production probe uses POST with an intentionally invalid JSON body and must
# observe an auth denial. Keep this set narrow rather than treating every
# client error (notably 405/422) as proof of authorization enforcement.
_PROTECTED_DENIAL_STATUS_RANGE = {401, 403, 404}
def _mapping_value(mapping: Mapping[str, Any], names: Sequence[str]) -> Any:
"""Read a small set of Docker/runner spellings without accepting aliases blindly."""
for name in names:
if name in mapping:
return mapping[name]
lowered = {name.lower() for name in names}
for key, value in mapping.items():
if isinstance(key, str) and key.lower() in lowered:
return value
return None
def _bool_mapping_value(mapping: Mapping[str, Any], names: Sequence[str]) -> bool | None:
value = _mapping_value(mapping, names)
return value if isinstance(value, bool) else None
def _http_status(mapping: Mapping[str, Any]) -> int | None:
value = _mapping_value(mapping, ("status_code", "http_status", "response_status", "status"))
if isinstance(value, bool) or not isinstance(value, int):
return None
return value if 100 <= value <= 599 else None
def _protected_route_record(value: Any) -> dict[str, Any] | None:
"""Normalize one non-mutating protected-route probe result.
A transport response is useful evidence only when the runner also records
that the unauthenticated request was denied and did not mutate state. A
bare boolean is deliberately rejected, because it cannot distinguish a
denied 401/405 from a successful mutating request.
"""
if not isinstance(value, Mapping):
return None
status = _http_status(value)
reachable = _bool_mapping_value(
value,
("reachable", "transport_reachable", "transport_ok", "transport"),
)
if reachable is None and status is not None:
reachable = True
authorized = _bool_mapping_value(value, ("authorized", "allowed", "accepted"))
denied = _bool_mapping_value(
value,
("denied", "unauthorized", "protected", "expected_denial", "denied_by_auth"),
)
if denied is None and authorized is not None:
denied = not authorized
if denied is None and status is not None:
denied = status in _PROTECTED_DENIAL_STATUS_RANGE
mutating = _bool_mapping_value(
value,
("mutating", "mutation", "mutation_succeeded", "mutation_success", "write_succeeded", "mutated"),
)
# An auth denial is an explicit non-mutating outcome for the malformed POST
# probe. A 2xx/5xx response still requires an explicit mutating=false fact.
if mutating is None and denied is True and status in _PROTECTED_DENIAL_STATUS_RANGE:
mutating = False
return {
"reachable": reachable,
"denied": denied,
"mutating": mutating,
"status_code": status,
}
def _protected_route_summary(value: Any) -> tuple[dict[str, Any], bool]:
"""Evaluate aggregate or per-target protected-route observations."""
per_target = False
raw_records: list[Any]
aggregate_defaults: dict[str, Any] = {}
if isinstance(value, Mapping):
aggregate_defaults = {
key: item
for key, item in value.items()
if key in {
"reachable",
"transport_reachable",
"transport_ok",
"transport",
"denied",
"unauthorized",
"protected",
"expected_denial",
"denied_by_auth",
"authorized",
"allowed",
"accepted",
"mutating",
"mutation",
"mutation_succeeded",
"mutation_success",
"write_succeeded",
"mutated",
"status_code",
"http_status",
"response_status",
"status",
}
}
nested = _mapping_value(value, ("targets", "target_results", "routes"))
if isinstance(nested, Mapping):
raw_records = list(nested.values())
per_target = True
elif isinstance(nested, Sequence) and not isinstance(nested, (str, bytes)):
if all(isinstance(item, Mapping) for item in nested):
raw_records = list(nested)
else:
# Some runners report only the target names and put the
# aggregate facts beside them. Preserve the requirement that
# both configured routes were probed while reusing those facts.
raw_records = [dict(aggregate_defaults) for _ in nested]
per_target = True
elif value and all(isinstance(item, Mapping) for item in value.values()):
# A URL-keyed mapping is a convenient runner representation.
raw_records = list(value.values())
per_target = True
else:
raw_records = [value]
elif isinstance(value, Sequence) and not isinstance(value, (str, bytes)):
raw_records = list(value)
per_target = True
else:
raw_records = []
records = [_protected_route_record(item) for item in raw_records]
if per_target and aggregate_defaults:
records = [
_protected_route_record({**aggregate_defaults, **item}) if isinstance(item, Mapping) else record
for item, record in zip(raw_records, records)
]
declared_count = value.get("target_count") if isinstance(value, Mapping) else None
declared_count_ok = declared_count is None or (
isinstance(declared_count, int) and not isinstance(declared_count, bool) and declared_count == len(records)
)
target_count_ok = bool(records) and declared_count_ok and (not per_target or len(records) >= 2)
pass_value = bool(
target_count_ok
and all(
record is not None
and record["reachable"] is True
and record["denied"] is True
and record["mutating"] is False
and record["status_code"] in _PROTECTED_DENIAL_STATUS_RANGE
for record in records
)
)
summary = {
"reachable": all(record is not None and record["reachable"] is True for record in records) if records else None,
"denied": all(record is not None and record["denied"] is True for record in records) if records else None,
"mutating": any(record is not None and record["mutating"] is True for record in records) if records else None,
"target_count": len(records),
"per_target": per_target,
"pass": pass_value,
}
return summary, pass_value
def evaluate_connectivity_checks(
observed: Mapping[str, Any],
*,
require_protected_route_evidence: bool = False,
) -> dict[str, Any]:
"""Evaluate route facts; optional production probes fail closed when present."""
checks: dict[str, Any] = {}
failed: list[str] = []
if not isinstance(observed, Mapping):
observed = {}
for name, expected in _CONNECTIVITY_EXPECTATIONS.items():
value = _probe_value(observed.get(name))
passed = value is expected
checks[name] = {"observed": value, "expected": expected, "pass": passed}
if not passed:
failed.append(name)
if _PROTECTED_ROUTE_KEY in observed or require_protected_route_evidence:
summary, passed = _protected_route_summary(observed.get(_PROTECTED_ROUTE_KEY))
checks[_PROTECTED_ROUTE_KEY] = {
"observed": summary,
"expected": {"reachable": True, "denied": True, "mutating": False},
"pass": passed,
}
if not passed:
failed.append(_PROTECTED_ROUTE_KEY)
return {"schema": f"{SCHEMA_VERSION}.connectivity", "ok": not failed, "checks": checks, "failed": failed}
def _nested(mapping: Mapping[str, Any], *keys: str) -> Any:
value: Any = mapping
for key in keys:
if not isinstance(value, Mapping):
return None
value = value.get(key)
return value
def _name(observation: Mapping[str, Any]) -> str | None:
value = observation.get("Name") or observation.get("name")
return value.lstrip("/") if isinstance(value, str) and value.lstrip("/") else None
def _id(observation: Mapping[str, Any]) -> str | None:
value = observation.get("Id") or observation.get("ID") or observation.get("id")
return value if isinstance(value, str) and value else None
def _pid(observation: Mapping[str, Any]) -> int | None:
value = _nested(observation, "State", "Pid")
if value is None:
value = observation.get("Pid") or observation.get("pid")
try:
value = int(value)
except (TypeError, ValueError):
return None
return value if value > 0 else None
def _running(observation: Mapping[str, Any]) -> bool:
return _nested(observation, "State", "Running") is True or observation.get("running") is True
def _labels_from(observation: Mapping[str, Any]) -> Mapping[str, Any]:
value = _nested(observation, "Config", "Labels")
if not isinstance(value, Mapping):
value = observation.get("Labels")
return value if isinstance(value, Mapping) else {}
def _network_mode(observation: Mapping[str, Any]) -> str | None:
value = _nested(observation, "HostConfig", "NetworkMode")
if value is None:
value = observation.get("NetworkMode") or observation.get("network_mode")
return value if isinstance(value, str) else None
def _network(observation: Mapping[str, Any], name: str) -> Mapping[str, Any] | None:
values = _nested(observation, "NetworkSettings", "Networks")
if not isinstance(values, Mapping):
values = observation.get("Networks")
value = values.get(name) if isinstance(values, Mapping) else None
return value if isinstance(value, Mapping) else None
def _mounts(observation: Mapping[str, Any]) -> list[tuple[str, str]]:
result: list[tuple[str, str]] = []
values = observation.get("Mounts")
if isinstance(values, Sequence) and not isinstance(values, (str, bytes)):
for item in values:
if isinstance(item, Mapping):
source = item.get("Source") or item.get("source")
destination = item.get("Destination") or item.get("destination") or item.get("Target")
if isinstance(source, str) and isinstance(destination, str):
result.append((source, destination))
values = _nested(observation, "HostConfig", "Binds")
if isinstance(values, Sequence) and not isinstance(values, (str, bytes)):
for item in values:
if isinstance(item, str) and ":" in item:
source, destination = item.split(":", 1)[:2]
result.append((source, destination))
return result
def _bindings(observation: Mapping[str, Any], port: int) -> list[Mapping[str, Any]]:
values = _nested(observation, "NetworkSettings", "Ports")
if not isinstance(values, Mapping):
values = observation.get("Ports")
value = values.get(f"{port}/tcp") if isinstance(values, Mapping) else None
return list(value) if isinstance(value, Sequence) and not isinstance(value, (str, bytes)) else []
def _digests(observation: Mapping[str, Any]) -> set[str]:
values = _nested(observation, "Config", "RepoDigests")
if not isinstance(values, Sequence) or isinstance(values, (str, bytes)):
values = observation.get("RepoDigests")
result: set[str] = set()
for item in values or ():
if isinstance(item, str) and "@" in item:
result.add(item.split("@", 1)[1])
image = observation.get("Image")
if isinstance(image, str) and _DIGEST.fullmatch(image):
result.add(image)
config_image = _nested(observation, "Config", "Image")
if isinstance(config_image, str) and _DIGEST.fullmatch(config_image):
result.add(config_image)
return result

File diff suppressed because it is too large Load diff

View file

@ -17,7 +17,32 @@ from dataclasses import dataclass, field
from typing import Any, Iterable, Mapping, Protocol, Sequence
from urllib.parse import urlsplit
SCHEMA_VERSION = "ouroboros.benchmark.cybergym.sidecar_attestation.v1"
from devtools.benchmarks.cybergym.cybergym_observations import (
SCHEMA_VERSION, # noqa: F401
_CONNECTIVITY_EXPECTATIONS, # noqa: F401
_DIGEST, # noqa: F401
_PROTECTED_DENIAL_STATUS_RANGE, # noqa: F401
_PROTECTED_ROUTE_KEY, # noqa: F401
_bindings, # noqa: F401
_bool_mapping_value, # noqa: F401
_digests, # noqa: F401
_http_status, # noqa: F401
_id, # noqa: F401
_labels_from, # noqa: F401
_mapping_value, # noqa: F401
_mounts, # noqa: F401
_name, # noqa: F401
_nested, # noqa: F401
_network, # noqa: F401
_network_mode, # noqa: F401
_pid, # noqa: F401
_probe_value, # noqa: F401
_protected_route_record, # noqa: F401
_protected_route_summary, # noqa: F401
_running, # noqa: F401
evaluate_connectivity_checks, # noqa: F401
)
NETWORK_NAME = "cybergym-internal"
DOCKER_HOST_ENV = "DOCKER_HOST"
LOOPBACK_HOST = "127.0.0.1"
@ -31,7 +56,6 @@ _SAFE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,127}$")
_SAFE_NAME = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.-]{0,127}$")
_SAFE_DNS = re.compile(r"^[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$")
_ENV_NAME = re.compile(r"^[A-Za-z_][A-Za-z0-9_]{0,127}$")
_DIGEST = re.compile(r"^sha256:[0-9a-f]{64}$")
# The upstream CyberGym distribution documents a public example key. Keep
# only its non-reversible digest prefix here; the key itself must never enter
@ -386,34 +410,6 @@ def build_connectivity_probe_plan(plan: NetworkPlan) -> tuple[dict[str, Any], ..
)
_CONNECTIVITY_EXPECTATIONS = {
"agent_to_server": True,
"verifier_to_private": True,
"agent_to_public": True,
"agent_to_verifier": False,
"agent_socket_visible": False,
}
def _probe_value(value: Any) -> bool | None:
if isinstance(value, bool):
return value
if isinstance(value, Mapping):
for key in ("reachable", "visible", "ok", "success", "passed"):
if isinstance(value.get(key), bool):
return value[key]
return None
_PROTECTED_ROUTE_KEY = "agent_to_server_protected"
# A wrong HTTP method (405) reaches Starlette routing before FastAPI's private
# API-key dependency, so it is not evidence that the route is protected. The
# production probe uses POST with an intentionally invalid JSON body and must
# observe an auth denial. Keep this set narrow rather than treating every
# client error (notably 405/422) as proof of authorization enforcement.
_PROTECTED_DENIAL_STATUS_RANGE = {401, 403, 404}
def _is_rootless_security_option(value: Any) -> bool:
if not isinstance(value, str):
return False
@ -421,195 +417,6 @@ def _is_rootless_security_option(value: Any) -> bool:
return option in {"rootless", "name=rootless"} or option.endswith("=rootless")
def _mapping_value(mapping: Mapping[str, Any], names: Sequence[str]) -> Any:
"""Read a small set of Docker/runner spellings without accepting aliases blindly."""
for name in names:
if name in mapping:
return mapping[name]
lowered = {name.lower() for name in names}
for key, value in mapping.items():
if isinstance(key, str) and key.lower() in lowered:
return value
return None
def _bool_mapping_value(mapping: Mapping[str, Any], names: Sequence[str]) -> bool | None:
value = _mapping_value(mapping, names)
return value if isinstance(value, bool) else None
def _http_status(mapping: Mapping[str, Any]) -> int | None:
value = _mapping_value(mapping, ("status_code", "http_status", "response_status", "status"))
if isinstance(value, bool) or not isinstance(value, int):
return None
return value if 100 <= value <= 599 else None
def _protected_route_record(value: Any) -> dict[str, Any] | None:
"""Normalize one non-mutating protected-route probe result.
A transport response is useful evidence only when the runner also records
that the unauthenticated request was denied and did not mutate state. A
bare boolean is deliberately rejected, because it cannot distinguish a
denied 401/405 from a successful mutating request.
"""
if not isinstance(value, Mapping):
return None
status = _http_status(value)
reachable = _bool_mapping_value(
value,
("reachable", "transport_reachable", "transport_ok", "transport"),
)
if reachable is None and status is not None:
reachable = True
authorized = _bool_mapping_value(value, ("authorized", "allowed", "accepted"))
denied = _bool_mapping_value(
value,
("denied", "unauthorized", "protected", "expected_denial", "denied_by_auth"),
)
if denied is None and authorized is not None:
denied = not authorized
if denied is None and status is not None:
denied = status in _PROTECTED_DENIAL_STATUS_RANGE
mutating = _bool_mapping_value(
value,
("mutating", "mutation", "mutation_succeeded", "mutation_success", "write_succeeded", "mutated"),
)
# An auth denial is an explicit non-mutating outcome for the malformed POST
# probe. A 2xx/5xx response still requires an explicit mutating=false fact.
if mutating is None and denied is True and status in _PROTECTED_DENIAL_STATUS_RANGE:
mutating = False
return {
"reachable": reachable,
"denied": denied,
"mutating": mutating,
"status_code": status,
}
def _protected_route_summary(value: Any) -> tuple[dict[str, Any], bool]:
"""Evaluate aggregate or per-target protected-route observations."""
per_target = False
raw_records: list[Any]
aggregate_defaults: dict[str, Any] = {}
if isinstance(value, Mapping):
aggregate_defaults = {
key: item
for key, item in value.items()
if key in {
"reachable",
"transport_reachable",
"transport_ok",
"transport",
"denied",
"unauthorized",
"protected",
"expected_denial",
"denied_by_auth",
"authorized",
"allowed",
"accepted",
"mutating",
"mutation",
"mutation_succeeded",
"mutation_success",
"write_succeeded",
"mutated",
"status_code",
"http_status",
"response_status",
"status",
}
}
nested = _mapping_value(value, ("targets", "target_results", "routes"))
if isinstance(nested, Mapping):
raw_records = list(nested.values())
per_target = True
elif isinstance(nested, Sequence) and not isinstance(nested, (str, bytes)):
if all(isinstance(item, Mapping) for item in nested):
raw_records = list(nested)
else:
# Some runners report only the target names and put the
# aggregate facts beside them. Preserve the requirement that
# both configured routes were probed while reusing those facts.
raw_records = [dict(aggregate_defaults) for _ in nested]
per_target = True
elif value and all(isinstance(item, Mapping) for item in value.values()):
# A URL-keyed mapping is a convenient runner representation.
raw_records = list(value.values())
per_target = True
else:
raw_records = [value]
elif isinstance(value, Sequence) and not isinstance(value, (str, bytes)):
raw_records = list(value)
per_target = True
else:
raw_records = []
records = [_protected_route_record(item) for item in raw_records]
if per_target and aggregate_defaults:
records = [
_protected_route_record({**aggregate_defaults, **item}) if isinstance(item, Mapping) else record
for item, record in zip(raw_records, records)
]
declared_count = value.get("target_count") if isinstance(value, Mapping) else None
declared_count_ok = declared_count is None or (
isinstance(declared_count, int) and not isinstance(declared_count, bool) and declared_count == len(records)
)
target_count_ok = bool(records) and declared_count_ok and (not per_target or len(records) >= 2)
pass_value = bool(
target_count_ok
and all(
record is not None
and record["reachable"] is True
and record["denied"] is True
and record["mutating"] is False
and record["status_code"] in _PROTECTED_DENIAL_STATUS_RANGE
for record in records
)
)
summary = {
"reachable": all(record is not None and record["reachable"] is True for record in records) if records else None,
"denied": all(record is not None and record["denied"] is True for record in records) if records else None,
"mutating": any(record is not None and record["mutating"] is True for record in records) if records else None,
"target_count": len(records),
"per_target": per_target,
"pass": pass_value,
}
return summary, pass_value
def evaluate_connectivity_checks(
observed: Mapping[str, Any],
*,
require_protected_route_evidence: bool = False,
) -> dict[str, Any]:
"""Evaluate route facts; optional production probes fail closed when present."""
checks: dict[str, Any] = {}
failed: list[str] = []
if not isinstance(observed, Mapping):
observed = {}
for name, expected in _CONNECTIVITY_EXPECTATIONS.items():
value = _probe_value(observed.get(name))
passed = value is expected
checks[name] = {"observed": value, "expected": expected, "pass": passed}
if not passed:
failed.append(name)
if _PROTECTED_ROUTE_KEY in observed or require_protected_route_evidence:
summary, passed = _protected_route_summary(observed.get(_PROTECTED_ROUTE_KEY))
checks[_PROTECTED_ROUTE_KEY] = {
"observed": summary,
"expected": {"reachable": True, "denied": True, "mutating": False},
"pass": passed,
}
if not passed:
failed.append(_PROTECTED_ROUTE_KEY)
return {"schema": f"{SCHEMA_VERSION}.connectivity", "ok": not failed, "checks": checks, "failed": failed}
def _validate_image(value: str) -> str:
value = _text(value, "image", max_len=1024)
if value.lower() in {"latest", "placeholder", "changeme", "example", "example/image"} or "${" in value:
@ -1014,106 +821,6 @@ class SidecarExpectation:
raise SidecarConfigurationError("publish_host_port must be a boolean")
def _nested(mapping: Mapping[str, Any], *keys: str) -> Any:
value: Any = mapping
for key in keys:
if not isinstance(value, Mapping):
return None
value = value.get(key)
return value
def _name(observation: Mapping[str, Any]) -> str | None:
value = observation.get("Name") or observation.get("name")
return value.lstrip("/") if isinstance(value, str) and value.lstrip("/") else None
def _id(observation: Mapping[str, Any]) -> str | None:
value = observation.get("Id") or observation.get("ID") or observation.get("id")
return value if isinstance(value, str) and value else None
def _pid(observation: Mapping[str, Any]) -> int | None:
value = _nested(observation, "State", "Pid")
if value is None:
value = observation.get("Pid") or observation.get("pid")
try:
value = int(value)
except (TypeError, ValueError):
return None
return value if value > 0 else None
def _running(observation: Mapping[str, Any]) -> bool:
return _nested(observation, "State", "Running") is True or observation.get("running") is True
def _labels_from(observation: Mapping[str, Any]) -> Mapping[str, Any]:
value = _nested(observation, "Config", "Labels")
if not isinstance(value, Mapping):
value = observation.get("Labels")
return value if isinstance(value, Mapping) else {}
def _network_mode(observation: Mapping[str, Any]) -> str | None:
value = _nested(observation, "HostConfig", "NetworkMode")
if value is None:
value = observation.get("NetworkMode") or observation.get("network_mode")
return value if isinstance(value, str) else None
def _network(observation: Mapping[str, Any], name: str) -> Mapping[str, Any] | None:
values = _nested(observation, "NetworkSettings", "Networks")
if not isinstance(values, Mapping):
values = observation.get("Networks")
value = values.get(name) if isinstance(values, Mapping) else None
return value if isinstance(value, Mapping) else None
def _mounts(observation: Mapping[str, Any]) -> list[tuple[str, str]]:
result: list[tuple[str, str]] = []
values = observation.get("Mounts")
if isinstance(values, Sequence) and not isinstance(values, (str, bytes)):
for item in values:
if isinstance(item, Mapping):
source = item.get("Source") or item.get("source")
destination = item.get("Destination") or item.get("destination") or item.get("Target")
if isinstance(source, str) and isinstance(destination, str):
result.append((source, destination))
values = _nested(observation, "HostConfig", "Binds")
if isinstance(values, Sequence) and not isinstance(values, (str, bytes)):
for item in values:
if isinstance(item, str) and ":" in item:
source, destination = item.split(":", 1)[:2]
result.append((source, destination))
return result
def _bindings(observation: Mapping[str, Any], port: int) -> list[Mapping[str, Any]]:
values = _nested(observation, "NetworkSettings", "Ports")
if not isinstance(values, Mapping):
values = observation.get("Ports")
value = values.get(f"{port}/tcp") if isinstance(values, Mapping) else None
return list(value) if isinstance(value, Sequence) and not isinstance(value, (str, bytes)) else []
def _digests(observation: Mapping[str, Any]) -> set[str]:
values = _nested(observation, "Config", "RepoDigests")
if not isinstance(values, Sequence) or isinstance(values, (str, bytes)):
values = observation.get("RepoDigests")
result: set[str] = set()
for item in values or ():
if isinstance(item, str) and "@" in item:
result.add(item.split("@", 1)[1])
image = observation.get("Image")
if isinstance(image, str) and _DIGEST.fullmatch(image):
result.add(image)
config_image = _nested(observation, "Config", "Image")
if isinstance(config_image, str) and _DIGEST.fullmatch(config_image):
result.add(config_image)
return result
def _container_report(
observation: Mapping[str, Any],
expected: SidecarExpectation,

View file

@ -2829,10 +2829,9 @@ def _maybe_inject_nanny_economics_reminder(
# says so instead of implying an activity that never happened.
_baseline_known = isinstance(getattr(ctx, "_nanny_delegate_baseline", None), dict)
if _baseline_known:
# Supervision/coordination rounds COUNT toward this burn (charter).
since_phrase = (
"since your last act of delegation (delegate_start / "
"schedule_subagent); supervision rounds count toward it"
"schedule_subagent), supervision included"
)
else:
since_phrase = "since this task started (no act of delegation yet)"
@ -5711,7 +5710,7 @@ def _nanny_finalization_message(
return (
"⚠️ NANNY_METERED_OVERRUN: your delegated run(s) succeeded, but you have "
f"since spent {_nanny_burn_phrase(rounds, cost)} beyond your last act of "
"delegation (supervision rounds count toward it). "
"delegation (supervision included). "
"A successful run is verified and integrated, not rebuilt. If "
"the remaining work is substantive, delegate it (a new delegate_start); "
"if you are wrapping up, keep the wrap-up short and account for the "
@ -5719,11 +5718,7 @@ def _nanny_finalization_message(
)
started = int(evidence.get("delegated_runs_started") or 0)
if not started and (evidence.get("evidence_read_failed") or not evidence):
# Zero attempts is an ACCUSATION and needs positively-established
# evidence: an unreadable custody log (or a failed read above) proves
# nothing (scope finding on a5e59bdf). A CONFIGURED actor still gets
# its own non-accusatory typed message (unknown over unreadable
# custody — never a fresh-start recipe); legacy nannies stay silent.
# Unreadable custody (a5e59bdf): configured -> unknown message; legacy -> silence.
from ouroboros.subagent_bootstrap import configured_actor_finalization_message
_actor = configured_actor_finalization_message(

View file

@ -124,6 +124,13 @@ BAND_BASELINE_PATHS = (
)
BAND_PATHS = {
"devtools/benchmarks/cybergym/cybergym_adapter.py": "Stateful campaign layer after the protocol split (ratchet heal); shrink next touch.",
"devtools/benchmarks/cybergym/cybergym_docker.py": "Docker runtime layer of the executor split: one container-machinery seam.",
"devtools/benchmarks/cybergym/cybergym_executor.py": "Executor assembly after docker/lifecycle/wire splits (ratchet heal); shrink next touch.",
"devtools/benchmarks/cybergym/cybergym_lifecycle.py": "Run/settle lifecycle layer of the executor split: one accounting seam.",
"devtools/benchmarks/cybergym/cybergym_protocol.py": "Stateless protocol layer of the adapter split: constants, validators, provenance.",
"devtools/benchmarks/cybergym/cybergym_sidecar.py": "Sidecar attestation core after the observations split (ratchet heal); shrink next touch.",
"devtools/benchmarks/cybergym/run_cybergym.py": "CyberGym launcher is one submit-shaped entry point (drift heal); split when a second arm lands.",
"devtools/benchmarks/swe_bench_pro/e1v2/run_pro.py": None,
"devtools/benchmarks/terminal_bench/harbor_installed_agent.py": None,
"devtools/benchmarks/terminal_bench/run_tb.py": None,
@ -169,6 +176,7 @@ BAND_PATHS = {
"tests/test_claude_code_gateway.py": None,
"tests/test_commit_gate.py": None,
"tests/test_contracts.py": None,
"tests/test_cybergym_protocol.py": "CyberGym protocol suite arrived in one piece with the benchmark (drift heal); split when the next protocol family lands.",
"tests/test_delegated_skill_payload.py": "Sol scope-review fix batch: P1 trust probes (forged index, symlinked git metadata), P2 golden-E2E review close and schema/docs pins joined the existing R1+gate-fix payload suite.",
"tests/test_evolution_redesign.py": None,
"tests/test_external_workspace_access.py": "Deliverables custody matrix covers nested roots, executor aliases, Presence ceilings, and direct shell target semantics.",
@ -214,7 +222,7 @@ BYTE_BASELINE_DEBT = {
BYTE_DEBT = {
"ouroboros/loop.py": 312847,
"tests/test_delegated_subagent_transport.py": 320571,
"tests/test_devtools_benchmarks.py": 328282,
"web/modules/chat.js": 224689,
"tests/test_delegated_subagent_transport.py": 320568,
"tests/test_devtools_benchmarks.py": 328116,
"web/modules/chat.js": 224549,
}

View file

@ -7,7 +7,6 @@ the live smoke remains an operator action documented by the benchmark.
from __future__ import annotations
import gzip
import hashlib
import io
import json
@ -20,12 +19,7 @@ import threading
import pytest
from devtools.benchmarks.cybergym import cybergym_executor as executor_module
from devtools.benchmarks.cybergym.cybergym_adapter import (
CAPABILITY_FINAL_POC_MISSING,
BudgetLedger,
final_poc_record,
run_campaign,
)
from devtools.benchmarks.cybergym.cybergym_adapter import final_poc_record
from devtools.benchmarks.cybergym.cybergym_executor import (
CommandResult,
CyberGymExecutor,
@ -33,12 +27,8 @@ from devtools.benchmarks.cybergym.cybergym_executor import (
ExecutorFailure,
_bind_container_image,
_install_workspace_backend_alias,
_parse_json_stdout,
_require_exact_effort,
_reuse_directory_observation,
_safe_extract,
_served_telemetry,
_validate_verify_response,
)
from devtools.benchmarks.cybergym.cybergym_sidecar import required_resource_labels
from ouroboros.headless import write_workspace_patch_artifacts
@ -939,54 +929,6 @@ def test_task_body_is_opaque_and_preserves_network_contract(tmp_path):
assert str(task_dir) not in guidance
def test_provider_probe_checks_exact_model_without_server_search(monkeypatch, tmp_path):
monkeypatch.setenv("OPENROUTER_API_KEY", "test-openrouter-key")
captured = {}
def http(method, url, *, body=None, headers=None, timeout=None):
if method == "GET" and url.endswith("/models"):
return {
"data": [{
"id": "deepseek/deepseek-v4-flash-0731",
"context_length": 1_310_720,
"supported_parameters": ["reasoning", "tools"],
}]
}
if method == "GET" and url.endswith("/key"):
return {"data": {"limit_remaining": 100}}
assert method == "POST"
captured["body"] = body
return {
"id": "response-1",
"model": "deepseek/deepseek-v4-flash-0731",
"provider": "OpenInference",
"choices": [{"message": {"content": "OK"}}],
"usage": {
"prompt_tokens": 12,
"completion_tokens": 3,
"cost": 0.006,
"cost_estimated": False,
},
}
executor = CyberGymExecutor(
_config(
tmp_path,
provider_probe=True,
expected_data_sha256="a" * 64,
expected_binary_sha256="b" * 64,
http_runner=http,
)
)
executor._probe_provider() # noqa: SLF001 - provider boundary assertion
assert captured["body"]["messages"] == [{"role": "user", "content": "Reply with OK."}]
assert "tools" not in captured["body"]
assert executor.provider_observation["observed_model"] == (
"deepseek/deepseek-v4-flash-0731"
)
def test_reused_directory_observation_rechecks_small_manifest_only(tmp_path):
payload_root = tmp_path / "payload"
payload_root.mkdir()
@ -1114,150 +1056,6 @@ def test_readiness_rejects_openapi_without_private_submit_fix(tmp_path, monkeypa
executor.start()
def test_verify_response_requires_success_body_and_designated_poc():
good = {"message": "All 1 PoCs for this agent_id have been verified", "poc_ids": ["poc-1"]}
assert _validate_verify_response(good, expected_poc_id="poc-1") == good
with pytest.raises(ExecutorFailure, match="HTTP 500"):
_validate_verify_response({"status_code": 500, "body": {"detail": "failed"}})
with pytest.raises(ExecutorFailure, match="poc_ids"):
_validate_verify_response({"message": "ok", "poc_ids": []})
with pytest.raises(ExecutorFailure, match="designated poc_id"):
_validate_verify_response(good, expected_poc_id="other")
def test_observed_effort_must_be_exactly_high():
assert _require_exact_effort("high") == "high"
for value in ("", "High", "max", None):
with pytest.raises(ExecutorFailure, match="exactly high"):
_require_exact_effort(value)
def test_served_telemetry_prefers_authoritative_trace_refs_over_requested_fields():
payload = {
"model": "requested/not-served",
"reasoning_effort": "high",
"trace_refs": {
"llm_call_refs": [
{"resolved_model": "deepseek/deepseek-v4-flash-0731", "provider": "provider-a"}
]
},
}
observed = _served_telemetry(payload)
assert observed["observed_model"] == "deepseek/deepseek-v4-flash-0731"
assert observed["observed_provider"] == "provider-a"
assert observed["trace_call_count"] == 1
assert observed["effort_source"] == "runtime_requested_field"
def test_served_telemetry_rejects_incomplete_or_mixed_trace_identity():
with pytest.raises(ExecutorFailure, match="incomplete served-call"):
_served_telemetry({"trace_refs": {"llm_call_refs": [{"provider": "provider-a"}]}})
with pytest.raises(ExecutorFailure, match="mixed served models"):
_served_telemetry(
{
"trace_refs": {
"llm_call_refs": [
{"resolved_model": "model-a", "provider": "provider-a"},
{"resolved_model": "model-b", "provider": "provider-a"},
]
}
}
)
def test_served_telemetry_reads_verified_response_wire_effort(tmp_path):
drive = tmp_path / "drive"
calls = drive / "observability" / "calls" / "opaque"
calls.mkdir(parents=True)
wire = {
"requested_effort": "high",
"applied_effort": "high",
"attempt_id": "attempt-1",
"candidate_sha256": "a" * 64,
}
blob_raw = json.dumps(
{"usage": {"request_wire": wire, "response_provider": "backend-a"}},
sort_keys=True,
).encode("utf-8")
blob_path = drive / "observability" / "blobs" / ("b" * 64 + ".json.gz")
blob_path.parent.mkdir(parents=True)
blob_path.write_bytes(gzip.compress(blob_raw))
blob_ref = {
"path": str(blob_path),
"sha256": hashlib.sha256(blob_raw).hexdigest(),
"size": len(blob_raw),
"kind": "json",
"encoding": "gzip",
}
manifest_raw = json.dumps(
{
"task_id": "opaque",
"call_id": "llm-1_response",
"llm_call_id": "llm-1",
"full_payload_ref": blob_ref,
},
sort_keys=True,
).encode("utf-8")
manifest_path = calls / "llm-1_response.json"
manifest_path.write_bytes(manifest_raw)
manifest_ref = {
"path": str(manifest_path),
"sha256": hashlib.sha256(manifest_raw).hexdigest(),
"call_id": "llm-1_response",
}
observed = _served_telemetry(
{
"reasoning_effort": "low",
"trace_refs": {
"llm_call_refs": [
{
"llm_call_id": "llm-1",
"resolved_model": "deepseek/deepseek-v4-flash-0731",
"provider": "provider-a",
"response_ref": manifest_ref,
}
]
},
},
allowed_roots=(drive,),
)
assert observed["observed_effort"] == "high"
assert observed["observed_provider"] == "backend-a"
assert observed["observed_provider_attempts"] == ["backend-a"]
assert observed["provider_distribution"] == {"backend-a": 1}
assert observed["effort_source"] == "served_response_wire"
assert observed["response_wire_effort_count"] == 1
assert observed["response_wire_provider_count"] == 1
def test_submit_stdout_parser_accepts_preceding_prose_and_multiline_json():
parsed = _parse_json_stdout('notice\n{\n "task_id": "opaque1234",\n "poc_id": "poc-1"\n}\n')
assert parsed == {"task_id": "opaque1234", "poc_id": "poc-1"}
def test_private_query_rejects_http_and_body_errors(tmp_path, monkeypatch):
config = _config(tmp_path)
monkeypatch.setenv("CYBERGYM_API_KEY", "test-secret-value")
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"status_code": 404, "body": {"detail": "Record not found"}},
)
)
with pytest.raises(ExecutorFailure, match="HTTP 404"):
executor._private_query("agent-" + "a" * 24, "arvo:1")
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"status_code": 200, "body": {"error": {"message": "bad"}}},
)
)
with pytest.raises(ExecutorFailure, match="error object"):
executor._private_query("agent-" + "a" * 24, "arvo:1")
def test_runtime_attestation_reinspects_immutable_ids_before_gateway_boundary(tmp_path, monkeypatch):
config = _config(tmp_path)
executor = CyberGymExecutor(config)
@ -1697,560 +1495,6 @@ def test_settled_workspace_cleanup_uses_exact_id_and_postcondition(tmp_path, mon
assert name not in executor._task_containers
def test_private_query_accepts_nested_items_wrapper(tmp_path, monkeypatch):
config = _config(tmp_path)
monkeypatch.setenv("CYBERGYM_API_KEY", "test-secret-value")
record = {"task_id": "arvo:1", "poc_id": "poc-1", "poc_hash": "a" * 64}
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"pocs": {"items": [record]}},
)
)
assert executor._private_query("agent-" + "a" * 24, "arvo:1") == [record]
def test_submit_response_binds_poc_id_not_nonexistent_hash_and_keeps_exit_code(tmp_path):
config = _config(tmp_path)
task_dir = config.run_root / "task"
task_dir.mkdir()
(task_dir / "final.poc").write_bytes(b"poc-bytes")
(task_dir / "submit.sh").write_text("TASK_ID=opaque1234\n", encoding="utf-8")
executor_name = "workspace"
executor_id = "c" * 64
submit_calls = []
def command(argv, *, cwd=None, env=None, timeout=None):
submit_calls.append(list(argv))
return CommandResult(
0,
json.dumps(
{
"task_id": "opaque1234",
"poc_id": "poc-1",
"exit_code": 71,
"output": "known",
# Upstream does not define a response hash; an incidental
# field must not override the local marker binding.
"hash": "not-the-poc-hash",
}
),
"",
)
executor = CyberGymExecutor(dataclasses_replace(config, command_runner=command))
executor._task_containers[executor_name] = executor_id
response, digest, masked = executor._submit_final( # noqa: SLF001 - boundary contract assertion
type("Task", (), {"task_id": "arvo:1"})(), task_dir, "workspace"
)
assert response["poc_id"] == "poc-1"
assert response["exit_code"] == 71
assert digest == hashlib.sha256(b"poc-bytes").hexdigest()
assert masked == "opaque1234"
assert executor_id in submit_calls[0]
assert executor_name not in submit_calls[0]
def test_unknown_gateway_attempt_blocks_campaign_cleanup(tmp_path):
config = _config(tmp_path)
calls = []
def command(*args, **kwargs):
calls.append(args)
raise AssertionError("cleanup must not run while gateway custody is unknown")
executor = CyberGymExecutor(dataclasses_replace(config, command_runner=command))
executor.started = True
executor.server_id = "server-123"
executor.network_id = "network-123"
executor._task_containers = {"workspace-agent-aaaaaaaaaaaaaaaaaaaaaaaa": "workspace-123"}
executor._gateway_attempts = {
"cybergym-attempt": {
"gateway_task_id": "cybergym-attempt",
"status": "admission_unknown",
"checkpoint": str(config.run_root / "checkpoint.json"),
}
}
report = executor.close()
assert report["ok"] is False
assert report["status"] == "custody_pending"
assert executor.custody_blocked is True
assert executor.server_id == "server-123"
assert executor.network_id == "network-123"
assert calls == []
assert (config.run_root / "custody_pending.json").is_file()
def test_gateway_admission_transport_error_registers_durable_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
seen = {}
def failing_http(*args, **kwargs):
seen.update(kwargs)
raise ExecutorFailure("HTTP POST transport failed")
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=failing_http))
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-opaque-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="transport failed"):
executor._gateway_wait(body, checkpoint)
assert "cybergym-opaque-attempt" in executor._gateway_attempts
assert executor._gateway_attempts["cybergym-opaque-attempt"]["status"] == "admission_unknown"
assert seen["headers"]["Idempotency-Key"].startswith("cybergym-")
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["custody_required"] is True
assert saved["status"] == "admission_unknown"
def test_gateway_definitive_admission_rejection_releases_phantom_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {
"status_code": 400,
"body": {"detail": "invalid task"},
},
)
)
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-rejected-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="HTTP 400"):
executor._gateway_wait(body, checkpoint)
assert executor._gateway_attempts == {}
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "admission_rejected"
assert saved["custody_required"] is False
def test_gateway_malformed_admission_keeps_unknown_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=lambda *args, **kwargs: {})
)
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-malformed-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="no task id"):
executor._gateway_wait(body, checkpoint)
assert "cybergym-malformed-attempt" in executor._gateway_attempts
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "admission_unknown_response"
assert saved["custody_required"] is True
def test_gateway_waits_for_final_cost_after_completed_status(tmp_path):
config = _config(tmp_path, provider_probe=False, task_timeout_sec=10)
task_id = "cybergym-cost-pending"
calls = []
status_rows = iter(
(
{
"task_id": task_id,
"status": "completed",
"result": {"cost_final": False},
},
{
"task_id": task_id,
"status": "completed",
"result": {"cost_final": True},
},
)
)
def http(method, url, **kwargs):
calls.append(method)
if method == "POST":
return {"task_id": task_id, "status": "scheduled"}
return next(status_rows)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None)
)
result = executor._gateway_wait(
{"task_id": task_id, "description": "test"},
config.run_root / "checkpoint.json",
)
assert result["result"]["cost_final"] is True
assert calls == ["POST", "GET", "GET"]
def test_gateway_cost_finality_conflict_keeps_polling(tmp_path):
config = _config(tmp_path, provider_probe=False, task_timeout_sec=10)
task_id = "cybergym-cost-conflict"
calls = []
status_rows = iter(
(
{
"task_id": task_id,
"status": "completed",
"cost_final": True,
"cost_breakdown": {"cost_final": False},
},
{
"task_id": task_id,
"status": "completed",
"cost_final": True,
"cost_breakdown": {"cost_final": True},
},
)
)
def http(method, _url, **_kwargs):
calls.append(method)
if method == "POST":
return {"task_id": task_id, "status": "scheduled"}
return next(status_rows)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None)
)
result = executor._gateway_wait( # noqa: SLF001 - accounting contract
{"task_id": task_id, "description": "test"},
config.run_root / "checkpoint.json",
)
assert result["cost_breakdown"]["cost_final"] is True
assert calls == ["POST", "GET", "GET"]
def _stub_terminal_task_executor(tmp_path, monkeypatch, gateway_result):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
monkeypatch.setattr(executor, "_generate", lambda *_args, **_kwargs: None)
monkeypatch.setattr(
executor_module,
"_install_workspace_backend_alias",
lambda *_args, **_kwargs: None,
)
monkeypatch.setattr(executor, "_workspace", lambda *_args, **_kwargs: "container-a")
monkeypatch.setattr(
executor,
"_task_body",
lambda task, *_args, **_kwargs: {"task_id": "cybergym-" + task.task_id.replace(":", "-")},
)
monkeypatch.setattr(
executor, "_gateway_wait", lambda *_args, **_kwargs: dict(gateway_result)
)
monkeypatch.setattr(
executor,
"_cleanup_workspace_container",
lambda *_args, **_kwargs: {"status": "verified"},
)
return config, executor
def test_fair_terminal_missing_marker_is_typed_and_settles_cost(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 185_217,
"completion_tokens": 754,
"cost_usd": 0.019249,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.019249,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:47101"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "failed"
assert rows[0]["capability_outcome"] == CAPABILITY_FINAL_POC_MISSING
assert rows[0]["final_submission_success"] is False
assert rows[0]["prompt_tokens"] == 185_217
assert rows[0]["completion_tokens"] == 754
assert rows[0]["cost_usd"] == pytest.approx(0.019249)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.019249)
assert projection.unresolved_upper_bound_usd == 0
def test_terminal_telemetry_failure_preserves_settled_cost(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["lifecycle"] == "post_gateway_evaluation_failed"
assert rows[0]["cost_usd"] == pytest.approx(0.25)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_missing_marker_with_failed_execution_stays_infra(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "infra_failed"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["capability_outcome"] == ""
assert rows[0]["final_submission_success"] is None
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_cleanup_diagnostic_failure_does_not_erase_terminal_cost(
tmp_path, monkeypatch
):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
def cleanup_failed(*_args, **_kwargs):
raise ExecutorFailure("cleanup failed")
original_write_json = executor_module._write_json
def fail_cleanup_report(path, value):
if pathlib.Path(path).name == "workspace_cleanup.json":
raise OSError("cleanup report failed")
return original_write_json(path, value)
monkeypatch.setattr(executor, "_cleanup_workspace_container", cleanup_failed)
monkeypatch.setattr(executor_module, "_write_json", fail_cleanup_report)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "failed"
assert rows[0]["capability_outcome"] == CAPABILITY_FINAL_POC_MISSING
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_pre_gateway_failures_settle_zero_and_do_not_block_next_task(
tmp_path, monkeypatch
):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
def fail_generation(*_args, **_kwargs):
raise ExecutorFailure("generation failed")
monkeypatch.setattr(executor, "_generate", fail_generation)
rows = run_campaign(
["arvo:1", "arvo:2"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=1,
)
assert [row["status"] for row in rows] == ["infra_failed", "infra_failed"]
assert all(row["cost_usd"] == 0 for row in rows)
assert all(row["cost_status"] == "known_no_dispatch" for row in rows)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=1).projection()
assert projection.settled_usd == 0
assert projection.unresolved_upper_bound_usd == 0
assert projection.can_dispatch is True
def test_post_admission_status_error_is_not_reclassified_as_zero_cost(
tmp_path, monkeypatch
):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
monkeypatch.setattr(executor, "_generate", lambda *_args, **_kwargs: None)
monkeypatch.setattr(
executor_module,
"_install_workspace_backend_alias",
lambda *_args, **_kwargs: None,
)
monkeypatch.setattr(executor, "_workspace", lambda *_args, **_kwargs: "container-a")
monkeypatch.setattr(
executor,
"_task_body",
lambda task, *_args, **_kwargs: {"task_id": "cybergym-" + task.task_id.replace(":", "-")},
)
def status_failed(*_args, **_kwargs):
raise ExecutorFailure("Ouroboros task status returned HTTP 404")
monkeypatch.setattr(executor, "_gateway_wait", status_failed)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["cost_usd"] is None
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == 0
assert projection.unresolved_upper_bound_usd is None
assert projection.can_dispatch is False
def test_cancel_503_recovers_terminal_gateway_payload(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-503"
terminal = {
"task_id": task_id,
"status": "failed",
"cost_usd": 0.060914,
"accounted_upper_bound_usd": 0.060914,
"unresolved_upper_bound_usd": 0.020062,
"cost_final": False,
}
calls = []
responses = iter(
(
{"status_code": 503, "body": {"detail": "teardown still live"}},
{"status_code": 200, "body": terminal},
)
)
def http(method, _url, **_kwargs):
calls.append(method)
return next(responses)
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
checkpoint = config.run_root / "checkpoint.json"
result = executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001
assert result == terminal
assert calls == ["POST", "GET"]
assert task_id not in executor._gateway_attempts
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "failed"
assert saved["cancel_status_code"] == 503
assert saved["result"]["accounted_upper_bound_usd"] == pytest.approx(0.060914)
def test_cancel_auth_failure_does_not_fallback_to_get(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-auth"
calls = []
def http(method, _url, **_kwargs):
calls.append(method)
return {"status_code": 401, "body": {"detail": "unauthorized"}}
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
with pytest.raises(ExecutorFailure, match="cancellation request failed"):
executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001
assert calls == ["POST"]
assert task_id in executor._gateway_attempts
def test_cancel_503_with_get_failure_keeps_custody_block(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-no-terminal"
calls = []
def http(method, _url, **_kwargs):
calls.append(method)
if method == "POST":
return {"status_code": 503, "body": {"detail": "teardown still live"}}
raise ExecutorFailure("status transport failed")
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
with pytest.raises(ExecutorFailure, match="status transport failed"):
executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001
assert calls == ["POST", "GET"]
assert task_id in executor._gateway_attempts
def dataclasses_replace(config, **changes):
import dataclasses

View file

@ -0,0 +1,780 @@
"""CyberGym gateway wire-protocol, served-telemetry, and accounting tests.
Split from ``tests/test_cybergym_executor.py`` along the HTTP/gateway seam:
provider probe, submit/verify/private-query wire parsing, served-telemetry
validation, gateway admission/cancel custody, and campaign cost accounting.
Shared fixtures (``_config``, ``dataclasses_replace``) are imported from the
original module; executor-lifecycle tests remain there.
"""
from __future__ import annotations
import gzip
import hashlib
import json
import pathlib
import pytest
from devtools.benchmarks.cybergym import cybergym_executor as executor_module
from devtools.benchmarks.cybergym.cybergym_adapter import (
CAPABILITY_FINAL_POC_MISSING,
BudgetLedger,
run_campaign,
)
from devtools.benchmarks.cybergym.cybergym_executor import (
CommandResult,
CyberGymExecutor,
ExecutorFailure,
_parse_json_stdout,
_require_exact_effort,
_served_telemetry,
_validate_verify_response,
)
from tests.test_cybergym_executor import _config, dataclasses_replace
def test_provider_probe_checks_exact_model_without_server_search(monkeypatch, tmp_path):
monkeypatch.setenv("OPENROUTER_API_KEY", "test-openrouter-key")
captured = {}
def http(method, url, *, body=None, headers=None, timeout=None):
if method == "GET" and url.endswith("/models"):
return {
"data": [{
"id": "deepseek/deepseek-v4-flash-0731",
"context_length": 1_310_720,
"supported_parameters": ["reasoning", "tools"],
}]
}
if method == "GET" and url.endswith("/key"):
return {"data": {"limit_remaining": 100}}
assert method == "POST"
captured["body"] = body
return {
"id": "response-1",
"model": "deepseek/deepseek-v4-flash-0731",
"provider": "OpenInference",
"choices": [{"message": {"content": "OK"}}],
"usage": {
"prompt_tokens": 12,
"completion_tokens": 3,
"cost": 0.006,
"cost_estimated": False,
},
}
executor = CyberGymExecutor(
_config(
tmp_path,
provider_probe=True,
expected_data_sha256="a" * 64,
expected_binary_sha256="b" * 64,
http_runner=http,
)
)
executor._probe_provider() # noqa: SLF001 - provider boundary assertion
assert captured["body"]["messages"] == [{"role": "user", "content": "Reply with OK."}]
assert "tools" not in captured["body"]
assert executor.provider_observation["observed_model"] == (
"deepseek/deepseek-v4-flash-0731"
)
def test_verify_response_requires_success_body_and_designated_poc():
good = {"message": "All 1 PoCs for this agent_id have been verified", "poc_ids": ["poc-1"]}
assert _validate_verify_response(good, expected_poc_id="poc-1") == good
with pytest.raises(ExecutorFailure, match="HTTP 500"):
_validate_verify_response({"status_code": 500, "body": {"detail": "failed"}})
with pytest.raises(ExecutorFailure, match="poc_ids"):
_validate_verify_response({"message": "ok", "poc_ids": []})
with pytest.raises(ExecutorFailure, match="designated poc_id"):
_validate_verify_response(good, expected_poc_id="other")
def test_observed_effort_must_be_exactly_high():
assert _require_exact_effort("high") == "high"
for value in ("", "High", "max", None):
with pytest.raises(ExecutorFailure, match="exactly high"):
_require_exact_effort(value)
def test_served_telemetry_prefers_authoritative_trace_refs_over_requested_fields():
payload = {
"model": "requested/not-served",
"reasoning_effort": "high",
"trace_refs": {
"llm_call_refs": [
{"resolved_model": "deepseek/deepseek-v4-flash-0731", "provider": "provider-a"}
]
},
}
observed = _served_telemetry(payload)
assert observed["observed_model"] == "deepseek/deepseek-v4-flash-0731"
assert observed["observed_provider"] == "provider-a"
assert observed["trace_call_count"] == 1
assert observed["effort_source"] == "runtime_requested_field"
def test_served_telemetry_rejects_incomplete_or_mixed_trace_identity():
with pytest.raises(ExecutorFailure, match="incomplete served-call"):
_served_telemetry({"trace_refs": {"llm_call_refs": [{"provider": "provider-a"}]}})
with pytest.raises(ExecutorFailure, match="mixed served models"):
_served_telemetry(
{
"trace_refs": {
"llm_call_refs": [
{"resolved_model": "model-a", "provider": "provider-a"},
{"resolved_model": "model-b", "provider": "provider-a"},
]
}
}
)
def test_served_telemetry_reads_verified_response_wire_effort(tmp_path):
drive = tmp_path / "drive"
calls = drive / "observability" / "calls" / "opaque"
calls.mkdir(parents=True)
wire = {
"requested_effort": "high",
"applied_effort": "high",
"attempt_id": "attempt-1",
"candidate_sha256": "a" * 64,
}
blob_raw = json.dumps(
{"usage": {"request_wire": wire, "response_provider": "backend-a"}},
sort_keys=True,
).encode("utf-8")
blob_path = drive / "observability" / "blobs" / ("b" * 64 + ".json.gz")
blob_path.parent.mkdir(parents=True)
blob_path.write_bytes(gzip.compress(blob_raw))
blob_ref = {
"path": str(blob_path),
"sha256": hashlib.sha256(blob_raw).hexdigest(),
"size": len(blob_raw),
"kind": "json",
"encoding": "gzip",
}
manifest_raw = json.dumps(
{
"task_id": "opaque",
"call_id": "llm-1_response",
"llm_call_id": "llm-1",
"full_payload_ref": blob_ref,
},
sort_keys=True,
).encode("utf-8")
manifest_path = calls / "llm-1_response.json"
manifest_path.write_bytes(manifest_raw)
manifest_ref = {
"path": str(manifest_path),
"sha256": hashlib.sha256(manifest_raw).hexdigest(),
"call_id": "llm-1_response",
}
observed = _served_telemetry(
{
"reasoning_effort": "low",
"trace_refs": {
"llm_call_refs": [
{
"llm_call_id": "llm-1",
"resolved_model": "deepseek/deepseek-v4-flash-0731",
"provider": "provider-a",
"response_ref": manifest_ref,
}
]
},
},
allowed_roots=(drive,),
)
assert observed["observed_effort"] == "high"
assert observed["observed_provider"] == "backend-a"
assert observed["observed_provider_attempts"] == ["backend-a"]
assert observed["provider_distribution"] == {"backend-a": 1}
assert observed["effort_source"] == "served_response_wire"
assert observed["response_wire_effort_count"] == 1
assert observed["response_wire_provider_count"] == 1
def test_submit_stdout_parser_accepts_preceding_prose_and_multiline_json():
parsed = _parse_json_stdout('notice\n{\n "task_id": "opaque1234",\n "poc_id": "poc-1"\n}\n')
assert parsed == {"task_id": "opaque1234", "poc_id": "poc-1"}
def test_private_query_rejects_http_and_body_errors(tmp_path, monkeypatch):
config = _config(tmp_path)
monkeypatch.setenv("CYBERGYM_API_KEY", "test-secret-value")
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"status_code": 404, "body": {"detail": "Record not found"}},
)
)
with pytest.raises(ExecutorFailure, match="HTTP 404"):
executor._private_query("agent-" + "a" * 24, "arvo:1")
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"status_code": 200, "body": {"error": {"message": "bad"}}},
)
)
with pytest.raises(ExecutorFailure, match="error object"):
executor._private_query("agent-" + "a" * 24, "arvo:1")
def test_private_query_accepts_nested_items_wrapper(tmp_path, monkeypatch):
config = _config(tmp_path)
monkeypatch.setenv("CYBERGYM_API_KEY", "test-secret-value")
record = {"task_id": "arvo:1", "poc_id": "poc-1", "poc_hash": "a" * 64}
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {"pocs": {"items": [record]}},
)
)
assert executor._private_query("agent-" + "a" * 24, "arvo:1") == [record]
def test_submit_response_binds_poc_id_not_nonexistent_hash_and_keeps_exit_code(tmp_path):
config = _config(tmp_path)
task_dir = config.run_root / "task"
task_dir.mkdir()
(task_dir / "final.poc").write_bytes(b"poc-bytes")
(task_dir / "submit.sh").write_text("TASK_ID=opaque1234\n", encoding="utf-8")
executor_name = "workspace"
executor_id = "c" * 64
submit_calls = []
def command(argv, *, cwd=None, env=None, timeout=None):
submit_calls.append(list(argv))
return CommandResult(
0,
json.dumps(
{
"task_id": "opaque1234",
"poc_id": "poc-1",
"exit_code": 71,
"output": "known",
# Upstream does not define a response hash; an incidental
# field must not override the local marker binding.
"hash": "not-the-poc-hash",
}
),
"",
)
executor = CyberGymExecutor(dataclasses_replace(config, command_runner=command))
executor._task_containers[executor_name] = executor_id
response, digest, masked = executor._submit_final( # noqa: SLF001 - boundary contract assertion
type("Task", (), {"task_id": "arvo:1"})(), task_dir, "workspace"
)
assert response["poc_id"] == "poc-1"
assert response["exit_code"] == 71
assert digest == hashlib.sha256(b"poc-bytes").hexdigest()
assert masked == "opaque1234"
assert executor_id in submit_calls[0]
assert executor_name not in submit_calls[0]
def test_unknown_gateway_attempt_blocks_campaign_cleanup(tmp_path):
config = _config(tmp_path)
calls = []
def command(*args, **kwargs):
calls.append(args)
raise AssertionError("cleanup must not run while gateway custody is unknown")
executor = CyberGymExecutor(dataclasses_replace(config, command_runner=command))
executor.started = True
executor.server_id = "server-123"
executor.network_id = "network-123"
executor._task_containers = {"workspace-agent-aaaaaaaaaaaaaaaaaaaaaaaa": "workspace-123"}
executor._gateway_attempts = {
"cybergym-attempt": {
"gateway_task_id": "cybergym-attempt",
"status": "admission_unknown",
"checkpoint": str(config.run_root / "checkpoint.json"),
}
}
report = executor.close()
assert report["ok"] is False
assert report["status"] == "custody_pending"
assert executor.custody_blocked is True
assert executor.server_id == "server-123"
assert executor.network_id == "network-123"
assert calls == []
assert (config.run_root / "custody_pending.json").is_file()
def test_gateway_admission_transport_error_registers_durable_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
seen = {}
def failing_http(*args, **kwargs):
seen.update(kwargs)
raise ExecutorFailure("HTTP POST transport failed")
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=failing_http))
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-opaque-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="transport failed"):
executor._gateway_wait(body, checkpoint)
assert "cybergym-opaque-attempt" in executor._gateway_attempts
assert executor._gateway_attempts["cybergym-opaque-attempt"]["status"] == "admission_unknown"
assert seen["headers"]["Idempotency-Key"].startswith("cybergym-")
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["custody_required"] is True
assert saved["status"] == "admission_unknown"
def test_gateway_definitive_admission_rejection_releases_phantom_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(
dataclasses_replace(
config,
http_runner=lambda *args, **kwargs: {
"status_code": 400,
"body": {"detail": "invalid task"},
},
)
)
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-rejected-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="HTTP 400"):
executor._gateway_wait(body, checkpoint)
assert executor._gateway_attempts == {}
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "admission_rejected"
assert saved["custody_required"] is False
def test_gateway_malformed_admission_keeps_unknown_custody(tmp_path):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=lambda *args, **kwargs: {})
)
checkpoint = config.run_root / "checkpoint.json"
body = {"task_id": "cybergym-malformed-attempt", "description": "test"}
with pytest.raises(ExecutorFailure, match="no task id"):
executor._gateway_wait(body, checkpoint)
assert "cybergym-malformed-attempt" in executor._gateway_attempts
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "admission_unknown_response"
assert saved["custody_required"] is True
def test_gateway_waits_for_final_cost_after_completed_status(tmp_path):
config = _config(tmp_path, provider_probe=False, task_timeout_sec=10)
task_id = "cybergym-cost-pending"
calls = []
status_rows = iter(
(
{
"task_id": task_id,
"status": "completed",
"result": {"cost_final": False},
},
{
"task_id": task_id,
"status": "completed",
"result": {"cost_final": True},
},
)
)
def http(method, url, **kwargs):
calls.append(method)
if method == "POST":
return {"task_id": task_id, "status": "scheduled"}
return next(status_rows)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None)
)
result = executor._gateway_wait(
{"task_id": task_id, "description": "test"},
config.run_root / "checkpoint.json",
)
assert result["result"]["cost_final"] is True
assert calls == ["POST", "GET", "GET"]
def test_gateway_cost_finality_conflict_keeps_polling(tmp_path):
config = _config(tmp_path, provider_probe=False, task_timeout_sec=10)
task_id = "cybergym-cost-conflict"
calls = []
status_rows = iter(
(
{
"task_id": task_id,
"status": "completed",
"cost_final": True,
"cost_breakdown": {"cost_final": False},
},
{
"task_id": task_id,
"status": "completed",
"cost_final": True,
"cost_breakdown": {"cost_final": True},
},
)
)
def http(method, _url, **_kwargs):
calls.append(method)
if method == "POST":
return {"task_id": task_id, "status": "scheduled"}
return next(status_rows)
executor = CyberGymExecutor(
dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None)
)
result = executor._gateway_wait( # noqa: SLF001 - accounting contract
{"task_id": task_id, "description": "test"},
config.run_root / "checkpoint.json",
)
assert result["cost_breakdown"]["cost_final"] is True
assert calls == ["POST", "GET", "GET"]
def _stub_terminal_task_executor(tmp_path, monkeypatch, gateway_result):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
monkeypatch.setattr(executor, "_generate", lambda *_args, **_kwargs: None)
monkeypatch.setattr(
executor_module,
"_install_workspace_backend_alias",
lambda *_args, **_kwargs: None,
)
monkeypatch.setattr(executor, "_workspace", lambda *_args, **_kwargs: "container-a")
monkeypatch.setattr(
executor,
"_task_body",
lambda task, *_args, **_kwargs: {"task_id": "cybergym-" + task.task_id.replace(":", "-")},
)
monkeypatch.setattr(
executor, "_gateway_wait", lambda *_args, **_kwargs: dict(gateway_result)
)
monkeypatch.setattr(
executor,
"_cleanup_workspace_container",
lambda *_args, **_kwargs: {"status": "verified"},
)
return config, executor
def test_fair_terminal_missing_marker_is_typed_and_settles_cost(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 185_217,
"completion_tokens": 754,
"cost_usd": 0.019249,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.019249,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:47101"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "failed"
assert rows[0]["capability_outcome"] == CAPABILITY_FINAL_POC_MISSING
assert rows[0]["final_submission_success"] is False
assert rows[0]["prompt_tokens"] == 185_217
assert rows[0]["completion_tokens"] == 754
assert rows[0]["cost_usd"] == pytest.approx(0.019249)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.019249)
assert projection.unresolved_upper_bound_usd == 0
def test_terminal_telemetry_failure_preserves_settled_cost(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["lifecycle"] == "post_gateway_evaluation_failed"
assert rows[0]["cost_usd"] == pytest.approx(0.25)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_missing_marker_with_failed_execution_stays_infra(tmp_path, monkeypatch):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "infra_failed"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["capability_outcome"] == ""
assert rows[0]["final_submission_success"] is None
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_cleanup_diagnostic_failure_does_not_erase_terminal_cost(
tmp_path, monkeypatch
):
gateway_result = {
"status": "completed",
"observed_model": "deepseek/deepseek-v4-flash-0731",
"observed_provider": "backend-a",
"reasoning_effort": "high",
"prompt_tokens": 100,
"completion_tokens": 10,
"cost_usd": 0.25,
"cost_final": True,
"cost_breakdown": {
"accounted_upper_bound_usd": 0.25,
"cost_final": True,
},
"outcome_axes": {"execution": {"status": "ok"}},
}
config, executor = _stub_terminal_task_executor(
tmp_path, monkeypatch, gateway_result
)
def cleanup_failed(*_args, **_kwargs):
raise ExecutorFailure("cleanup failed")
original_write_json = executor_module._write_json
def fail_cleanup_report(path, value):
if pathlib.Path(path).name == "workspace_cleanup.json":
raise OSError("cleanup report failed")
return original_write_json(path, value)
monkeypatch.setattr(executor, "_cleanup_workspace_container", cleanup_failed)
monkeypatch.setattr(executor_module, "_write_json", fail_cleanup_report)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "failed"
assert rows[0]["capability_outcome"] == CAPABILITY_FINAL_POC_MISSING
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == pytest.approx(0.25)
assert projection.unresolved_upper_bound_usd == 0
def test_pre_gateway_failures_settle_zero_and_do_not_block_next_task(
tmp_path, monkeypatch
):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
def fail_generation(*_args, **_kwargs):
raise ExecutorFailure("generation failed")
monkeypatch.setattr(executor, "_generate", fail_generation)
rows = run_campaign(
["arvo:1", "arvo:2"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=1,
)
assert [row["status"] for row in rows] == ["infra_failed", "infra_failed"]
assert all(row["cost_usd"] == 0 for row in rows)
assert all(row["cost_status"] == "known_no_dispatch" for row in rows)
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=1).projection()
assert projection.settled_usd == 0
assert projection.unresolved_upper_bound_usd == 0
assert projection.can_dispatch is True
def test_post_admission_status_error_is_not_reclassified_as_zero_cost(
tmp_path, monkeypatch
):
config = _config(tmp_path, provider_probe=False)
executor = CyberGymExecutor(config)
monkeypatch.setattr(executor, "start", lambda: None)
monkeypatch.setattr(executor, "_generate", lambda *_args, **_kwargs: None)
monkeypatch.setattr(
executor_module,
"_install_workspace_backend_alias",
lambda *_args, **_kwargs: None,
)
monkeypatch.setattr(executor, "_workspace", lambda *_args, **_kwargs: "container-a")
monkeypatch.setattr(
executor,
"_task_body",
lambda task, *_args, **_kwargs: {"task_id": "cybergym-" + task.task_id.replace(":", "-")},
)
def status_failed(*_args, **_kwargs):
raise ExecutorFailure("Ouroboros task status returned HTTP 404")
monkeypatch.setattr(executor, "_gateway_wait", status_failed)
rows = run_campaign(
["arvo:1"],
run_root=config.run_root,
executor=executor.run_task,
estimated_cost_usd=1,
budget_cap_usd=2,
)
assert rows[0]["status"] == "infra_failed"
assert rows[0]["cost_usd"] is None
projection = BudgetLedger(config.run_root / "claims.jsonl", cap_usd=2).projection()
assert projection.settled_usd == 0
assert projection.unresolved_upper_bound_usd is None
assert projection.can_dispatch is False
def test_cancel_503_recovers_terminal_gateway_payload(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-503"
terminal = {
"task_id": task_id,
"status": "failed",
"cost_usd": 0.060914,
"accounted_upper_bound_usd": 0.060914,
"unresolved_upper_bound_usd": 0.020062,
"cost_final": False,
}
calls = []
responses = iter(
(
{"status_code": 503, "body": {"detail": "teardown still live"}},
{"status_code": 200, "body": terminal},
)
)
def http(method, _url, **_kwargs):
calls.append(method)
return next(responses)
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
checkpoint = config.run_root / "checkpoint.json"
result = executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001
assert result == terminal
assert calls == ["POST", "GET"]
assert task_id not in executor._gateway_attempts
saved = json.loads(checkpoint.read_text(encoding="utf-8"))
assert saved["status"] == "failed"
assert saved["cancel_status_code"] == 503
assert saved["result"]["accounted_upper_bound_usd"] == pytest.approx(0.060914)
def test_cancel_auth_failure_does_not_fallback_to_get(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-auth"
calls = []
def http(method, _url, **_kwargs):
calls.append(method)
return {"status_code": 401, "body": {"detail": "unauthorized"}}
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
with pytest.raises(ExecutorFailure, match="cancellation request failed"):
executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001
assert calls == ["POST"]
assert task_id in executor._gateway_attempts
def test_cancel_503_with_get_failure_keeps_custody_block(tmp_path):
config = _config(tmp_path, poll_interval_sec=0)
task_id = "cybergym-cancel-no-terminal"
calls = []
def http(method, _url, **_kwargs):
calls.append(method)
if method == "POST":
return {"status_code": 503, "body": {"detail": "teardown still live"}}
raise ExecutorFailure("status transport failed")
executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http))
executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion
"gateway_task_id": task_id,
"status": "submitted",
}
with pytest.raises(ExecutorFailure, match="status transport failed"):
executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001
assert calls == ["POST", "GET"]
assert task_id in executor._gateway_attempts