diff --git a/netx_api/collection_service.py b/netx_api/collection_service.py index 69f263b..97a1311 100644 --- a/netx_api/collection_service.py +++ b/netx_api/collection_service.py @@ -217,6 +217,7 @@ def list_collection_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> db.refresh(row) if str(row.status or "") == "running": sync_job_progress(db, job_id) + finalize_collection_job(db, job_id) db.refresh(row) job_ids = [str(x.id) for x in rows] output_counts = _output_counts_for_jobs(db, job_ids) @@ -237,6 +238,7 @@ def get_collection_job(db: Session, job_id: str) -> dict[str, Any]: db.refresh(job) if str(job.status or "") == "running": sync_job_progress(db, job_id) + finalize_collection_job(db, job_id) db.refresh(job) return { "job": job_to_out(job, output_count=_output_count_for_job(db, job_id)).model_dump(), @@ -282,6 +284,33 @@ def _active_runs(runs: list[NeCollectionRun]) -> bool: return any(str(r.status or "") in ("pending", "running") for r in runs) +def _runs_in_progress(runs: list[NeCollectionRun]) -> bool: + """True when a worker is actively collecting (not merely queued as pending).""" + return any(str(r.status or "") == "running" for r in runs) + + +def _blocking_running_job_ids(db: Session, ne_ids: list[str], *, exclude_job_id: str) -> list[str]: + if not ne_ids: + return [] + other_jobs = ( + db.query(NeCollectionJob) + .filter(NeCollectionJob.status == "running", NeCollectionJob.id != exclude_job_id) + .all() + ) + blocked: list[str] = [] + ne_set = set(ne_ids) + for other in other_jobs: + other_id = str(other.id) + overlap = ( + db.query(NeCollectionRun.ne_id) + .filter(NeCollectionRun.job_id == other_id, NeCollectionRun.ne_id.in_(list(ne_set))) + .first() + ) + if overlap: + blocked.append(other_id) + return blocked + + def pause_collection_job(db: Session, job_id: str) -> CollectionJobOut: job = db.get(NeCollectionJob, job_id) if not job: @@ -370,17 +399,22 @@ def start_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, Co 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): + if str(job.status or "") == "running" or _runs_in_progress(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") + ne_ids = list({str(r.ne_id) for r in runs if str(r.ne_id or "").strip()}) + blocking = _blocking_running_job_ids(db, ne_ids, exclude_job_id=job_id) + if blocking: + raise HTTPException(status_code=409, detail=f"collection_ne_busy: {blocking[0][:12]}") + if str(job.status or "") == "pending": run_ids = [str(r.id) for r in runs if str(r.status or "") == "pending"] if not run_ids: raise HTTPException(status_code=400, detail="collection_nothing_to_start") - return _start_job_retry(db, job, job_id, run_ids, commands, reset_all_counts=False) + return _start_job_retry(db, job, job_id, run_ids, commands, reset_all_counts=True) retry_ids = _reset_runs_for_retry( db, @@ -404,11 +438,15 @@ def retry_failed_collection_job( 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): + if str(job.status or "") == "running" or _runs_in_progress(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") + ne_ids = list({str(r.ne_id) for r in runs if str(r.ne_id or "").strip()}) + blocking = _blocking_running_job_ids(db, ne_ids, exclude_job_id=job_id) + if blocking: + raise HTTPException(status_code=409, detail=f"collection_ne_busy: {blocking[0][:12]}") 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) diff --git a/netx_api/ne_collect_runner.py b/netx_api/ne_collect_runner.py index 4e1761d..882c53d 100644 --- a/netx_api/ne_collect_runner.py +++ b/netx_api/ne_collect_runner.py @@ -129,7 +129,7 @@ def _claim_run(job_id: str, run_id: str) -> bool: if job_status != "running": return False if run_status == "running": - return True + return False if run_status != "pending": return False run.status = "running" @@ -149,7 +149,7 @@ def _run_single(job_id: str, run_id: str, commands: list[str]) -> None: try: run = db.get(NeCollectionRun, run_id) st = str(run.status or "") if run else "" - if st in ("cancelled", "success", "fail"): + if st in ("cancelled", "success", "fail", "running"): sync_job_progress(db, job_id) finalize_collection_job(db, job_id) return diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 35c5ac0..c1050fd 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -60,6 +60,8 @@ const en = { added: "Added", create: "Create job", creating: "Creating…", + expand: "Expand", + collapse: "Collapse", meta: "{{ne}} NE(s) selected · {{cmd}} command(s)", }, eligible: { @@ -85,6 +87,8 @@ const en = { retryFailedDone: "Retrying failed devices", deleted: "Job deleted", nothingToRetry: "No failed devices to retry", + neBusy: "Some NEs are busy in another collection job; wait for it to finish before starting", + jobRunning: "This job is already running", confirmDelete: "Delete this collection job and its log files?", jobs: { title: "Collection jobs", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index ec08d50..db65dd3 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -60,6 +60,8 @@ const zh = { added: "已添加", create: "创建任务", creating: "创建中…", + expand: "展开", + collapse: "收起", meta: "已选 {{ne}} 台网元 · {{cmd}} 条命令", }, eligible: { @@ -85,6 +87,8 @@ const zh = { retryFailedDone: "已开始重采失败网元", deleted: "任务已删除", nothingToRetry: "没有可重试的失败网元", + neBusy: "部分网元正在被其他采集任务占用,请等待其完成后再开始", + jobRunning: "该任务已在执行中", confirmDelete: "确定删除该采集任务?相关日志文件将一并删除。", jobs: { title: "采集任务", diff --git a/web/src/index.css b/web/src/index.css index 8db5e31..6a24b80 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -909,6 +909,14 @@ pre { margin-right: auto; } +.collect-create-panel__head { + margin-bottom: 0; +} + +.collect-create-panel--collapsed { + padding-bottom: 16px; +} + .collect-create-panel h3 { margin: 0; font-size: 15px; diff --git a/web/src/pages/CollectPage.tsx b/web/src/pages/CollectPage.tsx index b616f92..011a5cb 100644 --- a/web/src/pages/CollectPage.tsx +++ b/web/src/pages/CollectPage.tsx @@ -23,6 +23,13 @@ import { formatSystemTime } from "../utils/time"; const POLL_MS = 2000; const ELIGIBLE_PAGE_SIZE = 20; +function collectionErrorMessage(err: unknown, t: (key: string) => string): string { + const raw = String(err); + if (raw.includes("collection_ne_busy")) return t("collect.neBusy"); + if (raw.includes("collection_job_running")) return t("collect.jobRunning"); + return raw; +} + export function CollectPage() { const { t } = useI18n(); const { showOk, showError } = useToast(); @@ -36,6 +43,7 @@ export function CollectPage() { const [nePage, setNePage] = useState(1); const [jobPage, setJobPage] = useState(1); const [expandedJobId, setExpandedJobId] = useState(""); + const [createOpen, setCreateOpen] = useState(false); const selectedIds = useMemo(() => Object.keys(selectedMap), [selectedMap]); const selectedList = useMemo(() => Object.values(selectedMap), [selectedMap]); @@ -98,7 +106,7 @@ export function CollectPage() { showOk(t("collect.paused")); await invalidateJobs(job.id); }, - onError: (err) => showError(String(err)), + onError: (err) => showError(collectionErrorMessage(err, t)), }); const startJobMutation = useMutation({ @@ -108,7 +116,7 @@ export function CollectPage() { setExpandedJobId(job.id); await invalidateJobs(job.id); }, - onError: (err) => showError(String(err)), + onError: (err) => showError(collectionErrorMessage(err, t)), }); const retryFailedMutation = useMutation({ @@ -140,6 +148,7 @@ export function CollectPage() { }), onSuccess: async (job) => { showOk(t("collect.created", { id: job.id })); + setCreateOpen(false); setExpandedJobId(job.id); await queryClient.invalidateQueries({ queryKey: queryKeys.neCollectionsAll }); await queryClient.invalidateQueries({ queryKey: queryKeys.neCollectionDetail(job.id) }); @@ -199,12 +208,23 @@ export function CollectPage() { return (
{t("collect.create.hint")}
++ {createOpen + ? t("collect.create.hint") + : t("collect.create.meta", { ne: selectedIds.length, cmd: commandLines })} +
+