From f1ed4fe7533db36cfa304232b1b28bd5acea1859 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 28 May 2026 10:36:47 +0800 Subject: [PATCH] feat(collect): draft jobs, merged create UI, and unified start Create collection jobs without auto-running; add /start for first run and re-run. Merge create form with filterable NE picker and selected list. Co-authored-by: Cursor --- netx_api/collection_job_state.py | 4 +- netx_api/collection_router.py | 38 +++-- netx_api/collection_service.py | 53 ++++-- web/src/constants/queryKeys.ts | 2 +- web/src/i18n/en.ts | 24 ++- web/src/i18n/zh.ts | 24 ++- web/src/index.css | 58 +++++++ web/src/pages/CollectPage.tsx | 281 +++++++++++++++++++------------ web/src/services/api.ts | 6 +- 9 files changed, 343 insertions(+), 147 deletions(-) diff --git a/netx_api/collection_job_state.py b/netx_api/collection_job_state.py index d550eee..bb3ad07 100644 --- a/netx_api/collection_job_state.py +++ b/netx_api/collection_job_state.py @@ -67,8 +67,8 @@ 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": + # Pending while job is running/queued, or draft job not started yet. + if str(job.status or "") in ("running", "pending"): continue anchor = job.started_at or job.created_at limit = pending_stale_sec diff --git a/netx_api/collection_router.py b/netx_api/collection_router.py index e67cdbd..a7206fc 100644 --- a/netx_api/collection_router.py +++ b/netx_api/collection_router.py @@ -6,7 +6,7 @@ from sqlalchemy.orm import Session from .collection_service import ( build_collection_job_zip, - create_and_start_collection, + create_collection, delete_collection_job, get_collection_job, list_collection_jobs, @@ -15,6 +15,7 @@ from .collection_service import ( pause_collection_job, resolve_run_output_path, restart_collection_job, + start_collection_job, retry_failed_collection_job, ) from .collection_schemas import CollectionJobCreate @@ -29,25 +30,15 @@ router = APIRouter(prefix="/v1/ne-collections", tags=["ne-collections"]) def api_eligible_ne( page: int = Query(default=1, ge=1), page_size: int = Query(default=200, ge=1, le=500), + keyword: str = Query(default=""), db: Session = Depends(get_db), ): - return list_eligible_ne(db, page=page, page_size=page_size) + return list_eligible_ne(db, page=page, page_size=page_size, keyword=keyword) @router.post("") -def api_create_collection( - body: CollectionJobCreate, - background_tasks: BackgroundTasks, - db: Session = Depends(get_db), -): - out, payload = create_and_start_collection(db, body) - background_tasks.add_task( - dispatch_collection_runs, - payload["job_id"], - payload["run_ids"], - payload["commands"], - ) - return out.model_dump() +def api_create_collection(body: CollectionJobCreate, db: Session = Depends(get_db)): + return create_collection(db, body).model_dump() @router.get("") @@ -103,12 +94,29 @@ def api_pause_collection(job_id: str, db: Session = Depends(get_db)): return pause_collection_job(db, job_id).model_dump() +@router.post("/{job_id}/start") +def api_start_collection( + job_id: str, + background_tasks: BackgroundTasks, + db: Session = Depends(get_db), +): + out, payload = start_collection_job(db, job_id) + background_tasks.add_task( + dispatch_collection_runs, + payload["job_id"], + payload["run_ids"], + payload["commands"], + ) + return out.model_dump() + + @router.post("/{job_id}/restart") def api_restart_collection( job_id: str, background_tasks: BackgroundTasks, db: Session = Depends(get_db), ): + """Alias of /start for backward compatibility.""" out, payload = restart_collection_job(db, job_id) background_tasks.add_task( dispatch_collection_runs, diff --git a/netx_api/collection_service.py b/netx_api/collection_service.py index 9e41382..69f263b 100644 --- a/netx_api/collection_service.py +++ b/netx_api/collection_service.py @@ -111,8 +111,25 @@ def run_to_out(row: NeCollectionRun) -> CollectionRunOut: ) -def list_eligible_ne(db: Session, *, page: int = 1, page_size: int = 200) -> dict[str, Any]: +def list_eligible_ne( + db: Session, + *, + page: int = 1, + page_size: int = 200, + keyword: str = "", +) -> dict[str, Any]: stmt = db.query(ManagedNE).filter(ManagedNE.connect_status == "pass") + kw = str(keyword or "").strip() + if kw: + like = f"%{kw}%" + stmt = stmt.filter( + or_( + ManagedNE.name.ilike(like), + ManagedNE.ip_address.ilike(like), + ManagedNE.vendor.ilike(like), + ManagedNE.device_type.ilike(like), + ) + ) total = int(stmt.count()) rows = ( stmt.order_by(ManagedNE.name.asc()) @@ -135,9 +152,7 @@ def list_eligible_ne(db: Session, *, page: int = 1, page_size: int = 200) -> dic return {"total": total, "page": page, "page_size": page_size, "items": items} -def create_and_start_collection( - db: Session, body: CollectionJobCreate -) -> tuple[CollectionJobOut, CollectionSchedulePayload]: +def create_collection(db: Session, body: CollectionJobCreate) -> CollectionJobOut: commands = _parse_commands(body.commands) if not commands: raise HTTPException(status_code=400, detail="commands_empty") @@ -168,16 +183,15 @@ def create_and_start_collection( job = NeCollectionJob( title=str(body.title or "").strip() or f"collect-{now.strftime('%Y%m%d-%H%M%S')}", commands="\n".join(commands), - status="running", + status="pending", ne_count=len(ne_rows), created_at=now, - started_at=now, - last_run_at=now, + started_at=None, + last_run_at=None, ) db.add(job) db.flush() - run_ids: list[str] = [] for ne in ne_rows: run = NeCollectionRun( job_id=str(job.id), @@ -187,16 +201,9 @@ def create_and_start_collection( status="pending", ) db.add(run) - run_ids.append(str(run.id)) db.commit() db.refresh(job) - - payload: CollectionSchedulePayload = { - "job_id": str(job.id), - "run_ids": run_ids, - "commands": commands, - } - return job_to_out(job, output_count=0), payload + return job_to_out(job, output_count=0) def list_collection_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict[str, Any]: @@ -355,7 +362,8 @@ def _start_job_retry( return job_to_out(job, output_count=_output_count_for_job(db, job_id)), payload -def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]: +def start_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]: + """Start a draft job or re-run all NEs after pause/completion.""" job = db.get(NeCollectionJob, job_id) if not job: raise HTTPException(status_code=404, detail="collection_job_not_found") @@ -367,6 +375,13 @@ def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, commands = _parse_commands(str(job.commands or "")) if not commands: raise HTTPException(status_code=400, detail="commands_empty") + + 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) + retry_ids = _reset_runs_for_retry( db, job_id, @@ -376,6 +391,10 @@ def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, return _start_job_retry(db, job, job_id, retry_ids, commands, reset_all_counts=True) +def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]: + return start_collection_job(db, job_id) + + def retry_failed_collection_job( db: Session, job_id: str ) -> tuple[CollectionJobOut, CollectionSchedulePayload]: diff --git a/web/src/constants/queryKeys.ts b/web/src/constants/queryKeys.ts index 9318727..ff62286 100644 --- a/web/src/constants/queryKeys.ts +++ b/web/src/constants/queryKeys.ts @@ -15,7 +15,7 @@ export const queryKeys = { managedNe: (keyword: string, vendor: string, connectStatus: string, page: number, pageSize: number) => ["managedNe", keyword, vendor, connectStatus, page, pageSize] as const, collectionEligibleNeAll: ["collectionEligibleNe"] as const, - collectionEligibleNe: (page: number) => ["collectionEligibleNe", page] as const, + collectionEligibleNe: (page: number, keyword: string) => ["collectionEligibleNe", page, keyword] as const, neCollectionsAll: ["neCollections"] as const, neCollections: (page: number) => ["neCollections", page] as const, neCollectionDetail: (jobId: string) => ["neCollection", jobId] as const, diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 94d5d10..5ad60fd 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -42,6 +42,25 @@ const en = { langEn: "English", }, collect: { + create: { + title: "New collection job", + hint: "Configure commands and NEs, then create the job. Click Start in the job list to run collection.", + selectedTitle: "Selected NEs", + selectedEmpty: "Filter below and add NEs one by one.", + selectedCount: "{{count}} selected", + clearSelected: "Clear selection", + remove: "Remove", + pickTitle: "Eligible NEs (connectivity passed)", + pickHint: "Only NEs with connect_status=pass. Run connectivity test in NE Management first.", + pickEmpty: "No matching NEs.", + filterKeyword: "Filter", + filterKeywordPh: "Name / IP / vendor / device type", + add: "Add", + added: "Added", + create: "Create job", + creating: "Creating…", + meta: "{{ne}} NE(s) selected · {{cmd}} command(s)", + }, eligible: { title: "Eligible NEs (connectivity passed)", hint: "Only NEs with connect_status=pass. Run connectivity test in NE Management first.", @@ -58,7 +77,8 @@ const en = { start: "Start collection", starting: "Submitting…", }, - started: "Collection job started: {{id}}", + created: "Collection job created: {{id}}", + started: "Collection started: {{id}}", paused: "Job paused", restarted: "Collection restarted", retryFailedDone: "Retrying failed devices", @@ -80,6 +100,8 @@ const en = { expand: "Details", collapse: "Hide", pause: "Pause", + start: "Start", + starting: "Starting…", restart: "Restart", retryFailed: "Retry failed", downloadResults: "Download results", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index f44941a..94bc2f7 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -42,6 +42,25 @@ const zh = { langEn: "English", }, collect: { + create: { + title: "新建采集任务", + hint: "配置命令与网元后创建任务;创建后需在任务列表中点击「开始」执行采集。", + selectedTitle: "已选择网元", + selectedEmpty: "请从下方筛选并逐个添加网元。", + selectedCount: "共 {{count}} 台", + clearSelected: "清空已选", + remove: "移除", + pickTitle: "可选网元(连通性已通过)", + pickHint: "仅展示 connect_status=pass 的网元。请先在「网元管理」完成连通性测试。", + pickEmpty: "没有匹配的网元。", + filterKeyword: "筛选", + filterKeywordPh: "名称 / IP / 厂商 / 设备类型", + add: "添加", + added: "已添加", + create: "创建任务", + creating: "创建中…", + meta: "已选 {{ne}} 台网元 · {{cmd}} 条命令", + }, eligible: { title: "可选网元(连通性已通过)", hint: "仅展示 connect_status=pass 的网元。请先在「网元管理」完成连通性测试。", @@ -58,7 +77,8 @@ const zh = { start: "开始采集", starting: "提交中…", }, - started: "采集任务已启动:{{id}}", + created: "采集任务已创建:{{id}}", + started: "采集已开始:{{id}}", paused: "任务已暂停", restarted: "已重新开始采集", retryFailedDone: "已开始重采失败网元", @@ -80,6 +100,8 @@ const zh = { expand: "详情", collapse: "收起", pause: "暂停", + start: "开始", + starting: "启动中…", restart: "重新开始", retryFailed: "重采失败网元", downloadResults: "下载结果", diff --git a/web/src/index.css b/web/src/index.css index 913b1f6..8db5e31 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -909,6 +909,64 @@ pre { margin-right: auto; } +.collect-create-panel h3 { + margin: 0; + font-size: 15px; + font-weight: 600; + color: #334155; +} + +.collect-selected-block, +.collect-pick-block { + margin-top: 20px; + padding-top: 16px; + border-top: 1px solid #e2e8f0; +} + +.collect-selected-block__head { + display: flex; + flex-wrap: wrap; + align-items: center; + gap: 10px 16px; + margin-bottom: 10px; +} + +.collect-selected-list { + display: flex; + flex-direction: column; + gap: 8px; + max-height: 220px; + overflow: auto; + padding: 4px 0; +} + +.collect-selected-chip { + display: flex; + align-items: center; + justify-content: space-between; + gap: 12px; + padding: 8px 12px; + background: #f8fafc; + border: 1px solid #e2e8f0; + border-radius: 8px; +} + +.collect-selected-chip__main { + display: flex; + flex-direction: column; + gap: 2px; + min-width: 0; +} + +.collect-selected-chip__meta { + font-size: 12px; + color: #64748b; +} + +.collect-pick-filters { + margin-bottom: 12px; +} + .collect-cmd-preview { font-size: 12px; background: #f8fafc; diff --git a/web/src/pages/CollectPage.tsx b/web/src/pages/CollectPage.tsx index 0b77d1e..ea18732 100644 --- a/web/src/pages/CollectPage.tsx +++ b/web/src/pages/CollectPage.tsx @@ -8,7 +8,7 @@ import { fetchEligibleNe, fetchNeCollections, pauseCollectionJob, - restartCollectionJob, + startCollectionJob, retryFailedCollectionJob, collectionJobDownloadUrl, collectionRunDownloadUrl, @@ -20,6 +20,9 @@ import type { CollectionJobDetail, CollectionJobItem, EligibleNeItem } from "../ import { pageCount } from "../utils/display"; import { formatSystemTime } from "../utils/time"; +const POLL_MS = 2000; +const ELIGIBLE_PAGE_SIZE = 20; + export function CollectPage() { const { t } = useI18n(); const { showOk, showError } = useToast(); @@ -27,17 +30,19 @@ export function CollectPage() { const [commands, setCommands] = useState(""); const [title, setTitle] = useState(""); - const [selected, setSelected] = useState([]); + const [selectedMap, setSelectedMap] = useState>({}); + const [neKeyword, setNeKeyword] = useState(""); const [nePage, setNePage] = useState(1); const [jobPage, setJobPage] = useState(1); const [expandedJobId, setExpandedJobId] = useState(""); - const POLL_MS = 2000; - const ELIGIBLE_PAGE_SIZE = 20; + const selectedIds = useMemo(() => Object.keys(selectedMap), [selectedMap]); + const selectedList = useMemo(() => Object.values(selectedMap), [selectedMap]); const eligibleQuery = useQuery({ - queryKey: queryKeys.collectionEligibleNe(nePage), - queryFn: () => fetchEligibleNe({ page: nePage, pageSize: ELIGIBLE_PAGE_SIZE }), + queryKey: queryKeys.collectionEligibleNe(nePage, neKeyword), + queryFn: () => + fetchEligibleNe({ page: nePage, pageSize: ELIGIBLE_PAGE_SIZE, keyword: neKeyword }), staleTime: 5000, }); @@ -95,10 +100,10 @@ export function CollectPage() { onError: (err) => showError(String(err)), }); - const restartMutation = useMutation({ - mutationFn: restartCollectionJob, + const startJobMutation = useMutation({ + mutationFn: startCollectionJob, onSuccess: async (job) => { - showOk(t("collect.restarted")); + showOk(t("collect.started", { id: job.id })); setExpandedJobId(job.id); await invalidateJobs(job.id); }, @@ -125,15 +130,15 @@ export function CollectPage() { onError: (err) => showError(String(err)), }); - const startMutation = useMutation({ + const createMutation = useMutation({ mutationFn: () => createNeCollection({ title: title.trim(), commands, - ne_ids: selected, + ne_ids: selectedIds, }), onSuccess: async (job) => { - showOk(t("collect.started", { id: job.id })); + showOk(t("collect.created", { id: job.id })); setExpandedJobId(job.id); await queryClient.invalidateQueries({ queryKey: queryKeys.neCollectionsAll }); await queryClient.invalidateQueries({ queryKey: queryKeys.neCollectionDetail(job.id) }); @@ -141,18 +146,21 @@ export function CollectPage() { onError: (err) => showError(String(err)), }); - const items = eligibleQuery.data?.items ?? []; - const allSelected = items.length > 0 && items.every((x) => selected.includes(x.id)); - - const toggleAll = () => { - if (allSelected) { - const ids = new Set(items.map((x) => x.id)); - setSelected((prev) => prev.filter((id) => !ids.has(id))); - } else { - setSelected((prev) => [...new Set([...prev, ...items.map((x) => x.id)])]); - } + const addNe = (row: EligibleNeItem) => { + setSelectedMap((prev) => (prev[row.id] ? prev : { ...prev, [row.id]: row })); }; + const removeNe = (id: string) => { + setSelectedMap((prev) => { + const next = { ...prev }; + delete next[id]; + return next; + }); + }; + + const clearSelected = () => setSelectedMap({}); + + const eligibleItems = eligibleQuery.data?.items ?? []; const neTotal = eligibleQuery.data?.total ?? 0; const nePages = pageCount(neTotal, ELIGIBLE_PAGE_SIZE); @@ -168,77 +176,20 @@ export function CollectPage() { [commands], ); + const actionPending = + pauseMutation.isPending || + startJobMutation.isPending || + retryFailedMutation.isPending || + deleteMutation.isPending; + return (
-
-
-
-

{t("collect.eligible.title")}

-

{t("collect.eligible.hint")}

-
- +
+
+

{t("collect.create.title")}

+

{t("collect.create.hint")}

- {eligibleQuery.isLoading ?

{t("common.refreshing")}

: null} - {!eligibleQuery.isLoading && items.length === 0 ? ( -

{t("collect.eligible.empty")}

- ) : ( - - - - - - - - - - - - {items.map((row: EligibleNeItem) => ( - - - - - - - - ))} - -
- - {t("managedNe.col.name")}{t("managedNe.col.vendor")}{t("managedNe.col.ip")}{t("managedNe.col.connect")}
- - setSelected((prev) => - prev.includes(row.id) ? prev.filter((x) => x !== row.id) : [...prev, row.id], - ) - } - /> - {row.name || row.ip_address}{row.vendor}{row.ip_address} - {row.connect_status} -
- )} - {neTotal > 0 ? ( -
-
{t("common.pagerMeta", { total: neTotal, page: nePage, pages: nePages })}
-
- - -
-
- ) : null} -
-
-

{t("collect.form.title")}

-

{t("collect.form.commandsHint")}