mirror of
https://github.com/hansjone/netx.git
synced 2026-10-10 19:20:46 +08:00
Add compare batch cancel and restart recovery so stuck runs unblock.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
9adf9f2764
commit
958942c5b6
8 changed files with 532 additions and 50 deletions
|
|
@ -159,6 +159,17 @@ def run_api_startup() -> None:
|
|||
)
|
||||
except Exception:
|
||||
_log.exception("startup: biz_state collect recovery failed")
|
||||
try:
|
||||
from .biz_state.compare_service import recover_interrupted_compares_on_startup
|
||||
|
||||
cmp_rec = recover_interrupted_compares_on_startup(db)
|
||||
if cmp_rec.get("runs"):
|
||||
_log.info(
|
||||
"startup: cancelled %s interrupted biz compare run(s)",
|
||||
cmp_rec.get("runs"),
|
||||
)
|
||||
except Exception:
|
||||
_log.exception("startup: biz compare recovery failed")
|
||||
try:
|
||||
from .port_traffic_migrate import backfill_port_traffic_series
|
||||
|
||||
|
|
|
|||
|
|
@ -1868,13 +1868,17 @@ def _set_run_progress(
|
|||
progress.update(extra)
|
||||
prev["progress"] = progress
|
||||
run.summary_json = prev
|
||||
run.status = "running"
|
||||
title = progress["sheet_title"] or progress["sheet_id"] or ""
|
||||
run.message = (
|
||||
f"{phase} {sheet_index}/{sheet_total}"
|
||||
+ (f" · {title}" if title else "")
|
||||
+ f" · {elapsed_ms // 1000}s"
|
||||
)[:1024]
|
||||
# Re-read status from DB — cancel may have been committed by another session
|
||||
# (UI cancel / startup recovery) while this worker still holds a stale "running".
|
||||
db.expire(run, ["status", "message"])
|
||||
if str(run.status or "") != "cancelled":
|
||||
run.status = "running"
|
||||
title = progress["sheet_title"] or progress["sheet_id"] or ""
|
||||
run.message = (
|
||||
f"{phase} {sheet_index}/{sheet_total}"
|
||||
+ (f" · {title}" if title else "")
|
||||
+ f" · {elapsed_ms // 1000}s"
|
||||
)[:1024]
|
||||
db.commit()
|
||||
|
||||
|
||||
|
|
@ -1937,10 +1941,18 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]:
|
|||
prev["sheets"] = list(sheet_metas)
|
||||
prev.update({k: agg[k] for k in agg})
|
||||
run.summary_json = prev
|
||||
db.expire(run, ["status"])
|
||||
# Never resurrect cancelled while publishing incremental sheet metas
|
||||
if str(run.status or "") == "cancelled":
|
||||
db.commit()
|
||||
return
|
||||
db.commit()
|
||||
|
||||
try:
|
||||
for idx, sheet in enumerate(sheets_cfg, start=1):
|
||||
if _run_is_cancelled(db, run_id):
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
return get_run(db, run.id)
|
||||
sid = sheet_key(sheet)
|
||||
# Mark current sheet running in the sidebar list
|
||||
for meta in sheet_metas:
|
||||
|
|
@ -2091,6 +2103,13 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]:
|
|||
fail_diffs.clear()
|
||||
ok_diffs.clear()
|
||||
db.commit()
|
||||
if _run_is_cancelled(db, run_id):
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
return get_run(db, run.id)
|
||||
|
||||
if _run_is_cancelled(db, run_id):
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
return get_run(db, run.id)
|
||||
|
||||
duration_ms = int((time.monotonic() - started_mono) * 1000)
|
||||
top_fields = sorted(
|
||||
|
|
@ -2117,6 +2136,9 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]:
|
|||
},
|
||||
"sheets": sheet_metas,
|
||||
}
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
if str(run.status or "") == "cancelled":
|
||||
return get_run(db, run.id)
|
||||
run.status = "success"
|
||||
run.summary_json = summary_payload
|
||||
run.mapping_stats_json = mapping_by_metric
|
||||
|
|
@ -2129,34 +2151,123 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]:
|
|||
db.commit()
|
||||
return get_run(db, run.id)
|
||||
except HTTPException as exc:
|
||||
run.status = "failed"
|
||||
run.message = str(getattr(exc, "detail", "") or exc)[:1024]
|
||||
prev = dict(run.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "failed"
|
||||
prog["elapsed_ms"] = int((time.monotonic() - started_mono) * 1000)
|
||||
prev["progress"] = prog
|
||||
prev["duration_ms"] = prog["elapsed_ms"]
|
||||
run.summary_json = prev
|
||||
db.commit()
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
if str(run.status or "") != "cancelled":
|
||||
run.status = "failed"
|
||||
run.message = str(getattr(exc, "detail", "") or exc)[:1024]
|
||||
prev = dict(run.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "failed"
|
||||
prog["elapsed_ms"] = int((time.monotonic() - started_mono) * 1000)
|
||||
prev["progress"] = prog
|
||||
prev["duration_ms"] = prog["elapsed_ms"]
|
||||
run.summary_json = prev
|
||||
db.commit()
|
||||
raise
|
||||
except Exception as exc:
|
||||
_log.exception("compare run failed run=%s job=%s", run_id, j.id)
|
||||
run.status = "failed"
|
||||
run.message = str(exc)[:1024]
|
||||
prev = dict(run.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "failed"
|
||||
prog["elapsed_ms"] = int((time.monotonic() - started_mono) * 1000)
|
||||
prev["progress"] = prog
|
||||
prev["duration_ms"] = prog["elapsed_ms"]
|
||||
run.summary_json = prev
|
||||
db.commit()
|
||||
run = db.get(BizCompareRun, run_id) or run
|
||||
if str(run.status or "") != "cancelled":
|
||||
run.status = "failed"
|
||||
run.message = str(exc)[:1024]
|
||||
prev = dict(run.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "failed"
|
||||
prog["elapsed_ms"] = int((time.monotonic() - started_mono) * 1000)
|
||||
prev["progress"] = prog
|
||||
prev["duration_ms"] = prog["elapsed_ms"]
|
||||
run.summary_json = prev
|
||||
db.commit()
|
||||
raise
|
||||
|
||||
|
||||
_INTERRUPT_MARK = "interrupted_by_restart"
|
||||
|
||||
|
||||
def _job_has_active_run(db: Session, job_id: str) -> BizCompareRun | None:
|
||||
return (
|
||||
db.query(BizCompareRun)
|
||||
.filter(
|
||||
BizCompareRun.job_id == str(job_id or ""),
|
||||
BizCompareRun.status.in_(("running", "queued")),
|
||||
)
|
||||
.order_by(BizCompareRun.created_at.desc())
|
||||
.first()
|
||||
)
|
||||
|
||||
|
||||
def _run_is_cancelled(db: Session, run_id: str) -> bool:
|
||||
"""Re-read status so user/startup cancel is visible to the worker thread."""
|
||||
db.expire_all()
|
||||
r = db.get(BizCompareRun, run_id)
|
||||
return bool(r and str(r.status or "") == "cancelled")
|
||||
|
||||
|
||||
def recover_interrupted_compares_on_startup(db: Session) -> dict[str, Any]:
|
||||
"""Mark orphaned running/queued compare runs as cancelled after process restart.
|
||||
|
||||
In-memory job locks die with the process; without this, the UI stays on
|
||||
「比对中」and blocks a new run.
|
||||
"""
|
||||
now = _utcnow()
|
||||
rows = (
|
||||
db.query(BizCompareRun)
|
||||
.filter(BizCompareRun.status.in_(("running", "queued")))
|
||||
.all()
|
||||
)
|
||||
n = 0
|
||||
for r in rows:
|
||||
r.status = "cancelled"
|
||||
msg = str(r.message or "").strip()
|
||||
if _INTERRUPT_MARK not in msg:
|
||||
r.message = f"{msg} | {_INTERRUPT_MARK}".strip(" |")[:1024]
|
||||
prev = dict(r.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "cancelled"
|
||||
prog["elapsed_ms"] = int(prog.get("elapsed_ms") or 0)
|
||||
prev["progress"] = prog
|
||||
# Mark in-flight sheet placeholders so UI does not show fake pass
|
||||
sheets = list(prev.get("sheets") or [])
|
||||
for sh in sheets:
|
||||
st = str(sh.get("status") or "")
|
||||
if st in ("pending", "running", "queued"):
|
||||
sh["status"] = "cancelled"
|
||||
prev["sheets"] = sheets
|
||||
r.summary_json = prev
|
||||
n += 1
|
||||
if n:
|
||||
db.commit()
|
||||
_log.info("startup: cancelled %s interrupted compare run(s)", n)
|
||||
return {"runs": n, "at": now.isoformat() + "Z"}
|
||||
|
||||
|
||||
def cancel_compare_run(db: Session, run_id: str) -> dict[str, Any]:
|
||||
"""Cancel a running/queued compare so a new run can start."""
|
||||
r = db.get(BizCompareRun, run_id)
|
||||
if not r:
|
||||
raise HTTPException(status_code=404, detail="run_not_found")
|
||||
st = str(r.status or "")
|
||||
if st not in ("running", "queued"):
|
||||
return get_run(db, run_id)
|
||||
r.status = "cancelled"
|
||||
msg = str(r.message or "").strip()
|
||||
r.message = f"{msg} | cancelled_by_user".strip(" |")[:1024]
|
||||
prev = dict(r.summary_json or {})
|
||||
prog = dict(prev.get("progress") or {})
|
||||
prog["phase"] = "cancelled"
|
||||
prev["progress"] = prog
|
||||
for sh in list(prev.get("sheets") or []):
|
||||
if str(sh.get("status") or "") in ("pending", "running", "queued"):
|
||||
sh["status"] = "cancelled"
|
||||
r.summary_json = prev
|
||||
db.commit()
|
||||
return get_run(db, run_id)
|
||||
|
||||
|
||||
def run_compare(db: Session, job_id: str, *, force_after_batch_id: str = "") -> dict[str, Any]:
|
||||
"""Synchronous compare (auto-compare / tests). Blocks the caller until done."""
|
||||
if _job_has_active_run(db, job_id):
|
||||
raise HTTPException(status_code=409, detail="compare_already_running")
|
||||
lock = _job_compare_lock(job_id)
|
||||
if not lock.acquire(blocking=False):
|
||||
raise HTTPException(status_code=409, detail="compare_already_running")
|
||||
|
|
@ -2183,6 +2294,8 @@ def enqueue_compare(
|
|||
Returns immediately so the HTTP worker / UI stay responsive. Poll
|
||||
``GET /compare/runs/{id}`` for progress (``summary.progress``).
|
||||
"""
|
||||
if _job_has_active_run(db, job_id):
|
||||
raise HTTPException(status_code=409, detail="compare_already_running")
|
||||
lock = _job_compare_lock(job_id)
|
||||
if not lock.acquire(blocking=False):
|
||||
raise HTTPException(status_code=409, detail="compare_already_running")
|
||||
|
|
|
|||
|
|
@ -746,6 +746,12 @@ def api_get_run(run_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
|||
return cmp_svc.get_run(db, run_id)
|
||||
|
||||
|
||||
@router.post("/compare/runs/{run_id}/cancel")
|
||||
def api_cancel_run(run_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
"""Cancel a stuck/running compare so a new run can start."""
|
||||
return cmp_svc.cancel_compare_run(db, run_id)
|
||||
|
||||
|
||||
@router.delete("/compare/runs/{run_id}")
|
||||
def api_delete_run(run_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return cmp_svc.delete_run(db, run_id)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue