From d9b9157fa4b736f60ebb665a1f8dd5dea876b269 Mon Sep 17 00:00:00 2001 From: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com> Date: Fri, 25 Sep 2026 01:22:52 +0300 Subject: [PATCH] fix(cowork): bound meter blindness and retain interrupted activity --- .../benchmarks/cowork_bench/METHODOLOGY.md | 41 +- devtools/benchmarks/cowork_bench/README.md | 21 +- .../cowork_bench/audit_cowork_bench.py | 32 +- devtools/benchmarks/cowork_bench/campaign.py | 10 +- .../cowork_bench/run_cowork_bench.py | 207 ++++++--- tests/test_cowork_bench_launcher.py | 51 ++- tests/test_cowork_bench_meter_blindness.py | 424 ++++++++++++++++++ 7 files changed, 699 insertions(+), 87 deletions(-) create mode 100644 tests/test_cowork_bench_meter_blindness.py diff --git a/devtools/benchmarks/cowork_bench/METHODOLOGY.md b/devtools/benchmarks/cowork_bench/METHODOLOGY.md index 2e58d48c6..fc1238a90 100644 --- a/devtools/benchmarks/cowork_bench/METHODOLOGY.md +++ b/devtools/benchmarks/cowork_bench/METHODOLOGY.md @@ -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:`: +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 diff --git a/devtools/benchmarks/cowork_bench/README.md b/devtools/benchmarks/cowork_bench/README.md index 22a5c194d..b7294d509 100644 --- a/devtools/benchmarks/cowork_bench/README.md +++ b/devtools/benchmarks/cowork_bench/README.md @@ -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:` 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 diff --git a/devtools/benchmarks/cowork_bench/audit_cowork_bench.py b/devtools/benchmarks/cowork_bench/audit_cowork_bench.py index 149536faa..5025d1878 100644 --- a/devtools/benchmarks/cowork_bench/audit_cowork_bench.py +++ b/devtools/benchmarks/cowork_bench/audit_cowork_bench.py @@ -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 diff --git a/devtools/benchmarks/cowork_bench/campaign.py b/devtools/benchmarks/cowork_bench/campaign.py index 236c0df9a..ea12a3703 100644 --- a/devtools/benchmarks/cowork_bench/campaign.py +++ b/devtools/benchmarks/cowork_bench/campaign.py @@ -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() diff --git a/devtools/benchmarks/cowork_bench/run_cowork_bench.py b/devtools/benchmarks/cowork_bench/run_cowork_bench.py index 042cc732f..ea9468a1e 100644 --- a/devtools/benchmarks/cowork_bench/run_cowork_bench.py +++ b/devtools/benchmarks/cowork_bench/run_cowork_bench.py @@ -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 diff --git a/tests/test_cowork_bench_launcher.py b/tests/test_cowork_bench_launcher.py index a7f2b8432..74a7db6ff 100644 --- a/tests/test_cowork_bench_launcher.py +++ b/tests/test_cowork_bench_launcher.py @@ -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 diff --git a/tests/test_cowork_bench_meter_blindness.py b/tests/test_cowork_bench_meter_blindness.py new file mode 100644 index 000000000..d8da50b53 --- /dev/null +++ b/tests/test_cowork_bench_meter_blindness.py @@ -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]