2067 lines
75 KiB
Python
Executable File
2067 lines
75 KiB
Python
Executable File
#!/usr/bin/env python3
|
||
"""muse_choice_watcher.py — Per-pane daemon auto-approving Muse prompts.
|
||
|
||
Kinds (prompt -> key): native Muse approval menu -> "1", agent
|
||
interview menu -> "1", explicit-phrase request -> captured TOKEN,
|
||
lettered A/B/C choice -> "A", numbered (1)/(2) menu -> "1", y/n
|
||
line-end prompt -> "y".
|
||
Each prompt is answered at most once
|
||
(after stability + re-verify), with an hourly cap as backstop.
|
||
|
||
Design notes (from the agy watcher post-mortem):
|
||
* State is namespaced by (socket, pane): pidfiles and logs embed a socket
|
||
slug, so identical pane ids on different tmux sockets (e.g. %0 on both
|
||
`default` and `lte`) never collide.
|
||
* Logging never kills the loop: every log call and every poll iteration is
|
||
exception-guarded, so auto-answers continue even when logging or tmux
|
||
hiccups. Log rotation is best-effort.
|
||
* Answers are once-per-prompt: a prompt signature must be stable across N
|
||
polls, is answered at most once, and re-answering requires the prompt to
|
||
disappear and a new signature to appear. An hourly cap bounds runaways.
|
||
|
||
Usage:
|
||
muse_choice_watcher.py start --socket /tmp/tmux-1000/default --pane %37
|
||
muse_choice_watcher.py start-all [--dry-run]
|
||
muse_choice_watcher.py status
|
||
muse_choice_watcher.py stop --socket ... --pane ...
|
||
muse_choice_watcher.py match < pane.txt # print JSON verdict for tuning
|
||
"""
|
||
|
||
import argparse
|
||
import fcntl
|
||
import hashlib
|
||
import json
|
||
import os
|
||
import re
|
||
import signal
|
||
import subprocess
|
||
import sys
|
||
import time
|
||
from datetime import datetime, timezone
|
||
|
||
STATE_DIR = "/tmp"
|
||
FILE_PREFIX = "muse-choice-watcher"
|
||
REPO_ROOT = "/home/super/Projects/NetVM"
|
||
DESIRED_STATE_FILE = os.path.join(REPO_ROOT, ".state", "muse-choices.json")
|
||
CTL_LOG = os.path.join(REPO_ROOT, "box-ctl.jsonl")
|
||
KNOWN_SOCKETS = [
|
||
"/tmp/tmux-1000/default",
|
||
"/tmp/tmux-1000/lte",
|
||
"/tmp/tmux-muse.sock",
|
||
]
|
||
|
||
POLL_INTERVAL = 0.5
|
||
STABILITY_POLLS = 2
|
||
MAX_ANSWERS_PER_HOUR = 20
|
||
ANSWERED_TTL_SECONDS = 600 # identical prompt back after 10m => stuck, allow one recovery answer
|
||
CLAIM_TTL_SECONDS = 30 # concurrent-claim window: bounds wedge if winner dies pre-send
|
||
RULES_FILE = os.path.join(REPO_ROOT, "muse-choices-rules.json")
|
||
HOLD_WINDOW_SECONDS = 120 # D3: short hold window, then expire to approve
|
||
NEGATIVE_KEYS = {"muse-approval": "2", "yn": "n"} # D2 deny keys
|
||
QUESTION_KINDS = frozenset({"interview", "letter", "numbered", "explicit-phrase"})
|
||
_RULES_CACHE = {"key": None, "rules": []}
|
||
LOG_MAX_BYTES = 1_000_000
|
||
HEARTBEAT_SECONDS = 60
|
||
TAIL_WINDOW = 25 # only prompts in the last N lines count (no stale scrollback)
|
||
OPTION_SPAN = 12 # A..C option lines must fit within N lines (wrapped lines ok)
|
||
MAX_OPTION_LEN = 160 # option lines longer than this are ignored (prose guard)
|
||
|
||
# A. text / B) text / C: text / C - text (single capital letter + delimiter)
|
||
OPTION_RE = re.compile(r"^\s*([A-Z])\s*[.\)\-:]\s+\S")
|
||
# Cue that the lettered list is a choice awaiting reply (not prose).
|
||
CUE_RE = re.compile(
|
||
r"(choose|choice|choices|select|option|options|reply|feedback|"
|
||
r"which one|pick one|enter\s+[A-Z]\b|type\s+[A-Z]\b|press\s+[A-Z]\b|"
|
||
r"A\s*/\s*B|A\s*,\s*B|A-C|A,B,C|\(A\))",
|
||
re.IGNORECASE,
|
||
)
|
||
QUESTION_RE = re.compile(r"\?\s*$")
|
||
|
||
# y/n token must END the line: prompts awaiting input put the options last
|
||
# ("Proceed? (y/n)", "Overwrite? [y/N]"). Mid-line mentions are prose.
|
||
YN_RE = re.compile(r"([yY]/[nN]|\[[yY]/[nN]\])\s*[\]:)>]?\s*$")
|
||
YN_WINDOW = 6 # y/n prompt must sit in the last N content lines
|
||
# Numbered menus: (1) text / (2) text with a selection cue nearby.
|
||
NUMBERED_OPT_RE = re.compile(r"^\s*\((\d+)\)\s+\S")
|
||
NUMBERED_CUE_RE = re.compile(
|
||
r"(Option:|Selection:|choose|choice|select an? |enter (a )?number|pick a number)",
|
||
re.IGNORECASE,
|
||
)
|
||
# Native Muse approval dialog patterns live in MUSE_APPROVAL_*_RE below
|
||
# (single definition; the dialog shape is asserted by TestMuseApproval).# Native Muse TUI approval menu (the live auto-approve target):
|
||
# Would you like to run the following
|
||
# $ <command echo, often wrapped across lines>
|
||
# > 1. Yes, proceed (y)
|
||
# 2. No, and tell Muse Code what to ...
|
||
# The cue carries no trailing "?", and the cursor may render as an ascii
|
||
# ">" or the single right-pointing angle quote (U+203A, escaped below).
|
||
MUSE_APPROVAL_CUE_RE = re.compile(r"^\s*Would you like to\b", re.IGNORECASE)
|
||
MUSE_APPROVAL_YES_RE = re.compile("^\\s*[>\\u203a]?\\s*1\\.\\s+Yes\\b")
|
||
MUSE_APPROVAL_NO_RE = re.compile("^\\s*[>\\u203a]?\\s*2\\.\\s+No\\b")
|
||
# A decided block ("approval decision accepted/denied") must never match:
|
||
# answering it would inject stray keys after the decision already landed.
|
||
MUSE_APPROVAL_DECIDED_RE = re.compile(r"approval decision", re.IGNORECASE)
|
||
MUSE_APPROVAL_SPAN = 16 # cue..options span (wrapped command echo ok)
|
||
# Collapsed approval: long commands collapse, hiding the 1.Yes/2.No pair:
|
||
# Would you like to run the following
|
||
# $ <first rows...>
|
||
# ... N command rows omitted
|
||
# ctrl+o view full command
|
||
MUSE_COLLAPSED_ROWS_RE = re.compile(r"rows?\s+omitted", re.IGNORECASE)
|
||
MUSE_COLLAPSED_KEY_RE = re.compile(r"ctrl\s*\+\s*o\b", re.IGNORECASE)
|
||
|
||
# Agent interview UI (request_user_input rendered in-pane; observed live):
|
||
# <question...?>
|
||
# > 1. First option (Recommended) <desc...>
|
||
# 2. Second option <desc...>
|
||
# 3. Both ...
|
||
# 4. None of the above ...
|
||
# The `>`/`›` cursor on an option line proves a live selectable menu
|
||
# (prose lists never carry it); the question above ends with "?".
|
||
INTERVIEW_OPT_RE = re.compile("^\\s*([>\\u203a])?\\s*(\\d+)\\.\\s+\\S")
|
||
INTERVIEW_CUE_ABOVE = 8 # question must sit within N lines above option 1
|
||
|
||
# Explicit-phrase request (model asks the user to reply a magic word;
|
||
# observed live): "Reply ACCEPT to approve this text as written ..."
|
||
# Only a single ALL-CAPS token qualifies -- lowercase/prose after Reply
|
||
# never matches, which keeps quoted instructions and chat from firing.
|
||
EXPLICIT_RE = re.compile(r"^\s*(?i:Reply)\s+([A-Z][A-Z0-9_-]{1,11})\b")
|
||
EXPLICIT_WINDOW = 8 # request must sit in the last N content lines
|
||
|
||
# Runtime state sensing for external agents driving panes via send-keys.
|
||
# A working pane shows a running indicator ("- running (Ns ...", "Calling
|
||
# tools (...", or the "esc to interrupt" tail, often wrapped/edge-cut);
|
||
# an idle pane sits at the prompt glyph (U+276F, escaped below).
|
||
RUNNING_RE = re.compile("([\\u2014-] running \\(|Calling tools \\(|esc to\\b)")
|
||
OPEN_PROMPT_RE = re.compile("^\\s*\\u276f")
|
||
STATE_TAIL_WINDOW = 12 # state cues must sit in the last N content lines
|
||
|
||
|
||
def slug_socket(socket_path):
|
||
"""Short filesystem-safe id for a tmux socket: basename + hash suffix.
|
||
|
||
The hash suffix guards against two different socket paths sharing a
|
||
basename (e.g. /tmp/a/default vs /tmp/b/default).
|
||
"""
|
||
base = os.path.basename(socket_path.rstrip("/")) or "sock"
|
||
safe = re.sub(r"[^A-Za-z0-9_.-]+", "_", base)[:32]
|
||
digest = hashlib.sha1(socket_path.encode()).hexdigest()[:6]
|
||
return "%s-%s" % (safe, digest)
|
||
|
||
|
||
def clean_pane(pane_id):
|
||
return re.sub(r"[^A-Za-z0-9]+", "", pane_id)
|
||
|
||
|
||
def pidfile_for(socket_path, pane_id):
|
||
return os.path.join(
|
||
STATE_DIR, "%s-%s-%s.pid" % (FILE_PREFIX, slug_socket(socket_path), clean_pane(pane_id))
|
||
)
|
||
|
||
|
||
def logfile_for(socket_path, pane_id):
|
||
return os.path.join(
|
||
STATE_DIR, "%s-%s-%s.log" % (FILE_PREFIX, slug_socket(socket_path), clean_pane(pane_id))
|
||
)
|
||
|
||
|
||
def answered_file_for(socket_path, pane_id):
|
||
"""On-disk answered-sig store: restarts must not re-answer prompts."""
|
||
return os.path.join(
|
||
STATE_DIR, "%s-%s-%s.answered.json" % (FILE_PREFIX, slug_socket(socket_path), clean_pane(pane_id))
|
||
)
|
||
|
||
|
||
class WatcherLog:
|
||
"""Never-raising JSON-lines logger with best-effort rotation."""
|
||
|
||
def __init__(self, path):
|
||
self.path = path
|
||
|
||
def _rotate(self):
|
||
try:
|
||
if os.path.exists(self.path) and os.path.getsize(self.path) >= LOG_MAX_BYTES:
|
||
try:
|
||
if os.path.exists(self.path + ".1"):
|
||
os.remove(self.path + ".1")
|
||
except OSError:
|
||
pass
|
||
os.rename(self.path, self.path + ".1")
|
||
except Exception:
|
||
pass
|
||
|
||
def log(self, level, msg, **fields):
|
||
try:
|
||
self._rotate()
|
||
rec = {
|
||
"ts": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"),
|
||
"level": level,
|
||
"msg": msg,
|
||
}
|
||
rec.update(fields)
|
||
with open(self.path, "a") as f:
|
||
f.write(json.dumps(rec) + "\n")
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _audit_local(action, name, caller, extra):
|
||
"""Append a box-ctl.jsonl audit record without importing approvals."""
|
||
rec = {
|
||
"ts": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"),
|
||
"action": action,
|
||
"type": "muse-choice",
|
||
"name": name,
|
||
"caller": caller,
|
||
}
|
||
if extra:
|
||
rec.update(extra)
|
||
with open(CTL_LOG, "a") as f:
|
||
f.write(json.dumps(rec) + "\n")
|
||
|
||
|
||
def audit(action, name=None, caller="muse-choice-watcher", extra=None):
|
||
"""Feed an event to box (box-ctl.jsonl). Never raises.
|
||
|
||
Prefers approvals.log_box_ctl so the audit schema stays unified; falls
|
||
back to a direct append if the import fails.
|
||
"""
|
||
try:
|
||
try:
|
||
import approvals
|
||
approvals.log_box_ctl(action, name=name, caller=caller, extra=extra)
|
||
except Exception:
|
||
_audit_local(action, name, caller, extra or {})
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def get_desired():
|
||
"""Desired daemon state: {"enabled", "dry_run", ...}. Default: on.
|
||
|
||
Policy: auto-approve is on unless explicitly disabled (`box
|
||
muse-choices off`), the pane opted out via launch flags, or a human
|
||
intervened. A missing/unreadable state file (or a file without the
|
||
key) means undefined -> on.
|
||
"""
|
||
default = {"enabled": True, "dry_run": False}
|
||
try:
|
||
with open(DESIRED_STATE_FILE) as f:
|
||
data = json.load(f)
|
||
return {"enabled": bool(data.get("enabled", True)),
|
||
"dry_run": bool(data.get("dry_run", False)),
|
||
"updated_at": data.get("updated_at", ""),
|
||
"updated_by": data.get("updated_by", "")}
|
||
except Exception:
|
||
return dict(default)
|
||
|
||
|
||
def set_enabled(enabled, dry_run=None, by="box"):
|
||
"""Write desired state + audit the transition. Returns the new state."""
|
||
state = get_desired()
|
||
state["enabled"] = bool(enabled)
|
||
if dry_run is not None:
|
||
state["dry_run"] = bool(dry_run)
|
||
state["updated_at"] = datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||
state["updated_by"] = by
|
||
try:
|
||
parent = os.path.dirname(DESIRED_STATE_FILE)
|
||
if parent:
|
||
os.makedirs(parent, exist_ok=True)
|
||
tmp = DESIRED_STATE_FILE + ".tmp.%d" % os.getpid()
|
||
with open(tmp, "w") as f:
|
||
json.dump(state, f, indent=1)
|
||
os.replace(tmp, DESIRED_STATE_FILE)
|
||
except Exception:
|
||
pass
|
||
audit("muse-choice-enabled" if enabled else "muse-choice-disabled",
|
||
caller=by, extra={"dry_run": state["dry_run"]})
|
||
return state
|
||
|
||
|
||
def _option_markers(lines):
|
||
"""Return [(line_index, letter)] for lettered-option lines."""
|
||
out = []
|
||
for i, line in enumerate(lines):
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
m = OPTION_RE.match(line)
|
||
if m:
|
||
out.append((i, m.group(1)))
|
||
return out
|
||
|
||
|
||
def _ordered_run(markers):
|
||
"""Find A,B[,C...] in order. Returns marker sub-list or None.
|
||
|
||
Markers between run members (wrapped text has no markers, so any
|
||
interleaving lettered line breaks the run) must fit OPTION_SPAN.
|
||
"""
|
||
want = "ABCDEFGHIJKLMNOPQRSTUVWXYZ"
|
||
best = []
|
||
for start in range(len(markers)):
|
||
if markers[start][1] != "A":
|
||
continue
|
||
run = [markers[start]]
|
||
for idx in range(start + 1, len(markers)):
|
||
expected = want[len(run)]
|
||
if markers[idx][1] == expected:
|
||
run.append(markers[idx])
|
||
else:
|
||
break
|
||
if len(run) >= 2 and run[-1][0] - run[0][0] <= OPTION_SPAN:
|
||
if len(run) > len(best):
|
||
best = run
|
||
return best or None
|
||
|
||
|
||
def _sig_for(kind, parts):
|
||
src = kind + "\n" + "\n".join(parts)
|
||
return hashlib.sha1(src.encode()).hexdigest()[:16]
|
||
|
||
|
||
def _find_letter(window):
|
||
"""A/B/C lettered choice block -> {"sig", "kind", "key", ...} or None."""
|
||
run = _ordered_run(_option_markers(window))
|
||
if not run:
|
||
return None
|
||
start, end = run[0][0], run[-1][0]
|
||
lo = max(0, start - 4)
|
||
hi = min(len(window), end + 5)
|
||
cue = None
|
||
for i in range(lo, hi):
|
||
if CUE_RE.search(window[i]) or QUESTION_RE.search(window[i]):
|
||
cue = window[i].strip()
|
||
break
|
||
if cue is None:
|
||
return None
|
||
options = [window[i].strip()[:120] for i, _ in run]
|
||
cue = cue[:200]
|
||
return {"sig": _sig_for("letter", options + [cue]), "kind": "letter",
|
||
"key": "A", "options": options, "cue": cue,
|
||
"start": start, "end": end}
|
||
|
||
|
||
def _find_numbered(window):
|
||
"""(1)/(2) numbered menu with selection cue -> match or None."""
|
||
markers = []
|
||
for i, line in enumerate(window):
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
m = NUMBERED_OPT_RE.match(line)
|
||
if m:
|
||
markers.append((i, int(m.group(1))))
|
||
ones = [i for i, n in markers if n == 1]
|
||
twos = [i for i, n in markers if n == 2]
|
||
if not ones or not twos:
|
||
return None
|
||
start = ones[0]
|
||
end = next((i for i in twos if i > start), None)
|
||
if end is None or end - start > OPTION_SPAN:
|
||
return None
|
||
lo = max(0, start - 2)
|
||
hi = min(len(window), end + 4)
|
||
cue = None
|
||
for i in range(lo, hi):
|
||
if NUMBERED_CUE_RE.search(window[i]):
|
||
cue = window[i].strip()
|
||
break
|
||
if cue is None:
|
||
return None
|
||
options = [window[i].strip()[:120]
|
||
for i, n in markers if start <= i <= end]
|
||
cue = cue[:200]
|
||
return {"sig": _sig_for("numbered", options + [cue]), "kind": "numbered",
|
||
"key": "1", "options": options, "cue": cue,
|
||
"start": start, "end": end}
|
||
|
||
|
||
def _find_yn(window):
|
||
"""y/n prompt at the end of a recent line -> match or None."""
|
||
recent = window[-YN_WINDOW:]
|
||
base = len(window) - len(recent)
|
||
for j in range(len(recent) - 1, -1, -1):
|
||
line = recent[j]
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if YN_RE.search(line):
|
||
text = line.strip()[:200]
|
||
ctx = recent[j - 1].strip()[:120] if j > 0 else ""
|
||
i = base + j
|
||
return {"sig": _sig_for("yn", [text, ctx]), "kind": "yn",
|
||
"key": "y", "options": [text], "cue": text,
|
||
"start": i, "end": i}
|
||
return None
|
||
|
||
|
||
def _find_muse_approval(window):
|
||
"""Native Muse TUI approval menu -> match or None.
|
||
|
||
Requires the Would-you-like cue plus a 1.Yes/2.No pair below it.
|
||
A decided block ("approval decision ...") within the span never
|
||
matches, so an already-landed decision is never double-answered.
|
||
"""
|
||
cue_idx = None
|
||
for i, line in enumerate(window):
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if MUSE_APPROVAL_CUE_RE.match(line):
|
||
cue_idx = i
|
||
break
|
||
if cue_idx is None:
|
||
return None
|
||
yes_idx = no_idx = None
|
||
for i in range(cue_idx + 1, min(len(window), cue_idx + 1 + MUSE_APPROVAL_SPAN)):
|
||
line = window[i]
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if yes_idx is None and MUSE_APPROVAL_YES_RE.match(line):
|
||
yes_idx = i
|
||
elif MUSE_APPROVAL_NO_RE.match(line):
|
||
no_idx = i
|
||
break
|
||
if yes_idx is None or no_idx is None:
|
||
return None
|
||
for line in window[cue_idx:no_idx + 4]:
|
||
if MUSE_APPROVAL_DECIDED_RE.search(line):
|
||
return None
|
||
cue = window[cue_idx].strip()[:200]
|
||
options = [window[yes_idx].strip()[:120], window[no_idx].strip()[:120]]
|
||
# Sig covers the WHOLE block (cue + $ command + options): bare
|
||
# options+cue are byte-identical across every command approval, which
|
||
# used to collapse all of a pane's approvals into one sig and wedge
|
||
# every dialog after the first (observed live 2026-10-06).
|
||
context = [ln.strip()[:120] for ln in window[cue_idx:no_idx + 1]]
|
||
return {"sig": _sig_for("muse-approval", context),
|
||
"kind": "muse-approval", "key": "1", "options": options,
|
||
"cue": cue, "context": context,
|
||
"start": cue_idx, "end": no_idx}
|
||
|
||
|
||
def _find_collapsed_approval(window):
|
||
"""Collapsed approval menu (long command, options hidden) or None.
|
||
|
||
Requires the Would-you-like cue plus the collapsed markers ("N
|
||
command rows omitted" + "ctrl+o view full command"). Answers with
|
||
a bare Enter: that expands the block to the full 1/2 form (or
|
||
accepts outright), and the expanded form answers as a fresh
|
||
muse-approval prompt on later polls. (ctrl+o is NOT used: it opens
|
||
a modal pager that strands the session.)
|
||
"""
|
||
cue_idx = None
|
||
for i, line in enumerate(window):
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if MUSE_APPROVAL_CUE_RE.match(line):
|
||
cue_idx = i
|
||
break
|
||
if cue_idx is None:
|
||
return None
|
||
rows_idx = key_idx = None
|
||
for i in range(cue_idx + 1, min(len(window), cue_idx + 1 + MUSE_APPROVAL_SPAN)):
|
||
line = window[i]
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if rows_idx is None and MUSE_COLLAPSED_ROWS_RE.search(line):
|
||
rows_idx = i
|
||
if MUSE_COLLAPSED_KEY_RE.search(line):
|
||
key_idx = i
|
||
if rows_idx is None or key_idx is None:
|
||
return None
|
||
end = max(rows_idx, key_idx)
|
||
for line in window[cue_idx:end + 4]:
|
||
if MUSE_APPROVAL_DECIDED_RE.search(line):
|
||
return None
|
||
cue = window[cue_idx].strip()[:200]
|
||
options = [window[rows_idx].strip()[:120],
|
||
window[key_idx].strip()[:120]]
|
||
context = [ln.strip()[:120] for ln in window[cue_idx:end + 1]]
|
||
return {"sig": _sig_for("muse-approval-collapsed", context),
|
||
"kind": "muse-approval-collapsed", "key": "Enter",
|
||
"enter": False, "options": options, "cue": cue,
|
||
"context": context, "start": cue_idx, "end": end}
|
||
|
||
|
||
def _find_interview(window):
|
||
"""Agent interview menu (cursor + ordered 1./2./.. + question) or None.
|
||
|
||
Requires an ordered run starting at 1 with at least options 1 and 2,
|
||
a selection cursor on one of the run's lines, and a "?"-ended
|
||
question above. The cursor requirement is what separates a live
|
||
interview from a prose numbered list.
|
||
"""
|
||
markers = []
|
||
for i, line in enumerate(window):
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
m = INTERVIEW_OPT_RE.match(line)
|
||
if m:
|
||
markers.append((i, int(m.group(2)), bool(m.group(1))))
|
||
ones = [i for i, n, c in markers if n == 1]
|
||
if not ones:
|
||
return None
|
||
start = ones[0]
|
||
run = [start]
|
||
for i, n, c in markers:
|
||
if i > start and n == len(run) + 1 and i - start <= OPTION_SPAN:
|
||
run.append(i)
|
||
if len(run) < 2:
|
||
return None
|
||
if not any(c for i, n, c in markers if run[0] <= i <= run[-1]):
|
||
return None
|
||
question = None
|
||
for j in range(max(0, start - INTERVIEW_CUE_ABOVE), start):
|
||
if QUESTION_RE.search(window[j]):
|
||
question = window[j].strip()
|
||
if question is None:
|
||
return None
|
||
options = [window[i].strip()[:120] for i in run]
|
||
question = question[:200]
|
||
return {"sig": _sig_for("interview", options + [question]),
|
||
"kind": "interview", "key": "1", "options": options,
|
||
"cue": question, "start": start, "end": run[-1]}
|
||
|
||
|
||
def _find_explicit(window):
|
||
"""Explicit-phrase request ("Reply ACCEPT to ...") or None.
|
||
|
||
Scans the last EXPLICIT_WINDOW content lines bottom-up so the
|
||
freshest request wins. The key is the captured token itself.
|
||
"""
|
||
recent = window[-EXPLICIT_WINDOW:]
|
||
base = len(window) - len(recent)
|
||
for j in range(len(recent) - 1, -1, -1):
|
||
line = recent[j]
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
m = EXPLICIT_RE.match(line)
|
||
if m:
|
||
token = m.group(1)
|
||
text = line.strip()[:200]
|
||
i = base + j
|
||
return {"sig": _sig_for("explicit-phrase", [text]),
|
||
"kind": "explicit-phrase", "key": token,
|
||
"options": [text], "cue": text,
|
||
"start": i, "end": i}
|
||
return None
|
||
|
||
|
||
def find_choice_prompt(text, tail_window=TAIL_WINDOW):
|
||
"""Detect a Muse prompt awaiting reply.
|
||
|
||
Kinds (priority order): muse-approval (native Would-you-like
|
||
menu -> key "1"), muse-approval-collapsed (collapsed long command
|
||
-> bare Enter, expand-or-accept), interview (cursor + ordered
|
||
1./2. menu + question -> key "1"), explicit-phrase ("Reply TOKEN
|
||
to ..." -> key TOKEN), letter (A/B/C -> key "A"), numbered
|
||
((1)/(2) menu -> key "1"), yn (y/n line-end -> key "y").
|
||
Returns {"sig", "kind", "key", "options", "cue", "start", "end"}
|
||
or None. Only blocks in the last `tail_window` content lines are
|
||
eligible, so answered/stale prompts in scrollback never re-trigger.
|
||
"""
|
||
if not text:
|
||
return None
|
||
lines = text.split("\n")
|
||
# Trailing blank lines (tall panes, fresh prompts) must not push a live
|
||
# prompt out of the window; staleness is judged on content lines.
|
||
while lines and not lines[-1].strip():
|
||
lines.pop()
|
||
window = lines[-tail_window:]
|
||
if not window:
|
||
return None
|
||
return (_find_muse_approval(window) or _find_collapsed_approval(window)
|
||
or _find_interview(window) or _find_explicit(window)
|
||
or _find_letter(window) or _find_numbered(window)
|
||
or _find_yn(window))
|
||
|
||
|
||
def runtime_state(text):
|
||
"""Classify a Muse pane capture for agentic send-keys drivers.
|
||
|
||
Returns {"state", "match"} where state is one of approval-pending
|
||
(a choice prompt awaits reply), working (running indicator),
|
||
open-prompt (idle at the input prompt), or unknown. Priority is
|
||
approval-pending > working > open-prompt: a working pane still
|
||
renders its prompt footer, and an approval can arrive mid-run.
|
||
"""
|
||
match = find_choice_prompt(text)
|
||
if match is not None:
|
||
return {"state": "approval-pending", "match": match}
|
||
lines = (text or "").split("\n")
|
||
while lines and not lines[-1].strip():
|
||
lines.pop()
|
||
tail = lines[-STATE_TAIL_WINDOW:]
|
||
for line in tail:
|
||
if RUNNING_RE.search(line):
|
||
return {"state": "working", "match": None}
|
||
for line in tail:
|
||
if len(line) > MAX_OPTION_LEN:
|
||
continue
|
||
if OPEN_PROMPT_RE.match(line):
|
||
return {"state": "open-prompt", "match": None}
|
||
return {"state": "unknown", "match": None}
|
||
|
||
|
||
def _cmdline(pid):
|
||
"""Argv of a pid as a list (empty when unreadable). Never raises."""
|
||
try:
|
||
with open("/proc/%d/cmdline" % pid, "rb") as f:
|
||
raw = f.read().split(b"\0")
|
||
return [p.decode(errors="replace") for p in raw if p]
|
||
except Exception:
|
||
return []
|
||
|
||
|
||
def _child_pids(pid):
|
||
"""Direct children of pid via a bounded /proc scan. Never raises."""
|
||
out = []
|
||
try:
|
||
names = os.listdir("/proc")
|
||
except Exception:
|
||
return out
|
||
for name in names:
|
||
if not name.isdigit():
|
||
continue
|
||
try:
|
||
with open("/proc/%s/stat" % name) as f:
|
||
parts = f.read().rsplit(")", 1)[1].split()
|
||
if int(parts[1]) == pid:
|
||
out.append(int(name))
|
||
except Exception:
|
||
continue
|
||
return out
|
||
|
||
|
||
def muse_approval_flags(cmd_argv):
|
||
"""Approval posture of a muse argv.
|
||
|
||
Returns {"auto_approve", "flags", "profile", "mode", "bypass"}.
|
||
profile is the --permission-profile id (built-ins :read-only,
|
||
:standard, :unrestricted) or None; mode is "yolo" for --yolo, the
|
||
profile id when one was passed, else "default"; bypass is True
|
||
when tool-approval prompts are disabled at launch (yolo,
|
||
--disable-approval, or --approval-mode=never), in which case the
|
||
pane shows no approval dialogs (questions may still render, so the
|
||
watcher stays active).
|
||
"""
|
||
flags = []
|
||
argv = cmd_argv or []
|
||
if "--yolo" in argv:
|
||
flags.append("yolo")
|
||
if "--disable-approval" in argv:
|
||
flags.append("disable-approval")
|
||
if "--disable-sandbox" in argv:
|
||
flags.append("disable-sandbox")
|
||
if "--trust-workspace" in argv:
|
||
flags.append("trust-workspace")
|
||
profile = None
|
||
for i, arg in enumerate(argv):
|
||
if arg == "--approval-mode" and i + 1 < len(argv):
|
||
flags.append("approval-mode=%s" % argv[i + 1])
|
||
elif arg.startswith("--approval-mode="):
|
||
flags.append("approval-mode=%s" % arg.split("=", 1)[1])
|
||
elif arg == "--permission-profile" and i + 1 < len(argv):
|
||
profile = argv[i + 1]
|
||
flags.append("permission-profile=%s" % profile)
|
||
elif arg.startswith("--permission-profile="):
|
||
profile = arg.split("=", 1)[1]
|
||
flags.append("permission-profile=%s" % profile)
|
||
auto = ("yolo" in flags or "disable-approval" in flags
|
||
or "approval-mode=never" in flags)
|
||
if "yolo" in flags:
|
||
mode = "yolo"
|
||
elif profile is not None:
|
||
mode = profile
|
||
else:
|
||
mode = "default"
|
||
return {"auto_approve": auto, "flags": flags, "profile": profile,
|
||
"mode": mode, "bypass": auto}
|
||
|
||
|
||
def launch_opt_out(cmd_argv):
|
||
"""True when a muse argv explicitly opts out of watcher auto-approve.
|
||
|
||
Bare argv (no approval flags) means undefined -> the watcher
|
||
answers. An explicit --approval-mode other than never (on-request,
|
||
untrusted, ...) declares human-in-the-loop intent -> hold.
|
||
"""
|
||
argv = cmd_argv or []
|
||
for i, arg in enumerate(argv):
|
||
mode = None
|
||
if arg == "--approval-mode" and i + 1 < len(argv):
|
||
mode = argv[i + 1]
|
||
elif arg.startswith("--approval-mode="):
|
||
mode = arg.split("=", 1)[1]
|
||
if mode is not None and mode != "never":
|
||
return True
|
||
return False
|
||
|
||
|
||
def pane_muse_argv(socket_path, pane_id):
|
||
"""Argv of the muse process in a pane (empty when not found)."""
|
||
try:
|
||
r = _tmux(socket_path, "list-panes", "-a", "-F",
|
||
"#{pane_id} #{pane_pid}", timeout=5)
|
||
if r.returncode != 0:
|
||
return []
|
||
except Exception:
|
||
return []
|
||
pane_pid = None
|
||
for line in r.stdout.split("\n"):
|
||
parts = line.split()
|
||
if len(parts) == 2 and parts[0] == pane_id:
|
||
try:
|
||
pane_pid = int(parts[1])
|
||
except ValueError:
|
||
return []
|
||
if pane_pid is None:
|
||
return []
|
||
for cand in [pane_pid] + _child_pids(pane_pid):
|
||
argv = _cmdline(cand)
|
||
if argv and ("muse-bin" in argv[0] or "muse-code" in argv[0]):
|
||
return argv
|
||
return []
|
||
|
||
|
||
# Minimum pane geometry for reliable approval rendering. Empirically
|
||
# derived: a 35x7 tile drops approval text the matcher needs, while
|
||
# 35x35/36x35/71x27 panes answer cleanly (width 35 works when tall
|
||
# enough; capture uses -J so wrapping is width-independent). Below
|
||
# either bound the pane is flagged squeezed (see `box runtime layout` /
|
||
# `spread`).
|
||
MIN_APPROVAL_WIDTH = 35
|
||
MIN_APPROVAL_HEIGHT = 12
|
||
|
||
# Session naming convention (see NODES.md): <node>--<role>--<id>
|
||
# separates node runtimes on the shared stable server, e.g.
|
||
# pip--worker--01. Ad-hoc sessions carry no node and show "-".
|
||
NODE_NAMES = ("muse", "pip", "646", "opm", "def", "dev")
|
||
NODE_SESSION_RE = re.compile(r"^([A-Za-z0-9]+)--([A-Za-z0-9-]+)--([A-Za-z0-9]+)$")
|
||
|
||
|
||
def node_from_session(session_name):
|
||
"""Node owning a tmux session per the naming convention.
|
||
|
||
Returns the node name, or None for ad-hoc sessions and unknown
|
||
node prefixes. Pure function for supervision display.
|
||
"""
|
||
if not session_name:
|
||
return None
|
||
m = NODE_SESSION_RE.match(session_name)
|
||
if not m:
|
||
return None
|
||
node = m.group(1).lower()
|
||
return node if node in NODE_NAMES else None
|
||
|
||
|
||
def runtime_rows(socket_path):
|
||
"""One row per pane: identity, approval posture, live state, watcher.
|
||
|
||
Rows are JSON-serializable dicts for `box runtime list` and external
|
||
agents. Panes that vanish mid-scan are skipped, never fatal.
|
||
"""
|
||
rows = []
|
||
try:
|
||
r = _tmux(socket_path, "list-panes", "-a", "-F",
|
||
"#{session_name}\t#{window_index}\t#{pane_id}\t"
|
||
"#{pane_current_command}\t#{pane_pid}\t"
|
||
"#{pane_width}\t#{pane_height}", timeout=10)
|
||
if r.returncode != 0:
|
||
return rows
|
||
except Exception:
|
||
return rows
|
||
for line in r.stdout.split("\n"):
|
||
parts = line.split("\t")
|
||
if len(parts) != 7:
|
||
continue
|
||
session, window, pane_id, cmd, pid, width, height = parts
|
||
try:
|
||
pane_pid = int(pid)
|
||
except ValueError:
|
||
pane_pid = None
|
||
try:
|
||
pane_width, pane_height = int(width), int(height)
|
||
except ValueError:
|
||
pane_width = pane_height = None
|
||
squeezed = (pane_width is not None and pane_height is not None
|
||
and (pane_width < MIN_APPROVAL_WIDTH
|
||
or pane_height < MIN_APPROVAL_HEIGHT))
|
||
node = node_from_session(session)
|
||
is_muse = "muse-bin" in cmd or "muse-code" in cmd
|
||
posture = {"auto_approve": None, "flags": [], "profile": None,
|
||
"mode": None, "bypass": None}
|
||
if is_muse and pane_pid:
|
||
candidates = [pane_pid] + _child_pids(pane_pid)
|
||
for cand in candidates:
|
||
argv = _cmdline(cand)
|
||
if argv and ("muse-bin" in argv[0]
|
||
or "muse-code" in argv[0]):
|
||
posture = muse_approval_flags(argv)
|
||
break
|
||
text = capture_pane(socket_path, pane_id)
|
||
if text is None:
|
||
continue
|
||
st = runtime_state(text)
|
||
match = st["match"] or {}
|
||
watcher_pid = watcher_alive(socket_path, pane_id)
|
||
rows.append({
|
||
"socket": socket_path, "session": session,
|
||
"window": window, "pane": pane_id, "cmd": cmd,
|
||
"pid": pane_pid, "is_muse": is_muse,
|
||
"node": node,
|
||
"width": pane_width, "height": pane_height,
|
||
"squeezed": squeezed,
|
||
"auto_approve": posture["auto_approve"],
|
||
"approval_flags": posture["flags"],
|
||
"permission_mode": posture["mode"],
|
||
"permission_profile": posture["profile"],
|
||
"permission_bypass": posture["bypass"],
|
||
"state": st["state"],
|
||
"prompt_kind": match.get("kind"),
|
||
"prompt_key": match.get("key"),
|
||
"watcher_alive": watcher_pid is not None,
|
||
"watcher_pid": watcher_pid,
|
||
})
|
||
return rows
|
||
|
||
|
||
def spread_targets(rows):
|
||
"""Rows that `box runtime spread` would break into own windows.
|
||
|
||
A row is a target when it runs a coding runtime (muse today; other
|
||
harnesses follow the registry) and its geometry is squeezed below
|
||
the reliable-approval minimums. Pure function of rows for testing.
|
||
"""
|
||
return [r for r in rows
|
||
if r.get("is_muse") and r.get("squeezed")]
|
||
|
||
|
||
def all_runtime_rows(sockets=None):
|
||
"""runtime_rows across every known socket that exists."""
|
||
rows = []
|
||
for sock in sockets or KNOWN_SOCKETS:
|
||
if not os.path.exists(sock):
|
||
continue
|
||
rows.extend(runtime_rows(sock))
|
||
return rows
|
||
|
||
|
||
def pane_state(socket_path, pane_id):
|
||
"""Single-pane runtime row, or {"error": ...} when not found."""
|
||
for row in runtime_rows(socket_path):
|
||
if row["pane"] == pane_id:
|
||
return row
|
||
return {"error": "no_such_pane", "socket": socket_path,
|
||
"pane": pane_id}
|
||
|
||
|
||
class WatcherState:
|
||
"""Tracks prompt stability and once-per-prompt answering.
|
||
|
||
When persist_path is set, answered sigs survive restarts (atomic
|
||
JSON store): a daemon that restarts while an answered prompt is
|
||
still visible must not answer it again (observed live: same sig
|
||
re-answered minutes later, stray "1" landing in the input box).
|
||
"""
|
||
|
||
def __init__(self, persist_path=None):
|
||
self.pending_sig = None
|
||
self.stable_count = 0
|
||
self.answered_sigs = {}
|
||
self.answer_times = []
|
||
self.last_capped_sig = None
|
||
self._persist_path = persist_path
|
||
if persist_path:
|
||
self._load_answered()
|
||
|
||
def _load_answered(self):
|
||
try:
|
||
with open(self._persist_path) as f:
|
||
data = json.load(f)
|
||
except Exception:
|
||
return
|
||
if not isinstance(data, dict):
|
||
return
|
||
now = time.time()
|
||
for sig, ts in data.items():
|
||
if (isinstance(sig, str) and isinstance(ts, (int, float))
|
||
and now - ts < ANSWERED_TTL_SECONDS):
|
||
self.answered_sigs[sig] = ts
|
||
|
||
def _save_answered(self):
|
||
if not self._persist_path:
|
||
return
|
||
try:
|
||
tmp = "%s.tmp.%d" % (self._persist_path, os.getpid())
|
||
with open(tmp, "w") as f:
|
||
json.dump(self.answered_sigs, f)
|
||
os.replace(tmp, self._persist_path)
|
||
except Exception:
|
||
pass
|
||
|
||
def refresh_answered(self):
|
||
"""Merge on-disk answered sigs into memory. Never raises.
|
||
|
||
Cross-process once-per-prompt: a peer watcher that answered
|
||
after our startup persisted its sig; re-reading before we type
|
||
suppresses the duplicate ('11' in the input box).
|
||
"""
|
||
if not self._persist_path:
|
||
return
|
||
try:
|
||
with open(self._persist_path) as f:
|
||
data = json.load(f)
|
||
except Exception:
|
||
return
|
||
if not isinstance(data, dict):
|
||
return
|
||
now = time.time()
|
||
for sig, ts in data.items():
|
||
if (isinstance(sig, str) and isinstance(ts, (int, float))
|
||
and now - ts < ANSWERED_TTL_SECONDS):
|
||
if sig not in self.answered_sigs:
|
||
self.answered_sigs[sig] = ts
|
||
|
||
def try_claim(self, sig, now=None):
|
||
"""Atomically claim a sig for this process. True iff we won.
|
||
|
||
O_CREAT|O_EXCL makes the first claimer win even when two
|
||
watchers reach the same stable prompt in the same poll window;
|
||
the loser suppresses instead of double-typing. Claims expire
|
||
after CLAIM_TTL_SECONDS (concurrent window only; stuck-dialog
|
||
recovery is governed by the answered store's longer TTL), so a
|
||
winner that dies between claim and send wedges at most briefly.
|
||
Memory-only states (no persist path) always win. Never raises.
|
||
"""
|
||
if not self._persist_path:
|
||
return True
|
||
now = time.time() if now is None else now
|
||
base = os.path.basename(self._persist_path)
|
||
directory = os.path.dirname(self._persist_path) or STATE_DIR
|
||
# Prune expired claims for this pane (best-effort).
|
||
try:
|
||
for name in os.listdir(directory):
|
||
if not name.startswith(base + ".") or not name.endswith(".claim"):
|
||
continue
|
||
p = os.path.join(directory, name)
|
||
try:
|
||
with open(p) as f:
|
||
rec = json.load(f)
|
||
ts = float(rec.get("ts", 0))
|
||
except Exception:
|
||
ts = 0
|
||
try:
|
||
if now - ts >= CLAIM_TTL_SECONDS:
|
||
os.remove(p)
|
||
except OSError:
|
||
pass
|
||
except Exception:
|
||
pass
|
||
path = "%s.%s.claim" % (self._persist_path, sig)
|
||
try:
|
||
fd = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY)
|
||
except FileExistsError:
|
||
return False
|
||
except OSError:
|
||
return True # claim store unavailable: fail open, send once
|
||
try:
|
||
os.write(fd, json.dumps({"pid": os.getpid(),
|
||
"ts": now}).encode())
|
||
except Exception:
|
||
pass
|
||
try:
|
||
os.close(fd)
|
||
except Exception:
|
||
pass
|
||
return True
|
||
|
||
def _prune_times(self, now):
|
||
cutoff = now - 3600
|
||
self.answer_times = [t for t in self.answer_times if t >= cutoff]
|
||
|
||
def cap_reached(self, now):
|
||
self._prune_times(now)
|
||
return len(self.answer_times) >= MAX_ANSWERS_PER_HOUR
|
||
|
||
def observe(self, match, now):
|
||
"""Feed one poll's match (or None). Returns "answer" when the prompt
|
||
is stable, unanswered, and under the hourly cap."""
|
||
if match is None:
|
||
self.pending_sig = None
|
||
self.stable_count = 0
|
||
return "none"
|
||
sig = match["sig"]
|
||
if sig in self.answered_sigs:
|
||
if now - self.answered_sigs[sig] < ANSWERED_TTL_SECONDS:
|
||
return "none"
|
||
# Expired: an identical prompt still (or again) present long
|
||
# after its answer is a stuck dialog, not a duplicate -- allow
|
||
# one recovery answer instead of wedging forever.
|
||
del self.answered_sigs[sig]
|
||
if sig == self.pending_sig:
|
||
self.stable_count += 1
|
||
else:
|
||
self.pending_sig = sig
|
||
self.stable_count = 1
|
||
if self.stable_count < STABILITY_POLLS:
|
||
return "wait"
|
||
if self.cap_reached(now):
|
||
return "capped"
|
||
return "answer"
|
||
|
||
def record_answer(self, sig, now):
|
||
self.answered_sigs[sig] = now
|
||
# Bound memory: drop expired entries, then oldest-first past the cap.
|
||
cutoff = now - ANSWERED_TTL_SECONDS
|
||
for s in [s for s, t in self.answered_sigs.items() if t < cutoff]:
|
||
del self.answered_sigs[s]
|
||
if len(self.answered_sigs) > 200:
|
||
for s in sorted(self.answered_sigs,
|
||
key=lambda s: self.answered_sigs[s])[:-100]:
|
||
del self.answered_sigs[s]
|
||
self.answer_times.append(now)
|
||
self.pending_sig = None
|
||
self.stable_count = 0
|
||
self._save_answered()
|
||
|
||
|
||
def _tmux(socket_path, *args, timeout=5):
|
||
cmd = ["tmux", "-S", socket_path] + list(args)
|
||
return subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
|
||
|
||
|
||
def pane_exists(socket_path, pane_id):
|
||
try:
|
||
r = _tmux(socket_path, "list-panes", "-a", "-F", "#{pane_id}", timeout=5)
|
||
if r.returncode != 0:
|
||
return False
|
||
return pane_id in r.stdout.split()
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def capture_pane(socket_path, pane_id, history=80):
|
||
"""Capture pane text with wrapped rows joined (-J).
|
||
|
||
-J makes matching width-independent: narrow panes wrap the same
|
||
dialog onto more physical rows, which otherwise pushes cue/option
|
||
spans apart and breaks the matcher.
|
||
"""
|
||
try:
|
||
r = _tmux(socket_path, "capture-pane", "-p", "-J", "-t", pane_id,
|
||
"-S", "-%d" % history, timeout=5)
|
||
if r.returncode != 0:
|
||
return None
|
||
return r.stdout
|
||
except Exception:
|
||
return None
|
||
|
||
|
||
def send_answer(socket_path, pane_id, letter="A", enter=True):
|
||
"""Type the choice key (+ Enter unless enter=False). True on success."""
|
||
try:
|
||
r1 = _tmux(socket_path, "send-keys", "-t", pane_id, letter, timeout=5)
|
||
if r1.returncode != 0:
|
||
return False
|
||
if not enter:
|
||
return True
|
||
r2 = _tmux(socket_path, "send-keys", "-t", pane_id, "Enter", timeout=5)
|
||
return r2.returncode == 0
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def _load_rules(log=None):
|
||
"""Load the D0 rules dictionary, cached by (mtime, size). Never raises:
|
||
a missing or broken file means no rules (fail open per D4)."""
|
||
global _RULES_CACHE
|
||
try:
|
||
st = os.stat(RULES_FILE)
|
||
key = (st.st_mtime, st.st_size)
|
||
except OSError:
|
||
key = "missing"
|
||
if _RULES_CACHE["key"] == key:
|
||
return _RULES_CACHE["rules"]
|
||
rules = []
|
||
if key != "missing":
|
||
try:
|
||
with open(RULES_FILE) as f:
|
||
data = json.load(f)
|
||
rules = [r for r in data.get("rules", [])
|
||
if isinstance(r, dict) and r.get("id")]
|
||
except Exception as e:
|
||
if log is not None:
|
||
try:
|
||
log.log("warn", "rules file unreadable, failing open",
|
||
error="%r" % (e,))
|
||
except Exception:
|
||
pass
|
||
rules = []
|
||
_RULES_CACHE = {"key": key, "rules": rules}
|
||
return rules
|
||
|
||
|
||
def _as_list(v):
|
||
return v if isinstance(v, list) else [v]
|
||
|
||
|
||
def _rule_matches(rule, match):
|
||
kind = rule.get("kind")
|
||
if kind is not None and match["kind"] not in _as_list(kind):
|
||
return False
|
||
token = rule.get("token")
|
||
if token is not None:
|
||
if match["kind"] != "explicit-phrase":
|
||
return False
|
||
if match["key"] not in _as_list(token):
|
||
return False
|
||
cmd = rule.get("command")
|
||
if cmd is not None:
|
||
ctx = match.get("context") or match.get("options", [])
|
||
try:
|
||
if not re.search(cmd, "\n".join(ctx)):
|
||
return False
|
||
except re.error:
|
||
return False
|
||
text = rule.get("text")
|
||
if text is not None:
|
||
hay = "\n".join(match.get("options", []) + [match.get("cue") or ""])
|
||
try:
|
||
if not re.search(text, hay):
|
||
return False
|
||
except re.error:
|
||
return False
|
||
return True
|
||
|
||
|
||
def evaluate_rules(match, log=None):
|
||
"""First matching rule wins -> (decision, rule). No match -> approve.
|
||
|
||
D2: deny rules are skipped for question kinds (questions always
|
||
resolve top-choice). Collapsed approvals have no visible No option,
|
||
so deny downgrades to hold there rather than auto-approving
|
||
something a rule flagged risky.
|
||
"""
|
||
for rule in _load_rules(log=log):
|
||
if not _rule_matches(rule, match):
|
||
continue
|
||
d = rule.get("decision", "approve")
|
||
if d not in ("approve", "deny", "hold"):
|
||
d = "approve"
|
||
if d == "deny":
|
||
if match["kind"] in NEGATIVE_KEYS:
|
||
return "deny", rule
|
||
if match["kind"] in QUESTION_KINDS:
|
||
continue
|
||
downgraded = dict(rule)
|
||
downgraded["downgraded_from"] = "deny"
|
||
return "hold", downgraded
|
||
return d, rule
|
||
return "approve", None
|
||
|
||
|
||
def holdfile_for(socket_path, pane_id):
|
||
return os.path.join(
|
||
STATE_DIR, "%s-%s-%s.held.json" % (FILE_PREFIX, slug_socket(socket_path),
|
||
clean_pane(pane_id)))
|
||
|
||
|
||
def write_hold(socket_path, pane_id, rec):
|
||
try:
|
||
with open(holdfile_for(socket_path, pane_id), "w") as f:
|
||
json.dump(rec, f)
|
||
return True
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def read_hold(socket_path, pane_id):
|
||
try:
|
||
with open(holdfile_for(socket_path, pane_id)) as f:
|
||
rec = json.load(f)
|
||
return rec if isinstance(rec, dict) and rec.get("sig") else None
|
||
except Exception:
|
||
return None
|
||
|
||
|
||
def clear_hold(socket_path, pane_id):
|
||
try:
|
||
os.remove(holdfile_for(socket_path, pane_id))
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def list_holds():
|
||
"""All holdfiles, pruning ones long expired with no watcher to own them."""
|
||
out = []
|
||
try:
|
||
names = os.listdir(STATE_DIR)
|
||
except OSError:
|
||
return out
|
||
now = time.time()
|
||
for name in names:
|
||
if not name.startswith(FILE_PREFIX + "-") or not name.endswith(".held.json"):
|
||
continue
|
||
path = os.path.join(STATE_DIR, name)
|
||
try:
|
||
with open(path) as f:
|
||
rec = json.load(f)
|
||
except Exception:
|
||
continue
|
||
if not isinstance(rec, dict) or not rec.get("sig"):
|
||
continue
|
||
try:
|
||
expired_ago = now - float(rec.get("held_until", 0))
|
||
except (TypeError, ValueError):
|
||
expired_ago = 0
|
||
if expired_ago > 300:
|
||
try:
|
||
os.remove(path)
|
||
except OSError:
|
||
pass
|
||
continue
|
||
rec["_file"] = name
|
||
out.append(rec)
|
||
return out
|
||
|
||
|
||
def _peer_suppressed(state, sig, now, log):
|
||
"""True when a peer watcher already owns sig; suppress our send.
|
||
|
||
Checks the shared on-disk answered store first (peer answered
|
||
earlier and persisted), then attempts an atomic claim (peer racing
|
||
us in the same window). On suppression the pending prompt is
|
||
reset so later polls re-observe cleanly; answered-store hits are
|
||
already merged into memory, claim-loss is not recorded (peer's
|
||
send owns it, and its claim expires quickly if it dies). Never
|
||
raises; memory-only states never suppress.
|
||
"""
|
||
try:
|
||
state.refresh_answered()
|
||
except Exception:
|
||
pass
|
||
try:
|
||
ts = state.answered_sigs.get(sig)
|
||
if ts is not None and now - ts < ANSWERED_TTL_SECONDS:
|
||
try:
|
||
log.log("info", "duplicate suppressed (peer answered)",
|
||
sig=sig)
|
||
except Exception:
|
||
pass
|
||
state.pending_sig = None
|
||
state.stable_count = 0
|
||
return True
|
||
except Exception:
|
||
pass
|
||
try:
|
||
if not state.try_claim(sig, now):
|
||
try:
|
||
log.log("info", "duplicate suppressed (peer claimed)",
|
||
sig=sig)
|
||
except Exception:
|
||
pass
|
||
state.pending_sig = None
|
||
state.stable_count = 0
|
||
return True
|
||
except Exception:
|
||
pass
|
||
return False
|
||
|
||
|
||
def _check_hold(socket_path, pane_id, state, match, now, log):
|
||
"""Rule-hold gate. Returns None (fresh: evaluate rules), "released"
|
||
(hold over: approve WITHOUT re-evaluating, else expiry would
|
||
re-hold forever), "held", "denied", or "duplicate-suppressed".
|
||
|
||
The holdfile is the single source of truth (crash-safe,
|
||
box-visible): present + fresh => suppress; directive deny => deny
|
||
now; expired or operator-cleared => release back to approve.
|
||
"""
|
||
sig = match["sig"]
|
||
hf = read_hold(socket_path, pane_id)
|
||
if hf is None:
|
||
return None
|
||
if hf.get("sig") != sig:
|
||
clear_hold(socket_path, pane_id) # stale hold for a gone prompt
|
||
return None
|
||
if hf.get("directive") == "approve":
|
||
clear_hold(socket_path, pane_id)
|
||
log.log("info", "hold released by operator, answering", sig=sig)
|
||
return "released"
|
||
if hf.get("directive") == "deny":
|
||
neg = NEGATIVE_KEYS.get(match["kind"])
|
||
if neg is None:
|
||
hf["directive"] = None
|
||
write_hold(socket_path, pane_id, hf)
|
||
log.log("warn", "resolve-deny refused: D2 never denies questions",
|
||
sig=sig, kind=match["kind"])
|
||
return "held"
|
||
if _peer_suppressed(state, sig, now, log):
|
||
return "duplicate-suppressed"
|
||
ok = send_answer(socket_path, pane_id, neg, enter=True)
|
||
clear_hold(socket_path, pane_id)
|
||
state.record_answer(sig, now)
|
||
log.log("info" if ok else "error", "denied %s (operator resolve)" % neg,
|
||
sig=sig, ok=ok, kind=match["kind"], key=neg)
|
||
audit("muse-choice-denied",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id, "sig": sig,
|
||
"ok": ok, "kind": match["kind"], "key": neg,
|
||
"via": "resolve"})
|
||
return "denied"
|
||
try:
|
||
expired = now >= float(hf.get("held_until", 0))
|
||
except (TypeError, ValueError):
|
||
expired = True
|
||
if expired:
|
||
clear_hold(socket_path, pane_id)
|
||
log.log("info", "hold expired, releasing to approve", sig=sig)
|
||
audit("muse-choice-hold-expired",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id, "sig": sig,
|
||
"kind": match["kind"], "rule": hf.get("rule")})
|
||
return "released"
|
||
return "held"
|
||
|
||
|
||
def _poll_once(socket_path, pane_id, state, log, dry_run=False):
|
||
"""Run one poll iteration. Returns an outcome string ("gone" tells the
|
||
caller to exit). Logs only state transitions, never per-poll spam."""
|
||
if not pane_exists(socket_path, pane_id):
|
||
return "gone"
|
||
text = capture_pane(socket_path, pane_id)
|
||
if text is None:
|
||
log.log("warn", "capture failed, continuing")
|
||
return "capture-failed"
|
||
match = find_choice_prompt(text)
|
||
now = time.time()
|
||
gate = _check_hold(socket_path, pane_id, state, match, now, log) \
|
||
if match is not None else None
|
||
if gate in ("held", "denied", "duplicate-suppressed"):
|
||
if gate in ("denied", "duplicate-suppressed"):
|
||
state.pending_sig = None
|
||
state.stable_count = 0
|
||
return gate
|
||
skip_eval = (gate == "released")
|
||
verdict = state.observe(match, now)
|
||
if verdict == "wait" and state.stable_count == 1:
|
||
log.log("info", "prompt seen", kind=match["kind"], key=match["key"],
|
||
sig=match["sig"], text=(match["cue"] or "")[:200])
|
||
return "seen"
|
||
if verdict == "answer":
|
||
log.log("info", "prompt stable, answering", kind=match["kind"],
|
||
key=match["key"], sig=match["sig"])
|
||
# Re-verify the prompt is still present before typing.
|
||
confirm = capture_pane(socket_path, pane_id)
|
||
cmatch = find_choice_prompt(confirm) if confirm else None
|
||
if cmatch is None or cmatch["sig"] != match["sig"]:
|
||
log.log("info", "prompt vanished before answer, skipping",
|
||
sig=match["sig"])
|
||
state.pending_sig = None
|
||
state.stable_count = 0
|
||
return "vanished"
|
||
if launch_opt_out(pane_muse_argv(socket_path, pane_id)):
|
||
log.log("info", "held: pane opted out via launch flags",
|
||
sig=match["sig"], kind=match["kind"])
|
||
state.record_answer(match["sig"], now)
|
||
return "held"
|
||
if skip_eval:
|
||
decision, rule = "approve", None
|
||
via = "hold-released"
|
||
else:
|
||
decision, rule = evaluate_rules(match, log)
|
||
via = "direct"
|
||
key = match["key"]
|
||
if decision == "deny":
|
||
neg = NEGATIVE_KEYS[match["kind"]]
|
||
if dry_run:
|
||
log.log("info", "dry-run would deny %s" % neg,
|
||
sig=match["sig"], kind=match["kind"], key=neg,
|
||
rule=rule["id"] if rule else None)
|
||
state.record_answer(match["sig"], now)
|
||
return "dry-denied"
|
||
if _peer_suppressed(state, match["sig"], now, log):
|
||
return "duplicate-suppressed"
|
||
ok = send_answer(socket_path, pane_id, neg, enter=True)
|
||
state.record_answer(match["sig"], now)
|
||
log.log("info" if ok else "error", "denied %s" % neg,
|
||
sig=match["sig"], ok=ok, kind=match["kind"], key=neg,
|
||
rule=rule["id"] if rule else None)
|
||
audit("muse-choice-denied",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id,
|
||
"sig": match["sig"], "ok": ok, "kind": match["kind"],
|
||
"key": neg, "via": "rule",
|
||
"rule": rule["id"] if rule else None})
|
||
return "denied"
|
||
if decision == "hold":
|
||
until = now + HOLD_WINDOW_SECONDS
|
||
write_hold(socket_path, pane_id, {
|
||
"sig": match["sig"], "kind": match["kind"], "key": key,
|
||
"text": (match["cue"] or "")[:200],
|
||
"rule": rule["id"] if rule else None,
|
||
"reason": (rule.get("reason") if rule else None) or "",
|
||
"downgraded_from": (rule.get("downgraded_from")
|
||
if rule else None),
|
||
"held_until": until, "directive": None,
|
||
"socket": socket_path, "pane": pane_id})
|
||
log.log("info", "held by rule %s" % (rule["id"] if rule else "?"),
|
||
sig=match["sig"], kind=match["kind"],
|
||
rule=rule["id"] if rule else None,
|
||
reason=(rule.get("reason") if rule else None) or "")
|
||
if not dry_run:
|
||
audit("muse-choice-held",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id,
|
||
"sig": match["sig"], "kind": match["kind"],
|
||
"rule": rule["id"] if rule else None,
|
||
"held_until": until})
|
||
return "held"
|
||
if dry_run:
|
||
log.log("info", "dry-run would answer %s" % key,
|
||
sig=match["sig"], kind=match["kind"], key=key,
|
||
options=match["options"], cue=match["cue"])
|
||
state.record_answer(match["sig"], now)
|
||
return "dry-answered"
|
||
if _peer_suppressed(state, match["sig"], now, log):
|
||
return "duplicate-suppressed"
|
||
ok = send_answer(socket_path, pane_id, key,
|
||
enter=match.get("enter", True))
|
||
state.record_answer(match["sig"], now)
|
||
log.log("info" if ok else "error", "answered %s" % key,
|
||
sig=match["sig"], ok=ok, kind=match["kind"], key=key,
|
||
options=match["options"], cue=match["cue"])
|
||
audit("muse-choice-answered",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id,
|
||
"sig": match["sig"], "ok": ok, "kind": match["kind"],
|
||
"key": key, "options": match["options"],
|
||
"cue": match["cue"], "via": via,
|
||
"rule": rule["id"] if rule else None})
|
||
return "answered"
|
||
if verdict == "capped":
|
||
if state.last_capped_sig != match["sig"]:
|
||
state.last_capped_sig = match["sig"]
|
||
log.log("warn", "hourly answer cap reached, holding",
|
||
sig=match["sig"], kind=match["kind"])
|
||
return "capped"
|
||
if verdict == "none":
|
||
state.last_capped_sig = None
|
||
return "none"
|
||
return "waiting"
|
||
|
||
|
||
def _log_posture(socket_path, pane_id, log):
|
||
"""Log the pane's permission posture once at watcher start.
|
||
|
||
The watcher answers with per-choice logging in every mode (that
|
||
trail is the default path and informs policy); a bypass posture
|
||
(yolo / approval disabled) additionally gets a box audit record,
|
||
since the session then makes choices outside the trail. Never
|
||
raises.
|
||
"""
|
||
try:
|
||
posture = muse_approval_flags(pane_muse_argv(socket_path, pane_id))
|
||
except Exception:
|
||
return
|
||
try:
|
||
log.log("info", "pane posture", mode=posture["mode"],
|
||
bypass=posture["bypass"], flags=posture["flags"])
|
||
except Exception:
|
||
pass
|
||
try:
|
||
if posture["bypass"] or posture["mode"] not in ("default", None):
|
||
audit("muse-choice-posture",
|
||
name="%s:%s" % (os.path.basename(socket_path), pane_id),
|
||
extra={"socket": socket_path, "pane": pane_id,
|
||
"mode": posture["mode"],
|
||
"profile": posture["profile"],
|
||
"bypass": posture["bypass"],
|
||
"flags": posture["flags"]})
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def watch_loop(socket_path, pane_id, dry_run=False):
|
||
"""Main daemon loop. Returns only when the pane is gone or signalled."""
|
||
log = WatcherLog(logfile_for(socket_path, pane_id))
|
||
state = WatcherState(
|
||
persist_path=answered_file_for(socket_path, pane_id))
|
||
log.log("info", "watcher started", socket=socket_path, pane=pane_id,
|
||
dry_run=dry_run, pid=os.getpid())
|
||
_log_posture(socket_path, pane_id, log)
|
||
polls = 0
|
||
answers = 0
|
||
last_heartbeat = time.time()
|
||
while True:
|
||
try:
|
||
outcome = _poll_once(socket_path, pane_id, state, log,
|
||
dry_run=dry_run)
|
||
if outcome == "gone":
|
||
log.log("info", "pane gone, exiting", polls=polls,
|
||
answers=answers)
|
||
return 0
|
||
if outcome in ("answered", "dry-answered"):
|
||
answers += 1
|
||
polls += 1
|
||
if time.time() - last_heartbeat >= HEARTBEAT_SECONDS:
|
||
last_heartbeat = time.time()
|
||
log.log("info", "heartbeat", polls=polls, answers=answers,
|
||
pending=state.pending_sig)
|
||
except Exception as e:
|
||
try:
|
||
log.log("error", "poll iteration failed, continuing",
|
||
error="%r" % (e,))
|
||
except Exception:
|
||
pass
|
||
time.sleep(POLL_INTERVAL)
|
||
|
||
|
||
def _pid_alive(pid):
|
||
try:
|
||
os.kill(pid, 0)
|
||
return True
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def _pid_is_watcher(pid):
|
||
"""True if pid looks like our watcher (guards against pid reuse)."""
|
||
try:
|
||
with open("/proc/%d/cmdline" % pid, "rb") as f:
|
||
cmd = f.read().decode(errors="replace")
|
||
return "muse_choice_watcher" in cmd
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def _watch_procs(proc_root="/proc"):
|
||
"""Live processes running this module's foreground `watch` argv.
|
||
|
||
Returns [{"pid", "socket", "pane"}]. Exact argv match only -- it can
|
||
never misidentify another program, so callers may act on the result.
|
||
"""
|
||
out = []
|
||
try:
|
||
pids = os.listdir(proc_root)
|
||
except OSError:
|
||
return out
|
||
me = os.getpid()
|
||
for pid in pids:
|
||
if not pid.isdigit() or int(pid) == me:
|
||
continue
|
||
try:
|
||
with open(os.path.join(proc_root, pid, "cmdline"), "rb") as f:
|
||
parts = f.read().split(b"\0")
|
||
except (OSError, IOError):
|
||
continue
|
||
try:
|
||
args = [p.decode(errors="replace") for p in parts if p]
|
||
except Exception:
|
||
continue
|
||
if len(args) < 6:
|
||
continue
|
||
if not args[1].endswith("muse_choice_watcher.py") or args[2] != "watch":
|
||
continue
|
||
try:
|
||
sock = args[args.index("--socket") + 1]
|
||
pane = args[args.index("--pane") + 1]
|
||
except (ValueError, IndexError):
|
||
continue
|
||
out.append({"pid": int(pid), "socket": sock, "pane": pane})
|
||
return out
|
||
|
||
|
||
def is_running(socket_path, pane_id):
|
||
pidfile = pidfile_for(socket_path, pane_id)
|
||
try:
|
||
with open(pidfile) as f:
|
||
pid = int(f.read().strip())
|
||
except Exception:
|
||
return None
|
||
if _pid_alive(pid) and _pid_is_watcher(pid):
|
||
return pid
|
||
try:
|
||
os.remove(pidfile)
|
||
except OSError:
|
||
pass
|
||
return None
|
||
|
||
|
||
def watcher_alive(socket_path, pane_id):
|
||
"""Pid of the live watcher for a pane, pidfile or orphan.
|
||
|
||
is_running covers the normal pidfile case; _watch_procs catches an
|
||
identical watcher alive with a lost pidfile (/tmp cleaned under
|
||
it). Starters must consult this (not is_running alone) or they
|
||
spawn a second daemon onto the same pane (double answers). Never
|
||
raises; None means no live watcher.
|
||
"""
|
||
try:
|
||
pid = is_running(socket_path, pane_id)
|
||
if pid:
|
||
return pid
|
||
except Exception:
|
||
pass
|
||
try:
|
||
for proc in _watch_procs():
|
||
if proc["socket"] == socket_path and proc["pane"] == pane_id:
|
||
return proc["pid"]
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
|
||
def _daemonize():
|
||
"""Double-fork away from the controlling terminal (survives shell exit)."""
|
||
if os.fork() != 0:
|
||
os._exit(0)
|
||
os.setsid()
|
||
if os.fork() != 0:
|
||
os._exit(0)
|
||
devnull = os.open(os.devnull, os.O_RDWR)
|
||
os.dup2(devnull, 0)
|
||
os.dup2(devnull, 1)
|
||
os.dup2(devnull, 2)
|
||
if devnull > 2:
|
||
os.close(devnull)
|
||
|
||
|
||
_PIDFILE_LOCK_FH = None
|
||
|
||
|
||
def _claim_pidfile(pidfile):
|
||
"""Claim a pidfile for this process. Returns False if another
|
||
starter holds it. Kernel-enforced: the winner holds an exclusive
|
||
flock for its whole lifetime, so concurrent box + timer starts can
|
||
never pile two daemons onto one pane (double answers, '11' in the
|
||
input box). Call _release_pidfile() on exit."""
|
||
global _PIDFILE_LOCK_FH
|
||
me = os.getpid()
|
||
try:
|
||
fh = open(pidfile, "a+")
|
||
except OSError:
|
||
return False
|
||
try:
|
||
fcntl.flock(fh, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||
except (OSError, IOError):
|
||
fh.close()
|
||
return False
|
||
try:
|
||
fh.seek(0)
|
||
content = fh.read().strip()
|
||
other = int(content) if content else None
|
||
except Exception:
|
||
other = None
|
||
if other and other != me and _pid_alive(other) and _pid_is_watcher(other):
|
||
# Live foreign owner (e.g. a daemon from before locking existed,
|
||
# or a pid recycled into a watcher): yield to it.
|
||
try:
|
||
fcntl.flock(fh, fcntl.LOCK_UN)
|
||
except Exception:
|
||
pass
|
||
fh.close()
|
||
return False
|
||
try:
|
||
fh.seek(0)
|
||
fh.truncate()
|
||
fh.write(str(me))
|
||
fh.flush()
|
||
except Exception:
|
||
pass
|
||
if _PIDFILE_LOCK_FH is not None:
|
||
try:
|
||
_PIDFILE_LOCK_FH.close()
|
||
except Exception:
|
||
pass
|
||
_PIDFILE_LOCK_FH = fh
|
||
return True
|
||
|
||
|
||
def _release_pidfile(pidfile):
|
||
"""Drop the claim taken by _claim_pidfile: remove the pidfile only
|
||
if we still own it, then release the lock. Never raises."""
|
||
global _PIDFILE_LOCK_FH
|
||
try:
|
||
with open(pidfile) as f:
|
||
owner = f.read().strip()
|
||
except Exception:
|
||
owner = ""
|
||
if owner == str(os.getpid()):
|
||
try:
|
||
os.remove(pidfile)
|
||
except OSError:
|
||
pass
|
||
if _PIDFILE_LOCK_FH is not None:
|
||
try:
|
||
fcntl.flock(_PIDFILE_LOCK_FH, fcntl.LOCK_UN)
|
||
except Exception:
|
||
pass
|
||
try:
|
||
_PIDFILE_LOCK_FH.close()
|
||
except Exception:
|
||
pass
|
||
_PIDFILE_LOCK_FH = None
|
||
|
||
|
||
def start_watcher(socket_path, pane_id, dry_run=False):
|
||
pid = watcher_alive(socket_path, pane_id)
|
||
if pid:
|
||
return {"ok": False, "status": "already_running", "pid": pid}
|
||
if not pane_exists(socket_path, pane_id):
|
||
return {"ok": False, "status": "no_such_pane"}
|
||
_daemonize()
|
||
# Child continues here with a clean identity: re-exec ourselves as a
|
||
# foreground `watch` process so /proc cmdline always names this module
|
||
# (liveness checks and ps output stay correct no matter who spawned us --
|
||
# box, timer, or direct CLI).
|
||
args = [sys.executable, os.path.abspath(__file__), "watch",
|
||
"--socket", socket_path, "--pane", pane_id]
|
||
if dry_run:
|
||
args.append("--dry-run")
|
||
os.execvp(sys.executable, args)
|
||
os._exit(1) # unreachable; exec replaces the image
|
||
|
||
|
||
def stop_watcher(socket_path, pane_id, timeout=5):
|
||
pid = watcher_alive(socket_path, pane_id)
|
||
if not pid:
|
||
return {"ok": True, "status": "not_running"}
|
||
try:
|
||
os.kill(pid, signal.SIGTERM)
|
||
except ProcessLookupError:
|
||
pass
|
||
deadline = time.time() + timeout
|
||
while time.time() < deadline and _pid_alive(pid):
|
||
time.sleep(0.1)
|
||
if _pid_alive(pid):
|
||
try:
|
||
os.kill(pid, signal.SIGKILL)
|
||
except ProcessLookupError:
|
||
pass
|
||
try:
|
||
os.remove(pidfile_for(socket_path, pane_id))
|
||
except OSError:
|
||
pass
|
||
return {"ok": True, "status": "stopped", "pid": pid}
|
||
|
||
|
||
def muse_panes(socket_path):
|
||
"""Pane ids on a socket whose current command looks like Muse."""
|
||
try:
|
||
r = _tmux(socket_path, "list-panes", "-a", "-F",
|
||
"#{pane_id} #{pane_current_command}", timeout=10)
|
||
if r.returncode != 0:
|
||
return []
|
||
except Exception:
|
||
return []
|
||
out = []
|
||
for line in r.stdout.split("\n"):
|
||
parts = line.strip().split(None, 1)
|
||
if len(parts) == 2 and "muse-bin" in parts[1]:
|
||
out.append(parts[0])
|
||
return out
|
||
|
||
|
||
def _start_detached(socket_path, pane_id, dry_run=False):
|
||
"""Fork off a watcher via start_watcher (which daemonizes further).
|
||
Returns True if a watcher is running for the pane afterwards."""
|
||
pid = os.fork()
|
||
if pid == 0:
|
||
start_watcher(socket_path, pane_id, dry_run=dry_run)
|
||
os._exit(0)
|
||
os.waitpid(pid, 0)
|
||
time.sleep(0.2)
|
||
return watcher_alive(socket_path, pane_id) is not None
|
||
|
||
|
||
def start_all(dry_run=False, sockets=None):
|
||
results = []
|
||
for sock in sockets or KNOWN_SOCKETS:
|
||
if not os.path.exists(sock):
|
||
results.append({"socket": sock, "status": "no_socket"})
|
||
continue
|
||
panes = muse_panes(sock)
|
||
if not panes:
|
||
results.append({"socket": sock, "status": "no_muse_panes"})
|
||
for pane in panes:
|
||
if watcher_alive(sock, pane):
|
||
results.append({"socket": sock, "pane": pane,
|
||
"status": "already_running"})
|
||
continue
|
||
ok = _start_detached(sock, pane, dry_run=dry_run)
|
||
results.append({"socket": sock, "pane": pane,
|
||
"status": "started" if ok else "failed"})
|
||
return results
|
||
|
||
|
||
def stop_all():
|
||
results = []
|
||
stopped_pids = set()
|
||
prefix = FILE_PREFIX + "-"
|
||
try:
|
||
names = os.listdir(STATE_DIR)
|
||
except OSError:
|
||
names = []
|
||
for name in names:
|
||
if not name.startswith(prefix) or not name.endswith(".pid"):
|
||
continue
|
||
try:
|
||
with open(os.path.join(STATE_DIR, name)) as f:
|
||
pid = int(f.read().strip())
|
||
except Exception:
|
||
continue
|
||
if _pid_alive(pid) and _pid_is_watcher(pid):
|
||
try:
|
||
os.kill(pid, signal.SIGTERM)
|
||
except ProcessLookupError:
|
||
pass
|
||
stopped_pids.add(pid)
|
||
results.append({"pidfile": name, "pid": pid, "status": "stopped"})
|
||
try:
|
||
os.remove(os.path.join(STATE_DIR, name))
|
||
except OSError:
|
||
pass
|
||
# "off means off": also stop exact-argv watchers running pidfile-less.
|
||
for proc in _watch_procs():
|
||
if proc["pid"] in stopped_pids:
|
||
continue
|
||
try:
|
||
os.kill(proc["pid"], signal.SIGTERM)
|
||
except ProcessLookupError:
|
||
pass
|
||
results.append({"pidfile": None, "pid": proc["pid"],
|
||
"status": "stopped-orphan", "socket": proc["socket"],
|
||
"pane": proc["pane"]})
|
||
return results
|
||
|
||
|
||
def reconcile(sockets=None):
|
||
"""Enforce desired state: start missing watchers when enabled, stop all
|
||
when disabled, prune dead pidfiles. Audits only when it changed
|
||
something. Returns a summary dict."""
|
||
desired = get_desired()
|
||
started, already, pruned, failed = [], [], [], []
|
||
stopped = []
|
||
if desired["enabled"]:
|
||
for sock in sockets or KNOWN_SOCKETS:
|
||
if not os.path.exists(sock):
|
||
continue
|
||
for pane in muse_panes(sock):
|
||
if watcher_alive(sock, pane):
|
||
already.append("%s:%s" % (sock, pane))
|
||
continue
|
||
if _start_detached(sock, pane, dry_run=desired["dry_run"]):
|
||
started.append("%s:%s" % (sock, pane))
|
||
else:
|
||
failed.append("%s:%s" % (sock, pane))
|
||
for row in status_all():
|
||
if not row["alive"]:
|
||
try:
|
||
os.remove(os.path.join(STATE_DIR, row["pidfile"]))
|
||
pruned.append(row["pidfile"])
|
||
except OSError:
|
||
pass
|
||
else:
|
||
stopped = ["%s:%s" % (r.get("pidfile"), r.get("pid")) for r in stop_all()]
|
||
if started or stopped or pruned or failed:
|
||
audit("muse-choice-reconciled",
|
||
extra={"started": started, "stopped": stopped, "pruned": pruned,
|
||
"failed": failed})
|
||
return {"enabled": desired["enabled"], "dry_run": desired["dry_run"],
|
||
"started": started, "already": already,
|
||
"stopped": stopped, "pruned": pruned, "failed": failed}
|
||
|
||
|
||
def recent_answers(limit=10, scan_lines=3000):
|
||
"""Most recent muse-choice-answered audit records (newest last)."""
|
||
try:
|
||
with open(CTL_LOG) as f:
|
||
lines = f.readlines()[-scan_lines:]
|
||
except OSError:
|
||
return []
|
||
out = []
|
||
for line in lines:
|
||
try:
|
||
rec = json.loads(line)
|
||
except Exception:
|
||
continue
|
||
if rec.get("action") == "muse-choice-answered":
|
||
out.append(rec)
|
||
return out[-limit:]
|
||
|
||
|
||
def status_all():
|
||
rows = []
|
||
prefix = FILE_PREFIX + "-"
|
||
try:
|
||
names = sorted(os.listdir(STATE_DIR))
|
||
except OSError:
|
||
return rows
|
||
for name in names:
|
||
if not name.startswith(prefix) or not name.endswith(".pid"):
|
||
continue
|
||
pid = None
|
||
try:
|
||
with open(os.path.join(STATE_DIR, name)) as f:
|
||
pid = int(f.read().strip())
|
||
except Exception:
|
||
pass
|
||
alive = bool(pid) and _pid_alive(pid) and _pid_is_watcher(pid)
|
||
logname = name[:-4] + ".log"
|
||
logsize = None
|
||
try:
|
||
logsize = os.path.getsize(os.path.join(STATE_DIR, logname))
|
||
except OSError:
|
||
pass
|
||
rows.append({"pidfile": name, "pid": pid, "alive": alive,
|
||
"log": logname, "log_bytes": logsize})
|
||
owned = set(r["pid"] for r in rows if r["alive"] and r["pid"])
|
||
for proc in _watch_procs():
|
||
if proc["pid"] in owned:
|
||
continue
|
||
logname = os.path.basename(
|
||
logfile_for(proc["socket"], proc["pane"]))
|
||
try:
|
||
logsize = os.path.getsize(os.path.join(STATE_DIR, logname))
|
||
except OSError:
|
||
logsize = None
|
||
rows.append({"pidfile": None, "pid": proc["pid"], "alive": True,
|
||
"orphan": True, "socket": proc["socket"],
|
||
"pane": proc["pane"], "log": logname,
|
||
"log_bytes": logsize})
|
||
return rows
|
||
|
||
|
||
def main(argv=None):
|
||
ap = argparse.ArgumentParser(description="Muse A/B/C choice watcher")
|
||
sub = ap.add_subparsers(dest="cmd", required=True)
|
||
|
||
p = sub.add_parser("start", help="start watching one pane (daemonizes)")
|
||
p.add_argument("--socket", required=True)
|
||
p.add_argument("--pane", required=True)
|
||
p.add_argument("--dry-run", action="store_true",
|
||
help="log answers instead of sending keys")
|
||
|
||
p = sub.add_parser("stop", help="stop one pane watcher")
|
||
p.add_argument("--socket", required=True)
|
||
p.add_argument("--pane", required=True)
|
||
|
||
p = sub.add_parser("start-all", help="watch every Muse pane on known sockets")
|
||
p.add_argument("--dry-run", action="store_true")
|
||
p.add_argument("--socket", action="append", dest="sockets", default=None)
|
||
|
||
sub.add_parser("stop-all", help="stop all watchers")
|
||
sub.add_parser("status", help="list watcher state files + liveness")
|
||
|
||
p = sub.add_parser("on", help="enable auto-answers (desired state) + reconcile now")
|
||
p.add_argument("--dry-run", action="store_true",
|
||
help="log answers instead of sending keys")
|
||
p.add_argument("--by", default="cli", help="caller identity for audit")
|
||
|
||
p = sub.add_parser("off", help="disable auto-answers (desired state) + stop all now")
|
||
p.add_argument("--by", default="cli", help="caller identity for audit")
|
||
|
||
p = sub.add_parser("reconcile", help="enforce desired state (for timer)")
|
||
p.add_argument("--socket", action="append", dest="sockets", default=None)
|
||
|
||
p = sub.add_parser("watch", help="run one pane loop in foreground (daemon target)")
|
||
p.add_argument("--socket", required=True)
|
||
p.add_argument("--pane", required=True)
|
||
p.add_argument("--dry-run", action="store_true")
|
||
|
||
p = sub.add_parser("match", help="print JSON verdict for text on stdin")
|
||
p.add_argument("--tail-window", type=int, default=TAIL_WINDOW)
|
||
|
||
p = sub.add_parser("logs", help="tail one pane watcher log")
|
||
p.add_argument("--socket", required=True)
|
||
p.add_argument("--pane", required=True)
|
||
p.add_argument("-n", type=int, default=20)
|
||
|
||
p = sub.add_parser("state", help="print JSON runtime state for one pane")
|
||
p.add_argument("--socket", required=True)
|
||
p.add_argument("--pane", required=True)
|
||
|
||
args = ap.parse_args(argv)
|
||
if args.cmd == "start":
|
||
existing = watcher_alive(args.socket, args.pane)
|
||
if existing:
|
||
print(json.dumps({"ok": True, "status": "already_running",
|
||
"pid": existing,
|
||
"log": logfile_for(args.socket, args.pane)}))
|
||
return 0
|
||
# start_watcher never returns in the daemon child; the parent fork
|
||
# dance happens inside, so wrap: fork here, child calls start.
|
||
pid = os.fork()
|
||
if pid == 0:
|
||
start_watcher(args.socket, args.pane, dry_run=args.dry_run)
|
||
os._exit(0)
|
||
_, status = os.waitpid(pid, 0)
|
||
time.sleep(0.3)
|
||
running = watcher_alive(args.socket, args.pane)
|
||
print(json.dumps({"ok": running is not None, "pid": running,
|
||
"log": logfile_for(args.socket, args.pane)}))
|
||
return 0 if running else 1
|
||
if args.cmd == "stop":
|
||
print(json.dumps(stop_watcher(args.socket, args.pane)))
|
||
return 0
|
||
if args.cmd == "start-all":
|
||
print(json.dumps(start_all(dry_run=args.dry_run, sockets=args.sockets), indent=1))
|
||
return 0
|
||
if args.cmd == "stop-all":
|
||
print(json.dumps(stop_all(), indent=1))
|
||
return 0
|
||
if args.cmd == "status":
|
||
print(json.dumps({"desired": get_desired(), "watchers": status_all()},
|
||
indent=1))
|
||
return 0
|
||
if args.cmd == "on":
|
||
state = set_enabled(True, dry_run=args.dry_run, by=args.by)
|
||
res = reconcile()
|
||
print(json.dumps({"desired": state, "reconcile": res}, indent=1))
|
||
return 0
|
||
if args.cmd == "off":
|
||
state = set_enabled(False, by=args.by)
|
||
stopped = stop_all()
|
||
print(json.dumps({"desired": state, "stopped": stopped}, indent=1))
|
||
return 0
|
||
if args.cmd == "reconcile":
|
||
print(json.dumps(reconcile(sockets=args.sockets), indent=1))
|
||
return 0
|
||
if args.cmd == "watch":
|
||
# Refuse to pile on: an identical watcher may be alive with a lost
|
||
# pidfile (e.g. /tmp cleaned under it). Exact argv match, no guessing.
|
||
for proc in _watch_procs():
|
||
if proc["socket"] == args.socket and proc["pane"] == args.pane:
|
||
return 3
|
||
pidfile = pidfile_for(args.socket, args.pane)
|
||
if not _claim_pidfile(pidfile):
|
||
return 3
|
||
try:
|
||
return watch_loop(args.socket, args.pane, dry_run=args.dry_run)
|
||
finally:
|
||
_release_pidfile(pidfile)
|
||
if args.cmd == "match":
|
||
text = sys.stdin.read()
|
||
print(json.dumps(find_choice_prompt(text, tail_window=args.tail_window), indent=1))
|
||
return 0
|
||
if args.cmd == "logs":
|
||
path = logfile_for(args.socket, args.pane)
|
||
try:
|
||
with open(path) as f:
|
||
lines = f.readlines()
|
||
except OSError as e:
|
||
print("no log: %s (%s)" % (path, e))
|
||
return 1
|
||
for line in lines[-args.n:]:
|
||
print(line.rstrip("\n"))
|
||
return 0
|
||
if args.cmd == "state":
|
||
print(json.dumps(pane_state(args.socket, args.pane), indent=1))
|
||
return 0
|
||
return 1
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|