Heal zombie running projects stuck on Review forever.

If the event log already ends with pipeline_complete, flip meta to complete on project/status fetch so Progress no longer spins on a dead worker.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-09-13 16:45:07 +02:00
co-authored by Cursor
parent 4882c9206a
commit e124ac2de7
7 changed files with 112 additions and 20 deletions
+14 -4
View File
@@ -497,13 +497,23 @@ async def list_running_pipelines(request: Request):
if meta.status not in (proj_svc.STATUS_QUEUED, proj_svc.STATUS_RUNNING):
continue
# Sweeper: if the execution is in a terminal Cloud Run state,
# the worker is already gone. Flip status → error so the UI
# stops lying. Skip the sweep when execution_name is missing
# (worker may still be enqueueing).
# Sweeper: if the execution is in a terminal Cloud Run / local
# state, the worker is already gone. Flip status → error so the
# UI stops lying. Also heal projects whose event log already
# ends with pipeline_complete (finished, meta never flipped).
healed = proj_svc.heal_if_pipeline_finished(storage, uid, meta.id)
if healed is not None:
continue
exec_state = "unknown"
if meta.execution_name:
exec_state = job_runner.get_execution_state(meta.execution_name)
elif not job_runner.use_cloud_run_jobs():
# Local zombie: no execution_name but a dead pid file, or
# no live proc — treat as failed after the stale window.
exec_state = job_runner.get_execution_state(
f"local/projects/{meta.id}"
)
if exec_state in ("succeeded", "failed", "cancelled"):
# Allow a short grace period so we don't race the worker
# writing its own terminal status. updated may be stale
+11 -2
View File
@@ -509,8 +509,16 @@ async def events(project_id: str, request: Request):
@router.get("/pipeline/{project_id}/status")
async def status(project_id: str, request: Request):
"""Polling fallback — returns current project state."""
_, meta = await resolve_or_404(request, project_id)
"""Polling fallback — returns current project state.
Also heals zombie ``running``/``queued`` projects whose event log
already ends with ``pipeline_complete`` (worker died after finishing).
"""
storage = get_storage(request)
owner_user_id, meta = await resolve_or_404(request, project_id)
healed = proj_svc.heal_if_pipeline_finished(storage, owner_user_id, project_id)
if healed is not None:
meta = healed
return {
"status": meta.status,
"summary": meta.summary,
@@ -519,6 +527,7 @@ async def status(project_id: str, request: Request):
"placement_status": meta.placement_status,
"placement_state": meta.placement_state,
"placement_running": (meta.placement_status or "draft") in ("queued", "running"),
"healed": healed is not None,
}
+4 -2
View File
@@ -140,8 +140,10 @@ async def list_projects(request: Request):
@router.get("/projects/{project_id}")
async def get_project(project_id: str, request: Request):
_, meta = await resolve_or_404(request, project_id)
return meta.model_dump()
storage = get_storage(request)
owner_user_id, meta = await resolve_or_404(request, project_id)
healed = proj_svc.heal_if_pipeline_finished(storage, owner_user_id, project_id)
return (healed or meta).model_dump()
@router.delete("/projects/{project_id}")
+45
View File
@@ -281,6 +281,51 @@ def mark_stale_running(
return None
def heal_if_pipeline_finished(
storage: StorageBackend, user_id: str, project_id: str,
) -> ProjectMeta | None:
"""If meta says queued/running but events already ended with
``pipeline_complete``, flip status to ``complete``.
Covers zombies where the worker wrote the terminal event (and often
the report) then died before the meta transition — e.g. container
rebuild mid-shutdown. Returns updated meta, or ``None`` if no heal.
"""
meta = get_project(storage, user_id, project_id)
if meta is None or meta.status not in (STATUS_RUNNING, STATUS_QUEUED):
return None
events_prefix = f"{_project_prefix(user_id, project_id)}/events/"
try:
keys = storage.list_prefix(events_prefix)
except Exception:
return None
event_keys = sorted(
k for k in keys if k.endswith(".json") and "/events/" in k
)
if not event_keys:
return None
try:
last = storage.read_json(event_keys[-1])
except Exception:
return None
if (last or {}).get("event") != "pipeline_complete":
return None
summary = (last.get("data") or {}).get("summary")
try:
return transition_status(
storage, user_id, project_id,
from_status={STATUS_RUNNING, STATUS_QUEUED},
to_status=STATUS_COMPLETE,
summary=summary if isinstance(summary, dict) else meta.summary,
cancel_requested=False,
pipeline_state=None,
)
except StatusConflict:
return None
# --- CRUD ---