From 75e6e745a977f676065b1c007b12c1c4ceb02002 Mon Sep 17 00:00:00 2001 From: operator Date: Sun, 4 Oct 2026 23:02:04 +0000 Subject: [PATCH] feat(core-loop): implement dual-monitoring, fast gateway dispatch, dedicated pip tasks sidechat, and 2m cadence --- ACCOUNTS.md | 4 +- bin/cred-client.py | 33 ++++++ bin/meta-acct.py | 10 +- bin/self_main_loop.py | 202 ++++++++++++++++++++++++++++++--- bin/super-cli.py | 27 +++++ job-sidechats.json | 6 + systemd/self-main-loop.service | 5 + systemd/self-main-loop.timer | 7 ++ 8 files changed, 273 insertions(+), 21 deletions(-) create mode 100644 systemd/self-main-loop.service create mode 100644 systemd/self-main-loop.timer diff --git a/ACCOUNTS.md b/ACCOUNTS.md index 960bcdb..df89fb7 100644 --- a/ACCOUNTS.md +++ b/ACCOUNTS.md @@ -27,8 +27,8 @@ node name, chrome-box profile, API `--account`, and the agent's display name. | agent | node | profile | login_type | meta_label | email | phone_otp | instagram_linked | status | egress_ip | cdp_port | display_name | notes | |-------|------|---------|------------|------------|-------|-----------|------------------|--------|-----------|----------|--------------|-------| | muse | muse | muse | email_otp | ltd.pixels.ltd@gmail.com | ltd.pixels.ltd@gmail.com | no | unknown | active | 104.28.195.181 | 9410 | muse | Main dev agent. Logged in 2026-10-03 via email OTP on bl. | -| pip | pip | pip | phone_otp | piparada | - | yes | yes | active | 104.28.195.181 | 9420 | pip | Phone OTP login 2026-10-03. Renamed from 'Muse' to 'pip'. Session lost on browser restart 2026-10-03; needs re-auth. IG avatar visible in account selector. | -| 646 | 646 | 646 | phone_otp | Meta Account | - | yes | unknown | active | 104.28.195.181 | 9430 | 646 | Shares phone number with piparada's account. Logged in 2026-10-03 via phone OTP (first Meta Account option). Node created 2026-10-03 (warp-646, CDP 9430). Browser up, session active (verified 2026-10-03). | +| pip | pip | pip | phone_otp | piparada | io.antonio.parada@gmail.com | yes | yes | active | 104.28.195.181 | 9420 | pip | Phone OTP login. Linked with IG piparada, email io.antonio.parada@gmail.com. Verified active 2026-10-04. | +| 646 | 646 | 646 | phone_otp | Meta Account | - | yes | no | active | 104.28.195.181 | 9430 | 646 | Shares phone number with piparada's account. Logged in 2026-10-03 via phone OTP (first Meta Account option). Node created 2026-10-03 (warp-646, CDP 9430). Browser up, session active. | | def | def | def | email_otp | defnotabotnet@gmail.com | defnotabotnet@gmail.com | no | yes | active | 104.28.195.181 | 9450 | def | Full onboarding completed 2026-10-04; age verification cleared via Instagram linking (paradahub). Active chat session. | | opm | opm | opm | email_otp | Nico Parada | artglobal.cc@gmail.com | no | yes | active | 104.28.195.181 | 9440 | opm | Email changed from yourfriendnico@proton.me to artglobal.cc@gmail.com. Linked with IG auxfate. Browser up, session active. | | dev | dev | dev | email_otp | paradaproduced@gmail.com | paradaproduced@gmail.com | no | yes | active | 104.28.195.181 | 9460 | dev | Full onboarding completed 2026-10-04; unlocked /access gate via Meta Accounts Center IG linking (veryraremeta). Active chat session. | diff --git a/bin/cred-client.py b/bin/cred-client.py index 6cab12d..14dc330 100755 --- a/bin/cred-client.py +++ b/bin/cred-client.py @@ -266,6 +266,19 @@ class CredClient: "email_notified": email_result.get("sent") if email_result else False } + def audit_meta(self, node): + """Query Meta Accounts Center for linked profiles and security status.""" + meta_script = os.path.join(NETVM_DIR, "bin", "meta-acct.py") + cmd = ["sudo", "-n", "ip", "netns", "exec", f"warp-{node}", sys.executable, meta_script, "list-linked", node] + try: + res = subprocess.run(cmd, capture_output=True, text=True, timeout=15) + lines = res.stdout.strip().splitlines() + if lines: + return json.loads(lines[-1]) + return {"error": res.stderr.strip() or "No output from meta-acct.py"} + except Exception as e: + return {"error": str(e)} + def main(): common = argparse.ArgumentParser(add_help=False) @@ -296,12 +309,32 @@ def main(): p_link.add_argument("--node", required=True, help="Node label") p_link.add_argument("--notify", action="store_true", help="Send email alert to operator via local MTA") + # meta-audit + p_meta = sub.add_parser("meta-audit", parents=[common], help="Query Meta Accounts Center for linked profiles and security status") + p_meta.add_argument("--node", required=True, help="Node label") + # list sub.add_parser("list", parents=[common], help="List all registered nodes and vitality statuses") args = p.parse_args() client = CredClient() + if args.command == "meta-audit": + res = client.audit_meta(args.node) + if args.json: + print(json.dumps(res, indent=2)) + else: + print(f"\n=== META ACCOUNTS CENTER AUDIT: {args.node} ===") + if res.get("error"): + print(f"Error: {res['error']}", file=sys.stderr) + sys.exit(1) + print(f"Meta Account Email: {res.get('email') or '(none / phone-only)'}") + profiles = res.get("profiles", []) + print(f"Linked Profiles ({len(profiles)}):") + for p_info in profiles: + print(f" - [{p_info.get('type')}] {p_info.get('name')}") + sys.exit(0) + if args.command == "link-instagram": res = client.link_instagram(args.node, notify=args.notify) if args.json: diff --git a/bin/meta-acct.py b/bin/meta-acct.py index 7bbf76c..3dda330 100755 --- a/bin/meta-acct.py +++ b/bin/meta-acct.py @@ -13,7 +13,15 @@ Scope: Accounts Center linkage/security surface only. No writes. """ import json, sys, time, urllib.request, os -NETVM = os.environ.get("NETVM_DIR", os.path.expanduser("~/Projects/NetVM")) +NETVM = os.environ.get("NETVM_DIR") +if not NETVM: + script_dir = os.path.dirname(os.path.abspath(__file__)) + parent = os.path.dirname(script_dir) + if os.path.exists(os.path.join(parent, "NODES.md")): + NETVM = parent + else: + NETVM = "/home/super/Projects/NetVM" + AC_BASE = "https://accountscenter.meta.com" def cdp_port(agent): diff --git a/bin/self_main_loop.py b/bin/self_main_loop.py index b6da317..4c6a71a 100755 --- a/bin/self_main_loop.py +++ b/bin/self_main_loop.py @@ -54,8 +54,8 @@ DEFAULT_AGENTS = ["muse", "pip", "646", "opm"] # Prompting sidechat per agent — mirrors box-ctl.py NOTIFY_SIDECHATS. DEFAULT_PROMPT_SIDECHAT = { "646": "646 tasks", - "opm": "heartbeat", - "pip": "646-pip-coord", + "opm": "main-loop brain", + "pip": "pip tasks", "muse": "muse tasks", } DEFAULT_SENDER = "opm" # neutral sender, mirrors `box notify` @@ -119,7 +119,14 @@ def get_config(st): agents = cfg.get("agents") or list(DEFAULT_AGENTS) agents = [a for a in agents if a in DEFAULT_AGENTS] prompt = dict(DEFAULT_PROMPT_SIDECHAT) - prompt.update(cfg.get("prompt_sidechat") or {}) + user_prompt = cfg.get("prompt_sidechat") or {} + for a, sc in user_prompt.items(): + # Migrate obsolete defaults + if a == "pip" and sc in ("646-pip-coord", "heartbeat"): + continue + if a == "opm" and sc in ("heartbeat",): + continue + prompt[a] = sc enabled = dict(cfg.get("enabled") or {}) # Default any unlisted agent to enabled. for a in agents: @@ -206,8 +213,56 @@ def compose_digest(agent, new): return digest[:DIGEST_MAX] +def get_monitored_sidechats(agent): + """Return dict of {thread_uuid: name} for sidechats registered for this agent.""" + sc_file = os.path.join(BASE, "job-sidechats.json") + if not os.path.exists(sc_file): + return {} + try: + with open(sc_file, "r", encoding="utf-8") as f: + sc_data = json.load(f) + except Exception: + return {} + + monitored = {} + for name, item in sc_data.items(): + if not isinstance(item, dict): + continue + uuid = item.get("thread_uuid") or item.get("uuid") + if not uuid: + continue + item_agent = item.get("agent") + if name.endswith(f"@{agent}") or (item_agent == agent and "@" not in name): + if name.startswith("pipe-") or name.startswith("test-"): + continue + monitored[uuid] = name + return monitored + + def send_prompt(sender, agent, sidechat, digest): - """Post the digest to the agent's prompting sidechat via dm.py.""" + """Post the digest to the agent's prompting sidechat. + Tries fast direct gateway send via muse_hybrid first, falling back to dm.py. + """ + target_uuid = sidechat + try: + import dm + target_uuid = dm.resolve_sidechat_target(sidechat, agent) + except Exception: + pass + + # Try fast direct gateway send + if target_uuid and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target_uuid.lower()): + try: + import muse_hybrid + send_node = agent if agent in DEFAULT_AGENTS else sender + res, err = muse_hybrid.send_message(send_node, digest, thread_id=target_uuid, wait=0) + if res and not err: + return True, "sent (gateway)" + log("Gateway send fallback for %s/%s due to: %s" % (agent, sidechat, err or res)) + except Exception as e: + log("Gateway send exception for %s/%s: %s" % (agent, sidechat, e)) + + # Fallback to dm.py send cmd = [sys.executable, DM_PY, "send", "--agent", sender, "--to", agent, "--target", sidechat, digest] try: @@ -219,7 +274,108 @@ def send_prompt(sender, agent, sidechat, digest): if r.returncode != 0: tail = ((r.stderr or "") + (r.stdout or "")).strip()[-300:] return False, "dm.py send failed rc=%d: %s" % (r.returncode, tail) - return True, "sent" + return True, "sent (dm.py)" + + +def check_sidechats(agent, cfg, wm, threads_meta): + """Tier 1 & Tier 2 unread check on agent's registered sidechats.""" + import muse_hybrid + monitored = get_monitored_sidechats(agent) + prompt_sidechat = cfg["prompt_sidechat"].get(agent) + prompt_uuid = None + try: + import dm + prompt_uuid = dm.resolve_sidechat_target(prompt_sidechat, agent) + except Exception: + pass + + sc_watermarks = wm.setdefault("sidechats", {}) + threads_by_id = {t.get("session_id"): t for t in (threads_meta or []) if isinstance(t, dict)} + + new_activity = [] + + for uuid, sc_name in monitored.items(): + # Do not monitor prompt destination to avoid loopback + if prompt_uuid and uuid.lower() == prompt_uuid.lower(): + continue + + th_meta = threads_by_id.get(uuid) + if not th_meta: + continue + + last_updated = sc_watermarks.get(uuid, {}).get("updated") + curr_updated = th_meta.get("updated") + + # First run: if no watermark, initialize to current newest message without alerting + if uuid not in sc_watermarks: + msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=3) + last_id = msgs[-1]["message_id"] if (msgs and not err and msgs) else None + sc_watermarks[uuid] = {"last_id": last_id, "updated": curr_updated} + continue + + # Tier 1 fast check: if updated timestamp unchanged, skip + if last_updated and curr_updated and last_updated == curr_updated: + continue + + # Tier 2 check: fetch recent messages + msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=10) + if err or not msgs: + continue + + last_id = sc_watermarks.get(uuid, {}).get("last_id") + formatted = [{ + "id": m.get("message_id") or f"msg-{m.get('seq')}", + "from": {"name": m.get("role", "unknown")}, + "role": m.get("role", "unknown"), + "text": m.get("text", "") + } for m in msgs] + + new_msgs = new_messages(formatted, {"last_id": last_id}) + incoming = [m for m in new_msgs if m.get("role") != "assistant"] + + if incoming: + new_activity.append({ + "sidechat": sc_name, + "uuid": uuid, + "messages": incoming + }) + + # Advance watermark + sc_watermarks[uuid] = { + "last_id": formatted[-1]["id"] if formatted else None, + "updated": curr_updated + } + + return new_activity + + +def compose_dual_digest(agent, main_new, sc_activity): + """Compose concise digest combining main chat and sidechat activity.""" + lines = [] + if main_new: + lines.append("New messages in %s main chat (%d):" % (agent, len(main_new))) + for m in main_new[:3]: + text = (m.get("text") or "").strip().replace("\n", " ") + q = " (?)" if "?" in text else "" + hot = " (!)" if (OPERATOR_RE.search(text) or "urgent" in text.lower()) else "" + lines.append("- %s: %s%s%s" % (author_of(m), text[:PREVIEW_MAX], q, hot)) + if len(main_new) > 3: + lines.append("(+%d more)" % (len(main_new) - 3)) + + if sc_activity: + if lines: + lines.append("") + lines.append("Incoming sidechat activity for %s:" % agent) + for item in sc_activity: + sc_name = item["sidechat"] + msgs = item["messages"] + lines.append("- [%s] (%d new):" % (sc_name, len(msgs))) + for m in msgs[:2]: + text = (m.get("text") or "").strip().replace("\n", " ") + lines.append(" * %s: %s" % (author_of(m), text[:PREVIEW_MAX])) + + lines.append("Check chat via box when available.") + return "\n".join(lines)[:DIGEST_MAX] def do_check(only_agent=None): @@ -235,11 +391,15 @@ def do_check(only_agent=None): prompted = 0 errors = 0 + import muse_hybrid + for agent in agents: if not enabled.get(agent, True): results[agent] = {"ok": True, "new": 0, "disabled": True} continue wm = agents_state.get(agent) or {} + + # 1. Main chat check ok, payload = read_main_chat(agent, cfg["read_limit"]) if not ok: log("%s: READ FAILED: %s" % (agent, payload)) @@ -247,33 +407,39 @@ def do_check(only_agent=None): errors += 1 continue messages = payload - new = new_messages(messages, wm) + main_new = new_messages(messages, wm) newest = messages[-1] if messages else None - if not new: - # Silent: advance watermark to newest seen, no sidechat post. - if newest: - agents_state[agent] = {"last_id": newest.get("id"), - "last_ts": newest.get("ts")} + if newest: + wm["last_id"] = newest.get("id") + wm["last_ts"] = newest.get("ts") + + # 2. Sidechats dual-monitoring check + threads_meta, _ = muse_hybrid.get_threads(agent) + sc_activity = check_sidechats(agent, cfg, wm, threads_meta) + + agents_state[agent] = wm + + total_new = len(main_new) + sum(len(a["messages"]) for a in sc_activity) + if total_new == 0: results[agent] = {"ok": True, "new": 0} continue - digest = compose_digest(agent, new) + + digest = compose_dual_digest(agent, main_new, sc_activity) sidechat = cfg["prompt_sidechat"].get(agent) if not sidechat: log("%s: no prompting sidechat configured, skipping" % agent) results[agent] = {"ok": False, "error": "no prompting sidechat"} errors += 1 continue + sent, detail = send_prompt(cfg["sender"], agent, sidechat, digest) if sent: - agents_state[agent] = {"last_id": new[-1].get("id"), - "last_ts": new[-1].get("ts")} prompted += 1 - log("%s: prompted %s with %d new" % (agent, sidechat, len(new))) - results[agent] = {"ok": True, "new": len(new), "prompted": sidechat} + log("%s: prompted %s with %d new (%s)" % (agent, sidechat, total_new, detail)) + results[agent] = {"ok": True, "new": total_new, "prompted": sidechat, "detail": detail} else: - # Keep watermark: next iteration retries the same digest. log("%s: PROMPT SEND FAILED: %s" % (agent, detail)) - results[agent] = {"ok": False, "error": detail, "new": len(new)} + results[agent] = {"ok": False, "error": detail, "new": total_new} errors += 1 # Save under the state lock with a fresh reload: an enable/disable may diff --git a/bin/super-cli.py b/bin/super-cli.py index 412fba8..d60d503 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -1967,6 +1967,28 @@ def cmd_cred_link_instagram(args): print(f" Email Alert : {n_badge}") print("\n" + c_yellow(" → Tap either Tailscale link from your phone/browser to complete Meta age verification.") + "\n") +def cmd_cred_meta_audit(args): + import cred_client + client = cred_client.CredClient() + res = client.audit_meta(args.node) + if getattr(args, "json", False): + print(json.dumps(res, indent=2)) + else: + print("\n" + c_bold(f"=== META ACCOUNTS CENTER AUDIT: {args.node} ===") + "\n") + if res.get("error"): + print(c_err(f" Error: {res['error']}")) + sys.exit(1) + email_str = res.get('email') or c_dim('(none / phone-only)') + print(f" Meta Account Email: {c_cyan(email_str)}") + profiles = res.get("profiles", []) + print(f" Linked Profiles ({len(profiles)}):") + for p_info in profiles: + p_type = p_info.get("type", "unknown") + p_name = p_info.get("name", "unknown") + badge = c_green(f"[{p_type}]") if p_type == "instagram" else c_magenta(f"[{p_type}]") + print(f" - {badge} {c_bold(p_name)}") + print() + def cmd_cred_list(args): import cred_client client = cred_client.CredClient() @@ -3134,6 +3156,9 @@ def build_parser(): p_c_link.add_argument("--node", required=True, help="Node label") p_c_link.add_argument("--notify", action="store_true", help="Send email alert to operator via local MTA") + p_c_audit = cred_sub.add_parser("meta-audit", parents=[common], help="Audit Meta Accounts Center for linked profiles and email") + p_c_audit.add_argument("--node", required=True, help="Node label") + cred_sub.add_parser("list", parents=[common], help="List all registered nodes and vitality statuses") # Domain: HARVEST @@ -3415,6 +3440,8 @@ def main(): cmd_cred_status(args) elif act == "link-instagram": cmd_cred_link_instagram(args) + elif act == "meta-audit": + cmd_cred_meta_audit(args) else: parser.print_help() elif args.domain == "harvest": diff --git a/job-sidechats.json b/job-sidechats.json index 5567e49..dd4bbd1 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -62,6 +62,12 @@ "agent": "646", "created_at": "2026-10-04T18:30:21.863666+00:00" }, + "pip tasks": { + "thread_uuid": "5f33b1ab-a310-4cbb-9ccc-bb7bcd3d1b34", + "agent": "pip", + "description": "Dedicated pip task prompting sidechat", + "created_at": "2026-10-04T22:59:13+00:00" + }, "pipe-c3f686": { "thread_uuid": "afa977ba-c75b-4b3e-8c1f-f83f13fb287e", "agent": "opm", diff --git a/systemd/self-main-loop.service b/systemd/self-main-loop.service new file mode 100644 index 0000000..a7c1efb --- /dev/null +++ b/systemd/self-main-loop.service @@ -0,0 +1,5 @@ +[Unit] +Description=NetVM main-chat self-monitor loop +[Service] +Type=oneshot +ExecStart=/home/super/Projects/NetVM/bin/self_main_loop.py check diff --git a/systemd/self-main-loop.timer b/systemd/self-main-loop.timer new file mode 100644 index 0000000..bf3a705 --- /dev/null +++ b/systemd/self-main-loop.timer @@ -0,0 +1,7 @@ +[Unit] +Description=Run dual-monitoring self-loop every 2 minutes +[Timer] +OnCalendar=*:0/2 +Persistent=true +[Install] +WantedBy=timers.target