Coverage for src/ai_jury/runagent.py: 100%

231 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-30 06:29 +0000

1"""Single-agent role dispatch for orchestrators: ``jury run-agent`` (issue #661). 

2 

3ai-jury already owns the transport-agnostic provider runtime an orchestrator 

4needs — an adapter per vendor, per-vendor availability, the read-only flags, 

5timeouts, the typed error taxonomy — but until now it was reachable only through 

6a full panel run. So an orchestrator that wanted to dispatch *one* implementer, 

7gate reviewer or chair had to hand-write the argv itself, and hand-written argv 

8is exactly where a read-only guarantee gets lost. 

9 

10This module is the pure core of that command: role policy, agent resolution, 

11attribution, and the on-disk shape of a detached run. Everything that touches a 

12subprocess, the clock or the filesystem lives in ``cli._run_run_agent`` or is 

13passed in as a seam, so all of the below is unit-testable offline. 

14 

15The security-load-bearing piece is :func:`role_policy`. ``review``/``gate``/ 

16``chair`` are read-only *always* — the same invocation a panel run uses, with no 

17way to widen it — because those roles read attacker-controlled content. 

18``implement``/``fix`` are the only roles that can write, and only when the 

19operator passes ``--allow-write``. 

20""" 

21 

22from __future__ import annotations 

23 

24import contextlib 

25import json 

26import os 

27import re 

28import secrets 

29import time 

30from dataclasses import dataclass, replace 

31from pathlib import Path 

32 

33from .config import AGY_AGENT, DEFAULT_CONFIG, AgentSpec 

34 

35#: Stable schema identifier for the ``jury run-agent`` JSON result. Bump this 

36#: when a key changes meaning or disappears; ``tests/test_cli_run_agent.py`` 

37#: pins the whole shape key-by-key so that has to be deliberate. 

38SCHEMA_VERSION = "ai-jury.run-agent.v1" 

39 

40#: Every role an orchestrator can dispatch. 

41ROLES: tuple[str, ...] = ("implement", "review", "gate", "chair", "fix") 

42 

43#: The roles that may modify a working tree — and only with ``--allow-write``. 

44WRITE_ROLES = frozenset({"implement", "fix"}) 

45 

46#: The roles that are read-only unconditionally. These read attacker-controlled 

47#: content (a diff, a PR body, another agent's output), so no flag widens them. 

48READ_ONLY_ROLES: tuple[str, ...] = tuple(r for r in ROLES if r not in WRITE_ROLES) 

49 

50 

51# --- role policy ------------------------------------------------------------- 

52 

53 

54@dataclass(frozen=True) 

55class RolePolicy: 

56 """What one role is allowed to do (pure value). 

57 

58 ``refusal`` set means the run must not start at all (exit 2). ``warning`` is 

59 advisory — the run proceeds, but the operator asked for something that was 

60 not honoured. 

61 """ 

62 

63 role: str 

64 write: bool = False 

65 refusal: str | None = None 

66 warning: str | None = None 

67 

68 

69def role_policy(role: str, allow_write: bool = False) -> RolePolicy: 

70 """Resolve ``(role, allow_write)`` to a role policy (pure). 

71 

72 The one place the read-only/write decision is made, so no adapter has to 

73 re-derive it: 

74 

75 * ``review``/``gate``/``chair`` → read-only, always. ``--allow-write`` is 

76 **ignored** rather than obeyed (with a warning), because a reviewer of an 

77 attacker-controlled diff must not become write-capable by way of a flag — 

78 the whole point of the least-privilege posture ``privilege.py`` audits. 

79 * ``implement``/``fix`` → write-capable, but only with ``--allow-write``. 

80 Without it the run is refused: silently downgrading an implementer to a 

81 read-only agent produces a confident report of work that never happened. 

82 """ 

83 name = (role or "").strip().lower() 

84 if name not in ROLES: 

85 return RolePolicy( 

86 role=name, 

87 refusal=f"unknown role '{role}'; expected one of {', '.join(ROLES)}", 

88 ) 

89 if name in WRITE_ROLES: 

90 if not allow_write: 

