#!/usr/bin/env python3 """Timer-driven side-chat adoption: shared gravity library. Pure Python, no network, no bl dependencies — everything the timers need to decide *where* a message goes and *whether* its loop is alive. Components: GravityConfig knobs from gravity.json (fail-closed defaults) ThreadRegistry (recipient, purpose) -> thread_uuid, JSON-backed resolve_target() explicit target > registry hit > autocreate > fallback LoopState FIRING/LANDED/SEND_FAILED/SEEN/ANSWERED/NUDGED/ ESCALATED/CLOSED/BROKEN + transition helpers detect_breaks() loop-break taxonomy detectors over log-derived events actionable_digest() wrap a wake digest as a tracked loop message loop_health() per-agent ANSWERED+CLOSED / LANDED The live integrations (job-dispatch.py, dm.py, sidechat-wake.py) import the pure functions here; the patches/ directory shows the call-site diffs. """ import json import os import re import sys import time from pathlib import Path # ---------------------------------------------------------------- config DEFAULTS = { "job_default_sidechat": True, "job_sidechat_fallback": "main", "dm_prefer_sidechat_default": False, "dm_sidechat_fallback": "main", "wake_actionable": True, "wake_ack_timeout": 7200, "wake_ack_nudges": 1, "wake_ack_escalate": "opm", "thread_autocreate": True, "loop_silent_ticks": 3, "loop_health_threshold": 0.5, } VALID_AGENTS = ("muse", "pip", "646", "opm") 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}") class GravityConfig: """Knobs. Unknown keys rejected (typo guard); missing keys get fail-closed defaults (see DEFAULTS).""" def __init__(self, data=None): data = data or {} unknown = set(data) - set(DEFAULTS) if unknown: raise ValueError("unknown gravity keys: %s" % sorted(unknown)) self._d = dict(DEFAULTS) self._d.update(data) @classmethod def load(cls, path): with open(path) as f: return cls(json.load(f)) def get(self, key): return self._d[key] def __getitem__(self, key): return self._d[key] def as_dict(self): return dict(self._d) # ---------------------------------------------------------------- registry class ThreadRegistry: """(recipient, purpose) -> {thread_uuid, created_ts, last_verified_ts}. A thread UUID is only valid in the account whose browser created it, so the registry is keyed by recipient — never share UUIDs across accounts. Rotation: set()/record_rotation() keep the key pointed at the newest UUID and stash the old one under `:previous` for audit (same model as thread_lifecycle.record_rotation). """ def __init__(self, path=None, state=None): self.path = path self._s = dict(state or {}) @classmethod def load(cls, path): try: with open(path) as f: return cls(path=path, state=json.load(f)) except (OSError, ValueError): return cls(path=path) def save(self): if not self.path: raise ValueError("no path for registry save") tmp = self.path + ".tmp" with open(tmp, "w") as f: json.dump(self._s, f, indent=2) os.replace(tmp, self.path) def _key(self, recipient, purpose): return "%s:%s" % (recipient, purpose) def get(self, recipient, purpose): """Current thread UUID for (recipient, purpose), following any rotation chain. Returns None if unknown.""" return self.resolve(self._key(recipient, purpose)) def resolve(self, key): """Return the current UUID for a registry key. set()/record_rotation() always keep the key pointing at the newest UUID; the `:previous` entry is audit history, not a traversal chain. (Mirrors the guarantee in thread_lifecycle.record_rotation.) """ cur = self._s.get(key) if isinstance(cur, str) and UUID_RE.fullmatch(cur): return cur return None def previous(self, recipient, purpose): """The rotated-away UUID, if any (for diagnostics).""" return self._s.get(self._key(recipient, purpose) + ":previous") def set(self, recipient, purpose, thread_uuid, ts=None): if not UUID_RE.fullmatch(thread_uuid or ""): raise ValueError("not a thread UUID: %r" % thread_uuid) key = self._key(recipient, purpose) old = self._s.get(key) if old and old != thread_uuid: self._s[key + ":previous"] = old self._s[key] = thread_uuid if ts: self._s[key + ":updated_ts"] = ts def record_rotation(self, recipient, purpose, new_uuid, ts=None): """Alias for set() — rotation is just a set with history.""" self.set(recipient, purpose, new_uuid, ts=ts) def known_purposes(self, recipient): prefix = recipient + ":" out = [] for k in self._s: if k.startswith(prefix) and ":" not in k[len(prefix):]: out.append(k[len(prefix):]) return sorted(out) # ---------------------------------------------------------------- target resolution def resolve_target(recipient, purpose, cfg, registry, explicit_target=None, probe=None): """Decide where a timer-driven send goes. Order: explicit --target > registry hit (verified) > autocreate > configured fallback. Never raises for routing reasons — worst case is the fallback with a logged reason. probe(thread_uuid) -> bool: cheap reachability check (sidechat use). create() -> thread_uuid | None: injected by caller (needs browser). Returns (target, thread_uuid_or_None, reason). """ if explicit_target: tid = UUID_RE.search(explicit_target) return explicit_target, tid.group(0) if tid else None, "explicit" if recipient not in VALID_AGENTS: return cfg["dm_sidechat_fallback"], None, "unknown_recipient" uuid = registry.get(recipient, purpose) if uuid: if probe is None or probe(uuid): return uuid, uuid, "registry_hit" return cfg["dm_sidechat_fallback"], None, "registry_stale" # No mapping: autocreate or fallback. return None, None, "no_mapping_needs_create" def fallback_target(cfg): """Configured fallback when no sidechat is available. Fails closed toward visibility: main, never drop.""" return cfg.get("dm_sidechat_fallback") or "main" # ---------------------------------------------------------------- loop state machine # Terminal states: CLOSED, ESCALATED, BROKEN, SEND_FAILED. STATES = ("FIRING", "LANDED", "SEND_FAILED", "SEEN", "ANSWERED", "NUDGED", "ESCALATED", "CLOSED", "BROKEN") TERMINAL = frozenset(("SEND_FAILED", "ESCALATED", "CLOSED", "BROKEN")) # Allowed transitions. NUDGED can cycle (nudge 1..N) until ANSWERED or # ESCALATED; BROKEN is reachable from any non-terminal state. TRANSITIONS = { "FIRING": ("LANDED", "SEND_FAILED", "BROKEN"), "LANDED": ("SEEN", "ANSWERED", "NUDGED", "BROKEN"), "SEEN": ("ANSWERED", "NUDGED", "BROKEN"), "ANSWERED": ("CLOSED", "BROKEN"), "NUDGED": ("ANSWERED", "NUDGED", "ESCALATED", "BROKEN"), "SEND_FAILED": (), "ESCALATED": (), "CLOSED": (), "BROKEN": (), } class LoopState: """One timer-driven loop: timer tick -> tracked follow-up.""" def __init__(self, loop_id, agent, purpose, thread_uuid=None): self.loop_id = loop_id self.agent = agent self.purpose = purpose self.thread_uuid = thread_uuid self.state = "FIRING" self.nudges = 0 self.ticks_without_reply = 0 self.history = [("FIRING", None)] def transition(self, to, note=None): if to not in TRANSITIONS.get(self.state, ()): raise ValueError("illegal loop transition %s -> %s" % (self.state, to)) self.state = to if to == "NUDGED": self.nudges += 1 self.history.append((to, note)) return self @property def terminal(self): return self.state in TERMINAL @property def healthy(self): return self.state in ("ANSWERED", "CLOSED") def tick(self): """One scheduler tick with no reply observed.""" self.ticks_without_reply += 1 return self.ticks_without_reply # ---------------------------------------------------------------- break detectors def detect_breaks(loop, cfg, checks): """Loop-break taxonomy over caller-supplied check results. checks: dict with boolean-ish keys: thread_reachable, signature_ok, scheduler_alive, reply_observed Returns a list of break names (may be empty). """ breaks = [] if loop.terminal: return breaks if not checks.get("scheduler_alive", True): breaks.append("scheduler_death") if not checks.get("signature_ok", True): breaks.append("auth_rot") if not checks.get("thread_reachable", True): breaks.append("dead_thread") # signature rot: thread was rotated — the stored uuid is stale but a # :previous chain exists. Caller signals via thread_rotated=True and # should already have resolved; flag if not resolved. if checks.get("thread_rotated") and not checks.get("thread_resolved"): breaks.append("signature_rot") if (not checks.get("reply_observed") and loop.ticks_without_reply >= cfg["loop_silent_ticks"] and loop.state in ("LANDED", "SEEN", "NUDGED")): breaks.append("silent_agent") return breaks # ---------------------------------------------------------------- actionable digest def actionable_digest(digest_text, ack_line=True): """Wrap a wake digest so the timer tick becomes a tracked loop. Adds the ack line that the follow-up record's reply closes. The caller sends the result via `dm.py send --expect-reply --thread ` so a dm_followup record exists — without that send path this is just text. """ lines = digest_text.rstrip().split("\n") if ack_line: lines.append("") lines.append("Reply here to acknowledge (closes the loop).") return "\n".join(lines) def dm_send_argv(sender, recipient, target_uuid, message, cfg, purpose="wake"): """Build the dm.py argv that makes a timer send a *tracked* loop. Thread binding (--thread) is what lets nudges route back into the thread instead of falling back to DM. """ return [ "send", "--agent", sender, "--to", recipient, "--target", target_uuid, "--thread", target_uuid, "--expect-reply", "--reply-timeout", str(cfg["wake_ack_timeout"]), "--reply-nudges", str(cfg["wake_ack_nudges"]), "--reply-escalate", cfg["wake_ack_escalate"], "--tag", "purpose:%s" % purpose, message, ] # ---------------------------------------------------------------- loop health def loop_health(loops, threshold=0.5): """Per-agent loop health: (ANSWERED + CLOSED) / LANDED. loops: iterable of LoopState (or dicts with agent/state keys). Returns {agent: {"landed": n, "answered": n, "health": float|None, "below_threshold": bool}}. """ def _agent(l): return l.agent if isinstance(l, LoopState) else l.get("agent") def _state(l): return l.state if isinstance(l, LoopState) else l.get("state") # A loop counts as LANDED once it leaves FIRING via LANDED (or beyond). LANDED_OR_BEYOND = ("LANDED", "SEEN", "ANSWERED", "NUDGED", "ESCALATED", "CLOSED", "BROKEN") acc = {} for l in loops: a, s = _agent(l), _state(l) r = acc.setdefault(a, {"landed": 0, "answered": 0}) if s in LANDED_OR_BEYOND: r["landed"] += 1 if s in ("ANSWERED", "CLOSED"): r["answered"] += 1 out = {} for a, r in acc.items(): h = (r["answered"] / r["landed"]) if r["landed"] else None out[a] = {"landed": r["landed"], "answered": r["answered"], "health": h, "below_threshold": h is not None and h < threshold} return out # ---------------------------------------------------------------- loop reconstruction & fleet health NETVM_ROOT = "/home/super/Projects/NetVM" FOLLOWUPS_FILE = os.path.join(NETVM_ROOT, "followups.json") DM_LOG_FILE = os.path.join(NETVM_ROOT, "dm-log.jsonl") JOB_LOG_FILE = os.path.join(NETVM_ROOT, "job-log.jsonl") def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list: """Reconstruct active and recent loops from followups.json and dm-log.jsonl. Returns a list of dicts: loop_id, agent, sender, target, purpose, state, sent_at, deadline, nudges_sent, nudges_allowed, escalate_to, tags, summary """ loops = {} # loop_id -> dict # 1. Load active/persisted followups if os.path.exists(FOLLOWUPS_FILE): try: with open(FOLLOWUPS_FILE, "r") as f: fdata = json.load(f) if isinstance(fdata, dict): for did, rec in fdata.items(): recipient = rec.get("recipient") st = rec.get("status", "pending") # Map status to loop state if st == "pending": lstate = "NUDGED" if rec.get("nudges_sent", 0) > 0 else "LANDED" elif st == "escalated": lstate = "ESCALATED" elif st in ("resolved", "closed"): lstate = "CLOSED" else: lstate = st.upper() loops[did] = { "loop_id": did, "agent": recipient, "sender": rec.get("sender", "super"), "target": rec.get("target", "main"), "thread_uuid": rec.get("thread_uuid"), "purpose": rec.get("route") or "followup", "state": lstate, "sent_at": rec.get("sent_at", ""), "deadline": rec.get("deadline", ""), "timeout_s": rec.get("timeout_s", 3600), "nudges_sent": rec.get("nudges_sent", 0), "nudges_allowed": rec.get("nudges_allowed", 2), "escalate_to": rec.get("escalate_to", "opm"), "tags": {"reply:expected": True}, "summary": f"Follow-up for DM {did}", "source": "followups.json", } except Exception: pass # 2. Extract tracked DMs from dm-log.jsonl answers_map = {} # id or ref -> answer event dm_events = [] failed_ids = set() # DM ids that failed to send - never delivered if os.path.exists(DM_LOG_FILE): try: with open(DM_LOG_FILE, "r") as f: for line in f: line = line.strip() if not line: continue try: entry = json.loads(line) ev_type = entry.get("type", "") msg = entry.get("msg") or "" # Check for answers/results/acks # Note: "verified" means delivered, NOT answered - do not include it here if "[RESULT" in msg or "[ACK" in msg or "acknowledged" in msg: sender = entry.get("agent", "") # Extract referenced DM id if present m_ref = re.search(r"\[ref:([a-f0-9-]+)\]", msg) or re.search(r"\[(?:ACK|RESULT)\s+([a-f0-9-]+)", msg) if m_ref: answers_map[m_ref.group(1)] = entry if entry.get("id"): answers_map[entry["id"]] = entry # Track failed sends - these were never delivered if ev_type in ("send_failed", "failed"): if entry.get("id"): failed_ids.add(entry.get("id")) tags = entry.get("tags") or {} if tags.get("reply:expected") or "[reply:expected]" in msg: dm_events.append(entry) except Exception: continue except Exception: pass # Merge dm_events into loops for entry in dm_events: did = entry.get("id") or entry.get("msg_id") or "" if not did: continue if did in loops: continue # Already have active record from followups.json if did in failed_ids: continue # Send failed - never delivered, don't count against agent health tags = entry.get("tags") or {} recipient = entry.get("to") or entry.get("recipient") or "" sender = entry.get("agent") or entry.get("sender") or "super" target = entry.get("target") or "main" sent_at = entry.get("ts") or entry.get("sent_at") or "" timeout_s = int(tags.get("reply:timeout") or 3600) nudges_max = int(tags.get("reply:nudges") or 2) esc = tags.get("reply:escalate") or "opm" purpose = tags.get("route") or "dm" # Determine state if did in answers_map or f"ref:{did}" in answers_map: lstate = "CLOSED" else: # Check age lstate = "LANDED" loops[did] = { "loop_id": did, "agent": recipient, "sender": sender, "target": target, "thread_uuid": tags.get("thread"), "purpose": purpose, "state": lstate, "sent_at": sent_at, "deadline": "", "timeout_s": timeout_s, "nudges_sent": 0, "nudges_allowed": nudges_max, "escalate_to": esc, "tags": tags, "summary": entry.get("msg", "")[:80], "source": "dm-log.jsonl", } # Filter and sort result = list(loops.values()) if agent: result = [l for l in result if l.get("agent") == agent or l.get("sender") == agent] if status_filter: sf = status_filter.lower() if sf == "active" or sf == "pending": result = [l for l in result if l.get("state") in ("FIRING", "LANDED", "SEEN", "NUDGED", "ESCALATED")] elif sf in ("closed", "resolved"): result = [l for l in result if l.get("state") in ("CLOSED", "ANSWERED")] else: result = [l for l in result if l.get("state", "").lower() == sf] # Sort newest first result.sort(key=lambda x: x.get("sent_at", ""), reverse=True) return result[:limit] def get_fleet_loop_health(threshold=None) -> dict: """Calculate fleet loop health per agent and overall verdict.""" if threshold is None: try: from variables import Variables threshold = Variables().get_float("loop_health_threshold") except Exception: threshold = 0.5 loops = reconstruct_loops(limit=100) raw = loop_health(loops, threshold=threshold) out = {} for ag in VALID_AGENTS: info = raw.get(ag, {"landed": 0, "answered": 0, "health": None, "below_threshold": False}) landed = info.get("landed", 0) answered = info.get("answered", 0) h = info.get("health") if landed == 0: status = "IDLE" elif h is not None and h >= threshold: status = "HEALTHY" else: status = "DEGRADED" out[ag] = { "agent": ag, "landed": landed, "answered": answered, "health": h, "health_pct": f"{int(h * 100)}%" if h is not None else "-", "threshold": threshold, "below_threshold": info.get("below_threshold", False), "status": status, } total_landed = sum(v["landed"] for v in out.values()) total_answered = sum(v["answered"] for v in out.values()) overall_h = (total_answered / total_landed) if total_landed > 0 else None return { "agents": out, "summary": { "total_landed": total_landed, "total_answered": total_answered, "overall_health": overall_h, "overall_health_pct": f"{int(overall_h * 100)}%" if overall_h is not None else "-", "threshold": threshold, "healthy": overall_h is None or overall_h >= threshold, } } def diagnose_breaks() -> list: """Diagnose break taxonomy across intrinsic loops and support services.""" import subprocess breaks = [] # 1. Scheduler / Timers check try: r = subprocess.run(["systemctl", "--user", "is-system-running"], capture_output=True, text=True, timeout=5) sys_state = r.stdout.strip() if sys_state in ("offline", "stopped"): breaks.append({ "type": "scheduler_death", "severity": "CRITICAL", "component": "systemd", "detail": f"Systemd user instance is {sys_state}", "remedy": "Restart systemd user session or start timer jobs manually." }) except Exception as e: breaks.append({ "type": "scheduler_death", "severity": "WARNING", "component": "systemd", "detail": f"Could not check systemd status: {e}", "remedy": "Verify systemctl --user is available." }) # 2. SSH key signature check key_path = os.path.expanduser("~/.ssh/id_ed25519") if not os.path.exists(key_path): breaks.append({ "type": "auth_rot", "severity": "CRITICAL", "component": "ssh-keys", "detail": f"SSH signing key {key_path} not found", "remedy": "Generate ed25519 key at ~/.ssh/id_ed25519 for cryptographically signed DMs." }) # 3. Active follow-up loops check active_loops = reconstruct_loops(limit=20, status_filter="pending") now_ts = time.time() for l in active_loops: nudges_sent = l.get("nudges_sent", 0) nudges_max = l.get("nudges_allowed", 2) if nudges_sent >= nudges_max and l.get("state") == "ESCALATED": breaks.append({ "type": "silent_agent", "severity": "WARNING", "loop_id": l["loop_id"], "agent": l["agent"], "detail": f"Agent {l['agent']} silent after {nudges_sent}/{nudges_max} nudges for DM {l['loop_id']}", "remedy": f"Check agent {l['agent']} browser tab with 'super fleet status' or nudge via 'super dm send'." }) # 4. Check for agents held up on approvals try: import approvals fleet_apps = approvals.check_fleet_approvals() for app in fleet_apps: if app.get("has_pending"): node = app["node"] is_trusted = app.get("is_trusted", False) ip = app.get("target") or app.get("ip") or "unknown target" breaks.append({ "type": "approval_blocked", "severity": "WARNING" if is_trusted else "CRITICAL", "component": f"node:{node}", "agent": node, "detail": f"Agent {node} is held up on browser approval for {ip}", "remedy": f"Run 'box approvals auto' or 'box approvals allow {node}'." }) for w in app.get("input_waits") or []: node = app["node"] breaks.append({ "type": "input_wait", "severity": "WARNING", "component": f"node:{node}", "agent": node, "detail": f"Agent {node} task '{w.get('task')}' is waiting: {w.get('status')} ({w.get('when')})", "remedy": f"Open {node}'s task and answer it, or 'box approvals check --node {node}'." }) except Exception: pass return breaks def remediate_breaks(dry_run=False) -> dict: """Progressively auto-remediate soft loop breakages while escalating hard breakages. Soft breakages (auto-healed): - Pending followups that have received an answer in dm-log.jsonl or chat-history are resolved. - Pending followups with expired deadlines and nudges remaining are re-armed and swept immediately. Hard breakages (escalated loudly): - silent_agent (nudges exhausted, no response) - auth_rot (missing SSH signing keys) - scheduler_death (systemd user session offline) """ import subprocess from datetime import datetime, timezone remediated = [] escalated = [] # 1. Check diagnosed hard breaks first breaks = diagnose_breaks() for b in breaks: if b.get("severity") in ("CRITICAL", "WARNING"): escalated.append(b) # 2. Check followups.json for soft break healing f_path = Path("/home/super/Projects/NetVM/followups.json") f_modified = False rearm_sweeper = False if f_path.exists(): try: with open(f_path, "r") as f: fdata = json.load(f) except Exception: fdata = {} now_iso = datetime.now(timezone.utc).isoformat() # Build answer map from reconstruct_loops loops = reconstruct_loops(limit=200) answered_dms = { l["loop_id"]: l for l in loops if l.get("state") in ("ANSWERED", "CLOSED") } for dm_id, rec in fdata.items(): if rec.get("status") == "pending": # Check if it was actually answered in logs if dm_id in answered_dms: remediated.append({ "action": "auto_resolve_answered", "loop_id": dm_id, "agent": rec.get("recipient"), "detail": f"Follow-up {dm_id} received reply in log but was pending in followups.json. Marked resolved." }) if not dry_run: rec["status"] = "resolved" rec["resolved_at"] = now_iso rec["resolved_note"] = "auto-healed: reply detected in dm-log" f_modified = True continue # Check if deadline expired and nudges remaining dl_str = rec.get("deadline", "") nudges_sent = rec.get("nudges_sent", 0) nudges_allowed = rec.get("nudges_allowed", 2) is_expired = False if dl_str: try: dl_dt = datetime.fromisoformat(dl_str.replace("Z", "+00:00")) if datetime.now(timezone.utc) > dl_dt: is_expired = True except Exception: pass if is_expired and nudges_sent < nudges_allowed: remediated.append({ "action": "rearm_expired_nudge", "loop_id": dm_id, "agent": rec.get("recipient"), "detail": f"Deadline expired for {dm_id} ({nudges_sent}/{nudges_allowed} nudges). Re-arming immediate sweep." }) if not dry_run: rec["deadline"] = now_iso f_modified = True rearm_sweeper = True if f_modified and not dry_run: tmp = f"{f_path}.tmp.{os.getpid()}" with open(tmp, "w") as f: json.dump(fdata, f, indent=2) os.replace(tmp, f_path) if rearm_sweeper and not dry_run: 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) except Exception: pass # Auto-remediate trusted approval blocks try: import approvals fleet_apps = approvals.check_fleet_approvals() for app in fleet_apps: if app.get("has_pending") and app.get("is_trusted") and app.get("status") != "KEY_APPROVAL": node = app["node"] if not dry_run: approvals.allow_node_approval(node, always=True, caller="loop-remediate") remediated.append({ "type": "approval_auto_allowed", "agent": node, "target": app.get("ip"), "action": f"Auto-approved trusted browser request on {node} ({app.get('ip')})" }) except Exception: pass # 3. Alert on hard breakages if any exist if escalated and not dry_run: job_log = Path("/home/super/Projects/NetVM/job-log.jsonl") alert_msg = f"[HARD_BREAK_ALERT] Detected {len(escalated)} unresolvable loop failure(s): " + "; ".join( f"{b.get('type')} ({b.get('severity')}): {b.get('detail')}" for b in escalated[:3] ) try: with open(job_log, "a", encoding="utf-8") as jf: jf.write(json.dumps({ "ts": datetime.now(timezone.utc).isoformat(), "type": "hard_break_alert", "escalated_count": len(escalated), "items": escalated, "summary": alert_msg[:280] }) + "\n") except Exception: pass # Dispatch DM alert to opm dm_py = Path("/home/super/Projects/NetVM/bin/dm.py") if dm_py.exists(): try: subprocess.run([ sys.executable, str(dm_py), "send", "--agent", "super", "--to", "opm", "--target", "main", alert_msg[:800] ], timeout=15, capture_output=True) except Exception: pass return { "ok": True, "remediated": remediated, "escalated": escalated, "dry_run": dry_run, "count": len(remediated), }