Coverage for src/keel/activity.py: 100%

178 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-10-02 20:26 +0000

1"""Lightweight, additive **command-activity** records for live observability. 

2 

3The resumable ``checkpoint`` is the ship backbone's own artifact (s0–s12) and is 

4deliberately ship-shaped. Most other keel commands (``triage``, ``morning``, 

5``pr-loop`` …) run in the main checkout, never write a checkpoint, and so are 

6invisible to ``keel-visual``'s live board. 

7 

8This module adds a *separate, additive* channel: a per-run JSON record under 

9``.keel/activity/`` that any command's adapter can stamp as it moves through its 

10own flow phases (from :mod:`keel.flows`). It never touches the checkpoint 

11contract. Records are keyed by ``run_id`` (one file each), so two commands in the 

12same repo never clobber one another. 

13 

14Pure-core + thin I/O, mirroring :mod:`keel.checkpoint`: the builders/validators 

15are deterministic (stable ordering, no wall-clock, no randomness); only 

16read/write/remove touch the filesystem. 

17""" 

18 

19from __future__ import annotations 

20 

21import contextlib 

22import json 

23import re 

24import time 

25from collections.abc import Callable, Iterator 

26from pathlib import Path 

27from typing import Any 

28 

29from . import config as cfg 

30from . import flows, lock, workspace 

31 

32ACTIVITY_SCHEMA_VERSION = "keel.activity.v1" 

33RECORD_TYPE_ACTIVITY = "command_activity" 

34DEFAULT_ACTIVITY_DIR = ".keel/activity" 

35STATUSES = ("running", "done", "merged") 

36 

37#: Whether the phase the run reached was actually *passed*. Separate from 

38#: :data:`STATUSES` because "advanced to s8" and "cleared s8" are different facts and 

39#: were previously written identically (#636) — a run whose gates came back red 

40#: recorded as ``phase: s8, status: running``, carrying no failure signal at all, and 

41#: the board painted it as in-progress. ``None`` means the step reached this stamp 

42#: without a verdict to report (planning, a phase with nothing to pass), which is not 

43#: the same as passing. 

44VERDICTS = ("pass", "blocked") 

45 

46#: The field a record keeps its delegates' measured token counts in (#1373): an object 

47#: keyed by the delegate call's id, each value ``{"model", "prompt_tokens", 

48#: "completion_tokens"}``. Keyed rather than summed, for two reasons: one ship run makes 

49#: several delegate calls, often to different vendors, and a sum could only be priced at 

50#: one model; and a key makes recording the same call twice a no-op instead of a double 

51#: count. Absent until a delegate records into the run. 

52USAGE_FIELD = "delegate_usage" 

53 

54#: How long a writer waits for another writer of the same record: attempts x poll. 

55#: Holding the lock takes one read and one atomic write, so two seconds is ample; it is 

56#: bounded because a holder killed mid-write leaves its claim behind. 

57LOCK_ATTEMPTS = 40 

58LOCK_POLL_S = 0.05 

59 

60# A run_id reduces to this slug for its filename; anything else is rejected so a 

61# crafted run_id can never escape the activity directory. 

62_RUN_ID_SLUG = re.compile(r"[^a-z0-9._-]+") 

63 

64 

65class ActivityError(ValueError): 

66 """Raised when an activity record or path is malformed.""" 

67 

68 

69def activity_contract_as_dict() -> dict[str, Any]: 

70 """Return the stable activity-record contract consumed by adapters.""" 

71 return { 

72 "schema_version": ACTIVITY_SCHEMA_VERSION, 

73 "record_type": RECORD_TYPE_ACTIVITY, 

74 "dir": DEFAULT_ACTIVITY_DIR, 

75 "keyed_by": "run_id", 

76 "statuses": list(STATUSES), 

77 "additive": True, 

78 "touches_checkpoint": False, 

79 "phase_source": "keel.flows.flow_for(command)", 

80 "usage_field": USAGE_FIELD, 

81 } 

82 

83 

84def configured_activity_dir(config: cfg.ProjectConfig) -> tuple[str, str]: 

85 """Return the configured activity directory and its source.""" 

86 pack = config.policy_pack or {} 

87 reports = pack.get("reports") if isinstance(pack.get("reports"), dict) else {} 

88 value = reports.get("activity") 