91 return RolePolicy( 

92 role=name, 

93 refusal=( 

94 f"role '{name}' needs write access to be useful; re-run with " 

95 f"--allow-write to grant it. Read-only roles " 

96 f"({', '.join(READ_ONLY_ROLES)}) never need it." 

97 ), 

98 ) 

99 return RolePolicy(role=name, write=True) 

100 warning = None 

101 if allow_write: 

102 warning = ( 

103 f"--allow-write ignored for read-only role '{name}': " 

104 f"{'/'.join(READ_ONLY_ROLES)} always run under the vendor's read-only " 

105 f"invocation, the same one a panel review uses." 

106 ) 

107 return RolePolicy(role=name, write=False, warning=warning) 

108 

109 

110# --- agent resolution -------------------------------------------------------- 

111 

112#: Bare vendor tokens ``--agent`` accepts with no ``[[agent]]`` entry at all, so 

113#: ``jury run-agent --agent claude --role review`` works in a repo with no 

114#: ``jury.toml``. CLI vendors reuse the shipped default entry (including its 

115#: read-only ``extra_args``); the hosted-API vendors need only their env key. 

116BUILTIN_AGENTS: tuple[str, ...] = ( 

117 "claude", 

118 "codex", 

119 "agy", 

120 "anthropic-api", 

121 "openai-api", 

122 "google-api", 

123 "xai-api", 

124) 

125 

126#: A model id as the vendors actually spell them: ``gemini-3.8-flash-high``, 

127#: ``qwen2.5-coder:7b``, ``anthropic/claude-x``. Anchored, and required to start 

128#: with an alphanumeric so a token can never be read as a flag by the CLI it is 

129#: forwarded to. 

130_MODEL_TOKEN_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:/-]*$") 

131 

132 

133def parse_agent_token(token: str) -> tuple[str, str | None, str | None]: 

134 """Split ``name`` / ``name:model`` (pure). 

135 

136 Returns ``(name, model_or_None, error_or_None)``. Only the FIRST colon 

137 separates, because a model id may legitimately contain one (an Ollama 

138 ``:tag``, a Bedrock ``…-v1:0``). 

139 """ 

140 raw = (token or "").strip() 

141 if not raw: 

142 return "", None, "--agent must name an agent (e.g. claude, or codex:gpt-5.2)" 

143 name, sep, model = raw.partition(":") 

144 name = name.strip() 

145 if not name: 

146 return "", None, f"invalid --agent '{token}': the agent name is empty" 

147 if not sep: 

148 return name, None, None 

149 model = model.strip() 

150 if not _MODEL_TOKEN_RE.match(model): 

151 return ( 

152 name, 

153 None, 

154 f"invalid model '{model}' in --agent '{token}': expected letters, " 

155 f"digits, '.', '_', ':', '/' or '-', starting with a letter or digit", 

156 ) 

157 return name, model, None 

158 

159 

160def builtin_spec(name: str) -> AgentSpec | None: 

161 """The default :class:`AgentSpec` for a bare built-in vendor token (pure). 

162 

163 A CLI vendor reuses the entry from :data:`config.DEFAULT_CONFIG` — the same 

164 command and the same read-only ``extra_args`` a default panel would use — so 

165 a bare ``--agent claude`` and a configured ``[[agent]] name = "claude"`` 

166 cannot drift apart. ``agy`` is not in that panel (it cannot be confined for 

167 untrusted diffs), but naming it here is an explicit single dispatch, so it 

168 resolves to :data:`config.AGY_AGENT`, the entry ``jury init`` writes. 

169 """ 

170 key = (name or "").strip().lower() 

171 if key not in BUILTIN_AGENTS: 

172 return None 

173 for raw in [*DEFAULT_CONFIG["agent"], AGY_AGENT]: 

174 if raw["name"] == key: 

175 return AgentSpec( 

176 name=raw["name"], 

177 vendor=raw["vendor"], 

178 command=raw.get("command", ""), 

179 extra_args=list(raw.get("extra_args", [])), 

180 ) 

181 # A hosted-API vendor: no command, no extra_args, keyed by its env var. The 

