2026-10-04 16:56:15 +00:00
|
|
|
#!/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)
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:58:15 +00:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:56:15 +00:00
|
|
|
def get_pipeline(run_id):
|
|
|
|
|
data = load_pipelines()
|
2026-10-04 16:58:15 +00:00
|
|
|
if run_id in data:
|
|
|
|
|
return data[run_id]
|
|
|
|
|
for k, v in data.items():
|
|
|
|
|
if k.startswith(run_id):
|
|
|
|
|
return v
|
|
|
|
|
return None
|
2026-10-04 16:56:15 +00:00
|
|
|
|
|
|
|
|
|
|
|
|
|
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]
|
2026-10-04 16:58:15 +00:00
|
|
|
|