diff --git a/bin/thread-listener.py b/bin/thread-listener.py new file mode 100755 index 0000000..18bfabc --- /dev/null +++ b/bin/thread-listener.py @@ -0,0 +1,102 @@ +#!/usr/bin/env python3 +""" +THREAD listener: Auto-executes 646's THREAD commands + +Polls 646's main chat for THREAD commands, executes via headless. +Format: THREAD + +Watermark prevents duplicate execution. +""" +import subprocess +import time +import json +import re +from pathlib import Path + +API = "/home/super/Projects/NetVM/bin/muse-chat-api.py" +NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh" +WATERMARK_FILE = Path("/home/super/Projects/NetVM/bridge/thread-watermark.json") + +def load_watermark(): + if WATERMARK_FILE.exists(): + with open(WATERMARK_FILE) as f: + return json.load(f) + return {"last_processed": None} + +def save_watermark(data): + WATERMARK_FILE.parent.mkdir(parents=True, exist_ok=True) + with open(WATERMARK_FILE, 'w') as f: + json.dump(data, f, indent=2) + +def run(cmd, timeout=60): + result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout) + return result.stdout.strip() + +def get_646_messages(): + return run(f"{NETVM_EXEC} 646 -- python3 {API} --account 646 messages 5") + +def execute_thread(agent, target, message): + """Execute a THREAD command via headless.""" + print(f"THREAD: {agent} -> {target}: {message[:60]}...") + + # Navigate to target + if target == "main": + run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main") + else: + # Assume it's a chat_id, navigate to thread + run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {target}") + + time.sleep(2) + + # Send message with 646 attribution + safe = message.replace('"', '\\"').replace('$', '\\$').replace('`', '\\`')[:500] + run(f'{NETVM_EXEC} {agent} -- python3 {API} --account {agent} send "[646] {safe}"') + + # Back to main + run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main") + print("Done") + +def poll_once(): + wm = load_watermark() + msgs = get_646_messages() + + # Find THREAD commands (not yet processed) + # Simple: look for "THREAD " at start of a message part + parts = [p.strip() for p in msgs.split('---') if p.strip()] + + for part in parts: + # Match THREAD + m = re.match(r'THREAD\s+(\w+)\s+(\w+)\s+(.+)', part, re.DOTALL) + if m: + agent, target, message = m.groups() + # Create a simple ID for deduplication + cmd_id = f"{agent}:{target}:{message[:30]}" + if cmd_id == wm.get("last_processed"): + continue # Already processed + + # Validate agent + if agent not in ["muse", "pip", "646"]: + print(f"Unknown agent: {agent}, skipping") + continue + + try: + execute_thread(agent, target, message.strip()) + wm["last_processed"] = cmd_id + save_watermark(wm) + except Exception as e: + print(f"THREAD failed: {e}") + +if __name__ == '__main__': + import sys + if '--once' in sys.argv: + poll_once() + elif '--daemon' in sys.argv: + print("THREAD listener running (poll every 60s)...") + while True: + try: + poll_once() + except Exception as e: + print(f"Poll error: {e}") + time.sleep(60) + else: + print("Use --once or --daemon")