feat(hybrid-gateway): integrate muse-cli with Cloudflare netns isolation, symmetric sidechat routing, and 646-pip sync unblock

This commit is contained in:
operator
2026-10-04 22:54:01 +00:00
parent 3d5fbe5aeb
commit 1a271b1bbd
60 changed files with 5068 additions and 55 deletions
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
+74
View File
@@ -0,0 +1,74 @@
#!/usr/bin/env python3
"""
box-query.py — Lightweight client to query box.muse-dev.online APIs over HTTPS.
Works inside containers and nodes without direct SSH access:
Uses SSH signature authentication (?identity=bl&ts=...&sig=...) against
the VM Box API.
Usage:
box-query.py timers
box-query.py jobs
box-query.py agents
box-query.py dms [limit]
"""
import argparse
import json
import os
import subprocess
import sys
import time
import urllib.parse
import urllib.request
BOX_API_BASE = os.environ.get("BOX_API_BASE", "https://box.muse-dev.online")
BOX_SIGN_KEY = os.environ.get("BOX_SIGN_KEY", os.path.expanduser("~/.ssh/id_ed25519"))
def sign_request(identity, endpoint):
ts = str(int(time.time()))
payload = f"{ts}\n{endpoint}".encode()
try:
p = subprocess.run(
["ssh-keygen", "-Y", "sign", "-f", BOX_SIGN_KEY, "-n", "box"],
input=payload, capture_output=True, timeout=15)
if p.returncode != 0:
return None, None
return ts, p.stdout.decode()
except Exception:
return None, None
def query(endpoint):
ts, sig = sign_request("bl", endpoint)
if not ts or not sig:
sys.stderr.write("Failed to sign request (missing key or ssh-keygen error)\n")
sys.exit(1)
query_str = urllib.parse.urlencode({
"identity": "bl",
"ts": ts,
"sig": sig
})
url = f"{BOX_API_BASE}/api/box/{endpoint}?{query_str}"
req = urllib.request.Request(
url,
headers={"User-Agent": "NetVM-box-query/1.0 (container)"}
)
with urllib.request.urlopen(req, timeout=15) as resp:
data = resp.read().decode("utf-8")
try:
return json.loads(data)
except Exception:
return data
def main():
p = argparse.ArgumentParser(description="Query box.muse-dev.online APIs via HTTPS signature auth")
p.add_argument("endpoint", choices=["timers", "jobs", "agents", "dms", "health"])
p.add_argument("--json", action="store_true")
args = p.parse_args()
res = query(args.endpoint)
print(json.dumps(res, indent=2))
if __name__ == "__main__":
main()
+129
View File
@@ -0,0 +1,129 @@
#!/usr/bin/env python3
"""Chat rate metric: messages per minute for main chat and each side chat.
Polls muse-chat-api.py for message counts, calculates delta vs previous poll,
logs rates to a time-series file.
Usage: chat-rate.py --account <name> --cdp-port <port> [--interval 60]
"""
import json, subprocess, sys, time, os, argparse
from datetime import datetime, timezone
API = os.path.expanduser("~/Projects/NetVM/bin/muse-chat-api.py")
STATE_FILE = os.path.expanduser("~/Projects/NetVM/logs/chat-rate-state.json")
LOG_FILE = os.path.expanduser("~/Projects/NetVM/logs/chat-rate.log")
def run_api(account, *args):
"""Run muse-chat-api.py and return stdout."""
cmd = ["python3", API, "--account", account] + list(args)
# Note: cdp-port is baked into the account config, not passed here
result = subprocess.run(cmd, capture_output=True, text=True, timeout=60)
return result.stdout
def count_messages(text):
"""Count messages in API output. Messages are separated by '---'."""
# The API outputs messages separated by ---\n
parts = [p.strip() for p in text.split("---") if p.strip()]
# Filter out non-message lines (headers, etc.)
# Messages typically have substantial content
return len([p for p in parts if len(p) > 10])
def get_side_chats(account):
"""List side chat names."""
out = run_api(account, "sidechat", "list")
# Parse: names and timestamps separated by blank lines
# Format: "Side chats\n\n<name>\n\n<timestamp>\n\n<name>\n\n<timestamp>..."
lines = [l.strip() for l in out.split("\n") if l.strip()]
chats = []
# Skip header "Side chats", then pair up (name, timestamp)
lines = [l for l in lines if l != "Side chats"]
# Lines alternate: name, timestamp, name, timestamp...
for i in range(0, len(lines), 2):
if i < len(lines):
name = lines[i]
# Verify next is a timestamp (ends with m/h/d)
if i + 1 < len(lines) and lines[i+1][-1] in "mhd":
chats.append(name)
return chats
def get_chat_count(account, chat_name=None):
"""Get message count for main or a side chat."""
if chat_name:
# Switch to side chat, get messages, switch back
run_api(account, "sidechat", "use", chat_name)
out = run_api(account, "messages")
run_api(account, "sidechat", "main") # switch back
else:
out = run_api(account, "messages")
return count_messages(out)
def main():
p = argparse.ArgumentParser()
p.add_argument("--account", required=True)
p.add_argument("--interval", type=int, default=60, help="poll interval seconds")
p.add_argument("--once", action="store_true", help="single poll, no loop")
args = p.parse_args()
# Load previous state
prev = {}
if os.path.exists(STATE_FILE):
with open(STATE_FILE) as f:
prev = json.load(f)
def poll():
now = datetime.now(timezone.utc).isoformat()
counts = {}
# Main chat
try:
counts["main"] = get_chat_count(args.account)
except Exception as e:
print(f"main: error {e}", file=sys.stderr)
# Side chats
try:
sc_list = get_side_chats(args.account)
print(f"DEBUG: found {len(sc_list)} side chats", file=sys.stderr)
for sc in sc_list:
try:
counts[f"side:{sc}"] = get_chat_count(args.account, sc)
except Exception as e:
print(f"side:{sc}: error {e}", file=sys.stderr)
except Exception as e:
print(f"sidechat list: error {e}", file=sys.stderr)
# Calculate rates
results = []
for chat, count in counts.items():
rate = 0.0
if chat in prev:
prev_count, prev_time = prev[chat]
dt = (datetime.fromisoformat(now) - datetime.fromisoformat(prev_time)).total_seconds() / 60.0
if dt > 0:
rate = (count - prev_count) / dt
results.append((now, chat, count, round(rate, 2)))
prev[chat] = (count, now)
# Log
os.makedirs(os.path.dirname(LOG_FILE), exist_ok=True)
with open(LOG_FILE, "a") as f:
for ts, chat, count, rate in results:
f.write(f"{ts} {chat} count={count} rate={rate}/min\n")
# Save state
with open(STATE_FILE, "w") as f:
json.dump(prev, f)
# Print
for ts, chat, count, rate in results:
print(f"{chat}: {count} msgs, {rate}/min")
if args.once:
poll()
else:
while True:
poll()
time.sleep(args.interval)
if __name__ == "__main__":
main()
+58 -1
View File
@@ -68,9 +68,45 @@ healthy() {
|| { HEALTH_FAIL_REASON="CDP up but no page target in list"; return 1; }
echo "$list" | grep -E -q '"url": "https://muse\.ai' \
|| { HEALTH_FAIL_REASON="CDP up but not on muse.ai"; return 1; }
warp_egress_healthy || return 1
return 0
}
# Stage 5: Warp egress health (2026-10-04). A partitioned browser (tunnel down,
# CDP green) passes stages 1-4 while being unable to reach muse.ai. Check the
# WireGuard handshake age and probe egress from inside the node's netns.
# Sets HEALTH_FAIL_REASON with a distinct "warp egress down" prefix so
# root-cause analysis can distinguish partitions from Chromium crashes.
warp_egress_healthy() {
local wg_out hs_line age_s code
wg_out="$(sudo -n ip netns exec "warp-${PROFILE}" wg show 2>/dev/null)" \
|| { HEALTH_FAIL_REASON="warp egress down (wg show failed in warp-${PROFILE})"; return 1; }
hs_line="$(printf '%s\n' "$wg_out" | grep -i "latest handshake" | head -1)"
[ -n "$hs_line" ] \
|| { HEALTH_FAIL_REASON="warp egress down (no WireGuard handshake in warp-${PROFILE})"; return 1; }
# "latest handshake: 1 minute, 41 seconds ago" -> total seconds
age_s="$(printf '%s\n' "$hs_line" | python3 -c '
import sys, re
s = sys.stdin.read()
m = re.search(r"(\d+)\s*hour", s); h = int(m.group(1)) if m else 0
m = re.search(r"(\d+)\s*minute", s); mi = int(m.group(1)) if m else 0
m = re.search(r"(\d+)\s*second", s); se = int(m.group(1)) if m else 0
print(h*3600 + mi*60 + se)
' 2>/dev/null)"
{ [ -n "$age_s" ] && [ "$age_s" -ge 0 ]; } 2>/dev/null \
|| { HEALTH_FAIL_REASON="warp egress down (unparseable handshake: $hs_line)"; return 1; }
[ "$age_s" -le 180 ] \
|| { HEALTH_FAIL_REASON="warp egress down (handshake ${age_s}s old in warp-${PROFILE})"; return 1; }
code="$(sudo -n ip netns exec "warp-${PROFILE}" curl -m 5 -s -o /dev/null -w "%{http_code}" "https://muse.ai/" 2>/dev/null)" \
|| { HEALTH_FAIL_REASON="warp egress down (egress probe curl failed in warp-${PROFILE})"; return 1; }
# Any 2xx/3xx means we reached muse.ai infra ("/" 307-redirects to the
# auth flow). The probe tests egress connectivity, not page content.
case "$code" in
2*|3*) return 0 ;;
*) HEALTH_FAIL_REASON="warp egress down (egress probe HTTP $code in warp-${PROFILE})"; return 1 ;;
esac
}
if healthy; then
exit 0
fi
@@ -101,6 +137,23 @@ if [ -n "$recent_pid" ]; then
log "browser launched recently (pid $recent_pid), skipping relaunch (probably still starting)"
exit 0
fi
# Warp-partition recovery (2026-10-04): relaunching Chrome cannot fix a dead
# Warp tunnel — the new browser would fail the same egress check and the
# watchdog would loop. If the failure is warp-egress, restart the tunnel first
# (netvm-node-up.sh is idempotent); only fall through to the Chrome relaunch
# if the tunnel does not recover.
case "$HEALTH_FAIL_REASON" in
"warp egress down"*)
log "warp partition detected ($HEALTH_FAIL_REASON), restarting tunnel via netvm-node-up.sh"
sudo -n "$NETVM_BIN/netvm-node-up.sh" "$PROFILE" >>"$LOG" 2>&1 || true
sleep 5
if healthy; then
log "tunnel restart recovered warp egress, chrome relaunch not needed"
exit 0
fi
log "tunnel restart did not recover egress ($HEALTH_FAIL_REASON), proceeding with chrome relaunch"
;;
esac
log "unhealthy ($HEALTH_FAIL_REASON), relaunching chromebox"
# bracket trick so pkill never matches its own command line
pat="profiles/${PROFILE:0:${#PROFILE}-1}[${PROFILE: -1}]/"
@@ -125,9 +178,13 @@ rm -f "/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chro
# muse-chat-api.py uses — a relaunch on the wrong port looks healthy to the
# launcher but is unreachable to the API (observed 2026-10-03: pip relaunched
# on 9278 instead of 9420, watchdog looped on "relaunch FAILED").
# 9>&-: do NOT let the backgrounded launcher inherit the watchdog lock fd.
# Inherited flock fds wedge the lock forever (observed 2026-10-04: opm/muse
# launchers held their profile lock for 100+ min, every later watchdog run
# skipped as "another run in progress" — the watchdog was silently dead).
systemd-run --user --scope --unit="netvm-chrome-${PROFILE}-$(date +%s)" \
"$NETVM_BIN/netvm-chrome.sh" --headless --cdp-port "$CDP_PORT" "$PROFILE" "https://muse.ai" \
>>"$CHROME_LOG" 2>&1 < /dev/null &
>>"$CHROME_LOG" 2>&1 < /dev/null 9>&- &
# Retry with backoff (2026-10-04): cold starts (fresh egress IP, Cloudflare
# handshake) can take >25s for the page title to appear. A single check after
# 25s kills working-but-slow browsers. Try up to 4 times, 15s apart (~60s
+178
View File
@@ -0,0 +1,178 @@
#!/usr/bin/env python3
"""
crypt-server.py — Public cryptographic attestation and key directory for NetVM.
Serves https://crypt.muse-dev.online/ (via cloudflared / reverse proxy).
Endpoints:
GET / -> Service directory / health JSON
GET /health -> Health check
GET /keys/allowed_signers -> OpenSSH allowed_signers formatted file
GET /keys/{identity}.pub -> Individual public key
GET /proofs -> List known proof hashes / work order attestations
GET /proofs/{id} -> Retrieve proof envelope and SSH signature
POST /proofs -> Submit / register a signed proof record
"""
import argparse
import glob
import json
import os
import re
import ssl
import sys
from http.server import HTTPServer, BaseHTTPRequestHandler
REPO_DIR = "/home/super/Projects/NetVM"
SIGNERS_DIR = os.path.join(REPO_DIR, "dm-signers")
PROOFS_DIR = os.path.join(REPO_DIR, "var", "proofs")
os.makedirs(PROOFS_DIR, exist_ok=True)
class CryptHandler(BaseHTTPRequestHandler):
server_version = "crypt-attestation/1.0"
def log_message(self, format, *args):
sys.stderr.write(f"crypt-server: {self.client_address[0]} - {format % args}\n")
def _send(self, code, content, content_type="application/json"):
if isinstance(content, (dict, list)):
body = json.dumps(content, indent=2).encode("utf-8")
elif isinstance(content, str):
body = content.encode("utf-8")
else:
body = bytes(content)
self.send_response(code)
self.send_header("Content-Type", content_type)
self.send_header("Content-Length", str(len(body)))
self.send_header("Access-Control-Allow-Origin", "*")
self.end_headers()
self.wfile.write(body)
def do_GET(self):
path = self.path.split("?")[0].rstrip("/")
if not path:
path = "/"
if path in ("/", "/health"):
self._send(200, {
"service": "crypt.muse-dev.online",
"status": "active",
"mode": "public_attestation",
"endpoints": [
"/keys/allowed_signers",
"/keys/<identity>.pub",
"/proofs",
"/proofs/<id>"
]
})
return
if path == "/keys/allowed_signers":
allowed_path = os.path.join(SIGNERS_DIR, "allowed_signers")
if os.path.exists(allowed_path):
with open(allowed_path, "r") as f:
data = f.read()
self._send(200, data, content_type="text/plain; charset=utf-8")
else:
self._send(404, {"error": "allowed_signers not found"})
return
m_key = re.match(r"^/keys/([a-zA-Z0-9_\-\.]+)\.pub$", path)
if m_key:
ident = m_key.group(1)
pub_path = os.path.join(SIGNERS_DIR, f"{ident}.pub")
if os.path.exists(pub_path):
with open(pub_path, "r") as f:
data = f.read()
self._send(200, data, content_type="text/plain; charset=utf-8")
else:
self._send(404, {"error": f"Public key for {ident} not found"})
return
if path == "/proofs":
proof_files = glob.glob(os.path.join(PROOFS_DIR, "*.json"))
ids = [os.path.basename(p)[:-5] for p in proof_files]
self._send(200, {"proofs": sorted(ids)})
return
m_proof = re.match(r"^/proofs/([a-zA-Z0-9_\-]+)$", path)
if m_proof:
pid = m_proof.group(1)
pf = os.path.join(PROOFS_DIR, f"{pid}.json")
if os.path.exists(pf):
with open(pf, "r") as f:
data = json.load(f)
self._send(200, data)
else:
self._send(404, {"error": f"Proof {pid} not found"})
return
self._send(404, {"error": "not found"})
def do_POST(self):
path = self.path.split("?")[0].rstrip("/")
if path == "/proofs":
try:
length = int(self.headers.get("Content-Length", 0))
except ValueError:
length = 0
if length <= 0 or length > 65536:
self._send(400, {"error": "invalid content length"})
return
try:
data = json.loads(self.rfile.read(length))
except Exception:
self._send(400, {"error": "malformed JSON"})
return
pid = data.get("id")
if not pid or not re.match(r"^[a-zA-Z0-9_\-]+$", pid):
self._send(400, {"error": "missing or invalid proof id"})
return
pf = os.path.join(PROOFS_DIR, f"{pid}.json")
with open(pf, "w") as f:
json.dump(data, f, indent=2)
self._send(201, {"status": "stored", "id": pid, "url": f"https://crypt.muse-dev.online/proofs/{pid}"})
return
self._send(404, {"error": "not found"})
def main():
parser = argparse.ArgumentParser(description="Crypt attestation and key server")
parser.add_argument("--port", type=int, default=8446)
parser.add_argument("--host", default="100.123.153.75")
parser.add_argument("--ssl", action="store_true", help="Enable self-signed HTTPS")
args = parser.parse_args()
server = HTTPServer((args.host, args.port), CryptHandler)
if args.ssl:
cert_file = os.path.join(REPO_DIR, "ssl", "crypt-selfsigned.crt")
key_file = os.path.join(REPO_DIR, "ssl", "crypt-selfsigned.key")
os.makedirs(os.path.dirname(cert_file), exist_ok=True)
if not os.path.exists(cert_file):
import subprocess
subprocess.run([
"openssl", "req", "-x509", "-newkey", "rsa:2048",
"-keyout", key_file, "-out", cert_file,
"-days", "3650", "-nodes",
"-subj", "/CN=crypt.muse-dev.online"
], check=True, capture_output=True)
os.chmod(key_file, 0o600)
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
ctx.load_cert_chain(cert_file, key_file)
server.socket = ctx.wrap_socket(server.socket, server_side=True)
print(f"Crypt server running on {'https' if args.ssl else 'http'}://{args.host}:{args.port}", file=sys.stderr)
try:
server.serve_forever()
except KeyboardInterrupt:
print("\nShutting down crypt server", file=sys.stderr)
if __name__ == "__main__":
main()
+152
View File
@@ -0,0 +1,152 @@
#!/usr/bin/env python3
"""
Side-chat to main-chat work siphon — detection rules.
Monitors side chat messages and identifies "siphon-worthy" content:
work that should surface in main chat for visibility.
Categories:
COMPLETED - work finished, results ready
BLOCKER - something is stuck, needs intervention
DECISION - a decision is needed from the user/operator
ALERT - health/security/urgency signal
MILESTONE - significant progress checkpoint
Detection is purely pattern-based (raw Python, no AI).
Each rule returns (category, confidence, summary) or None.
"""
import re
from dataclasses import dataclass
from typing import Optional
@dataclass
class SiphonHit:
category: str # COMPLETED, BLOCKER, DECISION, ALERT, MILESTONE
confidence: float # 0.0 - 1.0
summary: str # one-line summary for main chat
thread_id: str # source side chat
message_id: str # source message
# Full message text is NOT stored here — main chat gets a summary
# plus a link back, never the full content (safety: no sensitive
# data siphoned verbatim).
# --- Keyword sets (configurable) ---
COMPLETED_PATTERNS = [
re.compile(r'\b(done|completed|finished|deployed|shipped|live|verified)\b', re.I),
re.compile(r'\b(all|tests?)\s+(pass|green|passing)\b', re.I),
re.compile(r'\[RESULT[^\]]*\]\s*OK', re.I),
re.compile(r'\b(merged|committed|pushed|published)\b', re.I),
]
BLOCKER_PATTERNS = [
re.compile(r'\b(blocked|stuck|failing|broken|down|error|failed)\b', re.I),
re.compile(r'\b(need|needs|waiting)\s+(your|approval|input|decision)\b', re.I),
re.compile(r'\b(can\'t|cannot|unable to)\b', re.I),
re.compile(r'\[RESULT[^\]]*\]\s*(FAIL|ERROR)', re.I),
]
DECISION_PATTERNS = [
re.compile(r'\b(should (i|we)|shall i|want me to)\b', re.I),
re.compile(r'\b(your call|needs? your|awaiting your)\b', re.I),
re.compile(r'\b(approve|approval)\b.*\?', re.I),
re.compile(r'^(yes|no)\s*\?\s*$', re.I),
]
ALERT_PATTERNS = [
re.compile(r'\b(security|vulnerability|breach|compromised|exploit)\b', re.I),
re.compile(r'\b(urgent|critical|emergency|asap)\b', re.I),
re.compile(r'\b(502|503|500)\b.*\b(error|down)\b', re.I),
re.compile(r'\b(ssh|tunnel).*\b(down|broken|failed)\b', re.I),
]
MILESTONE_PATTERNS = [
re.compile(r'\b(milestone|phase \d+ (complete|done)|shipped v)\b', re.I),
re.compile(r'\b(all \d+ (items? )?done)\b', re.I),
]
# Patterns that suppress siphoning (safety)
SUPPRESS_PATTERNS = [
re.compile(r'\b(password|secret|token|key|pin)\s*[:=]', re.I),
re.compile(r'-----BEGIN', re.I), # never siphon key material
re.compile(r'\[do not siphon\]', re.I), # explicit opt-out marker
]
def _match_score(text: str, patterns) -> float:
"""Return confidence based on how many patterns match."""
hits = sum(1 for p in patterns if p.search(text))
if hits == 0:
return 0.0
# Diminishing returns: 1 hit = 0.6, 2 = 0.8, 3+ = 0.95
return min(0.95, 0.6 + (hits - 1) * 0.2)
def _extract_summary(text: str, max_len: int = 120) -> str:
"""Extract a safe one-line summary. Strips to first meaningful line."""
# Take first non-empty line, truncate
for line in text.strip().split('\n'):
line = line.strip()
if line and len(line) > 10:
if len(line) > max_len:
return line[:max_len - 3] + '...'
return line
return text[:max_len]
def detect(text: str, thread_id: str, message_id: str,
min_confidence: float = 0.6) -> Optional[SiphonHit]:
"""
Check a side chat message for siphon-worthy content.
Returns SiphonHit or None.
"""
# Safety: suppress sensitive content
for p in SUPPRESS_PATTERNS:
if p.search(text):
return None
candidates = [
("COMPLETED", _match_score(text, COMPLETED_PATTERNS)),
("BLOCKER", _match_score(text, BLOCKER_PATTERNS)),
("DECISION", _match_score(text, DECISION_PATTERNS)),
("ALERT", _match_score(text, ALERT_PATTERNS)),
("MILESTONE", _match_score(text, MILESTONE_PATTERNS)),
]
# Sort by confidence descending; ALERT wins ties (safety: urgency first)
# Use negative confidence for descending, and ALERT as tiebreaker
candidates.sort(key=lambda x: (-x[1], 0 if x[0] == "ALERT" else 1))
best_cat, best_conf = candidates[0]
if best_conf < min_confidence:
return None
return SiphonHit(
category=best_cat,
confidence=best_conf,
summary=_extract_summary(text),
thread_id=thread_id,
message_id=message_id,
)
# --- Opt-out registry ---
_opt_out_threads: set = set()
def opt_out(thread_id: str):
"""Agent opts a side chat out of siphoning."""
_opt_out_threads.add(thread_id)
def opt_in(thread_id: str):
"""Re-enable siphoning for a side chat."""
_opt_out_threads.discard(thread_id)
def is_opted_out(thread_id: str) -> bool:
return thread_id in _opt_out_threads
+36 -4
View File
@@ -43,6 +43,14 @@ import urllib.request
import urllib.parse
from datetime import datetime, timezone, timedelta
# Shared per-agent rate limiter (jittered) — de-correlates fleet sends
# so the 4 nodes don't hit the shared egress IP in lockstep.
try:
from rate_limiter import rate_limit_wait
HAS_RATE_LIMITER = True
except ImportError:
HAS_RATE_LIMITER = False
API = "/home/super/Projects/NetVM/bin/muse-chat-api.py"
NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh"
VALID_AGENTS = ["muse", "pip", "646", "opm"]
@@ -222,8 +230,9 @@ SIDCHAT_ALIASES = {
# autoprovision, which registers the live UUID in job-sidechats.json.
# Do NOT re-add hardcoded UUIDs here; use job-sidechats.json instead.
def resolve_sidechat_target(target):
"""Resolve target alias or name to UUID dynamically from job-sidechats.json."""
def resolve_sidechat_target(target, agent=None):
"""Resolve target alias or name to UUID dynamically from job-sidechats.json.
Supports agent-specific scoped targets (e.g. '646-pip' resolving to recipient's local thread)."""
if not target or not str(target).strip():
raise ValueError(
"DM target must not be empty: specify a sidechat name/UUID, "
@@ -239,8 +248,26 @@ def resolve_sidechat_target(target):
try:
with open(sc_file, "r", encoding="utf-8") as f:
sc_data = json.load(f)
# 1. Check agent-scoped alias first (e.g. target:646-pip, agent:646)
if agent:
scoped_key = f"{target}@{agent}"
if scoped_key in sc_data:
val = sc_data[scoped_key]
if isinstance(val, dict):
return val.get("thread_uuid") or val.get("uuid") or target
elif isinstance(val, str):
return val
# 2. Check general key
val = sc_data.get(target)
if isinstance(val, dict):
# If mapped entry defines an agent, but recipient differs, check for agent match
if agent and val.get("agent") and val.get("agent") != agent:
# Look for an entry explicitly matching this agent
for k, item in sc_data.items():
if isinstance(item, dict) and item.get("agent") == agent and k.startswith(target):
return item.get("thread_uuid") or target
return val.get("thread_uuid") or val.get("uuid") or target
elif isinstance(val, str):
return val
@@ -541,7 +568,7 @@ def dm_send(agent, target, message, verify=True, raw=False,
pre_create_uuid = None # parked thread before autoprovision (creation check)
nav_is_uuid = False
# Resolve well-known aliases or dynamic thread mappings to UUIDs.
nav_target = resolve_sidechat_target(target)
nav_target = resolve_sidechat_target(target, recipient)
if nav_target != target:
log_event({"type": "alias_resolved", "id": msg_id, "target": target, "thread_uuid": nav_target,
"source": resolve_sidechat_source(target)})
@@ -646,6 +673,11 @@ def dm_send(agent, target, message, verify=True, raw=False,
f"actual_url={_ps_detail.get('actual_url')})", file=sys.stderr)
sys.exit(1)
# Per-agent rate limit (jittered): de-correlate fleet sends across the
# shared egress IP. Near no-op at normal pacing (sends already sleep).
if HAS_RATE_LIMITER:
rate_limit_wait(agent)
time.sleep(2)
# Send (raw mode: no truncation — signatures must survive intact)
@@ -830,7 +862,7 @@ def dm_read(agent, target, n=5, quiet=False, width=200):
print(f"ERROR: Unknown agent {agent}", file=sys.stderr)
sys.exit(1)
nav_target = resolve_sidechat_target(target)
nav_target = resolve_sidechat_target(target, agent)
if target == "main":
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main")
else:
+714
View File
@@ -0,0 +1,714 @@
#!/usr/bin/env python3
"""
exec-constrained.py — Constrained HTTPS script-execution endpoint for bl.
SSH-fallback: when SSH to bl is down, operators can still run a constrained
set of scripts (dm.py, job-dispatch, chromebox control) over HTTPS.
This is the hardened replacement for exec-server.py's /exec endpoint, which
executes ARBITRARY shell commands (subprocess.run(cmd, shell=True)). This
server NEVER takes a command string. Clients request a named OPERATION with
validated ARGUMENTS; the server maps op -> fixed argv. No shell. No
interpolation. No RCE.
Usage:
python3 exec-constrained.py --port 8444
Client:
curl -sk -X POST https://127.0.0.1:8444/exec \\
-H 'Authorization: Bearer <token>' \\
-H 'Content-Type: application/json' \\
-d '{"op": "dm.send", "args": {"agent": "opm", "to": "646",
"target": "main", "message": "hello"}}'
Or signature auth (no secret in transit):
payload=$(python3 -c "import json,time,secrets; print(json.dumps({
'op': 'dm.send',
'args': {'agent':'opm','to':'646','target':'main','message':'hi'},
'ts': time.time(), 'nonce': secrets.token_hex(16)}))")
sig=$(printf '%s' "$payload" | ssh-keygen -Y sign -f ~/.ssh/id_frontdoor -n exec-constrained)
curl -sk -X POST https://127.0.0.1:8444/exec \\
-H 'Content-Type: application/json' \\
-d "$(python3 -c "import json,sys; print(json.dumps({
'identity': 'operator-main',
'payload': sys.argv[1], 'signature': sys.argv[2]}))" "$payload" "$sig")"
Auth model: master bearer token + per-agent bearer tokens (same files as
exec-server.py, so migration is drop-in) OR ssh-keygen -Y signatures over
the op envelope (namespace 'exec-constrained'). See SECURITY.md.
"""
import argparse
import hashlib
import hmac
import json
import os
import re
import secrets
import subprocess
import sys
import tempfile
import threading
import time
from http.server import HTTPServer, BaseHTTPRequestHandler
import ssl
# ---------------------------------------------------------------- config
BIN_DIR = '/home/super/Projects/NetVM/bin'
JOBS_DIR = '/home/super/Projects/NetVM/jobs'
TOKEN_FILE = '/home/super/.exec-server-token'
TOKEN_DIR = '/home/super/.exec-tokens'
SIGNERS_FILE = '/home/super/.exec-signers'
NONCE_FILE = '/home/super/.exec-constrained-nonces'
AUDIT_LOG = '/home/super/.exec-constrained-audit.jsonl'
SIG_NAMESPACE = 'exec-constrained'
SIG_MAX_SKEW = 300
AGENTS = ('muse', 'pip', '646', 'opm')
TARGET_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9/_.-]{0,63}$')
JOB_RE = re.compile(r'^[a-z0-9][a-z0-9-]{0,63}$')
IDENT_RE = re.compile(r'^[a-z0-9-]+$')
NONCE_RE = re.compile(r'^[0-9a-fA-F]{16,128}$')
HEX_RE = re.compile(r'^[0-9a-f]{8,128}$')
MAX_BODY = 64 * 1024 # 64KB request cap
MAX_MESSAGE = 2000 # dm.py message cap
MAX_OUTPUT = 256 * 1024 # per-stream output cap
WORK_DIR = '/home/super'
# ------------------------------------------------------- rate limiting
class RateLimiter:
"""Token bucket per identity. Thread-safe."""
def __init__(self, rate_per_min=20, burst=5):
self.rate = rate_per_min / 60.0
self.burst = burst
self._buckets = {}
self._lock = threading.Lock()
def allow(self, ident):
now = time.monotonic()
with self._lock:
tokens, last = self._buckets.get(ident, (self.burst, now))
tokens = min(self.burst, tokens + (now - last) * self.rate)
if tokens >= 1.0:
self._buckets[ident] = (tokens - 1.0, now)
return True
self._buckets[ident] = (tokens, now)
return False
LIMITER = RateLimiter()
# ------------------------------------------------------------- auth
def get_token():
try:
with open(TOKEN_FILE) as f:
return f.read().strip()
except FileNotFoundError:
return None
def check_token(token):
"""Bearer token -> identity label, or None. Never logs the token."""
master = get_token()
if token and master and hmac.compare_digest(token, master):
return 'master'
try:
names = os.listdir(TOKEN_DIR)
except FileNotFoundError:
return None
for name in names:
if not IDENT_RE.fullmatch(name):
continue
p = os.path.join(TOKEN_DIR, name)
if not os.path.isfile(p):
continue
try:
with open(p) as f:
t = f.read().strip()
except OSError:
continue
if t and hmac.compare_digest(token, t):
return name
return None
def check_nonce(nonce):
now = time.time()
fresh = []
try:
with open(NONCE_FILE) as f:
for line in f:
parts = line.split()
if len(parts) != 2:
continue
n, t = parts
try:
if now - float(t) < 2 * SIG_MAX_SKEW:
fresh.append((n, t))
except ValueError:
pass
except FileNotFoundError:
pass
if any(n == nonce for n, _ in fresh):
return False
fresh.append((nonce, str(now)))
try:
with open(NONCE_FILE, 'w') as f:
for n, t in fresh:
f.write(f'{n} {t}\n')
os.chmod(NONCE_FILE, 0o600)
except OSError:
return False
return True
def check_signature(identity, payload, signature):
"""ssh-keygen -Y signature over {"op","args","ts","nonce"} -> identity/None."""
if not IDENT_RE.fullmatch(identity or ''):
return None
try:
data = json.loads(payload)
except Exception:
return None
if not isinstance(data, dict):
return None
op = data.get('op')
args = data.get('args')
ts = data.get('ts')
nonce = data.get('nonce')
if not isinstance(op, str) or op not in OPS:
return None
if not isinstance(args, dict):
return None
if not isinstance(nonce, str) or not NONCE_RE.fullmatch(nonce):
return None
try:
ts = float(ts)
except (TypeError, ValueError):
return None
if abs(time.time() - ts) > SIG_MAX_SKEW:
return None
if not check_nonce(nonce):
return None
sig_path = None
try:
with tempfile.NamedTemporaryFile('w', delete=False, suffix='.sig') as f:
f.write(signature if signature.endswith('\n') else signature + '\n')
sig_path = f.name
p = subprocess.run(
['ssh-keygen', '-Y', 'verify', '-f', SIGNERS_FILE, '-I', identity,
'-n', SIG_NAMESPACE, '-s', sig_path],
input=payload.encode(), capture_output=True, timeout=15)
return identity if p.returncode == 0 else None
except Exception:
return None
finally:
if sig_path:
try:
os.unlink(sig_path)
except OSError:
pass
# ------------------------------------------------------- allowlist
#
# Each op maps to a FIXED script with VALIDATED arguments. The client can
# never influence: the executable path, the subcommand, or any flag name.
# Only whitelisted argument VALUES flow through, each checked below.
#
# To add an op: add an entry here with a build() function. build() receives
# the validated args dict and returns an argv list (no shell). Arg spec is
# enforced by validate() before build() runs.
class OpError(Exception):
pass
def _clean_message(s):
if not isinstance(s, str) or not s.strip():
raise OpError('message must be a non-empty string')
if len(s) > MAX_MESSAGE:
raise OpError(f'message too long (max {MAX_MESSAGE})')
if any(ord(c) < 32 and c not in '\n\t' for c in s):
raise OpError('message contains control characters')
return s
def _agent(v):
if v not in AGENTS:
raise OpError(f'agent must be one of {AGENTS}')
return v
def _target(v):
if not isinstance(v, str) or not TARGET_RE.fullmatch(v):
raise OpError('target must match ^[A-Za-z0-9][A-Za-z0-9/_.-]{0,63}$')
return v
def _opt_int(v, lo, hi, name):
if v is None:
return None
if isinstance(v, bool) or not isinstance(v, int):
raise OpError(f'{name} must be an integer')
if not (lo <= v <= hi):
raise OpError(f'{name} must be {lo}..{hi}')
return v
def _job_name(v):
# Must exist in JOBS_DIR and match the safe pattern (no path traversal).
if not isinstance(v, str) or not JOB_RE.fullmatch(v):
raise OpError('job must match ^[a-z0-9][a-z0-9-]{0,63}$')
path = os.path.join(JOBS_DIR, v + '.json')
if not os.path.isfile(path):
raise OpError('unknown job')
return v
def _dm_send_build(a):
argv = [sys.executable, os.path.join(BIN_DIR, 'dm.py'), 'send',
'--agent', a['agent'], '--target', a['target']]
if a.get('to'):
argv += ['--to', a['to']]
if a.get('expect_reply'):
argv += ['--expect-reply']
if a.get('reply_timeout') is not None:
argv += ['--reply-timeout', str(a['reply_timeout'])]
if a.get('reply_nudges') is not None:
argv += ['--reply-nudges', str(a['reply_nudges'])]
if a.get('reply_escalate'):
argv += ['--reply-escalate', a['reply_escalate']]
if a.get('route'):
argv += ['--route', a['route']]
for t in a.get('tags', []):
argv += ['--tag', t]
argv.append(a['message'])
return argv
def _dm_send_validate(raw):
if not isinstance(raw, dict):
raise OpError('args must be an object')
allowed = {'agent', 'to', 'target', 'message', 'expect_reply',
'reply_timeout', 'reply_nudges', 'reply_escalate',
'route', 'tags'}
for k in raw:
if k not in allowed:
raise OpError(f'unknown arg: {k}')
a = {}
a['agent'] = _agent(raw.get('agent'))
if 'to' in raw and raw['to'] is not None:
a['to'] = _agent(raw['to'])
a['target'] = _target(raw.get('target'))
a['message'] = _clean_message(raw.get('message'))
a['expect_reply'] = bool(raw.get('expect_reply', False))
a['reply_timeout'] = _opt_int(raw.get('reply_timeout'), 60, 604800,
'reply_timeout')
a['reply_nudges'] = _opt_int(raw.get('reply_nudges'), 0, 10,
'reply_nudges')
esc = raw.get('reply_escalate')
if esc is not None:
a['reply_escalate'] = _agent(esc)
route = raw.get('route')
if route is not None:
if not isinstance(route, str) or not TARGET_RE.fullmatch(route):
raise OpError('route must match target pattern')
a['route'] = route
tags = raw.get('tags', [])
if not isinstance(tags, list) or len(tags) > 8:
raise OpError('tags must be a list of at most 8 strings')
clean_tags = []
for t in tags:
if not isinstance(t, str) or len(t) > 120:
raise OpError('tag must be a string <= 120 chars')
# Canonical tag vocabulary only (see DEPLOY-DECISIONS.md).
if not re.fullmatch(r'[a-z0-9_:-]+=[a-zA-Z0-9_.:/-]*', t) and \
not re.fullmatch(r'[a-z0-9_:-]+', t):
raise OpError(f'malformed tag: {t}')
clean_tags.append(t)
a['tags'] = clean_tags
return a
def _dm_thread_validate(raw):
# Same shape as dm.send but uses the thread subcommand.
return _dm_send_validate(raw)
def _dm_thread_build(a):
argv = _dm_send_build(a)
argv[2] = 'thread'
return argv
def _dm_read_validate(raw):
if not isinstance(raw, dict):
raise OpError('args must be an object')
allowed = {'agent', 'target', 'limit'}
for k in raw:
if k not in allowed:
raise OpError(f'unknown arg: {k}')
return {
'agent': _agent(raw.get('agent')),
'target': _target(raw.get('target')),
'limit': _opt_int(raw.get('limit', 20), 1, 100, 'limit') or 20,
}
def _dm_read_build(a):
return [sys.executable, os.path.join(BIN_DIR, 'dm.py'), 'read',
'--agent', a['agent'], '--target', a['target'],
'--limit', str(a['limit'])]
def _job_run_validate(raw):
if not isinstance(raw, dict):
raise OpError('args must be an object')
allowed = {'job'}
for k in raw:
if k not in allowed:
raise OpError(f'unknown arg: {k}')
return {'job': _job_name(raw.get('job'))}
def _job_run_build(a):
return [sys.executable, os.path.join(BIN_DIR, 'job-dispatch.py'), a['job']]
def _chat_messages_validate(raw):
if not isinstance(raw, dict):
raise OpError('args must be an object')
allowed = {'account', 'limit'}
for k in raw:
if k not in allowed:
raise OpError(f'unknown arg: {k}')
return {
'account': _agent(raw.get('account')),
'limit': _opt_int(raw.get('limit', 20), 1, 100, 'limit') or 20,
}
def _chat_messages_build(a):
return [sys.executable, os.path.join(BIN_DIR, 'muse-chat-api.py'),
'--account', a['account'], 'messages', '--limit', str(a['limit'])]
def _chat_send_validate(raw):
if not isinstance(raw, dict):
raise OpError('args must be an object')
allowed = {'account', 'message', 'thread'}
for k in raw:
if k not in allowed:
raise OpError(f'unknown arg: {k}')
a = {
'account': _agent(raw.get('account')),
'message': _clean_message(raw.get('message')),
}
if raw.get('thread') is not None:
th = raw['thread']
if not isinstance(th, str) or not HEX_RE.fullmatch(th):
raise OpError('thread must be a hex uuid')
a['thread'] = th
return a
def _chat_send_build(a):
argv = [sys.executable, os.path.join(BIN_DIR, 'muse-chat-api.py'),
'--account', a['account'], 'send', '--message', a['message']]
if a.get('thread'):
argv += ['--thread', a['thread']]
return argv
def _health_validate(raw):
if raw not in ({}, None):
raise OpError('health.check takes no args')
return {}
def _health_build(a):
return ['/bin/bash', os.path.join(BIN_DIR, 'agent-health.sh'), '--check']
# op -> {validate, build, timeout, side_effecting, description}
OPS = {
'dm.send': {
'validate': _dm_send_validate, 'build': _dm_send_build,
'timeout': 120, 'side_effecting': True,
'desc': 'Send a DM via dm.py (verified delivery)',
},
'dm.thread': {
'validate': _dm_thread_validate, 'build': _dm_thread_build,
'timeout': 120, 'side_effecting': True,
'desc': 'Send a threaded DM via dm.py',
},
'dm.read': {
'validate': _dm_read_validate, 'build': _dm_read_build,
'timeout': 60, 'side_effecting': False,
'desc': 'Read recent DMs (read-only)',
},
'job.run': {
'validate': _job_run_validate, 'build': _job_run_build,
'timeout': 300, 'side_effecting': True,
'desc': 'Run a job from the jobs directory',
},
'chat.messages': {
'validate': _chat_messages_validate, 'build': _chat_messages_build,
'timeout': 60, 'side_effecting': False,
'desc': 'Read recent chat messages (read-only)',
},
'chat.send': {
'validate': _chat_send_validate, 'build': _chat_send_build,
'timeout': 120, 'side_effecting': True,
'desc': 'Send a chat message via muse-chat-api.py',
},
'health.check': {
'validate': _health_validate, 'build': _health_build,
'timeout': 120, 'side_effecting': False,
'desc': 'Run agent-health.sh --check (read-only)',
},
'exec.ping': {
'validate': _health_validate,
'build': lambda a: ['/bin/echo', 'PONG'],
'timeout': 10, 'side_effecting': False,
'desc': 'Canary no-op for watchdogs',
},
}
# identity -> set of ops. 'master' may invoke everything. Unknown identities
# get the read-only subset. Per-agent tokens inherit their agent name as the
# identity; tighten per agent here as needed.
PERMISSIONS = {
'master': set(OPS),
'operator-main': set(OPS),
'operator-646': {'dm.send', 'dm.thread', 'dm.read', 'job.run',
'chat.messages', 'health.check', 'exec.ping'},
'operator-muse': {'dm.send', 'dm.read', 'chat.messages', 'exec.ping'},
'operator-pip': {'dm.send', 'dm.read', 'chat.messages', 'exec.ping'},
'exec-canary': {'exec.ping'},
}
DEFAULT_PERMS = {'dm.read', 'chat.messages', 'health.check', 'exec.ping'}
def permitted(ident, op):
perms = PERMISSIONS.get(ident, DEFAULT_PERMS)
return op in perms
# ------------------------------------------------------- audit log
_audit_lock = threading.Lock()
def audit(entry):
"""Append-only JSONL audit. Never includes tokens or signatures."""
entry = dict(entry)
entry['ts'] = time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())
line = json.dumps(entry, separators=(',', ':')) + '\n'
with _audit_lock:
try:
with open(AUDIT_LOG, 'a') as f:
f.write(line)
os.chmod(AUDIT_LOG, 0o600)
except OSError:
pass
# ------------------------------------------------------- HTTP handler
class Handler(BaseHTTPRequestHandler):
server_version = 'exec-constrained/1.0'
def log_message(self, format, *args):
sys.stderr.write(f'{self.client_address[0]} - {format % args}\n')
def _send(self, code, obj):
body = json.dumps(obj).encode()
self.send_response(code)
self.send_header('Content-Type', 'application/json')
self.send_header('Content-Length', str(len(body)))
self.end_headers()
self.wfile.write(body)
def do_GET(self):
if self.path == '/health':
self._send(200, {'status': 'ok', 'mode': 'constrained'})
elif self.path == '/ops':
# List ops (names + descriptions only; no arg details needed).
self._send(200, {'ops': [
{'op': k, 'desc': v['desc'],
'side_effecting': v['side_effecting']}
for k, v in sorted(OPS.items())]})
else:
self._send(404, {'error': 'not found'})
def do_POST(self):
if self.path == '/exec':
self._handle_exec()
else:
self._send(404, {'error': 'not found'})
def _handle_exec(self):
peer = self.client_address[0]
try:
length = int(self.headers.get('Content-Length', 0))
except ValueError:
length = 0
if length <= 0 or length > MAX_BODY:
self._send(413 if length > MAX_BODY else 400,
{'error': 'bad content length'})
return
try:
data = json.loads(self.rfile.read(length))
except Exception:
self._send(400, {'error': 'invalid json'})
return
if not isinstance(data, dict):
self._send(400, {'error': 'body must be an object'})
return
# --- auth: Bearer <redacted> header (preferred) or legacy body token,
# --- or ssh-keygen signature envelope.
ident = None
op = data.get('op')
args = data.get('args', {})
authz = self.headers.get('Authorization', '')
if authz.startswith('Bearer '):
token = authz[7:].strip()
ident = check_token(token)
elif data.get('token'):
ident = check_token(data['token'])
elif data.get('identity') and data.get('payload') \
and data.get('signature'):
# Signature envelope carries op/args; ignore body's op/args.
ident = check_signature(data['identity'], data['payload'],
data['signature'])
if ident:
try:
env = json.loads(data['payload'])
op, args = env.get('op'), env.get('args', {})
except Exception:
op, args = None, {}
if not ident:
audit({'event': 'auth_failed', 'peer': peer,
'op': str(op)[:64]})
self._send(401, {'error': 'unauthorized'})
return
# --- rate limit
if not LIMITER.allow(ident):
audit({'event': 'rate_limited', 'ident': ident, 'peer': peer,
'op': str(op)[:64]})
self._send(429, {'error': 'rate limited'})
return
# --- op + permission + args
if not isinstance(op, str) or op not in OPS:
self._send(400, {'error': 'unknown op',
'hint': 'GET /ops for the allowlist'})
return
if not permitted(ident, op):
audit({'event': 'forbidden', 'ident': ident, 'peer': peer,
'op': op})
self._send(403, {'error': 'forbidden for this identity'})
return
try:
spec = OPS[op]
clean = spec['validate'](args)
argv = spec['build'](clean)
except OpError as e:
self._send(400, {'error': f'bad args: {e}'})
return
except Exception as e:
self._send(500, {'error': 'internal'})
return
# --- execute (NO shell, argv only)
audit({'event': 'exec_start', 'ident': ident, 'peer': peer,
'op': op,
'args_sha': hashlib.sha256(
json.dumps(clean, sort_keys=True).encode()
).hexdigest()[:16]})
t0 = time.monotonic()
try:
p = subprocess.run(argv, capture_output=True, text=True,
timeout=spec['timeout'], cwd=WORK_DIR)
rc = p.returncode
out = p.stdout[-MAX_OUTPUT:]
err = p.stderr[-MAX_OUTPUT:]
ok = True
except subprocess.TimeoutExpired:
rc, out, err, ok = -1, '', 'timeout', False
except Exception as e:
rc, out, err, ok = -1, '', str(e)[:500], False
dt = round(time.monotonic() - t0, 2)
audit({'event': 'exec_done', 'ident': ident, 'peer': peer,
'op': op, 'rc': rc, 'duration_s': dt, 'ok': ok})
if ok:
self._send(200, {'rc': rc, 'stdout': out, 'stderr': err,
'duration_s': dt})
else:
self._send(500, {'error': err, 'rc': rc})
# ------------------------------------------------------- main
def main():
global TOKEN_FILE, TOKEN_DIR, SIGNERS_FILE, NONCE_FILE, AUDIT_LOG
global BIN_DIR, JOBS_DIR, WORK_DIR
ap = argparse.ArgumentParser()
ap.add_argument('--port', type=int, default=8444)
ap.add_argument('--host', default='127.0.0.1')
ap.add_argument('--token-file', default=TOKEN_FILE)
ap.add_argument('--token-dir', default=TOKEN_DIR)
ap.add_argument('--signers-file', default=SIGNERS_FILE)
ap.add_argument('--nonce-file', default=NONCE_FILE)
ap.add_argument('--audit-log', default=AUDIT_LOG)
ap.add_argument('--bin-dir', default=BIN_DIR)
ap.add_argument('--jobs-dir', default=JOBS_DIR)
ap.add_argument('--cert-file',
default='/home/super/.exec-constrained-cert.pem')
ap.add_argument('--key-file',
default='/home/super/.exec-constrained-key.pem')
ap.add_argument('--work-dir', default='/home/super')
args = ap.parse_args()
TOKEN_FILE, TOKEN_DIR = args.token_file, args.token_dir
SIGNERS_FILE, NONCE_FILE = args.signers_file, args.nonce_file
AUDIT_LOG = args.audit_log
BIN_DIR, JOBS_DIR = args.bin_dir, args.jobs_dir
WORK_DIR = args.work_dir
server = HTTPServer((args.host, args.port), Handler)
cert_file = args.cert_file
key_file = args.key_file
if not os.path.exists(cert_file):
print('Generating self-signed cert...', file=sys.stderr)
subprocess.run([
'openssl', 'req', '-x509', '-newkey', 'rsa:2048',
'-keyout', key_file, '-out', cert_file,
'-days', '3650', '-nodes',
'-subj', '/CN=bl-exec-constrained',
], check=True, capture_output=True)
os.chmod(key_file, 0o600)
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
ctx.load_cert_chain(cert_file, key_file)
server.socket = ctx.wrap_socket(server.socket, server_side=True)
print(f'Constrained exec listening on https://{args.host}:{args.port}/exec',
file=sys.stderr)
print(f'Allowlist: {", ".join(sorted(OPS))}', file=sys.stderr)
try:
server.serve_forever()
except KeyboardInterrupt:
print('\nShutting down', file=sys.stderr)
if __name__ == '__main__':
main()
+377
View File
@@ -0,0 +1,377 @@
#!/usr/bin/env python3
"""
HTTPS exec server for operator remote command execution.
Runs on bl (stable), accepts authenticated POST /exec, returns command output.
Usage:
python3 exec-server.py --port 8443 --token-file /home/super/.exec-token
646 (or any operator) can then:
curl -k -X POST https://100.123.153.75:8443/exec \
-H "Content-Type: application/json" \
-d '{"cmd": "dm.py send --agent 646 --to opm --target main \"hello\"", "token": "..."}'
Via the VM Caddy (no tailnet needed from the agent's container):
curl -k -X POST https://34-139-37-135.sslip.io/exec/exec \
-H "Content-Type: application/json" \
-d '{"cmd": "...", "token": "<per-agent token>"}'
Signature auth (preferred for agents): no secret crosses the wire at all.
The agent signs {"cmd","ts","nonce"} with their registered SSH key
(`ssh-keygen -Y sign -n exec-server`) and posts
{"identity","payload","signature"}. Verified against the signers file
with `ssh-keygen -Y verify`; ts must be within 300s and the nonce unused.
This is the same identity primitive as signed board posts, and it never
trips secret-handling guardrails because a signature is not a secret.
Auth: the master token (TOKEN_FILE) plus per-agent tokens, one file per
agent under TOKEN_DIR (0600). Each agent's token is individually revocable
by deleting its file. The using identity is logged (never the token).
"""
import argparse
import hashlib
import hmac
import json
import secrets
import subprocess
import sys
from http.server import HTTPServer, BaseHTTPRequestHandler
import ssl
# Token file - generated on first run if not exists
TOKEN_FILE = '/home/super/.exec-server-token'
# Per-agent tokens: TOKEN_DIR/<agent> contains that agent's token
TOKEN_DIR = '/home/super/.exec-tokens'
# ssh-keygen signature auth: public keys of fleet identities, one per line
# ("<identity> ssh-ed25519 AAAA..."), synced from the VM's allowed_signers.
SIGNERS_FILE = '/home/super/.exec-signers'
# Seen nonces for replay protection ("<nonce> <ts>" per line).
NONCE_FILE = '/home/super/.exec-nonces'
SIG_NAMESPACE = 'exec-server'
SIG_MAX_SKEW = 300 # seconds; nonces remembered for 2x this
def get_token():
"""Load or generate the auth token."""
try:
with open(TOKEN_FILE, 'r') as f:
return f.read().strip()
except FileNotFoundError:
token = secrets.token_hex(32)
with open(TOKEN_FILE, 'w') as f:
f.write(token)
# Secure permissions
import os
os.chmod(TOKEN_FILE, 0o600)
print(f'Generated new token in {TOKEN_FILE}', file=sys.stderr)
return token
def check_token(token):
"""Check token against the master token and per-agent tokens.
Returns the identity label ('master' or the agent filename), or None."""
if token and hmac.compare_digest(token, get_token()):
return 'master'
import os
try:
names = os.listdir(TOKEN_DIR)
except FileNotFoundError:
return None
for name in names:
p = os.path.join(TOKEN_DIR, name)
if not os.path.isfile(p):
continue
try:
with open(p) as f:
t = f.read().strip()
except OSError:
continue
if t and hmac.compare_digest(token, t):
return name
return None
def check_nonce(nonce):
"""True if the nonce was never used; records it. Prunes expired entries."""
import os
import time
now = time.time()
fresh = []
try:
with open(NONCE_FILE) as f:
for line in f:
parts = line.split()
if len(parts) != 2:
continue
n, t = parts
try:
if now - float(t) < 2 * SIG_MAX_SKEW:
fresh.append((n, t))
except ValueError:
pass
except FileNotFoundError:
pass
if any(n == nonce for n, _ in fresh):
return False
fresh.append((nonce, str(now)))
try:
with open(NONCE_FILE, 'w') as f:
for n, t in fresh:
f.write(f"{n} {t}\n")
os.chmod(NONCE_FILE, 0o600)
except OSError:
return False
return True
def check_signature(identity, payload, signature):
"""Verify an `ssh-keygen -Y` signature over the payload envelope.
Returns the identity on success, None on failure. The envelope must be
JSON {"cmd","ts","nonce"} with a fresh ts and an unused nonce."""
import json as _json
import os
import re
import subprocess
import tempfile
import time
if not re.fullmatch(r'[a-z0-9-]+', identity or ''):
return None
try:
data = _json.loads(payload)
except Exception:
return None
if not isinstance(data, dict):
return None
cmd = data.get('cmd')
ts = data.get('ts')
nonce = data.get('nonce')
if not isinstance(cmd, str) or not cmd:
return None
if not isinstance(nonce, str) or not re.fullmatch(r'[0-9a-fA-F]{16,128}', nonce):
return None
try:
ts = float(ts)
except (TypeError, ValueError):
return None
if abs(time.time() - ts) > SIG_MAX_SKEW:
return None
if not check_nonce(nonce):
return None
sig_path = None
try:
with tempfile.NamedTemporaryFile('w', delete=False, suffix='.sig') as f:
f.write(signature if signature.endswith('\n') else signature + '\n')
sig_path = f.name
p = subprocess.run(
['ssh-keygen', '-Y', 'verify', '-f', SIGNERS_FILE, '-I', identity,
'-n', SIG_NAMESPACE, '-s', sig_path],
input=payload.encode(), capture_output=True, timeout=15)
return identity if p.returncode == 0 else None
except Exception:
return None
finally:
if sig_path:
try:
os.unlink(sig_path)
except OSError:
pass
class ExecHandler(BaseHTTPRequestHandler):
def log_message(self, format, *args):
# Quiet logging, just to stderr
sys.stderr.write(f'{self.client_address[0]} - {format % args}\n')
def do_POST(self):
# Read body
content_length = int(self.headers.get('Content-Length', 0))
if content_length > 1024 * 1024: # 1MB max
self.send_response(413)
self.end_headers()
return
body = self.rfile.read(content_length)
try:
data = json.loads(body)
except json.JSONDecodeError:
self.send_response(400)
self.end_headers()
self.wfile.write(b'{"error": "invalid json"}')
return
if self.path == '/exec/rotate':
self.handle_rotate(data)
return
if self.path != '/exec':
self.send_response(404)
self.end_headers()
return
# Auth: bearer token OR ssh-keygen -Y signature (no secret in transit).
# Bearer: {"token": "...", "cmd": "..."}.
# Signature: {"identity": "...", "payload": "{\"cmd\":...,\"ts\":...,\"nonce\":...}",
# "signature": "<ssh-keygen -Y armor>"}. check_signature
# validates the envelope; cmd comes from the signed payload.
ident = None
cmd = ''
token = data.get('token', '')
if token:
ident = check_token(token)
cmd = data.get('cmd', '')
elif data.get('identity') and data.get('payload') and data.get('signature'):
ident = check_signature(data['identity'], data['payload'], data['signature'])
if ident:
try:
cmd = json.loads(data['payload']).get('cmd', '')
except Exception:
cmd = ''
if not ident:
self.send_response(401)
self.end_headers()
self.wfile.write(b'{"error": "unauthorized"}')
return
sys.stderr.write(f'exec as {ident} from {self.client_address[0]}\n')
# Get command
if not cmd or not isinstance(cmd, str):
self.send_response(400)
self.end_headers()
self.wfile.write(b'{"error": "missing cmd"}')
return
# Execute (with timeout)
timeout = min(data.get('timeout', 60), 300) # max 5 min
try:
result = subprocess.run(
cmd,
shell=True,
capture_output=True,
text=True,
timeout=timeout,
cwd='/home/super'
)
response = {
'stdout': result.stdout,
'stderr': result.stderr,
'rc': result.returncode,
}
except subprocess.TimeoutExpired:
response = {'error': 'timeout', 'rc': -1}
except Exception as e:
response = {'error': str(e), 'rc': -1}
# Send response
resp_body = json.dumps(response).encode()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_header('Content-Length', str(len(resp_body)))
self.end_headers()
self.wfile.write(resp_body)
def handle_rotate(self, data):
"""POST /exec/rotate - replace an agent's token with a fresh one.
The old token dies immediately; the new one is returned ONLY in the
response body (never logged). Lets an agent bootstrap from a
trust-root-delivered token and end up with one nobody else knows.
Agents may rotate only their own token; master may rotate any
agent's (not its own - that stays manual on bl)."""
import os
import re
token = data.get('token', '')
ident = check_token(token)
if not ident:
self.send_response(401)
self.end_headers()
self.wfile.write(b'{"error": "unauthorized"}')
return
target = data.get('agent', ident)
if not re.fullmatch(r'[a-z0-9-]+', target or ''):
self.send_response(400)
self.end_headers()
self.wfile.write(b'{"error": "bad agent name"}')
return
if target == 'master' or (ident != 'master' and target != ident):
self.send_response(403)
self.end_headers()
self.wfile.write(b'{"error": "forbidden"}')
return
path = os.path.join(TOKEN_DIR, target)
if not os.path.isfile(path):
self.send_response(404)
self.end_headers()
self.wfile.write(b'{"error": "no such agent token"}')
return
new_token = secrets.token_hex(32)
tmp = path + '.tmp'
with open(tmp, 'w') as f:
f.write(new_token + '\n')
os.chmod(tmp, 0o600)
os.replace(tmp, path)
sys.stderr.write(f'token rotated for {target} by {ident}\n')
resp_body = json.dumps({'token': new_token}).encode()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_header('Content-Length', str(len(resp_body)))
self.end_headers()
self.wfile.write(resp_body)
def do_GET(self):
if self.path == '/health':
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.end_headers()
self.wfile.write(b'{"status": "ok"}')
else:
self.send_response(404)
self.end_headers()
def main():
global TOKEN_FILE, TOKEN_DIR
parser = argparse.ArgumentParser()
parser.add_argument('--port', type=int, default=8443)
parser.add_argument('--host', default='0.0.0.0')
parser.add_argument('--token-file', default='/home/super/.exec-server-token')
parser.add_argument('--token-dir', default='/home/super/.exec-tokens')
parser.add_argument('--signers-file', default='/home/super/.exec-signers')
parser.add_argument('--nonce-file', default='/home/super/.exec-nonces')
args = parser.parse_args()
TOKEN_FILE = args.token_file
TOKEN_DIR = args.token_dir
global SIGNERS_FILE, NONCE_FILE
SIGNERS_FILE = args.signers_file
NONCE_FILE = args.nonce_file
# Ensure token exists
token = get_token()
print(f'Token: {token[:8]}... (full in {TOKEN_FILE})', file=sys.stderr)
server = HTTPServer((args.host, args.port), ExecHandler)
# Wrap with TLS (self-signed is fine for our use, we use -k)
# Generate self-signed cert if not exists
import os
cert_file = '/home/super/.exec-server-cert.pem'
key_file = '/home/super/.exec-server-key.pem'
if not os.path.exists(cert_file):
print('Generating self-signed cert...', file=sys.stderr)
subprocess.run([
'openssl', 'req', '-x509', '-newkey', 'rsa:2048',
'-keyout', key_file, '-out', cert_file,
'-days', '3650', '-nodes',
'-subj', '/CN=bl-exec-server'
], check=True, capture_output=True)
os.chmod(key_file, 0o600)
context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
context.load_cert_chain(cert_file, key_file)
server.socket = context.wrap_socket(server.socket, server_side=True)
print(f'Exec server listening on https://{args.host}:{args.port}/exec', file=sys.stderr)
print('Health check: https://<host>:<port>/health', file=sys.stderr)
try:
server.serve_forever()
except KeyboardInterrupt:
print('\nShutting down', file=sys.stderr)
if __name__ == '__main__':
main()
+55
View File
@@ -0,0 +1,55 @@
#!/bin/bash
# exec-sign.sh — call the bl exec-constrained server with SSH-signature auth.
# No bearer token, no secret crosses the wire: you sign the request envelope
# with your registered fleet key and the server verifies it against the
# signers file. A signature is not a secret, so this never trips
# secret-handling guardrails.
#
# Server: exec-constrained.py — named ops ONLY, no arbitrary shell.
# Signature namespace: exec-constrained
# Envelope: {"op","args","ts","nonce"} — op must be in the server allowlist:
# dm.send, dm.thread, dm.read, job.run, chat.messages, chat.send,
# health.check, exec.ping
# Nonce: 16+ hex chars, replay-protected server-side. ts: unix epoch, ±300s skew.
#
# Usage: exec-sign.sh <op> '<args-json>' [identity] [keyfile] [url]
# op a named op, e.g. exec.ping
# args-json JSON object of the op's arguments, e.g. '{}'
# identity defaults to operator-646 (must be a principal in the
# server's signers file)
# keyfile defaults to ~/.ssh/id_frontdoor
# url defaults to https://exec.muse-dev.online/exec (cloudflared).
# NOTE: the old VM Caddy /exec route to bl:8443 died with
# exec-server.py — do not point this at the sslip.io URL.
set -euo pipefail
OP="${1:?usage: exec-sign.sh <op> '<args-json>' [identity] [keyfile] [url]}"
# NOTE: do NOT write this as ${2:-{}} — bash matches the first } as the
# expansion's close brace and appends a literal } when $2 is set.
ARGS_JSON="${2-}"
if [ -z "$ARGS_JSON" ]; then ARGS_JSON='{}'; fi
IDENTITY="${3:-operator-646}"
KEY="${4:-$HOME/.ssh/id_frontdoor}"
URL="${5:-https://exec.muse-dev.online/exec}"
TS=$(date +%s)
NONCE=$(python3 -c "import secrets; print(secrets.token_hex(16))")
PAYLOAD=$(python3 -c "
import json, sys
op, args_json, ts, nonce = sys.argv[1:5]
args = json.loads(args_json)
if not isinstance(args, dict):
raise SystemExit('args-json must be a JSON object')
print(json.dumps({'op': op, 'args': args, 'ts': int(ts), 'nonce': nonce}))
" "$OP" "$ARGS_JSON" "$TS" "$NONCE")
SIG=$(printf '%s' "$PAYLOAD" | ssh-keygen -Y sign -f "$KEY" -n exec-constrained)
BODY=$(python3 -c "
import json, sys
ident, payload, sig = sys.argv[1:4]
print(json.dumps({'identity': ident, 'payload': payload, 'signature': sig}))
" "$IDENTITY" "$PAYLOAD" "$SIG")
# Cloudflare Bot Fight Mode blocks python-urllib POSTs (error 1010):
# use curl with a browser User-Agent instead.
curl -sS -X POST "$URL" \
-H 'Content-Type: application/json' \
-H 'User-Agent: Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/126.0.0.0 Safari/537.36' \
--data "$BODY"
+59
View File
@@ -0,0 +1,59 @@
#!/bin/bash
# exec-watch.sh — 15-min watchdog for the bl exec server's signature-auth path.
# Signs a canary op as the exec-canary identity and POSTs it to the
# local exec-constrained server. Records result in /home/super/.exec-watch-status.json
# and appends to logs/exec-watch.log. Failures are pull-based (status file +
# log); wire push alerting here if the fleet wants paging.
#
# Runs via systemd user timer exec-watch.timer (OnUnitActiveSec=15min).
# Server: exec-constrained.py — named ops ONLY, no arbitrary shell;
# signature namespace: exec-constrained, with {op,args,ts,nonce} envelope.
# NOTE: signers file (/home/super/.exec-signers) is synced MANUALLY from the
# VM's /srv/board/allowed_signers on identity renames/adds — bl cannot ssh
# back to the VM, so there is no pull sync. Operator step, documented in
# docs/TOKEN_POLICY.md.
set -uo pipefail
KEY=/home/super/.exec-canary
URL=https://100.123.153.75:8444/exec
STATUS=/home/super/.exec-watch-status.json
LOG=/home/super/Projects/NetVM/logs/exec-watch.log
ts=$(date +%s)
nonce=$(python3 -c "import secrets; print(secrets.token_hex(16))")
payload=$(python3 -c "import json,sys; print(json.dumps({'op':'exec.ping','args':{},'ts':int(sys.argv[1]),'nonce':sys.argv[2]}))" "$ts" "$nonce")
sig=$(printf '%s' "$payload" | ssh-keygen -Y sign -f "$KEY" -n exec-constrained 2>/dev/null)
if [ -z "${sig:-}" ]; then
result="sign-failed"
else
body=$(python3 -c "import json,sys; print(json.dumps({'identity': 'exec-canary', 'payload': sys.argv[1], 'signature': sys.argv[2]}))" "$payload" "$sig")
out=$(curl -sk -m 25 -X POST "$URL" -H 'Content-Type: application/json' -d "$body" 2>/dev/null)
if echo "$out" | grep -q '"rc": 0'; then
result="ok"
else
result="bad-response"
fi
fi
now_iso=$(date -u +%FT%TZ)
{
python3 - "$STATUS" "$now_iso" "$result" <<'PYEOF'
import json, sys
status_path, now_iso, result = sys.argv[1], sys.argv[2], sys.argv[3]
try:
st = json.load(open(status_path))
except Exception:
st = {}
if result == "ok":
st.update({"last_ok": now_iso, "last_fail": None,
"consecutive_failures": 0, "result": "ok"})
else:
st.update({"last_ok": st.get("last_ok"),
"last_fail": now_iso,
"consecutive_failures": st.get("consecutive_failures", 0) + 1,
"result": result})
json.dump(st, open(status_path, "w"))
print(f"[{now_iso}] exec-watch: {result} "
f"(consecutive_failures={st['consecutive_failures']})")
PYEOF
} >> "$LOG" 2>&1
[ "$result" = "ok" ]
+8
View File
@@ -0,0 +1,8 @@
#!/bin/bash
# flap-check.sh - count browser relaunches per profile in the last hour
# Called by the browser-flap-detector cron. Avoids nested SSH quoting hell.
cutoff=$(date -u -d "1 hour ago" +%Y-%m-%dT%H:%M:%S)
for prof in muse pip 646 opm; do
n=$(grep "\[$prof\]" /home/super/Projects/NetVM/chromebox-watchdog.log 2>/dev/null | grep "relaunch OK" | awk -v d="$cutoff" '{ts=substr($1,2,19); if (ts > d) c++} END {print c+0}')
echo "$prof:$n"
done
+54
View File
@@ -185,6 +185,60 @@ HEALTHY_AGENTS=""
esac
done
# --- Warp partition detection (2026-10-04) ---
# A partitioned node has a live browser + CDP but no internet egress: the
# chromebox watchdog sees a healthy browser while all automation fails.
# Condition id: partition:<node>. The detail names the node, the WireGuard
# handshake age, the egress probe result, and the timestamp, and says
# PARTITION explicitly so #lobby readers can tell a network partition from
# a browser crash at a glance. Anti-spam comes from the shared consecutive-
# failure state machine (2 consecutive failures before first page, re-page
# at most every 30 min).
WARP_PROBE_URL="${WARP_PROBE_URL:-https://1.1.1.1/cdn-cgi/trace}"
warp_partition_probe() { # <node> -> prints "<handshake_age_s|unknown> <ok|FAIL>"
local node="$1" iface epoch now age_s
iface=$(sudo -n ip netns exec "warp-$node" sh -c 'wg show interfaces 2>/dev/null | head -1')
now=$(date +%s)
epoch=$(sudo -n ip netns exec "warp-$node" wg show "$iface" latest-handshakes 2>/dev/null | awk '{print $2}')
case "$epoch" in ''|*[!0-9]*) age_s="unknown" ;; *) age_s=$(( now - epoch )) ;; esac
if sudo -n ip netns exec "warp-$node" curl -s -m 8 -o /dev/null "$WARP_PROBE_URL" 2>/dev/null; then
echo "$age_s ok"
else
echo "$age_s FAIL"
fi
}
"$BIN/netvm-registry.py" 2>/dev/null | while IFS=: read -r node port; do
[ -n "$node" ] && [ -n "$port" ] || continue
cond="partition:$node"
read -r hs_age probe_res < <(warp_partition_probe "$node")
ts=$(date -u +%FT%TZ)
if [ "$probe_res" = "ok" ]; then
failing=0
detail="warp egress restored for $node at $ts (probe $WARP_PROBE_URL ok)"
else
failing=1
detail="PARTITION $node: warp egress down at $ts (handshake ${hs_age}s ago, probe $WARP_PROBE_URL FAILED)"
fi
if injected "$cond"; then
failing=1
detail="PARTITION $node: warp egress down at $ts (handshake ${hs_age}s ago, probe $WARP_PROBE_URL FAILED) [INJECTED]"
fi
read -r action fails < <(state_machine "$cond" "$failing")
case "$action" in
ALERT_FIRST|ALERT_REALERT)
emit_record "ALERT" "$cond" "$detail" "$fails"
echo "$cond" >> "$STATE_DIR/.alerts.tmp"
;;
RECOVERY)
emit_record "RECOVERY" "$cond" "$detail" "$fails"
;;
SUPPRESSED)
log "$cond still critical x$fails — re-page suppressed by quiet hours ($QUIET_HOURS)"
;;
esac
done
# Recompute healthy agents in the main shell (pipeline subshell above can't export).
HEALTHY_AGENTS=""
"$BIN/netvm-registry.py" 2>/dev/null | while IFS=: read -r node port; do
+12
View File
@@ -0,0 +1,12 @@
#!/bin/bash
# fleet-status.sh - one-line health per profile: browser up/down + relay code
declare -A ports=( [muse]=9410 [pip]=9420 [646]=9430 [opm]=9440 )
declare -A veth=( [muse]=10.201.35.2 [pip]=10.201.87.2 [646]=10.201.202.2 [opm]=10.201.157.2 )
for p in muse pip 646 opm; do
port=${ports[$p]}
if pgrep -f "remote-debugging-port=$port" >/dev/null 2>&1; then b=UP; else b=DOWN; fi
r=$(curl -s -m 4 -o /dev/null -w "%{http_code}" "http://${veth[$p]}:$port/json/version" 2>/dev/null || echo 000)
echo "$p: browser=$b relay=$r"
done
echo ===
tail -20 /home/super/Projects/NetVM/chromebox-watchdog.log | grep -E "FAILED|unhealthy" | tail -4
+10 -1
View File
@@ -416,6 +416,7 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
# 2. Extract tracked DMs from dm-log.jsonl
answers_map = {} # id or ref -> answer event
dm_events = []
failed_ids = set() # DM ids that failed to send - never delivered
if os.path.exists(DM_LOG_FILE):
try:
with open(DM_LOG_FILE, "r") as f:
@@ -429,7 +430,8 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
msg = entry.get("msg") or ""
# Check for answers/results/acks
if "[RESULT" in msg or "[ACK" in msg or "acknowledged" in msg or ev_type == "verified":
# Note: "verified" means delivered, NOT answered - do not include it here
if "[RESULT" in msg or "[ACK" in msg or "acknowledged" in msg:
sender = entry.get("agent", "")
# Extract referenced DM id if present
m_ref = re.search(r"\[ref:([a-f0-9-]+)\]", msg) or re.search(r"\[(?:ACK|RESULT)\s+([a-f0-9-]+)", msg)
@@ -438,6 +440,11 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
if entry.get("id"):
answers_map[entry["id"]] = entry
# Track failed sends - these were never delivered
if ev_type in ("send_failed", "failed"):
if entry.get("id"):
failed_ids.add(entry.get("id"))
tags = entry.get("tags") or {}
if tags.get("reply:expected") or "[reply:expected]" in msg:
dm_events.append(entry)
@@ -453,6 +460,8 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
continue
if did in loops:
continue # Already have active record from followups.json
if did in failed_ids:
continue # Send failed - never delivered, don't count against agent health
tags = entry.get("tags") or {}
recipient = entry.get("to") or entry.get("recipient") or ""
+296
View File
@@ -0,0 +1,296 @@
#!/usr/bin/env python3
"""Identity variance testing framework.
Sends gentle, controlled probes from each NetVM node's Warp identity and
measures per-identity response behavior so that future rate limits can be
attributed (per-identity vs per-IP vs time-based vs random).
Design goals:
- Very gentle traffic: 2 probes per node per cycle, default 10-min cycle.
- Never logs secrets: only sha256 hashes of WireGuard private keys.
- JSONL results log: one line per probe, machine-readable.
Usage:
identity-variance-test.py probe # run one probe cycle over all nodes
identity-variance-test.py analyze [--since HOURS] [--log PATH]
Probes per node:
1. https://www.cloudflare.com/cdn-cgi/trace (identity + egress IP Cloudflare sees)
2. https://muse.ai/ (production-relevant landing page)
Exit codes: 0 ok, 1 partial (some nodes failed), 2 fatal.
"""
import argparse
import hashlib
import json
import os
import re
import subprocess
import sys
import time
from datetime import datetime, timezone
DEFAULT_LOG = os.path.expanduser(
"~/Projects/NetVM/.state/identity-variance/variance.jsonl"
)
NODES = ["muse", "pip", "646", "opm", "def"]
PROBES = [
("cf-trace", "https://www.cloudflare.com/cdn-cgi/trace"),
("muse-landing", "https://muse.ai/"),
]
CURL_TIMEOUT = 15
def sh(cmd, timeout=30):
"""Run cmd (list) and return (rc, stdout, stderr)."""
try:
p = subprocess.run(
cmd, capture_output=True, text=True, timeout=timeout
)
return p.returncode, p.stdout, p.stderr
except subprocess.TimeoutExpired as e:
return 124, (e.stdout or ""), "timeout"
except FileNotFoundError:
return 127, "", "command not found"
def node_netns_exists(node):
rc, out, _ = sh(["sudo", "-n", "ip", "netns", "list"])
return rc == 0 and f"warp-{node}" in out
def identity_hash(node):
"""Return truncated sha256 of the node's WireGuard private key (never the key)."""
rc, out, _ = sh(["sudo", "-n", "cat", f"/etc/netvm/{node}.conf"])
if rc == 0:
for line in out.splitlines():
m = re.match(r"\s*PrivateKey\s*=\s*(\S+)", line)
if m:
return hashlib.sha256(m.group(1).encode()).hexdigest()[:16]
return None
def probe_node(node, target_name, url):
"""One probe from inside the node's netns. Returns dict."""
# -w fields: http_code, time_total, size_download, remote_ip
fmt = "%{http_code} %{time_total} %{size_download} %{remote_ip}"
cmd = [
"sudo", "-n", "ip", "netns", "exec", f"warp-{node}",
"curl", "-s", "-o", "/dev/null", "-m", str(CURL_TIMEOUT),
"-w", fmt, url,
]
started = datetime.now(timezone.utc)
rc, out, err = sh(cmd, timeout=CURL_TIMEOUT + 10)
ended = datetime.now(timezone.utc)
rec = {
"ts": started.isoformat(),
"node": node,
"probe": target_name,
"url": url,
"curl_rc": rc,
}
parts = out.strip().split()
if rc == 0 and len(parts) == 4:
try:
rec["http_status"] = int(parts[0])
except ValueError:
rec["http_status"] = None
try:
rec["response_ms"] = round(float(parts[1]) * 1000, 1)
except ValueError:
rec["response_ms"] = None
try:
rec["bytes"] = int(parts[2])
except ValueError:
rec["bytes"] = None
rec["remote_ip"] = parts[3] if parts[3] != "0.0.0.0" else None
else:
rec["http_status"] = None
rec["response_ms"] = None
rec["bytes"] = None
rec["remote_ip"] = None
rec["curl_error"] = (err or "curl failed").strip()[:200]
rec["rate_limited"] = rec["http_status"] in (429,)
rec["blocked"] = rec["http_status"] in (403,)
rec["ok"] = rec["http_status"] is not None and 200 <= rec["http_status"] < 400
rec["elapsed_wall_ms"] = round(
(ended - started).total_seconds() * 1000, 1
)
return rec
def egress_ip(node):
"""Best-effort egress IP as seen from inside the netns."""
rc, out, _ = sh([
"sudo", "-n", "ip", "netns", "exec", f"warp-{node}",
"curl", "-s", "-m", "10", "https://api.ipify.org",
], timeout=20)
if rc == 0 and re.fullmatch(r"[0-9a-fA-F.:]+", out.strip()):
return out.strip()
return None
def run_probe_cycle(log_path):
os.makedirs(os.path.dirname(log_path), exist_ok=True)
results = []
any_ok, any_fail = False, False
idhash = {}
for node in NODES:
if not node_netns_exists(node):
results.append({
"ts": datetime.now(timezone.utc).isoformat(),
"node": node, "probe": "node-skip",
"ok": False, "note": "netns warp-%s missing" % node,
})
any_fail = True
continue
idhash[node] = identity_hash(node)
for target_name, url in PROBES:
rec = probe_node(node, target_name, url)
rec["identity_hash"] = idhash[node]
results.append(rec)
if rec["ok"]:
any_ok = True
else:
any_fail = True
time.sleep(1) # gentle pacing between probes
# Attach egress IP per node (one lookup per node, cached per cycle).
egress = {}
for node in {r["node"] for r in results if r.get("probe") != "node-skip"}:
egress[node] = egress_ip(node)
for rec in results:
if rec.get("probe") != "node-skip":
rec["egress_ip"] = egress.get(rec["node"])
with open(log_path, "a") as f:
for rec in results:
f.write(json.dumps(rec) + "\n")
print(json.dumps({
"cycle_ts": datetime.now(timezone.utc).isoformat(),
"log": log_path,
"records": len(results),
"ok": sum(1 for r in results if r.get("ok")),
"failed": sum(1 for r in results if not r.get("ok")),
"nodes": sorted({r["node"] for r in results}),
"egress": egress,
"identities": idhash,
}, indent=2))
if not any_ok:
return 2
return 1 if any_fail else 0
def analyze(log_path, since_hours=None):
if not os.path.exists(log_path):
print("no log yet at %s" % log_path, file=sys.stderr)
return 2
cutoff = None
if since_hours:
cutoff = time.time() - since_hours * 3600
per_node = {}
rate_limit_events = []
identity_changes = {}
egress_by_node = {}
with open(log_path) as f:
for line in f:
line = line.strip()
if not line:
continue
try:
r = json.loads(line)
except json.JSONDecodeError:
continue
if r.get("probe") == "node-skip":
continue
try:
ts = datetime.fromisoformat(r["ts"]).timestamp()
except (ValueError, KeyError):
continue
if cutoff and ts < cutoff:
continue
node = r["node"]
st = per_node.setdefault(node, {
"probes": 0, "ok": 0, "ms": [], "429": 0, "403": 0,
"statuses": {}, "identities": set(), "probes_by_target": {},
})
st["probes"] += 1
if r.get("ok"):
st["ok"] += 1
if r.get("response_ms") is not None:
st["ms"].append(r["response_ms"])
if r.get("rate_limited"):
st["429"] += 1
rate_limit_events.append((r["ts"], node, r["probe"]))
if r.get("blocked"):
st["403"] += 1
s = r.get("http_status")
st["statuses"][str(s)] = st["statuses"].get(str(s), 0) + 1
if r.get("identity_hash"):
st["identities"].add(r["identity_hash"])
t = r.get("probe")
st["probes_by_target"][t] = st["probes_by_target"].get(t, 0) + 1
if r.get("egress_ip"):
egress_by_node.setdefault(node, set()).add(r["egress_ip"])
def pct(vals, p):
if not vals:
return None
s = sorted(vals)
return round(s[min(len(s) - 1, int(p / 100 * len(s)))], 1)
report = {"log": log_path, "nodes": {}}
for node in sorted(per_node):
st = per_node[node]
report["nodes"][node] = {
"probes": st["probes"],
"ok_rate": round(st["ok"] / st["probes"], 3) if st["probes"] else 0,
"latency_ms": {
"mean": round(sum(st["ms"]) / len(st["ms"]), 1) if st["ms"] else None,
"p50": pct(st["ms"], 50),
"p99": pct(st["ms"], 99),
},
"http_429": st["429"],
"http_403": st["403"],
"statuses": st["statuses"],
"identity_rotations": max(0, len(st["identities"]) - 1),
"egress_ips": sorted(egress_by_node.get(node, set())),
}
if len(st["identities"]) > 1:
identity_changes[node] = sorted(st["identities"])
# Divergence analysis: did all nodes see the same fate at the same time?
report["rate_limit_events"] = [
{"ts": ts, "node": n, "probe": p} for ts, n, p in rate_limit_events[-50:]
]
report["identity_changes"] = identity_changes
print(json.dumps(report, indent=2))
return 0
def main():
ap = argparse.ArgumentParser(description="Identity variance testing framework")
sub = ap.add_subparsers(dest="cmd", required=True)
p_probe = sub.add_parser("probe", help="run one probe cycle")
p_probe.add_argument("--log", default=DEFAULT_LOG)
p_an = sub.add_parser("analyze", help="summarize variance data")
p_an.add_argument("--log", default=DEFAULT_LOG)
p_an.add_argument("--since", type=float, default=None,
help="only include last N hours")
args = ap.parse_args()
if args.cmd == "probe":
sys.exit(run_probe_cycle(args.log))
sys.exit(analyze(args.log, args.since))
if __name__ == "__main__":
main()
+50 -4
View File
@@ -222,10 +222,56 @@ def send_dm(agent, target, message, dry_run=False, followup_tags=None,
if HAS_RATE_LIMITER:
rate_limit_wait(agent)
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target]
+ (["--allow-main-chat"] if (target == "main" and allow_main_chat) else [])
+ (followup_tags or []) + [message])
# Cryptographic attestation: sign message with local SSH key
signed_payload = None
dm_sign_sh = NETVM_ROOT / "bin" / "dm-sign.sh"
priv_key = Path(os.path.expanduser("~/.ssh/id_ed25519"))
if dm_sign_sh.exists() and priv_key.exists():
try:
sign_res = subprocess.run(
[str(dm_sign_sh), "--from", "super", "--key", str(priv_key), message],
capture_output=True, text=True, timeout=10
)
if sign_res.returncode == 0 and "-----BEGIN SSH SIGNATURE-----" in sign_res.stdout:
signed_payload = sign_res.stdout.strip()
# Extract message id and register proof to crypt.muse-dev.online
id_m = re.search(r"\[id:([a-f0-9]+)\]", signed_payload)
proof_id = id_m.group(1) if id_m else None
if proof_id:
proof_data = {
"id": proof_id,
"signer": "super",
"target_agent": agent,
"target_conversation": target,
"raw_payload": signed_payload,
"ts": datetime.now(timezone.utc).isoformat()
}
try:
import urllib.request
req = urllib.request.Request(
"https://crypt.muse-dev.online/proofs",
data=json.dumps(proof_data).encode("utf-8"),
headers={"Content-Type": "application/json", "User-Agent": "job-dispatch/1.0"},
method="POST"
)
with urllib.request.urlopen(req, timeout=3) as resp:
pass
except Exception as pe:
sys.stderr.write(f"warning: proof registration to crypt.muse-dev.online failed: {pe}\n")
except Exception as se:
sys.stderr.write(f"warning: dm signing failed: {se}\n")
if signed_payload:
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target, "--raw"]
+ (["--allow-main-chat"] if (target == "main" and allow_main_chat) else [])
+ (followup_tags or []) + [signed_payload])
else:
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target]
+ (["--allow-main-chat"] if (target == "main" and allow_main_chat) else [])
+ (followup_tags or []) + [message])
result = subprocess.run(cmd, capture_output=True, text=True, timeout=120)
if result.returncode != 0:
+94
View File
@@ -0,0 +1,94 @@
#!/usr/bin/env python3
"""Thread keepalive timer entry point (runs on bl).
Reads keepalive-config.json, ticks every registered thread:
ensure reachable -> ping if idle past threshold -> recreate if dead.
Designed to run from a systemd timer (e.g. every 5 minutes). Each tick is
idempotent and cheap: threads that are alive and recently active are a
single navigate+verify (no ping sent).
Usage:
keepalive-timer.py [--config PATH] [--state PATH] [--dry-run] [--status]
--status print the registry with idle times and exit (no changes)
Exit codes: 0 = all ok, 1 = one or more threads failed, 2 = config error.
"""
import argparse
import json
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from keepalive import Keepalive
NETVM_ROOT = os.environ.get("NETVM_ROOT", "/home/super/Projects/NetVM")
DEFAULT_CONFIG = os.path.join(NETVM_ROOT, "keepalive-config.json")
DEFAULT_STATE = os.path.join(NETVM_ROOT, "keepalive-threads.json")
def load_config(path):
try:
with open(path) as f:
cfg = json.load(f)
except FileNotFoundError:
print(f"keepalive: config not found: {path}", file=sys.stderr)
return None
except Exception as e:
print(f"keepalive: bad config {path}: {e}", file=sys.stderr)
return None
threads = cfg.get("threads", [])
if not isinstance(threads, list):
print("keepalive: config 'threads' must be a list", file=sys.stderr)
return None
return threads
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--config", default=DEFAULT_CONFIG)
ap.add_argument("--state", default=DEFAULT_STATE)
ap.add_argument("--dry-run", action="store_true")
ap.add_argument("--status", action="store_true")
args = ap.parse_args()
ka = Keepalive(state_path=args.state, dry_run=args.dry_run)
if args.status:
print(json.dumps(ka.status(), indent=2))
return 0
threads = load_config(args.config)
if threads is None:
return 2
results = {}
for t in threads:
key = t.get("key")
agent = t.get("agent", "opm")
if not key or not t.get("enabled", True):
continue
try:
status = ka.tick(
key,
agent,
idle_threshold_s=int(t.get("idle_threshold_s", 3600)),
max_retries=int(t.get("max_retries", 3)),
ping_message=t.get("ping_message"),
)
except Exception as e:
status = f"error: {e}"
ka.log_event("keepalive_tick_error",
{"key": key, "agent": agent, "error": str(e)[:200]})
results[key] = status
print(f"keepalive: {key}@{agent} -> {status}")
failed = [k for k, v in results.items()
if v in ("failed",) or v.startswith("error")]
return 1 if failed else 0
if __name__ == "__main__":
sys.exit(main())
+345
View File
@@ -0,0 +1,345 @@
#!/usr/bin/env python3
"""Generic thread keepalive for bl side chats.
Pattern proven by the heartbeat job overnight: a state file maps a stable
key -> thread UUID, and each tick navigates directly to
https://muse.ai/thread/<uuid> (UUID reuse fix in muse-chat-api.py
cmd_sidechat_use) instead of name-based sidebar lookup.
This module generalizes that pattern:
- ANY side chat can register for keepalive (not just job sidechats).
- Each tick: ensure the thread is reachable; ping it if idle past threshold;
recreate it if the stored UUID is dead.
- State lives in keepalive-threads.json (atomic write via tmp+rename).
Usage:
from keepalive import Keepalive
ka = Keepalive(state_path="/home/super/Projects/NetVM/keepalive-threads.json")
ka.ensure("ops-watch", agent="opm") # reuse or create
ka.tick("ops-watch", agent="opm",
ping_message="[keepalive] ops-watch {ts}",
idle_threshold_s=3600) # ping if idle
The timer entry point (keepalive-timer.py) drives this from a JSON config.
"""
import json
import os
import re
import subprocess
import sys
import time
from datetime import datetime, timezone
# ---------------------------------------------------------------------------
# Paths (overridable for tests)
# ---------------------------------------------------------------------------
NETVM_ROOT = os.environ.get("NETVM_ROOT", "/home/super/Projects/NetVM")
NETVM_EXEC = os.path.join(NETVM_ROOT, "bin", "netvm-exec.sh")
CHAT_API = os.path.join(NETVM_ROOT, "bin", "muse-chat-api.py")
DEFAULT_STATE = os.path.join(NETVM_ROOT, "keepalive-threads.json")
DEFAULT_LOG = os.path.join(NETVM_ROOT, "keepalive-log.jsonl")
UUID_RE = re.compile(
r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-"
r"[0-9a-f]{4}-[0-9a-f]{12})"
)
def _utcnow():
return datetime.now(timezone.utc).isoformat()
def _ts():
return time.time()
# ---------------------------------------------------------------------------
# Registry
# ---------------------------------------------------------------------------
class Keepalive:
"""Thread keepalive registry + ensure/ping/check operations."""
def __init__(self, state_path=DEFAULT_STATE, log_path=DEFAULT_LOG,
dry_run=False):
self.state_path = state_path
self.log_path = log_path
self.dry_run = dry_run
# -- state ---------------------------------------------------------
def load_state(self):
if os.path.exists(self.state_path):
try:
with open(self.state_path) as f:
data = json.load(f)
return data if isinstance(data, dict) else {}
except Exception:
return {}
return {}
def save_state(self, state):
if self.dry_run:
return
tmp = self.state_path + ".tmp"
with open(tmp, "w") as f:
json.dump(state, f, indent=2)
os.replace(tmp, self.state_path)
def log_event(self, event_type, data):
if self.dry_run:
return
entry = {"ts": _utcnow(), "type": event_type, **data}
try:
with open(self.log_path, "a") as f:
f.write(json.dumps(entry) + "\n")
except Exception:
pass
# -- browser ops (mirrors job-dispatch.py; UUID navigation, not names) --
def _run(self, agent, *api_args, timeout=60):
"""Run muse-chat-api.py in the agent's netns via netvm-exec.sh."""
cmd = [NETVM_EXEC, agent, "--", "python3", CHAT_API,
"--account", agent] + list(api_args)
try:
r = subprocess.run(cmd, capture_output=True, text=True,
timeout=timeout)
return r.returncode, r.stdout.strip(), r.stderr.strip()
except Exception as e:
return -1, "", str(e)
def navigate_to_uuid(self, agent, thread_uuid):
"""Go directly to https://muse.ai/thread/<uuid>.
This is the heartbeat UUID-reuse fix: cmd_sidechat_use navigates by
URL for UUID args instead of searching sidebar titles by name.
Returns True if the command succeeded.
"""
if self.dry_run:
print(f"[DRY] navigate {agent} -> {thread_uuid}")
return True
rc, out, err = self._run(agent, "sidechat", "use", thread_uuid,
timeout=45)
return rc == 0
def current_thread_uuid(self, agent):
"""Read the browser's current URL and extract the thread UUID."""
if self.dry_run:
return None
rc, out, err = self._run(agent, "url", timeout=30)
if rc != 0:
return None
m = UUID_RE.search(out or "")
return m.group(1) if m else None
def create_sidechat(self, agent, name=None):
"""Create a new sidechat in the agent's account. Returns True."""
if self.dry_run:
print(f"[DRY] create sidechat for {agent}")
return True
rc, out, err = self._run(agent, "sidechat", "create", timeout=90)
return rc == 0 and "Created:" in (out or "")
def send_message(self, agent, message):
"""Send a message to the current chat. Returns True on success."""
if self.dry_run:
print(f"[DRY] send to {agent}: {message[:80]}")
return True
rc, out, err = self._run(agent, "send", message, timeout=60)
return rc == 0
# -- keepalive ops ---------------------------------------------------
def check(self, key, agent):
"""Verify the registered thread is reachable.
Navigates to the stored UUID and confirms the browser lands on it.
Returns (ok, thread_uuid). Updates last_verified_ts on success.
"""
state = self.load_state()
rec = state.get(key)
if not rec or not rec.get("thread_uuid"):
return False, None
uuid = rec["thread_uuid"]
if not self.navigate_to_uuid(agent, uuid):
self.log_event("keepalive_check_failed",
{"key": key, "agent": agent, "thread_uuid": uuid,
"reason": "navigate_failed"})
return False, uuid
cur = self.current_thread_uuid(agent)
if cur == uuid:
rec["last_verified_ts"] = _ts()
rec["consecutive_failures"] = 0
state[key] = rec
self.save_state(state)
self.log_event("keepalive_check_ok",
{"key": key, "agent": agent, "thread_uuid": uuid})
return True, uuid
self.log_event("keepalive_check_failed",
{"key": key, "agent": agent, "thread_uuid": uuid,
"reason": "url_mismatch", "current": cur})
return False, uuid
def ensure(self, key, agent, name=None):
"""Ensure the thread exists and is reachable; create if needed.
Returns (ok, thread_uuid). Mirrors the job-dispatch.py reuse-or-
create flow: try stored UUID first, fall back to creation, then
capture the new UUID from the browser URL.
"""
state = self.load_state()
rec = state.get(key)
if rec and rec.get("thread_uuid"):
ok, uuid = self.check(key, agent)
if ok:
return True, uuid
# Stored thread is dead — fall through to recreate.
print(f"keepalive: stored thread for {key} unreachable, "
f"recreating", file=sys.stderr)
# Create a fresh sidechat.
if not self.create_sidechat(agent, name=name):
self.log_event("keepalive_create_failed",
{"key": key, "agent": agent})
return False, None
# Capture the new thread UUID from the browser URL.
uuid = None
if not self.dry_run:
for _ in range(15):
time.sleep(1)
uuid = self.current_thread_uuid(agent)
if uuid:
break
else:
uuid = "dry-run-uuid"
if not uuid:
self.log_event("keepalive_create_failed",
{"key": key, "agent": agent,
"reason": "uuid_capture_failed"})
return False, None
now = _ts()
state[key] = {
"thread_uuid": uuid,
"agent": agent,
"created_ts": now,
"last_ping_ts": 0,
"last_verified_ts": now,
"last_activity_ts": now,
"ping_count": 0,
"consecutive_failures": 0,
}
self.save_state(state)
self.log_event("keepalive_created",
{"key": key, "agent": agent, "thread_uuid": uuid})
return True, uuid
def ping(self, key, agent, message=None):
"""Send a keepalive ping into the thread.
Navigates to the thread first (cheap no-op if already there),
then sends the ping message. Updates last_ping_ts / ping_count.
"""
state = self.load_state()
rec = state.get(key)
if not rec or not rec.get("thread_uuid"):
return False
uuid = rec["thread_uuid"]
if not self.navigate_to_uuid(agent, uuid):
return False
msg = (message or "[keepalive:{key}] tick {ts} {uuid}").format(
key=key, ts=_utcnow(), uuid=uuid)
if not self.send_message(agent, msg):
self.log_event("keepalive_ping_failed",
{"key": key, "agent": agent, "thread_uuid": uuid})
return False
now = _ts()
rec["last_ping_ts"] = now
rec["last_activity_ts"] = now
rec["ping_count"] = rec.get("ping_count", 0) + 1
state[key] = rec
self.save_state(state)
self.log_event("keepalive_ping",
{"key": key, "agent": agent, "thread_uuid": uuid,
"ping_count": rec["ping_count"]})
return True
def note_activity(self, key):
"""Record external activity (e.g. a job just sent to the thread)
so the idle timer doesn't ping unnecessarily."""
state = self.load_state()
rec = state.get(key)
if rec:
rec["last_activity_ts"] = _ts()
state[key] = rec
self.save_state(state)
def tick(self, key, agent, idle_threshold_s=3600, max_retries=3,
ping_message=None):
"""One keepalive tick for a registered thread.
- Ensures the thread exists (recreate if dead).
- Pings only if idle longer than idle_threshold_s.
- After max_retries consecutive failures, forces recreation.
Returns a status string: ok | pinged | recreated | failed.
"""
state = self.load_state()
rec = state.get(key)
if not rec or not rec.get("thread_uuid"):
ok, _ = self.ensure(key, agent)
return "recreated" if ok else "failed"
failures = rec.get("consecutive_failures", 0)
if failures >= max_retries:
# Force recreation: drop the dead UUID and re-ensure.
self.log_event("keepalive_force_recreate",
{"key": key, "agent": agent,
"failures": failures,
"old_uuid": rec.get("thread_uuid")})
rec["thread_uuid"] = None
state[key] = rec
self.save_state(state)
ok, _ = self.ensure(key, agent)
return "recreated" if ok else "failed"
ok, _ = self.check(key, agent)
if not ok:
state = self.load_state()
rec = state.get(key, {})
rec["consecutive_failures"] = failures + 1
state[key] = rec
self.save_state(state)
return "failed"
now = _ts()
idle_for = now - rec.get("last_activity_ts", 0)
if idle_for >= idle_threshold_s:
if self.ping(key, agent, message=ping_message):
return "pinged"
return "failed"
return "ok"
def unregister(self, key):
"""Remove a thread from keepalive (does not delete the sidechat)."""
state = self.load_state()
if key in state:
del state[key]
self.save_state(state)
self.log_event("keepalive_unregistered", {"key": key})
return True
return False
def status(self):
"""Return the full registry with computed idle times."""
state = self.load_state()
now = _ts()
out = {}
for key, rec in state.items():
r = dict(rec)
r["idle_s"] = int(now - rec.get("last_activity_ts", now))
out[key] = r
return out
+88
View File
@@ -0,0 +1,88 @@
#!/usr/bin/env python3
"""
Side-chat to main-chat work siphon — monitor loop.
Polls side chats for new messages, runs detection, siphons hits to main.
This is the integration point for bl. In production:
- list_sidechats() calls muse-chat-api.py or the sidechat manager
- get_messages() reads thread messages via CDP
- post_to_main() sends via muse-chat-api.py send to main chat
For the prototype, all three are injectable (see tests).
"""
import time
from typing import Callable, Dict, List
from detect import detect, is_opted_out
from siphon import siphon, RateLimiter
# Message shape: {"id": str, "text": str, "author": str, "ts": str}
Message = Dict[str, str]
def monitor_once(
list_sidechats: Callable[[], List[Dict[str, str]]],
get_messages: Callable[[str, str], List[Message]],
post_to_main: Callable[[str], bool],
watermarks: Dict[str, str],
limiter: RateLimiter = None,
min_confidence: float = 0.6,
) -> Dict[str, str]:
"""
One poll cycle. Returns updated watermarks.
list_sidechats: () -> [{"id": thread_id, "name": str, "agent": str}]
get_messages: (thread_id, since_msg_id) -> [messages newer than watermark]
post_to_main: (text) -> True on success
watermarks: {thread_id: last_seen_message_id}
"""
lim = limiter or RateLimiter()
new_marks = dict(watermarks)
for chat in list_sidechats():
tid = chat["id"]
agent = chat.get("agent", "unknown")
if is_opted_out(tid):
continue
since = watermarks.get(tid, "")
try:
messages = get_messages(tid, since)
except Exception:
continue # don't let one bad chat kill the cycle
for msg in messages:
mid = msg.get("id", "")
text = msg.get("text", "")
if not mid or not text:
continue
# Update watermark to newest seen
new_marks[tid] = mid
hit = detect(text, tid, mid, min_confidence)
if hit:
siphon(hit, agent, post_to_main, lim)
return new_marks
def monitor_loop(
list_sidechats,
get_messages,
post_to_main,
poll_interval: int = 60,
watermarks: Dict[str, str] = None,
):
"""Run forever. For production use with systemd timer instead."""
marks = watermarks or {}
limiter = RateLimiter()
while True:
marks = monitor_once(
list_sidechats, get_messages, post_to_main, marks, limiter
)
time.sleep(poll_interval)
+125 -1
View File
@@ -16,6 +16,9 @@ Usage:
muse-chat-api.py --account <agent> wait [timeout]
muse-chat-api.py --account <agent> approvals # check pending approvals
muse-chat-api.py --account <agent> upload <file> [--message txt] [--dry-run]
Exit codes: 0 ok, 1 error, 2 APPROVAL_NEEDED (human decision),
3 DOM_NOT_READY (browser not in expected chat state - safe to retry).
"""
import json, urllib.request, websocket, time, sys, argparse, importlib.util
import os
@@ -146,6 +149,120 @@ def check_approvals(ws):
return actions
# Exit code for "browser DOM not ready" — distinct from 1 (generic error)
# and 2 (APPROVAL_NEEDED, needs a human). DOM_NOT_READY is safe to retry
# after a few seconds: the React app may still be hydrating (S3 partial
# load) or the Warp tunnel may have just hiccupped. dm.py treats any
# non-"sent" send output as a failed attempt and retries, so this is a
# fail-closed retry signal, not a crash.
DOM_NOT_READY_EXIT = 3
# Shimmer spans (span.animate-pulse-light) exist in the settled state too
# (2 observed — avatar/media skeleton slots, not load progress). Only a
# mass-skeleton count indicates a stuck/broken render. Heuristic threshold
# from docs/DOM-EDGE-STATES.md 1.1; the settled-state signature below is
# the primary readiness signal.
SHIMMER_STORM_THRESHOLD = 30
def chat_ready_probe(ws):
"""One-shot DOM readiness probe. Returns a dict of observed state, or
{} if the page could not be evaluated at all."""
result = ev1(ws, """(() => {
const q = s => { try { return document.querySelector(s); } catch(e) { return null; } };
const qa = s => { try { return document.querySelectorAll(s); } catch(e) { return []; } };
const title = document.title || '';
const url = window.location.href || '';
const switcher = !!q('[data-testid="hatch-chat-switcher-trigger"]');
const buttons = qa('button').length;
const msglog = !!q('div[role="log"]');
const composer = !!(q('[contenteditable="true"]') ||
q('textarea[placeholder*="Message"]') ||
q('div[role="textbox"]'));
const shimmer = qa('span.animate-pulse-light').length;
let alertText = '';
try {
alertText = [...qa('[role="alert"]')]
.filter(e => e.tagName !== 'SCRIPT' && (e.innerText||'').trim().length > 0)
.slice(0, 2).map(e => e.innerText.slice(0, 160)).join(' | ');
} catch(e) {}
let onLine = true;
try { onLine = navigator.onLine; } catch(e) {}
return JSON.stringify({title, url, switcher, buttons, msglog,
composer, shimmer, alertText, onLine});
})()""")
try:
return json.loads(result) if result else {}
except Exception:
return {}
def classify_chat_state(p):
"""Classify a readiness probe. Returns (state, diagnosis). READY is the
only state a send may proceed from. Reference: docs/DOM-EDGE-STATES.md.
Note: 'Chat \u2014' is the 'Chat \u2014 <identity>' title prefix."""
if not p:
return ("NO_PROBE",
"CDP evaluate returned nothing - page blank, navigating, or CDP wedged")
title = p.get("title", "")
url = p.get("url", "")
if title == "muse.ai":
return ("LANDING",
"browser is on the muse.ai landing page, not in chat (state: LANDING)")
if not p.get("onLine", True):
return ("OFFLINE",
"navigator.onLine is false - browser reports no network")
if p.get("alertText"):
return ("REAL_ERROR",
"page shows an error banner: %s" % p["alertText"][:160])
if "/thread/new" in url:
# Legitimate only right after `sidechat create`: the new-thread
# composer is up and the title has settled. A stripped /thread/new
# (no title, no composer) is the S8 broken state.
if p.get("composer") and "Chat \u2014" in title:
return ("READY", "new-thread composer state (/thread/new)")
return ("S8_STRIPPED",
"stripped /thread/new page - no composer/title (state: S8)")
switcher = p.get("switcher", False)
buttons = p.get("buttons", 0)
if title == "Muse" and (not switcher or buttons < 30):
return ("S3_PARTIAL",
"React still hydrating: title 'Muse', %d buttons, switcher=%s (state: S3)"
% (buttons, switcher))
if not ("Chat \u2014" in title and switcher and buttons > 30 and p.get("msglog")):
return ("NOT_SETTLED",
"settled-state signature not met (title=%r switcher=%s buttons=%d msglog=%s)" %
(title[:40], switcher, buttons, p.get("msglog", False)))
if p.get("shimmer", 0) > SHIMMER_STORM_THRESHOLD:
return ("SHIMMER_STORM",
"mass skeleton render: %d shimmer spans (settled has ~2)" % p["shimmer"])
if not p.get("composer"):
return ("COMPOSER_MISSING",
"chat settled but no composer input found - send has nowhere to type")
return ("READY", "settled chat state (title=%r)" % title[:40])
def assert_chat_ready(ws, context="send"):
"""Fail-closed React-readiness gate ("matching chromebox").
Call before any DOM mutation (send, upload). On mismatch prints a
structured diagnosis to stderr and exits DOM_NOT_READY_EXIT (3) - a
retryable signal, distinct from APPROVAL_NEEDED (2, needs a human).
S3 partial load gets one 5s grace re-probe (transient render phase);
everything else fails immediately.
"""
probe = chat_ready_probe(ws)
state, detail = classify_chat_state(probe)
if state == "S3_PARTIAL":
time.sleep(5)
probe = chat_ready_probe(ws)
state, detail = classify_chat_state(probe)
if state != "READY":
print("DOM_NOT_READY [%s] %s" % (state, detail), file=sys.stderr)
print("DOM_NOT_READY probe: %s" % json.dumps(probe), file=sys.stderr)
sys.exit(DOM_NOT_READY_EXIT)
return probe
def cmd_approvals(ws):
"""Check and handle pending approvals."""
actions = check_approvals(ws)
@@ -160,7 +277,12 @@ def cmd_approvals(ws):
sys.exit(2)
def cmd_send(ws, message):
# Check approvals first
# React-readiness gate ("matching chromebox"): fail closed if the
# browser is not in the expected chat state (landing page,
# partial load, partition). Exit 3 = safe to retry, unlike
# APPROVAL_NEEDED (2, needs a human).
assert_chat_ready(ws, context="send")
# Check approvals
actions = check_approvals(ws)
for dialog, trusted, action in actions:
if not trusted:
@@ -479,6 +601,8 @@ def cmd_upload(ws, filepath, message=None, dry_run=False):
if not os.path.isfile(filepath):
print("ERROR: file not found: %s" % filepath, file=sys.stderr)
sys.exit(1)
# React-readiness gate ("matching chromebox") - see cmd_send.
assert_chat_ready(ws, context="upload")
actions = check_approvals(ws)
for dialog, trusted, action in actions:
if not trusted:
+49
View File
@@ -0,0 +1,49 @@
#!/usr/bin/env bash
# muse-cli-node <node> <muse-cli-args...>
# Runs muse-cli inside the node's isolated network namespace with its own
# dedicated Cloudflare WARP egress identity and isolated cookies/config.
# Auto-refreshes cookies via Chromium CDP if authentication fails.
set -euo pipefail
NETVM_BIN="$(cd "$(dirname "$0")" && pwd)"
NODE="${1:?usage: muse-cli-node <node> [args...]}"
shift
CONF_DIR="$HOME/.config/muse-cli/$NODE"
mkdir -p "$CONF_DIR"
run_cli() {
export MUSE_CONFIG_DIR="$CONF_DIR"
"$NETVM_BIN/netvm-exec.sh" "$NODE" -- python3 -c "
import os, sys
os.environ['MUSE_CONFIG_DIR'] = '$CONF_DIR'
import muse_cli.cli as cli
cli.CONFIG_DIR = '$CONF_DIR'
cli.CONFIG_FILE = os.path.join('$CONF_DIR', 'config.json')
cli.COOKIES_FILE = os.path.join('$CONF_DIR', 'cookies.txt')
try:
cli.main()
except SystemExit as e:
sys.exit(e.code)
" "$@"
}
# First attempt
set +e
OUTPUT=$(run_cli "$@" 2>&1)
RC=$?
set -e
# Detect authentication failure and trigger single auto-refresh retry
if [ $RC -ne 0 ] && echo "$OUTPUT" | grep -qiE "AuthError|hatch_sess|401|unauthorized"; then
echo "Notice: Auth failure detected for $NODE; refreshing cookies from Chromium CDP..." >&2
if "$NETVM_BIN/netvm-exec.sh" "$NODE" -- python3 "$NETVM_BIN/refresh-node-cookies.py" "$NODE"; then
echo "Notice: Cookies refreshed for $NODE. Retrying operation..." >&2
run_cli "$@"
exit $?
fi
fi
# If succeeded or failed on something else, emit output and exit with RC
echo "$OUTPUT"
exit $RC
+100
View File
@@ -0,0 +1,100 @@
#!/usr/bin/env python3
"""
muse_hybrid.py: Hybrid integration layer combining fast headless gateway actions
(via muse-cli-node with per-node Cloudflare WARP egress & cookies) with Chromebox
CDP DOM interactions.
Provides:
- muse_threads(agent) -> list of threads
- muse_history(agent, thread_id=None, limit=10) -> list of messages
- muse_send(agent, text, thread_id=None, wait=0) -> dict response
- muse_session_start(agent, title=None) -> dict response
"""
import subprocess
import json
import os
import sys
from pathlib import Path
BIN_DIR = Path(__file__).resolve().parent
MUSE_CLI_NODE = BIN_DIR / "muse-cli-node"
def is_node_configured(node):
conf_dir = Path.home() / ".config" / "muse-cli" / node
return (conf_dir / "cookies.txt").exists()
def run_muse_cli(node, args, timeout=30):
cmd = [str(MUSE_CLI_NODE), node] + args
res = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
return res.returncode, res.stdout, res.stderr
def get_threads(node):
"""Retrieve thread list via fast gateway. Returns (threads_list, error_str)."""
rc, stdout, stderr = run_muse_cli(node, ["threads"])
if rc != 0:
return None, stderr or stdout
try:
data = json.loads(stdout)
return data, None
except Exception as e:
return None, f"JSON parse error: {e}"
def get_history(node, thread_id=None, limit=15):
"""Retrieve message history via fast gateway. Returns (history_list, error_str)."""
args = ["history", "--limit", str(limit)]
if thread_id and thread_id != "main":
args.extend(["--thread", str(thread_id)])
rc, stdout, stderr = run_muse_cli(node, args)
if rc != 0:
return None, stderr or stdout
try:
data = json.loads(stdout)
return data, None
except Exception as e:
return None, f"JSON parse error: {e}"
def send_message(node, text, thread_id=None, wait=0):
"""Send message via fast gateway. Returns (result_dict, error_str)."""
args = ["send"]
if thread_id and thread_id != "main":
args.extend(["--thread", str(thread_id)])
if wait > 0:
args.extend(["--wait", str(wait)])
args.append(text)
rc, stdout, stderr = run_muse_cli(node, args, timeout=max(30, wait + 15))
if rc != 0:
return None, stderr or stdout
try:
data = json.loads(stdout)
return data, None
except Exception as e:
return {"raw": stdout}, None
if __name__ == "__main__":
if len(sys.argv) < 3:
print("Usage: muse_hybrid.py <node> <threads|history|send> [args...]")
sys.exit(1)
node = sys.argv[1]
action = sys.argv[2]
if action == "threads":
threads, err = get_threads(node)
if err:
print("Error:", err, file=sys.stderr)
sys.exit(1)
print(json.dumps(threads, indent=2))
elif action == "history":
tid = sys.argv[3] if len(sys.argv) > 3 else None
msgs, err = get_history(node, tid)
if err:
print("Error:", err, file=sys.stderr)
sys.exit(1)
print(json.dumps(msgs, indent=2))
elif action == "send":
txt = sys.argv[3]
tid = sys.argv[4] if len(sys.argv) > 4 else None
res, err = send_message(node, txt, tid)
if err:
print("Error:", err, file=sys.stderr)
sys.exit(1)
print(json.dumps(res, indent=2))
+2 -2
View File
@@ -29,10 +29,10 @@ fi
# 94x0 sequence (9410 muse, 9420 pip, 9430 646, 9440 opm, ...).
# Scan the registry and live listeners; first free wins.
used_ports() {
{ grep -oP '^\|\s*\K94\d0(?=\s*\|)' "$NODES_MD" 2>/dev/null || true; \
{ awk -F'|' 'NF>=5 && $5 ~ /94[0-9]0/ {gsub(/ /,"",$5); print $5}' "$NODES_MD" 2>/dev/null || true; \
ss -tln 2>/dev/null | grep -oP ':\K94\d0\b' || true; } | sort -u
}
PORT=9450
PORT=9410
while used_ports | grep -qx "$PORT"; do PORT=$((PORT+10)); done
[ "$PORT" -gt 9600 ] && { echo "port pool exhausted" >&2; exit 1; }
+88 -9
View File
@@ -8,15 +8,70 @@ Any .py script can use this to avoid hitting Muse rate limits:
Uses a token bucket per agent stored in /tmp (persists across invocations,
not across reboots). Default: 1 op per 3s sustained, burst of 5, max ~20/min.
Jitter (2026-10-04): each agent gets a DETERMINISTIC per-agent jitter factor
derived from sha256(salt + agent), in [1-jitter, 1+jitter] (default
jitter=0.3, i.e. +/-30%). Rationale: all 4 fleet nodes egress from a single
Cloudflare IP, so Cloudflare sees correlated traffic. When fleet timers fire
simultaneously, per-agent jitter drifts each node's send phase apart instead
of hitting in lockstep. Deterministic per agent = reproducible behavior;
different agents = de-correlated phases. Does not make limits stricter: the
mean interval is unchanged.
"""
import hashlib
import json
import time
import os
import random
import time
RATE_LIMIT_FILE = "/tmp/netvm-rate-limit.json"
RATE_LIMIT_INTERVAL = 3.0 # seconds between ops
RATE_LIMIT_INTERVAL = 3.0 # seconds between ops (base; jittered per agent)
RATE_LIMIT_BURST = 5
RATE_LIMIT_MAX_PER_MIN = 20
DEFAULT_JITTER = 0.3 # +/-30% deterministic per-agent interval jitter
JITTER_SEED_SALT = "netvm-rate-jitter:v1:"
def _jitter_factor(agent, jitter=DEFAULT_JITTER):
"""Deterministic per-agent multiplier in [1-jitter, 1+jitter].
Seeded by sha256(salt + agent): the same agent always gets the same
factor (reproducible), different agents get different factors
(de-correlated). Uses an isolated Random instance; global random
state is untouched.
"""
if jitter <= 0:
return 1.0
digest = hashlib.sha256((JITTER_SEED_SALT + agent).encode()).digest()
seed = int.from_bytes(digest[:8], "big")
rng = random.Random(seed)
return 1.0 + jitter * (rng.random() * 2.0 - 1.0)
def effective_interval(agent, interval=RATE_LIMIT_INTERVAL,
jitter=DEFAULT_JITTER):
"""The actual minimum op spacing for this agent after jitter."""
return interval * _jitter_factor(agent, jitter)
def rate_limit_policy(agent, interval=RATE_LIMIT_INTERVAL,
burst=RATE_LIMIT_BURST, jitter=DEFAULT_JITTER):
"""Return the effective rate-limit policy for an agent (audit/docs)."""
return {
"agent": agent,
"base_interval_s": interval,
"jitter": jitter,
"jitter_factor": round(_jitter_factor(agent, jitter), 4),
"effective_interval_s": round(
effective_interval(agent, interval, jitter), 3),
"burst": burst,
"max_per_min": RATE_LIMIT_MAX_PER_MIN,
"state_file": RATE_LIMIT_FILE,
"scope": ("per-agent buckets; state shared in one file, "
"keys namespaced by agent"),
}
def _load_state():
try:
@@ -25,6 +80,7 @@ def _load_state():
except:
return {}
def _save_state(state):
try:
# Atomic write via temp file
@@ -35,11 +91,17 @@ def _save_state(state):
except:
pass
def rate_limit_wait(agent, interval=RATE_LIMIT_INTERVAL, burst=RATE_LIMIT_BURST):
def rate_limit_wait(agent, interval=RATE_LIMIT_INTERVAL,
burst=RATE_LIMIT_BURST, jitter=DEFAULT_JITTER):
"""
Block until the agent is allowed to perform an operation.
Returns the time waited in seconds (0 if no wait needed).
The interval is jittered deterministically per agent so fleet nodes
don't send in lockstep.
"""
eff_interval = effective_interval(agent, interval, jitter)
state = _load_state()
now = time.time()
@@ -58,11 +120,11 @@ def rate_limit_wait(agent, interval=RATE_LIMIT_INTERVAL, burst=RATE_LIMIT_BURST)
recent = state.get(recent_key, [])
recent = [t for t in recent if time.time() - t < 60]
# Check interval limit
# Check interval limit (jittered per agent)
last = state.get(agent, 0)
now = time.time()
if now - last < interval:
wait = interval - (now - last)
if now - last < eff_interval:
wait = eff_interval - (now - last)
time.sleep(wait)
# Record this operation
@@ -73,10 +135,13 @@ def rate_limit_wait(agent, interval=RATE_LIMIT_INTERVAL, burst=RATE_LIMIT_BURST)
_save_state(state)
return 0
def rate_limit_check(agent):
def rate_limit_check(agent, interval=RATE_LIMIT_INTERVAL,
jitter=DEFAULT_JITTER):
"""
Non-blocking check. Returns (allowed: bool, wait_seconds: float).
"""
eff_interval = effective_interval(agent, interval, jitter)
state = _load_state()
now = time.time()
recent = state.get(f"{agent}_recent", [])
@@ -87,7 +152,21 @@ def rate_limit_check(agent):
return False, max(0, wait)
last = state.get(agent, 0)
if now - last < RATE_LIMIT_INTERVAL:
return False, RATE_LIMIT_INTERVAL - (now - last)
if now - last < eff_interval:
return False, eff_interval - (now - last)
return True, 0
if __name__ == "__main__":
import argparse
p = argparse.ArgumentParser(
description="Inspect NetVM per-agent rate-limit policy")
p.add_argument("--policy", nargs="*", default=["muse", "pip", "646", "opm"],
help="Show effective policy per agent")
p.add_argument("--no-jitter", action="store_true",
help="Show policy without jitter")
args = p.parse_args()
j = 0.0 if args.no_jitter else DEFAULT_JITTER
for a in args.policy:
print(json.dumps(rate_limit_policy(a, jitter=j), indent=2))
+82
View File
@@ -0,0 +1,82 @@
#!/usr/bin/env python3
"""
refresh-node-cookies.py: Extract fresh cookies from a running node's Chromium CDP
instance inside its isolated network namespace and save to ~/.config/muse-cli/<node>/cookies.txt.
"""
import sys
import os
import json
import urllib.request
from pathlib import Path
NODES = {
"muse": 9410,
"pip": 9420,
"646": 9430,
"opm": 9440,
"def": 9450,
}
def export_cookies(node):
port = NODES.get(node)
if not port:
print(f"Unknown node: {node}", file=sys.stderr)
return False
conf_dir = Path.home() / ".config" / "muse-cli" / node
conf_dir.mkdir(parents=True, exist_ok=True)
cookie_file = conf_dir / "cookies.txt"
cfg_file = conf_dir / "config.json"
# We import websocket inside netns execution
try:
import websocket
except ImportError:
print("websocket-client not installed", file=sys.stderr)
return False
try:
req = urllib.request.urlopen(f"http://127.0.0.1:{port}/json", timeout=3)
tabs = json.loads(req.read().decode())
if not tabs:
print(f"No tabs found on port {port}", file=sys.stderr)
return False
ws_url = tabs[0]["webSocketDebuggerUrl"]
ws = websocket.create_connection(ws_url, timeout=5)
ws.send(json.dumps({
"id": 1,
"method": "Network.getCookies",
"params": {"urls": ["https://muse.ai"]}
}))
res = json.loads(ws.recv())
ws.close()
cookies = res.get("result", {}).get("cookies", [])
if not cookies:
print(f"No cookies returned from CDP for {node}", file=sys.stderr)
return False
lines = [f"{c['name']}={c['value']}" for c in cookies]
cookie_file.write_text("; ".join(lines) + "\n")
cookie_file.chmod(0o600)
cfg_data = {}
if cfg_file.exists():
try:
cfg_data = json.loads(cfg_file.read_text())
except Exception:
pass
cfg_data["cookies_file"] = str(cookie_file)
cfg_file.write_text(json.dumps(cfg_data, indent=2))
cfg_file.chmod(0o600)
return True
except Exception as e:
print(f"Failed to export cookies for {node}: {e}", file=sys.stderr)
return False
if __name__ == "__main__":
if len(sys.argv) < 2:
print("Usage: refresh-node-cookies.py <node>", file=sys.stderr)
sys.exit(1)
success = export_cookies(sys.argv[1])
sys.exit(0 if success else 1)
+26 -6
View File
@@ -42,6 +42,8 @@ from datetime import datetime, timezone
BASE = "/home/super/Projects/NetVM"
BIN = os.path.join(BASE, "bin")
if BIN not in sys.path:
sys.path.insert(0, BIN)
BOX_CHAT = os.path.join(BIN, "box-chat.py")
DM_PY = os.path.join(BIN, "dm.py")
STATE_FILE = os.path.join(BASE, "self-main-loop-watermark.json")
@@ -132,7 +134,25 @@ def get_config(st):
def read_main_chat(agent, limit):
"""Read agent's muse.ai Main chat via box-chat.py. Returns (ok, messages|error)."""
"""Read agent's muse.ai Main chat. Tries fast headless gateway (muse_hybrid) first, falling back to box-chat.py."""
# Fast path: muse_hybrid over WebSocket Noise frame inside isolated netns
try:
import muse_hybrid
msgs, err = muse_hybrid.get_history(agent, thread_id=None, limit=limit)
if msgs and not err:
formatted = []
for m in msgs:
formatted.append({
"id": m.get("message_id") or f"msg-{m.get('seq')}",
"ts": None,
"from": {"name": m.get("role", "unknown")},
"text": m.get("text", "")
})
return True, formatted
except Exception:
pass
# Fallback to box-chat.py / CDP if gateway fails or returns empty
cmd = [sys.executable, BOX_CHAT, "thread-messages", agent, "main",
"--limit", str(limit)]
try:
@@ -172,16 +192,16 @@ def author_of(m):
def compose_digest(agent, new):
"""Concise, actionable digest for the prompting sidechat."""
lines = ["[main-loop] %d new in %s main chat:" % (len(new), agent)]
"""Concise, actionable digest for the prompting sidechat without synthetic ceremony."""
lines = ["New messages in %s main chat (%d):" % (agent, len(new))]
for m in new[:5]:
text = (m.get("text") or "").strip().replace("\n", " ")
q = " [?]" if "?" in text else ""
hot = " [!]" if (OPERATOR_RE.search(text) or "urgent" in text.lower()) else ""
q = " (?)" if "?" in text else ""
hot = " (!)" if (OPERATOR_RE.search(text) or "urgent" in text.lower()) else ""
lines.append("- %s: %s%s%s" % (author_of(m), text[:PREVIEW_MAX], q, hot))
if len(new) > 5:
lines.append("(+%d more)" % (len(new) - 5))
lines.append("Check main chat via box when you can.")
lines.append("Check main chat via box when available.")
digest = "\n".join(lines)
return digest[:DIGEST_MAX]
+256
View File
@@ -0,0 +1,256 @@
#!/usr/bin/env python3
"""
Siphon bl wiring — connects siphon/monitor.py to bl's chat infrastructure.
Implements the three injectable functions:
- list_sidechats(): threads to monitor (from state files)
- get_messages(thread_id, since): read via muse-chat-api.py
- post_to_main(text): send to opm's main chat via muse-chat-api.py
Run via systemd timer every 60s, or manually:
python3 siphon-bl.py --once
python3 siphon-bl.py --dry-run
"""
import argparse
import re
import json
import os
import subprocess
import sys
import time
# Add bin dir to path for siphon imports
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from monitor import monitor_once
from siphon import RateLimiter
NETVM_BIN = "/home/super/Projects/NetVM/bin"
CHAT_API = os.path.join(NETVM_BIN, "muse-chat-api.py")
# State files that map reuse_key -> thread UUID
STATE_FILES = [
"/home/super/Projects/NetVM/job-sidechats.json",
"/home/super/sidechat-wake/wake-sidechats.json",
"/home/super/Projects/NetVM/keepalive-threads.json",
]
WATERMARKS_FILE = "/home/super/Projects/NetVM/siphon-watermarks.json"
LOG_FILE = "/home/super/Projects/NetVM/siphon-bl.log"
def log(msg):
line = f"[{time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())}] {msg}"
print(line, flush=True)
try:
with open(LOG_FILE, "a") as f:
f.write(line + "\n")
except OSError:
pass
def run_chat_api(agent, *args, timeout=60):
"""Run muse-chat-api.py and return stdout."""
cmd = [sys.executable, CHAT_API, "--account", agent] + list(args)
result = subprocess.run(
cmd, capture_output=True, text=True, timeout=timeout
)
return result.stdout.strip(), result.returncode
def list_sidechats():
"""
Build list of threads to monitor from state files.
Returns: [{"id": thread_uuid, "name": reuse_key, "agent": agent}]
"""
chats = []
seen = set()
for sf in STATE_FILES:
try:
with open(sf) as f:
state = json.load(f)
except (OSError, json.JSONDecodeError):
continue
for key, val in state.items():
# Handle different state formats
if isinstance(val, dict):
uuid = val.get("thread_uuid") or val.get("uuid")
agent = val.get("agent", "opm")
elif isinstance(val, str):
uuid = val
agent = "opm"
else:
continue
if uuid and uuid not in seen:
seen.add(uuid)
chats.append({
"id": uuid,
"name": key,
"agent": agent,
})
return chats
def get_messages(thread_id, since_msg_id):
"""
Get messages from a thread newer than since_msg_id.
Uses muse-chat-api.py: sidechat use <uuid>, then messages.
Returns: [{"id": str, "text": str, "author": str, "ts": str}]
"""
# Find which agent owns this thread
agent = "opm" # default
for chat in list_sidechats():
if chat["id"] == thread_id:
agent = chat["agent"]
break
# Navigate to the thread
out, rc = run_chat_api(agent, "sidechat", "use", thread_id, timeout=30)
if rc != 0:
return []
# Get messages
out, rc = run_chat_api(agent, "messages", timeout=30)
if rc != 0:
return []
# Parse messages — format is "---" separated
messages = []
# Use timestamp + hash as message ID (no stable IDs from the API)
for i, chunk in enumerate(out.split("\n---\n")):
chunk = chunk.strip()
if not chunk or chunk == "Ok":
continue
# Create a stable-ish ID from content hash
import hashlib
mid = hashlib.md5(chunk.encode()).hexdigest()[:12]
# Skip if we've seen this (watermark comparison)
if since_msg_id and mid <= since_msg_id:
continue
messages.append({
"id": mid,
"text": chunk[:2000], # truncate long messages
"author": agent,
"ts": str(time.time()),
})
# Return to main chat
run_chat_api(agent, "sidechat", "main", timeout=15)
return messages
def post_to_main(text):
"""
Post siphoned summary to opm's main chat.
Returns True on success.
Tracked hits (text contains [reply:expected]) route via dm.py
--expect-reply --thread so a dm_followup record is created.
Untracked hits use the direct API path.
"""
# Tracked? Look for the follow-up tag the modulated siphon appends.
if "[reply:expected]" in text:
# Extract source thread UUID from the thread URL in the text.
m = re.search(r"https://muse\.ai/thread/([a-f0-9-]{36})", text)
thread_id = m.group(1) if m else None
if thread_id:
return post_to_main_tracked(text, thread_id)
# Tracked but no thread URL: fall through to direct (fail-open
# toward visibility).
log("tracked siphon hit without thread URL, using direct post")
# Untracked (or fallback): direct API send to main chat.
run_chat_api("opm", "sidechat", "main", timeout=15)
out, rc = run_chat_api("opm", "send", text, timeout=60)
return rc == 0
def post_to_main_tracked(text, thread_id):
"""
Post a tracked siphon hit via dm.py so a dm_followup record is
created. Returns True on success.
"""
import shlex
cmd = [
sys.executable,
os.path.join(NETVM_BIN, "dm.py"),
"send",
"--agent", "opm",
"--to", "opm",
"--target", "main",
"--expect-reply",
"--thread", thread_id,
text,
]
try:
result = subprocess.run(
cmd, capture_output=True, text=True, timeout=120
)
ok = "SENT and VERIFIED" in (result.stdout or "")
if not ok:
log(f"dm.py tracked post failed: {(result.stdout or '')[:200]}")
return ok
except Exception as e:
log(f"dm.py tracked post exception: {e}")
return False
def load_watermarks():
try:
with open(WATERMARKS_FILE) as f:
return json.load(f)
except (OSError, json.JSONDecodeError):
return {}
def save_watermarks(marks):
tmp = WATERMARKS_FILE + ".tmp"
with open(tmp, "w") as f:
json.dump(marks, f)
os.replace(tmp, WATERMARKS_FILE)
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--once", action="store_true",
help="Run one poll cycle and exit")
parser.add_argument("--dry-run", action="store_true",
help="Don't post, just show what would be siphoned")
parser.add_argument("--min-confidence", type=float, default=0.6)
args = parser.parse_args()
if args.dry_run:
# Dry run: show what would be detected without posting
def dry_post(text):
print(f"[DRY] would post to main:\n{text}\n")
return True
post_fn = dry_post
else:
post_fn = post_to_main
watermarks = load_watermarks()
limiter = RateLimiter()
chats = list_sidechats()
log(f"Monitoring {len(chats)} threads")
new_marks = monitor_once(
list_sidechats,
get_messages,
post_fn,
watermarks,
limiter,
min_confidence=args.min_confidence,
)
save_watermarks(new_marks)
log(f"Cycle complete. Watermarks: {len(new_marks)} threads tracked.")
if __name__ == "__main__":
main()
+158
View File
@@ -0,0 +1,158 @@
#!/usr/bin/env python3
"""
Side-chat to main-chat work siphon — siphon action.
When detection fires, post a summary to main chat with:
- Category badge
- One-line summary (never full message text)
- Link back to the source side chat thread
- Confidence score (for transparency)
Safety:
- Rate limited (max N siphons per hour per thread)
- Never posts full message content
- Respects opt-out registry
- Deduplicates (same message_id never siphoned twice)
"""
import time
from dataclasses import dataclass, field
from typing import Callable, Optional
from detect import SiphonHit, is_opted_out
# --- Follow-up modulation ---
#
# Wire the follow-up modulation table into the siphon so each hit gets
# the right follow-up policy:
# ALERT / BLOCKER / DECISION -> tracked, fast fuse for ALERT/BLOCKER
# COMPLETED / MILESTONE -> untracked (no nudge budget burned)
#
# modulate.py must be landed on bl before this runs (rollout step 1).
# If the import fails we degrade to the old behavior: post the summary
# with no follow-up tags (fail-closed toward visibility, not tracking).
try:
from modulate import for_siphon_hit, render_tags
_MODULATION_AVAILABLE = True
except ImportError: # pragma: no cover - deploy keeps modulate.py present
_MODULATION_AVAILABLE = False
for_siphon_hit = None
render_tags = None
def policy_for_hit(hit: SiphonHit):
"""Follow-up policy for a siphon hit, or None when untracked.
COMPLETED / MILESTONE hits return None (post the summary, create no
follow-up record). ALERT / BLOCKER / DECISION return a Policy whose
tags render into the canonical bracket vocabulary.
"""
if not _MODULATION_AVAILABLE:
return None
return for_siphon_hit(hit.category)
def is_tracked(hit: SiphonHit) -> bool:
"""True when this hit should create a follow-up record.
Callers that route tracked posts through dm.py --expect-reply (so a
dm_followup record is actually created) can use this to choose the
post path. Untracked hits post as plain summaries.
"""
return policy_for_hit(hit) is not None
# --- Rate limiting ---
@dataclass
class RateLimiter:
max_per_hour: int = 5
_timestamps: dict = field(default_factory=dict) # thread_id -> [ts, ...]
def allow(self, thread_id: str) -> bool:
now = time.time()
stamps = self._timestamps.get(thread_id, [])
# Prune older than 1 hour
stamps = [s for s in stamps if now - s < 3600]
if len(stamps) >= self.max_per_hour:
return False
stamps.append(now)
self._timestamps[thread_id] = stamps
return True
# --- Deduplication ---
_siphoned_ids: set = set()
def already_siphoned(message_id: str) -> bool:
return message_id in _siphoned_ids
def mark_siphoned(message_id: str):
_siphoned_ids.add(message_id)
# --- Siphon action ---
CATEGORY_EMOJI = {
"COMPLETED": "✅",
"BLOCKER": "🚧",
"DECISION": "❓",
"ALERT": "🚨",
"MILESTONE": "🎯",
}
def format_siphon(hit: SiphonHit, agent_name: str = "sidechat") -> str:
"""
Format a siphon message for main chat.
Never includes full message text — summary + link only.
"""
emoji = CATEGORY_EMOJI.get(hit.category, "📋")
thread_url = f"https://muse.ai/thread/{hit.thread_id}"
return (
f"{emoji} [{hit.category}] from {agent_name} side chat\n"
f"{hit.summary}\n"
f"→ {thread_url}\n"
f"(confidence {hit.confidence:.0%})"
)
def siphon(hit: SiphonHit,
agent_name: str,
post_to_main: Callable[[str], bool],
limiter: Optional[RateLimiter] = None) -> bool:
"""
Execute the siphon: post summary to main chat.
post_to_main: callable that posts text to main chat, returns True on success.
Returns True if siphoned, False if suppressed.
"""
# Safety checks
if is_opted_out(hit.thread_id):
return False
if already_siphoned(hit.message_id):
return False
lim = limiter or RateLimiter()
if not lim.allow(hit.thread_id):
return False
text = format_siphon(hit, agent_name)
# Follow-up modulation: tracked hits (ALERT/BLOCKER/DECISION) get
# the canonical follow-up tags appended — [reply:expected],
# [reply:timeout=N], [reply:nudges=N], [reply:escalate=X], and
# [input:siphon] for the audit trail. Untracked hits
# (COMPLETED/MILESTONE) post as plain summaries.
policy = policy_for_hit(hit)
if policy is not None:
text = text + "\n" + render_tags(policy)
ok = post_to_main(text)
if ok:
mark_siphoned(hit.message_id)
return ok
+61 -24
View File
@@ -1174,24 +1174,38 @@ def cmd_dm_files(args):
# ---------------------------------------------------------------------------
def cmd_thread_list(args):
agent = args.agent
cmd = ["python3", str(BIN_DIR / "box-chat.py"), "thread-list", agent]
res = subprocess.run(cmd, capture_output=True, text=True)
threads = []
# Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress)
try:
data = json.loads(res.stdout)
import muse_hybrid
gw_threads, gw_err = muse_hybrid.get_threads(agent)
if gw_threads and not gw_err:
for t in gw_threads:
threads.append({
"id": t.get("session_id", ""),
"kind": "thread" if t.get("thread") else "chat",
"title": t.get("title") or "(no title)",
"participants": [agent],
"last_message_at": t.get("updated", "")
})
except Exception:
print(c_red("Failed to parse box-chat.py response:") + f"\n{res.stdout}\n{res.stderr}")
return
pass
# Fallback to box-chat.py / CDP if gateway returned no threads
if not threads:
cmd = ["python3", str(BIN_DIR / "box-chat.py"), "thread-list", agent]
res = subprocess.run(cmd, capture_output=True, text=True)
try:
data = json.loads(res.stdout)
threads = data.get("threads", []) if data.get("ok") else []
except Exception:
threads = []
if args.json:
print(json.dumps(data, indent=2))
print(json.dumps({"ok": True, "threads": threads}, indent=2))
return
if not data.get("ok"):
print(c_red(f"Error fetching threads for {agent}: {data.get('error', 'unknown error')}"))
return
threads = data.get("threads", [])
print("\n" + c_bold(f"=== ACTIVE THREADS FOR {agent.upper()} ({len(threads)}) ===") + "\n")
headers = ["ID / UUID", "KIND", "TITLE", "PARTICIPANTS", "LAST ACTIVE"]
@@ -1213,25 +1227,39 @@ def cmd_thread_view(args):
agent = args.agent
thread_id = args.thread_id
limit = getattr(args, "limit", 15)
messages = []
cmd = ["python3", str(BIN_DIR / "box-chat.py"), "thread-messages", agent, thread_id, "--limit", str(limit)]
res = subprocess.run(cmd, capture_output=True, text=True)
# Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress)
try:
data = json.loads(res.stdout)
import muse_hybrid
gw_msgs, gw_err = muse_hybrid.get_history(agent, thread_id=thread_id, limit=limit)
if gw_msgs and not gw_err:
for m in gw_msgs:
role = m.get("role", "unknown")
messages.append({
"from": {"name": role, "role": role},
"text": m.get("text", "")
})
except Exception:
print(c_red("Failed to parse box-chat.py response:") + f"\n{res.stdout}\n{res.stderr}")
return
pass
# Fallback to box-chat.py / CDP if gateway returned no messages
if not messages:
cmd = ["python3", str(BIN_DIR / "box-chat.py"), "thread-messages", agent, thread_id, "--limit", str(limit)]
res = subprocess.run(cmd, capture_output=True, text=True)
try:
data = json.loads(res.stdout)
messages = data.get("messages", []) if data.get("ok") else []
except Exception:
messages = []
if not messages:
print(c_red(f"Error viewing thread {thread_id}: no messages returned"))
return
if args.json:
print(json.dumps(data, indent=2))
print(json.dumps({"ok": True, "messages": messages}, indent=2))
return
if not data.get("ok"):
print(c_red(f"Error viewing thread {thread_id}: {data.get('error', 'unknown error')}"))
return
messages = data.get("messages", [])
print("\n" + c_bold(f"=== THREAD {thread_id} ({agent}) — {len(messages)} messages ===") + "\n")
for msg in messages:
@@ -3280,6 +3308,11 @@ def build_parser():
p_va_rb.add_argument("name", help="Variable name")
p_va_rb.add_argument("--revision", default=None, help="Revision step (int) or timestamp")
# Domain: muse (fast headless gateway via muse-cli-node with isolated per-node Cloudflare WARP egress)
p_muse = subparsers.add_parser("muse", parents=[common], help="Direct headless gateway client (muse-cli-node)")
p_muse.add_argument("node", choices=VALID_NODES, help="Target agent node")
p_muse.add_argument("muse_args", nargs=argparse.REMAINDER, help="Arguments passed directly to muse-cli-node")
return parser
def main():
@@ -3446,6 +3479,10 @@ def main():
cmd_loop_strat(args)
elif args.domain == "vars":
cmd_loop_vars(args)
elif args.domain == "muse":
cmd = [str(BIN_DIR / "muse-cli-node"), args.node] + (args.muse_args or [])
res = subprocess.run(cmd)
sys.exit(res.returncode)
else:
parser.print_help()
+303
View File
@@ -0,0 +1,303 @@
#!/usr/bin/env python3
"""
Thread lifecycle manager for side chat rotation.
Evaluates rotation triggers (age, message count, explicit flag) for a reuse_key
and performs rotation when triggered. Integrates with:
- job-sidechats.json (state + lifecycle policy)
- sidechat_manager.py (thread operations)
- box follow-up system (thread_id migration)
Usage:
python3 thread_lifecycle.py check --reuse-key <key> [--state PATH]
python3 thread_lifecycle.py rotate --reuse-key <key> --reason <reason> [--state PATH]
The `check` command is designed to run at job dispatch time, before the reuse
check in job-dispatch.py. It returns:
exit 0: no rotation needed, reuse existing thread
exit 2: rotation triggered, new thread needed (caller should create fresh)
exit 1: error
This is a prototype. Not deployed.
"""
import argparse
import json
import os
import sys
import time
from datetime import datetime, timezone, timedelta
# Default lifecycle policy (applied when no per-key policy exists)
DEFAULT_POLICY = {
"max_age_hours": 168, # 7 days
"max_messages": 200,
"channel_type": "job",
}
# Channel-type defaults (override DEFAULT_POLICY)
CHANNEL_DEFAULTS = {
"heartbeat": {"max_age_hours": 24, "max_messages": 500, "channel_type": "heartbeat"},
"job": {"max_age_hours": 168, "max_messages": 200, "channel_type": "job"},
"coordination": {"max_age_hours": 720, "max_messages": 200, "channel_type": "coordination"},
"dm": {"max_age_hours": None, "max_messages": None, "channel_type": "dm"}, # never rotate
}
def _utcnow():
return datetime.now(timezone.utc)
def _parse_ts(ts_str):
"""Parse ISO-8601 UTC timestamp. Returns datetime or None."""
if not ts_str:
return None
try:
# Handle both "2026-10-04T12:00:00Z" and "... +00:00" forms
ts_str = ts_str.replace("Z", "+00:00")
dt = datetime.fromisoformat(ts_str)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
except (ValueError, TypeError):
return None
def load_state(state_path):
"""Load job-sidechats.json. Returns dict."""
if not os.path.exists(state_path):
return {}
with open(state_path) as f:
return json.load(f)
def save_state(state_path, state):
"""Save state atomically (write temp + rename)."""
tmp = state_path + ".tmp"
with open(tmp, "w") as f:
json.dump(state, f, indent=2)
os.rename(tmp, state_path)
def get_policy(state, reuse_key):
"""
Get lifecycle policy for a reuse_key.
Priority: per-key policy > channel-type default > global default.
Returns dict with max_age_hours, max_messages, channel_type.
"""
lifecycle = state.get("_lifecycle", {})
defaults = state.get("_lifecycle_defaults", {})
# Start with global defaults
policy = dict(DEFAULT_POLICY)
policy.update(defaults)
# Apply channel-type defaults if specified
key_policy = lifecycle.get(reuse_key, {})
channel_type = key_policy.get("channel_type")
if channel_type and channel_type in CHANNEL_DEFAULTS:
policy.update(CHANNEL_DEFAULTS[channel_type])
# Per-key overrides win
policy.update(key_policy)
return policy
def get_thread_age_hours(state, reuse_key):
"""
Get thread age in hours from creation.
Returns None if unknown (fail-open: age trigger skipped).
"""
# Prefer explicit creation timestamp
created_key = f"{reuse_key}:created_at"
created = _parse_ts(state.get(created_key))
if created:
delta = _utcnow() - created
return delta.total_seconds() / 3600
# Fall back to last rotation timestamp
rotated = _parse_ts(state.get(f"{reuse_key}:rotated_at"))
if rotated:
delta = _utcnow() - rotated
return delta.total_seconds() / 3600
return None
def check_rotation_needed(state, reuse_key, message_count=None, force=False):
"""
Evaluate rotation triggers for a reuse_key.
Args:
state: loaded state dict
reuse_key: the reuse key to check
message_count: current message count (None = unknown, skip count trigger)
force: explicit rotation request (topic change / manual)
Returns:
(needed: bool, reason: str, details: dict)
"""
policy = get_policy(state, reuse_key)
# DM channels never rotate
if policy.get("channel_type") == "dm":
return False, "dm_no_rotate", {}
# Explicit force always wins
if force:
return True, "explicit", {"requested": True}
# Age trigger
max_age = policy.get("max_age_hours")
if max_age is not None:
age_hours = get_thread_age_hours(state, reuse_key)
if age_hours is not None and age_hours >= max_age:
return True, "age", {
"age_hours": round(age_hours, 1),
"max_age_hours": max_age,
}
# Message count trigger
max_msg = policy.get("max_messages")
if max_msg is not None and message_count is not None:
if message_count >= max_msg:
return True, "message_count", {
"message_count": message_count,
"max_messages": max_msg,
}
return False, "none", {}
def resolve_thread_for_nudge(state, thread_id):
"""
Resolve a thread_id for nudge delivery, handling rotation.
If thread_id points to an archived/rotated thread, follow the
:previous chain to find the current UUID.
Returns the UUID to use (may be the input if no rotation found).
"""
# Build reverse map: old_uuid -> reuse_key
for key, value in state.items():
if key.startswith("_") or ":" in key:
continue
prev_key = f"{key}:previous"
if state.get(prev_key) == thread_id:
# This thread was rotated; use the current UUID
current = state.get(key)
if current and current != thread_id:
return current
return thread_id
def record_rotation(state, reuse_key, old_uuid, new_uuid, reason):
"""
Update state after rotation. Returns updated state dict.
Does not save — caller saves.
"""
now = _utcnow().strftime("%Y-%m-%dT%H:%M:%SZ")
prev_count = state.get(f"{reuse_key}:rotation_count", 0)
state[f"{reuse_key}:previous"] = old_uuid
state[reuse_key] = new_uuid
state[f"{reuse_key}:rotated_at"] = now
state[f"{reuse_key}:created_at"] = now # new thread creation baseline
state[f"{reuse_key}:rotation_reason"] = reason
state[f"{reuse_key}:rotation_count"] = prev_count + 1
return state
def cmd_check(args):
"""Check if rotation is needed. Exit 0=no, 2=yes, 1=error."""
state = load_state(args.state)
reuse_key = args.reuse_key
if reuse_key not in state:
print(f"No existing thread for {reuse_key}, no rotation check needed")
return 0
needed, reason, details = check_rotation_needed(
state, reuse_key,
message_count=args.messages,
force=args.force,
)
if needed:
print(json.dumps({"rotate": True, "reason": reason, "details": details}))
return 2
else:
print(json.dumps({"rotate": False, "reason": reason}))
return 0
def cmd_rotate(args):
"""Record a rotation (called after new thread is created)."""
state = load_state(args.state)
reuse_key = args.reuse_key
old_uuid = state.get(reuse_key)
if not old_uuid:
print(f"No existing thread for {reuse_key}", file=sys.stderr)
return 1
if not args.new_uuid:
print("--new-uuid required", file=sys.stderr)
return 1
state = record_rotation(state, reuse_key, old_uuid, args.new_uuid, args.reason)
save_state(args.state, state)
print(json.dumps({
"rotated": True,
"reuse_key": reuse_key,
"old_uuid": old_uuid,
"new_uuid": args.new_uuid,
"reason": args.reason,
}))
return 0
def cmd_resolve(args):
"""Resolve a thread_id through rotation chain (for nudge router)."""
state = load_state(args.state)
resolved = resolve_thread_for_nudge(state, args.thread_id)
print(json.dumps({
"input": args.thread_id,
"resolved": resolved,
"rotated": resolved != args.thread_id,
}))
return 0
def main():
p = argparse.ArgumentParser(description="Side chat thread lifecycle manager")
p.add_argument("--state", default=os.path.expanduser(
"~/Projects/NetVM/job-sidechats.json"),
help="Path to job-sidechats.json")
sub = p.add_subparsers(dest="cmd", required=True)
c = sub.add_parser("check", help="Check if rotation is needed")
c.add_argument("--reuse-key", required=True)
c.add_argument("--messages", type=int, default=None,
help="Current message count (None = skip count trigger)")
c.add_argument("--force", action="store_true",
help="Force rotation (explicit topic change)")
r = sub.add_parser("rotate", help="Record a completed rotation")
r.add_argument("--reuse-key", required=True)
r.add_argument("--new-uuid", required=True)
r.add_argument("--reason", default="manual")
v = sub.add_parser("resolve", help="Resolve thread_id through rotation")
v.add_argument("--thread-id", required=True)
args = p.parse_args()
if args.cmd == "check":
sys.exit(cmd_check(args))
elif args.cmd == "rotate":
sys.exit(cmd_rotate(args))
elif args.cmd == "resolve":
sys.exit(cmd_resolve(args))
if __name__ == "__main__":
main()