89 if isinstance(value, str) and value.strip(): 

90 return value, "policy_pack.reports.activity" 

91 return DEFAULT_ACTIVITY_DIR, "default" 

92 

93 

94def resolve_dir(root: str | Path, config: cfg.ProjectConfig) -> Path: 

95 """Resolve the activity directory under ``root`` and reject escapes.""" 

96 raw, _ = configured_activity_dir(config) 

97 path = Path(raw) 

98 if workspace.is_root_anchored(raw): 

99 raise ActivityError("activity dir must be relative to the project root") 

100 root_path = Path(root).resolve() 

101 resolved = (root_path / path).resolve() 

102 try: 

103 resolved.relative_to(root_path) 

104 except ValueError as exc: 

105 raise ActivityError("activity dir escapes the project root") from exc 

106 return resolved 

107 

108 

109def run_id_slug(run_id: str) -> str: 

110 """Reduce a run_id to a safe filename stem (lowercase ``[a-z0-9._-]``).""" 

111 if not isinstance(run_id, str) or not run_id.strip(): 

112 raise ActivityError("run_id must be a non-empty string") 

113 slug = _RUN_ID_SLUG.sub("-", run_id.strip().lower()).strip("-.") 

114 if not slug: 

115 raise ActivityError("run_id has no usable characters") 

116 return slug 

117 

118 

119def record_path(root: str | Path, config: cfg.ProjectConfig, run_id: str) -> Path: 

120 """Path of the activity record for ``run_id`` under ``root``.""" 

121 return resolve_dir(root, config) / f"{run_id_slug(run_id)}.json" 

122 

123 

124def _phase_ids(command: str) -> tuple[str, ...]: 

125 return tuple(phase.id for phase in flows.flow_for(command)) 

126 

127 

128def build_activity_record( 

129 *, 

130 command: str, 

131 run_id: str, 

132 phase: str, 

133 status: str = "running", 

134 verdict: str | None = None, 

135 issue: int | None = None, 

136 pr: int | None = None, 

137 note: str | None = None, 

138) -> dict[str, Any]: 

139 """Build one deterministic activity record, validating command + phase. 

140 

141 ``command`` must be a known :mod:`keel.flows` command and ``phase`` one of 

142 that command's flow phase ids. ``status`` is ``running``, ``done`` or 

143 ``merged`` (a real merge landed, distinct from a soft ``done``). 

144 

145 ``verdict`` (:data:`VERDICTS`) says whether the phase was **passed**, which 

146 ``status`` deliberately does not: a blocked gate run is still ``running`` in the 

147 board's sense — it advanced, it did not finish — and recording only that made a 

148 red gate indistinguishable from an in-progress one (#636). ``None`` = no verdict 

149 to report, which must not read as a pass. 

150 """ 

151 if not flows.is_known(command): 

152 raise ActivityError(f"unknown command: {command!r}") 

153 if phase not in _phase_ids(command): 

154 raise ActivityError(f"phase {phase!r} is not a {command} flow phase") 

155 if status not in STATUSES: 

156 raise ActivityError(f"unsupported status: {status!r}") 

157 if verdict is not None and verdict not in VERDICTS: 

158 raise ActivityError(f"unsupported verdict: {verdict!r}") 

159 return { 

160 "schema_version": ACTIVITY_SCHEMA_VERSION, 

161 "record_type": RECORD_TYPE_ACTIVITY, 

162 "command": command, 

163 "run_id": run_id, 

164 "phase": phase, 

165 "status": status, 

166 "verdict": verdict, 

167 "issue": issue, 

168 "pr": pr, 

169 "note": note, 

170 } 

171 

172 

173def validate_activity(record: Any) -> None: 

174 """Validate the stable activity-record shape.""" 

175 if not isinstance(record, dict): 

176 raise ActivityError("activity must be an object") 

177 if record.get("schema_version") != ACTIVITY_SCHEMA_VERSION: 

178 raise ActivityError("unsupported schema_version") 

179 if record.get("record_type") != RECORD_TYPE_ACTIVITY: 

180 raise ActivityError("unsupported record_type") 

181 command = record.get("command") 

182 if not isinstance(command, str) or not flows.is_known(command): 

183 raise ActivityError("unsupported command") 

184 if not isinstance(record.get("run_id"), str) or not record["run_id"].strip(): 

