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
« 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.
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.
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.
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"""
19from __future__ import annotations
21import contextlib
22import json
23import re
24import time
25from collections.abc import Callable, Iterator
26from pathlib import Path
27from typing import Any
29from . import config as cfg
30from . import flows, lock, workspace
32ACTIVITY_SCHEMA_VERSION = "keel.activity.v1"
33RECORD_TYPE_ACTIVITY = "command_activity"
34DEFAULT_ACTIVITY_DIR = ".keel/activity"
35STATUSES = ("running", "done", "merged")
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")
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"
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
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._-]+")
65class ActivityError(ValueError):
66 """Raised when an activity record or path is malformed."""
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 }
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"
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
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
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"
124def _phase_ids(command: str) -> tuple[str, ...]:
125 return tuple(phase.id for phase in flows.flow_for(command))
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.
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``).
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 }
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}")
206def _is_count(value: Any) -> bool:
207 return isinstance(value, int) and not isinstance(value, bool) and value >= 0
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
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.
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
252def carry_usage(record: dict[str, Any], existing: dict[str, Any] | None) -> dict[str, Any]:
253 """``record`` with the delegate counts ``existing`` already held.
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)}
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"
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
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"))
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))
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.
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)
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.
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"
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
379def read_all_activity(dir_path: str | Path) -> list[dict[str, Any]]:
380 """Every readable activity record in ``dir_path``, sorted by run_id.
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