2026-10-04 16:23:10 +00:00
|
|
|
#!/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
|
2026-10-04 20:06:58 +00:00
|
|
|
import re
|
2026-10-04 16:23:10 +00:00
|
|
|
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"
|
2026-10-04 16:34:25 +00:00
|
|
|
JOBS_DIR = NETVM_ROOT / "jobs"
|
2026-10-04 16:23:10 +00:00
|
|
|
FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
|
|
|
|
|
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
|
2026-10-04 20:06:58 +00:00
|
|
|
DM_LOG = NETVM_ROOT / "dm-log.jsonl"
|
2026-10-04 16:23:10 +00:00
|
|
|
DM_PY = BIN_DIR / "dm.py"
|
2026-10-04 16:34:25 +00:00
|
|
|
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
|
2026-10-04 16:23:10 +00:00
|
|
|
|
|
|
|
|
|
|
|
|
|
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):
|
2026-10-04 18:29:13 +00:00
|
|
|
"""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)."""
|
2026-10-04 16:23:10 +00:00
|
|
|
cmd = [
|
|
|
|
|
sys.executable,
|
|
|
|
|
str(DM_PY),
|
|
|
|
|
"send",
|
|
|
|
|
"--agent", sender,
|
|
|
|
|
"--to", recipient,
|
|
|
|
|
"--target", target,
|
2026-10-04 18:29:13 +00:00
|
|
|
] + (["--allow-main-chat"] if target == "main" else []) + [
|
2026-10-04 16:23:10 +00:00
|
|
|
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)
|
|
|
|
|
|
2026-10-04 20:06:58 +00:00
|
|
|
_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
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
|
2026-10-06 08:20:14 +00:00
|
|
|
_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 (<name>-YYYYMMDD-HHMMSS-<hex8>)."""
|
|
|
|
|
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": "<job-name>"} -> dispatch a fallback job, or
|
|
|
|
|
{"op": "<exec-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
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
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
|
2026-10-06 08:20:14 +00:00
|
|
|
fallbacks_count = 0
|
2026-10-04 16:23:10 +00:00
|
|
|
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")
|
|
|
|
|
|
2026-10-04 18:29:13 +00:00
|
|
|
# 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
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
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()
|
2026-10-04 20:06:58 +00:00
|
|
|
# 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"
|
2026-10-04 16:23:10 +00:00
|
|
|
modified = True
|
|
|
|
|
append_job_log({
|
|
|
|
|
"ts": utcnow_str(),
|
|
|
|
|
"type": "followup_nudged",
|
|
|
|
|
"dm_id": dm_id,
|
|
|
|
|
"nudge_num": nudge_num,
|
|
|
|
|
"recipient": recipient,
|
|
|
|
|
"target": delivery_target,
|
2026-10-04 20:06:58 +00:00
|
|
|
"thread_uuid": rec.get("thread_uuid"),
|
2026-10-04 16:23:10 +00:00
|
|
|
})
|
|
|
|
|
else:
|
|
|
|
|
print(f"Sweeper WARNING: nudge send failed: {out}", file=sys.stderr)
|
2026-10-04 18:29:13 +00:00
|
|
|
# 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,
|
|
|
|
|
})
|
2026-10-04 16:23:10 +00:00
|
|
|
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,
|
|
|
|
|
})
|
2026-10-04 16:34:25 +00:00
|
|
|
|
|
|
|
|
# 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
|
2026-10-06 08:20:14 +00:00
|
|
|
# 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
|
2026-10-04 16:23:10 +00:00
|
|
|
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,
|
2026-10-06 08:20:14 +00:00
|
|
|
"fallbacks": fallbacks_count,
|
2026-10-04 16:23:10 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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()
|