185 raise ActivityError("run_id must be a non-empty string") 

186 if record.get("phase") not in _phase_ids(command): 

187 raise ActivityError("unsupported phase") 

188 if record.get("status") not in STATUSES: 

189 raise ActivityError("unsupported status") 

190 # Absent is fine (older records, phases with nothing to pass); a *wrong* value is 

191 # not — a board that trusts this field must never read a typo as a pass. 

192 if record.get("verdict") is not None and record.get("verdict") not in VERDICTS: 

193 raise ActivityError("unsupported verdict") 

194 usage = record.get(USAGE_FIELD) 

195 if usage is not None: 

196 # The cost report prices whatever is here as measured, so a malformed entry is 

197 # refused at the door rather than read as a count. 

198 if not isinstance(usage, dict): 

199 raise ActivityError(f"{USAGE_FIELD} must be an object") 

200 for call_id, entry in usage.items(): 

201 issue = usage_entry_issue(entry) 

202 if issue is not None: 

203 raise ActivityError(f"{USAGE_FIELD}[{call_id!r}]: {issue}") 

204 

205 

206def _is_count(value: Any) -> bool: 

207 return isinstance(value, int) and not isinstance(value, bool) and value >= 0 

208 

209 

210def usage_entry_issue(entry: Any) -> str | None: 

211 """Why ``entry`` is not a well-formed delegate-usage entry, or ``None`` when it is.""" 

212 if not isinstance(entry, dict): 

213 return "entry must be an object" 

214 model = entry.get("model") 

215 if not isinstance(model, str) or not model.strip(): 

216 return "model must be a non-empty string" 

217 for name in ("prompt_tokens", "completion_tokens"): 

218 if not _is_count(entry.get(name)): 

219 return f"{name} must be a non-negative integer" 

220 return None 

221 

222 

223def with_delegate_usage( 

224 record: dict[str, Any], 

225 call_id: str, 

226 *, 

227 model: str, 

228 prompt_tokens: int, 

229 completion_tokens: int, 

230) -> dict[str, Any]: 

231 """A copy of ``record`` that also carries one delegate call's token counts. 

232 

233 Keyed by ``call_id``, so recording the same call again replaces its entry rather than 

234 adding it twice. The entry is validated here, not only on write, so a bad count fails 

235 before anything is built on it. 

236 """ 

237 entry = { 

238 "model": model, 

239 "prompt_tokens": prompt_tokens, 

240 "completion_tokens": completion_tokens, 

241 } 

242 issue = usage_entry_issue(entry) 

243 if issue is not None: 

244 raise ActivityError(f"{USAGE_FIELD}: {issue}") 

245 if not isinstance(call_id, str) or not call_id.strip(): 

246 raise ActivityError(f"{USAGE_FIELD}: call id must be a non-empty string") 

247 updated = dict(record) 

248 updated[USAGE_FIELD] = {**(record.get(USAGE_FIELD) or {}), call_id: entry} 

249 return updated 

250 

251 

252def carry_usage(record: dict[str, Any], existing: dict[str, Any] | None) -> dict[str, Any]: 

253 """``record`` with the delegate counts ``existing`` already held. 

254 

255 Every phase stamp rebuilds the record from :func:`build_activity_record`, which knows 

256 nothing of counts; without this, the next stamp after a delegate recorded would erase 

257 what it recorded (#1373). 

258 """ 

259 usage = (existing or {}).get(USAGE_FIELD) 

260 if not usage: 

261 return record 

262 return {**record, USAGE_FIELD: dict(usage)} 

263 

264 

265def encode_activity(record: dict[str, Any]) -> str: 

266 """Encode one activity record as stable JSON.""" 

267 validate_activity(record) 

268 return json.dumps(record, indent=2, sort_keys=True) + "\n" 

269 

270 

271def parse_activity(text: str) -> dict[str, Any]: 

272 """Parse and validate one activity record.""" 

273 try: 

274 record = json.loads(text) 

275 except json.JSONDecodeError as exc: 

276 raise ActivityError("invalid JSON") from exc 

277 validate_activity(record) 

278 return record 

279 

280 

281def read_activity(path: str | Path) -> dict[str, Any] | None: 

282 """Read one activity record; a missing file means no such run.""" 