182 # model comes from `vendor:model` (or a [[agent]] entry). 

183 return AgentSpec(name=key, vendor=key) 

184 

185 

186def resolve_agent(config, token: str) -> tuple[AgentSpec | None, str | None]: 

187 """Resolve ``--agent`` to a concrete :class:`AgentSpec` (pure). 

188 

189 Precedence: a ``[[agent]]`` entry whose ``name`` matches wins over a 

190 built-in vendor of the same name, so an operator who configured 

191 ``name = "claude"`` with their own model/flags gets exactly that. A 

192 ``name:model`` suffix overrides the resolved spec's model either way. 

193 

194 A configured agent is matched whether or not it is ``enabled``: that flag 

195 selects the *panel*, and ``run-agent`` is an explicit single dispatch — the 

196 caller named this agent on purpose. 

197 """ 

198 name, model, error = parse_agent_token(token) 

199 if error is not None: 

200 return None, error 

201 

202 spec = None 

203 for candidate in getattr(config, "agents", None) or []: 

204 if candidate.name == name: 

205 spec = candidate 

206 break 

207 if spec is None: 

208 spec = builtin_spec(name) 

209 if spec is None: 

210 return None, ( 

211 f"unknown agent '{name}': add a [[agent]] entry named '{name}' to " 

212 f"jury.toml, or use a built-in vendor ({', '.join(BUILTIN_AGENTS)})" 

213 ) 

214 if model: 

215 spec = replace(spec, model=model) 

216 return spec, None 

217 

218 

219def transport_for(spec: AgentSpec) -> str: 

220 """How this agent is reached: ``cli``, ``api`` or ``local`` (pure). 

221 

222 Deliberately the same vocabulary — and the same rule — as 

223 ``jury --doctor --json``, so an orchestrator reading both sees one answer. 

224 """ 

225 from .doctor import _transport 

226 

227 # Keyed on the adapter (#705), like doctor's own row: how a seat is reached 

228 # follows the protocol it is invoked through, not its vendor identity. 

229 return _transport(spec.adapter_key, spec.command, spec.endpoint) 

230 

231 

232# --- attribution ------------------------------------------------------------- 

233 

234#: Prefixes that name a TRANSPORT rather than a model, so ``anthropic-api:`` 

235#: in ``anthropic-api:claude-opus-4-5`` is stripped while the ``:7b`` in 

236#: ``qwen2.5-coder:7b`` is not. Mirrors keel's ``agents._TRANSPORT_PREFIXES``, 

237#: which is what consumes these labels. 

238_TRANSPORT_PREFIXES = frozenset( 

239 { 

240 "anthropic-api", 

241 "openai-api", 

242 "google-api", 

243 "xai-api", 

244 "openai-compatible", 

245 "ollama", 

246 "local", 

247 } 

248) 

249 

250 

251def strip_transport(model: str) -> str: 

252 """Drop a ``<transport>:`` prefix, leaving the vendor's own model id (pure). 

253 

254 Both ``ollama:qwen2.5:7b`` and ``anthropic-api:claude-opus-4-5`` carry a 

255 colon and only the second has the model on the right, so the colon is read 

256 by what sits on either side of it, never by position. A colon NOT preceded 

257 by a transport belongs to the model. 

258 """ 

259 text = (model or "").strip().lower() 

260 head, sep, tail = text.partition(":") 

261 if sep and tail and head in _TRANSPORT_PREFIXES: 

262 return tail 

263 return text 

264 

265 

266def model_base(model: str) -> str: 

