diff --git a/netx_api/config.py b/netx_api/config.py
index e69b099..b0024eb 100644
--- a/netx_api/config.py
+++ b/netx_api/config.py
@@ -82,6 +82,9 @@ class Settings(BaseSettings):
lldp_collect_scheduler_enabled: bool = True
lldp_collect_scheduler_tick_sec: int = 60
lldp_collect_startup_grace_sec: int = 3600
+ # Reclaim hung discover jobs (updated_at / created_at older than these).
+ lldp_collect_stale_run_sec: int = 7200
+ lldp_collect_pending_stale_sec: int = 300
# Port traffic monitoring (CLI rate bit/s samples)
port_traffic_scheduler_enabled: bool = True
port_traffic_scheduler_tick_sec: int = 15
diff --git a/netx_api/lldp_collect_schemas.py b/netx_api/lldp_collect_schemas.py
index eeedb6d..39e9e36 100644
--- a/netx_api/lldp_collect_schemas.py
+++ b/netx_api/lldp_collect_schemas.py
@@ -15,21 +15,25 @@ class LldpCollectTargetRef(BaseModel):
class LldpCollectPolicyOut(BaseModel):
enabled: bool = False
- interval_days: int = 1
+ interval_days: int = 1 # legacy view of interval_hours (ceil days)
+ interval_hours: int = 24
concurrency: int = 4
scope_mode: str = "all"
selected_targets: list[LldpCollectTargetRef] = Field(default_factory=list)
auto_add_unmatched: bool = True
+ history_keep: int = 30
updated_at: datetime | None = None
class LldpCollectPolicyUpdate(BaseModel):
enabled: bool | None = None
interval_days: int | None = Field(default=None, ge=1, le=365)
+ interval_hours: int | None = Field(default=None, ge=1, le=8760)
concurrency: int | None = Field(default=None, ge=1, le=32)
scope_mode: str | None = None
selected_targets: list[LldpCollectTargetRef] | None = None
auto_add_unmatched: bool | None = None
+ history_keep: int | None = Field(default=None, ge=0, le=200)
class LldpCollectJobSummary(BaseModel):
@@ -41,7 +45,8 @@ class LldpCollectJobSummary(BaseModel):
done: int = 0
edges_added: int = 0
edges_updated: int = 0
- edges_stale: int = 0
+ edges_stale: int = 0 # legacy alias of edges_missing
+ edges_missing: int = 0
error: str = ""
started_at: datetime | None = None
ended_at: datetime | None = None
@@ -53,7 +58,8 @@ class LldpCollectDashboardOut(BaseModel):
fabric_node_count: int = 0
fabric_edge_count: int = 0
fabric_edge_active: int = 0
- fabric_edge_stale: int = 0
+ fabric_edge_stale: int = 0 # legacy alias of fabric_edge_missing
+ fabric_edge_missing: int = 0
last_discover_at: datetime | None = None
running_job: LldpCollectJobSummary | None = None
last_job: LldpCollectJobSummary | None = None
diff --git a/netx_api/lldp_collect_service.py b/netx_api/lldp_collect_service.py
index 5e0f27c..99330bd 100644
--- a/netx_api/lldp_collect_service.py
+++ b/netx_api/lldp_collect_service.py
@@ -16,15 +16,29 @@ from .lldp_collect_schemas import (
)
from .models import LldpCollectPolicy, TopoDiscoverJob, TopoFabricStats
from .topology_schemas import FabricDiscoverRequest
-from .topology_service import get_discover_job, start_discover_job
+from .topology_service import (
+ get_discover_job,
+ prune_discover_jobs,
+ reclaim_stale_discover_jobs,
+ start_discover_job,
+)
POLICY_ID = 1
+DEFAULT_HISTORY_KEEP = 30
+MAX_INTERVAL_HOURS = 8760 # 365d
def _utcnow() -> datetime:
return datetime.utcnow()
+def _normalize_interval_hours(row: LldpCollectPolicy) -> int:
+ hours = int(getattr(row, "interval_hours", 0) or 0)
+ if hours <= 0:
+ hours = max(1, int(row.interval_days or 1)) * 24
+ return max(1, min(MAX_INTERVAL_HOURS, hours))
+
+
def ensure_policy(db: Session) -> LldpCollectPolicy:
row = db.get(LldpCollectPolicy, POLICY_ID)
if row is None:
@@ -32,15 +46,22 @@ def ensure_policy(db: Session) -> LldpCollectPolicy:
id=POLICY_ID,
enabled=False,
interval_days=1,
+ interval_hours=24,
concurrency=4,
scope_mode="all",
selected_targets=[],
auto_add_unmatched=True,
+ history_keep=DEFAULT_HISTORY_KEEP,
updated_at=_utcnow(),
)
db.add(row)
db.commit()
db.refresh(row)
+ # Heal legacy rows missing hours.
+ if int(getattr(row, "interval_hours", 0) or 0) <= 0:
+ row.interval_hours = max(1, int(row.interval_days or 1)) * 24
+ db.commit()
+ db.refresh(row)
return row
@@ -56,13 +77,20 @@ def _policy_out(row: LldpCollectPolicy) -> LldpCollectPolicyOut:
if src not in {"managed", "ume"}:
src = "managed"
refs.append(LldpCollectTargetRef(source=src, id=tid))
+ keep = getattr(row, "history_keep", None)
+ if keep is None:
+ keep = DEFAULT_HISTORY_KEEP
+ hours = _normalize_interval_hours(row)
+ days = max(1, min(365, (hours + 23) // 24))
return LldpCollectPolicyOut(
enabled=bool(row.enabled),
- interval_days=int(row.interval_days or 1),
+ interval_days=days,
+ interval_hours=hours,
concurrency=int(row.concurrency or 4),
scope_mode="selected" if str(row.scope_mode or "") == "selected" else "all",
selected_targets=refs,
auto_add_unmatched=bool(row.auto_add_unmatched),
+ history_keep=max(0, min(200, int(keep))),
updated_at=row.updated_at,
)
@@ -76,8 +104,14 @@ def update_policy(db: Session, body: LldpCollectPolicyUpdate) -> LldpCollectPoli
data = body.model_dump(exclude_unset=True)
if "enabled" in data and data["enabled"] is not None:
row.enabled = bool(data["enabled"])
- if "interval_days" in data and data["interval_days"] is not None:
- row.interval_days = max(1, min(365, int(data["interval_days"])))
+ if "interval_hours" in data and data["interval_hours"] is not None:
+ hours = max(1, min(MAX_INTERVAL_HOURS, int(data["interval_hours"])))
+ row.interval_hours = hours
+ row.interval_days = max(1, min(365, (hours + 23) // 24))
+ elif "interval_days" in data and data["interval_days"] is not None:
+ days = max(1, min(365, int(data["interval_days"])))
+ row.interval_days = days
+ row.interval_hours = days * 24
if "concurrency" in data and data["concurrency"] is not None:
row.concurrency = max(1, min(32, int(data["concurrency"])))
if "scope_mode" in data and data["scope_mode"] is not None:
@@ -104,9 +138,12 @@ def update_policy(db: Session, body: LldpCollectPolicyUpdate) -> LldpCollectPoli
row.selected_targets = cleaned
if "auto_add_unmatched" in data and data["auto_add_unmatched"] is not None:
row.auto_add_unmatched = bool(data["auto_add_unmatched"])
+ if "history_keep" in data and data["history_keep"] is not None:
+ row.history_keep = max(0, min(200, int(data["history_keep"])))
row.updated_at = _utcnow()
db.commit()
db.refresh(row)
+ prune_discover_jobs(db, keep=int(row.history_keep or 0))
return _policy_out(row)
@@ -123,6 +160,7 @@ def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None:
edges_added=int(job.edges_added or 0),
edges_updated=int(job.edges_updated or 0),
edges_stale=int(job.edges_stale or 0),
+ edges_missing=int(job.edges_stale or 0),
error=job.error or "",
started_at=job.started_at,
ended_at=job.ended_at,
@@ -131,6 +169,7 @@ def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None:
def has_running_job(db: Session) -> TopoDiscoverJob | None:
+ reclaim_stale_discover_jobs(db)
return (
db.query(TopoDiscoverJob)
.filter(TopoDiscoverJob.status.in_(["pending", "running"]))
@@ -149,41 +188,57 @@ def last_finished_job(db: Session) -> TopoDiscoverJob | None:
def next_due_at(db: Session, policy: LldpCollectPolicy) -> datetime | None:
+ """Due time based on last *scheduled* successful collect only (manual must not reset)."""
if not policy.enabled:
return None
- days = max(1, int(policy.interval_days or 1))
+ hours = _normalize_interval_hours(policy)
last = (
db.query(TopoDiscoverJob)
- .filter(TopoDiscoverJob.status == "done", TopoDiscoverJob.ended_at.isnot(None))
+ .filter(
+ TopoDiscoverJob.status == "done",
+ TopoDiscoverJob.trigger_mode == "schedule",
+ TopoDiscoverJob.ended_at.isnot(None),
+ )
.order_by(TopoDiscoverJob.ended_at.desc())
.first()
)
if last is None or last.ended_at is None:
return _utcnow()
- return last.ended_at + timedelta(days=days)
+ return last.ended_at + timedelta(hours=hours)
def build_discover_request(policy: LldpCollectPolicy) -> FabricDiscoverRequest:
concurrency = max(1, min(32, int(policy.concurrency or 4)))
auto_add = bool(policy.auto_add_unmatched)
if str(policy.scope_mode or "") == "selected":
- ne_ids: list[str] = []
+ managed_ids: list[str] = []
+ ume_ids: list[str] = []
for raw in policy.selected_targets or []:
- if isinstance(raw, dict):
- tid = str(raw.get("id") or "").strip()
- if tid:
- ne_ids.append(tid)
- if not ne_ids:
+ if not isinstance(raw, dict):
+ continue
+ tid = str(raw.get("id") or "").strip()
+ if not tid:
+ continue
+ src = str(raw.get("source") or "managed").strip().lower() or "managed"
+ if src == "ume":
+ ume_ids.append(tid)
+ else:
+ managed_ids.append(tid)
+ if not managed_ids and not ume_ids:
raise HTTPException(status_code=400, detail="no_selected_targets")
return FabricDiscoverRequest(
scope="ne_ids",
- ne_ids=ne_ids,
+ ne_ids=[],
+ managed_ne_ids=managed_ids,
+ ume_ne_ids=ume_ids,
concurrency=concurrency,
auto_add_unmatched=auto_add,
)
return FabricDiscoverRequest(
scope="all_inventory",
ne_ids=[],
+ managed_ne_ids=[],
+ ume_ne_ids=[],
concurrency=concurrency,
auto_add_unmatched=auto_add,
)
@@ -195,6 +250,7 @@ def start_collect(db: Session, *, trigger_mode: str = "manual") -> dict:
policy = ensure_policy(db)
body = build_discover_request(policy)
job = start_discover_job(db, body, trigger_mode=trigger_mode)
+ prune_discover_jobs(db, keep=int(getattr(policy, "history_keep", DEFAULT_HISTORY_KEEP) or 0))
return {"ok": True, "job": job.model_dump()}
@@ -209,6 +265,7 @@ def get_dashboard(db: Session) -> LldpCollectDashboardOut:
fabric_edge_count=int(stats.edge_count if stats else 0),
fabric_edge_active=int(stats.edge_active if stats else 0),
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),
diff --git a/netx_api/main.py b/netx_api/main.py
index 51b1e3d..dd5bc43 100644
--- a/netx_api/main.py
+++ b/netx_api/main.py
@@ -824,6 +824,20 @@ def on_startup() -> None:
_schedule_log.exception("startup: auth/port_traffic/topology schema migration failed")
_reset_runtime_pause_flags()
_fail_stale_running_sync_jobs_on_startup()
+ try:
+ from .topology_service import reclaim_stale_discover_jobs
+
+ db_topo = SessionLocal()
+ try:
+ closed = reclaim_stale_discover_jobs(db_topo, force_all_open=True)
+ if closed:
+ _schedule_log.warning(
+ "startup: closed %s orphaned topology discover jobs", closed
+ )
+ finally:
+ db_topo.close()
+ except Exception:
+ _schedule_log.exception("startup: topology discover job cleanup failed")
if _needs_startup_alarm_sync_before_ws():
begin_startup_alarm_sync_gate()
_schedule_log.info(
diff --git a/netx_api/models.py b/netx_api/models.py
index 9abdc9a..1d0693d 100644
--- a/netx_api/models.py
+++ b/netx_api/models.py
@@ -503,11 +503,14 @@ class LldpCollectPolicy(Base):
id: Mapped[int] = mapped_column(Integer, primary_key=True, default=1)
enabled: Mapped[bool] = mapped_column(Boolean, default=False)
- interval_days: Mapped[int] = mapped_column(Integer, default=1)
+ interval_days: Mapped[int] = mapped_column(Integer, default=1) # legacy; prefer interval_hours
+ interval_hours: Mapped[int] = mapped_column(Integer, default=24)
concurrency: Mapped[int] = mapped_column(Integer, default=4)
scope_mode: Mapped[str] = mapped_column(String(32), default="all") # all | selected
selected_targets: Mapped[list] = mapped_column(_JsonType, default=list)
auto_add_unmatched: Mapped[bool] = mapped_column(Boolean, default=True)
+ # Keep N finished discover jobs (items + raw_preview); 0 = keep none finished.
+ history_keep: Mapped[int] = mapped_column(Integer, default=30)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
diff --git a/netx_api/topology_migrate.py b/netx_api/topology_migrate.py
index 492659c..0e3cbc2 100644
--- a/netx_api/topology_migrate.py
+++ b/netx_api/topology_migrate.py
@@ -31,13 +31,21 @@ def ensure_topology_schema(conn: Connection) -> None:
# Best-effort column add for existing DBs (create_all won't alter).
alter_stmts: list[str] = []
if dialect.startswith("postgres"):
- alter_stmts.append(
- "ALTER TABLE topo_discover_job ADD COLUMN IF NOT EXISTS trigger_mode VARCHAR(32) DEFAULT 'manual'"
+ alter_stmts.extend(
+ [
+ "ALTER TABLE topo_discover_job ADD COLUMN IF NOT EXISTS trigger_mode VARCHAR(32) DEFAULT 'manual'",
+ "ALTER TABLE lldp_collect_policy ADD COLUMN IF NOT EXISTS history_keep INTEGER DEFAULT 30",
+ "ALTER TABLE lldp_collect_policy ADD COLUMN IF NOT EXISTS interval_hours INTEGER DEFAULT 24",
+ ]
)
elif dialect.startswith("sqlite"):
# SQLite: ignore if column already exists.
- alter_stmts.append(
- "ALTER TABLE topo_discover_job ADD COLUMN trigger_mode VARCHAR(32) DEFAULT 'manual'"
+ alter_stmts.extend(
+ [
+ "ALTER TABLE topo_discover_job ADD COLUMN trigger_mode VARCHAR(32) DEFAULT 'manual'",
+ "ALTER TABLE lldp_collect_policy ADD COLUMN history_keep INTEGER DEFAULT 30",
+ "ALTER TABLE lldp_collect_policy ADD COLUMN interval_hours INTEGER DEFAULT 24",
+ ]
)
for sql in alter_stmts:
try:
@@ -45,6 +53,30 @@ def ensure_topology_schema(conn: Connection) -> None:
except Exception:
_log.debug("topology alter skipped/failed: %s", sql[:80], exc_info=True)
+ # Backfill hours from legacy days when hours unset/zero.
+ try:
+ conn.execute(
+ text(
+ "UPDATE lldp_collect_policy SET interval_hours = "
+ "CASE WHEN COALESCE(interval_hours, 0) <= 0 "
+ "THEN GREATEST(1, COALESCE(interval_days, 1)) * 24 "
+ "ELSE interval_hours END"
+ )
+ )
+ except Exception:
+ try:
+ # SQLite has no GREATEST in older builds — use MAX.
+ conn.execute(
+ text(
+ "UPDATE lldp_collect_policy SET interval_hours = "
+ "CASE WHEN COALESCE(interval_hours, 0) <= 0 "
+ "THEN MAX(1, COALESCE(interval_days, 1)) * 24 "
+ "ELSE interval_hours END"
+ )
+ )
+ except Exception:
+ _log.debug("interval_hours backfill skipped", exc_info=True)
+
if not dialect.startswith("postgres"):
return
stmts = [
diff --git a/netx_api/topology_router.py b/netx_api/topology_router.py
index 5881482..2ae0853 100644
--- a/netx_api/topology_router.py
+++ b/netx_api/topology_router.py
@@ -66,6 +66,7 @@ def api_fabric_edges(
layer: str = "physical",
status: str = "",
source: str = "",
+ keyword: str = "",
page: int = Query(default=1, ge=1),
page_size: int = Query(default=100, ge=1, le=2000),
db: Session = Depends(get_db),
@@ -76,6 +77,7 @@ def api_fabric_edges(
layer=layer,
status=status,
source=source,
+ keyword=keyword,
page=page,
page_size=page_size,
)
diff --git a/netx_api/topology_schemas.py b/netx_api/topology_schemas.py
index 1bfb89d..c8dadbb 100644
--- a/netx_api/topology_schemas.py
+++ b/netx_api/topology_schemas.py
@@ -32,18 +32,24 @@ class FabricEdgeOut(BaseModel):
b_node_id: str
a_port: str = ""
b_port: str = ""
+ a_name: str = ""
+ b_name: str = ""
+ a_ip: str = ""
+ b_ip: str = ""
source: str = "lldp"
status: str = "active"
attrs: dict[str, Any] = Field(default_factory=dict)
discovered_at: datetime | None = None
last_seen_at: datetime | None = None
+ updated_at: datetime | None = None
class FabricSummaryOut(BaseModel):
node_count: int = 0
edge_count: int = 0
edge_active: int = 0
- edge_stale: int = 0
+ edge_stale: int = 0 # legacy alias of edge_missing
+ edge_missing: int = 0
last_discover_at: datetime | None = None
updated_at: datetime | None = None
@@ -59,7 +65,12 @@ class FabricDiscoverRequest(BaseModel):
"""Start LLDP discovery into fabric (no CDP)."""
scope: str = Field(default="ne_ids", description="all_inventory | ne_ids")
- ne_ids: list[str] = Field(default_factory=list)
+ ne_ids: list[str] = Field(
+ default_factory=list,
+ description="Legacy mixed ids (managed first, then ume). Prefer managed_ne_ids/ume_ne_ids.",
+ )
+ managed_ne_ids: list[str] = Field(default_factory=list)
+ ume_ne_ids: list[str] = Field(default_factory=list)
auto_add_unmatched: bool = Field(
default=True,
description="Create SSH placeholder ManagedNEs for LLDP neighbors not in inventory",
@@ -105,7 +116,8 @@ class FabricDiscoverJobOut(BaseModel):
done: int = 0
edges_added: int = 0
edges_updated: int = 0
- edges_stale: int = 0
+ edges_stale: int = 0 # legacy alias of edges_missing
+ edges_missing: int = 0
error: str = ""
started_at: datetime | None = None
ended_at: datetime | None = None
diff --git a/netx_api/topology_service.py b/netx_api/topology_service.py
index 1de375b..aff09f7 100644
--- a/netx_api/topology_service.py
+++ b/netx_api/topology_service.py
@@ -5,7 +5,7 @@ from __future__ import annotations
import re
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
-from datetime import datetime
+from datetime import datetime, timedelta
from typing import Any
from uuid import uuid4
@@ -15,9 +15,11 @@ from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from .cli_resolve import get_default_profile, infer_device_type_vendor
+from .config import settings
from .db import SessionLocal
from .device_types import LLDP_DISCOVERED_NE_SOURCE, WEBCRT_NE_SOURCE
from .models import (
+ LldpCollectPolicy,
ManagedNE,
TopoDiscoverJob,
TopoDiscoverJobItem,
@@ -164,10 +166,16 @@ def _node_out(n: TopoFabricNode) -> FabricNodeOut:
)
-def _edge_out(e: TopoFabricEdge) -> FabricEdgeOut:
+def _edge_out(
+ e: TopoFabricEdge,
+ *,
+ nodes_by_id: dict[str, TopoFabricNode] | None = None,
+) -> FabricEdgeOut:
src = str(e.source or "lldp").strip().lower() or "lldp"
if src == "stale":
src = "lldp"
+ a_node = (nodes_by_id or {}).get(e.a_node_id)
+ b_node = (nodes_by_id or {}).get(e.b_node_id)
return FabricEdgeOut(
id=e.id,
layer=e.layer or "physical",
@@ -175,14 +183,27 @@ def _edge_out(e: TopoFabricEdge) -> FabricEdgeOut:
b_node_id=e.b_node_id,
a_port=e.a_port or "",
b_port=e.b_port or "",
+ a_name=(a_node.name if a_node else "") or "",
+ b_name=(b_node.name if b_node else "") or "",
+ a_ip=(a_node.ip if a_node else "") or "",
+ b_ip=(b_node.ip if b_node else "") or "",
source=src,
status=_normalize_edge_status(e.status or "active"),
attrs=dict(e.attrs or {}),
discovered_at=e.discovered_at,
last_seen_at=e.last_seen_at,
+ updated_at=e.updated_at,
)
+def _nodes_by_ids(db: Session, ids: set[str]) -> dict[str, TopoFabricNode]:
+ clean = {str(i).strip() for i in ids if str(i or "").strip()}
+ if not clean:
+ return {}
+ rows = db.query(TopoFabricNode).filter(TopoFabricNode.id.in_(list(clean))).all()
+ return {r.id: r for r in rows}
+
+
def _normalize_endpoints(
a_id: str, b_id: str, a_port: str, b_port: str
) -> tuple[str, str, str, str]:
@@ -315,6 +336,7 @@ def get_fabric_summary(db: Session) -> FabricSummaryOut:
edge_count=row.edge_count,
edge_active=row.edge_active,
edge_stale=row.edge_stale,
+ edge_missing=row.edge_stale,
last_discover_at=row.last_discover_at,
updated_at=row.updated_at,
)
@@ -363,6 +385,7 @@ def list_fabric_edges(
layer: str = "physical",
status: str = "",
source: str = "",
+ keyword: str = "",
page: int = 1,
page_size: int = PAGE_DEFAULT,
) -> dict[str, Any]:
@@ -375,11 +398,33 @@ def list_fabric_edges(
if nid:
q = q.filter(or_(TopoFabricEdge.a_node_id == nid, TopoFabricEdge.b_node_id == nid))
st = str(status or "").strip().lower()
- if st:
+ if st in _EDGE_STATUS_MISSING_COMPAT:
+ q = q.filter(TopoFabricEdge.status.in_(list(_EDGE_STATUS_MISSING_COMPAT)))
+ elif st:
q = q.filter(TopoFabricEdge.status == st)
src = str(source or "").strip().lower()
if src:
+ if src == "stale":
+ src = "lldp"
q = q.filter(TopoFabricEdge.source == src)
+ kw = str(keyword or "").strip()
+ if kw:
+ like = f"%{kw}%"
+ matched_ids = [
+ r.id
+ for r in db.query(TopoFabricNode.id)
+ .filter(or_(TopoFabricNode.name.ilike(like), TopoFabricNode.ip.ilike(like)))
+ .limit(2000)
+ .all()
+ ]
+ if not matched_ids:
+ return {"total": 0, "page": page, "page_size": page_size, "items": []}
+ q = q.filter(
+ or_(
+ TopoFabricEdge.a_node_id.in_(matched_ids),
+ TopoFabricEdge.b_node_id.in_(matched_ids),
+ )
+ )
total = int(q.count())
rows = (
q.order_by(TopoFabricEdge.updated_at.desc())
@@ -387,11 +432,12 @@ def list_fabric_edges(
.limit(page_size)
.all()
)
+ node_map = _nodes_by_ids(db, {e.a_node_id for e in rows} | {e.b_node_id for e in rows})
return {
"total": total,
"page": page,
"page_size": page_size,
- "items": [_edge_out(e).model_dump() for e in rows],
+ "items": [_edge_out(e, nodes_by_id=node_map).model_dump() for e in rows],
}
@@ -1105,7 +1151,11 @@ def _pick_managed_ne(
def ensure_lldp_discovered_managed_ne(
- db: Session, *, remote_name: str = "", remote_ip: str = ""
+ db: Session,
+ *,
+ remote_name: str = "",
+ remote_ip: str = "",
+ placeholder_by_name: dict[str, ManagedNE] | None = None,
) -> ManagedNE:
"""SSH placeholder ManagedNE for an LLDP neighbor not in inventory.
@@ -1117,18 +1167,24 @@ def ensure_lldp_discovered_managed_ne(
ip_hint = str(remote_ip or "").strip()[:128]
now = _utcnow()
- # Reuse existing LLDP placeholder by normalized hostname.
- if name_key:
+ cache = placeholder_by_name
+ if cache is None:
+ cache = {}
for ne in (
db.query(ManagedNE)
.filter(ManagedNE.source == LLDP_DISCOVERED_NE_SOURCE)
.all()
):
- if _norm_host(ne.name or "") == name_key:
- if ip_hint and not str(ne.source_ref or "").strip():
- ne.source_ref = ip_hint
- ne.updated_at = now
- return ne
+ nk = _norm_host(ne.name or "")
+ if nk and nk not in cache:
+ cache[nk] = ne
+
+ if name_key and name_key in cache:
+ ne = cache[name_key]
+ if ip_hint and not str(ne.source_ref or "").strip():
+ ne.source_ref = ip_hint
+ ne.updated_at = now
+ return ne
row = ManagedNE(
id=uuid4().hex,
@@ -1151,42 +1207,106 @@ def ensure_lldp_discovered_managed_ne(
)
db.add(row)
db.flush()
+ if name_key:
+ cache[name_key] = row
+ if placeholder_by_name is not None and name_key:
+ placeholder_by_name[name_key] = row
return row
+class _FabricPeerIndex:
+ """In-memory name/IP index for one discover target (avoids O(nodes) per neighbor)."""
+
+ def __init__(self, db: Session, self_id: str) -> None:
+ self.db = db
+ self.self_id = self_id
+ self.by_ip: dict[str, list[TopoFabricNode]] = {}
+ self.by_name: dict[str, list[TopoFabricNode]] = {}
+ self.placeholder_by_name: dict[str, ManagedNE] = {}
+ for n in db.query(TopoFabricNode).filter(TopoFabricNode.id != self_id).all():
+ ip = str(n.ip or "").strip()
+ if ip:
+ self.by_ip.setdefault(ip, []).append(n)
+ nk = _norm_host(n.name or "")
+ if nk:
+ self.by_name.setdefault(nk, []).append(n)
+ for ne in (
+ db.query(ManagedNE).filter(ManagedNE.source == LLDP_DISCOVERED_NE_SOURCE).all()
+ ):
+ nk = _norm_host(ne.name or "")
+ if nk and nk not in self.placeholder_by_name:
+ self.placeholder_by_name[nk] = ne
+
+ def _best(self, matched: list[TopoFabricNode]) -> TopoFabricNode:
+ matched.sort(key=lambda n: _fabric_match_score(self.db, n), reverse=True)
+ return matched[0]
+
+ def match(self, hit: NeighborHit) -> TopoFabricNode | None:
+ name_key = _norm_host(hit.remote_name)
+ ip_key = str(hit.remote_ip or "").strip()
+ matched: list[TopoFabricNode] = []
+ if ip_key:
+ matched.extend(self.by_ip.get(ip_key) or [])
+ if name_key:
+ for n in self.by_name.get(name_key) or []:
+ if n not in matched:
+ matched.append(n)
+ if matched:
+ return self._best(matched)
+
+ if ip_key:
+ ne = _pick_managed_ne(self.db, ip=ip_key)
+ if ne is not None:
+ node = ensure_fabric_node_for_managed(self.db, ne)
+ self._remember(node)
+ return node
+ ume = (
+ self.db.query(UmeInventoryNE)
+ .filter(UmeInventoryNE.ip_address == ip_key)
+ .first()
+ )
+ if ume is not None:
+ node = ensure_fabric_node_for_ume(self.db, ume)
+ self._remember(node)
+ return node
+ if name_key:
+ ne = _pick_managed_ne(self.db, name_key=name_key)
+ if ne is not None:
+ node = ensure_fabric_node_for_managed(self.db, ne)
+ self._remember(node)
+ return node
+ return None
+
+ def _remember(self, node: TopoFabricNode) -> None:
+ if not node or node.id == self.self_id:
+ return
+ ip = str(node.ip or "").strip()
+ if ip:
+ bucket = self.by_ip.setdefault(ip, [])
+ if node not in bucket:
+ bucket.append(node)
+ nk = _norm_host(node.name or "")
+ if nk:
+ bucket = self.by_name.setdefault(nk, [])
+ if node not in bucket:
+ bucket.append(node)
+
+ def ensure_placeholder(self, *, remote_name: str, remote_ip: str) -> TopoFabricNode:
+ placeholder = ensure_lldp_discovered_managed_ne(
+ self.db,
+ remote_name=remote_name,
+ remote_ip=remote_ip,
+ placeholder_by_name=self.placeholder_by_name,
+ )
+ peer = ensure_fabric_node_for_managed(self.db, placeholder)
+ self._remember(peer)
+ return peer
+
+
def _match_hit_to_fabric_node(
db: Session, hit: NeighborHit, *, self_id: str
) -> TopoFabricNode | None:
- name_key = _norm_host(hit.remote_name)
- ip_key = str(hit.remote_ip or "").strip()
- candidates = db.query(TopoFabricNode).filter(TopoFabricNode.id != self_id).all()
-
- matched: list[TopoFabricNode] = []
- for n in candidates:
- names = {_norm_host(n.name or "")}
- ips = {str(n.ip or "").strip()}
- if ip_key and ip_key in ips:
- matched.append(n)
- continue
- if name_key and name_key in names:
- matched.append(n)
- if matched:
- matched.sort(key=lambda n: _fabric_match_score(db, n), reverse=True)
- return matched[0]
-
- # Inventory not yet in fabric (or missed due to concurrent insert).
- if ip_key:
- ne = _pick_managed_ne(db, ip=ip_key)
- if ne is not None:
- return ensure_fabric_node_for_managed(db, ne)
- ume = db.query(UmeInventoryNE).filter(UmeInventoryNE.ip_address == ip_key).first()
- if ume is not None:
- return ensure_fabric_node_for_ume(db, ume)
- if name_key:
- ne = _pick_managed_ne(db, name_key=name_key)
- if ne is not None:
- return ensure_fabric_node_for_managed(db, ne)
- return None
+ return _FabricPeerIndex(db, self_id).match(hit)
def _retarget_fabric_edges(db: Session, *, from_id: str, to_id: str) -> None:
@@ -1247,6 +1367,7 @@ def merge_duplicate_fabric_nodes(db: Session) -> dict[str, int]:
"""Collapse duplicate fabric nodes (same managed/ume/name/ip) onto inventory canonicals."""
nodes = db.query(TopoFabricNode).order_by(TopoFabricNode.created_at.asc()).all()
merged = 0
+ placeholders_removed = 0
# 1) Same managed_ne_id / ume_ne_id (constraint may be missing on old DBs).
by_managed: dict[str, list[TopoFabricNode]] = {}
@@ -1320,6 +1441,53 @@ def merge_duplicate_fabric_nodes(db: Session) -> dict[str, int]:
continue
_absorb(canon, [o])
+ # 2b) LLDP placeholders (score=2) → real inventory (score>=3) by hostname / seen mgmt IP.
+ # Placeholders have managed_ne_id so they are NOT orphans; absorb + drop empty ManagedNE.
+ db.flush()
+ nodes = db.query(TopoFabricNode).all()
+ reals = [n for n in nodes if _fabric_match_score(db, n) >= 3]
+ placeholders = [n for n in nodes if _fabric_match_score(db, n) == 2]
+ real_by_name: dict[str, TopoFabricNode] = {}
+ real_by_ip: dict[str, TopoFabricNode] = {}
+ for n in sorted(reals, key=lambda x: _fabric_match_score(db, x), reverse=True):
+ nk = _norm_host(n.name or "")
+ if nk and nk not in real_by_name:
+ real_by_name[nk] = n
+ ip = str(n.ip or "").strip()
+ if ip and ip not in real_by_ip:
+ real_by_ip[ip] = n
+ for p in placeholders:
+ if db.get(TopoFabricNode, p.id) is None:
+ continue
+ canon = None
+ nk = _norm_host(p.name or "")
+ ip = str(p.ip or "").strip()
+ seen_ip = ""
+ mid = str(p.managed_ne_id or "").strip()
+ ph_ne = db.get(ManagedNE, mid) if mid else None
+ if ph_ne is not None:
+ seen_ip = str(ph_ne.source_ref or "").strip()
+ if ip and ip in real_by_ip:
+ canon = real_by_ip[ip]
+ elif seen_ip and seen_ip in real_by_ip:
+ canon = real_by_ip[seen_ip]
+ elif nk and nk in real_by_name:
+ canon = real_by_name[nk]
+ if canon is None or canon.id == p.id:
+ continue
+ _absorb(canon, [p])
+ db.flush()
+ # Drop placeholder ManagedNE if nothing else references it.
+ if ph_ne is not None and str(ph_ne.source or "").strip().lower() == LLDP_DISCOVERED_NE_SOURCE:
+ still = (
+ db.query(TopoFabricNode)
+ .filter(TopoFabricNode.managed_ne_id == ph_ne.id)
+ .count()
+ )
+ if still == 0:
+ db.delete(ph_ne)
+ placeholders_removed += 1
+
# 3) WebCRT session hosts sharing an IP with a real inventory fabric node.
db.flush()
nodes = db.query(TopoFabricNode).all()
@@ -1338,10 +1506,10 @@ def merge_duplicate_fabric_nodes(db: Session) -> dict[str, int]:
canon = max(real, key=lambda n: _fabric_match_score(db, n))
_absorb(canon, webcrtish)
- if merged:
+ if merged or placeholders_removed:
db.commit()
refresh_fabric_stats(db)
- return {"merged": merged}
+ return {"merged": merged, "placeholders_removed": placeholders_removed}
def _raw_preview(raw: str, *, limit: int = _RAW_PREVIEW_MAX) -> str:
@@ -1396,6 +1564,7 @@ def _job_out(db: Session, job: TopoDiscoverJob, *, include_items: bool = True) -
edges_added=int(job.edges_added or 0),
edges_updated=int(job.edges_updated or 0),
edges_stale=int(job.edges_stale or 0),
+ edges_missing=int(job.edges_stale or 0),
error=job.error or "",
started_at=job.started_at,
ended_at=job.ended_at,
@@ -1410,6 +1579,25 @@ def get_discover_job(db: Session, job_id: str) -> FabricDiscoverJobOut:
return _job_out(db, job)
+def _ume_target_dict(db: Session, uid: str, default_profile: Any) -> dict[str, str] | None:
+ ume = db.query(UmeInventoryNE).filter(UmeInventoryNE.ne_id == uid).one_or_none()
+ if ume is None:
+ return None
+ if default_profile is not None:
+ dtype, vendor = infer_device_type_vendor(str(ume.ne_type or ""), default_profile)
+ else:
+ dtype, vendor = "zte_zxros", (ume.vendor or "ZTE")
+ name = (ume.host_name or ume.ne_name or ume.user_label or ume.ip_address or uid).strip()
+ return {
+ "ne_id": uid,
+ "ume_ne_id": uid,
+ "ne_name": name,
+ "ne_ip": ume.ip_address or "",
+ "vendor": vendor or (ume.vendor or "ZTE"),
+ "device_type": dtype or "zte_zxros",
+ }
+
+
def _resolve_scan_targets(
db: Session, body: FabricDiscoverRequest
) -> list[dict[str, str]]:
@@ -1429,6 +1617,42 @@ def _resolve_scan_targets(
}
)
return targets
+
+ managed_ids = [str(x).strip() for x in (body.managed_ne_ids or []) if str(x).strip()]
+ ume_ids = [str(x).strip() for x in (body.ume_ne_ids or []) if str(x).strip()]
+ if managed_ids or ume_ids:
+ seen: set[str] = set()
+ for mid in managed_ids:
+ if mid in seen:
+ continue
+ ne = db.get(ManagedNE, mid)
+ if ne is None:
+ continue
+ seen.add(mid)
+ targets.append(
+ {
+ "ne_id": ne.id,
+ "ume_ne_id": "",
+ "ne_name": ne.name or "",
+ "ne_ip": ne.ip_address or "",
+ "vendor": ne.vendor or "",
+ "device_type": ne.device_type or "",
+ }
+ )
+ for uid in ume_ids:
+ key = f"ume:{uid}"
+ if key in seen:
+ continue
+ row = _ume_target_dict(db, uid, default_profile)
+ if row is None:
+ continue
+ seen.add(key)
+ targets.append(row)
+ if not targets:
+ raise HTTPException(status_code=400, detail="ne_ids_required")
+ return targets
+
+ # Legacy mixed ne_ids: prefer ManagedNE, leftover treated as UME.
filter_ids = {str(x).strip() for x in (body.ne_ids or []) if str(x).strip()}
if not filter_ids:
raise HTTPException(status_code=400, detail="ne_ids_required")
@@ -1447,27 +1671,36 @@ def _resolve_scan_targets(
)
filter_ids.discard(mid)
for uid in list(filter_ids):
- ume = db.query(UmeInventoryNE).filter(UmeInventoryNE.ne_id == uid).one_or_none()
- if ume is None:
- continue
- if default_profile is not None:
- dtype, vendor = infer_device_type_vendor(str(ume.ne_type or ""), default_profile)
- else:
- dtype, vendor = "zte_zxros", (ume.vendor or "ZTE")
- name = (ume.host_name or ume.ne_name or ume.user_label or ume.ip_address or uid).strip()
- targets.append(
- {
- "ne_id": uid,
- "ume_ne_id": uid,
- "ne_name": name,
- "ne_ip": ume.ip_address or "",
- "vendor": vendor or (ume.vendor or "ZTE"),
- "device_type": dtype or "zte_zxros",
- }
- )
+ row = _ume_target_dict(db, uid, default_profile)
+ if row is not None:
+ targets.append(row)
return targets
+def prune_discover_jobs(db: Session, *, keep: int = 30) -> int:
+ """Delete finished discover jobs beyond ``keep`` (newest kept). Open jobs always retained."""
+ keep = max(0, min(200, int(keep)))
+ finished = (
+ db.query(TopoDiscoverJob)
+ .filter(TopoDiscoverJob.status.in_(["done", "failed"]))
+ .order_by(TopoDiscoverJob.created_at.desc())
+ .all()
+ )
+ to_drop = finished if keep == 0 else finished[keep:]
+ if not to_drop:
+ return 0
+ dropped = 0
+ for job in to_drop:
+ db.query(TopoDiscoverJobItem).filter(TopoDiscoverJobItem.job_id == job.id).delete(
+ synchronize_session=False
+ )
+ db.delete(job)
+ dropped += 1
+ if dropped:
+ db.commit()
+ return dropped
+
+
def _discover_one_target(
target: dict[str, str],
*,
@@ -1558,17 +1791,16 @@ def _discover_one_target(
unmatched: list[dict[str, str]] = []
touched: list[str] = []
replaced: list[str] = []
+ peer_index = _FabricPeerIndex(db, fabric_node.id)
for hit in hits:
- peer = _match_hit_to_fabric_node(db, hit, self_id=fabric_node.id)
+ peer = peer_index.match(hit)
if peer is None:
if auto_add_unmatched and (hit.remote_name or hit.remote_ip):
# Not in inventory → SSH placeholder ManagedNE (empty IP/creds).
- placeholder = ensure_lldp_discovered_managed_ne(
- db,
+ peer = peer_index.ensure_placeholder(
remote_name=(hit.remote_name or "").strip(),
remote_ip=(hit.remote_ip or "").strip(),
)
- peer = ensure_fabric_node_for_managed(db, placeholder)
peer.attrs = dict(peer.attrs or {})
peer.attrs["from_lldp_unmatched"] = True
peer.last_seen_at = now
@@ -1751,6 +1983,13 @@ def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None:
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
+ pass
except Exception as exc: # noqa: BLE001
db.rollback()
job = db.get(TopoDiscoverJob, job_id)
@@ -1766,12 +2005,87 @@ def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None:
_RUNNING_JOBS.discard(job_id)
+def reclaim_stale_discover_jobs(
+ db: Session,
+ *,
+ force_all_open: bool = False,
+ now: datetime | None = None,
+) -> int:
+ """Mark orphaned / hung discover jobs as failed so scheduling can proceed.
+
+ - ``force_all_open``: process restart — all pending/running rows are dead.
+ - Otherwise: pending older than pending_stale_sec, or running with stale updated_at.
+ """
+ now = now or _utcnow()
+ run_sec = max(60, int(getattr(settings, "lldp_collect_stale_run_sec", 7200) or 7200))
+ pend_sec = max(30, int(getattr(settings, "lldp_collect_pending_stale_sec", 300) or 300))
+ open_jobs = (
+ db.query(TopoDiscoverJob)
+ .filter(TopoDiscoverJob.status.in_(["pending", "running"]))
+ .all()
+ )
+ if not open_jobs:
+ return 0
+ closed = 0
+ for job in open_jobs:
+ status = str(job.status or "")
+ if force_all_open:
+ reason = "stale_running_reset_on_startup"
+ elif status == "pending":
+ created = job.created_at or job.updated_at or now
+ if created > now - timedelta(seconds=pend_sec):
+ continue
+ reason = "pending_stale_timeout"
+ else:
+ touched = job.updated_at or job.started_at or job.created_at or now
+ if touched > now - timedelta(seconds=run_sec):
+ continue
+ reason = "running_stale_timeout"
+ job.status = "failed"
+ job.ended_at = now
+ job.updated_at = now
+ msg = str(job.error or "").strip()
+ job.error = (msg + ("; " if msg else "") + reason)[:1024]
+ closed += 1
+ with _JOB_LOCK:
+ _RUNNING_JOBS.discard(job.id)
+ if closed:
+ db.commit()
+ return closed
+
+
def start_discover_job(
db: Session,
body: FabricDiscoverRequest,
*,
trigger_mode: str = "manual",
) -> FabricDiscoverJobOut:
+ reclaim_stale_discover_jobs(db)
+ # Serialize multi-worker starts via singleton policy row lock (PG/SQLite FOR UPDATE).
+ pol = db.get(LldpCollectPolicy, 1)
+ if pol is None:
+ pol = LldpCollectPolicy(
+ id=1,
+ enabled=False,
+ interval_days=1,
+ interval_hours=24,
+ concurrency=4,
+ scope_mode="all",
+ selected_targets=[],
+ auto_add_unmatched=True,
+ history_keep=30,
+ updated_at=_utcnow(),
+ )
+ db.add(pol)
+ db.commit()
+ db.query(LldpCollectPolicy).filter(LldpCollectPolicy.id == 1).with_for_update().one()
+ if (
+ db.query(TopoDiscoverJob)
+ .filter(TopoDiscoverJob.status.in_(["pending", "running"]))
+ .first()
+ is not None
+ ):
+ raise HTTPException(status_code=409, detail="lldp_collect_already_running")
scope = str(body.scope or "ne_ids").strip().lower() or "ne_ids"
if scope not in {"all_inventory", "ne_ids"}:
raise HTTPException(status_code=400, detail="invalid_scope")
@@ -1779,11 +2093,18 @@ def start_discover_job(
if trig not in {"manual", "schedule", "topology"}:
trig = "manual"
now = _utcnow()
+ # Persist explicit source lists when present; keep legacy ne_ids for older clients.
+ stored_ids = list(body.ne_ids or [])
+ if body.managed_ne_ids or body.ume_ne_ids:
+ stored_ids = [
+ *(f"managed:{x}" for x in (body.managed_ne_ids or []) if str(x).strip()),
+ *(f"ume:{x}" for x in (body.ume_ne_ids or []) if str(x).strip()),
+ ]
job = TopoDiscoverJob(
id=uuid4().hex,
scope=scope,
trigger_mode=trig,
- ne_ids_json=list(body.ne_ids or []),
+ ne_ids_json=stored_ids,
status="pending",
total=0,
done=0,
diff --git a/tests/test_lldp_collect.py b/tests/test_lldp_collect.py
index 1c44fff..8f33631 100644
--- a/tests/test_lldp_collect.py
+++ b/tests/test_lldp_collect.py
@@ -3,26 +3,36 @@
from __future__ import annotations
import unittest
+from datetime import datetime, timedelta
from uuid import uuid4
from netx_api.db import Base, SessionLocal, engine
from netx_api.lldp_collect_schemas import LldpCollectPolicyUpdate
from netx_api.lldp_collect_service import (
+ build_discover_request,
ensure_policy,
get_dashboard,
+ has_running_job,
next_due_at,
update_policy,
)
-from netx_api.models import LldpCollectPolicy
+from netx_api.models import LldpCollectPolicy, ManagedNE, TopoDiscoverJob, TopoDiscoverJobItem
+from netx_api.topology_service import prune_discover_jobs, reclaim_stale_discover_jobs
class LldpCollectTests(unittest.TestCase):
@classmethod
def setUpClass(cls) -> None:
Base.metadata.create_all(bind=engine)
+ from netx_api.topology_migrate import ensure_topology_schema
+
+ with engine.begin() as conn:
+ ensure_topology_schema(conn)
def setUp(self) -> None:
self.db = SessionLocal()
+ self.db.query(TopoDiscoverJobItem).delete()
+ self.db.query(TopoDiscoverJob).delete()
self.db.query(LldpCollectPolicy).delete()
self.db.commit()
@@ -41,6 +51,7 @@ class LldpCollectTests(unittest.TestCase):
dash = get_dashboard(self.db)
self.assertFalse(dash.policy.enabled)
self.assertIsNone(dash.next_due_at)
+ self.assertEqual(dash.policy.history_keep, 30)
def test_policy_enable_updates(self) -> None:
ensure_policy(self.db)
@@ -48,18 +59,189 @@ class LldpCollectTests(unittest.TestCase):
self.db,
LldpCollectPolicyUpdate(
enabled=True,
- interval_days=2,
+ interval_hours=48,
concurrency=6,
scope_mode="all",
auto_add_unmatched=True,
+ history_keep=5,
),
)
self.assertTrue(out.enabled)
+ self.assertEqual(out.interval_hours, 48)
self.assertEqual(out.interval_days, 2)
self.assertEqual(out.concurrency, 6)
+ self.assertEqual(out.history_keep, 5)
due = next_due_at(self.db, ensure_policy(self.db))
self.assertIsNotNone(due)
+ def test_next_due_uses_hours_and_ignores_manual(self) -> None:
+ policy = ensure_policy(self.db)
+ policy.enabled = True
+ policy.interval_hours = 6
+ policy.interval_days = 1
+ self.db.commit()
+ now = datetime.utcnow()
+ sched = TopoDiscoverJob(
+ id=uuid4().hex,
+ scope="all_inventory",
+ trigger_mode="schedule",
+ status="done",
+ ended_at=now - timedelta(hours=1),
+ created_at=now - timedelta(hours=2),
+ updated_at=now - timedelta(hours=1),
+ )
+ manual = TopoDiscoverJob(
+ id=uuid4().hex,
+ scope="all_inventory",
+ trigger_mode="manual",
+ status="done",
+ ended_at=now,
+ created_at=now,
+ updated_at=now,
+ )
+ self.db.add(sched)
+ self.db.add(manual)
+ self.db.commit()
+ due = next_due_at(self.db, ensure_policy(self.db))
+ self.assertIsNotNone(due)
+ assert due is not None
+ self.assertEqual(due, sched.ended_at + timedelta(hours=6))
+
+ def test_start_discover_rejects_second_while_running(self) -> None:
+ from netx_api.topology_schemas import FabricDiscoverRequest
+ from netx_api.topology_service import start_discover_job
+
+ now = datetime.utcnow()
+ running = TopoDiscoverJob(
+ id=uuid4().hex,
+ scope="all_inventory",
+ trigger_mode="manual",
+ status="running",
+ created_at=now,
+ updated_at=now,
+ started_at=now,
+ )
+ self.db.add(running)
+ ensure_policy(self.db)
+ self.db.commit()
+ with self.assertRaises(Exception) as ctx:
+ start_discover_job(self.db, FabricDiscoverRequest(scope="ne_ids", ne_ids=["x"]))
+ detail = getattr(ctx.exception, "detail", str(ctx.exception))
+ self.assertEqual(detail, "lldp_collect_already_running")
+
+ def test_build_request_respects_source(self) -> None:
+ suffix = uuid4().hex[:8]
+ ne = ManagedNE(
+ id=f"m-{suffix}",
+ name=f"M-{suffix}",
+ ip_address=f"203.0.113.{(int(suffix[:2], 16) % 200) + 1}",
+ vendor="Cisco",
+ device_type="cisco_ios",
+ )
+ self.db.add(ne)
+ policy = ensure_policy(self.db)
+ policy.scope_mode = "selected"
+ # Same id string marked as ume should NOT resolve via managed path when building lists.
+ policy.selected_targets = [
+ {"source": "managed", "id": ne.id},
+ {"source": "ume", "id": f"ume-{suffix}"},
+ ]
+ self.db.commit()
+ req = build_discover_request(ensure_policy(self.db))
+ self.assertEqual(req.managed_ne_ids, [ne.id])
+ self.assertEqual(req.ume_ne_ids, [f"ume-{suffix}"])
+ self.assertEqual(req.ne_ids, [])
+ self.db.delete(ne)
+ self.db.commit()
+
+ def test_prune_discover_jobs_keeps_newest(self) -> None:
+ now = datetime.utcnow()
+ ids: list[str] = []
+ for i in range(5):
+ jid = uuid4().hex
+ ids.append(jid)
+ self.db.add(
+ TopoDiscoverJob(
+ id=jid,
+ scope="all_inventory",
+ trigger_mode="manual",
+ status="done",
+ created_at=now - timedelta(minutes=5 - i),
+ updated_at=now - timedelta(minutes=5 - i),
+ ended_at=now - timedelta(minutes=5 - i),
+ )
+ )
+ self.db.add(
+ TopoDiscoverJobItem(
+ id=uuid4().hex,
+ job_id=jid,
+ ne_name=f"n{i}",
+ raw_preview="x" * 20,
+ created_at=now,
+ )
+ )
+ # Keep one open job — must survive prune.
+ open_id = uuid4().hex
+ self.db.add(
+ TopoDiscoverJob(
+ id=open_id,
+ scope="all_inventory",
+ trigger_mode="manual",
+ status="running",
+ created_at=now,
+ updated_at=now,
+ )
+ )
+ self.db.commit()
+ dropped = prune_discover_jobs(self.db, keep=2)
+ self.assertEqual(dropped, 3)
+ left = {r.id for r in self.db.query(TopoDiscoverJob).all()}
+ self.assertIn(open_id, left)
+ self.assertEqual(len(left), 3) # 2 finished + 1 running
+
+ def test_reclaim_force_all_open_on_startup(self) -> None:
+ now = datetime.utcnow()
+ job = TopoDiscoverJob(
+ id=uuid4().hex,
+ scope="all_inventory",
+ trigger_mode="schedule",
+ status="running",
+ total=10,
+ 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, 1)
+ self.db.refresh(job)
+ self.assertEqual(job.status, "failed")
+ self.assertIn("stale_running_reset_on_startup", job.error or "")
+ self.assertIsNone(has_running_job(self.db))
+
+ def test_reclaim_running_by_stale_updated_at(self) -> None:
+ old = datetime.utcnow() - timedelta(hours=5)
+ job = TopoDiscoverJob(
+ id=uuid4().hex,
+ scope="ne_ids",
+ trigger_mode="manual",
+ status="running",
+ total=3,
+ done=0,
+ created_at=old,
+ updated_at=old,
+ started_at=old,
+ )
+ self.db.add(job)
+ self.db.commit()
+ closed = reclaim_stale_discover_jobs(self.db)
+ self.assertEqual(closed, 1)
+ self.db.refresh(job)
+ self.assertEqual(job.status, "failed")
+ self.assertIn("running_stale_timeout", job.error or "")
+
if __name__ == "__main__":
unittest.main()
diff --git a/tests/test_topology.py b/tests/test_topology.py
index d58f9f1..f7420b0 100644
--- a/tests/test_topology.py
+++ b/tests/test_topology.py
@@ -87,6 +87,10 @@ class FabricTopologyTests(unittest.TestCase):
@classmethod
def setUpClass(cls) -> None:
Base.metadata.create_all(bind=engine)
+ from netx_api.topology_migrate import ensure_topology_schema
+
+ with engine.begin() as conn:
+ ensure_topology_schema(conn)
def setUp(self) -> None:
self.db = SessionLocal()
@@ -482,6 +486,112 @@ Management Addresses:
self.db.commit()
return fa, fb, ne_a, ne_b
+ def test_merge_lldp_placeholder_into_real_inventory(self) -> None:
+ suffix = uuid4().hex[:8]
+ # Real inventory NE + fabric node.
+ real = ManagedNE(
+ id=f"real-{suffix}",
+ name=f"R1-{suffix}",
+ ip_address=f"198.51.100.{(int(suffix[:2], 16) % 80) + 10}",
+ vendor="Cisco",
+ device_type="cisco_ios",
+ )
+ self.db.add(real)
+ self.db.commit()
+ fr = svc.ensure_fabric_node_for_managed(self.db, real)
+ # LLDP placeholder with same hostname key + seen mgmt IP.
+ ph = svc.ensure_lldp_discovered_managed_ne(
+ self.db,
+ remote_name=f"R1-{suffix}",
+ remote_ip=real.ip_address,
+ )
+ fp = svc.ensure_fabric_node_for_managed(self.db, ph)
+ self.db.commit()
+ # Edge hanging off placeholder should retarget to real.
+ edge, _ = svc.upsert_fabric_edge(
+ self.db,
+ a_node_id=fr.id,
+ b_node_id=fp.id,
+ a_port="Gi0/0",
+ b_port="Gi0/1",
+ source="lldp",
+ )
+ # Need a third node so edge isn't self-loop after merge… actually A=real B=placeholder
+ # after absorb B→A becomes self-loop and edge is deleted. Use external peer.
+ peer_ne = ManagedNE(
+ id=f"peer-{suffix}",
+ name=f"P-{suffix}",
+ ip_address=f"198.51.100.{(int(suffix[2:4], 16) % 80) + 100}",
+ vendor="Cisco",
+ device_type="cisco_ios",
+ )
+ self.db.add(peer_ne)
+ self.db.commit()
+ fpeer = svc.ensure_fabric_node_for_managed(self.db, peer_ne)
+ edge2, _ = svc.upsert_fabric_edge(
+ self.db,
+ a_node_id=fp.id,
+ b_node_id=fpeer.id,
+ a_port="Gi1/0",
+ b_port="Gi1/1",
+ source="lldp",
+ )
+ self.db.commit()
+ edge2_id = edge2.id
+
+ out = svc.merge_duplicate_fabric_nodes(self.db)
+ self.assertGreaterEqual(out["merged"], 1)
+ self.assertGreaterEqual(out.get("placeholders_removed", 0), 1)
+ self.db.expire_all()
+ self.assertIsNone(self.db.get(TopoFabricNode, fp.id))
+ self.assertIsNone(self.db.get(ManagedNE, ph.id))
+ # Edge from placeholder→peer should now be real→peer.
+ moved = self.db.get(TopoFabricEdge, edge2_id)
+ self.assertIsNotNone(moved)
+ assert moved is not None
+ ends = {moved.a_node_id, moved.b_node_id}
+ self.assertEqual(ends, {fr.id, fpeer.id})
+
+ self.db.delete(real)
+ self.db.delete(peer_ne)
+ self.db.commit()
+
+ def test_list_fabric_edges_missing_filter_and_names(self) -> None:
+ suffix = uuid4().hex[:8]
+ fa, fb, ne_a, ne_b = self._pair_nodes(suffix)
+ edge, _ = svc.upsert_fabric_edge(
+ self.db,
+ a_node_id=fa.id,
+ b_node_id=fb.id,
+ a_port="Gi0/0",
+ b_port="Gi0/1",
+ source="lldp",
+ )
+ self.db.commit()
+ svc._apply_missing_and_purge(
+ self.db, scanned_ok={fa.id}, touched_edge_ids=set()
+ )
+ self.db.commit()
+
+ missing = svc.list_fabric_edges(self.db, status="missing", page=1, page_size=50)
+ self.assertGreaterEqual(missing["total"], 1)
+ hit = next(i for i in missing["items"] if i["id"] == edge.id)
+ self.assertEqual(hit["status"], "missing")
+ self.assertTrue(hit["a_name"] or hit["b_name"])
+ self.assertGreaterEqual(int((hit.get("attrs") or {}).get("miss_count") or 0), 1)
+
+ by_kw = svc.list_fabric_edges(
+ self.db, keyword=ne_a.name[:6], status="missing", page=1, page_size=50
+ )
+ self.assertTrue(any(i["id"] == edge.id for i in by_kw["items"]))
+
+ active = svc.list_fabric_edges(self.db, status="active", page=1, page_size=50)
+ self.assertFalse(any(i["id"] == edge.id for i in active["items"]))
+
+ self.db.delete(ne_a)
+ self.db.delete(ne_b)
+ self.db.commit()
+
def test_edge_missing_after_one_absent_cycle(self) -> None:
suffix = uuid4().hex[:8]
fa, fb, ne_a, ne_b = self._pair_nodes(suffix)
diff --git a/web/src/constants/queryKeys.ts b/web/src/constants/queryKeys.ts
index 95fed4f..76b8e03 100644
--- a/web/src/constants/queryKeys.ts
+++ b/web/src/constants/queryKeys.ts
@@ -57,6 +57,9 @@ export const queryKeys = {
lldpCollectJobs: (page: number) => ["lldpCollectJobs", page] as const,
lldpCollectJobAll: ["lldpCollectJob"] as const,
lldpCollectJob: (jobId: string) => ["lldpCollectJob", jobId] as const,
+ fabricEdgesAll: ["fabricEdges"] as const,
+ fabricEdges: (status: string, keyword: string, page: number) =>
+ ["fabricEdges", status, keyword, page] as const,
configSyncDashboard: ["configSyncDashboard"] as const,
configSyncPolicy: ["configSyncPolicy"] as const,
configSyncCyclesAll: ["configSyncCycles"] as const,
diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts
index 5cfbe01..694e317 100644
--- a/web/src/i18n/en.ts
+++ b/web/src/i18n/en.ts
@@ -236,7 +236,11 @@ const en = {
enabled: "Enable scheduled collect",
autoAddUnmatched: "Create SSH placeholder NEs for unmatched neighbors",
intervalDays: "Interval (days)",
+ interval: "Collect interval",
+ intervalUnitDays: "days",
+ intervalUnitHours: "hours",
concurrency: "Concurrency",
+ historyKeep: "Jobs to keep",
scope: "Scope",
scopeAll: "All managed NEs",
scopeSelected: "Selected NEs",
@@ -245,6 +249,13 @@ const en = {
targetKeywordPh: "Name / IP / vendor",
jobsTitle: "Collect jobs",
jobDetailTitle: "Job details",
+ edgesTitle: "Fabric links",
+ edgeKeywordPh: "Endpoint name or IP",
+ edgeStatusAll: "All statuses",
+ edgeStatusActive: "Active",
+ edgeStatusMissing: "Missing",
+ missCount: "Miss cycles",
+ replacedBy: "Replaced",
kpi: {
nodes: "Fabric nodes",
edges: "Active links",
@@ -268,6 +279,11 @@ const en = {
ended: "Ended",
neighbors: "Neighbors",
error: "Error",
+ aSide: "Side A",
+ bSide: "Side B",
+ aPort: "Port A",
+ bPort: "Port B",
+ lastSeen: "Last seen",
},
},
configSync: {
@@ -1309,6 +1325,7 @@ const en = {
discoverLocatePeer: "Locate peer",
discoverLocateNe: "Locate NE",
discoverRawPreview: "Raw output preview (truncated when long; parsing uses full output)",
+ discoverGoLldp: "Details → LLDP links",
findNode: "Find",
findNodePh: "Type to search: name / IP / vendor",
findNoMatch: "No matching node on the canvas",
@@ -1316,7 +1333,7 @@ const en = {
locateOnCanvas: "Click to locate on canvas",
openPortTraffic: "Open port traffic",
removeStale: "Clear missing ({{count}})",
- removeStaleHint: "Delete red-marked missing edges",
+ removeStaleHint: "Manually delete red missing edges; scheduled collect also purges after 4 consecutive misses",
staleRemoved: "Removed {{count}} missing edge(s)",
connectHint: "Connect mode: click source node, then target (Esc or Select to exit)",
},
diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts
index 7939377..f6ac710 100644
--- a/web/src/i18n/zh.ts
+++ b/web/src/i18n/zh.ts
@@ -236,7 +236,11 @@ const zh = {
enabled: "启用周期调度",
autoAddUnmatched: "未纳管邻居建 SSH 占位网元",
intervalDays: "周期(天)",
+ interval: "采集周期",
+ intervalUnitDays: "天",
+ intervalUnitHours: "小时",
concurrency: "并发",
+ historyKeep: "任务保留数",
scope: "范围",
scopeAll: "全部纳管网元",
scopeSelected: "指定网元",
@@ -245,6 +249,13 @@ const zh = {
targetKeywordPh: "名称 / IP / 厂商",
jobsTitle: "采集任务",
jobDetailTitle: "任务明细",
+ edgesTitle: "Fabric 链路",
+ edgeKeywordPh: "本端 / 对端名称或 IP",
+ edgeStatusAll: "全部状态",
+ edgeStatusActive: "活跃",
+ edgeStatusMissing: "未发现",
+ missCount: "未发现周期",
+ replacedBy: "已被替换",
kpi: {
nodes: "Fabric 网元",
edges: "活跃链路",
@@ -268,6 +279,11 @@ const zh = {
ended: "结束",
neighbors: "邻居",
error: "错误",
+ aSide: "本端",
+ bSide: "对端",
+ aPort: "本端口",
+ bPort: "对端口",
+ lastSeen: "最近看见",
},
},
configSync: {
@@ -1305,6 +1321,7 @@ const zh = {
discoverLocatePeer: "定位对端",
discoverLocateNe: "定位本端",
discoverRawPreview: "采集原文预览(过长会截断,不影响解析)",
+ discoverGoLldp: "查看详情 → LLDP 链路",
findNode: "定位",
findNodePh: "输入即搜:名称 / IP / 厂商",
findNoMatch: "画布上没有匹配的节点",
@@ -1312,7 +1329,7 @@ const zh = {
locateOnCanvas: "点击定位到画布",
openPortTraffic: "打开端口流量",
removeStale: "清除未发现 ({{count}})",
- removeStaleHint: "删除红色标记的未发现链路",
+ removeStaleHint: "手动删除红色未发现链路;周期采集连续 4 次未发现也会自动清理",
staleRemoved: "已删除 {{count}} 条未发现链路",
connectHint: "连线模式:依次点击源节点 → 目标节点(Esc 或切回选择退出)",
},
diff --git a/web/src/pages/TopologyPage.tsx b/web/src/pages/TopologyPage.tsx
index d7d69ab..8a73dac 100644
--- a/web/src/pages/TopologyPage.tsx
+++ b/web/src/pages/TopologyPage.tsx
@@ -1,4 +1,5 @@
-import { createContext, useCallback, useContext, useEffect, useMemo, useRef, useState } from "react";
+import { createContext, useCallback, useContext, useEffect, useMemo, useRef, useState } from "react";
+import { Link } from "react-router-dom";
import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query";
import {
ReactFlow,
@@ -324,7 +325,7 @@ const BUILTIN_EDGE_DEFAULTS: EdgeDefaults = {
function sourceKind(source: string): EdgeSourceKind {
const src = (source || "manual").toLowerCase();
- if (src === "stale") return "stale";
+ if (src === "stale" || src === "missing") return "stale";
if (src === "lldp" || src === "cdp") return "discovered";
return "manual";
}
@@ -514,12 +515,6 @@ export function TopologyPage() {
const [searchHitIds, setSearchHitIds] = useState
+ + {t("topology.discoverGoLldp")} + +
> ) : null} @@ -2562,316 +2515,6 @@ export function TopologyPage() { ) : null} - {discoverListOpen ? ( -- {t("topology.discoverSummary") - .replace("{{scanned}}", String(discoverSummary.scanned)) - .replace("{{added}}", String(discoverSummary.added)) - .replace("{{updated}}", String(discoverSummary.updated)) - .replace("{{stale}}", String(discoverSummary.stale)) - .replace("{{failed}}", String(discoverSummary.failed))} -
-{t("topology.discoverListEmpty")}
- ) : ( -- {t("topology.discoverShowLimited") - .replace("{{shown}}", String(DISCOVER_LIST_CAP)) - .replace("{{total}}", String(discoverFiltered.length))} -
- ) : null} -- {[discoverDetail.ne_ip, discoverDetail.command].filter(Boolean).join(SEP)} -
-- {discoverDetail.error || t("topology.discoverNeFail")} -
- ) : null} - {discoverDetail.parser_stub ? ( -- {t("topology.discoverParserStub").replace( - "{{parser}}", - discoverDetail.parser_key || "unknown", - )} -
- ) : null} - -{t("topology.discoverLinksEmpty")}
- ) : ( -| {t("topology.discoverColPeer")} | -{t("topology.discoverColLocalPort")} | -{t("topology.discoverColRemotePort")} | -{t("topology.discoverColProtocol")} | -{t("topology.discoverColAction")} | -- |
|---|---|---|---|---|---|
|
-
- {link.peer_name || link.peer_ip || link.peer_ne_id || "?"}
- {link.peer_ip ? {link.peer_ip} : null}
-
- |
- {link.local_port || "?"} | -{link.remote_port || "?"} | -{(link.protocol || "?").toUpperCase()} | -- - {link.action === "added" - ? t("topology.discoverLinkAdded") - : link.action === "updated" - ? t("topology.discoverLinkUpdated") - : link.action === "kept_manual" - ? t("topology.discoverLinkKeptManual") - : link.action || "?"} - - | -- {link.peer_node_id ? ( - - ) : null} - | -
{t("topology.discoverUnmatchedEmpty")}
- ) : ( -| {t("topology.discoverColRemote")} | -{t("topology.discoverColLocalPort")} | -{t("topology.discoverColRemotePort")} | -
|---|---|---|
| - {(u.remote_name || u.remote_ip || "?").trim()} - {u.remote_ip && u.remote_name ? ` (${u.remote_ip})` : ""} - | -{u.local_port || "?"} | -{u.remote_port || "?"} | -
{discoverDetail.raw_preview}
- >
- ) : null}
-
-