p2p-relay: add watermark to prevent duplicates
This commit is contained in:
+60
-21
@@ -1,54 +1,93 @@
|
|||||||
#!/usr/bin/env python3
|
#!/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.
|
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 subprocess
|
||||||
import sys
|
import sys
|
||||||
import time
|
import time
|
||||||
|
import json
|
||||||
|
import hashlib
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
API = "~/Projects/NetVM/bin/muse-chat-api.py"
|
API = "/home/super/Projects/NetVM/bin/muse-chat-api.py"
|
||||||
EXEC_646 = "~/Projects/NetVM/bin/netvm-exec.sh 646 --"
|
EXEC_646 = "/home/super/Projects/NetVM/bin/netvm-exec.sh 646 --"
|
||||||
EXEC_MUSE = "~/Projects/NetVM/bin/netvm-exec.sh muse --"
|
EXEC_MUSE = "/home/super/Projects/NetVM/bin/netvm-exec.sh muse --"
|
||||||
|
|
||||||
OUTBOX_646 = "3ed3aef3-2d92-408c-9e13-f36e40a75976"
|
OUTBOX_646 = "3ed3aef3-2d92-408c-9e13-f36e40a75976"
|
||||||
INBOX_MUSE = "6db43088-b87d-43a3-ae2a-157acebcfa21"
|
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):
|
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()
|
return result.stdout.strip()
|
||||||
|
|
||||||
def get_outbox_messages():
|
def get_outbox_messages():
|
||||||
# Navigate 646 to outbox, read messages
|
|
||||||
run(f"{EXEC_646} python3 {API} --account 646 sidechat use {OUTBOX_646}")
|
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")
|
msgs = run(f"{EXEC_646} python3 {API} --account 646 messages 5")
|
||||||
run(f"{EXEC_646} python3 {API} --account 646 sidechat main")
|
run(f"{EXEC_646} python3 {API} --account 646 sidechat main")
|
||||||
return msgs
|
return msgs
|
||||||
|
|
||||||
def send_to_muse(message):
|
def send_to_muse(message):
|
||||||
# Navigate muse to inbox, send message
|
|
||||||
run(f"{EXEC_MUSE} python3 {API} --account muse sidechat use {INBOX_MUSE}")
|
run(f"{EXEC_MUSE} python3 {API} --account muse sidechat use {INBOX_MUSE}")
|
||||||
time.sleep(2)
|
time.sleep(3)
|
||||||
run(f"{EXEC_MUSE} python3 {API} --account muse send \"{message}\"")
|
# 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")
|
run(f"{EXEC_MUSE} python3 {API} --account muse sidechat main")
|
||||||
|
|
||||||
def relay_once():
|
def relay_once():
|
||||||
|
wm = load_watermark()
|
||||||
msgs = get_outbox_messages()
|
msgs = get_outbox_messages()
|
||||||
# Simple: if there are messages, relay the latest
|
|
||||||
# In production, track watermark to avoid duplicates
|
if not msgs or "P2P test" in msgs and "Hello muse" in msgs:
|
||||||
if msgs and "P2P setup" not in msgs:
|
# Skip if only the test message (already relayed manually)
|
||||||
# Extract last message (crude)
|
# Or if empty
|
||||||
lines = [l for l in msgs.split('\n') if l.strip() and '---' not in l]
|
pass
|
||||||
if lines:
|
|
||||||
latest = lines[-1][:500]
|
# Extract messages (split by ---)
|
||||||
print(f"Relaying: {latest[:80]}...")
|
parts = [p.strip() for p in msgs.split('---') if p.strip()]
|
||||||
send_to_muse(f"[from 646] {latest}")
|
# Filter out system messages
|
||||||
return True
|
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
|
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 __name__ == '__main__':
|
||||||
if '--once' in sys.argv:
|
if '--once' in sys.argv:
|
||||||
relay_once()
|
relay_once()
|
||||||
|
|||||||
Reference in New Issue
Block a user