feat(core): node veth idempotency, multi-node cdp ports, brain last ids persistence, and pulse interval variable

This commit is contained in:
operator
2026-10-05 15:59:18 +00:00
parent 94d6502289
commit d529d9ebab
6 changed files with 117 additions and 22 deletions
+1 -1
View File
@@ -59,7 +59,7 @@ WARN_AFTER = 10.0 # log a warning when a waiter waits this long
POLL_INTERVAL = 0.05 # ticket/slot poll cadence POLL_INTERVAL = 0.05 # ticket/slot poll cadence
QUEUE_DIR = "/tmp/cdp-queue" QUEUE_DIR = "/tmp/cdp-queue"
VALID_NODES = ("muse", "pip", "646", "opm") VALID_NODES = ("muse", "pip", "646", "opm", "dev", "def")
class QueueTimeout(Exception): class QueueTimeout(Exception):
+58 -2
View File
@@ -135,10 +135,64 @@ def cmd_login_activity(port):
section = body[idx:idx+1500] if idx >= 0 else "" section = body[idx:idx+1500] if idx >= 0 else ""
return {"where_logged_in": section} return {"where_logged_in": section}
def cmd_link_instagram(agent, port):
"""Automate Meta Accounts Center linking flow and second-click OAuth handoff."""
import websocket
tab = get_ac_tab(port)
ws = websocket.create_connection(tab["webSocketDebuggerUrl"], timeout=15)
try:
# Step 1: Navigate to manage accounts
ws.send(json.dumps({"id": 1, "method": "Page.navigate", "params": {"url": f"{AC_BASE}/manage/"}}))
ws.recv()
time.sleep(3)
# Step 2: Look for 'Add profiles and devices' button or check if already linked
expr_find_add = """(() => {
const btns = [...document.querySelectorAll("div[role='button'], button, a")];
const addBtn = btns.find(b => /add profiles|add accounts/i.test((b.innerText||"").trim()));
if (addBtn) {
addBtn.click();
return {status: "clicked_add", text: addBtn.innerText};
}
return {status: "no_add_button"};
})()"""
ws.send(json.dumps({"id": 2, "method": "Runtime.evaluate", "params": {"expression": expr_find_add, "returnByValue": True}}))
res2 = json.loads(ws.recv()).get("result", {}).get("result", {}).get("value", {})
time.sleep(3)
# Step 3: Check dialog / popup for Instagram option or FXCAL handoff
expr_handle_dialog = """(() => {
const btns = [...document.querySelectorAll("div[role='button'], button, a")];
// Check for Add Instagram button or Continue/Confirm
const igBtn = btns.find(b => /instagram/i.test((b.innerText||"").trim()) && /add|connect/i.test((b.innerText||"").trim()));
if (igBtn) {
igBtn.click();
return {status: "clicked_instagram_option"};
}
const confirmBtn = btns.find(b => /continue|confirm|yes, finish/i.test((b.innerText||"").trim()));
if (confirmBtn) {
confirmBtn.click();
return {status: "clicked_confirm", text: confirmBtn.innerText};
}
return {status: "dialog_scanned", current_url: window.location.href};
})()"""
ws.send(json.dumps({"id": 3, "method": "Runtime.evaluate", "params": {"expression": expr_handle_dialog, "returnByValue": True}}))
res3 = json.loads(ws.recv()).get("result", {}).get("result", {}).get("value", {})
time.sleep(2)
return {
"status": "success",
"step1_add": res2,
"step2_dialog": res3,
"current_url": tab.get("url")
}
finally:
ws.close()
def main(): def main():
if len(sys.argv) != 3: if len(sys.argv) < 3:
print(__doc__.strip().split("\n")[0]) print(__doc__.strip().split("\n")[0])
print("Usage: meta-acct.py <list-linked|security-status|login-activity> <agent>") print("Usage: meta-acct.py <list-linked|security-status|login-activity|link-instagram> <agent>")
sys.exit(2) sys.exit(2)
cmd, agent = sys.argv[1], sys.argv[2] cmd, agent = sys.argv[1], sys.argv[2]
port = cdp_port(agent) port = cdp_port(agent)
@@ -149,6 +203,8 @@ def main():
out = cmd_security_status(port) out = cmd_security_status(port)
elif cmd == "login-activity": elif cmd == "login-activity":
out = cmd_login_activity(port) out = cmd_login_activity(port)
elif cmd == "link-instagram":
out = cmd_link_instagram(agent, port)
else: else:
raise SystemExit(f"meta-acct: unknown command '{cmd}'") raise SystemExit(f"meta-acct: unknown command '{cmd}'")
except Exception as e: except Exception as e:
+3
View File
@@ -24,6 +24,9 @@ if ! nsexec ip link show "$VPEER" >/dev/null 2>&1; then
nsexec ip addr add "${PEER_IP}/30" dev "$VPEER" 2>/dev/null || true nsexec ip addr add "${PEER_IP}/30" dev "$VPEER" 2>/dev/null || true
nsexec ip link set "$VPEER" up nsexec ip link set "$VPEER" up
fi fi
# Idempotent: re-assert host-side GW addr even if veth pre-existed (2026-10-04: def/dev lost theirs)
ip addr add "${GW}/30" dev "$VETH" 2>/dev/null || true
ip link set "$VETH" up
nsexec ip link set lo up nsexec ip link set lo up
# host NAT + forwarding for the veth subnet # host NAT + forwarding for the veth subnet
+19 -3
View File
@@ -574,9 +574,10 @@ def read_brain_messages(agent, sidechat_name, since_ts):
Resolves sidechat_name -> UUID via dm (never hardcoded; UUIDs rotate), Resolves sidechat_name -> UUID via dm (never hardcoded; UUIDs rotate),
reads via the same muse_hybrid primitive the loop uses, returns reads via the same muse_hybrid primitive the loop uses, returns
[{"sender", "text", "ts"}]. Tracks last message_id per (agent, name); [{"sender", "text", "ts"}]. Tracks last message_id per (agent, name)
first run anchors at newest with no backfill (same policy as in the module cache, which do_check() persists to the watermark JSON
new_messages for main chat). (brain.brain_last_ids) across ticks. First run anchors at newest with
no backfill (same policy as new_messages for main chat).
""" """
try: try:
import dm import dm
@@ -635,6 +636,14 @@ def do_check(only_agent=None):
try: try:
from brain import BrainWorkspace from brain import BrainWorkspace
brain = BrainWorkspace(STATE_FILE, read_fn=read_brain_messages) brain = BrainWorkspace(STATE_FILE, read_fn=read_brain_messages)
# Restore persisted brain message IDs so !loop commands posted
# between ticks are not missed (module cache is per-process).
for k, v in (brain.state.get("brain_last_ids") or {}).items():
try:
ag, nm = k.split("|", 1)
_BRAIN_LAST_IDS[(ag, nm)] = v
except ValueError:
pass
intake = brain.intake() intake = brain.intake()
if intake.get("commands"): if intake.get("commands"):
log("brain: %d commands, %d acks, %d ignored-senders" % ( log("brain: %d commands, %d acks, %d ignored-senders" % (
@@ -787,6 +796,13 @@ def do_check(only_agent=None):
watermarks=wm_epochs, watermarks=wm_epochs,
health=health, health=health,
) )
# Persist brain message IDs for the next tick.
try:
brain.state["brain_last_ids"] = {
"%s|%s" % k: v for k, v in _BRAIN_LAST_IDS.items()
if v is not None}
except Exception:
pass
brain.save() brain.save()
except Exception as e: except Exception as e:
log("brain post_thinking failed (non-fatal): %r" % e) log("brain post_thinking failed (non-fatal): %r" % e)
+32 -16
View File
@@ -45,42 +45,59 @@ def get_tailscale_dns():
pass pass
return "bl.tailfb5960.ts.net" return "bl.tailfb5960.ts.net"
def get_node_cdp_port(node):
nodes_file = os.path.join(NETVM_DIR, "NODES.md")
pinned = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9460}
if node in pinned:
return pinned[node]
try:
with open(nodes_file) as f:
for line in f:
if not line.strip().startswith("|"): continue
cells = [c.strip() for c in line.strip("|").split("|")]
if len(cells) >= 4 and cells[0] == node and cells[3].isdigit():
return int(cells[3])
except Exception:
pass
return 9450
def fetch_node_ig_link(node): def fetch_node_ig_link(node):
"""Query CDP within node netns to get fresh Meta Accounts Center OAuth URL.""" """Query CDP within node netns to get fresh Meta Accounts Center OAuth URL."""
script = """ port = get_node_cdp_port(node)
script = f"""
import json, websocket, sys, urllib.request import json, websocket, sys, urllib.request
try: try:
tabs = json.loads(urllib.request.urlopen("http://127.0.0.1:9450/json").read()) tabs = json.loads(urllib.request.urlopen("http://127.0.0.1:{port}/json").read())
except Exception as e: except Exception as e:
print(json.dumps({"error": f"CDP unreachable: {e}"})) print(json.dumps({{"error": f"CDP unreachable: {{e}}"}}))
sys.exit(0) sys.exit(0)
muse_tab = next((t for t in tabs if "muse.ai" in t.get("url", "") and t.get("type") == "page"), None) muse_tab = next((t for t in tabs if "muse.ai" in t.get("url", "") and t.get("type") == "page"), None)
if not muse_tab: if not muse_tab:
print(json.dumps({"error": "No muse.ai tab open"})) print(json.dumps({{"error": "No muse.ai tab open"}}))
sys.exit(0) sys.exit(0)
try: try:
ws = websocket.create_connection(muse_tab["webSocketDebuggerUrl"], timeout=5) ws = websocket.create_connection(muse_tab["webSocketDebuggerUrl"], timeout=5)
expr = '''(async()=>{ expr = \'\'\'(async()=>{{
try { try {{
const r = await fetch('/api/hatch/age-confirmation/linking-web-auth?account_type=instagram', { const r = await fetch('/api/hatch/age-confirmation/linking-web-auth?account_type=instagram', {{
headers: {'Accept': 'application/json'} headers: {{'Accept': 'application/json'}}
}); }});
return await r.json(); return await r.json();
} catch(e) { return {error: String(e)}; } }} catch(e) {{ return {{error: String(e)}}; }}
})()''' }})()\'\'\'
ws.send(json.dumps({"id": 1, "method": "Runtime.evaluate", "params": {"expression": expr, "awaitPromise": True, "returnByValue": True}})) ws.send(json.dumps({{"id": 1, "method": "Runtime.evaluate", "params": {{"expression": expr, "awaitPromise": True, "returnByValue": True}}}}))
while True: while True:
msg = json.loads(ws.recv()) msg = json.loads(ws.recv())
if msg.get("id") == 1: if msg.get("id") == 1:
val = msg.get("result", {}).get("result", {}).get("value", {}) val = msg.get("result", {{}}).get("result", {{}}).get("value", {{}})
print(json.dumps(val)) print(json.dumps(val))
break break
ws.close() ws.close()
except Exception as e: except Exception as e:
print(json.dumps({"error": f"CDP evaluate failed: {e}"})) print(json.dumps({{"error": f"CDP evaluate failed: {{e}}"}}))
""" """
cmd = ["sudo", "-n", "ip", "netns", "exec", f"warp-{node}", sys.executable, "-c", script] cmd = ["sudo", "-n", "ip", "netns", "exec", f"warp-{node}", sys.executable, "-c", script]
try: try:
@@ -96,8 +113,7 @@ except Exception as e:
def check_node_cleared(node): def check_node_cleared(node):
"""Check if node has transitioned past access/verification into active chat.""" """Check if node has transitioned past access/verification into active chat."""
checker = os.path.join(NETVM_DIR, "bin", "accounts-health.py") checker = os.path.join(NETVM_DIR, "bin", "accounts-health.py")
# Lookup port from registry or default to 9450 for def port = get_node_cdp_port(node)
port = 9450
cmd = ["sudo", "-n", "ip", "netns", "exec", f"warp-{node}", sys.executable, checker, str(port)] cmd = ["sudo", "-n", "ip", "netns", "exec", f"warp-{node}", sys.executable, checker, str(port)]
try: try:
res = subprocess.run(cmd, capture_output=True, text=True, timeout=10) res = subprocess.run(cmd, capture_output=True, text=True, timeout=10)
+4
View File
@@ -464,6 +464,10 @@ _BUILTIN_SPECS = {
"default": 7200, "type": "int", "min": 60, "max": 604800, "unit": "seconds", "default": 7200, "type": "int", "min": 60, "max": 604800, "unit": "seconds",
"description": "Timeout before nudging unacknowledged wake digests (2h default)" "description": "Timeout before nudging unacknowledged wake digests (2h default)"
}, },
"pulse_interval_m": {
"default": 30, "type": "int", "min": 5, "max": 1440, "unit": "minutes",
"description": "Cadence in minutes between autonomous operator pulse cycles"
},
} }
# Built-in fallbacks # Built-in fallbacks