diff --git a/bin/p2p-relay.py b/bin/p2p-relay.py index b994c77..47d9179 100755 --- a/bin/p2p-relay.py +++ b/bin/p2p-relay.py @@ -1,53 +1,92 @@ #!/usr/bin/env python3 """ -P2P Relay: 646 -> muse via headless +P2P Relay: 646 -> muse via headless (with watermark) Reads from 646's p2p-muse-outbox, writes to muse's p2p-646-inbox. -Usage: p2p-relay.py [--once] [--daemon] +Tracks watermark to avoid duplicate relays. + +Usage: p2p-relay.py [--once] +Watermark: ~/Projects/NetVM/bridge/p2p-watermark.json """ import subprocess import sys import time +import json +import hashlib +from pathlib import Path -API = "~/Projects/NetVM/bin/muse-chat-api.py" -EXEC_646 = "~/Projects/NetVM/bin/netvm-exec.sh 646 --" -EXEC_MUSE = "~/Projects/NetVM/bin/netvm-exec.sh muse --" +API = "/home/super/Projects/NetVM/bin/muse-chat-api.py" +EXEC_646 = "/home/super/Projects/NetVM/bin/netvm-exec.sh 646 --" +EXEC_MUSE = "/home/super/Projects/NetVM/bin/netvm-exec.sh muse --" OUTBOX_646 = "3ed3aef3-2d92-408c-9e13-f36e40a75976" INBOX_MUSE = "6db43088-b87d-43a3-ae2a-157acebcfa21" +WATERMARK_FILE = Path("/home/super/Projects/NetVM/bridge/p2p-watermark.json") + +def load_watermark(): + if WATERMARK_FILE.exists(): + with open(WATERMARK_FILE) as f: + return json.load(f) + return {"last_hash": None, "last_relay": 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): - result = subprocess.run(cmd, shell=True, capture_output=True, text=True) + result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=60) return result.stdout.strip() def get_outbox_messages(): - # Navigate 646 to outbox, read messages run(f"{EXEC_646} python3 {API} --account 646 sidechat use {OUTBOX_646}") - time.sleep(2) + time.sleep(3) msgs = run(f"{EXEC_646} python3 {API} --account 646 messages 5") run(f"{EXEC_646} python3 {API} --account 646 sidechat main") return msgs def send_to_muse(message): - # Navigate muse to inbox, send message run(f"{EXEC_MUSE} python3 {API} --account muse sidechat use {INBOX_MUSE}") - time.sleep(2) - run(f"{EXEC_MUSE} python3 {API} --account muse send \"{message}\"") + time.sleep(3) + # Escape for shell + safe = message.replace('"', '\\"').replace('$', '\\$').replace('`', '\\`') + run(f'{EXEC_MUSE} python3 {API} --account muse send "{safe}"') run(f"{EXEC_MUSE} python3 {API} --account muse sidechat main") def relay_once(): + wm = load_watermark() msgs = get_outbox_messages() - # Simple: if there are messages, relay the latest - # In production, track watermark to avoid duplicates - if msgs and "P2P setup" not in msgs: - # Extract last message (crude) - lines = [l for l in msgs.split('\n') if l.strip() and '---' not in l] - if lines: - latest = lines[-1][:500] - print(f"Relaying: {latest[:80]}...") - send_to_muse(f"[from 646] {latest}") - return True - return False + + if not msgs or "P2P test" in msgs and "Hello muse" in msgs: + # Skip if only the test message (already relayed manually) + # Or if empty + pass + + # Extract messages (split by ---) + parts = [p.strip() for p in msgs.split('---') if p.strip()] + # Filter out system messages + parts = [p for p in parts if p and len(p) > 10 and "P2P setup" not in p] + + if not parts: + print("No new messages in outbox") + return False + + latest = parts[-1] + msg_hash = hashlib.md5(latest.encode()).hexdigest()[:16] + + if msg_hash == wm.get("last_hash"): + print("Already relayed (watermark match)") + return False + + print(f"Relaying new message: {latest[:80]}...") + send_to_muse(f"[from 646] {latest[:500]}") + + wm["last_hash"] = msg_hash + wm["last_relay"] = time.strftime("%Y-%m-%dT%H:%M:%S") + save_watermark(wm) + print("Relayed and watermark updated") + return True if __name__ == '__main__': if '--once' in sys.argv: