Files
box/bin/followup-sweeper.py
T

377 lines
15 KiB
Python
Executable File

#!/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
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
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
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,
}
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()