267 """Strip a model id to a coarse **family + major** label (pure). 

268 

269 ``qwen2.5:7b`` → ``qwen``, ``gemma2`` → ``gemma``, ``gpt-5.5`` → ``gpt-5``, 

270 ``anthropic-api:claude-opus-4-5`` → ``claude-opus-4-5``. The label has to 

271 stay stable across a vendor's point releases, or every model bump would 

272 fork the attribution history of otherwise-identical work. 

273 

274 **The grouping is deliberately coarse and deliberately uneven**, and it is 

275 not a bug to be smoothed out here. The rule is byte-for-byte keel's 

276 ``agents.model_base`` (ship #2036) because the two projects write the same 

277 ``model:<base>`` label onto the same issues — a "better" rule on one side 

278 only would silently split one project's history in half. Consequences worth 

279 knowing before you rely on the label: 

280 

281 * A tier or effort suffix **collapses**: ``gemini-3.8-flash``, 

282 ``gemini-3.8-flash-high`` and ``gemini-3.8-pro`` all become ``gemini-3``, 

283 because everything after the first hyphen is cut at the first ``.``. 

284 * A vendor that spells its version with hyphens instead **keeps** it: 

285 ``claude-opus-4-5`` and ``claude-opus-4-6`` stay distinct. 

286 

287 So the label answers "roughly which family ran this", not "exactly which 

288 model". When the exact id matters, read the ``model`` field of the result 

289 document, which carries it verbatim. 

290 """ 

291 text = strip_transport(model) 

292 if not text: 

293 return "" 

294 text = text.split(":", 1)[0] # an Ollama :tag 

295 if "-" in text: 

296 # A hyphenated family: keep <word>-<major>, drop the .minor. 

297 head, _, tail = text.partition("-") 

298 return f"{head}-{tail.split('.', 1)[0]}" 

299 # Otherwise drop the trailing numeric run (digits and dots). 

300 i = len(text) 

301 while i > 0 and (text[i - 1].isdigit() or text[i - 1] == "."): 

302 i -= 1 

303 return text[:i] 

304 

305 

306def attribution(vendor: str, model: str | None = None) -> dict: 

307 """Who did the work: ``{vendor, model, label}`` (pure). 

308 

309 ``label`` is the space-joined pair an orchestrator applies verbatim — 

310 ``agent:<vendor>`` plus a versionless ``model:<base>`` when a model is known. 

311 Split it on whitespace to get the individual labels back. 

312 """ 

313 parts = [f"agent:{vendor}"] 

314 base = model_base(model or "") 

315 if base: 

316 parts.append(f"model:{base}") 

317 return {"vendor": vendor, "model": model or None, "label": " ".join(parts)} 

318 

319 

320# --- result shape ------------------------------------------------------------ 

321 

322 

323def result_dict(spec: AgentSpec, role: str, result) -> dict: 

324 """Project an :class:`~ai_jury.adapters.AgentResult` onto the v1 export (pure). 

325 

326 Every key is always present, with a stable type, so a consumer can read the 

327 document without probing for optional fields. 

328 

329 ``model`` is **the id this invocation sent** (:attr:`AgentResult.model`, 

330 #709) whenever the run recorded one, and only otherwise the configured 

331 ``spec.model``. Those two differ exactly where reporting the wrong one 

332 matters: an ``effort`` level encoded as a model-id suffix, or an adapter 

333 that consulted the vendor's live listing and fell back. This export said 

334 ``gemini-3-pro`` for a run that sent ``gemini-3-pro-high`` (#722), so an 

335 orchestrator reading it back saw what it had asked for rather than what 

336 answered. ``attribution`` is derived from the same string, because two 

337 fields naming one model must not name two. 

338 """ 

339 from .adapters import ERR_TIMEOUT 

340 

341 model = (getattr(result, "model", "") or "").strip() or spec.model 

342 return { 

343 "schema_version": SCHEMA_VERSION, 

344 "ok": bool(result.ok), 

345 "agent": spec.name, 

346 "vendor": spec.vendor, 

347 "model": model or None, 

348 "role": role, 

349 "transport": transport_for(spec), 

350 "text": result.output or "", 

351 "exit_code": result.exit_code, 

352 "duration_s": round(float(result.duration_s), 3), 

353 "timed_out": result.error_code == ERR_TIMEOUT, 

354 "error_code": result.error_code, 

355 "error": result.error, 

356 "attribution": attribution(spec.vendor, model), 

357 } 

358 

359 

360# --- detached runs ----------------------------------------------------------- 

361 

362#: Subdirectory of the cache dir that holds detached-run state. The cache dir is 

363#: already the project's "derived local state" location, already documented as 

364#: sensitive, and already overridable with ``--cache-dir``/``$JURY_CACHE_DIR``. 

365RUNS_DIR_NAME = "run-agent" 

366 

367#: A run id has to be safe as a filename: no separators, no traversal, bounded. 

368_RUN_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$") 

369 

370#: Status values a state file can carry. 

371STATUS_RUNNING = "running" 

372STATUS_DONE = "done" 

373#: A run still marked running whose child process is gone (see run_summary). 

374STATUS_LOST = "lost" 

375 

376 

377def new_run_id(token_fn=None) -> str: 

378 """A fresh run id (12 hex characters).""" 

379 return (token_fn or secrets.token_hex)(6) 

380 

381 

382class RunIdError(ValueError): 

383 """An identifier that must never be turned into a path (issue #661). 

384 

385 Its own type so a caller can tell "you named a run badly" apart from the 

386 ``ValueError`` a corrupt JSON body raises. 

387 """ 

388 

389 

390def check_run_id(run_id) -> str | None: 

391 """None when ``run_id`` is usable as a filename, else why not (pure).""" 

392 if not isinstance(run_id, str) or not _RUN_ID_RE.match(run_id): 

393 return ( 

394 f"invalid run id '{run_id}': use letters, digits, '.', '_' or '-' " 

395 f"(max 64 characters, starting with a letter or digit)" 

396 ) 

397 return None 

398 

399 

400def runs_dir(cache_dir=None) -> Path: 

401 """The directory holding detached-run state files.""" 

402 from .cache import default_cache_dir 

403 

404 base = Path(cache_dir) if cache_dir else default_cache_dir() 

405 return base / RUNS_DIR_NAME 

406 

407 

408def run_path(run_id, suffix: str, cache_dir=None) -> Path: 

409 """The path of one run's file, or raise :class:`RunIdError` (issue #661). 

410 

411 **The only place a run id becomes a path**, so every caller — the detach 

412 parent, the ``--_child`` writer, ``--wait``, ``--status`` — inherits the 

413 same check rather than each remembering to make it. Validating in the 

414 parent alone was not enough: the child re-parses its own ``--run-id`` and 

415 wrote wherever it pointed, so ``--run-id ../../PWNED`` (or an absolute 

416 path) put a file holding the agent's full output outside the cache 

417 directory entirely. 

418 

419 Two barriers, because one of them can be loosened by a future edit to a 

420 regex: the id must match :data:`_RUN_ID_RE` (no separator, no leading dot, 

421 so neither ``..`` nor an absolute path can be spelled), and the resolved 

422 file must still sit directly inside the runs directory. 

423 """ 

424 problem = check_run_id(run_id) 

425 if problem: 

426 raise RunIdError(problem) 

427 directory = runs_dir(cache_dir) 

428 path = directory / f"{run_id}{suffix}" 

429 # Belt and braces: `..` and `/` are already unspellable above, but a path 

430 # that escaped anyway must not be written to. Compared without touching the 

431 # filesystem, so a missing directory is not an error here. 

432 if os.path.normpath(path.parent) != os.path.normpath(directory): 

433 raise RunIdError(f"invalid run id '{run_id}': resolves outside {directory}") 

434 return path 

435 

436 

437def state_path(run_id, cache_dir=None) -> Path: 

438 return run_path(run_id, ".json", cache_dir) 

439 

440 

441def output_path(run_id, cache_dir=None) -> Path: 

442 return run_path(run_id, ".out", cache_dir) 

443 

444 

445def write_state(run_id, state: dict, cache_dir=None) -> Path: 

446 """Write a run's state file atomically, 0600, in a 0700 directory. 

447 

448 A state file holds the agent's full output, which is derived from the 

449 prompt — the same trust level as the diff a review sees — so it is written 

450 with the same restrictive permissions the result cache uses. The rename is 

451 atomic so a ``--wait`` poller never reads a half-written document. 

452 

453 Raises :class:`RunIdError` before creating anything when the id is not a 

454 safe filename; a write to an unvalidated location must fail loudly. 

455 """ 

