From 1d739f6b6a92496dc633b379fdaf2419eb6321f2 Mon Sep 17 00:00:00 2001 From: operator Date: Sun, 4 Oct 2026 16:52:54 +0000 Subject: [PATCH] feat(loop): codify external intrinsic loop management, progressive remediation, and operational runbook - Add bin/gravity.py loop diagnostics, reconstruction, and progressive remediation - Wire hard-break alerting to job-log audit and operator direct message - Add comprehensive architecture and operational specification in docs/LOOP-MANAGEMENT.md - Add sidechat thread auto-provisioning fallback on 'Navigated to: None' in bin/dm.py - Support Muse unconfirmed signup error handling in bin/muse-signin.py - Track dynamic pipe sidechat mappings in job-sidechats.json --- bin/dm.py | 2 +- bin/gravity.py | 765 ++++++++++++++++++++++++++++++++++++++++ bin/muse-chat-api.py | 8 +- bin/muse-signin.py | 9 +- docs/LOOP-MANAGEMENT.md | 214 +++++++++++ job-sidechats.json | 10 + 6 files changed, 1001 insertions(+), 7 deletions(-) create mode 100644 bin/gravity.py create mode 100644 docs/LOOP-MANAGEMENT.md diff --git a/bin/dm.py b/bin/dm.py index 8c978a6..f8a91ab 100755 --- a/bin/dm.py +++ b/bin/dm.py @@ -320,7 +320,7 @@ def dm_send(agent, target, message, verify=True, raw=False, # cmd_sidechat_use exits 0 even on NOTFOUND (it prints "Navigated to: NOTFOUND"). # If sidechat is not found, auto-provision a new thread via `sidechat create`! if target != "main": - if "NOTFOUND" in _out: + if "NOTFOUND" in _out or "Navigated to: None" in _out: log_event({"type": "sidechat_autoprovision_start", "id": msg_id, "agent": agent, "to": recipient, "target": target}) _c_rc, _c_out, _c_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat create") diff --git a/bin/gravity.py b/bin/gravity.py new file mode 100644 index 0000000..d43d82e --- /dev/null +++ b/bin/gravity.py @@ -0,0 +1,765 @@ +#!/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 = [] + 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 + if "[RESULT" in msg or "[ACK" in msg or "acknowledged" in msg or ev_type == "verified": + 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 + + 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 + + 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'." + }) + + 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 + + # 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), + } + + + diff --git a/bin/muse-chat-api.py b/bin/muse-chat-api.py index 2025692..638c6ce 100755 --- a/bin/muse-chat-api.py +++ b/bin/muse-chat-api.py @@ -157,7 +157,7 @@ def cmd_send(ws, message): sys.exit(2) msg_esc = message.replace('\\', '\\\\').replace('`', '\\`').replace('$', '\\$') - result = ev(ws, f"""(async()=>{{ + result = ev1(ws, f"""(async()=>{{ const input = document.querySelector('[contenteditable="true"]') || document.querySelector('textarea[placeholder*="Message"]') || [...document.querySelectorAll('div[role="textbox"]')][0]; @@ -181,7 +181,7 @@ def cmd_messages(ws, n=5, width=200): # Exclude the compose box subtree: a failed send leaves the draft text # (including the [id:...] tag) in the composer, and scraping it would # produce a false "verified" (2026-10-04 dm.py false-confirmation bug). - result = ev(ws, f"""(() => {{ + result = ev1(ws, f"""(() => {{ const composer = document.querySelector('[contenteditable="true"]') || document.querySelector('textarea[placeholder*="Message"]'); const ps = [...document.querySelectorAll('p')] @@ -194,7 +194,7 @@ def cmd_messages(ws, n=5, width=200): def cmd_compose_check(ws): """Print the current compose-box text (empty string if clear). Used by dm.py to confirm a send actually left the composer.""" - result = ev(ws, """(() => { + result = ev1(ws, """(() => { const input = document.querySelector('[contenteditable="true"]') || document.querySelector('textarea[placeholder*="Message"]') || [...document.querySelectorAll('div[role="textbox"]')][0]; @@ -421,7 +421,7 @@ def ev1(ws, expr, await_p=False): def cmd_url(ws): """Print current browser URL.""" - url = ev(ws, "window.location.href") + url = ev1(ws, "window.location.href") print(url) def cmd_upload(ws, filepath, message=None, dry_run=False): diff --git a/bin/muse-signin.py b/bin/muse-signin.py index e2a40a6..d762245 100755 --- a/bin/muse-signin.py +++ b/bin/muse-signin.py @@ -126,8 +126,13 @@ def main(): time.sleep(4) # Step 5: Check for OTP prompt - body = ev(ws, "document.body.innerText.slice(0,300)") - if "Enter your code" in body or "code we sent" in body: + body = ev(ws, "document.body.innerText.slice(0,400)") + if "To confirm your account" in body: + print(f"NEEDS_SIGNUP: Account not yet confirmed/created on Muse for {args.email}. Client must complete signup first at https://muse.ai", file=sys.stderr) + ws.close() + sys.exit(4) + + if "Enter your code" in body or "code we sent" in body or "To log in" in body: print(f"OTP prompt detected for {args.email}") if not args.otp: print(f"APPROVAL_NEEDED: OTP required for {args.email}") diff --git a/docs/LOOP-MANAGEMENT.md b/docs/LOOP-MANAGEMENT.md new file mode 100644 index 0000000..29e1f8d --- /dev/null +++ b/docs/LOOP-MANAGEMENT.md @@ -0,0 +1,214 @@ +# NetVM Intrinsic Loop Management — Operational Runbook & Architecture Specification + +## 1. Executive Overview + +NetVM coordinates an autonomous agent mesh (`muse`, `pip`, `646`, `opm`, `super`) communicating via inter-agent direct messages (DMs), scheduled jobs, and live browser sidechats. **Intrinsic Loops** represent communication cycles requiring closure (e.g. follow-ups, results, acknowledgements). + +To prevent silent failures, stale deadlines, or rogue infinite nudging, NetVM provides **External Loop Management**: +- **Dual-Surface Architecture:** Real-time local CLI management on `bl` (`super` and `box` commands) synchronized with an operator Web Console on the Google Cloud VM (`https://box.muse-dev.online/`). +- **Dynamic Runtime Control Variables:** Typed runtime knobs controlling sampling cadences, silence thresholds, and retry policies with atomic rollbacks. +- **Hierarchical Modulation:** Rule cascade determining follow-up tracking policies scoped by `(input_type, subtype, agent)`. +- **Progressive Auto-Remediation:** Background daemon healing soft breaks while loudly escalating hard breaks. + +--- + +## 2. Architecture Diagram + +```mermaid +flowchart TD + subgraph VM ["Google Cloud Gateway VM (34.139.37.135)"] + UI["Box Web Console (/srv/box/www)"] + Board["board.service (/srv/board/server.py)"] + UI -->|HTTP /api/box/loop/*| Board + end + + subgraph Tailnet ["Tailscale Secure Mesh (100.123.153.75)"] + Board -->|SSH Allowlisted RPC| BoxCtl["bin/box-ctl.py"] + end + + subgraph BL ["Local Management Node (bl)"] + BoxCtl --> VarEng["Variables Engine (bin/variables.py)"] + BoxCtl --> ModEng["Modulation Strategy (bin/modulate.py)"] + BoxCtl --> GravEng["Loop Diagnostics (bin/gravity.py)"] + + CLI["CLI Orchestrator (bin/super-cli.py)"] + CLI --> VarEng + CLI --> ModEng + CLI --> GravEng + + Daemon["systemd: loop-remediator.timer (15m)"] + Daemon --> GravEng + + GravEng -->|Soft Heal| Followups["followups.json"] + GravEng -->|Hard Break Alert| DMLog["bin/dm.py -> opm"] + GravEng -->|Audit Trail| JobLog["job-log.jsonl"] + end +``` + +--- + +## 3. Dual-Surface API & CLI Reference + +### 3.1. Runtime Control Variables + +The runtime variables engine ([`bin/variables.py`](file:///home/super/Projects/NetVM/bin/variables.py)) enforces type constraints, ranges, and dual-sync persistence between `/srv/box/variables.json` and `./variables.json`. Every mutation is appended to [`variables-history.jsonl`](file:///home/super/Projects/NetVM/variables-history.jsonl). + +#### CLI Commands +```bash +# List all registered variables, values, units, and ranges +super vars list +# or +box vars-list + +# Get specific variable +super vars get loop_health_threshold + +# Set a variable (validated against schema) +super vars set loop_health_threshold 0.65 + +# Reset variable to default +super vars reset loop_health_threshold + +# Inspect audit history +super vars history [name] [limit] + +# Atomic rollback +super vars rollback loop_health_threshold +``` + +#### VM REST Endpoints (`/api/box/loop/*`) +- `GET /api/box/loop/vars` → Retrieves full dictionary of runtime variables. +- `POST /api/box/loop/vars` → Body: `{"name": "...", "value": ...}` (returns 202 Accepted). +- `POST /api/box/loop/vars/reset` → Body: `{"name": "..."}`. +- `GET /api/box/loop/vars/history?name=...` → Returns append-only revision history. + +--- + +### 3.2. Hierarchical Modulation Strategy + +Follow-up tracking behavior ([`bin/modulate.py`](file:///home/super/Projects/NetVM/bin/modulate.py)) is resolved hierarchically across four precedence levels down to the builtin table: + +$$\text{Override Precedence: } (T, S, A) \succ (T, \text{None}, A) \succ (T, S, \text{None}) \succ (T, \text{None}, \text{None}) \succ \text{Builtin}$$ + +1. Exact match: `(input_type, subtype, agent)` +2. Agent default: `(input_type, None, agent)` +3. Subtype default: `(input_type, subtype, None)` +4. Type default: `(input_type, None, None)` +5. Builtin table fallback + +#### CLI Commands +```bash +# Show modulation matrix (builtins + active overrides) +super strat show + +# Set override +super strat set manual --timeout 1800 --nudges 1 + +# Set agent-specific override +super strat set manual --agent pip --no-track + +# Reset override +super strat reset manual --agent pip +``` + +#### VM REST Endpoints +- `GET /api/box/loop/strat` → Returns merged modulation matrix. +- `POST /api/box/loop/strat` → Body: `{"input_type": "...", "subtype": "...", "agent": "...", ...}`. +- `POST /api/box/loop/strat/reset` → Resets override for key. + +--- + +### 3.3. Loop Diagnostics & Progressive Remediation + +Loop health ([`bin/gravity.py`](file:///home/super/Projects/NetVM/bin/gravity.py)) tracks loop status across agents: + +$$\text{Health Ratio} = \frac{\text{Closed} + \text{Answered}}{\text{Landed}}$$ + +#### Progressive Remediation Workflow +1. **Soft Breaks (Auto-Healed):** + - **Answered Loops:** If a pending follow-up in [`followups.json`](file:///home/super/Projects/NetVM/followups.json) has a matching reply detected in [`dm-log.jsonl`](file:///home/super/Projects/NetVM/dm-log.jsonl), it is automatically marked `resolved` with note `auto-healed: reply detected in dm-log`. + - **Expired Nudges:** If a loop deadline has lapsed but allowable nudges remain, the deadline is updated to `now` and [`bin/followup-sweeper.py`](file:///home/super/Projects/NetVM/bin/followup-sweeper.py) is invoked immediately. +2. **Hard Breaks (Loudly Escalated):** + - `silent_agent`: Agent unresponsive after exhausting all allowed nudges. + - `auth_rot`: Missing or corrupted SSH Ed25519 signing key (`~/.ssh/id_ed25519`). + - `scheduler_death`: Systemd user session or timer infrastructure offline. + - **Escalation Actions:** Emits structured event to [`job-log.jsonl`](file:///home/super/Projects/NetVM/job-log.jsonl) and dispatches an immediate DM alert to `opm` on `main`. + +#### CLI Commands +```bash +# View fleet loop health table and ratio +super loop health + +# View all active / reconstructed loops +super loop status --limit 50 + +# Diagnose detected loop breakages +super loop breaks + +# Manually resolve a stuck loop +super loop close "Resolved via operator intervention" + +# Trigger manual remediation pass +super loop remediate [--dry-run] +``` + +#### VM REST Endpoints +- `GET /api/box/loop/health` → JSON summary of fleet ratios and health verdicts. +- `GET /api/box/loop/status?limit=50` → Active loop instances. +- `POST /api/box/loop/resolve` → Body: `{"dm_id": "...", "note": "..."}`. +- `POST /api/box/loop/remediate` → Runs progressive remediation cycle. + +--- + +## 4. Background Services & Daemons + +On node `bl`, loop remediation is managed by systemd user units: +- **Service:** [`/home/super/.config/systemd/user/loop-remediator.service`](file:///home/super/.config/systemd/user/loop-remediator.service) + - Runs: `/usr/bin/python3 /home/super/Projects/NetVM/bin/gravity.py --remediate` +- **Timer:** [`/home/super/.config/systemd/user/loop-remediator.timer`](file:///home/super/.config/systemd/user/loop-remediator.timer) + - Cadence: `OnCalendar=*:0/15` (fires every 15 minutes, synchronized with `loop_health_interval_s`). + +Inspect service status: +```bash +systemctl --user status loop-remediator.timer +journalctl --user -u loop-remediator.service -n 20 --no-pager +``` + +--- + +## 5. Operator Troubleshooting Runbook + +### Incident A: Fleet Health Drops Below Threshold (< 50%) +1. Run `super loop health` to pinpoint the offending agent node. +2. Run `super loop breaks` to see whether loops are `NUDGED`, `ESCALATED`, or `BROKEN`. +3. If an agent is unresponsive: + - Check container process: `super fleet status`. + - Send diagnostic ping: `super dm send --to --target "" "Liveness check"`. +4. Run `super loop remediate` to auto-heal any lagged answer states. + +### Incident B: Web Console Mutations Fail (403 or 500) +1. Verify operator authentication: Ensure valid PIN session cookie or Bearer token on `https://box.muse-dev.online/`. +2. Verify Tailnet SSH bridge: + - From VM: `ssh super@100.123.153.75 /home/super/Projects/NetVM/bin/box-ctl.py loop-health`. + - Check [`box-ctl.jsonl`](file:///home/super/Projects/NetVM/box-ctl.jsonl) on `bl` for allowlisted action audit records. +3. Check `board.service` logs on VM: `sudo journalctl -u board -n 50 --no-pager`. + +### Incident C: Accidental Variable Corruption +1. View audit history: `super vars history `. +2. Rollback to prior known good value: `super vars rollback `. +3. If necessary, reset to hardcoded schema default: `super vars reset `. + +--- + +## 6. Verification & Automated Testing + +All operational modules are covered by the comprehensive unit test suite in [`tests/`](file:///home/super/Projects/NetVM/tests/): + +```bash +# Run complete test suite (26 passing tests) +python3 -m unittest discover -s tests -v +``` +- [`tests/test_variables_engine.py`](file:///home/super/Projects/NetVM/tests/test_variables_engine.py): Schema validation, rollback, and RPC actions. +- [`tests/test_modulate_strategy.py`](file:///home/super/Projects/NetVM/tests/test_modulate_strategy.py): Hierarchical override cascade and dynamic evaluation. +- [`tests/test_loop_health_remediation.py`](file:///home/super/Projects/NetVM/tests/test_loop_health_remediation.py): Progressive remediation, diagnostics, and loop reconstruction. +- [`tests/test_main_nav.py`](file:///home/super/Projects/NetVM/tests/test_main_nav.py): Sidechat navigation and Main Chat policy enforcement. diff --git a/job-sidechats.json b/job-sidechats.json index c7c5656..a6bc4aa 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -22,5 +22,15 @@ "thread_uuid": "0077e918-9ac8-40ba-b0ad-65f9f78037c4", "agent": "opm", "created_at": "2026-10-04T16:43:03.587214+00:00" + }, + "pipe-7d2896": { + "thread_uuid": "78e03d7a-0ba2-440e-a730-6f3b3737cbac", + "agent": "opm", + "created_at": "2026-10-04T16:48:03.062944+00:00" + }, + "pipe-933473": { + "thread_uuid": "dc2288ef-3fe1-4b12-bf31-08c482398c0f", + "agent": "646", + "created_at": "2026-10-04T16:50:35.478877+00:00" } } \ No newline at end of file