Simplify CyberGym settlement and consolidate gateway custody

This commit is contained in:
Ouroboros 2026-09-18 00:48:43 +03:00
parent 2fd8ed5f78
commit 26159dc644
13 changed files with 1078 additions and 1117 deletions

View file

@ -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"),

View file

@ -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,

View file

@ -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]:

View file

@ -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

View file

@ -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)

View file

@ -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.",

View file

@ -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"

View file

@ -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"

View file

@ -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:

View file

@ -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()

View file

@ -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

View file

@ -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

View file

@ -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)