283 activity_path = Path(path) 

284 if not activity_path.exists(): 

285 return None 

286 return parse_activity(activity_path.read_text(encoding="utf-8")) 

287 

288 

289def write_activity(path: str | Path, record: dict[str, Any]) -> None: 

290 """Write one validated activity record atomically.""" 

291 workspace.write_text_atomic(path, encode_activity(record)) 

292 

293 

294@contextlib.contextmanager 

295def record_lock( 

296 lock_root: str | Path, 

297 path: str | Path, 

298 *, 

299 owner: str, 

300 attempts: int | None = None, 

301 _sleep: Callable[[float], None] | None = None, 

302) -> Iterator[bool]: 

303 """Serialise the read-modify-write of one activity record; yields whether it is held. 

304 

305 Detached reviewers finish independently, and two recording at once would each read 

306 the record without the other's counts and the second write would drop the first's. 

307 The claim is :mod:`keel.lock`'s atomic ``mkdir``. Bounded: after ``attempts`` tries — 

308 or on an I/O error claiming it — this yields ``False`` and the caller decides whether 

309 to go ahead unlocked (a phase stamp must land) or skip (a count may not be guessed). 

310 """ 

311 # Resolved at call time, not bound as defaults, so the bound and the sleep can be 

312 # changed where they are defined rather than at every caller. 

313 attempts = LOCK_ATTEMPTS if attempts is None else attempts 

314 sleep = time.sleep if _sleep is None else _sleep 

315 resource = f"activity-{Path(path).stem}" 

316 held = False 

317 for attempt in range(attempts): 

318 try: 

319 held = lock.claim_resource(lock_root, resource, owner=owner).granted 

320 except OSError: 

321 break 

322 if held or attempt + 1 == attempts: 

323 break 

324 sleep(LOCK_POLL_S) 

325 try: 

326 yield held 

327 finally: 

328 if held: 

329 with contextlib.suppress(OSError): 

330 lock.release_resource(lock_root, resource, owner=owner, best_effort=True) 

331 

332 

333def record_delegate_usage( 

334 path: str | Path, 

335 lock_root: str | Path, 

336 *, 

337 call_id: str, 

338 model: str, 

339 prompt_tokens: int, 

340 completion_tokens: int, 

341 attempts: int | None = None, 

342 _sleep: Callable[[float], None] | None = None, 

343) -> str: 

344 """Add one delegate call's counts to the run's existing record, under the record lock. 

345 

346 Returns ``recorded``; ``no-record`` when the run has no record to add to (a count is 

347 never the thing that creates one — it has no phase); or ``busy`` when the lock could 

348 not be taken, in which case nothing is written rather than risking a lost update. 

349 Raises :class:`ActivityError` / ``OSError`` for a malformed record or a failed write. 

350 """ 

351 with record_lock(lock_root, path, owner=call_id, attempts=attempts, _sleep=_sleep) as held: 

352 if not held: 

353 return "busy" 

354 existing = read_activity(path) 

355 if existing is None: 

356 return "no-record" 

357 write_activity( 

358 path, 

359 with_delegate_usage( 

360 existing, 

361 call_id, 

362 model=model, 

363 prompt_tokens=prompt_tokens, 

364 completion_tokens=completion_tokens, 

365 ), 

366 ) 

367 return "recorded" 

368 

369 

370def remove_activity(path: str | Path) -> bool: 

371 """Delete an activity record. Returns ``True`` if a file was removed.""" 

372 activity_path = Path(path) 

373 if not activity_path.exists(): 

374 return False 

375 activity_path.unlink() 

376 return True 

377 

378 

379def read_all_activity(dir_path: str | Path) -> list[dict[str, Any]]: 

380 """Every readable activity record in ``dir_path``, sorted by run_id. 

381 

382 Fail-soft: an unreadable or malformed file is skipped, never raised — one bad 

383 record must not blank the board. The directory missing yields ``[]``. 

384 """ 

385 directory = Path(dir_path) 

386 if not directory.is_dir(): 

387 return [] 

388 records: list[dict[str, Any]] = [] 

389 for entry in sorted(directory.glob("*.json")): 

390 try: 

391 record = parse_activity(entry.read_text(encoding="utf-8")) 

392 except (ActivityError, OSError): 

393 continue 

394 records.append(record) 

395 return records