diff --git a/netx_api/app_startup.py b/netx_api/app_startup.py index 650ad04..42bd032 100644 --- a/netx_api/app_startup.py +++ b/netx_api/app_startup.py @@ -69,18 +69,18 @@ def run_api_startup() -> None: ume_support._reset_runtime_pause_flags() ume_support._fail_stale_running_sync_jobs_on_startup() try: - from .topology_service import bootstrap_topology_tree, reclaim_stale_discover_jobs + from .topology_service import bootstrap_topology_tree, recover_lldp_discover_on_startup db_topo = SessionLocal() try: bootstrap_topology_tree(db_topo) - closed = reclaim_stale_discover_jobs(db_topo, force_all_open=True) - if closed: - _log.warning("startup: closed %s orphaned topology discover jobs", closed) + resumed = recover_lldp_discover_on_startup(db_topo) + if resumed: + _log.info("startup: resumed %s interrupted LLDP discover job(s)", resumed) finally: db_topo.close() except Exception: - _log.exception("startup: topology discover job cleanup failed") + _log.exception("startup: topology discover job recovery failed") if ume_support._needs_startup_alarm_sync_before_ws(): begin_startup_alarm_sync_gate() diff --git a/netx_api/lldp_collect_router.py b/netx_api/lldp_collect_router.py index 2811084..56fd363 100644 --- a/netx_api/lldp_collect_router.py +++ b/netx_api/lldp_collect_router.py @@ -14,7 +14,10 @@ from .lldp_collect_service import ( get_job_detail, get_policy, list_jobs, + pause_collect, + resume_collect, start_collect, + stop_collect, update_policy, ) @@ -43,6 +46,21 @@ def api_start(db: Session = Depends(get_db)) -> dict[str, Any]: return start_collect(db, trigger_mode="manual") +@router.post("/jobs/{job_id}/pause") +def api_pause_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + return pause_collect(db, job_id) + + +@router.post("/jobs/{job_id}/resume") +def api_resume_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + return resume_collect(db, job_id) + + +@router.post("/jobs/{job_id}/stop") +def api_stop_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + return stop_collect(db, job_id) + + @router.get("/jobs") def api_list_jobs( page: int = Query(default=1, ge=1), diff --git a/netx_api/lldp_collect_service.py b/netx_api/lldp_collect_service.py index eec6504..dcb23e4 100644 --- a/netx_api/lldp_collect_service.py +++ b/netx_api/lldp_collect_service.py @@ -18,10 +18,13 @@ from .models import LldpCollectPolicy, TopoDiscoverJob, TopoFabricStats from .topology_schemas import FabricDiscoverRequest from .topology_service import ( get_discover_job, + pause_discover_job, prune_discover_jobs, reclaim_stale_discover_jobs, refresh_fabric_stats, + resume_discover_job, start_discover_job, + stop_discover_job, ) POLICY_ID = 1 @@ -170,10 +173,11 @@ def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None: def has_running_job(db: Session) -> TopoDiscoverJob | None: + """Active job including paused — blocks starting a new collect.""" reclaim_stale_discover_jobs(db) return ( db.query(TopoDiscoverJob) - .filter(TopoDiscoverJob.status.in_(["pending", "running"])) + .filter(TopoDiscoverJob.status.in_(["pending", "running", "paused"])) .order_by(TopoDiscoverJob.created_at.desc()) .first() ) @@ -182,7 +186,7 @@ def has_running_job(db: Session) -> TopoDiscoverJob | None: def last_finished_job(db: Session) -> TopoDiscoverJob | None: return ( db.query(TopoDiscoverJob) - .filter(TopoDiscoverJob.status.in_(["done", "failed"])) + .filter(TopoDiscoverJob.status.in_(["done", "failed", "cancelled"])) .order_by(TopoDiscoverJob.created_at.desc()) .first() ) @@ -255,6 +259,18 @@ def start_collect(db: Session, *, trigger_mode: str = "manual") -> dict: return {"ok": True, "job": job.model_dump()} +def pause_collect(db: Session, job_id: str) -> dict: + return pause_discover_job(db, job_id).model_dump() + + +def resume_collect(db: Session, job_id: str) -> dict: + return resume_discover_job(db, job_id).model_dump() + + +def stop_collect(db: Session, job_id: str) -> dict: + return stop_discover_job(db, job_id).model_dump() + + def get_dashboard(db: Session) -> LldpCollectDashboardOut: policy = ensure_policy(db) running = has_running_job(db) diff --git a/netx_api/ops_tasks_service.py b/netx_api/ops_tasks_service.py index 2d37739..8b03047 100644 --- a/netx_api/ops_tasks_service.py +++ b/netx_api/ops_tasks_service.py @@ -1,4 +1,7 @@ -"""Unified live task overview across NetX runners.""" +"""Unified live task overview across NetX runners. + +Titles are language-neutral subjects; the UI prefixes localized kind labels. +""" from __future__ import annotations @@ -41,6 +44,7 @@ def _item( started_at: datetime | None = None, updated_at: datetime | None = None, progress: str = "", + inflight: int = 0, detail: str = "", href: str = "", ) -> dict[str, Any]: @@ -54,6 +58,7 @@ def _item( "started_at": _iso(started_at), "updated_at": _iso(updated_at), "progress": progress, + "inflight": int(inflight or 0), "detail": detail, "href": href, } @@ -106,7 +111,7 @@ def _port_traffic_items(db: Session, actors: dict[str, str]) -> list[dict[str, A _item( kind="port_traffic", id=did, - title=f"端口流量 · {name}", + title=name, status=status, trigger="schedule" if actor == "scheduler" else "manual", actor=actor, @@ -144,13 +149,14 @@ def _config_sync_items(db: Session, actors: dict[str, str]) -> list[dict[str, An _item( kind="config_sync", id=cid, - title=f"配置同步 · {cid[:8]}", + title=cid[:8], status=str(row.status or "pending"), trigger=trigger, actor=actor, started_at=row.started_at or row.created_at, updated_at=row.ended_at or row.started_at or row.created_at, - progress=f"{done}/{planned}" + (f" · running {running}" if running else ""), + progress=f"{done}/{planned}", + inflight=int(running), detail=str(row.error_message or "")[:240], href="/network/tasks/config-sync", ) @@ -174,8 +180,11 @@ def _collection_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any .filter(NeCollectionRun.job_id == jid, NeCollectionRun.status == "running") .count() ) - title = str(row.title or "").strip() or f"采集任务 · {jid[:8]}" + title = str(row.title or "").strip() or jid[:8] actor = actors.get(jid) or "—" + ok = int(row.success_count or 0) + fail = int(row.fail_count or 0) + total = int(row.ne_count or 0) items.append( _item( kind="ne_collect", @@ -186,11 +195,8 @@ def _collection_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any actor=actor, started_at=row.started_at or row.created_at, updated_at=row.last_run_at or row.ended_at or row.started_at or row.created_at, - progress=( - f"ok {int(row.success_count or 0)} / fail {int(row.fail_count or 0)}" - f" / total {int(row.ne_count or 0)}" - + (f" · running {running}" if running else "") - ), + progress=f"{ok}/{fail}/{total}", + inflight=int(running), detail=str(row.error_message or "")[:240], href="/network/tasks/collect", ) @@ -201,7 +207,7 @@ def _collection_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any def _lldp_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]: rows = ( db.query(TopoDiscoverJob) - .filter(TopoDiscoverJob.status.in_(("pending", "running"))) + .filter(TopoDiscoverJob.status.in_(("pending", "running", "paused"))) .order_by(TopoDiscoverJob.created_at.desc()) .limit(20) .all() @@ -219,7 +225,7 @@ def _lldp_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]: _item( kind="lldp_discover", id=jid, - title=f"LLDP 发现 · {trigger}", + title=trigger, status=str(row.status or "pending"), trigger=trigger, actor=actor, @@ -255,13 +261,13 @@ def _ume_sync_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]] _item( kind="ume_sync", id=jid, - title=f"UME 同步 · {domain}", + title=domain, status="running", trigger=trigger, actor=actor, started_at=row.started_at, updated_at=row.started_at, - progress=f"pull {pulled} · +{inserted} ~{updated}", + progress=f"{pulled}/+{inserted}/~{updated}", detail=str(row.error_message or "")[:240], href="/ume", ) @@ -286,7 +292,7 @@ def _ne_connect_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any _item( kind="ne_connect", id=f"managed:{nid}", - title=f"连通性测试 · {name}", + title=name, status="testing", trigger="manual", actor=actors.get(nid) or "—", @@ -310,7 +316,7 @@ def _ne_connect_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any _item( kind="ne_connect", id=f"ume:{uid}", - title=f"连通性测试 · UME {uid[:12]}", + title=f"UME {uid[:12]}", status="testing", trigger="manual", actor=actors.get(uid) or "—", @@ -338,7 +344,7 @@ def _ume_runtime_items() -> list[dict[str, Any]]: _item( kind="ume_runtime", id=task, - title=f"UME · {task}", + title=task, status=status, trigger="system", actor="system", @@ -372,28 +378,23 @@ def _webcrt_items(actors: dict[str, str]) -> list[dict[str, Any]]: progress = "" if lifecycle == "connecting": elapsed_ms = row.get("elapsed_ms") - progress = f"{int(elapsed_ms)} ms" if isinstance(elapsed_ms, int) else "logging in" - detail = "authenticating" + progress = f"{int(elapsed_ms)} ms" if isinstance(elapsed_ms, int) else "" elif lifecycle == "ready": - progress = "attached" if row.get("connect_ms") is not None: - detail = f"connect {int(row['connect_ms'])} ms" + progress = f"{int(row['connect_ms'])} ms" elif lifecycle == "detached": - progress = "detached" deadline = row.get("detach_deadline") if isinstance(deadline, (int, float)) and deadline > 0: left = max(0, int(deadline - now)) - detail = f"grace {left}s" - else: - detail = "awaiting reconnect / close" + progress = f"{left}s" + detail = str(row.get("connect_error") or "")[:240] elif lifecycle == "error": - progress = "error" detail = str(row.get("connect_error") or "")[:240] items.append( _item( kind="webcrt", id=sid, - title=f"WebCRT · {name}", + title=name, status=lifecycle, trigger="manual", actor=actor, diff --git a/netx_api/topology_discover.py b/netx_api/topology_discover.py index ee7463a..ee8b718 100644 --- a/netx_api/topology_discover.py +++ b/netx_api/topology_discover.py @@ -11,8 +11,12 @@ from .topology_discover_common import ( ) from .topology_discover_jobs import ( _run_discover_job, + pause_discover_job, reclaim_stale_discover_jobs, + recover_lldp_discover_on_startup, + resume_discover_job, start_discover_job, + stop_discover_job, ) from .topology_discover_scan import ( _apply_discover_hits, @@ -30,7 +34,11 @@ __all__ = [ "_run_discover_job", "_ume_target_dict", "get_discover_job", + "pause_discover_job", "prune_discover_jobs", "reclaim_stale_discover_jobs", + "recover_lldp_discover_on_startup", + "resume_discover_job", "start_discover_job", + "stop_discover_job", ] diff --git a/netx_api/topology_discover_common.py b/netx_api/topology_discover_common.py index 3a0e592..3f8a5e1 100644 --- a/netx_api/topology_discover_common.py +++ b/netx_api/topology_discover_common.py @@ -224,7 +224,7 @@ def prune_discover_jobs(db: Session, *, keep: int = 30) -> int: keep = max(0, min(200, int(keep))) finished = ( db.query(TopoDiscoverJob) - .filter(TopoDiscoverJob.status.in_(["done", "failed"])) + .filter(TopoDiscoverJob.status.in_(["done", "failed", "cancelled"])) .order_by(TopoDiscoverJob.created_at.desc()) .all() ) diff --git a/netx_api/topology_discover_jobs.py b/netx_api/topology_discover_jobs.py index 7b7a6aa..9921fff 100644 --- a/netx_api/topology_discover_jobs.py +++ b/netx_api/topology_discover_jobs.py @@ -1,9 +1,10 @@ -"""Discover job lifecycle: start, background run, stale reclaim.""" +"""Discover job lifecycle: start, pause/resume/stop, background run, stale reclaim.""" from __future__ import annotations import logging import threading -from concurrent.futures import ThreadPoolExecutor, as_completed +import time +from concurrent.futures import FIRST_COMPLETED, ThreadPoolExecutor, wait from datetime import datetime, timedelta from uuid import uuid4 @@ -29,18 +30,212 @@ from .topology_schemas import FabricDiscoverJobOut, FabricDiscoverRequest _log = logging.getLogger("netx.topology.discover") +_ACTIVE_STATUSES = ("pending", "running", "paused") +_FINISHED_STATUSES = ("done", "failed", "cancelled") -def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None: + +def _target_key(t: dict) -> str: + ume = str(t.get("ume_ne_id") or "").strip() + if ume: + return f"ume:{ume}" + return f"managed:{str(t.get('ne_id') or '').strip()}" + + +def _item_key(item: TopoDiscoverJobItem) -> str: + ume = str(item.ume_ne_id or "").strip() + if ume: + return f"ume:{ume}" + return f"managed:{str(item.ne_id or '').strip()}" + + +def _read_job_status(db: Session, job_id: str) -> str: + job = db.get(TopoDiscoverJob, job_id) + if job is None: + return "cancelled" + return str(job.status or "") + + +def _request_from_job(db: Session, job: TopoDiscoverJob) -> FabricDiscoverRequest: + """Rebuild a discover request from a stored job (resume after worker death).""" + pol = db.get(LldpCollectPolicy, 1) + concurrency = max(1, min(32, int(getattr(pol, "concurrency", None) or 4))) + auto_add = bool(getattr(pol, "auto_add_unmatched", True)) if pol is not None else True + scope = str(job.scope or "ne_ids").strip().lower() or "ne_ids" + stored = list(job.ne_ids_json or []) + managed_ids: list[str] = [] + ume_ids: list[str] = [] + legacy: list[str] = [] + for raw in stored: + s = str(raw or "").strip() + if not s: + continue + if s.startswith("managed:"): + managed_ids.append(s[len("managed:") :]) + elif s.startswith("ume:"): + ume_ids.append(s[len("ume:") :]) + else: + legacy.append(s) + if scope == "all_inventory": + return FabricDiscoverRequest( + scope="all_inventory", + ne_ids=[], + managed_ne_ids=[], + ume_ne_ids=[], + concurrency=concurrency, + auto_add_unmatched=auto_add, + trigger_mode=str(job.trigger_mode or "manual"), + ) + if managed_ids or ume_ids: + return FabricDiscoverRequest( + scope="ne_ids", + ne_ids=[], + managed_ne_ids=managed_ids, + ume_ne_ids=ume_ids, + concurrency=concurrency, + auto_add_unmatched=auto_add, + trigger_mode=str(job.trigger_mode or "manual"), + ) + return FabricDiscoverRequest( + scope="ne_ids", + ne_ids=legacy, + managed_ne_ids=[], + ume_ne_ids=[], + concurrency=concurrency, + auto_add_unmatched=auto_add, + trigger_mode=str(job.trigger_mode or "manual"), + ) + + +def _record_item( + db: Session, + job: TopoDiscoverJob, + job_id: str, + result: dict, + *, + added: int, + updated: int, +) -> tuple[int, int]: + item = TopoDiscoverJobItem( + id=uuid4().hex, + job_id=job_id, + ne_id=str(result.get("ne_id") or ""), + ume_ne_id=str(result.get("ume_ne_id") or ""), + fabric_node_id=str(result.get("fabric_node_id") or ""), + ne_name=str(result.get("ne_name") or "")[:256], + ne_ip=str(result.get("ne_ip") or "")[:128], + ok=bool(result.get("ok")), + command=str(result.get("command") or "")[:256], + neighbors=int(result.get("neighbors") or 0), + edges_added=int(result.get("edges_added") or 0), + edges_updated=int(result.get("edges_updated") or 0), + unmatched_count=int(result.get("unmatched_count") or 0), + unmatched_json=list(result.get("unmatched") or []), + parser_key=str(result.get("parser_key") or "")[:64], + parser_stub=bool(result.get("parser_stub")), + error=str(result.get("error") or "")[:1024], + raw_preview=str(result.get("raw_preview") or ""), + created_at=_utcnow(), + ) + db.add(item) + added += int(result.get("edges_added") or 0) + updated += int(result.get("edges_updated") or 0) + job.done = int(job.done or 0) + 1 + job.edges_added = added + job.edges_updated = updated + job.updated_at = _utcnow() + db.commit() + if int(job.done or 0) % 50 == 0: + try: + refresh_fabric_stats(db) + except Exception: # noqa: BLE001 + db.rollback() + return added, updated + + +def _finalize_success( + db: Session, + job_id: str, + *, + added: int, + updated: int, + stale: int, +) -> None: + stats = db.get(TopoFabricStats, "global") + if stats is None: + stats = TopoFabricStats(id="global") + db.add(stats) + stats.last_discover_at = _utcnow() + db.commit() + merge_duplicate_fabric_nodes(db) + refresh_fabric_stats(db) + + job = db.get(TopoDiscoverJob, job_id) + if job is None: + return + # Stop/pause may have won the race while we were finishing. + if str(job.status or "") in ("cancelled", "paused"): + return + job.status = "done" + job.ended_at = _utcnow() + job.updated_at = job.ended_at + job.edges_added = added + job.edges_updated = updated + job.edges_stale = stale + db.commit() + try: + from .lldp_collect_service import DEFAULT_HISTORY_KEEP, ensure_policy + + keep = int(getattr(ensure_policy(db), "history_keep", DEFAULT_HISTORY_KEEP) or 0) + prune_discover_jobs(db, keep=keep) + except Exception: # noqa: BLE001 + _log.warning("prune_discover_jobs failed job=%s", job_id, exc_info=True) + + +def _finalize_cancelled(db: Session, job_id: str, *, added: int, updated: int) -> None: + try: + refresh_fabric_stats(db) + except Exception: # noqa: BLE001 + db.rollback() + job = db.get(TopoDiscoverJob, job_id) + if job is None: + return + if str(job.status or "") != "cancelled": + job.status = "cancelled" + job.error = (str(job.error or "").strip() or "stopped_by_user")[:1024] + job.ended_at = _utcnow() + job.updated_at = job.ended_at + job.edges_added = added + job.edges_updated = updated + db.commit() + try: + from .lldp_collect_service import DEFAULT_HISTORY_KEEP, ensure_policy + + keep = int(getattr(ensure_policy(db), "history_keep", DEFAULT_HISTORY_KEEP) or 0) + prune_discover_jobs(db, keep=keep) + except Exception: # noqa: BLE001 + _log.warning("prune_discover_jobs failed job=%s", job_id, exc_info=True) + + +def _run_discover_job( + job_id: str, + body: FabricDiscoverRequest, + *, + resume: bool = False, +) -> None: db = SessionLocal() try: job = db.get(TopoDiscoverJob, job_id) if job is None: return - job.status = "running" - job.started_at = _utcnow() - job.updated_at = job.started_at + if str(job.status or "") == "cancelled": + return + if str(job.status or "") != "paused": + job.status = "running" + if not job.started_at: + job.started_at = _utcnow() + job.updated_at = _utcnow() try: - targets = _resolve_scan_targets(db, body) + all_targets = _resolve_scan_targets(db, body) except HTTPException as exc: job.status = "failed" job.error = str(exc.detail or "resolve_failed")[:1024] @@ -48,10 +243,27 @@ def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None: job.updated_at = job.ended_at db.commit() return - job.total = len(targets) + + prior_items = ( + db.query(TopoDiscoverJobItem).filter(TopoDiscoverJobItem.job_id == job_id).all() + ) + done_keys = {_item_key(it) for it in prior_items if _item_key(it) not in {"managed:", "ume:"}} + scanned_ok: set[str] = { + str(it.fabric_node_id) + for it in prior_items + if it.ok and str(it.fabric_node_id or "").strip() + } + touched_edges: set[str] = set() + # After worker death we lost in-memory touched edges — skip miss to avoid false marks. + skip_miss = bool(resume) + + if not prior_items: + job.total = len(all_targets) + elif int(job.total or 0) <= 0: + job.total = len(all_targets) db.commit() - # Reduce cross-worker races on self nodes before concurrent SSH/apply. + targets = [t for t in all_targets if _target_key(t) not in done_keys] try: _preensure_discover_targets(db, targets) except Exception: # noqa: BLE001 @@ -60,104 +272,146 @@ def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None: from .cli_budget import clamp_cli_workers concurrency = clamp_cli_workers(int(body.concurrency or 4)) - added = 0 - updated = 0 - stale = 0 - scanned_ok: set[str] = set() - touched_edges: set[str] = set() + added = int(job.edges_added or 0) + updated = int(job.edges_updated or 0) + stale = int(job.edges_stale or 0) + remaining = list(targets) + in_flight: dict = {} + cancelled = False with ThreadPoolExecutor(max_workers=concurrency) as pool: - futs = { - pool.submit( - _discover_one_target, t, auto_add_unmatched=bool(body.auto_add_unmatched) - ): t - for t in targets - } - for fut in as_completed(futs): - result = fut.result() - item = TopoDiscoverJobItem( - id=uuid4().hex, - job_id=job_id, - ne_id=str(result.get("ne_id") or ""), - ume_ne_id=str(result.get("ume_ne_id") or ""), - fabric_node_id=str(result.get("fabric_node_id") or ""), - ne_name=str(result.get("ne_name") or "")[:256], - ne_ip=str(result.get("ne_ip") or "")[:128], - ok=bool(result.get("ok")), - command=str(result.get("command") or "")[:256], - neighbors=int(result.get("neighbors") or 0), - edges_added=int(result.get("edges_added") or 0), - edges_updated=int(result.get("edges_updated") or 0), - unmatched_count=int(result.get("unmatched_count") or 0), - unmatched_json=list(result.get("unmatched") or []), - parser_key=str(result.get("parser_key") or "")[:64], - parser_stub=bool(result.get("parser_stub")), - error=str(result.get("error") or "")[:1024], - raw_preview=str(result.get("raw_preview") or ""), - created_at=_utcnow(), - ) - db.add(item) - added += int(result.get("edges_added") or 0) - updated += int(result.get("edges_updated") or 0) - if result.get("ok") and result.get("scanned_node_id"): - scanned_ok.add(str(result["scanned_node_id"])) - for eid in result.get("touched_edge_ids") or []: - touched_edges.add(str(eid)) - # Cutover edges already marked missing — skip same-job miss bump. - for eid in result.get("replaced_edge_ids") or []: - touched_edges.add(str(eid)) - job.done = int(job.done or 0) + 1 + while remaining or in_flight: + db.expire_all() + status = _read_job_status(db, job_id) + + if status == "cancelled": + cancelled = True + for fut in list(in_flight): + fut.cancel() + # Drain in-flight that already started (cancel is best-effort). + while in_flight: + done_set, _ = wait(in_flight.keys(), return_when=FIRST_COMPLETED) + for fut in done_set: + tgt = in_flight.pop(fut, None) + if fut.cancelled(): + continue + try: + result = fut.result() + except Exception as exc: # noqa: BLE001 + result = { + "ne_id": str((tgt or {}).get("ne_id") or ""), + "ume_ne_id": str((tgt or {}).get("ume_ne_id") or ""), + "ne_name": str((tgt or {}).get("ne_name") or ""), + "ne_ip": str((tgt or {}).get("ne_ip") or ""), + "ok": False, + "error": str(exc)[:1024], + } + job = db.get(TopoDiscoverJob, job_id) + if job is None: + break + added, updated = _record_item( + db, job, job_id, result, added=added, updated=updated + ) + break + + if status == "paused": + if not in_flight: + time.sleep(0.5) + continue + # Let in-flight finish, but do not submit more. + elif status in ("running", "pending"): + if status == "pending": + job = db.get(TopoDiscoverJob, job_id) + if job is not None: + job.status = "running" + job.updated_at = _utcnow() + db.commit() + while remaining and len(in_flight) < concurrency: + if _read_job_status(db, job_id) not in ("running", "pending"): + break + t = remaining.pop(0) + fut = pool.submit( + _discover_one_target, + t, + auto_add_unmatched=bool(body.auto_add_unmatched), + ) + in_flight[fut] = t + else: + # Unexpected terminal status. + cancelled = status == "cancelled" + break + + if not in_flight: + if status == "paused": + continue + break + + done_set, _ = wait(in_flight.keys(), timeout=0.5, return_when=FIRST_COMPLETED) + if not done_set: + continue + for fut in done_set: + tgt = in_flight.pop(fut, None) + if fut.cancelled(): + continue + try: + result = fut.result() + except Exception as exc: # noqa: BLE001 + result = { + "ne_id": str((tgt or {}).get("ne_id") or ""), + "ume_ne_id": str((tgt or {}).get("ume_ne_id") or ""), + "ne_name": str((tgt or {}).get("ne_name") or ""), + "ne_ip": str((tgt or {}).get("ne_ip") or ""), + "ok": False, + "error": str(exc)[:1024], + } + job = db.get(TopoDiscoverJob, job_id) + if job is None: + cancelled = True + break + added, updated = _record_item( + db, job, job_id, result, added=added, updated=updated + ) + if result.get("ok") and result.get("scanned_node_id"): + scanned_ok.add(str(result["scanned_node_id"])) + for eid in result.get("touched_edge_ids") or []: + touched_edges.add(str(eid)) + for eid in result.get("replaced_edge_ids") or []: + touched_edges.add(str(eid)) + + if cancelled or _read_job_status(db, job_id) == "cancelled": + _finalize_cancelled(db, job_id, added=added, updated=updated) + return + + # May still be paused with no remaining work — treat as done. + db.expire_all() + status = _read_job_status(db, job_id) + if status == "paused" and remaining: + # Worker exiting while paused with work left — keep paused for later resume. + job = db.get(TopoDiscoverJob, job_id) + if job is not None: job.edges_added = added job.edges_updated = updated job.updated_at = _utcnow() db.commit() - # Keep Fabric KPI cache roughly in sync during long runs (dashboard polls job). - if int(job.done or 0) % 50 == 0: - try: - refresh_fabric_stats(db) - except Exception: # noqa: BLE001 - db.rollback() + return - # Absent on a successfully scanned endpoint → missing; purge after N cycles. - if scanned_ok: + if scanned_ok and not skip_miss: newly_missing, purged = _apply_missing_and_purge( db, scanned_ok=scanned_ok, touched_edge_ids=touched_edges, ) stale = newly_missing + purged - job.edges_stale = stale - db.commit() + job = db.get(TopoDiscoverJob, job_id) + if job is not None: + job.edges_stale = stale + db.commit() - stats = db.get(TopoFabricStats, "global") - if stats is None: - stats = TopoFabricStats(id="global") - db.add(stats) - stats.last_discover_at = _utcnow() - db.commit() - merge_duplicate_fabric_nodes(db) - refresh_fabric_stats(db) - - job = db.get(TopoDiscoverJob, job_id) - if job is not None: - job.status = "done" - job.ended_at = _utcnow() - job.updated_at = job.ended_at - job.edges_added = added - job.edges_updated = updated - job.edges_stale = stale - db.commit() - try: - from .lldp_collect_service import DEFAULT_HISTORY_KEEP, ensure_policy - - keep = int(getattr(ensure_policy(db), "history_keep", DEFAULT_HISTORY_KEEP) or 0) - prune_discover_jobs(db, keep=keep) - except Exception: # noqa: BLE001 - _log.warning("prune_discover_jobs failed job=%s", job_id, exc_info=True) + _finalize_success(db, job_id, added=added, updated=updated, stale=stale) except Exception as exc: # noqa: BLE001 db.rollback() job = db.get(TopoDiscoverJob, job_id) - if job is not None: + if job is not None and str(job.status or "") not in ("cancelled", "paused"): job.status = "failed" job.error = str(exc)[:1024] job.ended_at = _utcnow() @@ -177,7 +431,7 @@ def reclaim_stale_discover_jobs( ) -> int: """Mark orphaned / hung discover jobs as failed so scheduling can proceed. - - ``force_all_open``: process restart — all pending/running rows are dead. + - ``force_all_open``: process restart — pending/running rows are dead (paused kept for resume). - Otherwise: pending older than pending_stale_sec, or running with stale updated_at. """ now = now or _utcnow() @@ -245,7 +499,7 @@ def start_discover_job( db.query(LldpCollectPolicy).filter(LldpCollectPolicy.id == 1).with_for_update().one() if ( db.query(TopoDiscoverJob) - .filter(TopoDiscoverJob.status.in_(["pending", "running"])) + .filter(TopoDiscoverJob.status.in_(list(_ACTIVE_STATUSES))) .first() is not None ): @@ -288,3 +542,169 @@ def start_discover_job( ) thread.start() return _job_out(db, job, include_items=False) + + +def pause_discover_job(db: Session, job_id: str) -> FabricDiscoverJobOut: + job = db.get(TopoDiscoverJob, str(job_id or "").strip()) + if job is None: + raise HTTPException(status_code=404, detail="discover_job_not_found") + if str(job.status or "") not in ("running", "pending"): + raise HTTPException(status_code=400, detail="job_not_running") + job.status = "paused" + job.updated_at = _utcnow() + db.commit() + db.refresh(job) + return _job_out(db, job, include_items=False) + + +def resume_discover_job(db: Session, job_id: str) -> FabricDiscoverJobOut: + job = db.get(TopoDiscoverJob, str(job_id or "").strip()) + if job is None: + raise HTTPException(status_code=404, detail="discover_job_not_found") + if str(job.status or "") != "paused": + raise HTTPException(status_code=400, detail="job_not_paused") + other = ( + db.query(TopoDiscoverJob) + .filter( + TopoDiscoverJob.id != job.id, + TopoDiscoverJob.status.in_(["pending", "running"]), + ) + .first() + ) + if other is not None: + raise HTTPException(status_code=409, detail="lldp_collect_already_running") + remaining = int(job.total or 0) - int(job.done or 0) + if int(job.total or 0) > 0 and remaining <= 0: + raise HTTPException(status_code=400, detail="no_pending_targets") + job.status = "running" + job.updated_at = _utcnow() + db.commit() + db.refresh(job) + + with _JOB_LOCK: + alive = job.id in _RUNNING_JOBS + if alive: + return _job_out(db, job, include_items=False) + + body = _request_from_job(db, job) + with _JOB_LOCK: + _RUNNING_JOBS.add(job.id) + thread = threading.Thread( + target=_run_discover_job, + args=(job.id, body), + kwargs={"resume": True}, + name=f"topo-discover-{job.id[:8]}", + daemon=True, + ) + thread.start() + return _job_out(db, job, include_items=False) + + +def stop_discover_job(db: Session, job_id: str) -> FabricDiscoverJobOut: + """Cancel remaining work and close the job (running/paused/pending).""" + job = db.get(TopoDiscoverJob, str(job_id or "").strip()) + if job is None: + raise HTTPException(status_code=404, detail="discover_job_not_found") + if str(job.status or "") not in _ACTIVE_STATUSES: + raise HTTPException(status_code=400, detail="job_not_active") + now = _utcnow() + job.status = "cancelled" + job.error = "stopped_by_user" + job.ended_at = now + job.updated_at = now + db.commit() + db.refresh(job) + # If no worker is attached (paused after restart), close immediately for clients. + with _JOB_LOCK: + alive = job.id in _RUNNING_JOBS + if not alive: + try: + refresh_fabric_stats(db) + except Exception: # noqa: BLE001 + db.rollback() + try: + from .lldp_collect_service import DEFAULT_HISTORY_KEEP, ensure_policy + + keep = int(getattr(ensure_policy(db), "history_keep", DEFAULT_HISTORY_KEEP) or 0) + prune_discover_jobs(db, keep=keep) + except Exception: # noqa: BLE001 + _log.warning("prune_discover_jobs after stop failed job=%s", job_id, exc_info=True) + return _job_out(db, job, include_items=False) + + +def recover_lldp_discover_on_startup(db: Session) -> int: + """Resume interrupted LLDP discover after process restart (config-sync style). + + - Keep the newest active job; mark older actives failed. + - ``paused`` stays paused (no auto dispatch) but still occupies the slot. + - ``pending`` / ``running`` are re-spawned with ``resume=True`` for remaining targets. + """ + actives = ( + db.query(TopoDiscoverJob) + .filter(TopoDiscoverJob.status.in_(list(_ACTIVE_STATUSES))) + .order_by(TopoDiscoverJob.created_at.asc()) + .all() + ) + if not actives: + return 0 + + primary = actives[-1] + now = _utcnow() + for stale in actives[:-1]: + _log.warning( + "lldp discover recovery closing older active job=%s (keep=%s)", + stale.id, + primary.id, + ) + stale.status = "failed" + stale.ended_at = now + stale.updated_at = now + msg = str(stale.error or "").strip() + stale.error = (msg + ("; " if msg else "") + "superseded_active_job")[:1024] + with _JOB_LOCK: + _RUNNING_JOBS.discard(stale.id) + db.commit() + db.refresh(primary) + + if str(primary.status or "") == "paused": + _log.info("lldp discover recovery job=%s stays paused (blocks new jobs)", primary.id) + return 0 + + total = int(primary.total or 0) + done = int(primary.done or 0) + if total > 0 and done >= total: + primary.status = "done" + primary.ended_at = now + primary.updated_at = now + db.commit() + _log.info("lldp discover recovery job=%s already complete", primary.id) + return 0 + + primary.status = "running" + if not primary.started_at: + primary.started_at = now + primary.updated_at = now + db.commit() + db.refresh(primary) + + body = _request_from_job(db, primary) + with _JOB_LOCK: + if primary.id in _RUNNING_JOBS: + _log.info("lldp discover recovery job=%s already has worker", primary.id) + return 0 + _RUNNING_JOBS.add(primary.id) + thread = threading.Thread( + target=_run_discover_job, + args=(primary.id, body), + kwargs={"resume": True}, + name=f"topo-discover-{primary.id[:8]}", + daemon=True, + ) + thread.start() + _log.info( + "lldp discover recovery resumed job=%s done=%s/%s", + primary.id, + done, + total, + ) + return 1 diff --git a/netx_api/topology_service.py b/netx_api/topology_service.py index 589741d..9dc29e4 100644 --- a/netx_api/topology_service.py +++ b/netx_api/topology_service.py @@ -8,9 +8,13 @@ from .topology_common import ( from .topology_discover import ( _apply_discover_hits, get_discover_job, + pause_discover_job, prune_discover_jobs, reclaim_stale_discover_jobs, + recover_lldp_discover_on_startup, + resume_discover_job, start_discover_job, + stop_discover_job, ) from .topology_fabric import ( _apply_missing_and_purge, @@ -79,13 +83,17 @@ __all__ = [ "merge_duplicate_fabric_nodes", "patch_view_edge_style", "patch_view_positions", + "pause_discover_job", "populate_view", "project_fabric_neighbors_to_view", "prune_discover_jobs", "reclaim_stale_discover_jobs", + "recover_lldp_discover_on_startup", "refresh_fabric_stats", "remove_view_nodes", + "resume_discover_job", "start_discover_job", + "stop_discover_job", "update_folder", "update_view", "upsert_fabric_edge", diff --git a/tests/test_lldp_collect.py b/tests/test_lldp_collect.py index 3390281..ff05aeb 100644 --- a/tests/test_lldp_collect.py +++ b/tests/test_lldp_collect.py @@ -14,6 +14,7 @@ from netx_api.lldp_collect_service import ( ensure_policy, get_dashboard, has_running_job, + last_finished_job, next_due_at, update_policy, ) @@ -299,6 +300,151 @@ class LldpCollectTests(unittest.TestCase): self.db.delete(ume_dup) self.db.commit() + def test_pause_resume_stop_job_control(self) -> None: + from netx_api.topology_discover_jobs import ( + pause_discover_job, + resume_discover_job, + stop_discover_job, + ) + from netx_api.topology_schemas import FabricDiscoverRequest + from netx_api.topology_service import start_discover_job + + now = datetime.utcnow() + job = TopoDiscoverJob( + id=uuid4().hex, + scope="ne_ids", + trigger_mode="manual", + ne_ids_json=["managed:m1"], + status="running", + total=10, + done=3, + created_at=now, + updated_at=now, + started_at=now, + ) + self.db.add(job) + ensure_policy(self.db) + self.db.commit() + + paused = pause_discover_job(self.db, job.id) + self.assertEqual(paused.status, "paused") + self.assertIsNotNone(has_running_job(self.db)) + + with self.assertRaises(Exception) as ctx: + start_discover_job(self.db, FabricDiscoverRequest(scope="ne_ids", ne_ids=["x"])) + self.assertEqual(getattr(ctx.exception, "detail", None), "lldp_collect_already_running") + + # Resume without live worker re-spawns thread; mark done so it has no work and + # avoid racing real SSH — use stop path for terminal instead when remaining=0. + job.done = 10 + self.db.commit() + with self.assertRaises(Exception) as ctx2: + resume_discover_job(self.db, job.id) + self.assertEqual(getattr(ctx2.exception, "detail", None), "no_pending_targets") + + job.done = 3 + job.status = "paused" + self.db.commit() + stopped = stop_discover_job(self.db, job.id) + self.assertEqual(stopped.status, "cancelled") + self.assertEqual(stopped.error, "stopped_by_user") + self.assertIsNone(has_running_job(self.db)) + last = last_finished_job(self.db) + self.assertIsNotNone(last) + assert last is not None + self.assertEqual(last.id, job.id) + self.assertEqual(last.status, "cancelled") + + def test_paused_survives_startup_reclaim(self) -> None: + now = datetime.utcnow() + job = TopoDiscoverJob( + id=uuid4().hex, + scope="all_inventory", + trigger_mode="manual", + status="paused", + total=5, + done=1, + created_at=now, + updated_at=now, + started_at=now, + ) + self.db.add(job) + self.db.commit() + closed = reclaim_stale_discover_jobs(self.db, force_all_open=True) + self.assertEqual(closed, 0) + self.db.refresh(job) + self.assertEqual(job.status, "paused") + + def test_recover_resumes_interrupted_running_job(self) -> None: + from unittest.mock import patch + + from netx_api.topology_discover_jobs import recover_lldp_discover_on_startup + + now = datetime.utcnow() + older = TopoDiscoverJob( + id=uuid4().hex, + scope="all_inventory", + trigger_mode="manual", + status="running", + total=10, + done=1, + created_at=now - timedelta(minutes=5), + updated_at=now - timedelta(minutes=5), + started_at=now - timedelta(minutes=5), + ) + primary = TopoDiscoverJob( + id=uuid4().hex, + scope="ne_ids", + trigger_mode="manual", + ne_ids_json=["managed:m1"], + status="running", + total=10, + done=3, + created_at=now, + updated_at=now, + started_at=now, + ) + self.db.add(older) + self.db.add(primary) + ensure_policy(self.db) + self.db.commit() + + with patch("netx_api.topology_discover_jobs.threading.Thread") as thread_cls: + thread_cls.return_value.start = lambda: None + n = recover_lldp_discover_on_startup(self.db) + self.assertEqual(n, 1) + self.db.refresh(older) + self.db.refresh(primary) + self.assertEqual(older.status, "failed") + self.assertIn("superseded_active_job", older.error or "") + self.assertEqual(primary.status, "running") + thread_cls.assert_called_once() + kwargs = thread_cls.call_args.kwargs + self.assertTrue(kwargs.get("kwargs", {}).get("resume")) + + def test_recover_keeps_paused_without_auto_resume(self) -> None: + from netx_api.topology_discover_jobs import recover_lldp_discover_on_startup + + now = datetime.utcnow() + job = TopoDiscoverJob( + id=uuid4().hex, + scope="all_inventory", + trigger_mode="manual", + status="paused", + total=8, + done=2, + created_at=now, + updated_at=now, + started_at=now, + ) + self.db.add(job) + self.db.commit() + n = recover_lldp_discover_on_startup(self.db) + self.assertEqual(n, 0) + self.db.refresh(job) + self.assertEqual(job.status, "paused") + self.assertIsNotNone(has_running_job(self.db)) + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 64b1f06..632f515 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -276,6 +276,13 @@ const en = { title: "LLDP links", collectNow: "Collect now", started: "LLDP collect started", + pause: "Pause", + resume: "Resume", + stop: "Stop", + confirmStop: "Stop this job and skip remaining NEs?", + paused: "Paused", + resumed: "Resumed", + stopped: "Collect stopped", policyTitle: "Collect policy", policySaved: "Policy saved", enabled: "Enable scheduled collect", @@ -1555,6 +1562,7 @@ const en = { hint: "Live view of running or pending background work. Actor prefers recent audit entries; scheduled work shows as scheduler/system.", empty: "No tasks to show right now.", open: "Open", + inflight: "running {{n}}", kpiActive: "Active / in-flight", kpiTotal: "Total rows", kpiKinds: "By kind", @@ -1578,6 +1586,13 @@ const en = { ume_runtime: "UME runtime", webcrt: "WebCRT", }, + trigger: { + manual: "Manual", + schedule: "Schedule", + topology: "Topology", + system: "System", + retry_failed: "Retry failed", + }, status: { collecting: "Collecting", running: "Running", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 9aafbfa..7c2078b 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -273,6 +273,13 @@ const zh = { title: "LLDP 链路", collectNow: "立即采集", started: "已启动 LLDP 采集", + pause: "暂停", + resume: "继续", + stop: "停止", + confirmStop: "停止后将取消未扫描的网元,确定停止本轮采集?", + paused: "已暂停", + resumed: "已继续", + stopped: "已停止采集", policyTitle: "采集策略", policySaved: "策略已保存", enabled: "启用周期调度", @@ -1546,6 +1553,7 @@ const zh = { hint: "汇总当前系统中正在运行或挂起的各类后台任务;发起人优先取近期审计记录,调度类任务显示为 scheduler/system。", empty: "当前没有可展示的任务。", open: "打开", + inflight: "进行中 {{n}}", kpiActive: "活跃/在途", kpiTotal: "条目总数", kpiKinds: "按类型", @@ -1569,6 +1577,13 @@ const zh = { ume_runtime: "UME 后台", webcrt: "WebCRT", }, + trigger: { + manual: "手动", + schedule: "调度", + topology: "拓扑", + system: "系统", + retry_failed: "重试失败", + }, status: { collecting: "采集中", running: "运行中", diff --git a/web/src/pages/audit/TaskOverviewPage.tsx b/web/src/pages/audit/TaskOverviewPage.tsx index 11417c7..93eed65 100644 --- a/web/src/pages/audit/TaskOverviewPage.tsx +++ b/web/src/pages/audit/TaskOverviewPage.tsx @@ -2,7 +2,7 @@ import { useMemo } from "react"; import { Link } from "react-router-dom"; import { useQuery } from "@tanstack/react-query"; import { useI18n } from "../../i18n"; -import { fetchOpsTasks } from "../../services/api"; +import { fetchOpsTasks, type OpsTaskItem } from "../../services/api"; import { formatSystemTime } from "../../utils/time"; const POLL_MS = 4000; @@ -16,18 +16,47 @@ function statusTone(status: string): string { return "other"; } -function kindLabel(kind: string, t: (k: string) => string): string { +function kindLabel(kind: string, t: (k: string, vars?: Record) => string): string { const key = `audit.tasks.kind.${kind}`; const tr = t(key); return tr === key ? kind : tr; } -function statusLabel(status: string, t: (k: string) => string): string { +function statusLabel(status: string, t: (k: string, vars?: Record) => string): string { const key = `audit.tasks.status.${status}`; const tr = t(key); return tr === key ? status : tr; } +function triggerLabel(trigger: string, t: (k: string, vars?: Record) => string): string { + const raw = String(trigger || "").trim(); + if (!raw || raw === "—") return "—"; + const key = `audit.tasks.trigger.${raw}`; + const tr = t(key); + return tr === key ? raw : tr; +} + +function taskTitle(row: OpsTaskItem, t: (k: string, vars?: Record) => string): string { + const kind = kindLabel(row.kind, t); + const subject = String(row.title || "").trim(); + if (!subject) return kind; + // Avoid "Kind · Kind · x" if an old backend still prefixed the localized kind. + if (subject === kind || subject.startsWith(`${kind} · `) || subject.startsWith(`${kind}·`)) { + return subject; + } + return `${kind} · ${subject}`; +} + +function progressLabel(row: OpsTaskItem, t: (k: string, vars?: Record) => string): string { + const base = String(row.progress || "").trim(); + const inflight = Number(row.inflight || 0); + if (inflight > 0) { + const extra = t("audit.tasks.inflight", { n: inflight }); + return base ? `${base} · ${extra}` : extra; + } + return base || "—"; +} + export function TaskOverviewPage() { const { t } = useI18n(); const query = useQuery({ @@ -111,7 +140,7 @@ export function TaskOverviewPage() { {kindLabel(row.kind, t)} -
{row.title}
+
{taskTitle(row, t)}
{row.detail ? (
{row.detail} @@ -124,8 +153,8 @@ export function TaskOverviewPage() { {row.actor || "—"} - {row.trigger || "—"} - {row.progress || "—"} + {triggerLabel(row.trigger, t)} + {progressLabel(row, t)} {formatSystemTime(row.started_at) || "—"} {formatSystemTime(row.updated_at) || "—"} diff --git a/web/src/pages/network/LldpLinksPage.tsx b/web/src/pages/network/LldpLinksPage.tsx index cc8c396..ebe6a02 100644 --- a/web/src/pages/network/LldpLinksPage.tsx +++ b/web/src/pages/network/LldpLinksPage.tsx @@ -6,7 +6,10 @@ import { fetchLldpCollectDashboard, fetchLldpCollectJob, fetchLldpCollectJobs, + pauseLldpCollectJob, + resumeLldpCollectJob, startLldpCollect, + stopLldpCollectJob, updateLldpCollectPolicy, } from "../../services/api"; import { queryKeys } from "../../constants/queryKeys"; @@ -55,7 +58,10 @@ export function LldpLinksPage() { staleTime: 1000, refetchInterval: (q) => { const running = q.state.data?.running_job; - return running && (running.status === "running" || running.status === "pending") ? POLL_MS : false; + return running && + (running.status === "running" || running.status === "pending" || running.status === "paused") + ? POLL_MS + : false; }, }); @@ -187,6 +193,33 @@ export function LldpLinksPage() { onError: (err) => showError(String(err)), }); + const pauseMut = useMutation({ + mutationFn: (id: string) => pauseLldpCollectJob(id), + onSuccess: async () => { + showOk(t("lldpLinks.paused")); + await refresh(); + }, + onError: (err) => showError(String(err)), + }); + + const resumeMut = useMutation({ + mutationFn: (id: string) => resumeLldpCollectJob(id), + onSuccess: async () => { + showOk(t("lldpLinks.resumed")); + await refresh(); + }, + onError: (err) => showError(String(err)), + }); + + const stopMut = useMutation({ + mutationFn: (id: string) => stopLldpCollectJob(id), + onSuccess: async () => { + showOk(t("lldpLinks.stopped")); + await refresh(); + }, + onError: (err) => showError(String(err)), + }); + const dash = dashQuery.data; const running = dash?.running_job; const last = dash?.last_job; @@ -228,6 +261,28 @@ export function LldpLinksPage() { > {t("lldpLinks.collectNow")} + {running?.status === "running" || running?.status === "pending" ? ( + + ) : null} + {running?.status === "paused" ? ( + + ) : null} + {running && + (running.status === "running" || running.status === "paused" || running.status === "pending") ? ( + + ) : null}
@@ -536,6 +591,7 @@ export function LldpLinksPage() { {t("lldpLinks.col.missingDelta")} {t("lldpLinks.col.started")} {t("lldpLinks.col.ended")} + @@ -559,7 +615,19 @@ export function LldpLinksPage() { {job.id.slice(0, 8)} {job.trigger_mode} {job.scope} - {job.status} + + + {job.status} + + {job.done}/{job.total} @@ -569,12 +637,45 @@ export function LldpLinksPage() { {job.edges_missing ?? job.edges_stale ?? 0} {job.started_at ? formatSystemTime(job.started_at) : "—"} {job.ended_at ? formatSystemTime(job.ended_at) : "—"} + +
+ {job.status === "running" || job.status === "pending" ? ( + + ) : null} + {job.status === "paused" ? ( + + ) : null} + {job.status === "running" || job.status === "paused" || job.status === "pending" ? ( + + ) : null} +
+ ); })} {!jobs.length ? ( - + {t("common.empty")} diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 77ee68c..c59a2b6 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -215,6 +215,7 @@ export type OpsTaskItem = { started_at: string | null; updated_at: string | null; progress: string; + inflight?: number; detail: string; href: string; }; @@ -1144,6 +1145,15 @@ export const updateLldpCollectPolicy = (body: Partial) => export const startLldpCollect = () => apiPost<{ ok: boolean; job: TopologyDiscoverJob }>("/v1/topology/lldp-collect/start", {}); +export const pauseLldpCollectJob = (jobId: string) => + apiPost(`/v1/topology/lldp-collect/jobs/${encodeURIComponent(jobId)}/pause`, {}); + +export const resumeLldpCollectJob = (jobId: string) => + apiPost(`/v1/topology/lldp-collect/jobs/${encodeURIComponent(jobId)}/resume`, {}); + +export const stopLldpCollectJob = (jobId: string) => + apiPost(`/v1/topology/lldp-collect/jobs/${encodeURIComponent(jobId)}/stop`, {}); + export const fetchLldpCollectJobs = (params: { page?: number; pageSize?: number }) => { const p = new URLSearchParams(); p.set("page", String(Math.max(1, Number(params.page || 1))));