#!/usr/bin/env python3 """ followup-sweeper.py — Autonomous follow-up deadline tracking and nudge sweeper. Monitors pending follow-ups in followups.json, delivers progressive nudges to recipients when deadlines expire (in-thread first, Main Chat on final nudge), and executes terminal escalations to opm when all nudges are exhausted. Usage: python3 followup-sweeper.py --once python3 followup-sweeper.py --loop --interval 30 """ import argparse import json import os import re import subprocess import sys import time from datetime import datetime, timezone, timedelta from pathlib import Path # Paths NETVM_ROOT = Path("/home/super/Projects/NetVM") BIN_DIR = NETVM_ROOT / "bin" JOBS_DIR = NETVM_ROOT / "jobs" FOLLOWUPS_FILE = NETVM_ROOT / "followups.json" JOB_LOG = NETVM_ROOT / "job-log.jsonl" DM_LOG = NETVM_ROOT / "dm-log.jsonl" DM_PY = BIN_DIR / "dm.py" DISPATCH_PY = BIN_DIR / "job-dispatch.py" sys.path.insert(0, str(BIN_DIR)) try: import pipeline_engine HAS_PIPELINE = True except ImportError: HAS_PIPELINE = False def utcnow_dt(): return datetime.now(timezone.utc) def utcnow_str(): return utcnow_dt().isoformat() def parse_iso(ts_str): if not ts_str: return None try: ts_clean = ts_str.replace("Z", "+00:00") dt = datetime.fromisoformat(ts_clean) if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) return dt except Exception: return None def load_followups(): if not FOLLOWUPS_FILE.exists(): return {} try: with open(FOLLOWUPS_FILE, "r", encoding="utf-8") as f: return json.load(f) except Exception: return {} def save_followups(data): tmp_path = f"{FOLLOWUPS_FILE}.tmp.{os.getpid()}" with open(tmp_path, "w", encoding="utf-8") as f: json.dump(data, f, indent=2) os.replace(tmp_path, FOLLOWUPS_FILE) def append_job_log(entry): os.makedirs(os.path.dirname(os.path.abspath(JOB_LOG)), exist_ok=True) with open(JOB_LOG, "a", encoding="utf-8") as f: f.write(json.dumps(entry) + "\n") def send_dm(sender, recipient, target, text): """Dispatch a DM via dm.py. The sweeper's delivery target is explicit followup state (or explicit in-code escalation routing), so pass the main-chat opt-in when the target is main (sidechat-first policy).""" cmd = [ sys.executable, str(DM_PY), "send", "--agent", sender, "--to", recipient, "--target", target, ] + (["--allow-main-chat"] if target == "main" else []) + [ text, ] try: res = subprocess.run(cmd, capture_output=True, text=True, timeout=90) return res.returncode == 0, res.stdout.strip() or res.stderr.strip() except Exception as e: return False, str(e) _DM_ID_RE = re.compile(r"\bDM ([0-9a-f]{8})\b") _THREAD_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}$", re.IGNORECASE, ) def _resolve_nudge_thread_uuid(nudge_output): """Parse the DM id from dm.py's stdout, scan the tail (last 5000 lines) of dm-log.jsonl for that id's sidechat_autoprovisioned event, falling back to the sent event's tags.thread. Returns the UUID or None.""" m = _DM_ID_RE.search(nudge_output or "") if not m: return None nudge_id = m.group(1) fallback = None try: with open(DM_LOG, "r", encoding="utf-8") as f: lines = f.readlines() except FileNotFoundError: return None for line in lines[-5000:]: line = line.strip() if not line: continue try: e = json.loads(line) except Exception: continue if e.get("id") != nudge_id: continue if e.get("type") == "sidechat_autoprovisioned": uuid = e.get("thread_uuid") or "" if _THREAD_UUID_RE.fullmatch(uuid): return uuid elif e.get("type") == "sent": cand = ((e.get("tags") or {}).get("thread")) or "" if _THREAD_UUID_RE.fullmatch(cand): fallback = cand 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: return {"status": "ok", "pending": 0, "nudges_sent": 0, "escalations": 0} now = utcnow_dt() nudges_count = 0 escalations_count = 0 fallbacks_count = 0 modified = False for dm_id, rec in list(followups.items()): if rec.get("status") != "pending": continue deadline_dt = parse_iso(rec.get("deadline")) if not deadline_dt or now < deadline_dt: continue # Deadline has expired! nudges_sent = rec.get("nudges_sent", 0) nudges_allowed = rec.get("nudges_allowed", 2) sender = rec.get("sender", "opm") recipient = rec.get("recipient") orig_target = rec.get("target", "main") thread_uuid = rec.get("thread_uuid") # Ghost-followup fail-fast (2026-10-04, operator-main): a null # thread_uuid means the sidechat was never provisioned, so the # harvester can never match a reply. Nudging is pointless -- flag # for manual triage once instead of burning the nudge budget and # escalating a ghost. Main-chat followups are unaffected (the # harvester matches those by target). if thread_uuid is None and orig_target != "main" and not rec.get("needs_review"): rec["needs_review"] = True rec["status"] = "needs_review" rec["review_reason"] = ( "ghost: unresolvable, manual triage " f"(thread_uuid null, target={orig_target}; " "reply can never auto-resolve)" ) modified = True append_job_log({ "ts": utcnow_str(), "type": "followup_ghost_suppressed", "dm_id": dm_id, "recipient": recipient, "target": orig_target, "reason": "ghost: unresolvable, manual triage", }) print(f"Sweeper: suppressing ghost followup {dm_id} " f"(thread_uuid null, target={orig_target}) -> needs_review", file=sys.stderr) continue if nudges_sent < nudges_allowed: # Deliver next nudge nudge_num = nudges_sent + 1 is_final = (nudge_num == nudges_allowed) # Routing: In-thread first, Main on final nudge delivery_target = "main" if is_final else orig_target nudge_text = ( f"[nudge {nudge_num}/{nudges_allowed}] [ref:{dm_id}] " f"Reminder: awaiting reply to request sent at {rec.get('sent_at', 'earlier')}." ) if is_final and orig_target != "main": nudge_text += f" (Origin thread: {orig_target})" print(f"Sweeper: Sending nudge {nudge_num}/{nudges_allowed} to {recipient}/{delivery_target}...") if not dry_run: ok, out = send_dm(sender, recipient, delivery_target, nudge_text) if ok: nudges_count += 1 rec["nudges_sent"] = nudge_num rec["last_nudge_at"] = utcnow_str() # Calculate interval for next nudge: proportional to timeout or default 10m timeout_s = rec.get("timeout_s", 1800) step_s = max(300, timeout_s // (nudges_allowed + 1)) rec["deadline"] = (now + timedelta(seconds=step_s)).isoformat() # C1: backfill the thread the nudge actually landed in. landed_uuid = _resolve_nudge_thread_uuid(out) if landed_uuid and rec.get("thread_uuid") != landed_uuid: rec["thread_uuid"] = landed_uuid print(f"Sweeper: backfilled thread_uuid={landed_uuid} " f"for followup {dm_id}", file=sys.stderr) # C2: record final-nudge routing for the harvester. if is_final: rec["final_nudge_target"] = "main" modified = True append_job_log({ "ts": utcnow_str(), "type": "followup_nudged", "dm_id": dm_id, "nudge_num": nudge_num, "recipient": recipient, "target": delivery_target, "thread_uuid": rec.get("thread_uuid"), }) else: print(f"Sweeper WARNING: nudge send failed: {out}", file=sys.stderr) # Failure-mode fix (2026-10-04): advance state on send # failure so a failing nudge is never re-fired every # timer tick. failed_sends is tracked separately from # nudges_sent so a delivery failure does not consume a # real nudge. Backoff: 5 min base, doubling per failure. failed = rec.get("failed_sends", 0) + 1 rec["failed_sends"] = failed rec["last_failure_at"] = utcnow_str() backoff_s = 300 * (2 ** min(failed - 1, 4)) # 5m,10m,20m,40m,80m cap rec["deadline"] = (now + timedelta(seconds=backoff_s)).isoformat() modified = True append_job_log({ "ts": utcnow_str(), "type": "followup_nudge_failed", "dm_id": dm_id, "nudge_num": nudge_num, "recipient": recipient, "target": delivery_target, "failed_sends": failed, "backoff_s": backoff_s, "error": str(out)[:200], }) # After 3 consecutive failures, stop retrying blindly and # flag for manual review (delivery may be uncertain or the # target may be permanently broken). if failed >= 3: rec["needs_review"] = True rec["review_reason"] = ( f"nudge send failed {failed} times consecutively " f"(last: {str(out)[:120]})" ) append_job_log({ "ts": utcnow_str(), "type": "followup_needs_review", "dm_id": dm_id, "recipient": recipient, "failed_sends": failed, }) else: nudges_count += 1 else: # All nudges exhausted: Terminal escalation escalate_to = rec.get("escalate_to", "opm") esc_text = ( f"[ESCALATION] Agent {recipient} failed to reply to DM {dm_id} " f"after {nudges_allowed} nudges. Target was: {orig_target} " f"(thread: {thread_uuid or 'n/a'}). Request sent: {rec.get('sent_at')}." ) print(f"Sweeper: Escalating expired follow-up {dm_id} to {escalate_to}...") if not dry_run: ok, out = send_dm("bl", escalate_to, "main", esc_text) rec["status"] = "escalated" rec["escalated_at"] = utcnow_str() escalations_count += 1 modified = True append_job_log({ "ts": utcnow_str(), "type": "followup_escalated", "dm_id": dm_id, "recipient": recipient, "escalated_to": escalate_to, }) # Check if this dm_id belongs to an active pipeline step if HAS_PIPELINE: run_entry, step_entry = pipeline_engine.record_step_timeout(dm_id) if run_entry and step_entry: j_name = step_entry.get("job_name") j_file = JOBS_DIR / f"{j_name}.json" if j_file.exists(): try: with open(j_file, "r", encoding="utf-8") as f: j_cfg = json.load(f) on_failure = j_cfg.get("on_failure") if on_failure and (JOBS_DIR / f"{on_failure}.json").exists(): env = os.environ.copy() env["CHAIN_PREV_JOB_ID"] = step_entry.get("job_id") env["CHAIN_PREV_RESULT"] = f"TIMEOUT: Agent {recipient} timed out after {nudges_allowed} nudges" env["CHAIN_PIPELINE_RUN_ID"] = run_entry.get("run_id") next_step_n = step_entry.get("step_n", 1) + 1 env["CHAIN_STEP_N"] = str(next_step_n) cmd = [sys.executable, str(DISPATCH_PY), on_failure, "--pipeline-run", run_entry.get("run_id"), "--step-n", str(next_step_n)] subprocess.Popen(cmd, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) else: 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 if modified and not dry_run: save_followups(followups) pending_count = sum(1 for r in followups.values() if r.get("status") == "pending") return { "status": "ok", "pending": pending_count, "nudges_sent": nudges_count, "escalations": escalations_count, "fallbacks": fallbacks_count, } def main(): parser = argparse.ArgumentParser(description="Autonomous follow-up deadline tracker and sweeper") parser.add_argument("--once", action="store_true", help="Run once and exit (default)") parser.add_argument("--loop", action="store_true", help="Run continuously in a daemon loop") parser.add_argument("--interval", type=int, default=60, help="Interval in seconds for loop (default 60)") parser.add_argument("--dry-run", action="store_true", help="Inspect without sending nudges or updating records") args = parser.parse_args() if not args.loop: stats = sweep_cycle(dry_run=args.dry_run) print(f"[{datetime.now(timezone.utc).strftime('%H:%M:%SZ')}] Sweep cycle: {stats['pending']} pending, {stats['nudges_sent']} nudges, {stats['escalations']} escalations.") return print(f"Starting follow-up sweeper loop (interval={args.interval}s)...") while True: try: stats = sweep_cycle(dry_run=args.dry_run) print(f"[{datetime.now(timezone.utc).strftime('%H:%M:%SZ')}] Sweep cycle: {stats['pending']} pending, {stats['nudges_sent']} nudges, {stats['escalations']} escalations.") except Exception as e: print(f"ERROR in sweeper loop: {e}", file=sys.stderr) time.sleep(args.interval) if __name__ == "__main__": main()