From a9f014f9fa340f7fc2e5368a2e2c2e15d3b4950e Mon Sep 17 00:00:00 2001 From: Muse Sidechat Date: Tue, 6 Oct 2026 08:20:14 +0000 Subject: [PATCH] feat: completion-enforcement loop (fallback, proof, emit-model, auditor) 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). --- bin/completion-audit.py | 312 ++++++++++++++++++++++++++ bin/followup-sweeper.py | 185 ++++++++++++++++ bin/gravity.py | 95 +++++++- bin/prompt_envelope.py | 5 +- bin/response-harvester.py | 97 +++++++- docs/JOB-SPEC.md | 80 +++++++ jobs/autonomy-pulse-646.json | 10 +- jobs/autonomy-pulse-opm.json | 10 +- jobs/autonomy-pulse-pip.json | 10 +- systemd/completion-audit.service | 10 + systemd/completion-audit.timer | 9 + tests/test_completion.py | 366 +++++++++++++++++++++++++++++++ 12 files changed, 1173 insertions(+), 16 deletions(-) create mode 100755 bin/completion-audit.py create mode 100644 systemd/completion-audit.service create mode 100644 systemd/completion-audit.timer create mode 100644 tests/test_completion.py diff --git a/bin/completion-audit.py b/bin/completion-audit.py new file mode 100755 index 0000000..591f1ca --- /dev/null +++ b/bin/completion-audit.py @@ -0,0 +1,312 @@ +#!/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()) diff --git a/bin/followup-sweeper.py b/bin/followup-sweeper.py index e0c8bc0..32270ff 100755 --- a/bin/followup-sweeper.py +++ b/bin/followup-sweeper.py @@ -144,6 +144,164 @@ def _resolve_nudge_thread_uuid(nudge_output): return fallback +_JOB_ID_RE = re.compile(r"^(.+)-(\d{8})-(\d{6})-([0-9a-f]{8})$") +_OP_NAME_RE = re.compile(r"^[a-z][a-z0-9_.]{0,63}$") +_JOB_NAME_RE = re.compile(r"^[a-z0-9-]{1,64}$") +_FALLBACK_RETRY_S = 3600 + + +def fallback_due(rec, now=None): + """True when a terminal followup should (re)attempt its on_no_result fallback. + + Fires once per record; a failed attempt may retry after _FALLBACK_RETRY_S. + Shared by the sweeper terminal path and gravity remediate so the two + firing paths can never double-execute. + """ + fb = rec.get("fallback") or {} + if fb.get("ran"): + return False + ts = fb.get("ts") + if not ts: + return True + try: + last = datetime.fromisoformat(str(ts).replace("Z", "+00:00")) + except Exception: + return True + if last.tzinfo is None: + last = last.replace(tzinfo=timezone.utc) + base = now or utcnow_dt() + return (base - last).total_seconds() >= _FALLBACK_RETRY_S + +_EXEC_OPS_MOD = None + + +def derive_job_name(job_id): + """Extract the job name from a dispatched job_id (-YYYYMMDD-HHMMSS-).""" + m = _JOB_ID_RE.match(job_id or "") + if not m: + return None + name = m.group(1) + if not (JOBS_DIR / f"{name}.json").exists(): + return None + return name + + +def load_job_fallback(job_name): + """Return (spec, error) for a job's on_no_result fallback. + + spec is None when the job declares none. Shape: + {"job": ""} -> dispatch a fallback job, or + {"op": "", "args": {...}} -> run one exec-constrained op. + """ + try: + with open(JOBS_DIR / f"{job_name}.json", "r", encoding="utf-8") as f: + cfg = json.load(f) + except Exception as e: + return None, f"unreadable job {job_name}: {e}" + spec = cfg.get("on_no_result") + if spec is None: + return None, None + if not isinstance(spec, dict) or set(spec) - {"job", "op", "args"}: + return None, "on_no_result must be an object with job|op (+args)" + if bool(spec.get("job")) == bool(spec.get("op")): + return None, "on_no_result needs exactly one of job|op" + if spec.get("job"): + jn = spec["job"] + if not isinstance(jn, str) or not _JOB_NAME_RE.fullmatch(jn): + return None, "on_no_result.job must be a valid job name" + if not (JOBS_DIR / f"{jn}.json").exists(): + return None, f"on_no_result.job {jn!r} does not exist" + else: + if not isinstance(spec.get("op"), str) or not _OP_NAME_RE.fullmatch(spec["op"]): + return None, "on_no_result.op must be a valid op name" + if "args" in spec and not isinstance(spec["args"], dict): + return None, "on_no_result.args must be an object" + return spec, None + + +def _load_exec_ops(): + global _EXEC_OPS_MOD + if _EXEC_OPS_MOD is None: + import importlib.util + mod_spec = importlib.util.spec_from_file_location( + "exec_constrained_sweeper", str(BIN_DIR / "exec-constrained.py")) + mod = importlib.util.module_from_spec(mod_spec) + mod_spec.loader.exec_module(mod) + _EXEC_OPS_MOD = mod + return _EXEC_OPS_MOD + + +def run_no_result_fallback(rec, dry_run=False): + """Execute a job's on_no_result fallback at terminal followup expiry. + + Returns an outcome dict; never raises (failures are outcome data so + one bad spec can't break the sweep). + """ + dm_id = rec.get("dm_id", "?") + outcome = {"dm_id": dm_id, "ran": False, "mode": None, + "configured": False, "detail": "no job fallback"} + try: + job_name = derive_job_name(rec.get("job_id")) + if not job_name: + outcome["detail"] = "no resolvable job_id" + return outcome + spec, err = load_job_fallback(job_name) + if err: + outcome.update(configured=True, detail=err) + return outcome + if spec is None: + return outcome + outcome["configured"] = True + if dry_run: + outcome.update(mode="dry_run", detail=json.dumps(spec)[:200]) + return outcome + if spec.get("job"): + env = os.environ.copy() + env["CHAIN_PREV_JOB_ID"] = rec.get("job_id", "") + env["CHAIN_PREV_RESULT"] = ( + f"TIMEOUT: Agent {rec.get('recipient')} gave no result; " + f"on_no_result fallback for job {job_name}") + cmd = [sys.executable, str(DISPATCH_PY), spec["job"]] + try: + p = subprocess.run(cmd, capture_output=True, text=True, + timeout=180, env=env) + except Exception as e: + outcome.update(mode="job", detail=f"dispatch exception: {e}") + return outcome + ok = p.returncode == 0 + outcome.update(ran=ok, mode="job", + detail=(f"dispatched {spec['job']}" if ok + else f"dispatch failed: {(p.stderr or p.stdout).strip()[:200]}")) + else: + mod = _load_exec_ops() + op = spec["op"] + op_spec = mod.OPS.get(op) + if op_spec is None: + outcome["detail"] = f"unknown op: {op}" + return outcome + args = dict(spec.get("args") or {}) + try: + clean = op_spec["validate"](args) + except Exception as e: + outcome["detail"] = f"op validation failed: {e}" + return outcome + argv = op_spec["build"](clean) + try: + p = subprocess.run(argv, capture_output=True, text=True, + timeout=op_spec.get("timeout", 120)) + except Exception as e: + outcome.update(mode="op", detail=f"op exception: {e}") + return outcome + ok = p.returncode == 0 + out = (p.stdout or p.stderr or "").strip() + outcome.update(ran=ok, mode="op", + detail=(f"{op} ok: {out[:200]}" if ok + else f"{op} failed rc={p.returncode}: {out[:200]}")) + except Exception as e: + outcome["detail"] = f"fallback exception: {e}" + return outcome + + def sweep_cycle(dry_run=False): followups = load_followups() if not followups: @@ -152,6 +310,7 @@ def sweep_cycle(dry_run=False): now = utcnow_dt() nudges_count = 0 escalations_count = 0 + fallbacks_count = 0 modified = False for dm_id, rec in list(followups.items()): @@ -334,6 +493,31 @@ def sweep_cycle(dry_run=False): pipeline_engine.fail_pipeline(run_entry.get("run_id"), "step_timed_out_without_fallback") except Exception: pass + # on_no_result fallback: the agent never replied, so run + # the job's declared server-side effect now (if any). + # fallback_due() dedupes against the gravity firing path. + fb = (run_no_result_fallback(rec) if fallback_due(rec) + else {"configured": False, "ran": False, "mode": None, + "detail": "fallback already ran"}) + if fb["configured"]: + rec["fallback"] = {"ran": fb["ran"], "mode": fb["mode"], + "detail": fb["detail"][:200], + "ts": utcnow_str()} + modified = True + append_job_log({ + "ts": utcnow_str(), + "type": ("fallback_executed" if fb["ran"] + else "fallback_failed"), + "dm_id": dm_id, + "recipient": recipient, + "job_id": rec.get("job_id"), + "mode": fb["mode"], + "detail": fb["detail"][:300], + }) + print(f"Sweeper: on_no_result fallback for {dm_id}: " + f"ran={fb['ran']} {fb['detail'][:120]}") + if fb["ran"]: + fallbacks_count += 1 else: escalations_count += 1 @@ -346,6 +530,7 @@ def sweep_cycle(dry_run=False): "pending": pending_count, "nudges_sent": nudges_count, "escalations": escalations_count, + "fallbacks": fallbacks_count, } diff --git a/bin/gravity.py b/bin/gravity.py index 397fa68..e399010 100644 --- a/bin/gravity.py +++ b/bin/gravity.py @@ -657,6 +657,62 @@ def diagnose_breaks() -> list: return breaks +def _load_sweeper_module(): + import importlib.util + mod_spec = importlib.util.spec_from_file_location( + "followup_sweeper_gravity", + str(Path(__file__).resolve().parent / "followup-sweeper.py")) + mod = importlib.util.module_from_spec(mod_spec) + mod_spec.loader.exec_module(mod) + return mod + + +def maybe_run_terminal_fallback(rec, now_iso, dry_run=False): + """Run a job's on_no_result fallback once at terminal followup expiry. + + Returns an outcome dict, or None when the record has no resolvable + fallback. Never raises. The sweeper terminal path shares the + fallback_due() guard, so the two firing paths can't double-execute. + """ + try: + sw = _load_sweeper_module() + except Exception as e: + return {"ran": False, "mode": None, + "detail": f"sweeper import failed: {e}"} + try: + if not sw.fallback_due(rec): + return None + job_name = sw.derive_job_name(rec.get("job_id")) + if not job_name: + return None + spec, err = sw.load_job_fallback(job_name) + if err or spec is None: + return None + if dry_run: + return {"ran": False, "mode": "dry_run", + "detail": json.dumps(spec)[:200]} + out = sw.run_no_result_fallback(rec) + rec["fallback"] = {"ran": out["ran"], "mode": out["mode"], + "detail": out["detail"][:200], "ts": now_iso} + try: + sw.append_job_log({ + "ts": now_iso, + "type": ("fallback_executed" if out["ran"] + else "fallback_failed"), + "dm_id": rec.get("dm_id"), + "recipient": rec.get("recipient"), + "job_id": rec.get("job_id"), + "mode": out["mode"], + "detail": out["detail"][:300], + }) + except Exception: + pass + return out + except Exception as e: + return {"ran": False, "mode": None, + "detail": f"fallback exception: {e}"} + + def remediate_breaks(dry_run=False) -> dict: """Progressively auto-remediate soft loop breakages while escalating hard breakages. @@ -746,6 +802,20 @@ def remediate_breaks(dry_run=False) -> dict: f_modified = True rearm_sweeper = True + # Terminal: nudges exhausted and still no reply. Run the + # job's on_no_result fallback (server-side guarantee). + if is_expired and nudges_sent >= nudges_allowed: + fb = maybe_run_terminal_fallback(rec, now_iso, dry_run) + if fb is not None: + remediated.append({ + "action": "terminal_fallback", + "loop_id": dm_id, + "agent": rec.get("recipient"), + "detail": f"on_no_result ran={fb.get('ran')}: {fb.get('detail', '')[:160]}", + }) + if not dry_run and "fallback" in rec: + f_modified = True + if f_modified and not dry_run: tmp = f"{f_path}.tmp.{os.getpid()}" with open(tmp, "w") as f: @@ -756,7 +826,10 @@ def remediate_breaks(dry_run=False) -> dict: sweeper_py = Path("/home/super/Projects/NetVM/bin/followup-sweeper.py") if sweeper_py.exists(): try: - subprocess.run([sys.executable, str(sweeper_py), "--once"], timeout=10) + # A single nudge send takes ~10s median; give the sweep + # room to finish or it dies mid-first-send every time. + subprocess.run([sys.executable, str(sweeper_py), "--once"], + timeout=300) except Exception: pass @@ -819,4 +892,24 @@ def remediate_breaks(dry_run=False) -> dict: } +def main(argv=None): + import argparse + ap = argparse.ArgumentParser(description="Loop gravity: reconcile and remediate followup loops") + ap.add_argument("--remediate", action="store_true", + help="Run remediate_breaks once (what loop-remediator.timer invokes)") + ap.add_argument("--dry-run", action="store_true", + help="Report actions without writing state or sending anything") + args = ap.parse_args(argv) + if not args.remediate: + ap.print_help() + return 2 + result = remediate_breaks(dry_run=args.dry_run) + print(json.dumps(result, indent=2)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) + + diff --git a/bin/prompt_envelope.py b/bin/prompt_envelope.py index 837e003..b7f81da 100644 --- a/bin/prompt_envelope.py +++ b/bin/prompt_envelope.py @@ -17,7 +17,10 @@ BOX_API = "https://box.muse-dev.online/api/box" def _response_rule(): return ( - "\nRESPONSE RULE: Execute your steps using [TOOL ...] directives or background tmux commands." + "\nRESPONSE RULE: Act by EMITTING [TOOL ...] / [DM ...] directive lines verbatim in your reply" + " — you do not run them yourself. The Box runtime on bl executes each directive" + " (this works from containers with no box CLI or tmux socket) and posts the result back here." + " Background tmux commands work too when you have a shell." " When complete, conclude your output with the [RESULT ...] line so the harvester records it.\n" ) diff --git a/bin/response-harvester.py b/bin/response-harvester.py index e8cc32a..1571422 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -301,6 +301,61 @@ def _scan_bracket_calls(text): return out +_PROOF_EVIDENCE_RE = re.compile( + r"sw-\d{8}-\d{6}-[0-9a-f]{4}" + r"|[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}" + r"|/(?:[\w.-]+/)+[\w.-]+" + r"|\b(?:swarm|timer|cron|job|thread|sidechat|slot)[-_ ]?(?:id|name|uuid)?\s*[:=]" + r"|\b\d+/\d+\s*(?:slots?|checks?|workers?)", + re.IGNORECASE) +_UUID_RE = re.compile(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}") + + +def result_has_evidence(result_text): + """True when a RESULT verdict carries checkable artifacts (IDs, paths, counts).""" + return bool(_PROOF_EVIDENCE_RE.search(result_text or "")) + + +def maybe_request_proof(agent, thread_id, job_id, result_text, dry_run=False): + """Ask for checkable evidence when a success RESULT has none. + + One-shot per (thread, job) via the nudge tracker. Returns True when a + proof followup was scheduled. + """ + if dry_run or not thread_id or not _UUID_RE.fullmatch(thread_id.lower()): + return False + if result_has_evidence(result_text): + return False + tracker = load_json_file(NUDGE_TRACKER_FILE) + rec = tracker.get(thread_id, {}) + done = rec.get("proof_jobs", []) + if job_id in done: + return False + ok, res = execute_agent_tool(agent, "followup.create", { + "agent": agent, + "in_m": 30, + "thread": thread_id, + "prompt": ( + f"[PROOF] Your [RESULT {job_id}] has no checkable evidence. " + f"Reply in this thread with the swarm/timer IDs, paths, or command output " + f"that prove the outcome — or say what is still missing."), + }) + if ok: + rec["proof_jobs"] = (done + [job_id])[-50:] + tracker[thread_id] = rec + save_json_file(NUDGE_TRACKER_FILE, tracker) + append_jsonl(JOB_LOG, { + "ts": utcnow(), + "type": "proof_requested", + "job_id": job_id, + "agent": agent, + "thread_id": thread_id, + }) + return True + sys.stderr.write(f"warning: proof followup failed for {job_id}: {res}\n") + return False + + def parse_tool_calls(text): """ Extract structured tool/exec calls from assistant messages. @@ -893,7 +948,9 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo thread_url = f"https://box.muse-dev.online/thread/{thread_id}" tool_hint = ( f"[Runtime Context: {thread_url}]\n" - f"Tools: [TOOL ] or curl -sk -X POST https://exec.muse-dev.online/exec\n" + f"Tools: EMIT one [TOOL ] line per action (you do not run it;" + f" the runtime executes it and replies here). curl -sk -X POST" + f" https://exec.muse-dev.online/exec works too.\n" f" • [TOOL tools.list {{}}] — discover every op dynamically\n" f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"\"}}] — spawn subagents\n" f" • [DM {{\"to\": \"\", \"target\": \"\", \"message\": \"\"}}] — send a DM\n" @@ -941,6 +998,11 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo sys.stderr.write(f"warning: failed to record swarm report: {se}\n") else: trigger_chain_next(job_id, result_text, success=not is_fail) + if not is_fail: + try: + maybe_request_proof(agent, thread_id, job_id, result_text) + except Exception as pe: + sys.stderr.write(f"warning: proof check failed: {pe}\n") archive_ephemeral_thread(agent, thread_id, job_id=job_id) clear_matching_followups(followups, agent, thread_id, mid, text, dry_run, job_id=job_id, verb="RESULT") @@ -951,7 +1013,8 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo dry_run, job_id=job_id, verb=verb) else: clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) - maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=dry_run) + maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, + dry_run=dry_run, acted=bool(tool_calls)) return new_messages, new_wm, job_results @@ -1392,10 +1455,13 @@ def check_and_archive_terminal_swarms(): -def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=False): +def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=False, + acted=False): """ If an agent replies conversationally in a sidechat backed by a job or follow-up without providing [RESULT ] or tool directives, deliver a terse 1-turn nudge footer. + When acted=True the agent DID emit directives but never closed: remind to close + with [RESULT] instead of rejecting the (good) action. """ if dry_run or not thread_id or thread_id == "main": return @@ -1440,14 +1506,23 @@ def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_ matching_job_id, matching_job_id, prompt_envelope.pick_profile(matching_job_id)) except Exception: _spawn = '[TOOL swarm.spawn {"count": 2, "task": "continue the job work"}]' - nudge_msg = ( - f"{_spawn}\n" - f"[STRICT ENFORCEMENT: Conversational commentary is rejected. Work requires active execution.]\n" - f"Thread Console: {thread_url}\n" - f"Emit executable tool calls now: [TOOL ] or curl against https://exec.muse-dev.online/exec\n" - f"When all operations are finished, close strictly with [RESULT {matching_job_id}] .\n" - f"{_spawn}" - ) + if acted: + nudge_msg = ( + f"Action received — now close the loop: reply with [RESULT {matching_job_id}] .\n" + f"Outcome needs checkable evidence (swarm/timer IDs, paths, or command output), not prose alone.\n" + f"Thread Console: {thread_url}" + ) + else: + nudge_msg = ( + f"{_spawn}\n" + f"[STRICT ENFORCEMENT: Conversational commentary is rejected. Work requires active execution.]\n" + f"Thread Console: {thread_url}\n" + f"EMIT tool calls verbatim in your reply — you do not run them yourself;" + f" the Box runtime on bl executes each directive and posts the result back here" + f" (works from containers with no box CLI). Or curl against https://exec.muse-dev.online/exec\n" + f"When all operations are finished, close strictly with [RESULT {matching_job_id}] .\n" + f"{_spawn}" + ) try: import muse_hybrid print(f"[{agent}] Injecting 1-turn strict nudge into {thread_name or thread_id[:8]} for job {matching_job_id}") diff --git a/docs/JOB-SPEC.md b/docs/JOB-SPEC.md index 37d70b6..83bb8f9 100644 --- a/docs/JOB-SPEC.md +++ b/docs/JOB-SPEC.md @@ -225,6 +225,86 @@ spam if a job misfires in a loop. - DMs are logged (dm-log.jsonl) for audit - Side chats are per-job, not shared across trust boundaries +## Completion Enforcement + +Sending a job DM does not complete work. Four mechanisms close the loop on bl: + +### 1. `on_no_result` fallback (server-side guarantee) + +A job may declare a fallback effect that runs when its followup reaches +terminal expiry with no agent result (`bin/followup-sweeper.py`): + +```json +"on_no_result": {"op": "swarm.spawn", "args": {"count": 2, "task": "...", "label": "..."}} +"on_no_result": {"job": ""} +``` + +- `op` form: validated + built through the `exec-constrained` op registry + (same validators the daemon uses), then executed as a subprocess. +- `job` form: dispatches the named job via `job-dispatch.py` with + `CHAIN_PREV_*` timeout context. +- Outcomes log as `fallback_executed` / `fallback_failed` in `job-log.jsonl` + and stamp `rec["fallback"]` on the followup record. Failures never break + the sweep. Jobs without the key behave exactly as before. +- Firing paths (either; never both): the sweeper terminal branch + (`nudges_sent >= nudges_allowed`, expired) and `gravity.py:remediate_breaks` + (same terminal condition, runs on `loop-remediator.timer` ~every 15m). + `fallback_due()` dedupes: fires once per record, retries a failed attempt + after 1h. Gravity also revives the local sweep loop (it previously had no + `__main__`, so the timer was a no-op) with a 300s re-arm budget (was 10s, + which strangled every sweep mid-first-send — median nudge send is ~10s). +- Split-brain note: the VM board sweeper (`box-request-sweeper.timer`) sends + the live `[NUDGE ]` DMs from remote `dm_followup` requests; the local + sweeper sends `[nudge N/M]` from `followups.json`. Both fire per deadline + until a reply resolves both sides. Accept the duplicate nudge as the cost + of a guaranteed local path; cross-system dedup is future work. +- Seeded on: `autonomy-pulse-646`, `autonomy-pulse-pip`, `autonomy-pulse-opm` + (each spawns 2 standing-work subagents if the agent naps through the pulse). +- Race note: set the job's followup timeout longer than the expected work + time, or a slow-but-working agent can double-fire alongside the fallback. + +### 2. NACK for directive-less replies (existing, sharpened) + +`response-harvester.py:maybe_nudge_untagged_sidechat` already rejects +conversational replies in job-backed threads (1-turn strict nudge, then a +5-minute escalation timer to opm; tracked in +`conversation-nudge-tracker.json`). Two refinements: + +- Acted-variant: when the agent emitted directives but never closed with + `[RESULT]`, the nudge acknowledges the action and demands the close + instead of crying "commentary rejected". +- Emit-model wording: envelope (`prompt_envelope._response_rule`), tool + hint, and strict nudge now state plainly that agents EMIT `[TOOL]` / + `[DM]` lines verbatim and the Box runtime on bl executes them — this + works from containers with no box CLI or tmux socket. (Root cause of the + 2026-10-06 zero-`tool_exec` stretch: agents believed they had to execute + tools locally and declined for lack of a "container equivalent".) + +### 3. Proof-of-result followups + +A success `[RESULT]` with no checkable artifact (swarm/timer IDs, UUIDs, +paths, `n/m` completion counts — see `result_has_evidence`) triggers a +one-shot `followup.create` (+30m, same thread) asking for the evidence, and +logs `proof_requested`. One per (thread, job). Failures and declines skip +proof (they already chain via `on_failure`). + +### 4. Completion auditor (15-minute timer) + +`bin/completion-audit.py` via `systemd/completion-audit.timer` +(`OnCalendar=*:4/15`, installed to `/etc/systemd/system`): + +- Computes the funnel per job family over `--hours` (default 24) from + `job-log.jsonl`: sent / dispatched / `tool_exec` / results, plus + `fallback_*`, `proof_requested`, and `job_failed` counts. +- Adds point-in-time swarm slot drain (flags running slots idle >90m) and + followup backlog (pending / overdue / escalated). +- Writes `logs/completion-audit-.json` + `logs/completion-audit-latest.json`. +- Posts the digest to `opm` / `heartbeat` only when DEGRADED + (silent families, failures, tool errors, stale slots, overdue/escalated + followups); when green, posts a heartbeat at most every 6h + (`logs/completion-audit-state.json`). Silent-when-healthy otherwise. +- Read-only except the digest DM and its own log/state files. + ## Future Expansions 1. **Conditional jobs**: Run Job B only if Job A succeeds with specific output diff --git a/jobs/autonomy-pulse-646.json b/jobs/autonomy-pulse-646.json index 1cd2696..8df1ae7 100644 --- a/jobs/autonomy-pulse-646.json +++ b/jobs/autonomy-pulse-646.json @@ -18,5 +18,13 @@ "route": "autonomy-pulse" }, "chain_next": null, - "on_failure": "alert" + "on_failure": "alert", + "on_no_result": { + "op": "swarm.spawn", + "args": { + "count": 2, + "task": "Standing pulse work for 646 (autonomy-pulse-646): execute your scope standing work: verify timers, check swarm results, act or close; report per-slot verdicts.", + "label": "autonomy-pulse-646-fallback" + } + } } diff --git a/jobs/autonomy-pulse-opm.json b/jobs/autonomy-pulse-opm.json index 1ab8844..bb500bc 100644 --- a/jobs/autonomy-pulse-opm.json +++ b/jobs/autonomy-pulse-opm.json @@ -18,5 +18,13 @@ "route": "autonomy-pulse" }, "chain_next": null, - "on_failure": "alert" + "on_failure": "alert", + "on_no_result": { + "op": "swarm.spawn", + "args": { + "count": 2, + "task": "Standing pulse work for opm (autonomy-pulse-opm): execute your scope standing work: verify timers, check swarm results, act or close; report per-slot verdicts.", + "label": "autonomy-pulse-opm-fallback" + } + } } diff --git a/jobs/autonomy-pulse-pip.json b/jobs/autonomy-pulse-pip.json index a644a72..a070888 100644 --- a/jobs/autonomy-pulse-pip.json +++ b/jobs/autonomy-pulse-pip.json @@ -18,5 +18,13 @@ "route": "autonomy-pulse" }, "chain_next": null, - "on_failure": "alert" + "on_failure": "alert", + "on_no_result": { + "op": "swarm.spawn", + "args": { + "count": 2, + "task": "Standing pulse work for pip (autonomy-pulse-pip): execute your scope standing work: verify timers, check swarm results, act or close; report per-slot verdicts.", + "label": "autonomy-pulse-pip-fallback" + } + } } diff --git a/systemd/completion-audit.service b/systemd/completion-audit.service new file mode 100644 index 0000000..6676ab4 --- /dev/null +++ b/systemd/completion-audit.service @@ -0,0 +1,10 @@ +[Unit] +Description=NetVM completion funnel auditor (15m digest) + +[Service] +Type=oneshot +User=super +WorkingDirectory=/home/super/Projects/NetVM +ExecStart=/usr/bin/python3 /home/super/Projects/NetVM/bin/completion-audit.py --hours 24 +StandardOutput=journal +StandardError=journal diff --git a/systemd/completion-audit.timer b/systemd/completion-audit.timer new file mode 100644 index 0000000..cba52af --- /dev/null +++ b/systemd/completion-audit.timer @@ -0,0 +1,9 @@ +[Unit] +Description=Run completion auditor every 15 minutes + +[Timer] +OnCalendar=*:4/15 +Persistent=true + +[Install] +WantedBy=timers.target diff --git a/tests/test_completion.py b/tests/test_completion.py new file mode 100644 index 0000000..5a12c46 --- /dev/null +++ b/tests/test_completion.py @@ -0,0 +1,366 @@ +"""Tests for the completion-enforcement loop: on_no_result fallbacks, +proof-of-result followups, the NACK acted-variant, and the auditor funnel.""" +import importlib.util +import json +import sys +import types +import unittest +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parent.parent + + +def _load(mod_name, rel_path): + spec = importlib.util.spec_from_file_location(mod_name, REPO_ROOT / rel_path) + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + return mod + + +sw = _load("sweeper_completion", "bin/followup-sweeper.py") +harv = _load("harvester_completion", "bin/response-harvester.py") +aud = _load("completion_audit_mod", "bin/completion-audit.py") +grav = _load("gravity_completion", "bin/gravity.py") + + +class DeriveJobName(unittest.TestCase): + def test_valid(self): + self.assertEqual( + sw.derive_job_name("autonomy-pulse-646-20261006-060000-ed9a26af"), + "autonomy-pulse-646", + ) + + def test_invalid_shapes(self): + for bad in (None, "", "nonsense", "job-2026-1-abc", + "nonexistent-job-20261006-060000-ed9a26af"): + self.assertIsNone(sw.derive_job_name(bad), bad) + + +class LoadJobFallback(unittest.TestCase): + def test_seeded_pulse_jobs(self): + for name in ("autonomy-pulse-646", "autonomy-pulse-pip", + "autonomy-pulse-opm"): + spec, err = sw.load_job_fallback(name) + self.assertIsNone(err, name) + self.assertIsNotNone(spec, name) + self.assertEqual(spec["op"], "swarm.spawn") + self.assertEqual(spec["args"]["count"], 2) + + def test_absent_is_none(self): + spec, err = sw.load_job_fallback("646-daily-checkin") + self.assertIsNone(spec) + self.assertIsNone(err) + + +class RunNoResultFallback(unittest.TestCase): + def test_unresolvable_rec(self): + out = sw.run_no_result_fallback({"dm_id": "x"}) + self.assertFalse(out["ran"]) + self.assertFalse(out["configured"]) + + def test_job_without_spec(self): + out = sw.run_no_result_fallback( + {"dm_id": "x", "job_id": "646-daily-checkin-20261006-090000-abcdef12", + "recipient": "646"}) + self.assertFalse(out["ran"]) + self.assertFalse(out["configured"]) + + def test_dry_run_marks_configured(self): + out = sw.run_no_result_fallback( + {"dm_id": "x", + "job_id": "autonomy-pulse-646-20261006-060000-ed9a26af", + "recipient": "646"}, + dry_run=True) + self.assertTrue(out["configured"]) + self.assertFalse(out["ran"]) + + def test_op_path_uses_validated_build(self): + calls = {} + + class FakeOpError(Exception): + pass + + def fake_validate(args): + calls["validated"] = dict(args) + return {"echo": "yes"} + + def fake_build(clean): + calls["built"] = clean + return ["/bin/echo", "fallback-ok"] + + fake_mod = types.SimpleNamespace( + OPS={"probe.op": {"validate": fake_validate, "build": fake_build, + "timeout": 10}}) + orig_exec, orig_load, orig_derive = ( + sw._load_exec_ops, sw.load_job_fallback, sw.derive_job_name) + sw._load_exec_ops = lambda: fake_mod + sw.load_job_fallback = lambda name: ( + ({"op": "probe.op", "args": {"a": 1}}, None)) + sw.derive_job_name = lambda jid: "autonomy-pulse-646" + try: + out = sw.run_no_result_fallback( + {"dm_id": "x", "job_id": "whatever", "recipient": "646"}) + finally: + sw._load_exec_ops, sw.load_job_fallback, sw.derive_job_name = ( + orig_exec, orig_load, orig_derive) + self.assertTrue(out["ran"], out) + self.assertEqual(out["mode"], "op") + self.assertEqual(calls["validated"], {"a": 1}) + self.assertEqual(calls["built"], {"echo": "yes"}) + self.assertIn("fallback-ok", out["detail"]) + + def test_unknown_op_is_outcome_not_raise(self): + orig_exec, orig_load, orig_derive = ( + sw._load_exec_ops, sw.load_job_fallback, sw.derive_job_name) + sw._load_exec_ops = lambda: types.SimpleNamespace(OPS={}) + sw.load_job_fallback = lambda name: ({"op": "nope.nope"}, None) + sw.derive_job_name = lambda jid: "autonomy-pulse-646" + try: + out = sw.run_no_result_fallback( + {"dm_id": "x", "job_id": "whatever", "recipient": "646"}) + finally: + sw._load_exec_ops, sw.load_job_fallback, sw.derive_job_name = ( + orig_exec, orig_load, orig_derive) + self.assertTrue(out["configured"]) + self.assertFalse(out["ran"]) + self.assertIn("unknown op", out["detail"]) + + +class FallbackDue(unittest.TestCase): + def test_fresh_record_due(self): + self.assertTrue(sw.fallback_due({"dm_id": "x"})) + + def test_ran_never_due(self): + self.assertFalse(sw.fallback_due( + {"fallback": {"ran": True, "ts": "2026-10-06T00:00:00+00:00"}})) + + def test_failed_recent_not_due(self): + from datetime import datetime, timezone, timedelta + ts = (datetime.now(timezone.utc) - timedelta(minutes=5)).isoformat() + self.assertFalse(sw.fallback_due( + {"fallback": {"ran": False, "ts": ts}})) + + def test_failed_old_due(self): + self.assertTrue(sw.fallback_due( + {"fallback": {"ran": False, "ts": "2026-10-05T00:00:00+00:00"}})) + + +class GravityFallback(unittest.TestCase): + def test_dry_run_configured(self): + rec = {"dm_id": "x", + "job_id": "autonomy-pulse-646-20261006-060000-ed9a26af", + "recipient": "646"} + out = grav.maybe_run_terminal_fallback( + rec, "2026-10-06T08:00:00+00:00", dry_run=True) + self.assertIsNotNone(out) + self.assertFalse(out["ran"]) + self.assertNotIn("fallback", rec) + + def test_no_spec_returns_none(self): + rec = {"dm_id": "x", + "job_id": "646-daily-checkin-20261006-090000-abcdef12", + "recipient": "646"} + self.assertIsNone(grav.maybe_run_terminal_fallback( + rec, "2026-10-06T08:00:00+00:00", dry_run=True)) + + def test_already_ran_returns_none(self): + rec = {"dm_id": "x", + "job_id": "autonomy-pulse-646-20261006-060000-ed9a26af", + "recipient": "646", + "fallback": {"ran": True, "ts": "2026-10-06T07:00:00+00:00"}} + self.assertIsNone(grav.maybe_run_terminal_fallback( + rec, "2026-10-06T08:00:00+00:00", dry_run=True)) + + +class GravityEntry(unittest.TestCase): + def test_no_args_prints_help(self): + import io + from contextlib import redirect_stdout + buf = io.StringIO() + with redirect_stdout(buf): + rc = grav.main([]) + self.assertEqual(rc, 2) + self.assertIn("remediate", buf.getvalue()) + + def test_remediate_dry_run_returns_json(self): + import io + from contextlib import redirect_stdout + buf = io.StringIO() + with redirect_stdout(buf): + rc = grav.main(["--remediate", "--dry-run"]) + self.assertEqual(rc, 0) + data = json.loads(buf.getvalue()) + self.assertTrue(data["ok"]) + self.assertTrue(data["dry_run"]) + + +class ResultEvidence(unittest.TestCase): + def test_positives(self): + for text in ( + "done, swarm sw-20261006-060000-ab12 reported 2/2", + "wrote /tmp/out.json with 40 rows", + "timer id: pulse-15m restarted", + "thread 1dfb3199-2f99-446c-83e3-848ae2da0a12 swept", + "3/3 slots complete", + ): + self.assertTrue(harv.result_has_evidence(text), text) + + def test_negatives(self): + for text in ("OK all good", "done, nothing to report", "", None): + self.assertFalse(harv.result_has_evidence(text), repr(text)) + + +class ProofRequest(unittest.TestCase): + def _patch(self, tmp_path): + orig = (harv.NUDGE_TRACKER_FILE, harv.JOB_LOG, harv.execute_agent_tool) + harv.NUDGE_TRACKER_FILE = tmp_path / "tracker.json" + harv.JOB_LOG = tmp_path / "job-log.jsonl" + calls = [] + harv.execute_agent_tool = lambda a, op, args: ( + calls.append((a, op, args)), (True, "scheduled"))[1] + return orig, calls + + def test_bare_result_schedules_proof(self): + import tempfile + with tempfile.TemporaryDirectory() as td: + orig, calls = self._patch(Path(td)) + try: + ok = harv.maybe_request_proof( + "646", "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "job-1", "OK all good") + finally: + (harv.NUDGE_TRACKER_FILE, harv.JOB_LOG, + harv.execute_agent_tool) = orig + self.assertTrue(ok) + self.assertEqual(calls[0][1], "followup.create") + self.assertEqual(calls[0][2]["in_m"], 30) + self.assertIn("PROOF", calls[0][2]["prompt"]) + + def test_evidence_skips(self): + ok = harv.maybe_request_proof( + "646", "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "job-1", "done, swarm sw-20261006-060000-ab12") + self.assertFalse(ok) + + def test_bad_thread_and_dry_run_skip(self): + self.assertFalse(harv.maybe_request_proof( + "646", "main", "job-1", "OK")) + self.assertFalse(harv.maybe_request_proof( + "646", "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "job-1", "OK", dry_run=True)) + + def test_one_shot_per_job(self): + import tempfile + with tempfile.TemporaryDirectory() as td: + orig, calls = self._patch(Path(td)) + try: + kw = dict(agent="646", + thread_id="1dfb3199-2f99-446c-83e3-848ae2da0a12", + job_id="job-9", result_text="OK") + self.assertTrue(harv.maybe_request_proof(**kw)) + self.assertFalse(harv.maybe_request_proof(**kw)) + finally: + (harv.NUDGE_TRACKER_FILE, harv.JOB_LOG, + harv.execute_agent_tool) = orig + self.assertEqual(len(calls), 1) + + +class NackActedVariant(unittest.TestCase): + def test_acted_variant_text(self): + import tempfile + sent = [] + fake = types.ModuleType("muse_hybrid") + fake.send_message = lambda a, m, thread_id=None, wait=0: sent.append(m) + with tempfile.TemporaryDirectory() as td: + tp = Path(td) + (tp / "followups.json").write_text(json.dumps({ + "m1": {"status": "pending", + "thread_uuid": "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "job_id": "job-7", "recipient": "646"}})) + orig = (harv.NUDGE_TRACKER_FILE, harv.FOLLOWUPS_FILE, + sys.modules.get("muse_hybrid")) + harv.NUDGE_TRACKER_FILE = tp / "tracker.json" + harv.FOLLOWUPS_FILE = tp / "followups.json" + sys.modules["muse_hybrid"] = fake + try: + harv.maybe_nudge_untagged_sidechat( + "646", "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "646 tasks", "mid-1", "some prose", acted=True) + finally: + (harv.NUDGE_TRACKER_FILE, harv.FOLLOWUPS_FILE, + old_mod) = orig + if old_mod is None: + sys.modules.pop("muse_hybrid", None) + else: + sys.modules["muse_hybrid"] = old_mod + self.assertEqual(len(sent), 1) + self.assertIn("Action received", sent[0]) + self.assertIn("[RESULT job-7]", sent[0]) + self.assertNotIn("STRICT ENFORCEMENT", sent[0]) + + +class AuditorFunnel(unittest.TestCase): + def _events(self): + base = "2026-10-06T07:00:00+00:00" + return [ + {"ts": base, "type": "job_sent", + "job_id": "work-finder-20261006-070000-aaaaaaaa"}, + {"ts": base, "type": "job_dispatched", + "job_id": "work-finder-20261006-070000-aaaaaaaa"}, + {"ts": base, "type": "tool_exec", "op": "swarm.spawn", + "success": True}, + {"ts": base, "type": "job_result", + "job_id": "work-finder-20261006-070000-aaaaaaaa", + "success": True}, + {"ts": base, "type": "job_sent", + "job_id": "pulse-20261006-070000-bbbbbbbb"}, + {"ts": base, "type": "job_failed", + "job_id": "pulse-20261006-070000-bbbbbbbb"}, + ] + + def test_funnel_counts(self): + from datetime import datetime, timezone + fam, tools = aud.compute_funnel( + self._events(), datetime(2026, 10, 6, 6, 0, tzinfo=timezone.utc)) + self.assertEqual(fam["work-finder"]["sent"], 1) + self.assertEqual(fam["work-finder"]["results"], 1) + self.assertEqual(fam["pulse"]["failed"], 1) + self.assertEqual(tools["tools"]["total"], 1) + self.assertEqual(tools["tools"]["ok"], 1) + + def test_family_of(self): + self.assertEqual(aud.family_of("a-b-20261006-070000-aaaaaaaa"), "a-b") + self.assertEqual(aud.family_of("weird"), "weird") + + def test_digest_verdicts(self): + healthy = {"ts": "2026-10-06T07:00:00+00:00", "window_h": 6, + "totals": {"sent": 4, "dispatched": 4, "results": 4, + "ok": 4}, + "tools": {"tools": {"total": 3, "ok": 3}, "tool_errs": {}}, + "families": {}, "silent_families": [], "swarms": {}, + "stale_running": [], "followups": {"pending": 1}, + "degraded": False, "reasons": []} + out = aud.render_digest(healthy) + self.assertIn("HEALTHY", out) + self.assertIn("4 sent", out) + bad = dict(healthy, degraded=True, + reasons=["1 silent families: x"], + silent_families=["x"]) + self.assertIn("DEGRADED", aud.render_digest(bad)) + + def test_should_post_policy(self): + import tempfile + with tempfile.TemporaryDirectory() as td: + orig = aud.STATE_FILE + aud.STATE_FILE = Path(td) / "state.json" + try: + self.assertTrue(aud.should_post({"degraded": True})[0]) + ok, why = aud.should_post({"degraded": False}) + self.assertTrue(ok) + self.assertEqual(why, "heartbeat") + finally: + aud.STATE_FILE = orig + + +if __name__ == "__main__": + unittest.main()