Add LLDP pause/resume/stop with restart recovery, and localize ops task titles.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-04 22:23:44 +08:00
parent ba038a3af6
commit 8ded8eeb7d
14 changed files with 923 additions and 136 deletions

View file

@ -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()

View file

@ -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),

View file

@ -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)

View file

@ -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,

View file

@ -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",
]

View file

@ -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()
)

View file

@ -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

View file

@ -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",