diff --git a/bin/tmux_auto_approver.py b/bin/tmux_auto_approver.py index 9d2df03..db7e712 100755 --- a/bin/tmux_auto_approver.py +++ b/bin/tmux_auto_approver.py @@ -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 = { diff --git a/bin/tmux_server_watchdog.py b/bin/tmux_server_watchdog.py new file mode 100755 index 0000000..913eba6 --- /dev/null +++ b/bin/tmux_server_watchdog.py @@ -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()) diff --git a/systemd/tmux-server-watchdog.service b/systemd/tmux-server-watchdog.service new file mode 100644 index 0000000..aa51b94 --- /dev/null +++ b/systemd/tmux-server-watchdog.service @@ -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 diff --git a/systemd/tmux-server-watchdog.timer b/systemd/tmux-server-watchdog.timer new file mode 100644 index 0000000..1765196 --- /dev/null +++ b/systemd/tmux-server-watchdog.timer @@ -0,0 +1,10 @@ +[Unit] +Description=Run tmux server death-capture every minute + +[Timer] +OnBootSec=1min +OnUnitActiveSec=1min +Persistent=true + +[Install] +WantedBy=timers.target diff --git a/tests/test_box_runtime.py b/tests/test_box_runtime.py index 894424b..3afee99 100644 --- a/tests/test_box_runtime.py +++ b/tests/test_box_runtime.py @@ -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): diff --git a/tests/test_tmux_auto_approver.py b/tests/test_tmux_auto_approver.py index a600c0a..df7ce82 100644 --- a/tests/test_tmux_auto_approver.py +++ b/tests/test_tmux_auto_approver.py @@ -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() diff --git a/tests/test_tmux_server_watchdog.py b/tests/test_tmux_server_watchdog.py new file mode 100644 index 0000000..36fc0ed --- /dev/null +++ b/tests/test_tmux_server_watchdog.py @@ -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()