From dcdf4147a6f686e1c76868feedda45ea88052d8e Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 11 Aug 2026 11:25:57 +0800 Subject: [PATCH] Add LLDP collect retry for failed and weak-evidence NEs. Mirror config-sync retry_failed so operators can re-collect hard failures and empty/stub LLDP results without re-running the full scope. Co-authored-by: Cursor --- netx_api/lldp_collect_router.py | 10 +- netx_api/lldp_collect_schemas.py | 9 +- netx_api/lldp_collect_service.py | 206 +++++++++++++++++++++++- netx_api/topology_discover_jobs.py | 2 +- tests/test_lldp_collect.py | 170 ++++++++++++++++++- web/src/i18n/en.ts | 2 + web/src/i18n/zh.ts | 2 + web/src/pages/network/LldpLinksPage.tsx | 36 ++++- web/src/services/api.ts | 7 +- web/src/types.ts | 2 + 10 files changed, 423 insertions(+), 23 deletions(-) diff --git a/netx_api/lldp_collect_router.py b/netx_api/lldp_collect_router.py index 56fd363..78c1bd7 100644 --- a/netx_api/lldp_collect_router.py +++ b/netx_api/lldp_collect_router.py @@ -8,7 +8,7 @@ from fastapi import APIRouter, Depends, Query from sqlalchemy.orm import Session from .db import get_db -from .lldp_collect_schemas import LldpCollectPolicyUpdate +from .lldp_collect_schemas import LldpCollectPolicyUpdate, LldpCollectStartBody from .lldp_collect_service import ( get_dashboard, get_job_detail, @@ -16,7 +16,7 @@ from .lldp_collect_service import ( list_jobs, pause_collect, resume_collect, - start_collect, + start_collect_from_body, stop_collect, update_policy, ) @@ -42,8 +42,10 @@ def api_dashboard(db: Session = Depends(get_db)) -> dict[str, Any]: @router.post("/start") -def api_start(db: Session = Depends(get_db)) -> dict[str, Any]: - return start_collect(db, trigger_mode="manual") +def api_start( + body: LldpCollectStartBody | None = None, db: Session = Depends(get_db) +) -> dict[str, Any]: + return start_collect_from_body(db, body) @router.post("/jobs/{job_id}/pause") diff --git a/netx_api/lldp_collect_schemas.py b/netx_api/lldp_collect_schemas.py index 39e9e36..5bb4d33 100644 --- a/netx_api/lldp_collect_schemas.py +++ b/netx_api/lldp_collect_schemas.py @@ -3,7 +3,7 @@ from __future__ import annotations from datetime import datetime -from typing import Any +from typing import Any, Literal from pydantic import BaseModel, Field @@ -43,6 +43,8 @@ class LldpCollectJobSummary(BaseModel): status: str = "" total: int = 0 done: int = 0 + success_count: int = 0 + fail_count: int = 0 edges_added: int = 0 edges_updated: int = 0 edges_stale: int = 0 # legacy alias of edges_missing @@ -66,6 +68,11 @@ class LldpCollectDashboardOut(BaseModel): next_due_at: datetime | None = None +class LldpCollectStartBody(BaseModel): + mode: Literal["full", "retry_failed"] = "full" + job_id: str | None = None + + class LldpCollectStartOut(BaseModel): ok: bool = True job: dict[str, Any] = Field(default_factory=dict) diff --git a/netx_api/lldp_collect_service.py b/netx_api/lldp_collect_service.py index dcb23e4..5cd0b01 100644 --- a/netx_api/lldp_collect_service.py +++ b/netx_api/lldp_collect_service.py @@ -5,6 +5,7 @@ from __future__ import annotations from datetime import datetime, timedelta from fastapi import HTTPException +from sqlalchemy import or_ from sqlalchemy.orm import Session from .lldp_collect_schemas import ( @@ -12,9 +13,10 @@ from .lldp_collect_schemas import ( LldpCollectJobSummary, LldpCollectPolicyOut, LldpCollectPolicyUpdate, + LldpCollectStartBody, LldpCollectTargetRef, ) -from .models import LldpCollectPolicy, TopoDiscoverJob, TopoFabricStats +from .models import LldpCollectPolicy, TopoDiscoverJob, TopoDiscoverJobItem, TopoFabricStats from .topology_schemas import FabricDiscoverRequest from .topology_service import ( get_discover_job, @@ -31,6 +33,12 @@ POLICY_ID = 1 DEFAULT_HISTORY_KEEP = 30 MAX_INTERVAL_HOURS = 8760 # 365d +# Soft failures: CLI ran but produced no trustworthy LLDP evidence (parity with +# config-sync empty_config_output — retryable as "no valid new data"). +_WEAK_LLDP_ERRORS = frozenset( + {"parser_stub", "empty_cli_output", "vendor_or_device_type_required"} +) + def _utcnow() -> datetime: return datetime.utcnow() @@ -151,7 +159,72 @@ def update_policy(db: Session, body: LldpCollectPolicyUpdate) -> LldpCollectPoli return _policy_out(row) -def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None: +def item_needs_retry(item: TopoDiscoverJobItem) -> bool: + """True when the NE failed or produced no trustworthy LLDP evidence.""" + if not bool(item.ok): + return True + if bool(item.parser_stub): + return True + err = str(item.error or "").strip() + return err in _WEAK_LLDP_ERRORS + + +def _retryable_item_filter(): + return or_( + TopoDiscoverJobItem.ok.is_(False), + TopoDiscoverJobItem.parser_stub.is_(True), + TopoDiscoverJobItem.error.in_(list(_WEAK_LLDP_ERRORS)), + ) + + +def _count_job_outcomes(db: Session, job_id: str) -> tuple[int, int]: + """Return (success_count, fail_count) for a discover job's items.""" + items = ( + db.query(TopoDiscoverJobItem.ok, TopoDiscoverJobItem.parser_stub, TopoDiscoverJobItem.error) + .filter(TopoDiscoverJobItem.job_id == job_id) + .all() + ) + fail = 0 + for ok, stub, error in items: + if (not bool(ok)) or bool(stub) or str(error or "").strip() in _WEAK_LLDP_ERRORS: + fail += 1 + total = len(items) + return max(0, total - fail), fail + + +def _job_outcome_map(db: Session, job_ids: list[str]) -> dict[str, tuple[int, int]]: + """Batch (success_count, fail_count) for many jobs.""" + if not job_ids: + return {} + rows = ( + db.query( + TopoDiscoverJobItem.job_id, + TopoDiscoverJobItem.ok, + TopoDiscoverJobItem.parser_stub, + TopoDiscoverJobItem.error, + ) + .filter(TopoDiscoverJobItem.job_id.in_(job_ids)) + .all() + ) + totals: dict[str, int] = {jid: 0 for jid in job_ids} + fails: dict[str, int] = {jid: 0 for jid in job_ids} + for job_id, ok, stub, error in rows: + jid = str(job_id) + totals[jid] = totals.get(jid, 0) + 1 + if (not bool(ok)) or bool(stub) or str(error or "").strip() in _WEAK_LLDP_ERRORS: + fails[jid] = fails.get(jid, 0) + 1 + return { + jid: (max(0, totals.get(jid, 0) - fails.get(jid, 0)), fails.get(jid, 0)) + for jid in job_ids + } + + +def _job_summary( + job: TopoDiscoverJob | None, + *, + success_count: int = 0, + fail_count: int = 0, +) -> LldpCollectJobSummary | None: if job is None: return None return LldpCollectJobSummary( @@ -161,6 +234,8 @@ def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None: status=job.status or "", total=int(job.total or 0), done=int(job.done or 0), + success_count=int(success_count or 0), + fail_count=int(fail_count or 0), edges_added=int(job.edges_added or 0), edges_updated=int(job.edges_updated or 0), edges_stale=int(job.edges_stale or 0), @@ -172,6 +247,78 @@ def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None: ) +def _summarize_job(db: Session, job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None: + if job is None: + return None + ok_n, fail_n = _count_job_outcomes(db, job.id) + return _job_summary(job, success_count=ok_n, fail_count=fail_n) + + +def collect_retry_targets(db: Session, src: TopoDiscoverJob) -> tuple[list[str], list[str]]: + """Build managed/ume NE id lists from retryable items of a prior job.""" + items = ( + db.query(TopoDiscoverJobItem) + .filter(TopoDiscoverJobItem.job_id == src.id) + .order_by(TopoDiscoverJobItem.created_at.asc()) + .all() + ) + managed: list[str] = [] + ume: list[str] = [] + seen: set[str] = set() + for it in items: + if not item_needs_retry(it): + continue + ume_id = str(it.ume_ne_id or "").strip() + ne_id = str(it.ne_id or "").strip() + if ume_id: + key = f"ume:{ume_id}" + if key in seen: + continue + seen.add(key) + ume.append(ume_id) + elif ne_id: + key = f"managed:{ne_id}" + if key in seen: + continue + seen.add(key) + managed.append(ne_id) + return managed, ume + + +def find_retry_source_job(db: Session, job_id: str | None = None) -> TopoDiscoverJob: + """Resolve the source job that still has retryable NE items.""" + src_id = str(job_id or "").strip() + if src_id: + src = db.get(TopoDiscoverJob, src_id) + if src is None: + raise HTTPException(status_code=404, detail="job_not_found") + return src + # Newest finished job that still has failed / weak-evidence items. + candidate_ids = [ + row[0] + for row in ( + db.query(TopoDiscoverJobItem.job_id) + .filter(_retryable_item_filter()) + .distinct() + .all() + ) + ] + if not candidate_ids: + raise HTTPException(status_code=404, detail="no_failed_job") + src = ( + db.query(TopoDiscoverJob) + .filter( + TopoDiscoverJob.id.in_(candidate_ids), + TopoDiscoverJob.status.in_(["done", "failed", "cancelled"]), + ) + .order_by(TopoDiscoverJob.created_at.desc()) + .first() + ) + if src is None: + raise HTTPException(status_code=404, detail="no_failed_job") + return src + + def has_running_job(db: Session) -> TopoDiscoverJob | None: """Active job including paused — blocks starting a new collect.""" reclaim_stale_discover_jobs(db) @@ -249,16 +396,52 @@ def build_discover_request(policy: LldpCollectPolicy) -> FabricDiscoverRequest: ) -def start_collect(db: Session, *, trigger_mode: str = "manual") -> dict: +def start_collect( + db: Session, + *, + trigger_mode: str = "manual", + mode: str = "full", + job_id: str | None = None, +) -> dict: if has_running_job(db) is not None: raise HTTPException(status_code=409, detail="lldp_collect_already_running") policy = ensure_policy(db) - body = build_discover_request(policy) - job = start_discover_job(db, body, trigger_mode=trigger_mode) + collect_mode = str(mode or "full").strip().lower() or "full" + if collect_mode == "retry_failed": + src = find_retry_source_job(db, job_id) + managed_ids, ume_ids = collect_retry_targets(db, src) + if not managed_ids and not ume_ids: + raise HTTPException(status_code=400, detail="no_failed_targets") + concurrency = max(1, min(32, int(policy.concurrency or 4))) + body = FabricDiscoverRequest( + scope="ne_ids", + ne_ids=[], + managed_ne_ids=managed_ids, + ume_ne_ids=ume_ids, + concurrency=concurrency, + auto_add_unmatched=bool(policy.auto_add_unmatched), + ) + trig = "retry_failed" + else: + body = build_discover_request(policy) + trig = str(trigger_mode or "manual").strip().lower() or "manual" + if trig == "retry_failed": + trig = "manual" + job = start_discover_job(db, body, trigger_mode=trig) prune_discover_jobs(db, keep=int(getattr(policy, "history_keep", DEFAULT_HISTORY_KEEP) or 0)) return {"ok": True, "job": job.model_dump()} +def start_collect_from_body(db: Session, body: LldpCollectStartBody | None = None) -> dict: + payload = body or LldpCollectStartBody() + return start_collect( + db, + trigger_mode="manual", + mode=str(payload.mode or "full"), + job_id=payload.job_id, + ) + + def pause_collect(db: Session, job_id: str) -> dict: return pause_discover_job(db, job_id).model_dump() @@ -292,8 +475,8 @@ def get_dashboard(db: Session) -> LldpCollectDashboardOut: fabric_edge_stale=int(stats.edge_stale if stats else 0), fabric_edge_missing=int(stats.edge_stale if stats else 0), last_discover_at=stats.last_discover_at if stats else None, - running_job=_job_summary(running), - last_job=_job_summary(last), + running_job=_summarize_job(db, running), + last_job=_summarize_job(db, last), next_due_at=next_due_at(db, policy), ) @@ -304,11 +487,18 @@ def list_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict: q = db.query(TopoDiscoverJob).order_by(TopoDiscoverJob.created_at.desc()) total = int(q.count()) rows = q.offset((page - 1) * page_size).limit(page_size).all() + outcomes = _job_outcome_map(db, [r.id for r in rows]) + items: list[dict] = [] + for r in rows: + ok_n, fail_n = outcomes.get(r.id, (0, 0)) + summary = _job_summary(r, success_count=ok_n, fail_count=fail_n) + if summary is not None: + items.append(summary.model_dump()) return { "total": total, "page": page, "page_size": page_size, - "items": [_job_summary(r).model_dump() for r in rows if _job_summary(r)], + "items": items, } diff --git a/netx_api/topology_discover_jobs.py b/netx_api/topology_discover_jobs.py index f38707f..bf78a89 100644 --- a/netx_api/topology_discover_jobs.py +++ b/netx_api/topology_discover_jobs.py @@ -526,7 +526,7 @@ def start_discover_job( if scope not in {"all_inventory", "ne_ids"}: raise HTTPException(status_code=400, detail="invalid_scope") trig = str(trigger_mode or getattr(body, "trigger_mode", None) or "manual").strip().lower() or "manual" - if trig not in {"manual", "schedule", "topology"}: + if trig not in {"manual", "schedule", "topology", "retry_failed"}: trig = "manual" now = _utcnow() # Persist explicit source lists when present; keep legacy ne_ids for older clients. diff --git a/tests/test_lldp_collect.py b/tests/test_lldp_collect.py index ff05aeb..6cdcbda 100644 --- a/tests/test_lldp_collect.py +++ b/tests/test_lldp_collect.py @@ -445,6 +445,174 @@ class LldpCollectTests(unittest.TestCase): self.assertEqual(job.status, "paused") self.assertIsNotNone(has_running_job(self.db)) + def test_item_needs_retry_covers_hard_and_weak(self) -> None: + from netx_api.lldp_collect_service import item_needs_retry + + hard = TopoDiscoverJobItem(id=uuid4().hex, job_id="j", ok=False, error="exec_failed") + weak_empty = TopoDiscoverJobItem( + id=uuid4().hex, job_id="j", ok=True, error="empty_cli_output" + ) + weak_stub = TopoDiscoverJobItem( + id=uuid4().hex, job_id="j", ok=True, parser_stub=True, error="parser_stub" + ) + ok = TopoDiscoverJobItem(id=uuid4().hex, job_id="j", ok=True, error="") + self.assertTrue(item_needs_retry(hard)) + self.assertTrue(item_needs_retry(weak_empty)) + self.assertTrue(item_needs_retry(weak_stub)) + self.assertFalse(item_needs_retry(ok)) + + def test_collect_retry_targets_and_start_retry_failed(self) -> None: + from unittest.mock import patch + + from netx_api.lldp_collect_service import ( + collect_retry_targets, + find_retry_source_job, + start_collect, + ) + from netx_api.topology_schemas import FabricDiscoverJobOut + + now = datetime.utcnow() + job_id = uuid4().hex + older_ok = uuid4().hex + # Older finished job with a failure — should not be picked when a newer one has fails. + self.db.add( + TopoDiscoverJob( + id=older_ok, + scope="all_inventory", + trigger_mode="manual", + status="done", + total=1, + done=1, + created_at=now - timedelta(hours=2), + updated_at=now - timedelta(hours=2), + ended_at=now - timedelta(hours=2), + ) + ) + self.db.add( + TopoDiscoverJobItem( + id=uuid4().hex, + job_id=older_ok, + ne_id="old-fail", + ok=False, + error="timeout", + created_at=now - timedelta(hours=2), + ) + ) + self.db.add( + TopoDiscoverJob( + id=job_id, + scope="all_inventory", + trigger_mode="manual", + status="done", + total=3, + done=3, + created_at=now - timedelta(hours=1), + updated_at=now - timedelta(hours=1), + ended_at=now - timedelta(hours=1), + ) + ) + self.db.add( + TopoDiscoverJobItem( + id=uuid4().hex, + job_id=job_id, + ne_id="ok-ne", + ok=True, + error="", + created_at=now, + ) + ) + self.db.add( + TopoDiscoverJobItem( + id=uuid4().hex, + job_id=job_id, + ne_id="fail-ne", + ok=False, + error="cli_budget_unavailable", + created_at=now, + ) + ) + self.db.add( + TopoDiscoverJobItem( + id=uuid4().hex, + job_id=job_id, + ume_ne_id="ume-weak", + ok=True, + error="empty_cli_output", + created_at=now, + ) + ) + ensure_policy(self.db) + self.db.commit() + + src = find_retry_source_job(self.db) + self.assertEqual(src.id, job_id) + managed, ume = collect_retry_targets(self.db, src) + self.assertEqual(managed, ["fail-ne"]) + self.assertEqual(ume, ["ume-weak"]) + + dash = get_dashboard(self.db) + self.assertIsNotNone(dash.last_job) + assert dash.last_job is not None + self.assertEqual(dash.last_job.fail_count, 2) + self.assertEqual(dash.last_job.success_count, 1) + + fake_out = FabricDiscoverJobOut( + id=uuid4().hex, + scope="ne_ids", + trigger_mode="retry_failed", + status="pending", + total=2, + done=0, + ) + with patch( + "netx_api.lldp_collect_service.start_discover_job", return_value=fake_out + ) as start_mock: + out = start_collect(self.db, mode="retry_failed") + self.assertTrue(out["ok"]) + self.assertEqual(out["job"]["trigger_mode"], "retry_failed") + req = start_mock.call_args.args[1] + self.assertEqual(req.scope, "ne_ids") + self.assertEqual(req.managed_ne_ids, ["fail-ne"]) + self.assertEqual(req.ume_ne_ids, ["ume-weak"]) + self.assertEqual(start_mock.call_args.kwargs.get("trigger_mode"), "retry_failed") + + def test_start_retry_failed_no_targets(self) -> None: + from fastapi import HTTPException + + from netx_api.lldp_collect_service import start_collect + + now = datetime.utcnow() + job_id = uuid4().hex + self.db.add( + TopoDiscoverJob( + id=job_id, + scope="all_inventory", + trigger_mode="manual", + status="done", + total=1, + done=1, + created_at=now, + updated_at=now, + ended_at=now, + ) + ) + self.db.add( + TopoDiscoverJobItem( + id=uuid4().hex, + job_id=job_id, + ne_id="ok-only", + ok=True, + error="", + created_at=now, + ) + ) + ensure_policy(self.db) + self.db.commit() + with self.assertRaises(HTTPException) as ctx: + start_collect(self.db, mode="retry_failed") + self.assertEqual(ctx.exception.status_code, 404) + self.assertEqual(ctx.exception.detail, "no_failed_job") + if __name__ == "__main__": - unittest.main() + unittest.main() \ No newline at end of file diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 23f8c02..582b3f5 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -319,6 +319,7 @@ const en = { lldpLinks: { title: "LLDP links", collectNow: "Collect now", + retryFailed: "Retry failed", started: "LLDP collect started", pause: "Pause", resume: "Resume", @@ -396,6 +397,7 @@ const en = { manual: "Manual", schedule: "Scheduled", topology: "Topology canvas", + retry_failed: "Retry failed", }, jobScope: { all: "All NEs", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index d7820ae..f8d4d97 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -316,6 +316,7 @@ const zh = { lldpLinks: { title: "LLDP 链路", collectNow: "立即采集", + retryFailed: "一键重试失败", started: "已启动 LLDP 采集", pause: "暂停", resume: "继续", @@ -392,6 +393,7 @@ const zh = { manual: "手动", schedule: "周期", topology: "拓扑画布", + retry_failed: "重试失败", }, jobScope: { all: "全部网元", diff --git a/web/src/pages/network/LldpLinksPage.tsx b/web/src/pages/network/LldpLinksPage.tsx index de88ec1..5e6dec5 100644 --- a/web/src/pages/network/LldpLinksPage.tsx +++ b/web/src/pages/network/LldpLinksPage.tsx @@ -48,11 +48,23 @@ function lldpScopeLabel(t: (k: string) => string, scope: string): string { function lldpItemResultLabel( t: (k: string) => string, - it: { ok?: boolean; parser_stub?: boolean | string | null; unmatched_count?: number; unmatched?: unknown[] }, + it: { + ok?: boolean; + parser_stub?: boolean | string | null; + unmatched_count?: number; + unmatched?: unknown[]; + error?: string | null; + }, ): string { - if (!it.ok) return t("lldpLinks.itemResult.fail"); + const err = String(it.error || "").trim(); + const weak = + Boolean(it.parser_stub) || + err === "parser_stub" || + err === "empty_cli_output" || + err === "vendor_or_device_type_required"; + if (!it.ok || weak) return t("lldpLinks.itemResult.fail"); const unmatchedCount = it.unmatched_count ?? (it.unmatched?.length || 0); - if (it.parser_stub || unmatchedCount > 0) return t("lldpLinks.itemResult.warn"); + if (unmatchedCount > 0) return t("lldpLinks.itemResult.warn"); return t("lldpLinks.itemResult.ok"); } @@ -250,7 +262,7 @@ export function LldpLinksPage() { }); const startMut = useMutation({ - mutationFn: () => startLldpCollect(), + mutationFn: (mode: "full" | "retry_failed" = "full") => startLldpCollect({ mode }), onSuccess: async () => { showOk(t("lldpLinks.started")); await refresh(); @@ -322,10 +334,17 @@ export function LldpLinksPage() { type="button" className="btn-primary" disabled={Boolean(running) || startMut.isPending} - onClick={() => startMut.mutate()} + onClick={() => startMut.mutate("full")} > {t("lldpLinks.collectNow")} + {running?.status === "running" || running?.status === "pending" ? (