257 lines
7.4 KiB
Python
Executable File
257 lines
7.4 KiB
Python
Executable File
#!/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()
|