Files
box/bin/pipeline_engine.py
T
operator d6ca600394 feat(relay): add subagent spawn, thread tools, and container box client
- Upgrade exec-constrained.py with subagent.spawn, thread.list, thread.view, pipeline.run ops
- Grant full ops permissions to all fleet agent identities (646, pip, muse, opm)
- Implement bin/box-relay.sh zero-dependency client supporting bearer and SSH signature auth
- Add fast hybrid gateway path to dm.py for sub-2s verified deliveries
- Fix wait=0 handling in super-cli.py subagent deployments
- Add hourly check-in jobs and scheduler for 646, pip, muse
- Document agent tooling and relay APIs in docs/AGENT-TOOLING.md
2026-10-04 23:30:18 +00:00

265 lines
8.0 KiB
Python

#!/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]}"
# Fast headless gateway provisioning: create session server-side
thread_uuid = None
job_file = JOBS_DIR / f"{pipeline_name}.json"
root_agent = "opm"
if job_file.exists():
try:
with open(job_file) as f:
root_agent = json.load(f).get("agent", "opm")
except Exception:
pass
try:
import muse_hybrid
res, err = muse_hybrid.start_session(root_agent, title=target)
if res and not err:
thread_uuid = res.get("session_id")
# Register in job-sidechats.json so dm.py resolves it immediately
sc_file = NETVM_ROOT / "job-sidechats.json"
if sc_file.exists():
with open(sc_file, "r", encoding="utf-8") as f:
sc_data = json.load(f)
sc_data[target] = {
"thread_uuid": thread_uuid,
"agent": root_agent,
"created_at": utcnow()
}
tmp_sc = f"{sc_file}.tmp.{os.getpid()}"
with open(tmp_sc, "w", encoding="utf-8") as f:
json.dump(sc_data, f, indent=2)
os.replace(tmp_sc, sc_file)
except Exception as e:
sys.stderr.write(f"warning: fast gateway session creation failed: {e}\n")
entry = {
"run_id": run_id,
"pipeline_name": pipeline_name,
"created_at": utcnow(),
"updated_at": utcnow(),
"status": "running",
"target": target,
"thread_uuid": thread_uuid,
"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 stop_pipeline(run_id, reason="cancelled_by_operator"):
"""Manually cancel or stop an active pipeline run."""
data = load_pipelines()
# Support prefix matching
target_key = None
for k in data.keys():
if k == run_id or k.startswith(run_id):
target_key = k
break
if not target_key:
return None
entry = data[target_key]
entry["status"] = "canceled"
entry["canceled_at"] = utcnow()
entry["updated_at"] = utcnow()
entry["cancel_reason"] = reason
save_pipelines(data)
return entry
def prune_pipelines(max_age_hours=24):
"""Mark stale running pipeline runs older than max_age_hours as timed_out."""
data = load_pipelines()
now = datetime.now(timezone.utc)
pruned = []
for k, v in data.items():
if v.get("status") == "running":
created_str = v.get("created_at")
if created_str:
try:
dt = datetime.fromisoformat(created_str.replace("Z", "+00:00"))
diff_h = (now - dt).total_seconds() / 3600.0
if diff_h >= max_age_hours:
v["status"] = "timed_out"
v["completed_at"] = utcnow()
v["updated_at"] = utcnow()
v["timeout_reason"] = f"stale_exceeded_{max_age_hours}h"
pruned.append(k)
except Exception:
pass
if pruned:
save_pipelines(data)
return pruned
def get_pipeline(run_id):
data = load_pipelines()
if run_id in data:
return data[run_id]
for k, v in data.items():
if k.startswith(run_id):
return v
return None
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]