456 final = state_path(run_id, cache_dir) 

457 directory = final.parent 

458 directory.mkdir(parents=True, exist_ok=True) 

459 with contextlib.suppress(OSError): 

460 directory.chmod(0o700) 

461 tmp = directory / f"{final.name}.{secrets.token_hex(4)}.tmp" 

462 fd = os.open(tmp, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) 

463 with os.fdopen(fd, "w", encoding="utf-8") as fh: 

464 json.dump(state, fh, indent=2) 

465 fh.write("\n") 

466 tmp.replace(final) 

467 return final 

468 

469 

470#: Upper bound on a state file read. A result document is small; this only 

471#: bounds memory against a corrupt or hostile file in a shared cache dir. 

472_MAX_STATE_BYTES = 16 * 1024 * 1024 

473 

474 

475def read_state(run_id, cache_dir=None) -> dict | None: 

476 """Read one run's state, or None when it is missing/unreadable/corrupt. 

477 

478 An unsafe id is a miss, not a raise: a read is a lookup, and "there is no 

479 such run" is the honest answer for a name that could never have named one. 

480 """ 

481 try: 

482 path = state_path(run_id, cache_dir) 

483 with path.open("r", encoding="utf-8") as fh: 

484 raw = fh.read(_MAX_STATE_BYTES + 1) 

485 if len(raw) > _MAX_STATE_BYTES: 

486 return None 

487 data = json.loads(raw) 

488 except (OSError, ValueError): # RunIdError is a ValueError 

489 return None 

490 return data if isinstance(data, dict) else None 

491 

492 

493class _DefaultProbe: 

494 """Sentinel: resolve :func:`pid_alive` when the call happens. 

495 

496 Distinct from ``None`` on purpose. ``alive_fn=None`` used to mean "use the 

497 default" in :func:`wait_for_run`/:func:`list_runs` but "do not probe at all" 

498 in :func:`run_summary` — the same value with opposite meanings, one function 

499 apart. Now ``None`` means "do not probe" everywhere and this means "use the 

500 default", which cannot be confused for a caller-supplied probe. 

501 """ 

502 

503 def __repr__(self) -> str: # pragma: no cover - debugging aid 

504 return "<default probe>" 

505 

506 

507#: See :class:`_DefaultProbe`. 

508DEFAULT_PROBE = _DefaultProbe() 

509 

510 

511def _resolve_probe(alive_fn): 

512 """The probe to call: the module-level default, a caller's, or none.""" 

513 return pid_alive if alive_fn is DEFAULT_PROBE else alive_fn 

514 

515 

516def pid_alive(pid) -> bool | None: 

517 """Is that process still running? ``None`` when we cannot safely tell. 

518 

519 POSIX only. On Windows ``os.kill(pid, 0)`` does **not** mean "probe": for 

520 any signal other than the two console events CPython opens the process and 

521 calls ``TerminateProcess``, so a liveness check there would kill the very 

522 run it was asked about. Returning ``None`` (unknown) leaves the run 

523 reported as ``running``, which is the safe reading of "we don't know". 

524 

525 **A live pid is not proof of identity.** Pids are recycled, so a state file 

526 that outlives its process — across a reboot, or in a ``--cache-dir`` shared 

527 between machines — can name a pid some unrelated process now holds, and this 

528 reports ``True`` for it. That is why the answer is only ever used to demote 

529 ``running`` to ``lost`` on a definite ``False``: a wrong ``True`` leaves a 

530 finished-looking run reported as still running (recoverable, and the state 

531 file carries ``started_at`` so a reader can see how old the claim is), while 

532 a wrong ``False`` would declare a live run dead. Ruling out recycling needs 

533 the process start time, which the standard library does not expose 

534 portably — and a runtime dependency is not on the table for this. 

535 

536 ``bool`` is excluded explicitly because it subclasses ``int``: without that 

537 guard ``pid_alive(True)`` probes pid 1, which on POSIX raises 

538 ``PermissionError`` and so reports "alive" — while 

539 :func:`liveness_unknown_reason` calls the very same state document 

540 unrecorded. Two classifications of one state that disagree is a drift the 

541 guard costs one term to prevent. 

542 """ 

