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
« 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).
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.
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.
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"""
22from __future__ import annotations
24import contextlib
25import json
26import os
27import re
28import secrets
29import time
30from dataclasses import dataclass, replace
31from pathlib import Path
33from .config import AGY_AGENT, DEFAULT_CONFIG, AgentSpec
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"
40#: Every role an orchestrator can dispatch.
41ROLES: tuple[str, ...] = ("implement", "review", "gate", "chair", "fix")
43#: The roles that may modify a working tree — and only with ``--allow-write``.
44WRITE_ROLES = frozenset({"implement", "fix"})
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)
51# --- role policy -------------------------------------------------------------
54@dataclass(frozen=True)
55class RolePolicy:
56 """What one role is allowed to do (pure value).
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 """
63 role: str
64 write: bool = False
65 refusal: str | None = None
66 warning: str | None = None
69def role_policy(role: str, allow_write: bool = False) -> RolePolicy:
70 """Resolve ``(role, allow_write)`` to a role policy (pure).
72 The one place the read-only/write decision is made, so no adapter has to
73 re-derive it:
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)
110# --- agent resolution --------------------------------------------------------
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)
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._:/-]*$")
133def parse_agent_token(token: str) -> tuple[str, str | None, str | None]:
134 """Split ``name`` / ``name:model`` (pure).
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
160def builtin_spec(name: str) -> AgentSpec | None:
161 """The default :class:`AgentSpec` for a bare built-in vendor token (pure).
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)
186def resolve_agent(config, token: str) -> tuple[AgentSpec | None, str | None]:
187 """Resolve ``--agent`` to a concrete :class:`AgentSpec` (pure).
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.
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
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
219def transport_for(spec: AgentSpec) -> str:
220 """How this agent is reached: ``cli``, ``api`` or ``local`` (pure).
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
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)
232# --- attribution -------------------------------------------------------------
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)
251def strip_transport(model: str) -> str:
252 """Drop a ``<transport>:`` prefix, leaving the vendor's own model id (pure).
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
266def model_base(model: str) -> str:
267 """Strip a model id to a coarse **family + major** label (pure).
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.
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:
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.
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]
306def attribution(vendor: str, model: str | None = None) -> dict:
307 """Who did the work: ``{vendor, model, label}`` (pure).
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)}
320# --- result shape ------------------------------------------------------------
323def result_dict(spec: AgentSpec, role: str, result) -> dict:
324 """Project an :class:`~ai_jury.adapters.AgentResult` onto the v1 export (pure).
326 Every key is always present, with a stable type, so a consumer can read the
327 document without probing for optional fields.
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
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 }
360# --- detached runs -----------------------------------------------------------
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"
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}$")
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"
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)
382class RunIdError(ValueError):
383 """An identifier that must never be turned into a path (issue #661).
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 """
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
400def runs_dir(cache_dir=None) -> Path:
401 """The directory holding detached-run state files."""
402 from .cache import default_cache_dir
404 base = Path(cache_dir) if cache_dir else default_cache_dir()
405 return base / RUNS_DIR_NAME
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).
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.
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
437def state_path(run_id, cache_dir=None) -> Path:
438 return run_path(run_id, ".json", cache_dir)
441def output_path(run_id, cache_dir=None) -> Path:
442 return run_path(run_id, ".out", cache_dir)
445def write_state(run_id, state: dict, cache_dir=None) -> Path:
446 """Write a run's state file atomically, 0600, in a 0700 directory.
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.
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
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
475def read_state(run_id, cache_dir=None) -> dict | None:
476 """Read one run's state, or None when it is missing/unreadable/corrupt.
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
493class _DefaultProbe:
494 """Sentinel: resolve :func:`pid_alive` when the call happens.
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 """
503 def __repr__(self) -> str: # pragma: no cover - debugging aid
504 return "<default probe>"
507#: See :class:`_DefaultProbe`.
508DEFAULT_PROBE = _DefaultProbe()
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
516def pid_alive(pid) -> bool | None:
517 """Is that process still running? ``None`` when we cannot safely tell.
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".
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.
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
557def list_runs(cache_dir=None, alive_fn=DEFAULT_PROBE) -> list[dict]:
558 """Every recorded run, newest first, as summary dicts.
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.
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
584def run_summary(state: dict, alive_fn=None) -> dict:
585 """The listing projection of one state document (pure given *alive_fn*).
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 }
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).
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 }
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.
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:
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.
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
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
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
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)
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)``.
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.
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)