a9f014f9fa
Close the loop so dispatched work actually completes on bl: - on_no_result fallback in followup-sweeper (op + job forms via exec-constrained registry / job-dispatch), seeded on the three autonomy-pulse jobs; fallback_due() dedupes the gravity path - gravity.py: add __main__ entry (loop-remediator.timer was a no-op), 300s re-arm budget, fallback firing + stamp/skip logic - harvester: proof-of-result followups (result_has_evidence), acted-variant NACK, emit-model tool-hint wording - envelope: RESPONSE RULE states the emit model (agents EMIT directives verbatim; runtime executes; works from bare containers) - completion-audit.py + systemd 15-min timer: per-family funnel, swarm drain, followup backlog; digest DM when degraded, 6h heartbeat - tests/test_completion.py (29 tests), JOB-SPEC.md docs Tests: 67/67 focused green (completion + tool_calls).
313 lines
11 KiB
Python
Executable File
313 lines
11 KiB
Python
Executable File
#!/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 (<name>-YYYYMMDD-HHMMSS-<hex8>) 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())
|