Coverage for src/keel/delegaterun.py: 100%
308 statements
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-02 20:26 +0000
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-02 20:26 +0000
1"""Thin I/O: execute one planned delegate run, foreground or detached (#1012).
3The pure half — which argv, which framing, which effort spelling, which attribution — is
4:mod:`keel.delegate`. This module is the only place that *runs* a plan: one subprocess for
5the ``cli``/``profile`` transports, one loopback HTTP call for ``ollama``, one
6:func:`keel.api_delegate.generate` call for the hosted and OpenAI-compatible ones. Every
7edge is injectable (``_run``, ``_opener``, ``_env``, ``_now``, ``_read``) so the whole
8surface is unit-testable offline and deterministic.
10**Fail-soft, always.** A missing binary, a nonzero exit, a quota refusal, a timeout, an
11unparseable answer — each becomes ``ok: false`` with a machine-readable ``error_code`` in
12the same JSON document a success produces. Nothing here raises at an operator. The
13policy that reads those codes — do not retry a ``rate-limit``, refuse a non-tool provider
14on tier 3, fall back to the host agent — stays with the caller (ship s4/s7), because it is
15policy and this is a transport.
17**The detach primitive.** A delegated implementation runs for tens of minutes; a host
18LLM's turn does not. ``--detach`` spawns the same run as a background child in its own
19session, and the state file under ``.keel/state/delegate/<run-id>.json`` is authoritative:
20the parent may exit, the session may end, the operator may start a new one, and
21``keel delegate wait <run-id>`` still returns the result. That is the primitive an
22orchestrating agent uses instead of a sleep loop — the loop it would otherwise write burns
23its own context window on polling and cannot survive the turn ending, which is how a live
24run ended up with three reviewers finished and no verdict on the PR.
26A run id is validated before it becomes a path. It reaches this module from a CLI flag and
27names a file under the state directory; ``..`` or a separator there is refused rather than
28normalized, so ``wait`` on an unknown or hostile id fails closed instead of reading one.
29"""
31from __future__ import annotations
33import contextlib
34import datetime
35import json
36import os
37import re
38import subprocess # nosec B404
39import time
40import urllib.error
41import urllib.request
42import uuid
43from collections.abc import Callable
44from pathlib import Path
45from typing import Any
47from . import api_delegate, delegate, runner, workspace
48from .delegate import RunPlan
50SCHEMA_VERSION = "keel.delegate-run.v1"
52#: Where detached run state lives, relative to the project root. Inside the existing
53#: gitignored ``.keel/state/`` tree, so a run record is never committed.
54STATE_RELDIR = (".keel", "state", "delegate")
56#: Characters a run id may contain. Deliberately tight: the id becomes a file name.
57_RUN_ID_OK = frozenset("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789._-")
59#: Cap on a single HTTP response body, matching :mod:`keel.api_delegate`.
60_MAX_RESPONSE_BYTES = 50 * 1024 * 1024
62#: Patterns that mark a quota refusal in a CLI's own output. A hosted API answers 429 and
63#: :mod:`keel.api_delegate` classifies it; a CLI exits nonzero with prose, and the caller's
64#: no-retry-on-quota rule needs the same ``rate-limit`` code either way.
65#:
66#: ``429`` needs a status-shaped context, because three digits are not a status code
67#: (#1133). Measured over the 368 vendor transcripts this repository has on disk: every
68#: single occurrence of a bare ``429`` was something else — ``2026-09-05 16:33:40.429`` in
69#: an xcodebuild line, ``git diff 4293a56``, ``a4810728429ea3…``, ``publish.yml:429``, and
70#: a ``429`` in the gutter of a code listing. Six false matches, no true one; and the
71#: transcript that triggered the bug was several hundred KB, where a sha or a line number
72#: carrying those digits is a near certainty. The other markers are phrases, which do not
73#: have that problem.
74_RATE_LIMIT_MARKERS = (
75 re.compile(r"rate[ _-]?limit"),
76 re.compile(r"resource_exhausted"),
77 re.compile(r"quota exceeded"),
78 re.compile(r"usage limit"),
79 re.compile(r"too many requests"),
80 # `HTTP 429`, `status: 429`, `429 Too Many Requests`, `error 429` — never a bare 429.
81 re.compile(r"(?:http[/ ]|status[: ]+|code[: ]+|error[: ]+)429\b"),
82 re.compile(r"\b429\s+(?:too many|client error|error\b)"),
83)
85#: Phrases with which a vendor reports **its own** timeout, as distinct from the wall-clock
86#: bound keel imposes (exit 124, :attr:`keel.runner.CommandResult.timed_out`).
87#:
88#: One vendor's exact wording is measured — agy's ``timeout waiting for response``, seen at
89#: 298s under a keel ``--timeout 900`` — and the rest of this list is the same word in the
90#: shapes a CLI conventionally uses. That is deliberately short: a classifier tuned on one
91#: vendor's prose is how the defect this fixes arrived, so what makes matching prose safe
92#: here is not the list but :func:`failure_signal`, which keeps these patterns away from
93#: the delegate's *answer* and shows them only the vendor's own error.
94_VENDOR_TIMEOUT_MARKERS = (
95 re.compile(r"timeout waiting for"),
96 re.compile(r"\btimed[ -]out\b"),
97 re.compile(r"deadline exceeded"),
98 re.compile(r"\btimeout\b.*\bexceeded\b"),
99)
102class RunIdError(ValueError):
103 """A run id that may not become a path."""
106def check_run_id(run_id: str) -> str:
107 """Validate a run id and return it. Raises :class:`RunIdError` on anything unsafe."""
108 if not run_id or not _RUN_ID_OK.issuperset(run_id) or run_id.strip(".") == "":
109 raise RunIdError(
110 f"invalid run id {run_id!r}: use letters, digits, '.', '_' or '-' only "
111 "(a run id becomes a file name under .keel/state/delegate/)"
112 )
113 return run_id
116def new_run_id(
117 *, _clock: Callable[[], float] = time.time, _token: Callable[[], str] | None = None
118) -> str:
119 """A fresh, sortable run id: ``<epoch-seconds>-<8 hex>``.
121 Time-prefixed so ``keel delegate status`` lists runs in the order they started, and
122 random-suffixed so two runs started in the same second cannot collide.
123 """
124 token = _token() if _token is not None else uuid.uuid4().hex[:8]
125 return f"{int(_clock())}-{token}"
128def state_dir(root: str | Path = ".") -> Path:
129 """The detached-run state directory for ``root`` (not created)."""
130 return Path(root).joinpath(*STATE_RELDIR)
133def state_path(root: str | Path, run_id: str) -> Path:
134 """Path of one run's state document."""
135 return state_dir(root) / f"{check_run_id(run_id)}.json"
138def out_path(root: str | Path, run_id: str) -> Path:
139 """Path of one detached run's captured stdout+stderr."""
140 return state_dir(root) / f"{check_run_id(run_id)}.out"
143def pid_path(root: str | Path, run_id: str) -> Path:
144 """Path of one detached run's pid, kept **beside** the record rather than in it.
146 Two writers touch a detached run: the parent, which knows the pid, and the child,
147 which knows the result. If both wrote the same document the parent's write would be a
148 read-check-write over a file the child can replace at any instant, and the child's
149 terminal record is precisely what must never be lost. Separate files means neither
150 writer needs a lock, a guard, or a retry: the record is the child's alone, the pid
151 file is the parent's alone, and a reader that wants both reads both.
152 """
153 return state_dir(root) / f"{check_run_id(run_id)}.pid"
156def crashed_path(root: str | Path, run_id: str) -> Path:
157 """Path of the marker saying a run can no longer finish. A **sidecar**, like the pid.
159 This is the last read-check-write removed. Marking a crash by editing the record
160 meant a reader could load a ``running`` document, decide the run was gone, and write
161 ``crashed`` back over a ``done`` the child had produced in between — the same lost
162 update the pid file exists to avoid, on the same file, from a third direction. With
163 the marker beside the record the invariant becomes simple enough to state in one
164 line: **the record is written by the child alone.** The parent writes the pid, a
165 reaper writes this, and :func:`run_record` composes the three — with the child's own
166 terminal status always winning, because it is the only one that saw the result.
167 """
168 return state_dir(root) / f"{check_run_id(run_id)}.crashed"
171def write_pid(root: str | Path, run_id: str, pid: int) -> None:
172 """Record a detached child's pid. Best-effort: a run without one is bounded by its
173 ``deadline_at`` instead, so failing to write it must not fail the spawn."""
174 with contextlib.suppress(OSError):
175 workspace.write_text_atomic(pid_path(root, run_id), f"{pid}\n")
178def read_pid(root: str | Path, run_id: str) -> int | None:
179 """The recorded pid of a detached run, or ``None`` when there is none to read."""
180 try:
181 return int(pid_path(root, run_id).read_text(encoding="utf-8").strip())
182 except (OSError, ValueError, RunIdError):
183 return None
186def clear_sidecars(root: str | Path, run_id: str) -> None:
187 """Delete a run id's pid and crash markers. Called before a spawn reuses the id.
189 A reused ``--run-id`` is not exotic — an orchestrator naming a run after the issue it
190 is implementing will reuse it on the retry. Left behind, the previous run's pid file
191 pairs the *new* record with a **dead** pid, and `keel delegate status` (which reaps)
192 lands on a run that started milliseconds ago and marks it crashed; a stale crash
193 marker does the same thing without even needing the race. Both are cleared while the
194 new record is still being set up, before anything can observe the pair.
195 """
196 for path in (pid_path(root, run_id), crashed_path(root, run_id)):
197 with contextlib.suppress(OSError):
198 path.unlink(missing_ok=True)
201def read_crash(root: str | Path, run_id: str) -> dict[str, Any] | None:
202 """The crash marker for a run, or ``None`` when there is none."""
203 try:
204 data = json.loads(crashed_path(root, run_id).read_text(encoding="utf-8"))
205 except (OSError, ValueError, RunIdError):
206 return None
207 return data if isinstance(data, dict) else None
210def run_record(root: str | Path, run_id: str) -> dict[str, Any] | None:
211 """A run's record as a **reader** wants it: the child's document plus the sidecars.
213 Composed rather than stored, so no writer has to carry another writer's half. The
214 child's own terminal status wins over a crash marker: a marker only ever says "this
215 looked abandoned", while a ``done`` record says "the delegate answered", and the
216 second is the stronger claim even when the marker was written later.
217 """
218 record = load_state(root, run_id)
219 if record is None:
220 return None
221 record["pid"] = read_pid(root, run_id)
222 crash = read_crash(root, run_id)
223 if crash is not None and record.get("status") == "running":
224 record["status"] = "crashed"
225 record["finished_at"] = crash.get("finished_at")
226 record["result"] = _detached_failure(
227 record, code="lost", message=crash.get("reason", "the run recorded no result")
228 )
229 return record
232def _read_text(path: str) -> str:
233 return Path(path).read_text(encoding="utf-8")
236def _now_iso(clock: Callable[[], datetime.datetime]) -> str:
237 return clock().isoformat()
240def _utc_now() -> datetime.datetime:
241 return datetime.datetime.now(datetime.UTC)
244def result_document(
245 plan: RunPlan,
246 *,
247 ok: bool,
248 text: str = "",
249 exit_code: int | None = None,
250 duration_s: float = 0.0,
251 timed_out: bool = False,
252 error_code: str | None = None,
253 error: str | None = None,
254 usage: dict[str, int] | None = None,
255) -> dict[str, Any]:
256 """The one JSON document ``keel delegate run`` prints, success or failure.
258 ``exit_code`` is ``None`` for the HTTP transports: there is no process, and reporting
259 a synthetic ``1`` would let a caller mistake a refused API key for a crashed CLI.
261 ``usage`` is the vendor's own token count for the call (#1373) —
262 ``{"prompt_tokens": …, "completion_tokens": …}`` — and ``None`` whenever the transport
263 reported none. Only the ``api`` transport fills it today; a CLI's or Ollama's run
264 reports ``None``, which means "not counted", never "free".
265 """
266 return {
267 "schema_version": SCHEMA_VERSION,
268 "ok": ok,
269 "provider": plan.provider,
270 "vendor": plan.vendor,
271 "model": plan.model,
272 "role": plan.role,
273 "transport": plan.transport,
274 "text": text,
275 "exit_code": exit_code,
276 "duration_s": duration_s,
277 "timed_out": timed_out,
278 "error_code": error_code,
279 "error": error,
280 "attribution": dict(plan.attribution),
281 "read_only": plan.read_only,
282 # Reported beside `read_only` so a caller can refuse rather than discover
283 # afterwards that its "reviewer" held the implementer's write flags.
284 "read_only_backed": plan.read_only_backed,
285 "effort_applied": plan.effort_applied,
286 "warnings": list(plan.warnings),
287 "usage": dict(usage) if usage else None,
288 }
291def rate_limited(text: str) -> bool:
292 """Does this text read as a quota refusal? (pure, best-effort)
294 Give it a :func:`failure_signal`, not a transcript. The patterns are prose, and a
295 delegate's *answer* is prose about code — a review that says the word "rate limit",
296 a diff that contains a sha with 429 in it. Reading the whole stream is how a run that
297 timed out was reported as a quota refusal (#1133).
298 """
299 lowered = (text or "").lower()
300 return any(marker.search(lowered) for marker in _RATE_LIMIT_MARKERS)
303def vendor_timed_out(text: str) -> bool:
304 """Does this text read as the vendor reporting its own timeout? (pure, best-effort)
306 Distinct from :attr:`keel.runner.CommandResult.timed_out`, which is keel's wrapper
307 killing the process. Both are timeouts and both set ``timed_out`` on the contract —
308 the recovery is the same and the operator needs to be told the same thing — but only
309 one of them leaves an exit code of 124.
310 """
311 lowered = (text or "").lower()
312 return any(marker.search(lowered) for marker in _VENDOR_TIMEOUT_MARKERS)
315def final_stream_event(stdout: str) -> dict[str, Any] | None:
316 """agy's last ``{"event": "result", ...}`` frame, or ``None``.
318 The authoritative status of a stream-json run: ``status`` and ``error`` say what
319 happened, where ``response`` is only what the model said.
320 """
321 for line in reversed((stdout or "").splitlines()):
322 line = line.strip()
323 if not line.startswith("{"):
324 continue
325 try:
326 frame = json.loads(line)
327 except (ValueError, TypeError):
328 continue
329 if isinstance(frame, dict) and frame.get("event") == "result":
330 result = frame.get("result")
331 return result if isinstance(result, dict) else frame
332 return None
335def failure_signal(
336 *,
337 stderr: str,
338 stdout: str,
339 output: str = "",
340 stream_json: bool = False,
341 limit: int = 400,
342) -> str:
343 """The text a failure *classification* may read — never the delegate's answer.
345 ``CommandResult.output`` is stdout and stderr glued together, and this module already
346 says why that is a diagnostic rather than a parse: every agent CLI writes progress and
347 notices to stderr. Classifying on it is worse than parsing it. The whole transcript of
348 a long run is hundreds of KB of the model's own prose about code — including, in the
349 transcripts on disk here, git shas and file:line references carrying the digits 429,
350 and reviews that discuss keel's own rate-limit handling by name. #1133 is one such run:
351 the vendor said ``timeout waiting for response`` after 298s and keel reported
352 ``rate-limit``, which tells an operator to wait out a quota window that does not exist.
354 So the signal is bounded and it is the vendor's, not the model's:
356 * **stream-json with a final ``result`` frame** — that frame's ``status`` and
357 ``error``, and deliberately not its ``response``. stderr is *not* added: the frame
358 already said why the run failed, and every agent CLI writes progress there. Prepending
359 it unconditionally meant the same sentence this function refuses to read out of the
360 transcript flipped the classification when it appeared on stderr instead.
361 * **stream-json with no frame** — stderr alone. A crash or a stream truncated mid-write
362 leaves nothing authoritative, and the stdout tail there is the model's own prose, so
363 falling back to it would parse the transcript in exactly the case the frame was meant
364 to protect. Both of these were found by the gate review; the first cut kept the tail.
365 * **anything else** — stderr plus the tail of stdout, which is where a CLI that has no
366 structured output prints why it stopped.
368 The frame's fields are **labelled** (``status: …``), not concatenated bare. A numeric
369 ``status`` is the one ``429`` that really is a status code, and stripping it to a bare
370 token put it out of reach of the very pattern that requires status-shaped context —
371 self-defeating for the field this function chose to read. Also from the review.
373 ``output`` is the last resort, used only when a result carries neither separated
374 stream. :func:`keel.runner._result` always fills both, so this is for a caller that
375 built a :class:`~keel.runner.CommandResult` by hand — where the concatenation is the
376 only signal there is, and refusing to read it would classify every such failure as a
377 plain nonzero exit.
378 """
379 if not (stderr or stdout):
380 return _tail(output, limit)
381 if stream_json:
382 event = final_stream_event(stdout)
383 if event is None:
384 return stderr or ""
385 parts = [
386 f"{field}: {event.get(field)}" for field in ("status", "error") if event.get(field)
387 ]
388 # A frame that carries no error text has said nothing about *why*; stderr is then
389 # the only account there is, and discarding it would lose a real quota refusal.
390 if not event.get("error"):
391 parts.append(stderr or "")
392 return "\n".join(part for part in parts if part)
393 return "\n".join(part for part in (stderr or "", _tail(stdout, limit)) if part)
396def execute(
397 plan: RunPlan,
398 *,
399 _run: Callable[..., runner.CommandResult] | None = None,
400 _opener=None,
401 _env=os.environ,
402 _now: Callable[[], float] = time.monotonic,
403 _read: Callable[[str], str] = _read_text,
404) -> dict[str, Any]:
405 """Run one plan and return the JSON contract. Never raises.
407 ``_run`` defaults to :func:`keel.runner.run_argv`, resolved **at call time** rather
408 than bound as the default value: a default argument is evaluated once at import, so a
409 test patching ``keel.runner.run_argv`` would still reach the real subprocess — which
410 for this module means actually launching an agent CLI.
411 """
412 # Rebound onto the parameter rather than a fresh local: #879's AST sweep reads the
413 # keywords of every call spelled `_run(...)` / `_popen(...)`, and renaming the callee
414 # would take this spawn site out of the rule's sight.
415 _run = runner.run_argv if _run is None else _run
416 started = _now()
418 def finish(**kwargs) -> dict[str, Any]:
419 return result_document(plan, duration_s=round(_now() - started, 3), **kwargs)
421 try:
422 prompt = _read(plan.prompt_path)
423 except OSError as exc:
424 return finish(ok=False, error_code="no-prompt", error=str(exc))
425 if not prompt.strip():
426 return finish(
427 ok=False,
428 error_code="no-prompt",
429 error=f"{plan.prompt_path} is empty; a delegate with no brief produces nothing",
430 )
431 if plan.transport == "ollama":
432 return _run_ollama(plan, prompt, finish, opener=_opener)
433 if plan.transport == "api":
434 return _run_api(plan, prompt, finish, env=_env, opener=_opener)
435 return _run_argv(plan, prompt, finish, run=_run)
438def _run_argv(plan: RunPlan, prompt: str, finish, *, run) -> dict[str, Any]:
439 """The ``cli`` and ``profile`` transports: one subprocess, prompt on stdin.
441 The delegate's answer is **stdout alone**. ``CommandResult.output`` glues stderr onto
442 it, and every agent CLI writes progress, warnings and login notices there — so reading
443 the concatenation lets a run that produced no answer come back ``ok: true`` with the
444 noise as its "output", which downstream is a diff to apply or a review to post. The
445 concatenation stays for the diagnostic tail of a *failure*, where both halves help.
446 """
447 argv = list(plan.argv)
448 if plan.stdin_mode == delegate.STDIN_STREAM_JSON:
449 stdin_text: str | None = delegate.stream_json_frame(prompt)
450 elif plan.stdin_mode == delegate.STDIN_TEXT:
451 stdin_text = prompt
452 else:
453 # A profile's `prompt_mode: arg`. The planner cannot build this argv — it is
454 # handed a path, never the text — so the prompt is appended here, at the one
455 # seam that has both.
456 argv.append(prompt)
457 stdin_text = None
458 result = run(argv, cwd=plan.cwd, timeout=plan.timeout, stdin_text=stdin_text)
459 raw = result.stdout
460 text = delegate.parse_stream_json(raw) if plan.stdin_mode == delegate.STDIN_STREAM_JSON else raw
461 if result.timed_out:
462 return finish(
463 ok=False,
464 exit_code=result.code,
465 timed_out=True,
466 error_code="timeout",
467 error=f"{argv[0]} timed out after {plan.timeout}s",
468 text=text,
469 )
470 if getattr(result, "spawn_failed", False):
471 # The runner's own OSError signal, not exit 127: a command that really ran can
472 # exit 127 too, and reading the code alone made this branch unreachable.
473 return finish(
474 ok=False,
475 exit_code=result.code,
476 error_code="missing-binary",
477 error=f"{argv[0]} could not be executed: {result.stderr or result.output}",
478 )
479 if not result.ok:
480 # Classified on the vendor's own error, never on the delegate's answer (#1133),
481 # and timeout-before-quota: the two codes send an operator in opposite directions
482 # — `rate-limit` says wait out a window, `timeout` says retry now, smaller or
483 # longer — and a vendor that times out often says so in a sentence that also
484 # mentions its limits.
485 signal = failure_signal(
486 stderr=result.stderr,
487 stdout=result.stdout,
488 output=result.output,
489 stream_json=plan.stdin_mode == delegate.STDIN_STREAM_JSON,
490 )
491 vendor_timeout = vendor_timed_out(signal)
492 if vendor_timeout:
493 code = "timeout"
494 elif rate_limited(signal):
495 code = "rate-limit"
496 else:
497 code = "nonzero-exit"
498 return finish(
499 ok=False,
500 exit_code=result.code,
501 # True for either bound: keel's wall-clock wrapper (exit 124) or the vendor's
502 # own. The run did time out; which timer fired is in `exit_code` and `error`.
503 timed_out=vendor_timeout,
504 error_code=code,
505 # The tail still shows both streams — a diagnostic wants everything a reader
506 # might recognise, which is exactly what a decision must not read.
507 error=f"{argv[0]} exited {result.code}: {_tail(result.output)}",
508 text=text,
509 )
510 if not text.strip():
511 return finish(
512 ok=False,
513 exit_code=result.code,
514 error_code="empty-output",
515 error=f"{argv[0]} exited 0 and produced no output",
516 )
517 return finish(ok=True, exit_code=result.code, text=text)
520def _run_api(plan: RunPlan, prompt: str, finish, *, env, opener) -> dict[str, Any]:
521 """Hosted and OpenAI-compatible vendors: exactly one call through ``api_delegate``."""
522 request = plan.request or {}
523 result = api_delegate.generate(
524 request["vendor"],
525 request["model"],
526 prompt,
527 endpoint=request.get("endpoint"),
528 api_key_env=request.get("api_key_env"),
529 max_tokens=request.get("max_tokens", api_delegate.DEFAULT_MAX_TOKENS),
530 timeout=request.get("timeout", api_delegate.DEFAULT_TIMEOUT),
531 extra_payload=request.get("extra_payload") or None,
532 _env=env,
533 _opener=opener,
534 )
535 if not result.ok:
536 # A `bad-response` can still carry the vendor's counts: the call was made.
537 return finish(
538 ok=False, error_code=result.error_code, error=result.error, usage=result.usage
539 )
540 return finish(ok=True, text=result.text, usage=result.usage)
543def _run_ollama(plan: RunPlan, prompt: str, finish, *, opener) -> dict[str, Any]:
544 """The built-in local vendor: one POST to the hardcoded loopback ``/api/generate``."""
545 request = plan.request or {}
546 body = json.dumps(delegate.ollama_payload(request["model"], prompt)).encode("utf-8")
547 # The URL is delegate.OLLAMA_GENERATE_URL, a hardcoded constant — never config, never
548 # the registry, never model or prompt content.
549 http_request = urllib.request.Request( # nosec B310
550 delegate.OLLAMA_GENERATE_URL,
551 data=body,
552 headers={"content-type": "application/json"},
553 method="POST",
554 )
555 client = opener if opener is not None else api_delegate.build_http_only_opener()
556 try:
557 with client.open(http_request, timeout=plan.timeout) as response:
558 raw = response.read(_MAX_RESPONSE_BYTES).decode("utf-8", errors="replace")
559 status = getattr(response, "status", 200)
560 except urllib.error.HTTPError as exc:
561 code = "rate-limit" if exc.code == 429 else "http"
562 return finish(ok=False, error_code=code, error=f"HTTP {exc.code}")
563 except Exception as exc:
564 # Deliberately broad, as in providerprobe: urllib raises URLError, http.client
565 # exceptions (not OSError subclasses), socket timeouts and the address guard's
566 # own error. An operator sees "the local server did not answer", not a traceback.
567 return finish(ok=False, error_code="network", error=str(exc))
568 if not 200 <= status < 300:
569 return finish(ok=False, error_code="http", error=f"HTTP {status}")
570 try:
571 data = json.loads(raw)
572 except ValueError:
573 return finish(ok=False, error_code="bad-response", error="response is not valid JSON")
574 text = delegate.parse_ollama_response(data)
575 if not text:
576 return finish(
577 ok=False, error_code="bad-response", error="response carried no completion text"
578 )
579 return finish(ok=True, text=text)
582def _tail(text: str, limit: int = 400) -> str:
583 stripped = (text or "").strip()
584 return stripped[-limit:]
587# --------------------------------------------------------------------------- detach
590def load_state(root: str | Path, run_id: str) -> dict[str, Any] | None:
591 """One run's state document, or ``None`` when it does not exist / cannot be read."""
592 try:
593 path = state_path(root, run_id)
594 except RunIdError:
595 return None
596 try:
597 return json.loads(path.read_text(encoding="utf-8"))
598 except (OSError, ValueError):
599 return None
602def write_state(root: str | Path, record: dict[str, Any]) -> Path:
603 """Persist a run's state atomically, so a reader never sees a torn document."""
604 path = state_path(root, record["run_id"])
605 workspace.write_text_atomic(path, json.dumps(record, indent=2, sort_keys=True))
606 return path
609def list_runs(root: str | Path = ".") -> list[dict[str, Any]]:
610 """Every readable run record under the state directory, oldest id first."""
611 directory = state_dir(root)
612 records = []
613 try:
614 names = sorted(p.stem for p in directory.glob("*.json"))
615 except OSError: # pragma: no cover - glob on a readable parent does not fail
616 return []
617 for name in names:
618 record = run_record(root, name)
619 if record is not None:
620 records.append(record)
621 return records
624#: Extra wall-clock seconds a detached run gets beyond its own ``--timeout`` before
625#: ``wait`` calls it abandoned. The child enforces the timeout itself and then has to
626#: write its result; the grace is the room for that last write.
627DEADLINE_GRACE_S = 60
630def start_detached(
631 argv: list[str],
632 *,
633 root: str | Path,
634 run_id: str,
635 cwd: str | None = None,
636 timeout: int | None = None,
637 provider: str | None = None,
638 role: str | None = None,
639 _popen=None,
640 _clock: Callable[[], datetime.datetime] = _utc_now,
641) -> dict[str, Any]:
642 """Spawn ``argv`` as a background child and record it. Never raises.
644 The record is written **before** the spawn, so a ``wait`` issued immediately
645 afterwards always finds a file, and the parent never writes it again: the pid goes to
646 its own ``<run-id>.pid`` file (:func:`pid_path`). That is the fix for a lost update.
647 A guarded read-check-write from the parent was still a race — the child's terminal
648 record can land between the read and the write, and the parent would put ``running``
649 back over the result the caller is waiting for. Two files, one writer each, and no
650 window at all.
652 ``deadline_at`` is stamped from the run's own ``--timeout`` plus
653 :data:`DEADLINE_GRACE_S`. It is what stops a ``wait`` with no ``--timeout`` of its own
654 from blocking forever when the child dies without recording anything — a
655 ``SIGKILL``, an OOM kill, a reboot.
657 The child gets its own session where the platform has one, so killing the caller's
658 process group does not take the delegate with it — surviving the caller is the whole
659 point of ``--detach``.
661 ``_popen`` defaults to :class:`subprocess.Popen` resolved at call time, for the same
662 reason :func:`execute`'s ``_run`` does: a default bound at import cannot be patched,
663 and here the cost of missing the seam is a real agent CLI launched by a unit test.
664 """
665 _popen = subprocess.Popen if _popen is None else _popen
666 check_run_id(run_id)
667 started = _clock()
668 record: dict[str, Any] = {
669 "schema_version": SCHEMA_VERSION,
670 "run_id": run_id,
671 "provider": provider,
672 "role": role,
673 "started_at": started.isoformat(),
674 "timeout": timeout,
675 "deadline_at": (
676 None
677 if timeout is None
678 else (started + datetime.timedelta(seconds=timeout + DEADLINE_GRACE_S)).isoformat()
679 ),
680 "status": "running",
681 "argv": list(argv),
682 "out_path": str(out_path(root, run_id)),
683 "result": None,
684 }
685 # A conditional expression rather than an `if`: the assignment then executes on
686 # every platform, so the 100 % line gate holds on the Windows runner too, where
687 # `os.setsid` does not exist and sessions are not a thing.
688 kwargs: dict[str, Any] = {"start_new_session": True} if hasattr(os, "setsid") else {}
689 try:
690 # Creating the directory and opening the log are I/O like any other: an
691 # unwritable root (a read-only checkout, a root owned by another user) must come
692 # back as the same fail-soft contract every other failure does, not as a
693 # PermissionError traceback out of `keel delegate run --detach`.
694 state_dir(root).mkdir(parents=True, exist_ok=True)
695 write_state(root, record)
696 # Same reason the log is opened "w": a reused run id must inherit nothing from
697 # the run before it. See `clear_sidecars`.
698 clear_sidecars(root, run_id)
699 handle = out_path(root, run_id).open("w", encoding="utf-8")
700 except OSError as exc:
701 return _failed_spawn_record(root, record, exc)
702 try:
703 # argv is keel's own re-invocation (sys.executable -m keel), never a shell.
704 child = _popen(
705 list(argv),
706 cwd=cwd,
707 stdin=subprocess.DEVNULL,
708 stdout=handle,
709 stderr=subprocess.STDOUT,
710 **kwargs,
711 )
712 except OSError as exc:
713 return _failed_spawn_record(root, record, exc)
714 finally:
715 handle.close()
716 write_pid(root, run_id, child.pid)
717 record["pid"] = child.pid
718 return record
721def _failed_spawn_record(root: str | Path, record: dict[str, Any], exc: OSError) -> dict[str, Any]:
722 """Mark a run that never started, persisting the record when the disk allows it.
724 Best-effort on purpose: the failure being reported may itself be "this directory
725 cannot be written", and a second write would raise the very traceback the caller is
726 being spared. The returned document is the contract either way.
727 """
728 record["status"] = "done"
729 record["pid"] = None
730 record["result"] = _spawn_failure(record, exc)
731 with contextlib.suppress(OSError):
732 write_state(root, record)
733 return record
736def _spawn_failure(record: dict[str, Any], exc: OSError) -> dict[str, Any]:
737 return _detached_failure(record, code="spawn-failed", message=str(exc))
740def _detached_failure(record: dict[str, Any], *, code: str, message: str) -> dict[str, Any]:
741 """A full return contract for a detached run that never produced one itself.
743 Same keys as :func:`result_document`, so ``keel delegate wait`` hands back one shape
744 whether the delegate answered, failed, or vanished.
745 """
746 return {
747 "schema_version": SCHEMA_VERSION,
748 "ok": False,
749 "provider": record.get("provider"),
750 "vendor": None,
751 "model": None,
752 "role": record.get("role"),
753 "transport": None,
754 "text": "",
755 "exit_code": None,
756 "duration_s": 0.0,
757 "timed_out": code == "lost",
758 "error_code": code,
759 "error": message,
760 "attribution": {},
761 "read_only": None,
762 "read_only_backed": False,
763 "effort_applied": False,
764 "warnings": [],
765 "usage": None,
766 "run_id": record.get("run_id"),
767 "out_path": record.get("out_path"),
768 }
771def finish_detached(
772 root: str | Path,
773 run_id: str,
774 result: dict[str, Any],
775 *,
776 _clock: Callable[[], datetime.datetime] = _utc_now,
777) -> Path:
778 """Record a detached run's result, keeping whatever the parent already wrote."""
779 record = load_state(root, run_id) or {
780 "schema_version": SCHEMA_VERSION,
781 "run_id": run_id,
782 "pid": None,
783 "started_at": _now_iso(_clock),
784 "argv": [],
785 "out_path": str(out_path(root, run_id)),
786 }
787 record["status"] = "done"
788 record["finished_at"] = _now_iso(_clock)
789 record["result"] = result
790 return write_state(root, record)
793#: How often :func:`wait` re-reads the state file.
794DEFAULT_POLL_S = 0.5
797def process_is_alive(pid: int) -> bool:
798 """Is ``pid`` still a process we can see? Best-effort, never raises.
800 ``os.kill(pid, 0)`` sends no signal; it asks the kernel whether the pid exists and
801 whether we may signal it. ``PermissionError`` therefore means **alive** (it exists and
802 belongs to someone else), which is the opposite of what a bare ``except`` would say.
804 Two honest limits, both bounded by the run's own ``deadline_at`` rather than papered
805 over: a child not yet reaped by its parent is a zombie and still answers "alive", and
806 a recycled pid can answer "alive" for an unrelated process. Neither can make ``wait``
807 return a wrong *result* — only make it wait longer before giving up.
808 """
809 if not hasattr(os, "kill"): # pragma: no cover - POSIX everywhere keel runs
810 return True
811 try:
812 os.kill(pid, 0)
813 except ProcessLookupError:
814 return False
815 except (PermissionError, OSError):
816 return True
817 return True
820def _abandoned(
821 record: dict[str, Any],
822 *,
823 alive: Callable[[int], bool],
824 clock: Callable[[], datetime.datetime],
825) -> str | None:
826 """Why this still-``running`` record can never finish, or ``None`` while it might."""
827 pid = record.get("pid")
828 if isinstance(pid, int) and not alive(pid):
829 return (
830 f"the delegate process (pid {pid}) is gone and recorded no result; "
831 f"its output is in {record.get('out_path')}"
832 )
833 deadline = record.get("deadline_at")
834 if isinstance(deadline, str) and deadline:
835 try:
836 when = datetime.datetime.fromisoformat(deadline)
837 except ValueError:
838 return None
839 if clock() >= when:
840 return (
841 f"the run passed its own deadline ({deadline}) without recording a "
842 f"result; its output is in {record.get('out_path')}"
843 )
844 return None
847def _mark_crashed(
848 root: str | Path,
849 run_id: str,
850 reason: str,
851 *,
852 clock: Callable[[], datetime.datetime],
853) -> dict[str, Any] | None:
854 """Note that a run vanished, in its own file, and return the composed record.
856 No read, no check, no write-back: the marker goes to :func:`crashed_path` and
857 :func:`run_record` decides what it means. So a child that wrote ``done`` while the
858 liveness probe was running still reports ``done`` — the marker is simply ignored —
859 and nothing can be lost, because nothing is overwritten.
860 """
861 with contextlib.suppress(OSError):
862 workspace.write_text_atomic(
863 crashed_path(root, run_id),
864 json.dumps({"reason": reason, "finished_at": _now_iso(clock)}, indent=2),
865 )
866 return run_record(root, run_id)
869def reap_abandoned(
870 root: str | Path = ".",
871 *,
872 _alive: Callable[[int], bool] | None = None,
873 _clock: Callable[[], datetime.datetime] = _utc_now,
874) -> list[str]:
875 """Mark every run that can no longer finish, and return their ids.
877 The same liveness and deadline test :func:`wait` applies, run across the whole state
878 directory. Without it a run only stops claiming to be ``running`` when somebody
879 happens to ``wait`` on it — so ``keel delegate status``, the command an operator uses
880 precisely because they are *not* waiting, would be the one view that never told the
881 truth about a killed child.
882 """
883 alive = process_is_alive if _alive is None else _alive
884 reaped = []
885 for record in list_runs(root):
886 if record.get("status") != "running":
887 continue
888 reason = _abandoned(record, alive=alive, clock=_clock)
889 if reason is not None:
890 _mark_crashed(root, record["run_id"], reason, clock=_clock)
891 reaped.append(record["run_id"])
892 return reaped
895def wait(
896 root: str | Path,
897 run_id: str,
898 *,
899 timeout: float | None = None,
900 poll: float = DEFAULT_POLL_S,
901 _sleep: Callable[[float], None] = time.sleep,
902 _now: Callable[[], float] = time.monotonic,
903 _clock: Callable[[], datetime.datetime] = _utc_now,
904 _alive: Callable[[int], bool] | None = None,
905) -> tuple[dict[str, Any] | None, str | None]:
906 """Block until ``run_id`` finishes. Returns ``(result, error_code)``.
908 The state file is the authority — a *result* is only ever read from it, so it survives
909 the process that produced it, a new session, and a reboot.
911 Three ways this returns without a result, and all three are bounded:
913 * ``unknown-run`` — fails closed and immediately. An orchestrator that mistypes a run
914 id, or asks about a run from another checkout, must not get a wait that can only end
915 in a timeout.
916 * ``lost`` — the child is gone, or the run passed the deadline stamped at spawn, and
917 nothing was recorded. Without this a ``SIGKILL``ed child left ``running`` forever
918 and a ``wait`` with no ``--timeout`` blocked indefinitely. The record is marked
919 ``crashed`` so ``keel delegate status`` stops claiming it is running.
920 * ``timeout`` — the caller's own ``--timeout`` elapsed first. The run may still be
921 alive; nothing is marked.
923 A dead pid is re-checked against the file before being called lost, because the child
924 exits *after* writing its result and the two are observed in the other order often
925 enough to matter.
926 """
927 deadline = None if timeout is None else _now() + timeout
928 alive = process_is_alive if _alive is None else _alive
929 while True:
930 record = run_record(root, run_id)
931 if record is None:
932 return None, "unknown-run"
933 status = record.get("status")
934 if status == "done":
935 return record.get("result"), None
936 if status == "crashed":
937 return record.get("result"), "lost"
938 reason = _abandoned(record, alive=alive, clock=_clock)
939 if reason is not None:
940 record = _mark_crashed(root, run_id, reason, clock=_clock)
941 if record is None:
942 return None, "unknown-run"
943 if record.get("status") == "done":
944 return record.get("result"), None
945 return record.get("result"), "lost"
946 if deadline is not None and _now() >= deadline:
947 return None, "timeout"
948 _sleep(poll)
951def planning_failure(
952 provider: str,
953 role: str,
954 *,
955 code: str,
956 message: str,
957) -> dict[str, Any]:
958 """The JSON contract for a run that never started — a plan that could not be made.
960 Same keys as :func:`result_document` with the unresolved halves null, so a caller
961 parses one shape whether the provider was unknown, the model token unsafe, or the
962 delegate simply exited nonzero.
963 """
964 return {
965 "schema_version": SCHEMA_VERSION,
966 "ok": False,
967 "provider": provider,
968 "vendor": None,
969 "model": None,
970 "role": role,
971 "transport": None,
972 "text": "",
973 "exit_code": None,
974 "duration_s": 0.0,
975 "timed_out": False,
976 "error_code": code,
977 "error": message,
978 "attribution": {},
979 "read_only": None,
980 "read_only_backed": False,
981 "effort_applied": False,
982 "warnings": [],
983 "usage": None,
984 }