Add Reprocess to retry failed IC reviews without the create wizard.
Successful reviews are kept; skipped ICs such as a DeepSeek 400 are run again. Reprocess all re-reviews every IC while still using the library cache. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -16,6 +16,7 @@ from fastapi import APIRouter, HTTPException, Request
|
||||
from sse_starlette.sse import EventSourceResponse
|
||||
|
||||
from pydantic import BaseModel
|
||||
from typing import Literal
|
||||
|
||||
from backend.routers.deps import get_storage, resolve_or_404
|
||||
from backend.services import event_bridge, job_runner
|
||||
@@ -24,12 +25,22 @@ from backend.services import projects as proj_svc
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
VALID_REGEN_STAGES = {"derating"}
|
||||
_REPROCESS_OK_FROM = frozenset({
|
||||
proj_svc.STATUS_COMPLETE,
|
||||
proj_svc.STATUS_ERROR,
|
||||
proj_svc.STATUS_CANCELLED,
|
||||
})
|
||||
|
||||
|
||||
class RegenRequest(BaseModel):
|
||||
stages: list[str]
|
||||
|
||||
|
||||
class ReprocessRequest(BaseModel):
|
||||
"""``failed`` retries skipped / errored IC reviews; ``all`` re-reviews every IC."""
|
||||
mode: Literal["failed", "all"] = "failed"
|
||||
|
||||
|
||||
router = APIRouter(tags=["pipeline"])
|
||||
|
||||
|
||||
@@ -179,6 +190,75 @@ async def resume(project_id: str, request: Request):
|
||||
return {"status": "resumed", "project_id": project_id}
|
||||
|
||||
|
||||
@router.post("/pipeline/{project_id}/reprocess", status_code=202)
|
||||
async def reprocess(project_id: str, request: Request, req: ReprocessRequest | None = None):
|
||||
"""Re-run a finished project without the create wizard.
|
||||
|
||||
Keeps BOM, netlist, datasheets, and library extractions. ``failed``
|
||||
(default) skips ICs that already produced a review; ``all`` re-reviews
|
||||
every IC.
|
||||
"""
|
||||
from backend._version import PINSCOPE_VERSION
|
||||
|
||||
storage = get_storage(request)
|
||||
owner_user_id, meta = await resolve_or_404(request, project_id)
|
||||
if not meta.has_bom or not meta.has_netlist:
|
||||
raise HTTPException(400, "Upload BOM and netlist before reprocessing")
|
||||
if meta.status not in _REPROCESS_OK_FROM:
|
||||
raise HTTPException(
|
||||
409,
|
||||
f"Reprocess is for complete / error / cancelled runs "
|
||||
f"(status={meta.status}). Pause uses resume; a live run uses cancel.",
|
||||
)
|
||||
|
||||
body = req or ReprocessRequest()
|
||||
retry_failed = body.mode == "failed"
|
||||
keep_refs = (
|
||||
proj_svc.completed_review_refs_for_retry(storage, owner_user_id, project_id)
|
||||
if retry_failed else []
|
||||
)
|
||||
|
||||
try:
|
||||
proj_svc.transition_status(
|
||||
storage, owner_user_id, project_id,
|
||||
from_status=_REPROCESS_OK_FROM,
|
||||
to_status=proj_svc.STATUS_QUEUED,
|
||||
cancel_requested=False,
|
||||
execution_name=None,
|
||||
pipeline_state=None,
|
||||
pause_checkpoint=None,
|
||||
pause_reason=None,
|
||||
completed_review_refs=keep_refs,
|
||||
pinscope_version=PINSCOPE_VERSION,
|
||||
)
|
||||
except proj_svc.StatusConflict:
|
||||
raise HTTPException(409, "Pipeline already running or queued")
|
||||
|
||||
try:
|
||||
execution_name = job_runner.enqueue_pipeline(
|
||||
project_id, owner_user_id, resume=retry_failed, free=False,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("enqueue_pipeline (reprocess) failed for %s", project_id)
|
||||
proj_svc.update_project(
|
||||
storage, owner_user_id, project_id,
|
||||
status=proj_svc.STATUS_ERROR,
|
||||
pipeline_state={"error": "Failed to enqueue worker"},
|
||||
)
|
||||
raise HTTPException(503, "Failed to enqueue pipeline worker; please retry")
|
||||
|
||||
proj_svc.update_project(
|
||||
storage, owner_user_id, project_id, execution_name=execution_name,
|
||||
)
|
||||
return {
|
||||
"status": "reprocess_started",
|
||||
"project_id": project_id,
|
||||
"mode": body.mode,
|
||||
"resume": retry_failed,
|
||||
"kept_review_refs": keep_refs,
|
||||
}
|
||||
|
||||
|
||||
@router.post("/pipeline/{project_id}/restart", status_code=202)
|
||||
async def restart(project_id: str, request: Request):
|
||||
"""Admin-only: wipe per-project extractions and run the pipeline free."""
|
||||
|
||||
@@ -1644,15 +1644,17 @@ async def run_pipeline(
|
||||
the project is left in ``paused_insufficient_credits`` with a
|
||||
checkpoint so it can be resumed later.
|
||||
|
||||
When ``resume=True``, prior completed review refs and spent credits are
|
||||
restored from the project's ``pause_checkpoint`` so completed work is
|
||||
skipped on the next pass.
|
||||
When ``resume=True``, prior completed review refs are restored so
|
||||
already-reviewed ICs are skipped (paused credit resume, or user
|
||||
reprocess of failed reviews).
|
||||
|
||||
When ``free=True`` (admin-initiated rerun), every call runs through
|
||||
``ApiLogger(free=True)`` so ``credits_charged`` is zeroed, the credit
|
||||
gate is bypassed, and ``meta.total_cost_usd`` is preserved rather than
|
||||
incremented. The raw Anthropic cost is still captured in log entries.
|
||||
incremented. The raw Anthropic cost is still captured in log entries.
|
||||
"""
|
||||
ctx: PipelineContext | None = None
|
||||
api_logger: ApiLogger | None = None
|
||||
try:
|
||||
meta = proj_svc.get_project(storage, user_id, project_id)
|
||||
if not meta:
|
||||
@@ -1826,37 +1828,59 @@ async def run_pipeline(
|
||||
# asyncio.CancelledError can also arrive during local-dev
|
||||
# subprocess shutdown (SIGTERM). Both are handled the same way.
|
||||
try:
|
||||
extra: dict = {
|
||||
"pipeline_state": {"error": "Pipeline cancelled by user"},
|
||||
"cancel_requested": False,
|
||||
}
|
||||
if ctx is not None:
|
||||
extra["completed_review_refs"] = sorted(
|
||||
ctx.completed_review_refs, key=natural_sort_key,
|
||||
)
|
||||
extra["skipped_components"] = (
|
||||
[s.to_dict() for s in ctx.skipped] or None
|
||||
)
|
||||
proj_svc.transition_status(
|
||||
storage, user_id, project_id,
|
||||
from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED},
|
||||
to_status=proj_svc.STATUS_CANCELLED,
|
||||
pipeline_state={"error": "Pipeline cancelled by user"},
|
||||
cancel_requested=False,
|
||||
**extra,
|
||||
)
|
||||
except proj_svc.StatusConflict:
|
||||
pass
|
||||
broker.publish(project_id, "pipeline_cancelled", {"error": "Pipeline cancelled by user"})
|
||||
# Last-mile flush so partial billing is captured.
|
||||
try:
|
||||
api_logger.flush(storage, user_id, project_id) # type: ignore[has-type]
|
||||
if api_logger is not None:
|
||||
api_logger.flush(storage, user_id, project_id)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
logger.exception("Pipeline run crashed for project %s", project_id)
|
||||
try:
|
||||
extra = {
|
||||
"pipeline_state": {"error": str(e)},
|
||||
"cancel_requested": False,
|
||||
}
|
||||
if ctx is not None:
|
||||
extra["completed_review_refs"] = sorted(
|
||||
ctx.completed_review_refs, key=natural_sort_key,
|
||||
)
|
||||
extra["skipped_components"] = (
|
||||
[s.to_dict() for s in ctx.skipped] or None
|
||||
)
|
||||
proj_svc.transition_status(
|
||||
storage, user_id, project_id,
|
||||
from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED},
|
||||
to_status=proj_svc.STATUS_ERROR,
|
||||
pipeline_state={"error": str(e)},
|
||||
cancel_requested=False,
|
||||
**extra,
|
||||
)
|
||||
except proj_svc.StatusConflict:
|
||||
pass
|
||||
broker.publish(project_id, "pipeline_error", {"error": str(e)})
|
||||
try:
|
||||
api_logger.flush(storage, user_id, project_id) # type: ignore[has-type]
|
||||
if api_logger is not None:
|
||||
api_logger.flush(storage, user_id, project_id)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@@ -129,6 +129,37 @@ class ProjectMeta(BaseModel):
|
||||
cancel_requested: bool = False
|
||||
|
||||
|
||||
def completed_review_refs_for_retry(
|
||||
storage: StorageBackend, user_id: str, project_id: str,
|
||||
) -> list[str]:
|
||||
"""ICs that already finished review and should be skipped on reprocess.
|
||||
|
||||
Drops refs that failed (skipped_components / report.review_errors) so
|
||||
those ICs are tried again.
|
||||
"""
|
||||
meta = get_project(storage, user_id, project_id)
|
||||
if not meta:
|
||||
return []
|
||||
failed: set[str] = set()
|
||||
for item in meta.skipped_components or []:
|
||||
stage = (item.get("stage") or "")
|
||||
ident = (item.get("identifier") or "").strip()
|
||||
if ident and stage in ("validation", "review"):
|
||||
failed.add(ident)
|
||||
report_key = f"{_project_prefix(user_id, project_id)}/report.json"
|
||||
if storage.exists(report_key):
|
||||
try:
|
||||
report = storage.read_json(report_key)
|
||||
except Exception:
|
||||
report = {}
|
||||
for ref in (report.get("review_errors") or {}):
|
||||
if ref:
|
||||
failed.add(str(ref))
|
||||
from backend.pinscopex.utils import natural_sort_key
|
||||
kept = [r for r in (meta.completed_review_refs or []) if r and r not in failed]
|
||||
return sorted(kept, key=natural_sort_key)
|
||||
|
||||
|
||||
def _project_prefix(user_id: str, project_id: str) -> str:
|
||||
return f"users/{user_id}/projects/{project_id}"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user