543 if os.name == "nt" or not isinstance(pid, int) or isinstance(pid, bool) or pid <= 0: 

544 return None 

545 try: 

546 os.kill(pid, 0) 

547 except ProcessLookupError: 

548 return False 

549 except PermissionError: 

550 # Alive, and owned by somebody else. 

551 return True 

552 except OSError: 

553 return None 

554 return True 

555 

556 

557def list_runs(cache_dir=None, alive_fn=DEFAULT_PROBE) -> list[dict]: 

558 """Every recorded run, newest first, as summary dicts. 

559 

560 Summaries only: the full agent text stays in the state file, so ``--status`` 

561 on a busy machine prints a listing rather than every transcript it ever ran. 

562 

563 ``alive_fn`` defaults to :data:`DEFAULT_PROBE`, which resolves to 

564 :func:`pid_alive` at call time (see :func:`wait_for_run` for why that is not 

565 a default argument). Passing ``None`` means "do not probe", the same as it 

566 does everywhere else. 

567 """ 

568 probe = _resolve_probe(alive_fn) 

569 directory = runs_dir(cache_dir) 

570 try: 

571 names = sorted(p.stem for p in directory.glob("*.json")) 

572 except OSError: 

573 return [] 

574 runs = [] 

575 for run_id in names: 

576 state = read_state(run_id, cache_dir) 

577 if state is None: 

578 continue 

579 runs.append(run_summary(state, alive_fn=probe)) 

580 runs.sort(key=lambda r: (r.get("started_at") or 0, r.get("run_id") or ""), reverse=True) 

581 return runs 

582 

583 

584def run_summary(state: dict, alive_fn=None) -> dict: 

585 """The listing projection of one state document (pure given *alive_fn*). 

586 

587 A run still marked ``running`` whose child is definitively gone is reported 

588 as ``lost``: the child is killed, or crashed hard enough to skip even its 

589 own ``finally``, and reporting it as running forever is the one answer that 

590 is certainly wrong. Only a definitive ``False`` from *alive_fn* demotes it — 

591 unknown (no pid recorded, or a platform we cannot probe) stays ``running``. 

592 """ 

593 status = state.get("status") 

594 if status == STATUS_RUNNING and alive_fn is not None and alive_fn(state.get("pid")) is False: 

595 status = STATUS_LOST 

596 return { 

597 "run_id": state.get("run_id"), 

598 "status": status, 

599 "agent": state.get("agent"), 

600 "role": state.get("role"), 

601 "ok": state.get("ok"), 

602 "error_code": state.get("error_code"), 

603 "started_at": state.get("started_at"), 

604 "duration_s": state.get("duration_s"), 

605 "pid": state.get("pid"), 

606 } 

607 

608 

609def initial_state( 

610 run_id: str, spec: AgentSpec, role: str, now: float, pid: int | None = None 

611) -> dict: 

612 """The state file written before a detached child is spawned (pure). 

613 

614 ``pid`` and ``timeout_s`` are recorded so a later ``--status`` can tell a 

615 live run from an abandoned one, and so ``--wait`` can derive a deadline 

616 from what the run was actually given rather than a guess. 

617 """ 

618 return { 

619 "schema_version": SCHEMA_VERSION, 

620 "run_id": run_id, 

621 "status": STATUS_RUNNING, 

622 "agent": spec.name, 

623 "vendor": spec.vendor, 

624 "model": spec.model or None, 

625 "role": role, 

626 "started_at": now, 

627 "pid": pid, 

628 "timeout_s": spec.timeout, 

629 } 

630 

631 

632def liveness_unknown_reason(state, alive_fn=DEFAULT_PROBE) -> str | None: 

633 """Why a running run's liveness cannot be determined, or None when it can. 

634 

635 Two genuinely different situations reach the same ``pid_alive`` answer of 

636 ``None``, and telling a POSIX operator that "this platform cannot probe" 

637 when the truth is "the child has not claimed the run yet" is simply false: 

638 

639 * **No process id recorded.** The ordinary window between ``--detach`` 

640 writing the state file and the child claiming it — and also where a child 

641 died before it could claim. Happens on every platform. 

642 * **The platform cannot probe one.** Windows, where ``os.kill(pid, 0)`` 

643 terminates rather than probes. 

644 

645 A finished run has nothing to determine, so it returns None. 

646 """ 

