adfcd2e602
- box passkey [show|fetch] (+ muse passkey): documents VM-only passkey (/srv/box/passkey.txt, fallback /etc/netvm/passkey.txt on 34.139.37.135), probes VM over SSH with graceful fallback; --json supported. No secrets on bl. - approvals: request_key_approval / check_node_key_request; KEY_APPROVAL status surfaced in `box approvals check`; allow/deny resolve + audit to box-ctl.jsonl; never auto-approved. New `box approvals request-key <node> --reason`. - box lookup (summary|fleet|threads|unread|approvals|key|docs) and docs-lookup engine with lookup_internal/ database (docs_internal symlink). - muse-tmux: non-TTY attach falls back to scrollback capture; prune NameError fix. - box/muse passthrough for tmux/muse/docs; thread list/view alias + prefix resolve. - Docs: AGENTS.md, AGENT-TOOLING.md, BOX-WEB-SURFACE-GUIDE.md, README. - Tests: key-approval + passkey tests; sync stale sidechat UUIDs and manifest name. - .gitignore runtime trackers (subagent-sessions, conversation-nudge-tracker).
823 lines
30 KiB
Python
823 lines
30 KiB
Python
#!/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 = []
|
|
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"):
|
|
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),
|
|
}
|
|
|
|
|
|
|