diff --git a/netx_api/collection_job_state.py b/netx_api/collection_job_state.py index ec24223..d550eee 100644 --- a/netx_api/collection_job_state.py +++ b/netx_api/collection_job_state.py @@ -12,7 +12,18 @@ _TERMINAL = frozenset({"success", "fail", "cancelled"}) def _sync_job_counts(job: NeCollectionJob, runs: list[NeCollectionRun]) -> None: job.success_count = sum(1 for r in runs if str(r.status) == "success") - job.fail_count = sum(1 for r in runs if str(r.status) in ("fail", "cancelled")) + job.fail_count = sum(1 for r in runs if str(r.status) == "fail") + # cancelled rows are tracked separately via ne_count - success - fail - active + + +def sync_job_progress(db: Session, job_id: str) -> None: + """Refresh success/fail counters while a job is still running.""" + job = db.get(NeCollectionJob, job_id) + if not job or str(job.status or "") != "running": + return + runs = db.query(NeCollectionRun).filter(NeCollectionRun.job_id == job_id).all() + _sync_job_counts(job, runs) + db.commit() def finalize_collection_job(db: Session, job_id: str) -> None: @@ -56,6 +67,9 @@ def reconcile_stale_collection_job(db: Session, job_id: str) -> bool: if st in _TERMINAL: continue if st == "pending": + # Pending while job is running means queued in the worker pool, not stuck. + if str(job.status or "") == "running": + continue anchor = job.started_at or job.created_at limit = pending_stale_sec reason = "collection_pending_stale" diff --git a/netx_api/collection_router.py b/netx_api/collection_router.py index 3aa1a9d..2c297f8 100644 --- a/netx_api/collection_router.py +++ b/netx_api/collection_router.py @@ -15,6 +15,7 @@ from .collection_service import ( pause_collection_job, resolve_run_output_path, restart_collection_job, + retry_failed_collection_job, ) from .collection_schemas import CollectionJobCreate from .db import get_db @@ -95,6 +96,11 @@ def api_restart_collection(job_id: str, db: Session = Depends(get_db)): return restart_collection_job(db, job_id).model_dump() +@router.post("/{job_id}/retry-failed") +def api_retry_failed_collection(job_id: str, db: Session = Depends(get_db)): + return retry_failed_collection_job(db, job_id).model_dump() + + @router.delete("/{job_id}") def api_delete_collection(job_id: str, db: Session = Depends(get_db)): return delete_collection_job(db, job_id) diff --git a/netx_api/collection_service.py b/netx_api/collection_service.py index 419f0b1..ec46a53 100644 --- a/netx_api/collection_service.py +++ b/netx_api/collection_service.py @@ -14,7 +14,12 @@ from sqlalchemy import func, or_ from sqlalchemy.orm import Session from .models import ManagedNE, NeCollectionJob, NeCollectionRun -from .collection_job_state import finalize_collection_job, reconcile_stale_collection_job, _sync_job_counts +from .collection_job_state import ( + finalize_collection_job, + reconcile_stale_collection_job, + sync_job_progress, + _sync_job_counts, +) from .collection_schemas import CollectionJobCreate, CollectionJobOut, CollectionRunOut from .ne_collect_runner import schedule_collection_runs from .ne_collection_paths import clear_run_output_files, collection_data_root @@ -23,7 +28,7 @@ _log = logging.getLogger("netx.collection") def _now() -> datetime: - return datetime.utcnow() + return datetime.now() def _parse_commands(text: str) -> list[str]: @@ -93,7 +98,7 @@ def run_to_out(row: NeCollectionRun) -> CollectionRunOut: ne_name=str(row.ne_name or ""), ne_ip=str(row.ne_ip or ""), status=str(row.status or "pending"), - message=str(row.message or "")[:1000], + message=str(row.message or ""), output_rel_path=rel, has_output=bool(rel), started_at=row.started_at, @@ -188,9 +193,13 @@ def list_collection_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> total = int(stmt.count()) rows = stmt.order_by(NeCollectionJob.created_at.desc()).offset((page - 1) * page_size).limit(page_size).all() for row in rows: - if str(row.status or "") not in ("done", "failed"): - reconcile_stale_collection_job(db, str(row.id)) + job_id = str(row.id) + if str(row.status or "") not in ("done", "failed", "paused"): + reconcile_stale_collection_job(db, job_id) db.refresh(row) + if str(row.status or "") == "running": + sync_job_progress(db, job_id) + db.refresh(row) job_ids = [str(x.id) for x in rows] output_counts = _output_counts_for_jobs(db, job_ids) return { @@ -205,9 +214,12 @@ def get_collection_job(db: Session, job_id: str) -> dict[str, Any]: job = db.get(NeCollectionJob, job_id) if not job: raise HTTPException(status_code=404, detail="collection_job_not_found") - if str(job.status or "") not in ("done", "failed"): + if str(job.status or "") not in ("done", "failed", "paused"): reconcile_stale_collection_job(db, job_id) db.refresh(job) + if str(job.status or "") == "running": + sync_job_progress(db, job_id) + db.refresh(job) return { "job": job_to_out(job, output_count=_output_count_for_job(db, job_id)).model_dump(), } @@ -277,6 +289,57 @@ def pause_collection_job(db: Session, job_id: str) -> CollectionJobOut: return job_to_out(job, output_count=_output_count_for_job(db, job_id)) +def _reset_runs_for_retry( + db: Session, + job_id: str, + runs: list[NeCollectionRun], + *, + only_statuses: frozenset[str], +) -> list[str]: + retry_ids: list[str] = [] + for run in runs: + st = str(run.status or "") + if st in ("pending", "running") or st not in only_statuses: + continue + clear_run_output_files(job_id, str(run.id)) + run.status = "pending" + run.message = "" + run.output_rel_path = "" + run.started_at = None + run.ended_at = None + retry_ids.append(str(run.id)) + return retry_ids + + +def _start_job_retry( + db: Session, + job: NeCollectionJob, + job_id: str, + retry_ids: list[str], + commands: list[str], + *, + reset_all_counts: bool, +) -> CollectionJobOut: + if not retry_ids: + raise HTTPException(status_code=400, detail="collection_nothing_to_retry") + now = _now() + job.status = "running" + job.ended_at = None + job.error_message = "" + if reset_all_counts: + job.success_count = 0 + job.fail_count = 0 + else: + runs = db.query(NeCollectionRun).filter(NeCollectionRun.job_id == job_id).all() + _sync_job_counts(job, runs) + job.started_at = now + job.last_run_at = now + db.commit() + db.refresh(job) + schedule_collection_runs(job_id, retry_ids, commands) + return job_to_out(job, output_count=_output_count_for_job(db, job_id)) + + def restart_collection_job(db: Session, job_id: str) -> CollectionJobOut: job = db.get(NeCollectionJob, job_id) if not job: @@ -289,32 +352,29 @@ def restart_collection_job(db: Session, job_id: str) -> CollectionJobOut: commands = _parse_commands(str(job.commands or "")) if not commands: raise HTTPException(status_code=400, detail="commands_empty") - retry_ids: list[str] = [] - for run in runs: - st = str(run.status or "") - if st in ("pending", "running"): - continue - clear_run_output_files(job_id, str(run.id)) - run.status = "pending" - run.message = "" - run.output_rel_path = "" - run.started_at = None - run.ended_at = None - retry_ids.append(str(run.id)) - if not retry_ids: - raise HTTPException(status_code=400, detail="collection_nothing_to_retry") - now = _now() - job.status = "running" - job.ended_at = None - job.error_message = "" - job.success_count = 0 - job.fail_count = 0 - job.started_at = now - job.last_run_at = now - db.commit() - db.refresh(job) - schedule_collection_runs(job_id, retry_ids, commands) - return job_to_out(job, output_count=0) + retry_ids = _reset_runs_for_retry( + db, + job_id, + runs, + only_statuses=frozenset({"success", "fail", "cancelled"}), + ) + return _start_job_retry(db, job, job_id, retry_ids, commands, reset_all_counts=True) + + +def retry_failed_collection_job(db: Session, job_id: str) -> CollectionJobOut: + job = db.get(NeCollectionJob, job_id) + if not job: + raise HTTPException(status_code=404, detail="collection_job_not_found") + runs = db.query(NeCollectionRun).filter(NeCollectionRun.job_id == job_id).all() + if not runs: + raise HTTPException(status_code=400, detail="collection_no_runs") + if str(job.status or "") == "running" or _active_runs(runs): + raise HTTPException(status_code=400, detail="collection_job_running") + commands = _parse_commands(str(job.commands or "")) + if not commands: + raise HTTPException(status_code=400, detail="commands_empty") + retry_ids = _reset_runs_for_retry(db, job_id, runs, only_statuses=frozenset({"fail"})) + return _start_job_retry(db, job, job_id, retry_ids, commands, reset_all_counts=False) def delete_collection_job(db: Session, job_id: str) -> dict[str, bool]: diff --git a/netx_api/ne_collect_runner.py b/netx_api/ne_collect_runner.py index 3809b2e..172f243 100644 --- a/netx_api/ne_collect_runner.py +++ b/netx_api/ne_collect_runner.py @@ -2,6 +2,7 @@ from __future__ import annotations import logging import re +import traceback from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeout from datetime import datetime from pathlib import Path @@ -9,7 +10,7 @@ from typing import Any from netmiko import ConnectHandler -from .collection_job_state import finalize_collection_job +from .collection_job_state import finalize_collection_job, sync_job_progress from .config import settings from .db import SessionLocal from .models import ManagedNE, NeCollectionJob, NeCollectionRun @@ -30,6 +31,13 @@ def _executor_pool() -> ThreadPoolExecutor: return _executor +def _format_run_error(exc: BaseException) -> str: + head = f"{type(exc).__name__}: {exc}" + tb = traceback.format_exc().strip() + text = f"{head}\n{tb}" if tb else head + return text[:1020] + + def _safe_filename_part(text: str) -> str: s = re.sub(r'[<>:"/\\|?*]', "_", str(text or "").strip()) return s[:80] or "device" @@ -152,14 +160,15 @@ def _run_single(job_id: str, run_id: str, commands: list[str]) -> None: ended_at=finished_at, ) except CredentialCryptoError as exc: - _update_run(run_id, status="fail", message=str(exc)[:1000], ended_at=datetime.now()) + _update_run(run_id, status="fail", message=str(exc)[:1020], ended_at=datetime.now()) except Exception as exc: _log.exception("collection failed run=%s", run_id) - _update_run(run_id, status="fail", message=str(exc).split("\n")[0][:1000], ended_at=datetime.now()) + _update_run(run_id, status="fail", message=_format_run_error(exc), ended_at=datetime.now()) finally: db.close() db2 = SessionLocal() try: + sync_job_progress(db2, job_id) finalize_collection_job(db2, job_id) finally: db2.close() diff --git a/tests/test_collection_job_state.py b/tests/test_collection_job_state.py index 6a2d1d6..7d881e6 100644 --- a/tests/test_collection_job_state.py +++ b/tests/test_collection_job_state.py @@ -26,7 +26,7 @@ class CollectionJobFinalizeTests(unittest.TestCase): self.assertEqual(job.status, "paused") self.assertEqual(job.success_count, 1) - self.assertEqual(job.fail_count, 1) + self.assertEqual(job.fail_count, 0) self.assertIsNotNone(job.ended_at) db.commit.assert_called_once() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index a852cab..94d5d10 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -61,6 +61,7 @@ const en = { started: "Collection job started: {{id}}", paused: "Job paused", restarted: "Collection restarted", + retryFailedDone: "Retrying failed devices", deleted: "Job deleted", nothingToRetry: "No failed devices to retry", confirmDelete: "Delete this collection job and its log files?", @@ -80,6 +81,7 @@ const en = { collapse: "Hide", pause: "Pause", restart: "Restart", + retryFailed: "Retry failed", downloadResults: "Download results", delete: "Delete", autoRefresh: "Auto-refresh every 2s while jobs are running", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 3c160b5..f44941a 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -61,6 +61,7 @@ const zh = { started: "采集任务已启动:{{id}}", paused: "任务已暂停", restarted: "已重新开始采集", + retryFailedDone: "已开始重采失败网元", deleted: "任务已删除", nothingToRetry: "没有可重试的失败网元", confirmDelete: "确定删除该采集任务?相关日志文件将一并删除。", @@ -80,6 +81,7 @@ const zh = { collapse: "收起", pause: "暂停", restart: "重新开始", + retryFailed: "重采失败网元", downloadResults: "下载结果", delete: "删除", autoRefresh: "任务进行中,每 2 秒自动刷新", diff --git a/web/src/index.css b/web/src/index.css index 3a98027..913b1f6 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -957,6 +957,21 @@ pre { background: #fff; } +.collect-run-message { + max-width: 420px; + max-height: 160px; + overflow: auto; + white-space: pre-wrap; + word-break: break-word; + font-size: 12px; + line-height: 1.45; + color: #334155; + padding: 6px 8px; + background: #f8fafc; + border: 1px solid #e2e8f0; + border-radius: 6px; +} + @media (max-width: 1200px) { .cards { grid-template-columns: 1fr; diff --git a/web/src/pages/CollectPage.tsx b/web/src/pages/CollectPage.tsx index 8b75d41..0b77d1e 100644 --- a/web/src/pages/CollectPage.tsx +++ b/web/src/pages/CollectPage.tsx @@ -9,6 +9,7 @@ import { fetchNeCollections, pauseCollectionJob, restartCollectionJob, + retryFailedCollectionJob, collectionJobDownloadUrl, collectionRunDownloadUrl, } from "../services/api"; @@ -104,6 +105,16 @@ export function CollectPage() { onError: (err) => showError(String(err)), }); + const retryFailedMutation = useMutation({ + mutationFn: retryFailedCollectionJob, + onSuccess: async (job) => { + showOk(t("collect.retryFailedDone")); + setExpandedJobId(job.id); + await invalidateJobs(job.id); + }, + onError: (err) => showError(String(err)), + }); + const deleteMutation = useMutation({ mutationFn: deleteCollectionJob, onSuccess: async (_, jobId) => { @@ -294,10 +305,16 @@ export function CollectPage() { onToggle={() => setExpandedJobId(expandedJobId === job.id ? "" : job.id)} onPause={() => pauseMutation.mutate(job.id)} onRestart={() => restartMutation.mutate(job.id)} + onRetryFailed={() => retryFailedMutation.mutate(job.id)} onDelete={() => { if (window.confirm(t("collect.confirmDelete"))) deleteMutation.mutate(job.id); }} - actionPending={pauseMutation.isPending || restartMutation.isPending || deleteMutation.isPending} + actionPending={ + pauseMutation.isPending || + restartMutation.isPending || + retryFailedMutation.isPending || + deleteMutation.isPending + } /> ))} @@ -325,6 +342,7 @@ function JobRow({ onToggle, onPause, onRestart, + onRetryFailed, onDelete, actionPending, }: { @@ -334,12 +352,14 @@ function JobRow({ onToggle: () => void; onPause: () => void; onRestart: () => void; + onRetryFailed: () => void; onDelete: () => void; actionPending: boolean; }) { const { t } = useI18n(); const canPause = job.status === "running"; const canRestart = job.status !== "running"; + const canRetryFailed = job.status !== "running" && job.fail_count > 0; const canDelete = job.status !== "running"; const hasResults = (job.output_count ?? 0) > 0; @@ -370,6 +390,11 @@ function JobRow({ {t("collect.jobs.restart")} ) : null} + {canRetryFailed ? ( + + ) : null} ) : null} + {jobStatus !== "running" && failCount > 0 ? ( + + ) : null} {runsQuery.isLoading ?
{t("common.refreshing")}
: null} {!runsQuery.isLoading && runs.length === 0 ?{t("common.empty")}
: null} @@ -498,7 +541,15 @@ function JobRunsPanel({