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 <run_id>' and 'super pipeline prune [--max-age H]' into bin/super-cli.py
This commit is contained in:
+54
-1
@@ -158,9 +158,61 @@ def fail_pipeline(run_id, reason="failed"):
|
|||||||
save_pipelines(data)
|
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):
|
def get_pipeline(run_id):
|
||||||
data = load_pipelines()
|
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():
|
def list_active_pipelines():
|
||||||
@@ -176,3 +228,4 @@ def list_pipeline_history(limit=20):
|
|||||||
reverse=True,
|
reverse=True,
|
||||||
)
|
)
|
||||||
return sorted_runs[:limit]
|
return sorted_runs[:limit]
|
||||||
|
|
||||||
|
|||||||
@@ -2271,6 +2271,32 @@ def cmd_pipeline_history(args):
|
|||||||
print_table(headers, rows)
|
print_table(headers, rows)
|
||||||
print(c_dim("\n Commands: super pipeline status | super pipeline run <name>\n"))
|
print(c_dim("\n Commands: super pipeline status | super pipeline run <name>\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)
|
# 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 = 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_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
|
# Top-level domain aliases: strat and vars
|
||||||
p_strat_alias = subparsers.add_parser("strat", parents=[common], help="Shortcut for 'loop strat'")
|
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")
|
strat_alias_sub = p_strat_alias.add_subparsers(dest="strat_action")
|
||||||
@@ -3186,6 +3219,10 @@ def main():
|
|||||||
cmd_pipeline_run(args)
|
cmd_pipeline_run(args)
|
||||||
elif act == "history":
|
elif act == "history":
|
||||||
cmd_pipeline_history(args)
|
cmd_pipeline_history(args)
|
||||||
|
elif act == "stop":
|
||||||
|
cmd_pipeline_stop(args)
|
||||||
|
elif act == "prune":
|
||||||
|
cmd_pipeline_prune(args)
|
||||||
else:
|
else:
|
||||||
parser.print_help()
|
parser.print_help()
|
||||||
elif args.domain == "loop":
|
elif args.domain == "loop":
|
||||||
|
|||||||
Reference in New Issue
Block a user