diff --git a/devtools/benchmarks/cybergym/cybergym_adapter.py b/devtools/benchmarks/cybergym/cybergym_adapter.py index 1547e557e..1c3028269 100644 --- a/devtools/benchmarks/cybergym/cybergym_adapter.py +++ b/devtools/benchmarks/cybergym/cybergym_adapter.py @@ -148,14 +148,7 @@ class BudgetProjection: def as_dict(self) -> dict[str, Any]: return { - "cap_usd": self.cap_usd, - "settled_usd": self.settled_usd, - "reserved_usd": self.reserved_usd, - "unresolved_upper_bound_usd": self.unresolved_upper_bound_usd, - "projected_usd": self.projected_usd, - "available_usd": self.available_usd, - "can_dispatch": self.can_dispatch, - "reason": self.reason, + **dataclasses.asdict(self), "active_task_ids": list(self.active_task_ids), "active_attempt_ids": list(self.active_attempt_ids), } @@ -988,26 +981,17 @@ class BudgetLedger: ) if before_write is not None: before_write(overspend_error) - if overspend: - self._append( - { - "schema": LEDGER_SCHEMA, - "event": "overspend", - "attempt_id": attempt, - "cost_usd": cost, - "ts_unix": time.time(), - } - ) - raise overspend_error self._append( { "schema": LEDGER_SCHEMA, - "event": "settle", + "event": "overspend" if overspend else "settle", "attempt_id": attempt, "cost_usd": cost, "ts_unix": time.time(), } ) + if overspend: + raise overspend_error def record_campaign_cost(self, cost_usd: float, *, label: str = "provider_probe") -> dict[str, Any]: """Record a known campaign-level charge before task reservations. @@ -1439,67 +1423,17 @@ def run_campaign( if callback_contract is not None else (contract if isinstance(contract, Mapping) else None), ) - except BudgetOverspend as exc: - budget_refs = dict(outcome.get("artifact_refs") or {}) - budget_refs.setdefault("task_dir", str(task_dir)) - budget_refs.setdefault("claims", str(ledger.path)) - if claim is not None: - budget_refs.setdefault( - "checkpoint", - str( - safe_task_path( - root / "checkpoints", - task.task_id, - str(claim["attempt_id"]), - ) - / "gateway_checkpoint.json" - ), - ) - budget_refs.setdefault("custody_pending", str(root / "custody_pending.json")) - row = build_task_result_row( - task.task_id, - trials=outcome.get("trials") or (), - final_trial=outcome.get("final_trial"), - final_poc_sha256=str(outcome.get("final_poc_sha256") or ""), - status="infra_failed", - lifecycle="budget_refused", - level=task.level, - masked_id=str(outcome.get("masked_id") or ""), - masked_id_source=str(outcome.get("masked_id_source") or ""), - observed_provider=str(outcome.get("observed_provider") or ""), - observed_model=str(outcome.get("observed_model") or ""), - observed_effort=( - str(outcome.get("observed_effort") or "") - if str(outcome.get("observed_effort") or "").strip().lower() == "high" - else "" - ), - observed_effort_source=str(outcome.get("observed_effort_source") or ""), - prompt_tokens=outcome.get("prompt_tokens"), - completion_tokens=outcome.get("completion_tokens"), - cached_tokens=outcome.get("cached_tokens"), - cost_usd=outcome.get("cost_usd"), - cost_estimated=outcome.get("cost_estimated"), - cost_final=outcome.get("cost_final"), - cost_status=str(outcome.get("cost_status") or ""), - infra_reason="budget_overspend", - artifact_refs=budget_refs, - error=str(exc), - runtime_result=outcome.get("runtime_result"), - task_contract=callback_contract - if callback_contract is not None - else (contract if isinstance(contract, Mapping) else None), - attempt_id=str(claim["attempt_id"]) if claim else "", - ) - except BudgetRefused: - # Claim-time refusal: no ledger event exists and the task was - # never dispatched, so it must NOT become an infra row. The - # dispatch engine pauses admission, waits for in-flight - # settlements to free headroom, and either resumes or ends the - # campaign with BudgetCapReached (run 20260907T233516Z flushed - # 1145 undispatched tasks into infra rows here). - raise except Exception as exc: - if claim is not None: + overspend = isinstance(exc, BudgetOverspend) + if isinstance(exc, BudgetRefused) and not overspend: + # Claim-time refusal: no ledger event exists and the task was + # never dispatched, so it must NOT become an infra row. The + # dispatch engine pauses admission, waits for in-flight + # settlements to free headroom, and either resumes or ends the + # campaign with BudgetCapReached (run 20260907T233516Z flushed + # 1145 undispatched tasks into infra rows here). + raise + if claim is not None and not overspend: terminal_accounting = _terminal_gateway_accounting( outcome.get("runtime_result") ) @@ -1527,7 +1461,7 @@ def run_campaign( final_trial=outcome.get("final_trial"), final_poc_sha256=str(outcome.get("final_poc_sha256") or ""), status="infra_failed", - lifecycle="executor_failed", + lifecycle="budget_refused" if overspend else "executor_failed", level=task.level, masked_id=str(outcome.get("masked_id") or ""), masked_id_source=str(outcome.get("masked_id_source") or ""), @@ -1546,7 +1480,7 @@ def run_campaign( cost_estimated=outcome.get("cost_estimated"), cost_final=outcome.get("cost_final"), cost_status=str(outcome.get("cost_status") or ""), - infra_reason=type(exc).__name__, + infra_reason="budget_overspend" if overspend else type(exc).__name__, artifact_refs=failure_refs, error=str(exc), runtime_result=outcome.get("runtime_result"), diff --git a/devtools/benchmarks/cybergym/cybergym_custody.py b/devtools/benchmarks/cybergym/cybergym_custody.py index 0c2812b7a..5545c4256 100644 --- a/devtools/benchmarks/cybergym/cybergym_custody.py +++ b/devtools/benchmarks/cybergym/cybergym_custody.py @@ -1,30 +1,33 @@ -"""Gateway cancellation custody for the CyberGym executor. +"""Gateway admission, waiting, cancellation and terminal custody. -Split out of ``cybergym_lifecycle.py`` so the lifecycle layer stays inside -its size-ratchet band. ``_CustodyMixin`` holds the cancellation request and -the post-cancellation custody poll — the seam that decides when a gateway -task the launcher stopped waiting for is truly terminal — mixed into -``CyberGymExecutor`` and dispatched on ``self`` at runtime. Everything here -is gateway HTTP plus checkpoint persistence; it never touches containers or -child processes. +``_CustodyMixin`` keeps one gateway attempt registered from admission through +terminal observation and transfer to outer-write custody. Normal and cancelled +waits retain their distinct bounds; both feed the same terminal transfer. +The methods are mixed into ``CyberGymExecutor`` without another scheduler. +Everything here is gateway HTTP, checkpoint persistence and in-memory ownership; +it never operates on containers or child processes. """ from __future__ import annotations +import hashlib import pathlib import time import urllib.parse +import uuid from collections.abc import Mapping from typing import Any from devtools.benchmarks.cybergym.cybergym_adapter import _TERMINAL_GATEWAY_STATUSES -from devtools.benchmarks.cybergym.cybergym_docker import _write_json +from devtools.benchmarks.cybergym.cybergym_docker import _GATEWAY_TASK_ID, _write_json from devtools.benchmarks.cybergym.cybergym_wire import ( FINALIZATION_GRACE_SEC, GATEWAY_TRANSPORT_RETRY_BUDGET_SEC, ExecutorFailure, + GatewayAdmissionRejected, GatewayTransportError, HttpStatusError, _cost_is_pending, + _definitive_admission_rejection, _CostGraceTracker, _gateway_finalizing, _gateway_path, @@ -34,8 +37,273 @@ from devtools.benchmarks.cybergym.cybergym_wire import ( ) +# Gateway statuses under which the task has been admitted but has not started +# executing: no worker lane, no provider spend, no wall clock the agent can +# pace against. The launcher's task deadline starts when the task leaves this +# set (full1507 postmortem: a submit-anchored deadline cancelled healthy tasks +# after ~1 h of runtime because they had queued ~1 h behind a finalization +# backlog). The isolate's own ``OUROBOROS_TASK_ABS_CEILING_SEC`` bounds the +# RUNNING phase from the same moment; ``TASK_DEADLINE_GRACE_SEC`` keeps the +# launcher's cancel a backstop behind that server-side settle, not a race +# against it. +_QUEUED_GATEWAY_STATUSES = frozenset({"", "scheduled", "queued", "pending"}) +TASK_DEADLINE_GRACE_SEC = 300.0 + class _CustodyMixin: - """Cancellation and terminal-custody methods for the CyberGym executor.""" + """Own gateway attempt custody across normal and cancelled execution.""" + + def _terminalize_gateway_attempt(self, gateway_task_id: str) -> None: + """Atomically transfer a settled gateway attempt to outer-write custody.""" + with self._registry_condition: + entry = self._gateway_attempts.get(gateway_task_id) + if isinstance(entry, Mapping): + workspace_name = str(entry.get("workspace_name") or "") + if workspace_name: + self._terminal_uncommitted_workspaces[workspace_name] = { + "task_id": str(entry.get("task_id") or ""), + "attempt_id": str(entry.get("attempt_id") or ""), + } + self._gateway_attempts.pop(gateway_task_id, None) + + def probe_gateway_alive(self) -> bool: + """Liveness probe for the dispatch breaker: did the gateway answer? + + Any answer (even a non-2xx status) proves the transport is back; only + a transport-level failure keeps the campaign paused. + """ + + try: + self.config.http_runner( + "GET", + _gateway_path(self.config.ouroboros_url, "/api/health"), + timeout=15, + ) + except GatewayTransportError: + return False + except HttpStatusError: + return True + except Exception: # noqa: BLE001 - malformed body still means "answered" + return True + return True + + def _gateway_wait( + self, + body: Mapping[str, Any], + checkpoint: pathlib.Path, + *, + workspace_name: str = "", + task_id: str = "", + attempt_id: str = "", + ) -> Mapping[str, Any]: + requested_task_id = str(body.get("task_id") or "").strip() + owner_task_id = str(task_id) + owner_attempt_id = str(attempt_id) + # The gateway currently echoes the opaque caller task id. Register it + # before POST so a dropped response can still be treated as an + # admitted-or-unknown attempt and retained for manual reattachment. + pending_id = requested_task_id or ("pending-" + uuid.uuid4().hex) + idempotency_key = "cybergym-" + hashlib.sha256( + (pending_id + "\0" + str(body.get("actor_id") or "cybergym")).encode() + ).hexdigest() + self._gateway_attempts[pending_id] = { + "gateway_task_id": requested_task_id, + "status": "admission_pending", + "checkpoint": str(checkpoint), + "idempotency_key": idempotency_key, + "workspace_name": str(workspace_name), + "task_id": owner_task_id, + "attempt_id": owner_attempt_id, + } + try: + created = _unwrap_http_json( + self.config.http_runner( + "POST", + _gateway_path(self.config.ouroboros_url, "/api/tasks"), + body=body, + headers={"Idempotency-Key": idempotency_key}, + timeout=60, + ), + operation="Ouroboros task admission", + ) + except BaseException as exc: + rejected = _definitive_admission_rejection(exc) + status = "admission_rejected" if rejected else "admission_unknown" + entry = self._gateway_attempts.get(pending_id) + if entry is not None: + entry.update({"status": status, "error": type(exc).__name__}) + if rejected: + # A typed 4xx response is evidence that the gateway refused the + # request before scheduling it. Do not retain a phantom + # custody claim, but keep the redacted checkpoint for audit. + self._gateway_attempts.pop(pending_id, None) + _write_json( + checkpoint, + { + "gateway_task_id": requested_task_id or pending_id, + "status": status, + "custody_required": not rejected, + "idempotency_key": idempotency_key, + "error": type(exc).__name__, + }, + ) + if rejected: + raise GatewayAdmissionRejected(str(exc)) from exc + raise + task_id = str(created.get("task_id") or "").strip() + if not task_id or not _GATEWAY_TASK_ID.fullmatch(task_id): + self._gateway_attempts[pending_id]["status"] = "admission_unknown_response" + _write_json( + checkpoint, + { + "gateway_task_id": requested_task_id or pending_id, + "status": "admission_unknown_response", + "custody_required": True, + "idempotency_key": idempotency_key, + }, + ) + raise ExecutorFailure("Ouroboros gateway returned no task id") + if requested_task_id and task_id != requested_task_id: + self._gateway_attempts[pending_id].update( + {"gateway_task_id": task_id, "status": "admission_id_mismatch"} + ) + _write_json( + checkpoint, + { + "gateway_task_id": task_id, + "submitted_task_id": requested_task_id, + "status": "admission_id_mismatch", + "custody_required": True, + "idempotency_key": idempotency_key, + }, + ) + raise ExecutorFailure("Ouroboros gateway changed the submitted task id") + if pending_id != task_id: + self._gateway_attempts.pop(pending_id, None) + self._gateway_attempts[task_id] = { + "gateway_task_id": task_id, + "status": "submitted", + "checkpoint": str(checkpoint), + "idempotency_key": idempotency_key, + "workspace_name": str(workspace_name), + "task_id": owner_task_id, + "attempt_id": owner_attempt_id, + } + _write_json( + checkpoint, + { + "gateway_task_id": task_id, + "status": "submitted", + "idempotency_key": idempotency_key, + "body": {k: v for k, v in body.items() if k != "description"}, + }, + ) + # Two bounds, one active at a time: the queue-wait cap while the + # gateway still reports the task as not started, then the task + # deadline anchored at the first observed non-queued status. + queue_started = time.monotonic() + queue_wait_cap = queue_started + float(self.config.task_timeout_sec) + run_deadline: float | None = None + observed_start_at: str | None = None + latest: Mapping[str, Any] = created + cost_grace = _CostGraceTracker() + transport_deadline: float | None = None + finalization_grace_until: float | None = None + while True: + bound = run_deadline if run_deadline is not None else queue_wait_cap + if time.monotonic() >= bound: + # The worker is done and the server is finalizing artifacts: + # a finished, paid result is minutes away — wait for it (once, + # bounded) instead of cancelling it. + if finalization_grace_until is None and _gateway_finalizing(latest): + finalization_grace_until = time.monotonic() + FINALIZATION_GRACE_SEC + run_deadline = finalization_grace_until + _write_json(checkpoint, { + "gateway_task_id": task_id, + "status": _response_status(latest), + "result": dict(latest), + "deadline_basis": "finalization_grace", + "finalization_grace_sec": FINALIZATION_GRACE_SEC, + }) + continue + break + try: + latest = _unwrap_http_json( + self.config.http_runner( + "GET", + _gateway_path(self.config.ouroboros_url, "/api/tasks/" + urllib.parse.quote(task_id, safe="")), + timeout=60, + ), + operation="Ouroboros task status", + ) + except GatewayTransportError: + # A transient transport failure (an isolate event-loop stall + # starves the HTTP answer) must not kill a healthy paid task + # on the first error: ride it out within a bounded budget. + # Exhaustion re-raises so a dead gateway still produces the + # circuit-breaker row. + now = time.monotonic() + if now >= bound: + # The task's own deadline passed while the gateway was + # unreachable: stop polling and cancel it like a normal + # deadline exit instead of writing a transport row. + break + if transport_deadline is None: + transport_deadline = min( + bound, now + GATEWAY_TRANSPORT_RETRY_BUDGET_SEC + ) + if now >= transport_deadline: + raise + self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) + continue + transport_deadline = None + returned_id = str(latest.get("task_id") or "").strip() + if returned_id and returned_id != task_id: + raise ExecutorFailure("Ouroboros status response belongs to a different task") + status = _response_status(latest) + if run_deadline is None and status not in _QUEUED_GATEWAY_STATUSES: + run_deadline = ( + time.monotonic() + + float(self.config.task_timeout_sec) + + TASK_DEADLINE_GRACE_SEC + ) + observed_start_at = time.strftime( + "%Y-%m-%dT%H:%M:%SZ", time.gmtime() + ) + frame = { + "gateway_task_id": task_id, + "status": status, + "result": dict(latest), + "deadline_basis": ( + "finalization_grace" if finalization_grace_until is not None + else "observed_start" if run_deadline is not None else "queue_wait_cap" + ), + } + if observed_start_at is not None: + frame["observed_start_at"] = observed_start_at + _write_json(checkpoint, frame) + if status in _TERMINAL_GATEWAY_STATUSES: + # Root post-task accounting can publish ``completed`` before + # its durable cost roll-up is final; only the bounded + # abandoned-residue grace (cybergym_wire) releases such a + # frame early, with the residue disclosed on it. + if status == "completed" and _cost_is_pending(latest): + accepted = cost_grace.accept( + latest, + now=time.monotonic(), + wall_now=time.time(), + ) + if accepted is None: + self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) + continue + latest = accepted + _write_json(checkpoint, {"gateway_task_id": task_id, "status": status, "result": dict(latest)}) + self._terminalize_gateway_attempt(task_id) + return latest + self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) + # The task may still be running after the local wait expires. Ask the + # gateway to stop it and retain the original attempt until a terminal + # custody response is observed; never return a reusable task id here. + return self._cancel_gateway_task(task_id, checkpoint) def _poll_gateway_custody( self, diff --git a/devtools/benchmarks/cybergym/cybergym_lifecycle.py b/devtools/benchmarks/cybergym/cybergym_lifecycle.py index 079a001ab..79abfcf75 100644 --- a/devtools/benchmarks/cybergym/cybergym_lifecycle.py +++ b/devtools/benchmarks/cybergym/cybergym_lifecycle.py @@ -5,7 +5,7 @@ existing imports keep working) to keep each module inside the size ratchet. This layer sits above ``cybergym_docker`` (it imports a few docker helpers from there) and below the executor assembly; it never imports the executor, so no import cycle is introduced. ``_LifecycleMixin`` collects the provider/settings -probe, startup, gateway dispatch, submission, and cleanup-custody methods that +probe, startup, submission, and cleanup-custody methods that are mixed into ``CyberGymExecutor`` and dispatched on ``self`` at runtime. """ @@ -18,7 +18,6 @@ import pathlib import re import shutil import time -import urllib.parse import uuid from collections.abc import Mapping, Sequence from typing import Any @@ -62,19 +61,10 @@ from devtools.benchmarks.cybergym.cybergym_sidecar import ( from devtools.benchmarks.cybergym.cybergym_wire import ( _HEX64, _PROVIDER_ID, - FINALIZATION_GRACE_SEC, - GATEWAY_TRANSPORT_RETRY_BUDGET_SEC, ExecutorFailure, - GatewayAdmissionRejected, - GatewayTransportError, HttpStatusError, - _cost_is_pending, - _CostGraceTracker, - _definitive_admission_rejection, _gateway_fair_completion, - _gateway_finalizing, _gateway_has_tool_markup, - _gateway_path, _nonnegative_number, _positive_int, _require_exact_effort, @@ -88,19 +78,6 @@ from devtools.benchmarks.cybergym.cybergym_wire import ( ) from ouroboros.openrouter_attribution import OPENROUTER_APP_HEADERS -_SETTLED = frozenset({"completed", "failed", "cancelled", "rejected_duplicate"}) - -# Gateway statuses under which the task has been admitted but has not started -# executing: no worker lane, no provider spend, no wall clock the agent can -# pace against. The launcher's task deadline starts when the task leaves this -# set (full1507 postmortem: a submit-anchored deadline cancelled healthy tasks -# after ~1 h of runtime because they had queued ~1 h behind a finalization -# backlog). The isolate's own ``OUROBOROS_TASK_ABS_CEILING_SEC`` bounds the -# RUNNING phase from the same moment; ``TASK_DEADLINE_GRACE_SEC`` keeps the -# launcher's cancel a backstop behind that server-side settle, not a race -# against it. -_QUEUED_GATEWAY_STATUSES = frozenset({"", "scheduled", "queued", "pending"}) -TASK_DEADLINE_GRACE_SEC = 300.0 _MASKED_TASK_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{7,255}$") @@ -351,7 +328,7 @@ def _write_checkpoint_delivery( class _LifecycleMixin: - """Provider/settings/startup/gateway/cleanup lifecycle methods.""" + """Provider/settings/startup/submission/cleanup lifecycle methods.""" def _ensure_key(self) -> str: value = os.environ.get(self.config.api_key_env, "") @@ -871,259 +848,6 @@ class _LifecycleMixin: }, } - def _terminalize_gateway_attempt(self, gateway_task_id: str) -> None: - """Atomically transfer a settled gateway attempt to outer-write custody.""" - with self._registry_condition: - entry = self._gateway_attempts.get(gateway_task_id) - if isinstance(entry, Mapping): - workspace_name = str(entry.get("workspace_name") or "") - if workspace_name: - self._terminal_uncommitted_workspaces[workspace_name] = { - "task_id": str(entry.get("task_id") or ""), - "attempt_id": str(entry.get("attempt_id") or ""), - } - self._gateway_attempts.pop(gateway_task_id, None) - - def probe_gateway_alive(self) -> bool: - """Liveness probe for the dispatch breaker: did the gateway answer? - - Any answer (even a non-2xx status) proves the transport is back; only - a transport-level failure keeps the campaign paused. - """ - - try: - self.config.http_runner( - "GET", - _gateway_path(self.config.ouroboros_url, "/api/health"), - timeout=15, - ) - except GatewayTransportError: - return False - except HttpStatusError: - return True - except Exception: # noqa: BLE001 - malformed body still means "answered" - return True - return True - - def _gateway_wait( - self, - body: Mapping[str, Any], - checkpoint: pathlib.Path, - *, - workspace_name: str = "", - task_id: str = "", - attempt_id: str = "", - ) -> Mapping[str, Any]: - requested_task_id = str(body.get("task_id") or "").strip() - owner_task_id = str(task_id) - owner_attempt_id = str(attempt_id) - # The gateway currently echoes the opaque caller task id. Register it - # before POST so a dropped response can still be treated as an - # admitted-or-unknown attempt and retained for manual reattachment. - pending_id = requested_task_id or ("pending-" + uuid.uuid4().hex) - idempotency_key = "cybergym-" + hashlib.sha256( - (pending_id + "\0" + str(body.get("actor_id") or "cybergym")).encode() - ).hexdigest() - self._gateway_attempts[pending_id] = { - "gateway_task_id": requested_task_id, - "status": "admission_pending", - "checkpoint": str(checkpoint), - "idempotency_key": idempotency_key, - "workspace_name": str(workspace_name), - "task_id": owner_task_id, - "attempt_id": owner_attempt_id, - } - try: - created = _unwrap_http_json( - self.config.http_runner( - "POST", - _gateway_path(self.config.ouroboros_url, "/api/tasks"), - body=body, - headers={"Idempotency-Key": idempotency_key}, - timeout=60, - ), - operation="Ouroboros task admission", - ) - except BaseException as exc: - rejected = _definitive_admission_rejection(exc) - status = "admission_rejected" if rejected else "admission_unknown" - entry = self._gateway_attempts.get(pending_id) - if entry is not None: - entry.update({"status": status, "error": type(exc).__name__}) - if rejected: - # A typed 4xx response is evidence that the gateway refused the - # request before scheduling it. Do not retain a phantom - # custody claim, but keep the redacted checkpoint for audit. - self._gateway_attempts.pop(pending_id, None) - _write_json( - checkpoint, - { - "gateway_task_id": requested_task_id or pending_id, - "status": status, - "custody_required": not rejected, - "idempotency_key": idempotency_key, - "error": type(exc).__name__, - }, - ) - if rejected: - raise GatewayAdmissionRejected(str(exc)) from exc - raise - task_id = str(created.get("task_id") or "").strip() - if not task_id or not _GATEWAY_TASK_ID.fullmatch(task_id): - self._gateway_attempts[pending_id]["status"] = "admission_unknown_response" - _write_json( - checkpoint, - { - "gateway_task_id": requested_task_id or pending_id, - "status": "admission_unknown_response", - "custody_required": True, - "idempotency_key": idempotency_key, - }, - ) - raise ExecutorFailure("Ouroboros gateway returned no task id") - if requested_task_id and task_id != requested_task_id: - self._gateway_attempts[pending_id].update( - {"gateway_task_id": task_id, "status": "admission_id_mismatch"} - ) - _write_json( - checkpoint, - { - "gateway_task_id": task_id, - "submitted_task_id": requested_task_id, - "status": "admission_id_mismatch", - "custody_required": True, - "idempotency_key": idempotency_key, - }, - ) - raise ExecutorFailure("Ouroboros gateway changed the submitted task id") - if pending_id != task_id: - self._gateway_attempts.pop(pending_id, None) - self._gateway_attempts[task_id] = { - "gateway_task_id": task_id, - "status": "submitted", - "checkpoint": str(checkpoint), - "idempotency_key": idempotency_key, - "workspace_name": str(workspace_name), - "task_id": owner_task_id, - "attempt_id": owner_attempt_id, - } - _write_json( - checkpoint, - { - "gateway_task_id": task_id, - "status": "submitted", - "idempotency_key": idempotency_key, - "body": {k: v for k, v in body.items() if k != "description"}, - }, - ) - # Two bounds, one active at a time: the queue-wait cap while the - # gateway still reports the task as not started, then the task - # deadline anchored at the first observed non-queued status. - queue_started = time.monotonic() - queue_wait_cap = queue_started + float(self.config.task_timeout_sec) - run_deadline: float | None = None - observed_start_at: str | None = None - latest: Mapping[str, Any] = created - cost_grace = _CostGraceTracker() - transport_deadline: float | None = None - finalization_grace_until: float | None = None - while True: - bound = run_deadline if run_deadline is not None else queue_wait_cap - if time.monotonic() >= bound: - # The worker is done and the server is finalizing artifacts: - # a finished, paid result is minutes away — wait for it (once, - # bounded) instead of cancelling it. - if finalization_grace_until is None and _gateway_finalizing(latest): - finalization_grace_until = time.monotonic() + FINALIZATION_GRACE_SEC - run_deadline = finalization_grace_until - _write_json(checkpoint, { - "gateway_task_id": task_id, - "status": _response_status(latest), - "result": dict(latest), - "deadline_basis": "finalization_grace", - "finalization_grace_sec": FINALIZATION_GRACE_SEC, - }) - continue - break - try: - latest = _unwrap_http_json( - self.config.http_runner( - "GET", - _gateway_path(self.config.ouroboros_url, "/api/tasks/" + urllib.parse.quote(task_id, safe="")), - timeout=60, - ), - operation="Ouroboros task status", - ) - except GatewayTransportError: - # A transient transport failure (an isolate event-loop stall - # starves the HTTP answer) must not kill a healthy paid task - # on the first error: ride it out within a bounded budget. - # Exhaustion re-raises so a dead gateway still produces the - # circuit-breaker row. - now = time.monotonic() - if now >= bound: - # The task's own deadline passed while the gateway was - # unreachable: stop polling and cancel it like a normal - # deadline exit instead of writing a transport row. - break - if transport_deadline is None: - transport_deadline = min( - bound, now + GATEWAY_TRANSPORT_RETRY_BUDGET_SEC - ) - if now >= transport_deadline: - raise - self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) - continue - transport_deadline = None - returned_id = str(latest.get("task_id") or "").strip() - if returned_id and returned_id != task_id: - raise ExecutorFailure("Ouroboros status response belongs to a different task") - status = _response_status(latest) - if run_deadline is None and status not in _QUEUED_GATEWAY_STATUSES: - run_deadline = ( - time.monotonic() - + float(self.config.task_timeout_sec) - + TASK_DEADLINE_GRACE_SEC - ) - observed_start_at = time.strftime( - "%Y-%m-%dT%H:%M:%SZ", time.gmtime() - ) - frame = { - "gateway_task_id": task_id, - "status": status, - "result": dict(latest), - "deadline_basis": ( - "finalization_grace" if finalization_grace_until is not None - else "observed_start" if run_deadline is not None else "queue_wait_cap" - ), - } - if observed_start_at is not None: - frame["observed_start_at"] = observed_start_at - _write_json(checkpoint, frame) - if status in _SETTLED: - # Root post-task accounting can publish ``completed`` before - # its durable cost roll-up is final; only the bounded - # abandoned-residue grace (cybergym_wire) releases such a - # frame early, with the residue disclosed on it. - if status == "completed" and _cost_is_pending(latest): - accepted = cost_grace.accept( - latest, - now=time.monotonic(), - wall_now=time.time(), - ) - if accepted is None: - self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) - continue - latest = accepted - _write_json(checkpoint, {"gateway_task_id": task_id, "status": status, "result": dict(latest)}) - self._terminalize_gateway_attempt(task_id) - return latest - self.config.sleep(max(0.5, float(self.config.poll_interval_sec))) - # The task may still be running after the local wait expires. Ask the - # gateway to stop it and retain the original attempt until a terminal - # custody response is observed; never return a reusable task id here. - return self._cancel_gateway_task(task_id, checkpoint) - def _submit_final( self, task: TaskSpec, task_dir: pathlib.Path, container_name: str ) -> tuple[dict[str, Any], str, str]: diff --git a/docs/architecture/01-high-level-architecture.md b/docs/architecture/01-high-level-architecture.md index 5620dc4e9..12e02ddfb 100644 --- a/docs/architecture/01-high-level-architecture.md +++ b/docs/architecture/01-high-level-architecture.md @@ -481,7 +481,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de ### Devtools boundary -`devtools/` (including `devtools/benchmarks/cybergym/`) lives outside the runtime and package discovery: no runtime imports, normal review, artifacts in an external output root; a sentinel-marked isolated root suppresses rotation warnings. `devtools/e2e_live/` is the live E2E stand: K staggered isolated real servers running the owner-shaped scenarios SM1 (a design-system-consistent brand-accent change landed as a reviewed release, the product's own review policy shaping the work), SW1 and SK1 with acceptance over durable artifacts and a browser probe, admitted through the same seed gate and manifest seams as the benchmark launchers (DEVELOPMENT "Live E2E stand"). +`devtools/` (including `devtools/benchmarks/cybergym/`) lives outside the runtime and package discovery: no runtime imports, normal review, artifacts in an external output root; a sentinel-marked isolated root suppresses rotation warnings. CyberGym keeps budget/result settlement in `cybergym_adapter`; `cybergym_custody._CustodyMixin` owns gateway admission, normal and cancelled waits, and their shared terminal-custody transfer, while `cybergym_lifecycle._LifecycleMixin` owns startup, official-verifier delivery and cleanup. Both wait paths retain their distinct bounds around the same attempt identity rather than creating another scheduler. `devtools/e2e_live/` is the live E2E stand: K staggered isolated real servers running the owner-shaped scenarios SM1 (a design-system-consistent brand-accent change landed as a reviewed release, the product's own review policy shaping the work), SW1 and SK1 with acceptance over durable artifacts and a browser probe, admitted through the same seed gate and manifest seams as the benchmark launchers (DEVELOPMENT "Live E2E stand"). ### Gateway Boundary v1 diff --git a/docs/v7next/DATA_LAYOUT_INVENTORY.md b/docs/v7next/DATA_LAYOUT_INVENTORY.md index 46185f887..dce4105c2 100644 --- a/docs/v7next/DATA_LAYOUT_INVENTORY.md +++ b/docs/v7next/DATA_LAYOUT_INVENTORY.md @@ -2,7 +2,7 @@ Machine extraction of the `docs/ARCHITECTURE.md` "Data layout (`~/Ouroboros/`)" tree — the durable-file orientation carrier (this tree's counterpart of the reference PERSISTENCE_OWNERS derivation checklist) — regenerated by `python scripts/regenerate_inventories.py`. Do not edit. Every entry is probed against reality: repo entries must exist as tracked paths; data-plane entries must appear as a literal in the runtime sources that construct them. A durable file renamed or removed in code while its tree row survives = red (`tests/test_generated_inventories.py`). -Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 568-657; UTF-8 SHA-256 `391c814a062ec74ce6c2ba11f0f9e1c348a75c2078ab1b0303fb1670a86eaa25`. +Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 568-657; UTF-8 SHA-256 `c6bac8219d5df86a2d1441e7b6cbf0f43edfa464745091b505c0279357004e69`. - entries: **78** (code-ref: 71, repo-dir: 6, repo-path: 1) diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index 827bba8a2..c28f8bf61 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -99,13 +99,10 @@ 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/cybergym/cybergym_reconcile.py": "CyberGym recovery joins existing checkpoint, result, claim and cleanup authority without repeating an agent; one recovery owner retains that crash-window contract.", "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, @@ -181,7 +178,8 @@ BAND_PATHS = { "tests/test_available_subagents_runtime.py": "Configured-session route and legacy custody regressions retained after removing compulsory source-request production tests.", "tests/test_build_scripts.py": None, "tests/test_commit_gate.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_cybergym_dispatch.py": "CyberGym dispatch tests cover completion-order admission, transient gateway pauses and budget-refusal recovery through one existing fake campaign harness.", + "tests/test_cybergym_docker.py": "CyberGym workspace custody tests cover atomic gateway transfer, recovery and durable-result acknowledgement using the same attested fake container fixtures.", "tests/test_deep_review_slot.py": "\u04243 deep self-review suite: the deep_review row/endpoint half and the three-delivery half (retrieving executors, coverage, header, custody, availability) share one fixture set (repo/drive roots, scripted LLM, fake session executor); split at the config-vs-delivery seam once the row contract stops moving, not by size.", "tests/test_delegate_answer.py": "Entered the band by the #204 escalation-route pins (walk-up, schema and expiry-note source pins) on top of the phase-B interaction suite; one coherent delegated-question surface, split only when a natural seam appears", "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.", diff --git a/tests/test_cybergym_executor_wire.py b/tests/test_cybergym_executor_wire.py index 06ee84376..62192360f 100644 --- a/tests/test_cybergym_executor_wire.py +++ b/tests/test_cybergym_executor_wire.py @@ -2,7 +2,8 @@ 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. +validation and campaign cost accounting. Gateway custody and deadline tests +live in ``test_cybergym_gateway_custody.py``. Shared fixtures (``_config``, ``dataclasses_replace``) are imported from the original module; executor-lifecycle tests remain there. """ @@ -13,7 +14,6 @@ import gzip import hashlib import json import pathlib -import time as _time import pytest @@ -478,170 +478,6 @@ def test_delivery_checkpoint_prevents_duplicate_submit_and_verify( assert delivery["final_poc_sha256"] == digest -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) @@ -1030,343 +866,6 @@ def test_post_admission_status_error_is_not_reclassified_as_zero_cost( assert projection.can_dispatch is True -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 test_gateway_poll_rides_out_transient_transport_errors(tmp_path): - # A ~95 s isolate event-loop stall used to kill a healthy paid task via a - # single 60 s poll timeout. The poll now retries transport failures - # within a bounded budget; the terminal answer after the stall wins. - config = _config(tmp_path, provider_probe=False, task_timeout_sec=60) - task_id = "cybergym-transient-stall" - calls = [] - - def http(method, _url, **_kwargs): - calls.append(method) - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - if calls.count("GET") <= 3: - raise GatewayTransportError("HTTP GET transport failed") - return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - result = executor._gateway_wait( # noqa: SLF001 - transport recovery contract - {"task_id": task_id, "description": "test"}, - config.run_root / "checkpoint.json", - ) - - assert result["status"] == "completed" - assert calls == ["POST", "GET", "GET", "GET", "GET"] - - -def test_gateway_poll_transport_budget_exhaustion_still_fails(tmp_path, monkeypatch): - # The retry budget is bounded: a gateway that never answers within it - # still produces the circuit-breaker transport row. - config = _config(tmp_path, provider_probe=False, task_timeout_sec=3600) - task_id = "cybergym-dead-gateway" - monkeypatch.setattr( - "devtools.benchmarks.cybergym.cybergym_lifecycle.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", - 0.0, - ) - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - raise GatewayTransportError("HTTP GET transport failed") - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - with pytest.raises(GatewayTransportError): - executor._gateway_wait( # noqa: SLF001 - transport recovery contract - {"task_id": task_id, "description": "test"}, - config.run_root / "checkpoint.json", - ) - - -def test_cancel_rides_out_transient_transport_errors(tmp_path): - # The cancel intent usually lands server-side and only the response is - # starved by an isolate stall; a duplicate POST is idempotent. The - # bounded retry turns a stall into a delayed cancel instead of a - # written-off paid attempt. - config = _config(tmp_path, poll_interval_sec=0) - task_id = "cybergym-cancel-stall" - 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 = [] - - def http(method, _url, **_kwargs): - calls.append(method) - if method == "POST" and calls.count("POST") <= 2: - raise GatewayTransportError("HTTP POST transport failed") - if method == "POST": - return {"status_code": 200, "body": {"task_id": task_id, "status": "cancel_requested"}} - return {"status_code": 200, "body": terminal} - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - 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", "POST", "POST", "GET"] - assert task_id not in executor._gateway_attempts - - -def test_cancel_custody_get_rides_out_transient_transport_errors(tmp_path): - config = _config(tmp_path, poll_interval_sec=0) - task_id = "cybergym-cancel-poll-stall" - terminal = { - "task_id": task_id, - "status": "failed", - "cost_usd": 0.2, - "cost_final": True, - } - calls = [] - - def http(method, _url, **_kwargs): - calls.append(method) - if method == "POST": - return { - "status_code": 202, - "body": {"task_id": task_id, "status": "cancel_requested"}, - } - if calls.count("GET") <= 3: - raise GatewayTransportError("HTTP GET transport failed") - return {"status_code": 200, "body": terminal} - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion - "gateway_task_id": task_id, - "status": "submitted", - } - result = executor._cancel_gateway_task( # noqa: SLF001 - task_id, config.run_root / "checkpoint.json" - ) - - assert result == terminal - assert calls == ["POST", "GET", "GET", "GET", "GET"] - assert task_id not in executor._gateway_attempts - - -def test_cancel_custody_get_transport_budget_exhaustion_keeps_custody( - tmp_path, monkeypatch -): - config = _config(tmp_path, poll_interval_sec=0) - task_id = "cybergym-cancel-poll-dead" - monkeypatch.setattr( - "devtools.benchmarks.cybergym.cybergym_custody.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", - 0.0, - ) - - def http(method, _url, **_kwargs): - if method == "POST": - return { - "status_code": 202, - "body": {"task_id": task_id, "status": "cancel_requested"}, - } - raise GatewayTransportError("HTTP GET transport failed") - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion - "gateway_task_id": task_id, - "status": "submitted", - } - checkpoint = config.run_root / "checkpoint.json" - with pytest.raises(GatewayTransportError): - executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001 - - assert task_id in executor._gateway_attempts - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["status"] == "cancel_poll_error" - assert saved["cancel_error"] == "GatewayTransportError" - - -def test_cancel_transport_budget_exhaustion_keeps_custody(tmp_path, monkeypatch): - # A cancel that never gets through within the budget keeps the original - # fail-closed behaviour: typed checkpoint evidence and retained custody. - config = _config(tmp_path, poll_interval_sec=0) - task_id = "cybergym-cancel-dead" - monkeypatch.setattr( - "devtools.benchmarks.cybergym.cybergym_custody.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", - 0.0, - ) - - def http(method, _url, **_kwargs): - raise GatewayTransportError("HTTP POST transport failed") - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion - "gateway_task_id": task_id, - "status": "submitted", - } - checkpoint = config.run_root / "checkpoint.json" - with pytest.raises(GatewayTransportError): - executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001 - - assert task_id in executor._gateway_attempts - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["status"] == "cancel_request_failed" - assert saved["cancel_error"] == "GatewayTransportError" - - -def test_gateway_poll_transport_error_at_deadline_goes_to_cancel(tmp_path, monkeypatch): - # A transport failure riding into the task's own deadline must exit to the - # cancel path, not raise a transport row: the deadline, not the network, - # decides the task's fate. - config = _config(tmp_path, provider_probe=False, task_timeout_sec=1) - task_id = "cybergym-deadline-stall" - calls = [] - - def http(method, _url, **_kwargs): - calls.append(method) - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - raise GatewayTransportError("HTTP GET transport failed") - - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) - ) - sentinel = {"status": "cancelled", "via": "cancel_path"} - monkeypatch.setattr( - executor, "_cancel_gateway_task", lambda *_a, **_k: dict(sentinel) - ) - result = executor._gateway_wait( # noqa: SLF001 - transport recovery contract - {"task_id": task_id, "description": "test"}, - config.run_root / "checkpoint.json", - ) - - assert result == sentinel - assert "GET" in calls - - -def test_cancel_custody_window_covers_deadline_wave_settle(tmp_path, monkeypatch): - # The post-cancel custody poll must outlast a deadline-wave settle: the - # old ~34 s bound wrote off paid tasks whose cancel had already landed. - config = _config(tmp_path, poll_interval_sec=3.0) - task_id = "cybergym-cancel-window" - captured = {} - - def http(method, _url, **_kwargs): - return { - "status_code": 200, - "body": {"task_id": task_id, "status": "cancel_requested"}, - } - - executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http)) - executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion - "gateway_task_id": task_id, - "status": "submitted", - } - - def fake_custody(*_args, **kwargs): - captured["custody_seconds"] = kwargs.get("custody_seconds") - return {} - - monkeypatch.setattr(executor, "_poll_gateway_custody", fake_custody) - executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001 - - assert captured["custody_seconds"] >= 300.0 - - def test_explicit_final_with_excluded_vul_exit_and_missing_fix_records_failure(): """A determinate vul-excluded failure binds without a fix-side code.""" from devtools.benchmarks.cybergym.cybergym_adapter import build_task_result_row @@ -1414,138 +913,6 @@ def test_explicit_final_with_missing_vul_exit_still_refused(): ) -class _FakeClock: - """Deterministic ``time`` stand-in for the gateway wait loop.""" - - def __init__(self) -> None: - self.now = 10_000.0 - - def monotonic(self) -> float: - return self.now - - def time(self) -> float: - return 1_700_000_000.0 + self.now - - def sleep(self, seconds: float) -> None: - self.now += float(seconds) - - strftime = staticmethod(_time.strftime) - gmtime = staticmethod(_time.gmtime) - - -def _deadline_executor(tmp_path, monkeypatch, http, *, task_timeout_sec): - from devtools.benchmarks.cybergym import cybergym_lifecycle - - clock = _FakeClock() - monkeypatch.setattr(cybergym_lifecycle, "time", clock) - config = _config( - tmp_path, - provider_probe=False, - task_timeout_sec=task_timeout_sec, - poll_interval_sec=1.0, - ) - executor = CyberGymExecutor( - dataclasses_replace(config, http_runner=http, sleep=clock.sleep) - ) - sentinel = {"status": "cancelled", "via": "cancel_path"} - monkeypatch.setattr( - executor, "_cancel_gateway_task", lambda *_a, **_k: dict(sentinel) - ) - return executor, clock, config.run_root / "checkpoint.json", sentinel - - -def test_gateway_deadline_is_anchored_at_observed_run_start(tmp_path, monkeypatch): - # full1507: tasks queued ~1 h behind a finalization backlog were cancelled - # after ~1 h of runtime because the deadline was anchored at submit. The - # clock must start when the gateway first reports a non-queued status. - task_id = "cybergym-observed-start" - # Each poll costs 6 s + 1 s poll interval: ~14 s queued (inside the 20 s - # queue cap), then the run starts at ~21 s, past a submit-anchored 20 s - # deadline, and completes at ~35 s -- inside the observed-start deadline. - frames = iter( - ( - {"task_id": task_id, "status": "scheduled"}, - {"task_id": task_id, "status": "scheduled"}, - {"task_id": task_id, "status": "running"}, - {"task_id": task_id, "status": "running"}, - { - "task_id": task_id, - "status": "completed", - "result": {"cost_final": True}, - }, - ) - ) - clock_ref = {} - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - clock_ref["clock"].now += 6.0 - return next(frames) - - executor, clock, checkpoint, sentinel = _deadline_executor( - tmp_path, monkeypatch, http, task_timeout_sec=20 - ) - submitted_at = clock.now - clock_ref["clock"] = clock - result = executor._gateway_wait( # noqa: SLF001 - deadline contract - {"task_id": task_id, "description": "test"}, checkpoint - ) - - assert result["status"] == "completed", "queued time was charged against the deadline" - assert clock.now - submitted_at > 20, "scenario must outlive a submit-anchored deadline" - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["deadline_basis"] == "observed_start" - assert saved["observed_start_at"].endswith("Z") - - -def test_gateway_queue_wait_cap_still_bounds_a_never_started_task(tmp_path, monkeypatch): - task_id = "cybergym-queue-cap" - polls = [] - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - polls.append(_kwargs.get("timeout")) - return {"task_id": task_id, "status": "scheduled"} - - executor, clock, checkpoint, sentinel = _deadline_executor( - tmp_path, monkeypatch, http, task_timeout_sec=10 - ) - result = executor._gateway_wait( # noqa: SLF001 - deadline contract - {"task_id": task_id, "description": "test"}, checkpoint - ) - - assert result == sentinel - assert 9 <= len(polls) <= 12 - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["deadline_basis"] == "queue_wait_cap" - assert "observed_start_at" not in saved - - -def test_gateway_run_deadline_carries_grace_behind_server_ceiling(tmp_path, monkeypatch): - from devtools.benchmarks.cybergym.cybergym_lifecycle import TASK_DEADLINE_GRACE_SEC - - task_id = "cybergym-run-deadline" - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - return {"task_id": task_id, "status": "running"} - - executor, clock, checkpoint, sentinel = _deadline_executor( - tmp_path, monkeypatch, http, task_timeout_sec=10 - ) - started = clock.now - result = executor._gateway_wait( # noqa: SLF001 - deadline contract - {"task_id": task_id, "description": "test"}, checkpoint - ) - - assert result == sentinel - elapsed = clock.now - started - assert 10 + TASK_DEADLINE_GRACE_SEC <= elapsed < 10 + TASK_DEADLINE_GRACE_SEC + 2.0 - - # r9 (2026-09-04): nine runs that finished on their own, 40-90 min before the # deadline, with a final message and no marker, were typed infrastructure # (FinalPocRefused) because the runtime marked execution `degraded` over the @@ -1628,97 +995,3 @@ def test_fair_completion_reads_outermost_execution_axis(): "result": {"outcome_axes": {"execution": {"status": "failed"}}}} assert _gateway_fair_completion(payload) == (True, "execution_ok") assert _gateway_fair_completion({"result": {"task_result": _envelope({"status": "failed"})}}) == (False, "execution_failed") - -# r9 (2026-09-04): nine finished, paid tasks were cancelled by the launcher at -# deadline+grace while the server was still finalizing their workspace -# artifacts (projected `running` / `artifact_status=finalizing`); the cancel path -# re-ran the same finalization and the 300 s custody window expired 1-15 min -# before the completed results landed. A finalizing frame now buys the task a -# bounded finalization grace on both the wait and the custody paths. - - -def test_gateway_wait_grants_finalization_grace_instead_of_cancelling(tmp_path, monkeypatch): - from devtools.benchmarks.cybergym.cybergym_wire import FINALIZATION_GRACE_SEC - - task_id = "cybergym-finalizing" - finalizing = { - "task_id": task_id, "status": "running", "artifact_status": "finalizing", - "child_status": "completed", "cost_final": True, - } - state = {"polls": 0} - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - state["polls"] += 1 - # Runs to its deadline, then finalizes for a while, then delivers. - if state["polls"] < 12: - return {"task_id": task_id, "status": "running"} - if state["polls"] < 20: - return dict(finalizing) - return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} - - executor, clock, checkpoint, sentinel = _deadline_executor( - tmp_path, monkeypatch, http, task_timeout_sec=5 - ) - result = executor._gateway_wait( # noqa: SLF001 - deadline contract - {"task_id": task_id, "description": "test"}, checkpoint - ) - assert result["status"] == "completed", "a finalizing task must not be cancelled" - assert result != sentinel - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["status"] == "completed" - assert FINALIZATION_GRACE_SEC > 60 - - -def test_gateway_wait_finalization_grace_is_granted_once(tmp_path, monkeypatch): - from devtools.benchmarks.cybergym import cybergym_lifecycle - - task_id = "cybergym-finalizing-forever" - - def http(method, _url, **_kwargs): - if method == "POST": - return {"task_id": task_id, "status": "scheduled"} - return {"task_id": task_id, "status": "running", "artifact_status": "finalizing", "child_status": "completed"} - - monkeypatch.setattr(cybergym_lifecycle, "FINALIZATION_GRACE_SEC", 20.0) - executor, clock, checkpoint, sentinel = _deadline_executor( - tmp_path, monkeypatch, http, task_timeout_sec=5 - ) - started = clock.now - result = executor._gateway_wait( # noqa: SLF001 - deadline contract - {"task_id": task_id, "description": "test"}, checkpoint - ) - assert result == sentinel, "finalization that outlives the grace is still cancelled" - elapsed = clock.now - started - assert 5 + 300 + 20 <= elapsed < 5 + 300 + 20 + 3 - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["deadline_basis"] == "finalization_grace" - - -def test_custody_poll_outlasts_its_window_while_the_task_is_finalizing(tmp_path, monkeypatch): - from devtools.benchmarks.cybergym import cybergym_custody - - task_id = "cybergym-custody-finalizing" - clock = _FakeClock() - monkeypatch.setattr(cybergym_custody, "time", clock) - monkeypatch.setattr(cybergym_custody, "FINALIZATION_GRACE_SEC", 120.0) - state = {"polls": 0} - - def http(method, _url, **_kwargs): - state["polls"] += 1 - clock.now += 10.0 - if state["polls"] < 8: # 80 s of finalizing: past a 30 s custody window - return {"task_id": task_id, "status": "running", "artifact_status": "finalizing", "child_status": "completed"} - return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} - - config = _config(tmp_path, provider_probe=False, poll_interval_sec=1.0) - executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http, sleep=clock.sleep)) - monkeypatch.setattr(executor, "_terminalize_gateway_attempt", lambda *_a, **_k: None) - checkpoint = config.run_root / "custody.json" - result = executor._poll_gateway_custody( # noqa: SLF001 - custody contract - task_id, checkpoint, cancel_response=None, cancel_status_code=503, custody_seconds=30.0 - ) - assert result["status"] == "completed" - saved = json.loads(checkpoint.read_text(encoding="utf-8")) - assert saved["status"] == "completed" diff --git a/tests/test_cybergym_gateway_custody.py b/tests/test_cybergym_gateway_custody.py new file mode 100644 index 000000000..4b045c66f --- /dev/null +++ b/tests/test_cybergym_gateway_custody.py @@ -0,0 +1,752 @@ +"""Gateway admission, cancellation, transport and deadline custody tests. + +These drive the normal and cancelled gateway paths through the concrete executor, +using the existing shared configuration fixture and a deterministic wait clock. +""" + +from __future__ import annotations + +import json +import time as _time + +import pytest + +from devtools.benchmarks.cybergym.cybergym_executor import ( + CyberGymExecutor, + ExecutorFailure, +) +from devtools.benchmarks.cybergym.cybergym_wire import GatewayTransportError +from tests.test_cybergym_executor import _config, dataclasses_replace + +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 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 test_gateway_poll_rides_out_transient_transport_errors(tmp_path): + # A ~95 s isolate event-loop stall used to kill a healthy paid task via a + # single 60 s poll timeout. The poll now retries transport failures + # within a bounded budget; the terminal answer after the stall wins. + config = _config(tmp_path, provider_probe=False, task_timeout_sec=60) + task_id = "cybergym-transient-stall" + calls = [] + + def http(method, _url, **_kwargs): + calls.append(method) + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + if calls.count("GET") <= 3: + raise GatewayTransportError("HTTP GET transport failed") + return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + result = executor._gateway_wait( # noqa: SLF001 - transport recovery contract + {"task_id": task_id, "description": "test"}, + config.run_root / "checkpoint.json", + ) + + assert result["status"] == "completed" + assert calls == ["POST", "GET", "GET", "GET", "GET"] + + +def test_gateway_poll_transport_budget_exhaustion_still_fails(tmp_path, monkeypatch): + # The retry budget is bounded: a gateway that never answers within it + # still produces the circuit-breaker transport row. + config = _config(tmp_path, provider_probe=False, task_timeout_sec=3600) + task_id = "cybergym-dead-gateway" + monkeypatch.setattr( + "devtools.benchmarks.cybergym.cybergym_custody.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", + 0.0, + ) + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + raise GatewayTransportError("HTTP GET transport failed") + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + with pytest.raises(GatewayTransportError): + executor._gateway_wait( # noqa: SLF001 - transport recovery contract + {"task_id": task_id, "description": "test"}, + config.run_root / "checkpoint.json", + ) + + +def test_cancel_rides_out_transient_transport_errors(tmp_path): + # The cancel intent usually lands server-side and only the response is + # starved by an isolate stall; a duplicate POST is idempotent. The + # bounded retry turns a stall into a delayed cancel instead of a + # written-off paid attempt. + config = _config(tmp_path, poll_interval_sec=0) + task_id = "cybergym-cancel-stall" + 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 = [] + + def http(method, _url, **_kwargs): + calls.append(method) + if method == "POST" and calls.count("POST") <= 2: + raise GatewayTransportError("HTTP POST transport failed") + if method == "POST": + return {"status_code": 200, "body": {"task_id": task_id, "status": "cancel_requested"}} + return {"status_code": 200, "body": terminal} + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + 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", "POST", "POST", "GET"] + assert task_id not in executor._gateway_attempts + + +def test_cancel_custody_get_rides_out_transient_transport_errors(tmp_path): + config = _config(tmp_path, poll_interval_sec=0) + task_id = "cybergym-cancel-poll-stall" + terminal = { + "task_id": task_id, + "status": "failed", + "cost_usd": 0.2, + "cost_final": True, + } + calls = [] + + def http(method, _url, **_kwargs): + calls.append(method) + if method == "POST": + return { + "status_code": 202, + "body": {"task_id": task_id, "status": "cancel_requested"}, + } + if calls.count("GET") <= 3: + raise GatewayTransportError("HTTP GET transport failed") + return {"status_code": 200, "body": terminal} + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion + "gateway_task_id": task_id, + "status": "submitted", + } + result = executor._cancel_gateway_task( # noqa: SLF001 + task_id, config.run_root / "checkpoint.json" + ) + + assert result == terminal + assert calls == ["POST", "GET", "GET", "GET", "GET"] + assert task_id not in executor._gateway_attempts + + +def test_cancel_custody_get_transport_budget_exhaustion_keeps_custody( + tmp_path, monkeypatch +): + config = _config(tmp_path, poll_interval_sec=0) + task_id = "cybergym-cancel-poll-dead" + monkeypatch.setattr( + "devtools.benchmarks.cybergym.cybergym_custody.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", + 0.0, + ) + + def http(method, _url, **_kwargs): + if method == "POST": + return { + "status_code": 202, + "body": {"task_id": task_id, "status": "cancel_requested"}, + } + raise GatewayTransportError("HTTP GET transport failed") + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion + "gateway_task_id": task_id, + "status": "submitted", + } + checkpoint = config.run_root / "checkpoint.json" + with pytest.raises(GatewayTransportError): + executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001 + + assert task_id in executor._gateway_attempts + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["status"] == "cancel_poll_error" + assert saved["cancel_error"] == "GatewayTransportError" + + +def test_cancel_transport_budget_exhaustion_keeps_custody(tmp_path, monkeypatch): + # A cancel that never gets through within the budget keeps the original + # fail-closed behaviour: typed checkpoint evidence and retained custody. + config = _config(tmp_path, poll_interval_sec=0) + task_id = "cybergym-cancel-dead" + monkeypatch.setattr( + "devtools.benchmarks.cybergym.cybergym_custody.GATEWAY_TRANSPORT_RETRY_BUDGET_SEC", + 0.0, + ) + + def http(method, _url, **_kwargs): + raise GatewayTransportError("HTTP POST transport failed") + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion + "gateway_task_id": task_id, + "status": "submitted", + } + checkpoint = config.run_root / "checkpoint.json" + with pytest.raises(GatewayTransportError): + executor._cancel_gateway_task(task_id, checkpoint) # noqa: SLF001 + + assert task_id in executor._gateway_attempts + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["status"] == "cancel_request_failed" + assert saved["cancel_error"] == "GatewayTransportError" + + +def test_gateway_poll_transport_error_at_deadline_goes_to_cancel(tmp_path, monkeypatch): + # A transport failure riding into the task's own deadline must exit to the + # cancel path, not raise a transport row: the deadline, not the network, + # decides the task's fate. + config = _config(tmp_path, provider_probe=False, task_timeout_sec=1) + task_id = "cybergym-deadline-stall" + calls = [] + + def http(method, _url, **_kwargs): + calls.append(method) + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + raise GatewayTransportError("HTTP GET transport failed") + + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=lambda _seconds: None) + ) + sentinel = {"status": "cancelled", "via": "cancel_path"} + monkeypatch.setattr( + executor, "_cancel_gateway_task", lambda *_a, **_k: dict(sentinel) + ) + result = executor._gateway_wait( # noqa: SLF001 - transport recovery contract + {"task_id": task_id, "description": "test"}, + config.run_root / "checkpoint.json", + ) + + assert result == sentinel + assert "GET" in calls + + +def test_cancel_custody_window_covers_deadline_wave_settle(tmp_path, monkeypatch): + # The post-cancel custody poll must outlast a deadline-wave settle: the + # old ~34 s bound wrote off paid tasks whose cancel had already landed. + config = _config(tmp_path, poll_interval_sec=3.0) + task_id = "cybergym-cancel-window" + captured = {} + + def http(method, _url, **_kwargs): + return { + "status_code": 200, + "body": {"task_id": task_id, "status": "cancel_requested"}, + } + + executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http)) + executor._gateway_attempts[task_id] = { # noqa: SLF001 - custody assertion + "gateway_task_id": task_id, + "status": "submitted", + } + + def fake_custody(*_args, **kwargs): + captured["custody_seconds"] = kwargs.get("custody_seconds") + return {} + + monkeypatch.setattr(executor, "_poll_gateway_custody", fake_custody) + executor._cancel_gateway_task(task_id, config.run_root / "checkpoint.json") # noqa: SLF001 + + assert captured["custody_seconds"] >= 300.0 + + + + +class _FakeClock: + """Deterministic ``time`` stand-in for the gateway wait loop.""" + + def __init__(self) -> None: + self.now = 10_000.0 + + def monotonic(self) -> float: + return self.now + + def time(self) -> float: + return 1_700_000_000.0 + self.now + + def sleep(self, seconds: float) -> None: + self.now += float(seconds) + + strftime = staticmethod(_time.strftime) + gmtime = staticmethod(_time.gmtime) + + +def _deadline_executor(tmp_path, monkeypatch, http, *, task_timeout_sec): + from devtools.benchmarks.cybergym import cybergym_custody + + clock = _FakeClock() + monkeypatch.setattr(cybergym_custody, "time", clock) + config = _config( + tmp_path, + provider_probe=False, + task_timeout_sec=task_timeout_sec, + poll_interval_sec=1.0, + ) + executor = CyberGymExecutor( + dataclasses_replace(config, http_runner=http, sleep=clock.sleep) + ) + sentinel = {"status": "cancelled", "via": "cancel_path"} + monkeypatch.setattr( + executor, "_cancel_gateway_task", lambda *_a, **_k: dict(sentinel) + ) + return executor, clock, config.run_root / "checkpoint.json", sentinel + + +def test_gateway_deadline_is_anchored_at_observed_run_start(tmp_path, monkeypatch): + # full1507: tasks queued ~1 h behind a finalization backlog were cancelled + # after ~1 h of runtime because the deadline was anchored at submit. The + # clock must start when the gateway first reports a non-queued status. + task_id = "cybergym-observed-start" + # Each poll costs 6 s + 1 s poll interval: ~14 s queued (inside the 20 s + # queue cap), then the run starts at ~21 s, past a submit-anchored 20 s + # deadline, and completes at ~35 s -- inside the observed-start deadline. + frames = iter( + ( + {"task_id": task_id, "status": "scheduled"}, + {"task_id": task_id, "status": "scheduled"}, + {"task_id": task_id, "status": "running"}, + {"task_id": task_id, "status": "running"}, + { + "task_id": task_id, + "status": "completed", + "result": {"cost_final": True}, + }, + ) + ) + clock_ref = {} + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + clock_ref["clock"].now += 6.0 + return next(frames) + + executor, clock, checkpoint, sentinel = _deadline_executor( + tmp_path, monkeypatch, http, task_timeout_sec=20 + ) + submitted_at = clock.now + clock_ref["clock"] = clock + result = executor._gateway_wait( # noqa: SLF001 - deadline contract + {"task_id": task_id, "description": "test"}, checkpoint + ) + + assert result["status"] == "completed", "queued time was charged against the deadline" + assert clock.now - submitted_at > 20, "scenario must outlive a submit-anchored deadline" + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["deadline_basis"] == "observed_start" + assert saved["observed_start_at"].endswith("Z") + + +def test_gateway_queue_wait_cap_still_bounds_a_never_started_task(tmp_path, monkeypatch): + task_id = "cybergym-queue-cap" + polls = [] + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + polls.append(_kwargs.get("timeout")) + return {"task_id": task_id, "status": "scheduled"} + + executor, clock, checkpoint, sentinel = _deadline_executor( + tmp_path, monkeypatch, http, task_timeout_sec=10 + ) + result = executor._gateway_wait( # noqa: SLF001 - deadline contract + {"task_id": task_id, "description": "test"}, checkpoint + ) + + assert result == sentinel + assert 9 <= len(polls) <= 12 + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["deadline_basis"] == "queue_wait_cap" + assert "observed_start_at" not in saved + + +def test_gateway_run_deadline_carries_grace_behind_server_ceiling(tmp_path, monkeypatch): + from devtools.benchmarks.cybergym.cybergym_custody import TASK_DEADLINE_GRACE_SEC + + task_id = "cybergym-run-deadline" + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + return {"task_id": task_id, "status": "running"} + + executor, clock, checkpoint, sentinel = _deadline_executor( + tmp_path, monkeypatch, http, task_timeout_sec=10 + ) + started = clock.now + result = executor._gateway_wait( # noqa: SLF001 - deadline contract + {"task_id": task_id, "description": "test"}, checkpoint + ) + + assert result == sentinel + elapsed = clock.now - started + assert 10 + TASK_DEADLINE_GRACE_SEC <= elapsed < 10 + TASK_DEADLINE_GRACE_SEC + 2.0 + + + + +# r9 (2026-09-04): nine finished, paid tasks were cancelled by the launcher at +# deadline+grace while the server was still finalizing their workspace +# artifacts (projected `running` / `artifact_status=finalizing`); the cancel path +# re-ran the same finalization and the 300 s custody window expired 1-15 min +# before the completed results landed. A finalizing frame now buys the task a +# bounded finalization grace on both the wait and the custody paths. + + +def test_gateway_wait_grants_finalization_grace_instead_of_cancelling(tmp_path, monkeypatch): + from devtools.benchmarks.cybergym.cybergym_wire import FINALIZATION_GRACE_SEC + + task_id = "cybergym-finalizing" + finalizing = { + "task_id": task_id, "status": "running", "artifact_status": "finalizing", + "child_status": "completed", "cost_final": True, + } + state = {"polls": 0} + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + state["polls"] += 1 + # Runs to its deadline, then finalizes for a while, then delivers. + if state["polls"] < 12: + return {"task_id": task_id, "status": "running"} + if state["polls"] < 20: + return dict(finalizing) + return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} + + executor, clock, checkpoint, sentinel = _deadline_executor( + tmp_path, monkeypatch, http, task_timeout_sec=5 + ) + result = executor._gateway_wait( # noqa: SLF001 - deadline contract + {"task_id": task_id, "description": "test"}, checkpoint + ) + assert result["status"] == "completed", "a finalizing task must not be cancelled" + assert result != sentinel + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["status"] == "completed" + assert FINALIZATION_GRACE_SEC > 60 + + +def test_gateway_wait_finalization_grace_is_granted_once(tmp_path, monkeypatch): + from devtools.benchmarks.cybergym import cybergym_custody + + task_id = "cybergym-finalizing-forever" + + def http(method, _url, **_kwargs): + if method == "POST": + return {"task_id": task_id, "status": "scheduled"} + return {"task_id": task_id, "status": "running", "artifact_status": "finalizing", "child_status": "completed"} + + monkeypatch.setattr(cybergym_custody, "FINALIZATION_GRACE_SEC", 20.0) + executor, clock, checkpoint, sentinel = _deadline_executor( + tmp_path, monkeypatch, http, task_timeout_sec=5 + ) + started = clock.now + result = executor._gateway_wait( # noqa: SLF001 - deadline contract + {"task_id": task_id, "description": "test"}, checkpoint + ) + assert result == sentinel, "finalization that outlives the grace is still cancelled" + elapsed = clock.now - started + assert 5 + 300 + 20 <= elapsed < 5 + 300 + 20 + 3 + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["deadline_basis"] == "finalization_grace" + + +def test_custody_poll_outlasts_its_window_while_the_task_is_finalizing(tmp_path, monkeypatch): + from devtools.benchmarks.cybergym import cybergym_custody + + task_id = "cybergym-custody-finalizing" + clock = _FakeClock() + monkeypatch.setattr(cybergym_custody, "time", clock) + monkeypatch.setattr(cybergym_custody, "FINALIZATION_GRACE_SEC", 120.0) + state = {"polls": 0} + + def http(method, _url, **_kwargs): + state["polls"] += 1 + clock.now += 10.0 + if state["polls"] < 8: # 80 s of finalizing: past a 30 s custody window + return {"task_id": task_id, "status": "running", "artifact_status": "finalizing", "child_status": "completed"} + return {"task_id": task_id, "status": "completed", "result": {"cost_final": True}} + + config = _config(tmp_path, provider_probe=False, poll_interval_sec=1.0) + executor = CyberGymExecutor(dataclasses_replace(config, http_runner=http, sleep=clock.sleep)) + monkeypatch.setattr(executor, "_terminalize_gateway_attempt", lambda *_a, **_k: None) + checkpoint = config.run_root / "custody.json" + result = executor._poll_gateway_custody( # noqa: SLF001 - custody contract + task_id, checkpoint, cancel_response=None, cancel_status_code=503, custody_seconds=30.0 + ) + assert result["status"] == "completed" + saved = json.loads(checkpoint.read_text(encoding="utf-8")) + assert saved["status"] == "completed" diff --git a/tests/test_cybergym_regrade.py b/tests/test_cybergym_regrade.py index d66a987f1..6f6c4a4a4 100644 --- a/tests/test_cybergym_regrade.py +++ b/tests/test_cybergym_regrade.py @@ -73,7 +73,7 @@ def test_inventory_is_append_only_and_counts_candidates(tmp_path): payload = write_regrade_inventory(output, [root]) assert payload["counts"] == {"selected": 1, "ready": 1, "no_final_poc": 0, "hash_mismatch": 0} - assert json.loads(output.read_text())["candidates"][0]["task_id"] == "arvo:1" + assert json.loads(output.read_text(encoding="utf-8"))["candidates"][0]["task_id"] == "arvo:1" class _FakeRegradeExecutor: diff --git a/tests/test_cybergym_root_cost_evidence.py b/tests/test_cybergym_root_cost_evidence.py index a715710a7..5856b7ccd 100644 --- a/tests/test_cybergym_root_cost_evidence.py +++ b/tests/test_cybergym_root_cost_evidence.py @@ -106,7 +106,7 @@ def test_real_ledger_checkpoint_gateway_and_disk_reach_grace(canonical_terminal, live = response.json() assert live["root_phase_checkpoint"]["accounting"] == stored["root_phase_checkpoint"]["accounting"] assert _abandoned_cost_residue_usd(live) == pytest.approx(0.05) - raw = json.loads((data / "task_results/cybergym-root.json").read_text()) + raw = json.loads((data / "task_results/cybergym-root.json").read_text(encoding="utf-8")) assert _abandoned_cost_residue_usd(raw) == pytest.approx(0.05) root = tmp_path / "executor" root.mkdir() diff --git a/tests/test_cybergym_runtime_mode.py b/tests/test_cybergym_runtime_mode.py index 0c7a273d9..8ee75ae87 100644 --- a/tests/test_cybergym_runtime_mode.py +++ b/tests/test_cybergym_runtime_mode.py @@ -14,11 +14,11 @@ def test_runtime_mode_reaches_applied_settings(tmp_path, mode): args = parse_args(["--runtime-mode", mode, "--budget-usd", "200", "--per-task-cost-usd", "5", "--workers", "32"]) template = tmp_path / "template.json" - template.write_text('{"OUROBOROS_RUNTIME_MODE": "advanced"}') + template.write_text('{"OUROBOROS_RUNTIME_MODE": "advanced"}', encoding="utf-8") output = tmp_path / "run" output.mkdir() path, metadata = _prepare_applied_settings(template, output, args) - applied = json.loads(path.read_text()) + applied = json.loads(path.read_text(encoding="utf-8")) assert applied["OUROBOROS_RUNTIME_MODE"] == mode assert metadata["runtime_mode"] == mode assert metadata["effective_overrides"]["OUROBOROS_RUNTIME_MODE"] == mode diff --git a/tests/test_cybergym_server.py b/tests/test_cybergym_server.py index 74bc3e2b3..92abd0a6c 100644 --- a/tests/test_cybergym_server.py +++ b/tests/test_cybergym_server.py @@ -465,7 +465,11 @@ def test_start_exposes_attested_base_url_and_closes(tmp_path): assert server.stopped is True -def test_state_dir_places_mutable_state_outside_run_root(tmp_path): +def test_state_dir_places_mutable_state_outside_run_root(tmp_path, monkeypatch): + # State placement/export is independent of the host filesystem probe. + monkeypatch.setattr( + "devtools.benchmarks.cybergym.cybergym_server._mount_fs_type", lambda _path: "ext4" + ) seed, commit = _seed_repo(tmp_path) state = tmp_path / "nvme-state" wrapper = CyberGymIsolatedServer( @@ -481,7 +485,8 @@ def test_state_dir_places_mutable_state_outside_run_root(tmp_path): assert wrapper.data_root == state.resolve() / "ouroboros-data" assert (wrapper.data_root / ".ouroboros_isolated_benchmark").is_file() assert wrapper.settings_path == wrapper.data_root / "settings.json" - assert wrapper.settings_path.stat().st_mode & 0o777 == 0o600 + if sys.platform != "win32": # Windows chmod does not expose POSIX owner-only mode bits. + assert wrapper.settings_path.stat().st_mode & 0o777 == 0o600 # The durable run root keeps the clone but not the mutable state tree. assert wrapper.clone_root.is_dir() assert not (wrapper.run_root / "ouroboros-data").exists() @@ -575,7 +580,11 @@ def test_mount_fs_type_longest_prefix_wins(): assert _mount_fs_type(pathlib.Path("/elsewhere"), "garbage line\n") == "" -def test_close_mirrors_audit_surface_to_run_root(tmp_path): +def test_close_mirrors_audit_surface_to_run_root(tmp_path, monkeypatch): + # State placement/export is independent of the host filesystem probe. + monkeypatch.setattr( + "devtools.benchmarks.cybergym.cybergym_server._mount_fs_type", lambda _path: "ext4" + ) seed, commit = _seed_repo(tmp_path) wrapper = CyberGymIsolatedServer( seed, @@ -596,7 +605,8 @@ def test_close_mirrors_audit_surface_to_run_root(tmp_path): for name in ("state", "logs", "task_results", "memory"): assert (mirror / name / "marker.txt").read_text(encoding="utf-8") == name assert not (mirror / "observability").exists() - assert (mirror / "settings.json").stat().st_mode & 0o777 == 0o600 + if sys.platform != "win32": # Windows chmod does not expose POSIX owner-only mode bits. + assert (mirror / "settings.json").stat().st_mode & 0o777 == 0o600 assert (mirror / ".ouroboros_isolated_benchmark").is_file() receipt = wrapper.state_export assert receipt["ok"] is True diff --git a/tests/test_root_accounting_snapshot.py b/tests/test_root_accounting_snapshot.py index d7633ca60..c1fac46be 100644 --- a/tests/test_root_accounting_snapshot.py +++ b/tests/test_root_accounting_snapshot.py @@ -78,7 +78,7 @@ def test_root_snapshot_reaches_saved_and_gateway_results_without_flattening_own_ assert stored["non_final_rows"] == 0 # Own-task openness is not tree openness. assert stored["cost_final"] is False assert "cost_estimated" not in stored - saved = json.loads((root / "task_results" / "root.json").read_text()) + saved = json.loads((root / "task_results" / "root.json").read_text(encoding="utf-8")) request = Request({ "type": "http", "method": "GET", "path": "/api/tasks/root", "path_params": {"task_id": "root"}, "query_string": b"", "headers": [], @@ -173,8 +173,10 @@ def test_checkpoint_refresh_updates_snapshot_but_stale_phase_patch_does_not(root assert projected["root_phase_checkpoint"] == after["root_phase_checkpoint"] -def test_weighted_compaction_preserves_checkpoint_counts(root): - from tests.fixtures_usage_compaction import _compact +def test_weighted_compaction_preserves_checkpoint_counts(root, monkeypatch): + from tests.fixtures_usage_compaction import _compact, age_fixture_clock + + age_fixture_clock(monkeypatch) for _ in range(6): reservation = _reserve(root)