feat(tmux): add server death watchdog daemon and multi-socket approver enhancements

This commit is contained in:
operator
2026-10-07 01:50:06 +00:00
parent f2527f183c
commit 90f4ef661a
7 changed files with 551 additions and 5 deletions
+42 -5
View File
@@ -285,13 +285,44 @@ def run_tmux_cmd(socket_path: str, *args: str, timeout: float = 3.0) -> Tuple[in
def capture_pane_text(socket_path: str, pane_id: str, lines: int = 30) -> str:
"""Capture recent lines from a pane."""
rc, out, _ = run_tmux_cmd(socket_path, "capture-pane", "-p", "-t", pane_id, "-S", f"-{lines}")
"""Capture recent lines from a pane, joining wrapped rows.
-J joins physical wrapped lines into logical lines so matching is
width-independent: narrow panes wrap the same dialog onto more
rows, which otherwise breaks cue/option regexes.
"""
rc, out, _ = run_tmux_cmd(socket_path, "capture-pane", "-p", "-J",
"-t", pane_id, "-S", f"-{lines}")
if rc == 0:
return out
return ""
MUSE_COMMAND_HINTS = ("muse-bin", "muse-code")
def should_defer_to_muse_watcher(socket_path: str, pane_id: str,
current_command: str) -> bool:
"""True when a per-pane muse watcher owns this pane.
Single-owner rule: muse_choice_watcher is authoritative for muse
panes (stability + re-verify + once-per-prompt + decided-block
guard). When its daemon is alive for this socket:pane, tmux must
skip the pane entirely, or both daemons answer the same prompt
within the same second ('11' + stray keys, observed live). Never
raises: import or liveness failures mean no owner, handle here.
"""
try:
cmd = current_command or ""
if not any(h in cmd for h in MUSE_COMMAND_HINTS):
return False
import muse_choice_watcher as mcw
alive = getattr(mcw, "watcher_alive", mcw.is_running)
return alive(socket_path, pane_id) is not None
except Exception:
return False
def gather_tmux_tally(state: Optional[AutoApproverState] = None) -> TmuxWorkerTally:
"""Scan all sockets and build a comprehensive tally of tmux workers."""
if state is None:
@@ -497,6 +528,9 @@ class AutoApproverRunner:
for p in tally.panes:
if not p.auto_approve:
continue
if should_defer_to_muse_watcher(p.socket, p.pane_id,
p.current_command):
continue
text = capture_pane_text(p.socket, p.pane_id, lines=30)
if not text:
@@ -516,9 +550,12 @@ class AutoApproverRunner:
continue
if verdict.matched and verdict.key:
# Deduplicate identical prompt to avoid infinite loop
# Deduplicate identical prompt to avoid infinite loop.
# Keyed by socket:pane: bare pane ids repeat on every
# tmux socket, so %1 on pip must not suppress %1 on opm.
sig = hashlib.sha1(f"{verdict.rule_id}:{verdict.excerpt}".encode()).hexdigest()
last_time, last_sig = self.recent_signatures.get(p.pane_id, (0, ""))
dedup_key = "%s:%s" % (p.socket, p.pane_id)
last_time, last_sig = self.recent_signatures.get(dedup_key, (0, ""))
if last_sig == sig and (time.time() - last_time) < 15.0:
continue # already handled recently
@@ -544,7 +581,7 @@ class AutoApproverRunner:
success = True # dry-run simulated
now = time.time()
self.recent_signatures[p.pane_id] = (now, sig)
self.recent_signatures[dedup_key] = (now, sig)
self.approval_counts.append(now)
event = {
+189
View File
@@ -0,0 +1,189 @@
#!/usr/bin/env python3
"""tmux_server_watchdog.py — Death-capture for tmux servers.
Runs on a 1-minute systemd timer. Remembers each known socket's server
pid; when a server dies or its pid changes without a witnessed death,
appends a forensics bundle (dmesg OOM/kill lines, memory, uptime,
journal tail) to logs/tmux-server-deaths.jsonl so the next "tmux
crashed" leaves evidence instead of a mystery.
Read-only against tmux itself: one `display-message -p` probe per
socket. Never raises; a watchdog must not need its own watchdog.
"""
import json
import os
import subprocess
import sys
from datetime import datetime, timezone
BIN_DIR = os.path.dirname(os.path.abspath(__file__))
REPO_ROOT = os.path.dirname(BIN_DIR)
sys.path.insert(0, BIN_DIR)
try:
from muse_choice_watcher import KNOWN_SOCKETS
except Exception:
KNOWN_SOCKETS = ["/tmp/tmux-1000/default"]
STATE_FILE = os.path.join(REPO_ROOT, ".state", "tmux-servers.json")
DEATH_LOG = os.path.join(REPO_ROOT, "logs", "tmux-server-deaths.jsonl")
def _now():
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def _run(cmd, timeout=10):
try:
r = subprocess.run(cmd, capture_output=True, text=True,
timeout=timeout)
return r.returncode, (r.stdout or "").strip()
except Exception as e:
return -1, "exec failed: %r" % (e,)
def probe(socket_path):
"""Server pid for a socket, or None when unreachable."""
rc, out = _run(["tmux", "-S", socket_path, "display-message",
"-p", "#{pid}"], timeout=10)
if rc != 0:
return None
try:
return int(out.strip().split()[0])
except (ValueError, IndexError):
return None
def collect_forensics(socket_path, last_pid):
"""Best-effort death evidence. Dict of strings, never raises."""
ev = {"ts": _now(), "socket": socket_path, "last_pid": last_pid}
rc, dmesg = _run(["dmesg"], timeout=10)
if rc != 0:
ev["dmesg"] = "unavailable: %s" % dmesg[:200]
else:
hits = [ln for ln in dmesg.split("\n")
if any(k in ln.lower() for k in
("oom", "killed process", "segfault", "tmux"))]
ev["dmesg_hits"] = hits[-15:]
_, ev["memory"] = _run(["free", "-m"], timeout=10)
_, ev["uptime"] = _run(["uptime"], timeout=10)
rc, journal = _run(["journalctl", "--user", "-n", "50"], timeout=10)
if rc == 0:
ev["journal_tmux"] = [ln for ln in journal.split("\n")
if "tmux" in ln.lower()][-10:]
else:
ev["journal_tmux"] = []
return ev
def read_state(path=None):
try:
with open(path or STATE_FILE) as f:
data = json.load(f)
return data if isinstance(data, dict) else {}
except Exception:
return {}
def write_state(state, path=None):
path = path or STATE_FILE
try:
parent = os.path.dirname(path)
if parent:
os.makedirs(parent, exist_ok=True)
tmp = "%s.tmp.%d" % (path, os.getpid())
with open(tmp, "w") as f:
json.dump(state, f, indent=1)
os.replace(tmp, path)
except Exception:
pass
def append_death(ev, path=None):
path = path or DEATH_LOG
try:
parent = os.path.dirname(path)
if parent:
os.makedirs(parent, exist_ok=True)
with open(path, "a") as f:
f.write(json.dumps(ev) + "\n")
except Exception:
pass
def evaluate(previous, probed):
"""Pure transition logic: (prev_state, {sock: pid|None}) ->
(new_state, events). Events: death | restart | started."""
new_state, events = {}, []
for sock, pid in sorted(probed.items()):
prev = (previous.get(sock) or {})
prev_pid = prev.get("pid")
if pid is None:
new_state[sock] = {"pid": None, "died": _now(),
"last_pid": prev_pid}
if prev_pid:
events.append({"type": "death", "socket": sock,
"last_pid": prev_pid})
else:
new_state[sock] = {"pid": pid, "since": _now()}
if prev_pid and prev_pid != pid:
# Changed with no witnessed death: restart inside one
# tick gap (or pid recycled under us). Treat as a
# restart, still worth a forensics note.
events.append({"type": "restart", "socket": sock,
"old_pid": prev_pid, "pid": pid})
elif not prev_pid and prev.get("died"):
events.append({"type": "started", "socket": sock,
"pid": pid})
elif not prev_pid and not prev:
events.append({"type": "started", "socket": sock,
"pid": pid})
return new_state, events
def check(sockets=None, dry_run=False):
"""Probe, transition state, log deaths. Returns summary dict."""
probed = {s: probe(s) for s in (sockets or KNOWN_SOCKETS)}
previous = read_state()
new_state, events = evaluate(previous, probed)
for ev in events:
if ev["type"] == "death":
bundle = collect_forensics(ev["socket"], ev["last_pid"])
bundle["event"] = "death"
if not dry_run:
append_death(bundle)
ev["forensics"] = bundle
elif ev["type"] == "restart":
bundle = collect_forensics(ev["socket"], ev["old_pid"])
bundle["event"] = "restart-gap-missed"
if not dry_run:
append_death(bundle)
ev["forensics"] = bundle
if not dry_run:
write_state(new_state)
return {"probed": probed, "events": events, "dry_run": dry_run}
def main(argv=None):
import argparse
ap = argparse.ArgumentParser(description="tmux server death-capture")
ap.add_argument("--sockets", nargs="*", default=None)
ap.add_argument("--dry-run", action="store_true")
ap.add_argument("--json", action="store_true")
args = ap.parse_args(argv)
try:
res = check(sockets=args.sockets, dry_run=args.dry_run)
except Exception as e:
print("watchdog failed: %r" % (e,), file=sys.stderr)
return 1
if args.json or args.dry_run:
print(json.dumps(res, indent=1, default=str))
else:
for ev in res["events"]:
print("%s: %s" % (ev["type"], ev["socket"]))
return 0
if __name__ == "__main__":
sys.exit(main())
+10
View File
@@ -0,0 +1,10 @@
[Unit]
Description=NetVM tmux server death-capture watchdog
After=network.target
[Service]
Type=oneshot
ExecStart=/usr/bin/python3 /home/super/Projects/NetVM/bin/tmux_server_watchdog.py
WorkingDirectory=/home/super/Projects/NetVM
StandardOutput=journal
StandardError=journal
+10
View File
@@ -0,0 +1,10 @@
[Unit]
Description=Run tmux server death-capture every minute
[Timer]
OnBootSec=1min
OnUnitActiveSec=1min
Persistent=true
[Install]
WantedBy=timers.target
+65
View File
@@ -131,6 +131,44 @@ class TestApprovalFlags(unittest.TestCase):
p = w.muse_approval_flags([])
self.assertFalse(p["auto_approve"])
def test_permission_profile(self):
p = w.muse_approval_flags(
["muse", "--permission-profile", ":unrestricted"])
self.assertEqual(p["profile"], ":unrestricted")
self.assertEqual(p["mode"], ":unrestricted")
self.assertIn("permission-profile=:unrestricted", p["flags"])
self.assertFalse(p["bypass"])
def test_permission_profile_equals(self):
p = w.muse_approval_flags(
["muse", "--permission-profile=:read-only"])
self.assertEqual(p["profile"], ":read-only")
self.assertEqual(p["mode"], ":read-only")
def test_yolo_mode_and_bypass(self):
p = w.muse_approval_flags(["muse", "--yolo"])
self.assertEqual(p["mode"], "yolo")
self.assertTrue(p["bypass"])
self.assertIsNone(p["profile"])
def test_bypass_without_profile(self):
p = w.muse_approval_flags(["muse", "--disable-approval"])
self.assertTrue(p["bypass"])
self.assertEqual(p["mode"], "default")
def test_default_mode(self):
p = w.muse_approval_flags(["muse"])
self.assertEqual(p["mode"], "default")
self.assertFalse(p["bypass"])
self.assertIsNone(p["profile"])
def test_sandbox_and_trust_flags(self):
p = w.muse_approval_flags(
["muse", "--disable-sandbox", "--trust-workspace"])
self.assertIn("disable-sandbox", p["flags"])
self.assertIn("trust-workspace", p["flags"])
self.assertFalse(p["bypass"])
def _tmux_result(returncode=0, stdout="", stderr=""):
r = mock.Mock()
@@ -238,6 +276,33 @@ class TestRuntimeRows(unittest.TestCase):
self.assertEqual(rows[0]["height"], 7)
self.assertTrue(rows[1]["squeezed"])
def test_rows_carry_permission_posture(self):
patches = self._patched(
captures={"%37": STATE_OPEN, "%38": STATE_SHELL},
children={2880158: [2881158]},
cmdlines={2881158: ["/home/super/.local/bin/muse-bin-1.4.3",
"--permission-profile", ":unrestricted",
"--disable-sandbox"]})
with patches[0], patches[1], patches[2], patches[3], patches[4]:
rows = w.runtime_rows("/tmp/sock")
self.assertEqual(rows[0]["permission_mode"], ":unrestricted")
self.assertEqual(rows[0]["permission_profile"], ":unrestricted")
self.assertFalse(rows[0]["permission_bypass"])
self.assertIn("permission-profile=:unrestricted",
rows[0]["approval_flags"])
self.assertIsNone(rows[1]["permission_mode"])
def test_narrow_but_tall_not_squeezed(self):
# Empirical sizing: 35-wide panes answer cleanly when tall
# enough; only short height drops dialog text. Width minimum
# must not flag working panes as squeezed.
listing = "muse\t1\t%37\tmuse-bin-1.4\t2880158\t35\t35\n"
patches = self._patched(
tmux_stdout=listing, captures={"%37": STATE_OPEN})
with patches[0], patches[1], patches[2], patches[3], patches[4]:
rows = w.runtime_rows("/tmp/sock")
self.assertFalse(rows[0]["squeezed"])
class TestNodeFromSession(unittest.TestCase):
def test_conforming_sessions(self):
+123
View File
@@ -220,5 +220,128 @@ class TestTallyGathering(unittest.TestCase):
self.assertEqual(tally.panes[1].agent_node, "dev")
class TestMuseDeferral(unittest.TestCase):
"""tmux approver must defer muse panes owned by muse_choice_watcher.
Live double-answer regression: both daemons answered the same
Would-you-like prompt within the same second (box-ctl + tmux audit
overlap on %40/%0/%2), producing '11' + stray keys in the input box.
When a per-pane muse watcher is alive, tmux must skip the pane.
"""
def setUp(self):
self.temp_dir = tempfile.TemporaryDirectory()
self.orig_state = tmux_auto_approver.STATE_FILE
self.orig_audit = tmux_auto_approver.AUDIT_LOG_FILE
tmux_auto_approver.STATE_FILE = Path(self.temp_dir.name) / "test_state.json"
tmux_auto_approver.AUDIT_LOG_FILE = Path(self.temp_dir.name) / "test_audit.jsonl"
def tearDown(self):
tmux_auto_approver.STATE_FILE = self.orig_state
tmux_auto_approver.AUDIT_LOG_FILE = self.orig_audit
self.temp_dir.cleanup()
def _muse_pane(self, pane_id="%37", socket="/tmp/tmux-1000/default"):
return TmuxPaneInfo(
socket=socket, session="muse", window_idx=0, pane_id=pane_id,
pane_pid=12345, current_command="muse-bin", active=True,
attached=True, title="muse terminal", agent_node="muse",
auto_approve=True,
)
def _tally(self, panes):
return TmuxWorkerTally(
total_sockets=1, total_sessions=1, total_panes=len(panes),
active_workers=len(panes), by_agent={}, panes=panes,
)
@patch("tmux_auto_approver.run_tmux_cmd")
@patch("tmux_auto_approver.capture_pane_text")
@patch("tmux_auto_approver.gather_tmux_tally")
def test_muse_pane_skipped_when_watcher_alive(
self, mock_tally, mock_capture, mock_tmux_cmd):
import muse_choice_watcher as mcw
mock_tally.return_value = self._tally([self._muse_pane()])
mock_capture.return_value = (
"Would you like to run the following?\n› 1. Yes, proceed (y)")
with patch.object(mcw, "is_running", return_value=99999):
runner = AutoApproverRunner(dry_run=True)
results = runner.run_once()
self.assertEqual(results, [])
mock_tmux_cmd.assert_not_called()
@patch("tmux_auto_approver.run_tmux_cmd")
@patch("tmux_auto_approver.capture_pane_text")
@patch("tmux_auto_approver.gather_tmux_tally")
def test_non_muse_pane_still_approved(
self, mock_tally, mock_capture, mock_tmux_cmd):
pane = TmuxPaneInfo(
socket="/tmp/tmux-pip.sock", session="worker", window_idx=0,
pane_id="%1", pane_pid=999, current_command="agent-worker",
active=True, attached=True, title="w", agent_node="pip",
auto_approve=True,
)
mock_tally.return_value = self._tally([pane])
mock_capture.return_value = (
"Would you like to run the following?\n› 1. Yes, proceed (y)")
runner = AutoApproverRunner(dry_run=True)
results = runner.run_once()
self.assertEqual(len(results), 1)
self.assertEqual(results[0]["action"], "DRY_RUN_MATCH")
class TestDedupSocketScoped(unittest.TestCase):
"""Dedup must be keyed by socket:pane, not bare pane id.
Same %N exists on every tmux socket; bare-pane dedup suppresses a
real prompt on socket B because socket A saw one (the %0-on-two-
sockets collision, tmux-side).
"""
def setUp(self):
self.temp_dir = tempfile.TemporaryDirectory()
self.orig_state = tmux_auto_approver.STATE_FILE
self.orig_audit = tmux_auto_approver.AUDIT_LOG_FILE
tmux_auto_approver.STATE_FILE = Path(self.temp_dir.name) / "test_state.json"
tmux_auto_approver.AUDIT_LOG_FILE = Path(self.temp_dir.name) / "test_audit.jsonl"
def tearDown(self):
tmux_auto_approver.STATE_FILE = self.orig_state
tmux_auto_approver.AUDIT_LOG_FILE = self.orig_audit
self.temp_dir.cleanup()
@patch("tmux_auto_approver.run_tmux_cmd")
@patch("tmux_auto_approver.capture_pane_text")
@patch("tmux_auto_approver.gather_tmux_tally")
def test_same_pane_id_on_two_sockets_both_approved(
self, mock_tally, mock_capture, mock_tmux_cmd):
def mk(sock):
return TmuxPaneInfo(
socket=sock, session="w", window_idx=0, pane_id="%1",
pane_pid=999, current_command="agent-worker", active=True,
attached=True, title="w", agent_node="pip",
auto_approve=True,
)
mock_tally.return_value = TmuxWorkerTally(
total_sockets=2, total_sessions=2, total_panes=2,
active_workers=2, by_agent={},
panes=[mk("/tmp/tmux-pip.sock"), mk("/tmp/tmux-opm.sock")],
)
mock_capture.return_value = (
"Would you like to run the following?\n› 1. Yes, proceed (y)")
runner = AutoApproverRunner(dry_run=True)
results = runner.run_once()
self.assertEqual(len(results), 2)
class TestCaptureJoinWrapped(unittest.TestCase):
def test_capture_joins_wrapped_lines(self):
with patch("tmux_auto_approver.run_tmux_cmd",
return_value=(0, "ok", "")) as m:
tmux_auto_approver.capture_pane_text("/tmp/s", "%1", lines=30)
args = m.call_args[0]
self.assertIn("-J", args)
if __name__ == "__main__":
unittest.main()
+112
View File
@@ -0,0 +1,112 @@
#!/usr/bin/env python3
"""test_tmux_server_watchdog.py — Death-capture transition logic.
Covers: steady state is quiet, pid change yields restart, alive->dead
yields a death bundle, dead->alive yields started, corrupt state file
is tolerated.
"""
import json
import sys
import unittest
from pathlib import Path
from unittest import mock
REPO_ROOT = Path("/home/super/Projects/NetVM")
BIN_DIR = REPO_ROOT / "bin"
sys.path.insert(0, str(BIN_DIR))
import tmux_server_watchdog as w
class TestEvaluate(unittest.TestCase):
def test_steady_alive_is_quiet(self):
prev = {"/s": {"pid": 100, "since": "t"}}
new, events = w.evaluate(prev, {"/s": 100})
self.assertEqual(events, [])
self.assertEqual(new["/s"]["pid"], 100)
def test_pid_change_is_restart(self):
prev = {"/s": {"pid": 100, "since": "t"}}
new, events = w.evaluate(prev, {"/s": 200})
self.assertEqual(len(events), 1)
self.assertEqual(events[0]["type"], "restart")
self.assertEqual(events[0]["old_pid"], 100)
def test_alive_to_dead_is_death(self):
prev = {"/s": {"pid": 100, "since": "t"}}
new, events = w.evaluate(prev, {"/s": None})
self.assertEqual(len(events), 1)
self.assertEqual(events[0]["type"], "death")
self.assertIsNone(new["/s"]["pid"])
def test_dead_stays_dead_is_quiet(self):
prev = {"/s": {"pid": None, "died": "t", "last_pid": 100}}
_, events = w.evaluate(prev, {"/s": None})
self.assertEqual(events, [])
def test_dead_to_alive_is_started(self):
prev = {"/s": {"pid": None, "died": "t", "last_pid": 100}}
_, events = w.evaluate(prev, {"/s": 300})
self.assertEqual(len(events), 1)
self.assertEqual(events[0]["type"], "started")
def test_unknown_socket_first_seen_is_started(self):
_, events = w.evaluate({}, {"/s": 300})
self.assertEqual(events[0]["type"], "started")
class TestCheck(unittest.TestCase):
def _iso(self, td):
state = str(td / "servers.json")
log = str(td / "deaths.jsonl")
p1 = mock.patch.object(w, "STATE_FILE", state)
p2 = mock.patch.object(w, "DEATH_LOG", log)
return p1, p2, state, log
def test_death_writes_bundle(self):
import tempfile
with tempfile.TemporaryDirectory() as td:
td = Path(td)
p1, p2, state, log = self._iso(td)
with open(state, "w") as f:
json.dump({"/s": {"pid": 100, "since": "t"}}, f)
with p1, p2, \
mock.patch.object(w, "probe", return_value=None), \
mock.patch.object(w, "collect_forensics",
return_value={"ts": "t", "socket": "/s",
"last_pid": 100}):
res = w.check(sockets=["/s"])
self.assertEqual(res["events"][0]["type"], "death")
bundle = json.loads(open(log).read().strip())
self.assertEqual(bundle["event"], "death")
self.assertEqual(bundle["last_pid"], 100)
self.assertIsNone(json.load(open(state))["/s"]["pid"])
def test_dry_run_writes_nothing(self):
import tempfile
with tempfile.TemporaryDirectory() as td:
td = Path(td)
p1, p2, state, log = self._iso(td)
with p1, p2, \
mock.patch.object(w, "probe", return_value=100):
res = w.check(sockets=["/s"], dry_run=True)
self.assertEqual(res["events"][0]["type"], "started")
self.assertFalse(Path(state).exists())
self.assertFalse(Path(log).exists())
def test_corrupt_state_tolerated(self):
import tempfile
with tempfile.TemporaryDirectory() as td:
td = Path(td)
p1, p2, state, _ = self._iso(td)
with open(state, "w") as f:
f.write("{{{nope")
with p1, p2, \
mock.patch.object(w, "probe", return_value=100):
res = w.check(sockets=["/s"])
self.assertEqual(res["events"][0]["type"], "started")
if __name__ == "__main__":
unittest.main()