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
This commit is contained in:
+765
@@ -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 `<key>: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 <uuid>` 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),
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user