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

1"""Thin I/O: execute one planned delegate run, foreground or detached (#1012). 

2 

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. 

9 

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. 

16 

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. 

25 

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""" 

30 

31from __future__ import annotations 

32 

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 

46 

47from . import api_delegate, delegate, runner, workspace 

48from .delegate import RunPlan 

49 

50SCHEMA_VERSION = "keel.delegate-run.v1" 

51 

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") 

55 

56#: Characters a run id may contain. Deliberately tight: the id becomes a file name. 

57_RUN_ID_OK = frozenset("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789._-") 

58 

59#: Cap on a single HTTP response body, matching :mod:`keel.api_delegate`. 

60_MAX_RESPONSE_BYTES = 50 * 1024 * 1024 

61 

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) 

84 

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) 

100 

101 

102class RunIdError(ValueError): 

103 """A run id that may not become a path.""" 

104 

105 

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 

114 

115 

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>``. 

120 

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}" 

126 

127 

128def state_dir(root: str | Path = ".") -> Path: 

129 """The detached-run state directory for ``root`` (not created).""" 

130 return Path(root).joinpath(*STATE_RELDIR) 

131 

132 

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" 

136 

137 

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" 

141 

142 

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. 

145 

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" 

154 

155 

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. 

158 

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" 

169 

170 

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") 

176 

177 

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 

184 

185 

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. 

188 

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) 

199 

200 

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 

208 

209 

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. 

212 

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 

230 

231 

232def _read_text(path: str) -> str: 

233 return Path(path).read_text(encoding="utf-8") 

234 

235 

236def _now_iso(clock: Callable[[], datetime.datetime]) -> str: 

237 return clock().isoformat() 

238 

239 

240def _utc_now() -> datetime.datetime: 

241 return datetime.datetime.now(datetime.UTC) 

242 

243 

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. 

257 

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. 

260 

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 } 

289 

290 

291def rate_limited(text: str) -> bool: 

292 """Does this text read as a quota refusal? (pure, best-effort) 

293 

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) 

301 

302 

303def vendor_timed_out(text: str) -> bool: 

304 """Does this text read as the vendor reporting its own timeout? (pure, best-effort) 

305 

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) 

313 

314 

315def final_stream_event(stdout: str) -> dict[str, Any] | None: 

316 """agy's last ``{"event": "result", ...}`` frame, or ``None``. 

317 

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 

333 

334 

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. 

344 

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. 

353 

354 So the signal is bounded and it is the vendor's, not the model's: 

355 

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. 

367 

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. 

372 

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) 

394 

395 

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. 

406 

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() 

417 

418 def finish(**kwargs) -> dict[str, Any]: 

419 return result_document(plan, duration_s=round(_now() - started, 3), **kwargs) 

420 

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) 

436 

437 

438def _run_argv(plan: RunPlan, prompt: str, finish, *, run) -> dict[str, Any]: 

439 """The ``cli`` and ``profile`` transports: one subprocess, prompt on stdin. 

440 

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) 

518 

519 

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) 

541 

542 

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) 

580 

581 

582def _tail(text: str, limit: int = 400) -> str: 

583 stripped = (text or "").strip() 

584 return stripped[-limit:] 

585 

586 

587# --------------------------------------------------------------------------- detach 

588 

589 

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 

600 

601 

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 

607 

608 

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 

622 

623 

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 

628 

629 

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. 

643 

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. 

651 

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. 

656 

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``. 

660 

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 

719 

720 

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. 

723 

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 

734 

735 

736def _spawn_failure(record: dict[str, Any], exc: OSError) -> dict[str, Any]: 

737 return _detached_failure(record, code="spawn-failed", message=str(exc)) 

738 

739 

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. 

742 

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 } 

769 

770 

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) 

791 

792 

793#: How often :func:`wait` re-reads the state file. 

794DEFAULT_POLL_S = 0.5 

795 

796 

797def process_is_alive(pid: int) -> bool: 

798 """Is ``pid`` still a process we can see? Best-effort, never raises. 

799 

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. 

803 

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 

818 

819 

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 

845 

846 

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. 

855 

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) 

867 

868 

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. 

876 

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 

893 

894 

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)``. 

907 

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. 

910 

911 Three ways this returns without a result, and all three are bounded: 

912 

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. 

922 

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) 

949 

950 

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. 

959 

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 }