#!/usr/bin/env python3 """Completion auditor: prove work gets done, or say exactly where it stalls. Runs on a 15-minute systemd timer (systemd/completion-audit.*). Reads job-log.jsonl, swarms.json, and followups.json; computes the completion funnel per job family plus swarm drain and followup backlog; writes a JSON report under logs/ and posts a compact digest to the ops heartbeat sidechat when degraded (or a heartbeat summary every 6h when green). Read-only except the digest DM and its own log/state files. Exit 0 always on a completed audit; tracebacks (real errors) fail the timer visibly. """ import argparse import json import os import re import subprocess import sys from collections import Counter, defaultdict from datetime import datetime, timezone, timedelta from pathlib import Path REPO_ROOT = Path(__file__).resolve().parent.parent JOB_LOG = REPO_ROOT / "job-log.jsonl" SWARMS_FILE = REPO_ROOT / "swarms.json" FOLLOWUPS_FILE = REPO_ROOT / "followups.json" LOGS_DIR = REPO_ROOT / "logs" STATE_FILE = LOGS_DIR / "completion-audit-state.json" _JOB_ID_RE = re.compile(r"^(.+)-(\d{8})-(\d{6})-([0-9a-f]{8})$") HEARTBEAT_INTERVAL_H = 6 STALE_RUNNING_MIN = 90 SILENT_MIN_SENT = 3 def utcnow(): return datetime.now(timezone.utc) def family_of(job_id): """Strip the dispatch suffix (-YYYYMMDD-HHMMSS-) to the family.""" m = _JOB_ID_RE.match(job_id or "") return m.group(1) if m else (job_id or "?") def parse_ts(ts): try: t = datetime.fromisoformat(str(ts)) except Exception: return None if t.tzinfo is None: t = t.replace(tzinfo=timezone.utc) return t def compute_funnel(events, cutoff): """Aggregate job-log events since cutoff. Returns (families, tools) where families maps family -> counters and tools holds global tool_exec stats. Pure over the event list. """ families = defaultdict(lambda: Counter()) tools = Counter() tool_errs = Counter() for e in events: t = parse_ts(e.get("ts")) if t is None or t < cutoff: continue ty = e.get("type") if ty == "job_sent": families[family_of(e.get("job_id"))]["sent"] += 1 elif ty == "job_dispatched": families[family_of(e.get("job_id"))]["dispatched"] += 1 elif ty == "tool_exec": tools["total"] += 1 if e.get("success"): tools["ok"] += 1 else: tools["fail"] += 1 tool_errs[e.get("op", "?")] += 1 elif ty == "job_result": fam = family_of(e.get("job_id")) families[fam]["results"] += 1 families[fam]["ok" if e.get("success") else "fail"] += 1 elif ty == "job_failed": families[family_of(e.get("job_id"))]["failed"] += 1 elif ty == "fallback_executed": families[family_of(e.get("job_id"))]["fallback_ok"] += 1 elif ty == "fallback_failed": families[family_of(e.get("job_id"))]["fallback_fail"] += 1 elif ty == "proof_requested": families[family_of(e.get("job_id"))]["proofs"] += 1 return families, {"tools": tools, "tool_errs": tool_errs} def swarm_drain(now): """Status counts + stale-running slots from swarms.json.""" try: data = json.load(open(SWARMS_FILE)) except Exception: return {"error": "swarms.json unreadable"}, [] values = data.values() if isinstance(data, dict) else data status = Counter() stale = [] for s in values: if not isinstance(s, dict): continue for sl in s.get("slots", []) or []: status[sl.get("status", "?")] += 1 if sl.get("status") == "running": upd = parse_ts(sl.get("updated_ts")) if upd and (now - upd) > timedelta(minutes=STALE_RUNNING_MIN): stale.append({ "swarm": s.get("swarm_id"), "slot": sl.get("slot"), "agent": sl.get("agent_id"), "idle_min": int((now - upd).total_seconds() // 60), }) return {"slots": dict(status)}, stale def followup_backlog(now): """Pending/overdue/escalated counts from followups.json.""" try: data = json.load(open(FOLLOWUPS_FILE)) except Exception: return {"error": "followups.json unreadable"} values = data.values() if isinstance(data, dict) else data out = Counter() for r in values: if not isinstance(r, dict): continue st = r.get("status", "?") out[st] += 1 if st == "pending": dl = parse_ts(r.get("deadline")) if dl and dl < now: out["overdue"] += 1 return dict(out) def build_report(hours): now = utcnow() cutoff = now - timedelta(hours=hours) events = [] try: with open(JOB_LOG) as f: for line in f: line = line.strip() if not line: continue try: events.append(json.loads(line)) except Exception: continue except FileNotFoundError: pass families, tools = compute_funnel(events, cutoff) fam = {k: dict(v) for k, v in sorted(families.items())} swarm, stale = swarm_drain(now) backlog = followup_backlog(now) totals = Counter() for v in fam.values(): for k, n in v.items(): totals[k] += n silent = sorted( k for k, v in fam.items() if v.get("sent", 0) >= SILENT_MIN_SENT and v.get("results", 0) == 0) degraded_reasons = [] if silent: degraded_reasons.append(f"{len(silent)} silent families: {', '.join(silent[:5])}") if totals.get("failed"): degraded_reasons.append(f"{totals['failed']} job_failed") if tools["tools"].get("fail"): top = tools["tool_errs"].most_common(3) degraded_reasons.append( "tool errors: " + ", ".join(f"{op}x{n}" for op, n in top)) if totals.get("fallback_fail"): degraded_reasons.append(f"{totals['fallback_fail']} fallback_failed") if stale: degraded_reasons.append(f"{len(stale)} running slots idle >{STALE_RUNNING_MIN}m") if backlog.get("overdue"): degraded_reasons.append(f"{backlog['overdue']} overdue followups") if backlog.get("escalated"): degraded_reasons.append(f"{backlog['escalated']} escalated followups") return { "ts": now.isoformat(), "window_h": hours, "totals": dict(totals), "tools": {k: dict(v) if isinstance(v, Counter) else v for k, v in tools.items()}, "families": fam, "silent_families": silent, "swarms": swarm, "stale_running": stale[:10], "followups": backlog, "degraded": bool(degraded_reasons), "reasons": degraded_reasons, } def render_digest(rep): t = rep["totals"] tools = rep["tools"].get("tools", {}) lines = [ f"Completion audit ({rep['window_h']}h, {rep['ts'][:16]}Z)", f"funnel: {t.get('sent', 0)} sent / {t.get('dispatched', 0)} dispatched / " f"{tools.get('total', 0)} tool_exec / {t.get('results', 0)} results " f"({t.get('ok', 0)} ok)", ] if rep["silent_families"]: lines.append("silent: " + ", ".join(rep["silent_families"][:6])) bits = [] if t.get("failed"): bits.append(f"{t['failed']} job_failed") if tools.get("fail"): bits.append(f"{tools['fail']} tool errors") if t.get("fallback_ok") or t.get("fallback_fail"): bits.append(f"fallback {t.get('fallback_ok', 0)} ok / {t.get('fallback_fail', 0)} fail") if t.get("proofs"): bits.append(f"{t['proofs']} proof reqs") if bits: lines.append("flags: " + ", ".join(bits)) sw = rep["swarms"].get("slots", {}) if sw: lines.append("swarms now: " + " / ".join(f"{v} {k}" for k, v in sorted(sw.items()))) if rep["stale_running"]: lines.append(f"stale running: {len(rep['stale_running'])} slots (see report)") fb = rep["followups"] if fb and "error" not in fb: lines.append( f"followups now: {fb.get('pending', 0)} pending / {fb.get('overdue', 0)} " f"overdue / {fb.get('escalated', 0)} escalated") if rep["degraded"]: lines.append("verdict: DEGRADED — " + "; ".join(rep["reasons"][:3])) else: lines.append("verdict: HEALTHY — work flowing, results landing") return "\n".join(lines) def should_post(report): """Post on degraded, else heartbeat at most every HEARTBEAT_INTERVAL_H.""" if report["degraded"]: return True, "degraded" try: state = json.load(open(STATE_FILE)) last = parse_ts(state.get("last_heartbeat")) except Exception: last = None if last is None or (utcnow() - last) > timedelta(hours=HEARTBEAT_INTERVAL_H): return True, "heartbeat" return False, "green-quiet" def post_digest(digest): argv = [sys.executable, str(REPO_ROOT / "bin" / "dm.py"), "send", "--agent", "super", "--to", "opm", "--target", "heartbeat", digest] p = subprocess.run(argv, capture_output=True, text=True, timeout=120) return p.returncode == 0, (p.stdout or p.stderr or "").strip()[:300] def save_report(report): LOGS_DIR.mkdir(parents=True, exist_ok=True) stamp = report["ts"].replace("+00:00", "Z").replace(":", "") dated = LOGS_DIR / f"completion-audit-{stamp[:15]}.json" body = json.dumps(report, indent=2) dated.write_text(body, encoding="utf-8") latest = LOGS_DIR / "completion-audit-latest.json" tmp = LOGS_DIR / f".completion-audit-latest.tmp.{os.getpid()}" tmp.write_text(body, encoding="utf-8") os.replace(tmp, latest) return dated def main(): ap = argparse.ArgumentParser(description="Completion funnel auditor") ap.add_argument("--hours", type=int, default=24) ap.add_argument("--post", dest="post", action="store_true", default=True) ap.add_argument("--no-post", dest="post", action="store_false") ap.add_argument("--json", action="store_true", help="Print raw report JSON") args = ap.parse_args() report = build_report(args.hours) path = save_report(report) if args.json: print(json.dumps(report, indent=2)) else: print(render_digest(report)) print(f"\nreport: {path}") if not args.post: print("post: skipped (--no-post)") return 0 do_post, why = should_post(report) if not do_post: print(f"post: skipped ({why})") return 0 ok, detail = post_digest(render_digest(report)) print(f"post: {'delivered' if ok else 'FAILED'} ({why}) {detail[:120]}") if ok and why == "heartbeat": try: state = {} if STATE_FILE.exists(): state = json.loads(STATE_FILE.read_text(encoding="utf-8")) state["last_heartbeat"] = report["ts"] STATE_FILE.write_text(json.dumps(state, indent=2), encoding="utf-8") except Exception as e: print(f"warning: state save failed: {e}") return 0 if __name__ == "__main__": sys.exit(main())