#!/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]