From 2b19d2cb456cb8208c83354246d653a6a38efc77 Mon Sep 17 00:00:00 2001 From: operator Date: Sun, 4 Oct 2026 16:58:15 +0000 Subject: [PATCH] feat(pipeline): add stop and prune commands to pipeline engine and CLI - Add stop_pipeline and prune_pipelines to bin/pipeline_engine.py - Support prefix matching and custom cancellation reason - Wire 'super pipeline stop ' and 'super pipeline prune [--max-age H]' into bin/super-cli.py --- bin/pipeline_engine.py | 55 +++++++++++++++++++++++++++++++++++++++++- bin/super-cli.py | 37 ++++++++++++++++++++++++++++ 2 files changed, 91 insertions(+), 1 deletion(-) diff --git a/bin/pipeline_engine.py b/bin/pipeline_engine.py index 8ad85a3..dac79f8 100644 --- a/bin/pipeline_engine.py +++ b/bin/pipeline_engine.py @@ -158,9 +158,61 @@ def fail_pipeline(run_id, reason="failed"): 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() - return data.get(run_id) + 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(): @@ -176,3 +228,4 @@ def list_pipeline_history(limit=20): reverse=True, ) return sorted_runs[:limit] + diff --git a/bin/super-cli.py b/bin/super-cli.py index 0f3d046..64c5ab9 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -2271,6 +2271,32 @@ def cmd_pipeline_history(args): print_table(headers, rows) print(c_dim("\n Commands: super pipeline status | super pipeline run \n")) + +def cmd_pipeline_stop(args): + import pipeline_engine + run_id = args.run_id.strip() + reason = getattr(args, "reason", "cancelled_by_operator") or "cancelled_by_operator" + entry = pipeline_engine.stop_pipeline(run_id, reason=reason) + if not entry: + print(c_red(f"Error: Pipeline run '{run_id}' not found."), file=sys.stderr) + sys.exit(1) + print(c_green(f"✔ Pipeline run '{entry.get('run_id')}' stopped/canceled.")) + if getattr(args, "json", False): + print(json.dumps(entry, indent=2)) + + +def cmd_pipeline_prune(args): + import pipeline_engine + max_age = getattr(args, "max_age", 1) or 1 + pruned = pipeline_engine.prune_pipelines(max_age_hours=max_age) + if pruned: + print(c_green(f"✔ Pruned {len(pruned)} stale pipeline runs (> {max_age}h): {', '.join(pruned)}")) + else: + print(c_dim(f"✔ No stale running pipelines found (> {max_age}h).")) + if getattr(args, "json", False): + print(json.dumps({"pruned": pruned, "count": len(pruned)}, indent=2)) + + # --------------------------------------------------------------------------- # Domain: LOOP (Intrinsic Loop Strategy, Health, Break Taxonomy, and Control) # --------------------------------------------------------------------------- @@ -3019,6 +3045,13 @@ def build_parser(): p_pipe_hist = pipe_sub.add_parser("history", parents=[common], help="Show historical pipeline runs and outcomes") p_pipe_hist.add_argument("-n", "--limit", type=int, default=20, help="Number of records to display") + p_pipe_stop = pipe_sub.add_parser("stop", parents=[common], help="Cancel/stop an active pipeline run") + p_pipe_stop.add_argument("run_id", help="Pipeline run ID (prefix match supported)") + p_pipe_stop.add_argument("--reason", default="cancelled_by_operator", help="Cancellation reason") + + p_pipe_prune = pipe_sub.add_parser("prune", parents=[common], help="Prune/timeout stale running pipelines") + p_pipe_prune.add_argument("--max-age", type=float, default=1.0, help="Max age in hours (default 1.0)") + # Top-level domain aliases: strat and vars p_strat_alias = subparsers.add_parser("strat", parents=[common], help="Shortcut for 'loop strat'") strat_alias_sub = p_strat_alias.add_subparsers(dest="strat_action") @@ -3186,6 +3219,10 @@ def main(): cmd_pipeline_run(args) elif act == "history": cmd_pipeline_history(args) + elif act == "stop": + cmd_pipeline_stop(args) + elif act == "prune": + cmd_pipeline_prune(args) else: parser.print_help() elif args.domain == "loop":