#!/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 , 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()