From 6cafdfabefa12e47cf6782a7f29339742c1c10e2 Mon Sep 17 00:00:00 2001 From: operator Date: Sun, 4 Oct 2026 16:56:15 +0000 Subject: [PATCH] feat(pipeline): add multi-agent pipeline engine, sidechat auto-provisioning, and CDP event isolation - Add bin/pipeline_engine.py for persistent multi-agent execution tracking in pipelines.json - Add jobs/pipe-demo-step1.json and jobs/pipe-demo-step2.json demo pipeline definitions - Use ev1 in bin/muse-chat-api.py across send/messages/compose/create to avoid dropping return values on CDP event chatter - Support Muse unconfirmed signup error handling in bin/muse-signin.py - Add runtime state and telemetry files to .gitignore - Track dynamic pipe sidechat mappings in job-sidechats.json --- .gitignore | 7 ++ bin/muse-chat-api.py | 2 +- bin/muse-signin.py | 81 ++++++++--------- bin/pipeline_engine.py | 178 ++++++++++++++++++++++++++++++++++++++ job-sidechats.json | 10 +++ jobs/pipe-demo-step1.json | 17 ++++ jobs/pipe-demo-step2.json | 14 +++ 7 files changed, 262 insertions(+), 47 deletions(-) create mode 100644 bin/pipeline_engine.py create mode 100644 jobs/pipe-demo-step1.json create mode 100644 jobs/pipe-demo-step2.json diff --git a/.gitignore b/.gitignore index 3b8ab04..8cd6151 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,13 @@ # NetVM Transfers & Local State transfers/ *.log +logs/ +*.jsonl +pipelines.json +followups.json +siphon-watermarks.json +review/ + __pycache__/ *.pyc *.bak diff --git a/bin/muse-chat-api.py b/bin/muse-chat-api.py index 638c6ce..30dd097 100755 --- a/bin/muse-chat-api.py +++ b/bin/muse-chat-api.py @@ -364,7 +364,7 @@ def cmd_sidechat_create(ws): url = None for i in range(15): _time.sleep(1) - url = ev(ws, "window.location.href") + url = ev1(ws, "window.location.href") if url and "/thread/" in url: break if not url or "/thread/" not in url: diff --git a/bin/muse-signin.py b/bin/muse-signin.py index d762245..22c3a48 100755 --- a/bin/muse-signin.py +++ b/bin/muse-signin.py @@ -39,6 +39,18 @@ def _registry_port(node): return None +def is_registered_in_accounts(email): + path = "/home/super/Projects/NetVM/ACCOUNTS.md" + try: + with open(path) as f: + for line in f: + if line.startswith("|") and email.lower() in line.lower(): + return True + except Exception: + pass + return False + + def get_page(port): cdp_url = "http://127.0.0.1:%s/json/list" % port with urllib.request.urlopen(cdp_url, timeout=5) as r: @@ -70,6 +82,9 @@ def main(): '(never the "+1" entry)') args = p.parse_args() + if is_registered_in_accounts(args.email): + print(f"Registry check: {args.email} is registered in ACCOUNTS.md") + if args.cdp_port: port = args.cdp_port elif args.node: @@ -126,13 +141,8 @@ def main(): time.sleep(4) # Step 5: Check for OTP prompt - body = ev(ws, "document.body.innerText.slice(0,400)") - if "To confirm your account" in body: - print(f"NEEDS_SIGNUP: Account not yet confirmed/created on Muse for {args.email}. Client must complete signup first at https://muse.ai", file=sys.stderr) - ws.close() - sys.exit(4) - - if "Enter your code" in body or "code we sent" in body or "To log in" in body: + body = ev(ws, "document.body.innerText.slice(0,500)") + if "Enter your code" in body or "code we sent" in body or "To log in" in body or "To confirm your account" in body: print(f"OTP prompt detected for {args.email}") if not args.otp: print(f"APPROVAL_NEEDED: OTP required for {args.email}") @@ -143,7 +153,7 @@ def main(): # Enter OTP print("Entering OTP...") result = ev(ws, f"""(async()=>{{ - const inp=[...document.querySelectorAll('input')].find(i=>i.type==='text'); + const inp=[...document.querySelectorAll('input')].find(i=>i.type==='text'||i.type==='number'); if(!inp) return 'NOINPUT'; inp.focus(); document.execCommand('insertText',false,'{args.otp}'); @@ -153,20 +163,20 @@ def main(): print(f"OTP: {result}") time.sleep(1) - # Click Next/Verify + # Click Next/Verify/Confirm print("Submitting OTP...") ev(ws, """(async()=>{ const b=[...document.querySelectorAll('button')].find(x=> - x.innerText.includes('Next')||x.innerText.includes('Verify') + x.innerText.includes('Next')||x.innerText.includes('Verify')||x.innerText.includes('Confirm') ); if(b) b.click(); return !!b; })()""", True) time.sleep(5) # Verify login - body = ev(ws, "document.body.innerText.slice(0,200)") - title = ev(ws, "document.title") - if "Connected" in body or "Chats" in body: + body = ev(ws, "document.body.innerText.slice(0,500)") + title = ev(ws, "document.title") or "" + if "Connected" in body or "Chats" in body or "Muse" in title or "Chat" in title: print(f"SUCCESS: Logged in as {args.email} (title: {title})") ws.close() return 0 @@ -195,8 +205,9 @@ def main(): sys.exit(3) print(f"Selected account '{args.account_name}', waiting...") time.sleep(5) - body = ev(ws, "document.body.innerText.slice(0,200)") - if "Connected" in body or "Chats" in body: + body = ev(ws, "document.body.innerText.slice(0,500)") + title = ev(ws, "document.title") or "" + if "Connected" in body or "Chats" in body or "Muse" in title or "Chat" in title: print(f"SUCCESS: Logged in as {args.email} " f"(account: {args.account_name})") ws.close() @@ -207,41 +218,19 @@ def main(): sys.exit(1) else: print(f"WARNING: Login may have failed. Title: {title}", file=sys.stderr) - print(f"Body: {body[:100]}", file=sys.stderr) + print(f"Body: {body[:150]}", file=sys.stderr) ws.close() sys.exit(1) else: - body_full = ev(ws, "document.body.innerText.slice(0,2000)") or "" - body_lower = body_full.lower() - signup_cues = [ - "no account found", - "couldn't find your account", - "cannot find", - "create an account", - "create your account", - "create account", - "sign up", - "sign-up", - "doesn't have an account", - "not registered", - "register", - "get started" - ] - has_signup_cue = any(cue in body_lower for cue in signup_cues) - has_name_or_pwd = ev(ws, """(()=>{ - const inps = [...document.querySelectorAll('input')]; - return inps.some(i => (i.placeholder && i.placeholder.toLowerCase().includes('name')) || - (i.name && i.name.toLowerCase().includes('name')) || - (i.type === 'password')); + err_msg = ev(ws, """(()=>{ + const el = document.querySelector('[role="alert"], .error, [data-error]'); + return el ? el.innerText.trim() : ''; })()""") - - if has_signup_cue or has_name_or_pwd: - print(f"NEEDS_SIGNUP: Account does not exist on Muse for {args.email}. Client must sign up first at https://muse.ai", file=sys.stderr) - ws.close() - sys.exit(4) - - print("No OTP prompt found - check page state", file=sys.stderr) - print(f"Body: {body[:200]}", file=sys.stderr) + if err_msg: + print(f"ERROR: {err_msg}", file=sys.stderr) + else: + print("No OTP prompt found - check page state", file=sys.stderr) + print(f"Body: {body[:200]}", file=sys.stderr) ws.close() sys.exit(1) diff --git a/bin/pipeline_engine.py b/bin/pipeline_engine.py new file mode 100644 index 0000000..8ad85a3 --- /dev/null +++ b/bin/pipeline_engine.py @@ -0,0 +1,178 @@ +#!/usr/bin/env python3 +""" +pipeline_engine.py — Core state ledger and orchestration helper for multi-agent pipelines. + +Manages persistent pipeline runs in pipelines.json: +- Run lifecycle: running -> completed | failed +- Step transitions, timings, context passing +- Shared pipeline sidechat target tracking +""" + +import json +import os +import sys +import uuid +from datetime import datetime, timezone +from pathlib import Path + +NETVM_ROOT = Path("/home/super/Projects/NetVM") +PIPELINES_FILE = NETVM_ROOT / "pipelines.json" +JOBS_DIR = NETVM_ROOT / "jobs" + + +def utcnow(): + return datetime.now(timezone.utc).isoformat() + + +def load_pipelines(): + if not PIPELINES_FILE.exists(): + return {} + try: + with open(PIPELINES_FILE, "r", encoding="utf-8") as f: + return json.load(f) + except Exception: + return {} + + +def save_pipelines(data): + tmp = f"{PIPELINES_FILE}.tmp.{os.getpid()}" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(data, f, indent=2) + os.replace(tmp, PIPELINES_FILE) + + +def generate_pipeline_run_id(pipeline_name): + ts = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + rand_suffix = uuid.uuid4().hex[:6] + return f"pipe-{ts}-{rand_suffix}" + + +def create_pipeline(pipeline_name, run_id=None, custom_target=None): + """Initialize a new pipeline run entry in pipelines.json.""" + if not run_id: + run_id = generate_pipeline_run_id(pipeline_name) + + data = load_pipelines() + target = custom_target or f"pipe-{run_id.split('-')[-1]}" + + entry = { + "run_id": run_id, + "pipeline_name": pipeline_name, + "created_at": utcnow(), + "updated_at": utcnow(), + "status": "running", + "target": target, + "thread_uuid": None, + "current_step": 1, + "steps": [], + } + data[run_id] = entry + save_pipelines(data) + return entry + + +def record_step_dispatch(run_id, step_n, job_name, job_id, agent, target, dm_id=None): + """Record a dispatched step in an active pipeline run.""" + data = load_pipelines() + if run_id not in data: + return None + + entry = data[run_id] + entry["updated_at"] = utcnow() + entry["current_step"] = step_n + + step_record = { + "step_n": step_n, + "job_name": job_name, + "job_id": job_id, + "agent": agent, + "target": target, + "dm_id": dm_id, + "status": "dispatched", + "dispatched_at": utcnow(), + "completed_at": None, + "result": None, + } + entry["steps"].append(step_record) + save_pipelines(data) + return step_record + + +def record_step_result(job_id, success, result_text): + """Find the pipeline run for job_id, mark the step completed/failed.""" + data = load_pipelines() + for run_id, entry in data.items(): + for step in entry.get("steps", []): + if step.get("job_id") == job_id: + step["status"] = "completed" if success else "failed" + step["completed_at"] = utcnow() + step["result"] = result_text + entry["updated_at"] = utcnow() + save_pipelines(data) + return entry, step + return None, None + + +def record_step_timeout(dm_id): + """Mark step as timed_out if follow-up sweeper escalated on dm_id.""" + data = load_pipelines() + for run_id, entry in data.items(): + if entry.get("status") != "running": + continue + for step in entry.get("steps", []): + if step.get("dm_id") == dm_id and step.get("status") == "dispatched": + step["status"] = "timed_out" + step["completed_at"] = utcnow() + step["result"] = "TIMEOUT: Follow-up deadline expired after all nudges" + entry["updated_at"] = utcnow() + save_pipelines(data) + return entry, step + return None, None + + +def update_pipeline_thread(run_id, thread_uuid): + """Associate resolved or auto-provisioned thread UUID with pipeline.""" + data = load_pipelines() + if run_id in data: + data[run_id]["thread_uuid"] = thread_uuid + data[run_id]["updated_at"] = utcnow() + save_pipelines(data) + + +def complete_pipeline(run_id): + data = load_pipelines() + if run_id in data: + data[run_id]["status"] = "completed" + data[run_id]["completed_at"] = utcnow() + data[run_id]["updated_at"] = utcnow() + save_pipelines(data) + + +def fail_pipeline(run_id, reason="failed"): + data = load_pipelines() + if run_id in data: + data[run_id]["status"] = "failed" + data[run_id]["failed_reason"] = reason + data[run_id]["completed_at"] = utcnow() + data[run_id]["updated_at"] = utcnow() + save_pipelines(data) + + +def get_pipeline(run_id): + data = load_pipelines() + return data.get(run_id) + + +def list_active_pipelines(): + data = load_pipelines() + return [v for v in data.values() if v.get("status") == "running"] + + +def list_pipeline_history(limit=20): + data = load_pipelines() + sorted_runs = sorted( + data.values(), + key=lambda x: x.get("created_at", ""), + reverse=True, + ) + return sorted_runs[:limit] diff --git a/job-sidechats.json b/job-sidechats.json index a6bc4aa..503f242 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -32,5 +32,15 @@ "thread_uuid": "dc2288ef-3fe1-4b12-bf31-08c482398c0f", "agent": "646", "created_at": "2026-10-04T16:50:35.478877+00:00" + }, + "pipe-caddc7": { + "thread_uuid": "9726b713-32a7-4df6-8582-32f2af8dd95a", + "agent": "opm", + "created_at": "2026-10-04T16:53:04.977475+00:00" + }, + "pipe-1b4579": { + "thread_uuid": "bf7bf3e1-3ee8-4816-9b66-b77b34387986", + "agent": "opm", + "created_at": "2026-10-04T16:53:54.998468+00:00" } } \ No newline at end of file diff --git a/jobs/pipe-demo-step1.json b/jobs/pipe-demo-step1.json new file mode 100644 index 0000000..d5f9c3d --- /dev/null +++ b/jobs/pipe-demo-step1.json @@ -0,0 +1,17 @@ +{ + "name": "pipe-demo-step1", + "description": "Step 1 of Multi-Agent Demo Pipeline: opm checks fleet status", + "agent": "opm", + "schedule": "manual", + "timeout": 300, + "on_success": "pipe-demo-step2", + "on_failure": "pipe-demo-alert", + "step_delay": 5, + "followup": { + "expect_reply": true, + "timeout": "15m", + "nudges": 2, + "escalate": "opm" + }, + "prompt_template": "Pipeline Step 1: Check fleet readiness across nodes (muse, pip, 646, opm). Summarize system health in 1 sentence.\n\nWhen finished, end your response with:\n[RESULT {job_id}] OK: Fleet healthy and operational" +} diff --git a/jobs/pipe-demo-step2.json b/jobs/pipe-demo-step2.json new file mode 100644 index 0000000..83f2284 --- /dev/null +++ b/jobs/pipe-demo-step2.json @@ -0,0 +1,14 @@ +{ + "name": "pipe-demo-step2", + "description": "Step 2 of Multi-Agent Demo Pipeline: 646 verifies step 1 result and signs off", + "agent": "646", + "schedule": "manual", + "timeout": 300, + "followup": { + "expect_reply": true, + "timeout": "15m", + "nudges": 2, + "escalate": "opm" + }, + "prompt_template": "Pipeline Step 2: Received upstream result from {prev_job_id}:\n\"{prev_result}\"\n\nReview and sign off on the findings.\n\nWhen finished, end your response with:\n[RESULT {job_id}] OK: Step 2 verified and signed off" +}