mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
fix(cowork): bound meter blindness and retain interrupted activity
This commit is contained in:
parent
2f7a3c9e45
commit
d9b9157fa4
7 changed files with 699 additions and 87 deletions
|
|
@ -114,11 +114,32 @@ billing delay and work actually in flight; multiplying full task budgets would
|
|||
prematurely stop a 32-task, $2000 campaign at $1200 spent. With a $100 reserve,
|
||||
the campaign stop boundary is $1900 spent instead. The supervisor stops when the
|
||||
campaign remainder reaches that allowance, when the invocation bound is reached,
|
||||
or when the meter or disk reserve becomes unavailable. A suspect meter read
|
||||
gets at most three observations requesting cache revalidation within one shared
|
||||
15-second window; only a finite, nonnegative, nondecreasing value may update the campaign. Rejected
|
||||
values and errors remain in the monitor and unit log, including after shutdown.
|
||||
This confirmation does not rebase spending or ignore known budget exhaustion.
|
||||
when the disk reserve becomes unavailable, or when the meter stays blind for the
|
||||
recorded `--meter-blindness-sec` bound (default 30 seconds).
|
||||
|
||||
Blindness counts from the request time of the last reading that was accepted
|
||||
and saved to the campaign file. Read time, retry pauses and ledger/monitor writes
|
||||
all count; a rejected, non-finite, lower or late value never restarts the bound.
|
||||
A healthy meter is read every 15 seconds, or sooner when one read would
|
||||
otherwise not fit. Each read gets at most 5 seconds and requests cache
|
||||
revalidation; a failed or rejected read is retried after about 3 seconds within
|
||||
the same bound, so one slow read cannot consume it. Only a finite, nonnegative,
|
||||
nondecreasing value may update the campaign. Rejected values and errors remain in
|
||||
the monitor and unit log, including after shutdown. The bound does not rebase
|
||||
spending, change the ceiling, reserve or invocation bound, or ignore known
|
||||
budget exhaustion. At the bound, the runnable supervisor fences admission and
|
||||
stops the runner group and exact-label resources. Synchronous filesystem calls
|
||||
cannot be preempted: an overlong save or diagnostic write is checked immediately
|
||||
on return, never credited as a fresh window. Cleanup itself can take time.
|
||||
A final reading, bounded by the same
|
||||
interval, happens only after cleanup and settles accounting; it cannot admit work.
|
||||
|
||||
If the campaign file cannot be saved, the run stops with
|
||||
`campaign_persistence_failed`; if settlement still cannot be saved, the campaign
|
||||
keeps its unsettled `active_run`. A failed ledger or monitor write is reported on
|
||||
stderr and in later records, without stopping work while the campaign record
|
||||
and meter remain valid. An unexpected supervisor failure is recorded as
|
||||
`supervisor_error`, never as a finished runner or exhausted budget.
|
||||
These are configurable operator bounds. Billing can be delayed and paid calls may already be in flight;
|
||||
the monitor is **not a provider-enforced hard dollar cap**.
|
||||
|
||||
|
|
@ -145,6 +166,16 @@ after model work is a genuine failed attempt, not a new attempt entitlement.
|
|||
The result ledger retains every selected ID, including `not_attempted` entries.
|
||||
A timeout before task submission is an infrastructure failure. After submission,
|
||||
missing token telemetry does not prove that no paid/model work happened.
|
||||
`not_attempted` means neither a runner row nor adapter start evidence exists; the
|
||||
benchmark's pre-created dump and empty workspace are not evidence. A started task
|
||||
without an adapter summary is `infra_failed` with reason `interrupted:<cause>`:
|
||||
the run's stop reason, or `runner_exited` when the runner ended without a supervisor
|
||||
stop. A live snapshot retains the legacy `infra_failed`/`missing_adapter_summary`
|
||||
diagnostic bucket with `provisional: true`; it asserts no interruption or settled
|
||||
outcome and must not be scored as a finished result. Its
|
||||
`paid_activity` is `observed` only when a copied checkpoint holds a token-bearing
|
||||
`llm_usage` record and is otherwise `unknown`; no task cost is inferred. A runner
|
||||
row without a summary receives the same interruption cause after the run ends.
|
||||
Infrastructure recovery uses new roots and the identical configuration, seed and
|
||||
immutable image. With no explicit new selection, it preserves the original task
|
||||
selection, and always retains cumulative ancestry. At most two recovery passes
|
||||
|
|
|
|||
|
|
@ -96,12 +96,17 @@ it is not a provider-enforced spending cap.
|
|||
|
||||
## Monitor, recover and audit
|
||||
|
||||
A suspect budget-meter response is confirmed with at most three reads requesting
|
||||
cache revalidation inside one 15-second read window. Each request uses only the
|
||||
remaining time. Counters must still be finite, nonnegative and no lower than the last
|
||||
accepted value. `monitor.json` and the unit log retain rejected observations,
|
||||
the previous value, errors and any successful confirmation. Persistent meter
|
||||
failure stops the run; a known exhausted budget never waits for another poll.
|
||||
The budget meter is read every 15 seconds. `--meter-blindness-sec` (default 30,
|
||||
recorded in the manifest) bounds the time since the last accepted, saved
|
||||
reading. Each read gets at most 5 seconds; a failed or rejected read is retried
|
||||
about every 3 seconds inside that bound, and reaching it stops the run with
|
||||
`budget_meter_unavailable`. Counters must still be finite, nonnegative and no
|
||||
lower than the last accepted value. `monitor.json` and the unit log retain
|
||||
rejected observations, the previous value, errors and any successful
|
||||
confirmation. A known exhausted budget never waits for another poll. A failed
|
||||
campaign-file write stops the run with `campaign_persistence_failed`; a failed
|
||||
ledger or monitor write is reported on stderr as `cowork_diagnostic_write_failed`
|
||||
while valid work continues.
|
||||
|
||||
Cleanup rechecks the exact run label after removing containers and networks.
|
||||
A competing cleanup's already-removed response succeeds only when that fresh
|
||||
|
|
@ -111,7 +116,9 @@ failure and keeps the campaign's unsettled custody.
|
|||
`monitor.json` records progress, key-meter spending, disk headroom and stop
|
||||
reasons. `run_manifest.json` records the seed, benchmark, immutable image ID,
|
||||
selected IDs, recovery ancestry and applied configuration; `result_index.jsonl` retains every selected task,
|
||||
including failures and tasks not started. Task artifacts live below
|
||||
including failures and tasks not started. A started task without an adapter
|
||||
summary is `infra_failed` with `interrupted:<cause>` and a `paid_activity` of
|
||||
`observed` or `unknown`, as defined in the methodology. Task artifacts live below
|
||||
`bench/dumps/`, with sanitized runtime logs in each task's `ouroboros/` folder.
|
||||
A launcher exit code is not a task score. Task polling checkpoints the selected
|
||||
scrubbed logs before its status request, so an aborted run can retain partial
|
||||
|
|
|
|||
|
|
@ -105,6 +105,26 @@ def _number(value: Any) -> float | None:
|
|||
return number if math.isfinite(number) and number >= 0 else None
|
||||
|
||||
|
||||
def usage_tokens(record: dict[str, Any]) -> dict[str, int | None]:
|
||||
"""Token counts stated by one llm_usage record; ``None`` where it states none."""
|
||||
usage = record.get("usage")
|
||||
usage = usage if isinstance(usage, dict) else {}
|
||||
counts = {key: _number(record.get(key, usage.get(key)))
|
||||
for key in ("prompt_tokens", "completion_tokens", "cached_tokens")}
|
||||
return {key: None if value is None else int(value) for key, value in counts.items()}
|
||||
|
||||
|
||||
def token_bearing(tokens: dict[str, int | None]) -> bool:
|
||||
"""Positive evidence of a model call; zero or unstated tokens prove nothing."""
|
||||
return (tokens["prompt_tokens"] or 0) + (tokens["completion_tokens"] or 0) > 0
|
||||
|
||||
|
||||
def model_activity_observed(events: Path) -> bool:
|
||||
"""Whether a copied events log holds a token-bearing llm_usage record (first one wins)."""
|
||||
return any(record.get("type") == "llm_usage" and token_bearing(usage_tokens(record))
|
||||
for _line, record in _records(events, []))
|
||||
|
||||
|
||||
def _omission_count(value: Any) -> int:
|
||||
if isinstance(value, (list, dict)):
|
||||
return len(value)
|
||||
|
|
@ -138,16 +158,12 @@ def audit_task(task_dump: Path, ledger: dict[str, Any]) -> dict[str, Any]:
|
|||
activity["usage_records"] += 1
|
||||
usage = record.get("usage")
|
||||
usage = usage if isinstance(usage, dict) else {}
|
||||
token_values = {}
|
||||
for key in ("prompt_tokens", "completion_tokens", "cached_tokens"):
|
||||
value = _number(record.get(key, usage.get(key)))
|
||||
tokens = usage_tokens(record)
|
||||
for key, value in tokens.items():
|
||||
if value is None:
|
||||
gaps.append({"source": source.name, "line": line, "reason": f"unknown_{key}"})
|
||||
token_values[key] = int(value or 0)
|
||||
activity[key] += token_values[key]
|
||||
activity["nonempty_usage_records"] += int(
|
||||
token_values["prompt_tokens"] + token_values["completion_tokens"] > 0
|
||||
)
|
||||
activity[key] += value or 0
|
||||
activity["nonempty_usage_records"] += int(token_bearing(tokens))
|
||||
cost = _number(record.get("cost", usage.get("cost")))
|
||||
if cost is None or record.get("cost_known") is False:
|
||||
unknown_cost += 1
|
||||
|
|
|
|||
|
|
@ -32,6 +32,10 @@ class UsageCounterError(ValueError):
|
|||
f"observed={observed_usage!r}, previous={previous_usage!r}")
|
||||
|
||||
|
||||
class CampaignPersistenceError(RuntimeError):
|
||||
"""The durable campaign record was not written, so it no longer bounds spending."""
|
||||
|
||||
|
||||
def validate_usage(usage: float, previous_usage: float) -> None:
|
||||
if not math.isfinite(usage) or usage < 0 or usage < previous_usage:
|
||||
raise UsageCounterError(usage, previous_usage)
|
||||
|
|
@ -127,9 +131,13 @@ class CampaignBudget:
|
|||
|
||||
self.record.update({"spent_usd": self.spent, "remaining_usd": self.remaining,
|
||||
"observed_at": time.time()})
|
||||
write_json(self.path, self.record)
|
||||
try:
|
||||
write_json(self.path, self.record)
|
||||
except Exception as exc:
|
||||
raise CampaignPersistenceError(f"campaign record not saved: {type(exc).__name__}: {exc}") from exc
|
||||
|
||||
def observe(self, usage: float) -> None:
|
||||
# Only a finite, nonnegative, nondecreasing value reaches the durable record.
|
||||
validate_usage(usage, self.record["last_usage"])
|
||||
self.record["last_usage"] = usage
|
||||
self.save()
|
||||
|
|
|
|||
|
|
@ -44,7 +44,14 @@ from devtools.benchmarks.common.result_index import (
|
|||
)
|
||||
from devtools.benchmarks.common.run_roots import assert_outside_repo, repo_root_from_devtools, run_root, timestamp_run_id
|
||||
from devtools.benchmarks.common.secrets import credential_fingerprint
|
||||
from devtools.benchmarks.cowork_bench.campaign import CampaignBudget, campaign_lock, key_usage, validate_usage
|
||||
from devtools.benchmarks.cowork_bench.audit_cowork_bench import model_activity_observed
|
||||
from devtools.benchmarks.cowork_bench.campaign import (
|
||||
CampaignBudget,
|
||||
CampaignPersistenceError,
|
||||
campaign_lock,
|
||||
key_usage,
|
||||
validate_usage,
|
||||
)
|
||||
from devtools.benchmarks.cowork_bench.resource_limits import LABEL_KEY, prepare_resource_env
|
||||
from ouroboros.platform_layer import kill_process_group_id, terminate_process_group_id
|
||||
from ouroboros.process_custody import spawn_supervised
|
||||
|
|
@ -63,6 +70,16 @@ SECRET_NAME = "ouroboros_bench.secret.json"
|
|||
WEB_TOOLS = ("web_search", "browse_page", "browser_action", "youtube_transcript")
|
||||
DELEGATED_VISION_TOOLS = ("analyze_screenshot", "vlm_query")
|
||||
SETTLED_STATUSES = frozenset({"passed", "failed", "agent_failed"})
|
||||
# Meter cadence. The blindness bound itself is the operator's recorded
|
||||
# --meter-blindness-sec; one read's slice keeps a slow read from consuming it.
|
||||
DEFAULT_METER_BLINDNESS_SEC = 30.0
|
||||
METER_POLL_SEC = 15.0
|
||||
METER_RETRY_SEC = 3.0
|
||||
METER_READ_SLICE_SEC = 5.0
|
||||
# Adapter-owned artifacts that exist only after this task's container began work.
|
||||
# The benchmark pre-creates the dump and its empty workspace; neither is evidence.
|
||||
START_EVIDENCE = ("ouroboros", "applied_settings.json", "preprocess.log", "mcp_proxy.log",
|
||||
"ouroboros_server.log", "traj_log.json", "traj.json")
|
||||
|
||||
|
||||
def dump_dir_name(model: str) -> str:
|
||||
|
|
@ -256,20 +273,33 @@ def _load(path: pathlib.Path) -> dict[str, Any]:
|
|||
return data if isinstance(data, dict) else {}
|
||||
|
||||
|
||||
def ledger_row(task: str, task_dump: pathlib.Path, runner_row: dict[str, str]) -> dict[str, Any]:
|
||||
def ledger_row(task: str, task_dump: pathlib.Path, runner_row: dict[str, str], *, cause: str) -> dict[str, Any]:
|
||||
"""One denominator-preserving row. The runner's exit code and CSV are NOT the status: the
|
||||
adapter summary says how the agent phase ended and ``eval_res.json`` is the verdict."""
|
||||
adapter summary says how the agent phase ended and ``eval_res.json`` is the verdict.
|
||||
|
||||
Without that summary, only a task with no runner row and no start evidence is
|
||||
``not_attempted``. Otherwise it is an infrastructure row whose paid activity is
|
||||
``observed`` from token-bearing usage or else ``unknown``, never an invented cost.
|
||||
"""
|
||||
summary = _load(task_dump / "ouroboros_summary.json")
|
||||
eval_res = _load(task_dump / "eval_res.json")
|
||||
paths = {"task_dump": str(task_dump)}
|
||||
details = {"runner": runner_row, "adapter": summary}
|
||||
details: dict[str, Any] = {"runner": runner_row, "adapter": summary}
|
||||
if runner_row.get("status") == "pg_fail":
|
||||
return task_result_row(benchmark=BENCHMARK, instance_id=task, status="infra_failed",
|
||||
reason_code="pg_fail", output_paths=paths, details=details)
|
||||
if not summary:
|
||||
status = "infra_failed" if runner_row else "not_attempted"
|
||||
return task_result_row(benchmark=BENCHMARK, instance_id=task, status=status,
|
||||
reason_code="missing_adapter_summary" if runner_row else "missing_result",
|
||||
evidence = [name for name in START_EVIDENCE if (task_dump / name).exists()]
|
||||
if not (runner_row or evidence):
|
||||
return task_result_row(benchmark=BENCHMARK, instance_id=task, status="not_attempted",
|
||||
reason_code="missing_result", output_paths=paths, details=details)
|
||||
observed = model_activity_observed(task_dump / "ouroboros" / "events.jsonl")
|
||||
details.update({"start_evidence": evidence, "paid_activity": "observed" if observed else "unknown"})
|
||||
# A live diagnostic snapshot is not proof of interruption. The legacy infra
|
||||
# bucket stays provisional until the run stops; it is never a settled outcome.
|
||||
details["provisional"] = cause == "in_progress"
|
||||
return task_result_row(benchmark=BENCHMARK, instance_id=task, status="infra_failed",
|
||||
reason_code="missing_adapter_summary" if cause == "in_progress" else f"interrupted:{cause}",
|
||||
output_paths=paths, details=details)
|
||||
reason = str(summary.get("reason_code") or "")
|
||||
if summary.get("infra_failed"):
|
||||
|
|
@ -290,10 +320,14 @@ def ledger_row(task: str, task_dump: pathlib.Path, runner_row: dict[str, str]) -
|
|||
)
|
||||
|
||||
|
||||
def write_ledger(ledger_path: pathlib.Path, bench_dir: pathlib.Path, model: str, tasks: list[str]) -> dict[str, int]:
|
||||
def write_ledger(ledger_path: pathlib.Path, bench_dir: pathlib.Path, model: str, tasks: list[str],
|
||||
*, cause: str) -> dict[str, int]:
|
||||
"""``cause`` names why a started task without a summary did not finish: the run's stop
|
||||
reason, ``runner_exited``, or ``in_progress`` in a snapshot taken while the run is live."""
|
||||
runner_rows = read_summary_csv(bench_dir / "benchmark_logs")
|
||||
dumps = bench_dir / "dumps" / dump_dir_name(model)
|
||||
rows = [ledger_row(task, dumps / f"SingleUserTurn-{task}", runner_rows.get(task, {})) for task in tasks]
|
||||
rows = [ledger_row(task, dumps / f"SingleUserTurn-{task}", runner_rows.get(task, {}), cause=cause)
|
||||
for task in tasks]
|
||||
write_result_index(ledger_path, rows)
|
||||
counts: dict[str, int] = {}
|
||||
for row in rows:
|
||||
|
|
@ -420,24 +454,29 @@ def stop_process_group(proc: subprocess.Popen) -> None:
|
|||
proc.wait(timeout=15)
|
||||
|
||||
|
||||
def observe_campaign_usage(api_key: str, campaign: CampaignBudget,
|
||||
diagnostics: list[dict[str, Any]], *, phase: str) -> None:
|
||||
"""Confirm suspect meter reads within one 15-second window, never accept a lower counter."""
|
||||
deadline = time.monotonic() + 15.0
|
||||
def observe_campaign_usage(api_key: str, campaign: CampaignBudget, diagnostics: list[dict[str, Any]],
|
||||
*, phase: str, deadline: float) -> float:
|
||||
"""Read the meter until one value is accepted and saved, or until ``deadline``.
|
||||
|
||||
Each read gets a bounded slice and a failed or rejected read is retried after a short
|
||||
pause; rejections never move the deadline and a lower counter is never accepted.
|
||||
Returns the monotonic time at which the accepted read was requested, the conservative
|
||||
anchor of the next blindness bound. A campaign write failure propagates unretried.
|
||||
"""
|
||||
previous = campaign.record["last_usage"]
|
||||
last_error: Exception = TimeoutError("usage confirmation window expired")
|
||||
last_error: Exception = TimeoutError("meter blindness bound elapsed without a confirmed reading")
|
||||
rejected = False
|
||||
for attempt in range(1, 4):
|
||||
remaining = deadline - time.monotonic()
|
||||
if remaining <= 0:
|
||||
break
|
||||
attempt = 0
|
||||
while (remaining := deadline - time.monotonic()) > 0:
|
||||
attempt += 1
|
||||
requested_at = time.monotonic()
|
||||
value = None
|
||||
try:
|
||||
value = key_usage(api_key, timeout=remaining)
|
||||
value = key_usage(api_key, timeout=min(METER_READ_SLICE_SEC, remaining))
|
||||
if time.monotonic() > deadline:
|
||||
raise TimeoutError("usage confirmation arrived after its shared deadline")
|
||||
raise TimeoutError("meter reading arrived after the blindness bound")
|
||||
validate_usage(value, previous)
|
||||
except Exception as exc: # Provider/transport/parse failures share the bounded confirmation.
|
||||
except Exception as exc: # Provider/transport/parse failures share the same bound.
|
||||
last_error = exc
|
||||
value = getattr(exc, "observed_usage", value)
|
||||
observation = {
|
||||
|
|
@ -450,34 +489,62 @@ def observe_campaign_usage(api_key: str, campaign: CampaignBudget,
|
|||
print(json.dumps({"event": "cowork_meter_observation", **observation}), flush=True)
|
||||
rejected = True
|
||||
else:
|
||||
# Persistence failures are not suspect provider reads and are never retried here.
|
||||
# A persistence failure is not a suspect provider read and is never retried here.
|
||||
campaign.observe(value)
|
||||
if time.monotonic() >= deadline:
|
||||
# Accounting may have reached disk, but a slow save cannot renew
|
||||
# permission to spend after the previous blindness window expired.
|
||||
raise TimeoutError("campaign persistence completed after the blindness bound")
|
||||
if rejected:
|
||||
observation = {"observed_at": time.time(), "phase": phase, "attempt": attempt,
|
||||
"previous_usage": previous, "observed_usage": value, "accepted": True}
|
||||
diagnostics.append(observation)
|
||||
print(json.dumps({"event": "cowork_meter_observation", **observation}), flush=True)
|
||||
return
|
||||
if attempt < 3:
|
||||
delay = min(3.0, max(0.0, deadline - time.monotonic()))
|
||||
if delay:
|
||||
time.sleep(delay)
|
||||
return requested_at
|
||||
time.sleep(min(METER_RETRY_SEC, max(0.0, deadline - time.monotonic())))
|
||||
raise last_error
|
||||
|
||||
|
||||
def publish_diagnostics(root: pathlib.Path, bench_dir: pathlib.Path, args: argparse.Namespace,
|
||||
monitor: dict[str, Any], failures: dict[str, dict[str, Any]], *, cause: str) -> dict[str, int] | None:
|
||||
"""Refresh the ledger snapshot and monitor. Neither is spending authority: a failed write
|
||||
is disclosed on stderr and in later records, never a reason to stop valid work."""
|
||||
def disclose(artifact: str, exc: Exception) -> None:
|
||||
entry = failures.setdefault(artifact, {"count": 0})
|
||||
entry.update(count=entry["count"] + 1, error_type=type(exc).__name__, error=str(exc))
|
||||
# Until a later write succeeds, this stderr line is the only record of the failure.
|
||||
print(json.dumps({"event": "cowork_diagnostic_write_failed", "artifact": artifact, **entry}),
|
||||
file=sys.stderr, flush=True)
|
||||
|
||||
counts = None
|
||||
try:
|
||||
counts = write_ledger(root / "result_index.jsonl", bench_dir, args.model, args.selected_tasks, cause=cause)
|
||||
except Exception as exc:
|
||||
disclose("result_index.jsonl", exc)
|
||||
try:
|
||||
write_json(root / "monitor.json", {**monitor, "ledger_counts": counts, "diagnostic_write_failures": failures})
|
||||
except Exception as exc:
|
||||
disclose("monitor.json", exc)
|
||||
return counts
|
||||
|
||||
|
||||
def supervise_run(args: argparse.Namespace, command: list[str], bench_dir: pathlib.Path,
|
||||
run_env: dict[str, str], api_key: str, campaign: CampaignBudget) -> dict[str, Any]:
|
||||
run_env: dict[str, str], api_key: str, campaign: CampaignBudget,
|
||||
*, confirmed_at: float | None = None) -> dict[str, Any]:
|
||||
"""Own runner lifetime, budget meter and disk reserve until every owned resource is settled.
|
||||
|
||||
Stop before the campaign limit with a reserve for work already sent to providers. API
|
||||
billing can lag: the reserve is disclosed, never a claim of a provider-enforced cap.
|
||||
Meter loss stops new work and the run, rather than continuing with unknown spending.
|
||||
``confirmed_at`` is the request time of the saved startup reading. Once the meter has
|
||||
gone ``--meter-blindness-sec`` without another confirmed, saved reading, or the campaign
|
||||
record cannot be saved, new work and the run stop rather than continue with unknown spending.
|
||||
"""
|
||||
root = bench_dir.parent
|
||||
stop_file = pathlib.Path(run_env["COWORK_STOP_FILE"])
|
||||
# A task's lifetime budget is not a reservation of unsettled provider charges.
|
||||
# The CLI validates this explicit billing allowance as nonnegative.
|
||||
reserve = args.budget_reserve_usd
|
||||
blindness = args.meter_blindness_sec
|
||||
initial_spent = campaign.spent
|
||||
if campaign.remaining <= reserve:
|
||||
raise ValueError("campaign remaining budget does not cover the in-flight reserve")
|
||||
|
|
@ -486,13 +553,21 @@ def supervise_run(args: argparse.Namespace, command: list[str], bench_dir: pathl
|
|||
reason = ""
|
||||
meter_error = ""
|
||||
meter_diagnostics: list[dict[str, Any]] = []
|
||||
write_failures: dict[str, dict[str, Any]] = {}
|
||||
code = 1
|
||||
campaign.start(root)
|
||||
|
||||
confirmed_at = time.monotonic() if confirmed_at is None else confirmed_at
|
||||
def request_stop(signum, _frame):
|
||||
mark_stop(stop_file, f"signal_{signum}")
|
||||
|
||||
try:
|
||||
try:
|
||||
campaign.start(root)
|
||||
except CampaignPersistenceError:
|
||||
reason = "campaign_persistence_failed"
|
||||
raise
|
||||
if time.monotonic() >= confirmed_at + blindness:
|
||||
reason = "budget_meter_unavailable"
|
||||
raise TimeoutError("startup meter observation expired before runner spawn")
|
||||
for sig in (signal.SIGTERM, signal.SIGINT):
|
||||
old_handlers[sig] = signal.signal(sig, request_stop)
|
||||
with (root / "console.log").open("w", encoding="utf-8") as stream:
|
||||
|
|
@ -501,13 +576,15 @@ def supervise_run(args: argparse.Namespace, command: list[str], bench_dir: pathl
|
|||
stderr=subprocess.STDOUT, text=True)
|
||||
while proc.poll() is None:
|
||||
if stop_file.exists():
|
||||
reason = stop_file.read_text(encoding="utf-8").strip()
|
||||
reason = stop_file.read_text(encoding="utf-8").strip() or "stop_requested"
|
||||
break
|
||||
try:
|
||||
observe_campaign_usage(api_key, campaign, meter_diagnostics, phase="poll")
|
||||
confirmed_at = observe_campaign_usage(api_key, campaign, meter_diagnostics, phase="poll",
|
||||
deadline=confirmed_at + blindness)
|
||||
except CampaignPersistenceError as exc:
|
||||
meter_error, reason = type(exc).__name__, "campaign_persistence_failed"
|
||||
except Exception as exc:
|
||||
meter_error = type(exc).__name__
|
||||
reason = "budget_meter_unavailable"
|
||||
meter_error, reason = type(exc).__name__, "budget_meter_unavailable"
|
||||
free = shutil.disk_usage(args.resource_root).free
|
||||
root_free = shutil.disk_usage(pathlib.Path.home().anchor).free
|
||||
if not reason and (free < args.min_free_gib * 1024**3 or root_free < args.min_root_free_gib * 1024**3):
|
||||
|
|
@ -516,24 +593,35 @@ def supervise_run(args: argparse.Namespace, command: list[str], bench_dir: pathl
|
|||
reason = "campaign_budget_reserve"
|
||||
if not reason and campaign.spent - initial_spent >= args.budget_usd:
|
||||
reason = "run_budget"
|
||||
counts = write_ledger(root / "result_index.jsonl", bench_dir, args.model, args.selected_tasks)
|
||||
write_json(root / "monitor.json", {
|
||||
if reason:
|
||||
# Stop owned work at once; the final record is published after cleanup.
|
||||
mark_stop(stop_file, reason)
|
||||
break
|
||||
publish_diagnostics(root, bench_dir, args, {
|
||||
"observed_at": time.time(), "launcher_pid": os.getpid(), "runner_pid": proc.pid,
|
||||
"campaign_spent_usd": campaign.spent, "campaign_remaining_usd": campaign.remaining,
|
||||
"run_spent_usd": campaign.spent - initial_spent, "inflight_reserve_usd": reserve,
|
||||
"meter_blindness_sec": blindness,
|
||||
"meter_confirmed_age_sec": round(time.monotonic() - confirmed_at, 3),
|
||||
"disk_free_bytes": free, "root_free_bytes": root_free,
|
||||
"ledger_counts": counts, "stop_reason": reason, "meter_error": meter_error,
|
||||
"meter_diagnostics": meter_diagnostics,
|
||||
})
|
||||
if reason:
|
||||
"stop_reason": reason, "meter_error": meter_error, "meter_diagnostics": meter_diagnostics,
|
||||
}, write_failures, cause="in_progress")
|
||||
if time.monotonic() >= confirmed_at + blindness:
|
||||
reason, meter_error = "budget_meter_unavailable", "TimeoutError"
|
||||
mark_stop(stop_file, reason)
|
||||
break
|
||||
# Poll every 15 s, but early enough that one full read still fits in the bound.
|
||||
wait = confirmed_at + blindness - METER_READ_SLICE_SEC - time.monotonic()
|
||||
try:
|
||||
code = int(proc.wait(timeout=15))
|
||||
code = int(proc.wait(timeout=min(METER_POLL_SEC, max(0.0, wait))))
|
||||
except subprocess.TimeoutExpired:
|
||||
continue
|
||||
if proc.poll() is not None:
|
||||
code = int(proc.returncode)
|
||||
except BaseException:
|
||||
# An unexpected supervisor failure is neither a finished runner nor exhausted money.
|
||||
reason = reason or "supervisor_error"
|
||||
raise
|
||||
finally:
|
||||
# Also fence creation on unexpected exceptions. Cleanup must not race the shell
|
||||
# admitting the next task after a finished task returned its concurrency token.
|
||||
|
|
@ -545,20 +633,24 @@ def supervise_run(args: argparse.Namespace, command: list[str], bench_dir: pathl
|
|||
finally:
|
||||
for sig, handler in old_handlers.items():
|
||||
signal.signal(sig, handler)
|
||||
# Accounting only: the runner group and exact-label resources are already gone.
|
||||
try:
|
||||
observe_campaign_usage(api_key, campaign, meter_diagnostics, phase="final")
|
||||
observe_campaign_usage(api_key, campaign, meter_diagnostics, phase="final",
|
||||
deadline=time.monotonic() + blindness)
|
||||
except Exception as exc:
|
||||
meter_error = type(exc).__name__
|
||||
cause = reason or "runner_exited"
|
||||
campaign.finish(root, outcome=reason or "runner_finished", meter_error=meter_error)
|
||||
write_json(root / "monitor.json", {
|
||||
counts = publish_diagnostics(root, bench_dir, args, {
|
||||
"observed_at": time.time(), "finished": True, "runner_exit_code": code,
|
||||
"campaign_spent_usd": campaign.spent, "campaign_remaining_usd": campaign.remaining,
|
||||
"run_spent_usd": campaign.spent - initial_spent, "stop_reason": reason,
|
||||
"meter_error": meter_error, "meter_diagnostics": meter_diagnostics,
|
||||
"ledger_counts": write_ledger(root / "result_index.jsonl", bench_dir, args.model, args.selected_tasks),
|
||||
})
|
||||
}, write_failures, cause=cause)
|
||||
return {"stop_reason": reason, "runner_exit_code": code, "meter_error": meter_error,
|
||||
"meter_diagnostics": meter_diagnostics,
|
||||
"meter_diagnostics": meter_diagnostics, "interruption_cause": cause,
|
||||
"diagnostic_write_failures": write_failures,
|
||||
"ledger_counts": counts,
|
||||
"key_spend_usd": campaign.spent - initial_spent, "campaign_spent_usd": campaign.spent,
|
||||
"campaign_remaining_usd": campaign.remaining, "inflight_reserve_usd": reserve}
|
||||
|
||||
|
|
@ -593,6 +685,8 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
|
|||
parser.add_argument("--campaign-budget-usd", type=float, default=1000.0)
|
||||
parser.add_argument("--prior-spend-usd", type=float, default=0.0, help="already spent before the first campaign baseline")
|
||||
parser.add_argument("--budget-reserve-usd", type=float, default=100.0, help="unspent allowance for delayed/in-flight billing")
|
||||
parser.add_argument("--meter-blindness-sec", type=float, default=DEFAULT_METER_BLINDNESS_SEC,
|
||||
help="stop once this long passes without a confirmed, saved meter reading")
|
||||
parser.add_argument("--resource-root", default="", help="heavy storage filesystem; defaults to the run root")
|
||||
parser.add_argument("--min-free-gib", type=float, default=200.0)
|
||||
parser.add_argument("--min-root-free-gib", type=float, default=40.0)
|
||||
|
|
@ -616,6 +710,9 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
|
|||
parser.error("budgets and timeout must be positive")
|
||||
if min(args.min_free_gib, args.min_root_free_gib, args.budget_reserve_usd, args.prior_spend_usd) < 0:
|
||||
parser.error("reserves and prior spending cannot be negative")
|
||||
if not math.isfinite(args.meter_blindness_sec) or args.meter_blindness_sec <= METER_POLL_SEC + METER_READ_SLICE_SEC:
|
||||
parser.error(f"meter blindness bound must exceed one {METER_POLL_SEC:g}s poll plus one "
|
||||
f"{METER_READ_SLICE_SEC:g}s read")
|
||||
return args
|
||||
|
||||
|
||||
|
|
@ -661,6 +758,8 @@ def main(argv: list[str] | None = None) -> int:
|
|||
"max_steps": args.max_steps,
|
||||
"memory_mode": "empty",
|
||||
"base_image": args.base_image,
|
||||
"meter": {"blindness_sec": args.meter_blindness_sec, "poll_sec": METER_POLL_SEC,
|
||||
"retry_sec": METER_RETRY_SEC, "read_slice_sec": METER_READ_SLICE_SEC},
|
||||
},
|
||||
extra={"outcome": "started"},
|
||||
)
|
||||
|
|
@ -767,20 +866,26 @@ def main(argv: list[str] | None = None) -> int:
|
|||
manifest["harness"]["campaign_file"] = str(campaign_path)
|
||||
secret_path = bench_dir / "configs" / SECRET_NAME
|
||||
with campaign_lock(campaign_path):
|
||||
# The validated, saved startup reading anchors the first meter blindness bound.
|
||||
confirmed_at = time.monotonic()
|
||||
campaign = CampaignBudget(campaign_path, fingerprint=credential_fingerprint(api_key),
|
||||
ceiling=args.campaign_budget_usd, usage=key_usage(api_key),
|
||||
prior_spend=args.prior_spend_usd)
|
||||
secret_path.touch(mode=0o600)
|
||||
try:
|
||||
secret_path.write_text(json.dumps({"settings": {args.credential_setting: api_key}}), encoding="utf-8")
|
||||
result = supervise_run(args, command, bench_dir, run_env, api_key, campaign)
|
||||
result = supervise_run(args, command, bench_dir, run_env, api_key, campaign,
|
||||
confirmed_at=confirmed_at)
|
||||
finally:
|
||||
secret_path.write_text("{}", encoding="utf-8")
|
||||
counts = write_ledger(ledger_output, bench_dir, args.model, tasks)
|
||||
final.update(result)
|
||||
# The supervised final publication already owns this write and its errors.
|
||||
# Repeating it here would turn a diagnostic failure back into an exception.
|
||||
counts = result["ledger_counts"]
|
||||
stopped = bool(result["stop_reason"] or result["meter_error"])
|
||||
infra = counts.get("infra_failed", 0) + counts.get("not_attempted", 0)
|
||||
infra = counts is None or counts.get("infra_failed", 0) + counts.get("not_attempted", 0)
|
||||
code = 3 if stopped else (1 if result["runner_exit_code"] or infra else 0)
|
||||
final.update({**result, "outcome": result["stop_reason"] or ("infra_failed" if infra else "completed"),
|
||||
final.update({"outcome": result["stop_reason"] or ("infra_failed" if infra else "completed"),
|
||||
"exit_code": code, "ledger_counts": counts})
|
||||
return code
|
||||
|
||||
|
|
|
|||
|
|
@ -151,7 +151,7 @@ def test_ledger_separates_real_failures_from_recoverable_infrastructure(tmp_path
|
|||
write_json(tmp_path / "ouroboros_summary.json", summary)
|
||||
if evaluation:
|
||||
write_json(tmp_path / "eval_res.json", evaluation)
|
||||
row = launcher.ledger_row("task", tmp_path, runner)
|
||||
row = launcher.ledger_row("task", tmp_path, runner, cause="runner_exited")
|
||||
assert row["status"] == expected
|
||||
assert row["instance_id"] == "task"
|
||||
|
||||
|
|
@ -225,6 +225,18 @@ def test_manifest_metadata_matches_config_received_by_container(dry_launcher):
|
|||
assert not config["settings"].get("OUROBOROS_OR_PROVIDER")
|
||||
|
||||
|
||||
def test_meter_blindness_bound_is_validated_and_recorded_in_the_manifest(dry_launcher):
|
||||
out, argv = dry_launcher
|
||||
assert launcher.parse_args([]).meter_blindness_sec == 30.0
|
||||
for invalid in ("20", "-1", "nan", "inf"): # Must leave room for one 15 s poll plus one read.
|
||||
with pytest.raises(SystemExit):
|
||||
launcher.parse_args(["--meter-blindness-sec", invalid])
|
||||
assert launcher.main([*argv, "--meter-blindness-sec", "45"]) == 0
|
||||
manifest = json.loads((out / "run_manifest.json").read_text(encoding="utf-8"))
|
||||
assert manifest["harness"]["meter"] == {"blindness_sec": 45.0, "poll_sec": 15.0,
|
||||
"retry_sec": 3.0, "read_slice_sec": 5.0}
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def supervised(tmp_path, monkeypatch):
|
||||
bench = tmp_path / "bench"
|
||||
|
|
@ -263,8 +275,11 @@ def supervised(tmp_path, monkeypatch):
|
|||
monkeypatch.setattr(launcher, "stop_process_group", stop)
|
||||
monkeypatch.setattr(launcher, "remove_run_containers", cleanup)
|
||||
monkeypatch.setattr(launcher, "key_usage", lambda _key, **_kwargs: 100)
|
||||
# Retries wait on this simulated clock, so a meter outage reaches its bound instantly.
|
||||
clock = SimpleNamespace(now=0.0)
|
||||
monkeypatch.setattr(launcher, "time", SimpleNamespace(
|
||||
time=time.time, monotonic=time.monotonic, sleep=lambda _delay: None))
|
||||
time=lambda: clock.now, monotonic=lambda: clock.now,
|
||||
sleep=lambda delay: setattr(clock, "now", clock.now + delay)))
|
||||
monkeypatch.setattr(launcher.shutil, "disk_usage", lambda _path: SimpleNamespace(free=1024**4))
|
||||
return args, bench, env, budget, events, handlers, proc
|
||||
|
||||
|
|
@ -618,7 +633,8 @@ def test_paid_runner_uses_immutable_image_and_scrubs_ambient_alternate_keys(dry_
|
|||
monkeypatch.setattr(launcher, "key_headroom", lambda _key, **_kwargs: {"effective": 1000})
|
||||
monkeypatch.setattr(launcher, "key_usage", lambda _key, **_kwargs: 100)
|
||||
monkeypatch.setattr(launcher, "prepare_resource_env", lambda env, **_kwargs: dict(env))
|
||||
def supervise(args, command, bench, env, api_key, campaign):
|
||||
def supervise(args, command, bench, env, api_key, campaign, *, confirmed_at):
|
||||
assert isinstance(confirmed_at, float)
|
||||
assert command[-1] == "one"
|
||||
assert env["IMAGE"] == IMAGE_ID
|
||||
assert not {"LLM_API_KEY", "MODEL_API_KEY", "OPENROUTER_API_KEY"} & env.keys()
|
||||
|
|
@ -628,7 +644,10 @@ def test_paid_runner_uses_immutable_image_and_scrubs_ambient_alternate_keys(dry_
|
|||
dump = bench / "dumps" / launcher.dump_dir_name(args.model) / "SingleUserTurn-one"
|
||||
write_json(dump / "ouroboros_summary.json", {"bench_status": "success"})
|
||||
write_json(dump / "eval_res.json", {"pass": True})
|
||||
return {"stop_reason": "", "meter_error": "", "runner_exit_code": 0}
|
||||
counts = launcher.write_ledger(bench.parent / "result_index.jsonl", bench, args.model,
|
||||
args.selected_tasks, cause="runner_exited")
|
||||
return {"stop_reason": "", "meter_error": "", "runner_exit_code": 0,
|
||||
"interruption_cause": "runner_exited", "ledger_counts": counts}
|
||||
monkeypatch.setattr(launcher, "supervise_run", supervise)
|
||||
paid_argv = [item for item in argv if item != "--dry-run"]
|
||||
assert launcher.main([*paid_argv, "--campaign-file", str(out.parent / "campaign.json")]) == 0
|
||||
|
|
@ -675,7 +694,7 @@ def test_confirmation_never_accepts_an_invalid_numeric_usage(supervised, monkeyp
|
|||
values = iter([invalid, 101.0])
|
||||
monkeypatch.setattr(launcher, "key_usage", lambda _key, **_kwargs: next(values))
|
||||
diagnostics = []
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll")
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll", deadline=30.0)
|
||||
assert budget.record["last_usage"] == 101
|
||||
assert diagnostics[0]["accepted"] is False
|
||||
assert diagnostics[0]["previous_usage"] == 100
|
||||
|
|
@ -693,7 +712,8 @@ def test_persistent_bad_counter_stops_and_keeps_rejected_values_in_final_monitor
|
|||
result = launcher.supervise_run(args, ["fake-runner"], bench, env, "not-a-real-key", budget)
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert budget.record["last_usage"] == 100
|
||||
assert len(calls) == 6 # One bounded confirmation in the loop and one final settlement read.
|
||||
# Reads 3 s apart until the 30 s bound, then one bounded final settlement of the same shape.
|
||||
assert calls == ([5.0] * 9 + [3.0]) * 2
|
||||
assert {row["phase"] for row in result["meter_diagnostics"]} == {"poll", "final"}
|
||||
assert all(row["observed_usage"] == 99.0 and not row["accepted"] for row in result["meter_diagnostics"])
|
||||
final = json.loads((bench.parent / "monitor.json").read_text(encoding="utf-8"))
|
||||
|
|
@ -701,7 +721,7 @@ def test_persistent_bad_counter_stops_and_keeps_rejected_values_in_final_monitor
|
|||
assert final["meter_error"] == "UsageCounterError"
|
||||
|
||||
|
||||
def test_confirmation_shares_one_timeout_window_instead_of_three_full_timeouts(supervised, monkeypatch):
|
||||
def test_each_read_gets_a_slice_so_one_slow_read_cannot_consume_the_bound(supervised, monkeypatch):
|
||||
_args, _bench, _env, budget, _events, _handlers, _proc = supervised
|
||||
clock = SimpleNamespace(now=0.0)
|
||||
def advance(delay):
|
||||
|
|
@ -717,10 +737,10 @@ def test_confirmation_shares_one_timeout_window_instead_of_three_full_timeouts(s
|
|||
original = budget.path.read_bytes()
|
||||
diagnostics = []
|
||||
with pytest.raises(OSError, match="meter unavailable"):
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll")
|
||||
assert timeouts == [15.0, 5.0]
|
||||
assert clock.now == 15.0
|
||||
assert len(diagnostics) == 2
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll", deadline=30.0)
|
||||
assert timeouts == [5.0] * 4 # Reads at 0, 8, 16 and 24 s; the bound ends at exactly 30 s.
|
||||
assert clock.now == 30.0
|
||||
assert len(diagnostics) == 4
|
||||
assert budget.path.read_bytes() == original
|
||||
|
||||
|
||||
|
|
@ -800,8 +820,9 @@ def test_lagging_counter_gets_time_to_catch_up_without_weakening_monotonicity(su
|
|||
return 99.0 if clock.now < 3.0 else 101.0
|
||||
monkeypatch.setattr(launcher, "key_usage", provider)
|
||||
diagnostics = []
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll")
|
||||
assert reads == [(0.0, 15.0), (3.0, 12.0)]
|
||||
anchor = launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll", deadline=30.0)
|
||||
assert reads == [(0.0, 5.0), (3.0, 5.0)]
|
||||
assert anchor == 3.0 # The accepted read's request time, not the start of the confirmation.
|
||||
assert budget.record["last_usage"] == 101
|
||||
assert [(row["observed_usage"], row["accepted"]) for row in diagnostics] == [(99.0, False), (101.0, True)]
|
||||
assert clock.now == 3.0
|
||||
|
|
@ -857,8 +878,8 @@ def test_late_valid_confirmation_is_not_accepted_even_if_reader_returns_it(super
|
|||
return 101.0
|
||||
monkeypatch.setattr(launcher, "key_usage", late)
|
||||
diagnostics = []
|
||||
with pytest.raises(TimeoutError, match="shared deadline"):
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll")
|
||||
with pytest.raises(TimeoutError, match="after the blindness bound"):
|
||||
launcher.observe_campaign_usage("not-a-real-key", budget, diagnostics, phase="poll", deadline=4.0)
|
||||
assert budget.record["last_usage"] == 100
|
||||
assert diagnostics[0]["observed_usage"] == 101
|
||||
assert diagnostics[0]["accepted"] is False
|
||||
|
|
|
|||
424
tests/test_cowork_bench_meter_blindness.py
Normal file
424
tests/test_cowork_bench_meter_blindness.py
Normal file
|
|
@ -0,0 +1,424 @@
|
|||
"""Meter blindness, failure domains and interrupted-task rows through the real supervisor.
|
||||
|
||||
A simulated clock drives every read, retry pause and runner wait. The meter is a
|
||||
script, and the runner, its process-group stop and exact-label cleanup are recorded
|
||||
fakes, so no provider, Docker daemon or paid benchmark is reached.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import pathlib
|
||||
import subprocess
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from devtools.benchmarks.common import manifests
|
||||
from devtools.benchmarks.cowork_bench import campaign as budgets
|
||||
from devtools.benchmarks.cowork_bench import run_cowork_bench as launcher
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def forbid_live_services(monkeypatch):
|
||||
def forbidden(*_args, **_kwargs):
|
||||
pytest.fail("offline meter test tried to contact a provider or Docker")
|
||||
monkeypatch.setattr(budgets.urllib.request, "urlopen", forbidden)
|
||||
monkeypatch.setattr(launcher, "_docker", forbidden)
|
||||
|
||||
|
||||
def slow(sim, timeout):
|
||||
sim.now += timeout # The read consumes its whole slice, as in the #1259 incident.
|
||||
raise TimeoutError("meter HTTP read exceeded its wall-clock deadline")
|
||||
|
||||
|
||||
def offline(_sim, _timeout):
|
||||
raise OSError("meter unavailable")
|
||||
|
||||
|
||||
def regressed(_sim, _timeout):
|
||||
return 99.0
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def sim(tmp_path, monkeypatch):
|
||||
"""One run whose meter answers ``sim.script`` in order, then ``sim.default``."""
|
||||
state = SimpleNamespace(now=0.0, events=[], reads=[], script=[], runner_ends=None, campaign_writes=0,
|
||||
fail_campaign_write=lambda _index: False)
|
||||
state.default = lambda sim, _timeout: sim.budget.record["last_usage"]
|
||||
bench = tmp_path / "bench"
|
||||
bench.mkdir()
|
||||
args = launcher.parse_args([])
|
||||
args.resource_root = tmp_path
|
||||
args.docker_host = "unix:///owned.sock"
|
||||
args.selected_tasks = []
|
||||
args.min_free_gib = args.min_root_free_gib = 0
|
||||
env = {"COWORK_STOP_FILE": str(tmp_path / "resource_stop"), "COWORK_RUN_LABEL": "owned-run"}
|
||||
state.budget = budgets.CampaignBudget(tmp_path / "campaign.json", fingerprint="key-a", ceiling=1000, usage=100)
|
||||
state.args, state.bench, state.env = args, bench, env
|
||||
state.dumps = bench / "dumps" / launcher.dump_dir_name(args.model)
|
||||
state.stop_file = pathlib.Path(env["COWORK_STOP_FILE"])
|
||||
proc = SimpleNamespace(pid=424242, returncode=None)
|
||||
proc.poll = lambda: proc.returncode
|
||||
|
||||
def wait(*, timeout):
|
||||
if state.now > 3600:
|
||||
pytest.fail("the simulated run was never stopped")
|
||||
if state.runner_ends is not None and state.now + timeout >= state.runner_ends:
|
||||
state.now = max(state.now, state.runner_ends)
|
||||
proc.returncode = 0
|
||||
return 0
|
||||
state.now += timeout
|
||||
raise subprocess.TimeoutExpired("fake-runner", timeout)
|
||||
|
||||
def spawn(command, **kwargs):
|
||||
assert command == ["fake-runner"] and kwargs["purpose"] == "cowork-official-runner"
|
||||
state.events.append((state.now, "spawn"))
|
||||
return proc
|
||||
|
||||
def stop(owned):
|
||||
assert owned is proc
|
||||
state.events.append((state.now, "stop-group"))
|
||||
owned.returncode = -15
|
||||
|
||||
def cleanup(host, label):
|
||||
assert (host, label) == (args.docker_host, env["COWORK_RUN_LABEL"])
|
||||
state.events.append((state.now, "cleanup-owned"))
|
||||
|
||||
def meter(_key, *, timeout):
|
||||
state.reads.append((state.now, timeout))
|
||||
step = state.script.pop(0) if state.script else state.default
|
||||
return step(state, timeout) if callable(step) else step
|
||||
|
||||
real_write_json = manifests.write_json
|
||||
|
||||
def campaign_write(path, payload):
|
||||
if pathlib.Path(path) == state.budget.path:
|
||||
state.campaign_writes += 1
|
||||
if state.fail_campaign_write(state.campaign_writes):
|
||||
raise OSError(28, "No space left on device")
|
||||
real_write_json(path, payload)
|
||||
|
||||
proc.wait = wait
|
||||
monkeypatch.setattr(launcher, "spawn_supervised", spawn)
|
||||
monkeypatch.setattr(launcher, "stop_process_group", stop)
|
||||
monkeypatch.setattr(launcher, "remove_run_containers", cleanup)
|
||||
monkeypatch.setattr(launcher, "key_usage", meter)
|
||||
monkeypatch.setattr(manifests, "write_json", campaign_write) # The campaign's own writer.
|
||||
monkeypatch.setattr(launcher.signal, "signal", lambda _sig, _handler: "old-handler")
|
||||
monkeypatch.setattr(launcher, "time", SimpleNamespace(
|
||||
time=lambda: state.now, monotonic=lambda: state.now,
|
||||
sleep=lambda delay: setattr(state, "now", state.now + delay)))
|
||||
monkeypatch.setattr(launcher.shutil, "disk_usage", lambda _path: SimpleNamespace(free=1024**4))
|
||||
state.supervise = lambda: launcher.supervise_run(args, ["fake-runner"], bench, env, "not-a-real-key",
|
||||
state.budget)
|
||||
return state
|
||||
|
||||
|
||||
def lifecycle(sim) -> list[str]:
|
||||
"""Every owned effect happens once: one runner, one group stop, one exact-label cleanup."""
|
||||
return [name for _at, name in sim.events]
|
||||
|
||||
|
||||
def durable(sim) -> dict:
|
||||
return json.loads(sim.budget.path.read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def reads_between(sim, start, end) -> list[tuple[float, float]]:
|
||||
return [read for read in sim.reads if start <= read[0] < end]
|
||||
|
||||
|
||||
def test_healthy_meter_polls_every_fifteen_seconds_with_bounded_reads(sim):
|
||||
sim.runner_ends = 60.0
|
||||
sim.default = lambda s, _timeout: 100.0 + s.now / 10
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "" and result["meter_error"] == ""
|
||||
assert sim.reads == [(0.0, 5.0), (15.0, 5.0), (30.0, 5.0), (45.0, 5.0), (60.0, 5.0)]
|
||||
assert durable(sim)["last_usage"] == 106.0
|
||||
assert durable(sim)["runs"][-1]["outcome"] == "runner_finished"
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
def test_one_slow_read_uses_only_its_slice_and_a_second_read_confirms(sim):
|
||||
sim.runner_ends = 50.0
|
||||
sim.script = [101.0, slow, 102.0]
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == ""
|
||||
# The slow read at 15 s gets 5 s, not the whole remaining bound; a retry follows at 23 s.
|
||||
assert sim.reads[:3] == [(0.0, 5.0), (15.0, 5.0), (23.0, 5.0)]
|
||||
assert [(row["error_type"] if not row["accepted"] else "accepted") for row in result["meter_diagnostics"]] == [
|
||||
"TimeoutError", "accepted"]
|
||||
# The accepted read at 23 s anchors the next bound, so the run keeps polling after 30 s.
|
||||
assert (38.0, 5.0) in sim.reads
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
def test_immediate_transient_errors_retry_every_three_seconds_inside_the_bound(sim):
|
||||
sim.runner_ends = 40.0
|
||||
sim.script = [101.0, offline, offline, 102.0]
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == ""
|
||||
assert [at for at, _timeout in sim.reads[:4]] == [0.0, 15.0, 18.0, 21.0]
|
||||
assert durable(sim)["last_usage"] == 102.0
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("bound", [30.0, 45.0])
|
||||
def test_persistent_outage_stops_admission_and_run_at_the_bound_then_settles_bounded(sim, bound):
|
||||
sim.args.meter_blindness_sec = bound
|
||||
sim.script = [101.0]
|
||||
sim.default = offline
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert result["meter_error"] == "OSError"
|
||||
assert sim.stop_file.read_text(encoding="utf-8").strip() == "budget_meter_unavailable"
|
||||
# Blindness counts from the last confirmed, saved reading at 0 s, including waits and retries.
|
||||
assert sim.events == [(0.0, "spawn"), (bound, "stop-group"), (bound, "cleanup-owned")]
|
||||
# The final settlement is accounting only, after cleanup, and bounded by the same B.
|
||||
assert reads_between(sim, bound, 2 * bound) and sim.now == 2 * bound
|
||||
record = durable(sim)
|
||||
assert "active_run" not in record and record["last_usage"] == 101.0
|
||||
assert record["runs"][-1]["outcome"] == "budget_meter_unavailable"
|
||||
assert record["runs"][-1]["meter_error"] == "OSError"
|
||||
|
||||
|
||||
def test_repeated_counter_regression_never_resets_the_bound_or_lowers_spend(sim):
|
||||
sim.script = [101.0]
|
||||
sim.default = regressed
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert result["meter_error"] == "UsageCounterError"
|
||||
assert [at for at, _timeout in reads_between(sim, 1, 30)] == [15.0, 18.0, 21.0, 24.0, 27.0]
|
||||
assert (30.0, "stop-group") in sim.events
|
||||
rejected = [row for row in result["meter_diagnostics"] if not row["accepted"]]
|
||||
assert rejected and all(row["observed_usage"] == 99.0 and row["previous_usage"] == 101.0 for row in rejected)
|
||||
assert durable(sim)["last_usage"] == 101.0
|
||||
|
||||
|
||||
def test_late_success_is_rejected_and_final_settlement_records_it_after_cleanup(sim):
|
||||
def late(s, _timeout):
|
||||
s.now = 31.0 # A reader that returns a valid value only after the bound.
|
||||
return 150.0
|
||||
|
||||
sim.script = [101.0, late, 150.0]
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
late_row = result["meter_diagnostics"][0]
|
||||
assert late_row["observed_usage"] == 150.0 and late_row["accepted"] is False
|
||||
assert "after the blindness bound" in late_row["error"]
|
||||
assert sim.events == [(0.0, "spawn"), (31.0, "stop-group"), (31.0, "cleanup-owned")]
|
||||
# Only the post-cleanup settlement may record 150; it cannot authorize further work.
|
||||
assert sim.reads[-1] == (31.0, 5.0)
|
||||
assert result["campaign_spent_usd"] == 50.0
|
||||
assert durable(sim)["runs"][-1]["outcome"] == "budget_meter_unavailable"
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
def test_bound_counts_diagnostic_time_and_polls_early_enough_for_one_read(sim, monkeypatch):
|
||||
real_write_ledger = launcher.write_ledger
|
||||
|
||||
def slow_ledger(*args, **kwargs):
|
||||
sim.now += 16.0
|
||||
return real_write_ledger(*args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(launcher, "write_ledger", slow_ledger)
|
||||
sim.script = [101.0]
|
||||
sim.default = offline
|
||||
result = sim.supervise()
|
||||
# A 16 s diagnostic write shortens the next wait so one full read starts at 25 s.
|
||||
assert reads_between(sim, 1, 30) == [(25.0, 5.0), (28.0, 2.0)]
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert (30.0, "stop-group") in sim.events
|
||||
|
||||
|
||||
def test_transient_campaign_write_failure_stops_with_its_own_reason(sim):
|
||||
sim.runner_ends = 100.0
|
||||
sim.script = [101.0, 105.0]
|
||||
sim.fail_campaign_write = lambda index: index == 3 # start, first poll, then the failed save
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "campaign_persistence_failed"
|
||||
assert result["meter_error"] == "CampaignPersistenceError"
|
||||
assert sim.stop_file.read_text(encoding="utf-8").strip() == "campaign_persistence_failed"
|
||||
assert sim.events == [(0.0, "spawn"), (15.0, "stop-group"), (15.0, "cleanup-owned")]
|
||||
record = durable(sim)
|
||||
assert record["runs"][-1]["outcome"] == "campaign_persistence_failed"
|
||||
assert record["last_usage"] == 105.0 and "active_run" not in record
|
||||
|
||||
|
||||
def test_persistent_campaign_write_failure_keeps_custody_unsettled_after_exact_cleanup(sim):
|
||||
sim.runner_ends = 100.0
|
||||
sim.script = [101.0, 105.0]
|
||||
sim.fail_campaign_write = lambda index: index >= 3
|
||||
with pytest.raises(budgets.CampaignPersistenceError):
|
||||
sim.supervise()
|
||||
assert sim.stop_file.read_text(encoding="utf-8").strip() == "campaign_persistence_failed"
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
record = durable(sim)
|
||||
assert record["active_run"] == str(sim.bench.parent) and record["last_usage"] == 101.0
|
||||
sim.fail_campaign_write = lambda _index: False
|
||||
with pytest.raises(ValueError, match="unsettled custody"):
|
||||
budgets.CampaignBudget(sim.budget.path, fingerprint="key-a", ceiling=1000, usage=110)
|
||||
|
||||
|
||||
def test_diagnostic_write_failures_are_disclosed_and_valid_work_continues(sim, monkeypatch, capsys):
|
||||
def unavailable(*_args, **_kwargs):
|
||||
raise OSError(28, "No space left on device")
|
||||
|
||||
monkeypatch.setattr(launcher, "write_ledger", unavailable)
|
||||
monkeypatch.setattr(launcher, "write_json", unavailable)
|
||||
sim.runner_ends = 40.0
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "" and result["runner_exit_code"] == 0
|
||||
assert {name: row["count"] for name, row in result["diagnostic_write_failures"].items()} == {
|
||||
"result_index.jsonl": 4, "monitor.json": 4} # Polls at 0, 15 and 30 s plus the final record.
|
||||
stderr = [json.loads(line) for line in capsys.readouterr().err.splitlines() if line.startswith("{")]
|
||||
assert {row["artifact"] for row in stderr if row["event"] == "cowork_diagnostic_write_failed"} == {
|
||||
"result_index.jsonl", "monitor.json"}
|
||||
assert not (sim.bench.parent / "monitor.json").exists()
|
||||
assert durable(sim)["runs"][-1]["outcome"] == "runner_finished"
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
def test_unexpected_supervisor_failure_is_not_reported_as_a_finished_runner(sim, monkeypatch):
|
||||
polls = []
|
||||
|
||||
def disk(_path):
|
||||
polls.append(sim.now)
|
||||
if len(polls) > 2:
|
||||
raise PermissionError("disk probe denied")
|
||||
return SimpleNamespace(free=1024**4)
|
||||
|
||||
monkeypatch.setattr(launcher.shutil, "disk_usage", disk)
|
||||
with pytest.raises(PermissionError):
|
||||
sim.supervise()
|
||||
assert sim.stop_file.read_text(encoding="utf-8").strip() == "supervisor_error"
|
||||
assert durable(sim)["runs"][-1]["outcome"] == "supervisor_error"
|
||||
assert lifecycle(sim) == ["spawn", "stop-group", "cleanup-owned"]
|
||||
|
||||
|
||||
def test_empty_stop_request_is_named_rather_than_runner_finished(sim, monkeypatch):
|
||||
def spawn(_command, **_kwargs):
|
||||
sim.stop_file.write_text("", encoding="utf-8")
|
||||
sim.events.append((sim.now, "spawn"))
|
||||
return SimpleNamespace(pid=1, returncode=None, poll=lambda: None)
|
||||
|
||||
monkeypatch.setattr(launcher, "stop_process_group", lambda _proc: sim.events.append((sim.now, "stop-group")))
|
||||
monkeypatch.setattr(launcher, "spawn_supervised", spawn)
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "stop_requested"
|
||||
assert durable(sim)["runs"][-1]["outcome"] == "stop_requested"
|
||||
|
||||
|
||||
def write_task(dumps: pathlib.Path, task: str, files: dict[str, object]) -> None:
|
||||
root = dumps / f"SingleUserTurn-{task}"
|
||||
(root / "workspace").mkdir(parents=True) # The benchmark pre-creates both directories.
|
||||
for name, value in files.items():
|
||||
path = root / name
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
if isinstance(value, list):
|
||||
path.write_text("".join(json.dumps(row) + "\n" for row in value), encoding="utf-8")
|
||||
else:
|
||||
path.write_text(json.dumps(value), encoding="utf-8")
|
||||
|
||||
|
||||
USAGE = {"type": "llm_usage", "prompt_tokens": 1200, "completion_tokens": 30, "cost": 0.01}
|
||||
|
||||
|
||||
def test_meter_stop_marks_started_tasks_interrupted_with_honest_paid_activity(sim):
|
||||
write_task(sim.dumps, "paid", {"applied_settings.json": {}, "ouroboros/events.jsonl": [{"type": "task"}, USAGE]})
|
||||
write_task(sim.dumps, "admitted", {"applied_settings.json": {}, "ouroboros/events.jsonl": [{"type": "task"}]})
|
||||
write_task(sim.dumps, "precreated", {})
|
||||
write_task(sim.dumps, "finished", {"ouroboros_summary.json": {"bench_status": "success"},
|
||||
"eval_res.json": {"pass": True}})
|
||||
sim.args.selected_tasks = ["paid", "admitted", "precreated", "never", "finished"]
|
||||
sim.script = [101.0]
|
||||
sim.default = offline
|
||||
result = sim.supervise()
|
||||
assert result["interruption_cause"] == "budget_meter_unavailable"
|
||||
rows = {row["instance_id"]: row for row in map(json.loads, (sim.bench.parent / "result_index.jsonl")
|
||||
.read_text(encoding="utf-8").splitlines())}
|
||||
assert {task: (row["status"], row["reason_code"], row["details"].get("paid_activity"))
|
||||
for task, row in rows.items()} == {
|
||||
"paid": ("infra_failed", "interrupted:budget_meter_unavailable", "observed"),
|
||||
"admitted": ("infra_failed", "interrupted:budget_meter_unavailable", "unknown"),
|
||||
"precreated": ("not_attempted", "missing_result", None),
|
||||
"never": ("not_attempted", "missing_result", None),
|
||||
"finished": ("passed", "passed", None),
|
||||
}
|
||||
assert rows["paid"]["details"]["start_evidence"] == ["ouroboros", "applied_settings.json"]
|
||||
assert not any("cost" in key for row in rows.values() for key in row["details"])
|
||||
# Interrupted rows stay unsettled for an explicit recovery; the passed task is never repeated.
|
||||
assert launcher.settled_tasks([sim.bench.parent]) == {"finished"}
|
||||
|
||||
|
||||
def test_task_without_summary_after_the_runner_exited_names_that_cause(sim):
|
||||
write_task(sim.dumps, "cut", {"mcp_proxy.log": "", "ouroboros/events.jsonl": [USAGE]})
|
||||
sim.args.selected_tasks = ["cut"]
|
||||
sim.runner_ends = 20.0
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "" and result["interruption_cause"] == "runner_exited"
|
||||
row = json.loads((sim.bench.parent / "result_index.jsonl").read_text(encoding="utf-8"))
|
||||
assert (row["status"], row["reason_code"], row["details"]["paid_activity"]) == (
|
||||
"infra_failed", "interrupted:runner_exited", "observed")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("events,expected", [
|
||||
([{"type": "llm_usage", "prompt_tokens": 0, "completion_tokens": 0}], "unknown"),
|
||||
([{"type": "llm_usage", "usage": {"prompt_tokens": 7, "completion_tokens": 1}}], "observed"),
|
||||
([{"type": "llm_usage"}], "unknown"),
|
||||
])
|
||||
def test_only_token_bearing_usage_is_observed_paid_activity(tmp_path, events, expected):
|
||||
write_task(tmp_path, "task", {"ouroboros/events.jsonl": events})
|
||||
row = launcher.ledger_row("task", tmp_path / "SingleUserTurn-task", {}, cause="in_progress")
|
||||
assert (row["status"], row["reason_code"]) == ("infra_failed", "missing_adapter_summary")
|
||||
assert row["details"]["provisional"] is True
|
||||
assert row["details"]["paid_activity"] == expected
|
||||
|
||||
|
||||
def test_runner_row_without_summary_names_interruption_and_discloses_activity(tmp_path):
|
||||
write_task(tmp_path, "task", {"ouroboros/events.jsonl": [USAGE]})
|
||||
row = launcher.ledger_row("task", tmp_path / "SingleUserTurn-task", {"status": "unknown"}, cause="signal_15")
|
||||
assert (row["status"], row["reason_code"], row["details"]["paid_activity"]) == (
|
||||
"infra_failed", "interrupted:signal_15", "observed")
|
||||
assert row["details"]["provisional"] is False
|
||||
|
||||
|
||||
def test_slow_persistence_cannot_renew_an_expired_window(sim, monkeypatch):
|
||||
real_observe = sim.budget.observe
|
||||
|
||||
def slow_save(usage):
|
||||
if sim.now == 15.0:
|
||||
sim.now = 31.0
|
||||
real_observe(usage)
|
||||
|
||||
monkeypatch.setattr(sim.budget, "observe", slow_save)
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert sim.events == [(0.0, "spawn"), (31.0, "stop-group"), (31.0, "cleanup-owned")]
|
||||
assert [at for at, _ in sim.reads] == [0.0, 15.0, 31.0] # Last is accounting-only.
|
||||
|
||||
|
||||
def test_startup_persistence_cannot_spawn_with_an_expired_observation(sim, monkeypatch):
|
||||
real_start = sim.budget.start
|
||||
|
||||
def slow_start(root):
|
||||
sim.now = 31.0
|
||||
real_start(root)
|
||||
|
||||
monkeypatch.setattr(sim.budget, "start", slow_start)
|
||||
with pytest.raises(TimeoutError, match="before runner spawn"):
|
||||
sim.supervise()
|
||||
assert lifecycle(sim) == ["cleanup-owned"]
|
||||
assert sim.stop_file.read_text(encoding="utf-8").strip() == "budget_meter_unavailable"
|
||||
|
||||
|
||||
def test_overlong_diagnostics_stop_on_return_even_if_runner_finished(sim, monkeypatch):
|
||||
def slow_ledger(*_args, **_kwargs):
|
||||
sim.now += 31.0
|
||||
return {}
|
||||
|
||||
monkeypatch.setattr(launcher, "write_ledger", slow_ledger)
|
||||
result = sim.supervise()
|
||||
assert result["stop_reason"] == "budget_meter_unavailable"
|
||||
assert sim.events == [(0.0, "spawn"), (31.0, "stop-group"), (31.0, "cleanup-owned")]
|
||||
assert [at for at, _ in sim.reads] == [0.0, 31.0]
|
||||
Loading…
Add table
Add a link
Reference in a new issue