thread-listener: auto-execute 646 THREAD commands
This commit is contained in:
Executable
+102
@@ -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 <agent> <main|chat_id> <message>
|
||||
|
||||
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 <agent> <target> <message>
|
||||
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")
|
||||
Reference in New Issue
Block a user