647 state = state or {} 

648 if state.get("status") != STATUS_RUNNING: 

649 return None 

650 pid = state.get("pid") 

651 if not isinstance(pid, int) or isinstance(pid, bool) or pid <= 0: 

652 return "it has not recorded a process id yet" 

653 probe = _resolve_probe(alive_fn) 

654 if probe is not None and probe(pid) is None: 

655 return "this platform cannot check a process id without terminating it" 

656 return None 

657 

658 

659#: Head-room added to a run's own timeout to get a default ``--wait`` deadline: 

660#: the agent has until its timeout, and the child needs a moment after that to 

661#: write its state file. 

662WAIT_GRACE_S = 60 

663 

664#: Deadline used when a run has recorded no timeout of its own (no state file 

665#: yet, or one written by an older version). Bounded rather than infinite, so a 

666#: scripted `--wait` cannot hang a pipeline forever. 

667DEFAULT_WAIT_TIMEOUT_S = 3600 

668 

669 

670def default_wait_timeout(state) -> float: 

671 """The ``--wait`` deadline implied by a run's own timeout (pure).""" 

672 recorded = (state or {}).get("timeout_s") 

673 if isinstance(recorded, int) and not isinstance(recorded, bool) and recorded > 0: 

674 return float(recorded + WAIT_GRACE_S) 

675 return float(DEFAULT_WAIT_TIMEOUT_S) 

676 

677 

678def wait_for_run( 

679 run_id: str, 

680 cache_dir=None, 

681 timeout: float | None = None, 

682 poll_s: float = 0.25, 

683 sleep=time.sleep, 

684 clock=time.monotonic, 

685 alive_fn=DEFAULT_PROBE, 

686) -> tuple[dict | None, bool]: 

687 """Block until a detached run finishes. Returns ``(state, timed_out)``. 

688 

689 ``sleep``, ``clock`` and ``alive_fn`` are injected so the lifecycle is 

690 testable against a fake clock instead of real seconds; ``alive_fn`` defaults 

691 to :data:`DEFAULT_PROBE`, i.e. :func:`pid_alive` looked up when this is 

692 *called*, and ``None`` disables probing. A missing state file 

693 is treated as "not finished yet", not as an error: a caller may legitimately 

694 start waiting before the child has written anything. 

695 

696 A run whose recorded process is **definitively gone** returns immediately 

697 with its status mapped to :data:`STATUS_LOST`, rather than blocking to the 

698 deadline for a result that is never coming — the same judgement ``--status`` 

699 makes, applied here so the two cannot disagree. Before reporting that, the 

700 state is re-read once: the child may have written its terminal document in 

701 the window between our read and the probe, and a real answer always wins 

702 over an inference about a pid. 

703 """ 

704 # Resolved at CALL time, not bound as a default: a default argument would 

705 # capture this module's `pid_alive` when the def executes, and a test that 

706 # patches `runagent.pid_alive` would then patch nothing while appearing to 

707 # work — which is exactly how a Windows-only infinite wait reached CI green 

708 # on every other platform. 

709 probe = _resolve_probe(alive_fn) 

710 deadline = None if timeout is None else clock() + float(timeout) 

711 while True: 

712 state = read_state(run_id, cache_dir) 

713 if state is not None and state.get("status") != STATUS_RUNNING: 

714 return state, False 

715 if state is not None and probe is not None and probe(state.get("pid")) is False: 

716 confirmed = read_state(run_id, cache_dir) 

717 if confirmed is not None and confirmed.get("status") != STATUS_RUNNING: 

718 return confirmed, False 

719 lost = dict(confirmed or state) 

720 lost["status"] = STATUS_LOST 

721 return lost, False 

722 if deadline is not None and clock() >= deadline: 

723 return state, True 

724 sleep(poll_s)