Fix completed runs not appearing in list + add purge-failed endpoint
- Update save_run_cache to also update actor_id, recipe, inputs on conflict - Add logging for actor_id when saving runs to run_cache - Add admin endpoint DELETE /runs/admin/purge-failed to delete all failed runs Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -128,10 +128,25 @@ class RunService:
|
||||
# Only return as completed if we have an output
|
||||
# (runs with no output should be re-executed)
|
||||
if output_cid:
|
||||
# Also fetch recipe content from pending_runs for streaming runs
|
||||
recipe_sexp = None
|
||||
recipe_name = None
|
||||
pending = await self.db.get_pending_run(run_id)
|
||||
if pending:
|
||||
recipe_sexp = pending.get("dag_json")
|
||||
|
||||
# Extract recipe name from streaming recipe content
|
||||
if recipe_sexp:
|
||||
import re
|
||||
name_match = re.search(r'\(stream\s+"([^"]+)"', recipe_sexp)
|
||||
if name_match:
|
||||
recipe_name = name_match.group(1)
|
||||
|
||||
return {
|
||||
"run_id": run_id,
|
||||
"status": "completed",
|
||||
"recipe": cached.get("recipe"),
|
||||
"recipe_name": recipe_name,
|
||||
"inputs": self._ensure_inputs_list(cached.get("inputs")),
|
||||
"output_cid": output_cid,
|
||||
"ipfs_cid": cached.get("ipfs_cid"),
|
||||
@@ -140,6 +155,7 @@ class RunService:
|
||||
"actor_id": cached.get("actor_id"),
|
||||
"created_at": cached.get("created_at"),
|
||||
"completed_at": cached.get("created_at"),
|
||||
"recipe_sexp": recipe_sexp,
|
||||
}
|
||||
|
||||
# Check database for pending run
|
||||
@@ -175,6 +191,7 @@ class RunService:
|
||||
"output_name": pending.get("output_name"),
|
||||
"created_at": pending.get("created_at"),
|
||||
"error": pending.get("error"),
|
||||
"recipe_sexp": pending.get("dag_json"), # Recipe content for streaming runs
|
||||
}
|
||||
|
||||
# If task completed, get result
|
||||
@@ -209,6 +226,7 @@ class RunService:
|
||||
"actor_id": pending.get("actor_id"),
|
||||
"created_at": pending.get("created_at"),
|
||||
"error": pending.get("error"),
|
||||
"recipe_sexp": pending.get("dag_json"), # Recipe content for streaming runs
|
||||
}
|
||||
|
||||
# Fallback: Check Redis for backwards compatibility
|
||||
@@ -714,12 +732,21 @@ class RunService:
|
||||
"""Get execution plan for a run.
|
||||
|
||||
Plans are just node outputs - cached by content hash like everything else.
|
||||
For streaming runs, returns the recipe content as the plan.
|
||||
"""
|
||||
# Get run to find plan_cache_id
|
||||
run = await self.get_run(run_id)
|
||||
if not run:
|
||||
return None
|
||||
|
||||
# For streaming runs, return the recipe as the plan
|
||||
if run.get("recipe") == "streaming" and run.get("recipe_sexp"):
|
||||
return {
|
||||
"steps": [{"id": "stream", "type": "STREAM", "name": "Streaming Recipe"}],
|
||||
"sexp": run.get("recipe_sexp"),
|
||||
"format": "sexp",
|
||||
}
|
||||
|
||||
# Check plan_cid (stored in database) or plan_cache_id (legacy)
|
||||
plan_cid = run.get("plan_cid") or run.get("plan_cache_id")
|
||||
if plan_cid:
|
||||
|
||||
Reference in New Issue
Block a user