mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 03:10:46 +08:00
Add LLDP collect page with conservative fabric edge lifecycle.
Move scheduled discovery under Network → LLDP links, mark absent edges missing after one successful scan, purge only after four miss cycles, and expose unmatched/raw job detail. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
9674de221d
commit
d05ba0f1f7
30 changed files with 4713 additions and 1870 deletions
|
|
@ -77,6 +77,11 @@ class Settings(BaseSettings):
|
|||
config_sync_scheduler_tick_sec: int = 60
|
||||
# After process start / unexpected restart, wait before any scheduled sync.
|
||||
config_sync_startup_grace_sec: int = 3600
|
||||
# Fabric LLDP collect (network topology management).
|
||||
# Scheduler thread may run, but policy.enabled defaults False — no collect until operator turns it on.
|
||||
lldp_collect_scheduler_enabled: bool = True
|
||||
lldp_collect_scheduler_tick_sec: int = 60
|
||||
lldp_collect_startup_grace_sec: int = 3600
|
||||
# Port traffic monitoring (CLI rate bit/s samples)
|
||||
port_traffic_scheduler_enabled: bool = True
|
||||
port_traffic_scheduler_tick_sec: int = 15
|
||||
|
|
|
|||
|
|
@ -49,3 +49,6 @@ WEBCRT_DEVICE_TYPES: tuple[str, ...] = SUPPORTED_DEVICE_TYPES + ("linux", "gener
|
|||
|
||||
# ManagedNE.source value for sessions created via WebCRT Quick Connect.
|
||||
WEBCRT_NE_SOURCE = "webcrt"
|
||||
|
||||
# ManagedNE.source for LLDP-discovered peers not yet in inventory (SSH shell, empty creds).
|
||||
LLDP_DISCOVERED_NE_SOURCE = "lldp"
|
||||
|
|
|
|||
57
netx_api/lldp_collect_router.py
Normal file
57
netx_api/lldp_collect_router.py
Normal file
|
|
@ -0,0 +1,57 @@
|
|||
"""HTTP routes for network-management LLDP link collect."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
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_service import (
|
||||
get_dashboard,
|
||||
get_job_detail,
|
||||
get_policy,
|
||||
list_jobs,
|
||||
start_collect,
|
||||
update_policy,
|
||||
)
|
||||
|
||||
router = APIRouter(prefix="/v1/topology/lldp-collect", tags=["topology-lldp-collect"])
|
||||
|
||||
|
||||
@router.get("/policy")
|
||||
def api_get_policy(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_policy(db).model_dump()
|
||||
|
||||
|
||||
@router.put("/policy")
|
||||
def api_put_policy(
|
||||
body: LldpCollectPolicyUpdate, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return update_policy(db, body).model_dump()
|
||||
|
||||
|
||||
@router.get("/dashboard")
|
||||
def api_dashboard(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_dashboard(db).model_dump()
|
||||
|
||||
|
||||
@router.post("/start")
|
||||
def api_start(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return start_collect(db, trigger_mode="manual")
|
||||
|
||||
|
||||
@router.get("/jobs")
|
||||
def api_list_jobs(
|
||||
page: int = Query(default=1, ge=1),
|
||||
page_size: int = Query(default=20, ge=1, le=100),
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
return list_jobs(db, page=page, page_size=page_size)
|
||||
|
||||
|
||||
@router.get("/jobs/{job_id}")
|
||||
def api_get_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_job_detail(db, job_id)
|
||||
96
netx_api/lldp_collect_scheduler.py
Normal file
96
netx_api/lldp_collect_scheduler.py
Normal file
|
|
@ -0,0 +1,96 @@
|
|||
"""Background scheduler for periodic fabric LLDP collect."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
from .config import settings
|
||||
from .db import SessionLocal
|
||||
from .lldp_collect_service import (
|
||||
ensure_policy,
|
||||
has_running_job,
|
||||
next_due_at,
|
||||
start_collect,
|
||||
)
|
||||
|
||||
_log = logging.getLogger("netx.lldp_collect.scheduler")
|
||||
_stop = threading.Event()
|
||||
_thread: threading.Thread | None = None
|
||||
_BOOT_MONO = time.monotonic()
|
||||
|
||||
|
||||
def _utcnow() -> datetime:
|
||||
return datetime.utcnow()
|
||||
|
||||
|
||||
def startup_grace_remaining_sec() -> float:
|
||||
grace = max(0, int(getattr(settings, "lldp_collect_startup_grace_sec", 3600) or 0))
|
||||
elapsed = time.monotonic() - _BOOT_MONO
|
||||
return max(0.0, float(grace) - elapsed)
|
||||
|
||||
|
||||
def in_startup_grace() -> bool:
|
||||
return startup_grace_remaining_sec() > 0
|
||||
|
||||
|
||||
def try_start_scheduled_collect() -> str | None:
|
||||
if not bool(getattr(settings, "lldp_collect_scheduler_enabled", True)):
|
||||
return None
|
||||
if in_startup_grace():
|
||||
return None
|
||||
db = SessionLocal()
|
||||
try:
|
||||
policy = ensure_policy(db)
|
||||
if not policy.enabled:
|
||||
return None
|
||||
if has_running_job(db) is not None:
|
||||
return None
|
||||
due = next_due_at(db, policy)
|
||||
if due is not None and due > _utcnow():
|
||||
return None
|
||||
out = start_collect(db, trigger_mode="schedule")
|
||||
job_id = str((out.get("job") or {}).get("id") or "")
|
||||
_log.info("lldp_collect scheduled job started id=%s", job_id)
|
||||
return job_id or None
|
||||
except Exception as exc: # noqa: BLE001
|
||||
db.rollback()
|
||||
detail = getattr(exc, "detail", None)
|
||||
if detail in {"no_selected_targets", "lldp_collect_already_running"}:
|
||||
_log.info("lldp_collect schedule skip: %s", detail)
|
||||
return None
|
||||
_log.exception("lldp_collect schedule start failed")
|
||||
return None
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def _loop() -> None:
|
||||
tick = max(15, int(getattr(settings, "lldp_collect_scheduler_tick_sec", 60) or 60))
|
||||
grace = max(0, int(getattr(settings, "lldp_collect_startup_grace_sec", 3600) or 0))
|
||||
_log.info("lldp_collect scheduler started tick=%ss startup_grace=%ss", tick, grace)
|
||||
while not _stop.is_set():
|
||||
try:
|
||||
try_start_scheduled_collect()
|
||||
except Exception:
|
||||
_log.exception("lldp_collect scheduler tick failed")
|
||||
_stop.wait(tick)
|
||||
_log.info("lldp_collect scheduler stopped")
|
||||
|
||||
|
||||
def start_lldp_collect_scheduler() -> None:
|
||||
global _thread
|
||||
if not bool(getattr(settings, "lldp_collect_scheduler_enabled", True)):
|
||||
_log.info("lldp_collect scheduler disabled by settings")
|
||||
return
|
||||
if _thread is not None and _thread.is_alive():
|
||||
return
|
||||
_stop.clear()
|
||||
_thread = threading.Thread(target=_loop, name="lldp-collect-scheduler", daemon=True)
|
||||
_thread.start()
|
||||
|
||||
|
||||
def stop_lldp_collect_scheduler() -> None:
|
||||
_stop.set()
|
||||
65
netx_api/lldp_collect_schemas.py
Normal file
65
netx_api/lldp_collect_schemas.py
Normal file
|
|
@ -0,0 +1,65 @@
|
|||
"""Schemas for network-management LLDP link collect (policy + dashboard)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class LldpCollectTargetRef(BaseModel):
|
||||
source: str = "managed" # managed | ume
|
||||
id: str
|
||||
|
||||
|
||||
class LldpCollectPolicyOut(BaseModel):
|
||||
enabled: bool = False
|
||||
interval_days: int = 1
|
||||
concurrency: int = 4
|
||||
scope_mode: str = "all"
|
||||
selected_targets: list[LldpCollectTargetRef] = Field(default_factory=list)
|
||||
auto_add_unmatched: bool = True
|
||||
updated_at: datetime | None = None
|
||||
|
||||
|
||||
class LldpCollectPolicyUpdate(BaseModel):
|
||||
enabled: bool | None = None
|
||||
interval_days: int | None = Field(default=None, ge=1, le=365)
|
||||
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
|
||||
|
||||
|
||||
class LldpCollectJobSummary(BaseModel):
|
||||
id: str
|
||||
scope: str = ""
|
||||
trigger_mode: str = "manual"
|
||||
status: str = ""
|
||||
total: int = 0
|
||||
done: int = 0
|
||||
edges_added: int = 0
|
||||
edges_updated: int = 0
|
||||
edges_stale: int = 0
|
||||
error: str = ""
|
||||
started_at: datetime | None = None
|
||||
ended_at: datetime | None = None
|
||||
created_at: datetime | None = None
|
||||
|
||||
|
||||
class LldpCollectDashboardOut(BaseModel):
|
||||
policy: LldpCollectPolicyOut
|
||||
fabric_node_count: int = 0
|
||||
fabric_edge_count: int = 0
|
||||
fabric_edge_active: int = 0
|
||||
fabric_edge_stale: int = 0
|
||||
last_discover_at: datetime | None = None
|
||||
running_job: LldpCollectJobSummary | None = None
|
||||
last_job: LldpCollectJobSummary | None = None
|
||||
next_due_at: datetime | None = None
|
||||
|
||||
|
||||
class LldpCollectStartOut(BaseModel):
|
||||
ok: bool = True
|
||||
job: dict[str, Any] = Field(default_factory=dict)
|
||||
234
netx_api/lldp_collect_service.py
Normal file
234
netx_api/lldp_collect_service.py
Normal file
|
|
@ -0,0 +1,234 @@
|
|||
"""LLDP collect policy + dashboard (network topology management)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from fastapi import HTTPException
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from .lldp_collect_schemas import (
|
||||
LldpCollectDashboardOut,
|
||||
LldpCollectJobSummary,
|
||||
LldpCollectPolicyOut,
|
||||
LldpCollectPolicyUpdate,
|
||||
LldpCollectTargetRef,
|
||||
)
|
||||
from .models import LldpCollectPolicy, TopoDiscoverJob, TopoFabricStats
|
||||
from .topology_schemas import FabricDiscoverRequest
|
||||
from .topology_service import get_discover_job, start_discover_job
|
||||
|
||||
POLICY_ID = 1
|
||||
|
||||
|
||||
def _utcnow() -> datetime:
|
||||
return datetime.utcnow()
|
||||
|
||||
|
||||
def ensure_policy(db: Session) -> LldpCollectPolicy:
|
||||
row = db.get(LldpCollectPolicy, POLICY_ID)
|
||||
if row is None:
|
||||
row = LldpCollectPolicy(
|
||||
id=POLICY_ID,
|
||||
enabled=False,
|
||||
interval_days=1,
|
||||
concurrency=4,
|
||||
scope_mode="all",
|
||||
selected_targets=[],
|
||||
auto_add_unmatched=True,
|
||||
updated_at=_utcnow(),
|
||||
)
|
||||
db.add(row)
|
||||
db.commit()
|
||||
db.refresh(row)
|
||||
return row
|
||||
|
||||
|
||||
def _policy_out(row: LldpCollectPolicy) -> LldpCollectPolicyOut:
|
||||
refs: list[LldpCollectTargetRef] = []
|
||||
for raw in row.selected_targets or []:
|
||||
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 not in {"managed", "ume"}:
|
||||
src = "managed"
|
||||
refs.append(LldpCollectTargetRef(source=src, id=tid))
|
||||
return LldpCollectPolicyOut(
|
||||
enabled=bool(row.enabled),
|
||||
interval_days=int(row.interval_days or 1),
|
||||
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),
|
||||
updated_at=row.updated_at,
|
||||
)
|
||||
|
||||
|
||||
def get_policy(db: Session) -> LldpCollectPolicyOut:
|
||||
return _policy_out(ensure_policy(db))
|
||||
|
||||
|
||||
def update_policy(db: Session, body: LldpCollectPolicyUpdate) -> LldpCollectPolicyOut:
|
||||
row = ensure_policy(db)
|
||||
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 "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:
|
||||
mode = str(data["scope_mode"] or "").strip().lower()
|
||||
if mode not in {"all", "selected"}:
|
||||
raise HTTPException(status_code=400, detail="invalid_scope_mode")
|
||||
row.scope_mode = mode
|
||||
if "selected_targets" in data and data["selected_targets"] is not None:
|
||||
cleaned: list[dict[str, str]] = []
|
||||
for ref in data["selected_targets"] or []:
|
||||
if isinstance(ref, LldpCollectTargetRef):
|
||||
tid = ref.id.strip()
|
||||
src = ref.source.strip().lower() or "managed"
|
||||
elif isinstance(ref, dict):
|
||||
tid = str(ref.get("id") or "").strip()
|
||||
src = str(ref.get("source") or "managed").strip().lower() or "managed"
|
||||
else:
|
||||
continue
|
||||
if not tid:
|
||||
continue
|
||||
if src not in {"managed", "ume"}:
|
||||
src = "managed"
|
||||
cleaned.append({"source": src, "id": tid})
|
||||
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"])
|
||||
row.updated_at = _utcnow()
|
||||
db.commit()
|
||||
db.refresh(row)
|
||||
return _policy_out(row)
|
||||
|
||||
|
||||
def _job_summary(job: TopoDiscoverJob | None) -> LldpCollectJobSummary | None:
|
||||
if job is None:
|
||||
return None
|
||||
return LldpCollectJobSummary(
|
||||
id=job.id,
|
||||
scope=job.scope or "",
|
||||
trigger_mode=getattr(job, "trigger_mode", None) or "manual",
|
||||
status=job.status or "",
|
||||
total=int(job.total or 0),
|
||||
done=int(job.done 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),
|
||||
error=job.error or "",
|
||||
started_at=job.started_at,
|
||||
ended_at=job.ended_at,
|
||||
created_at=job.created_at,
|
||||
)
|
||||
|
||||
|
||||
def has_running_job(db: Session) -> TopoDiscoverJob | None:
|
||||
return (
|
||||
db.query(TopoDiscoverJob)
|
||||
.filter(TopoDiscoverJob.status.in_(["pending", "running"]))
|
||||
.order_by(TopoDiscoverJob.created_at.desc())
|
||||
.first()
|
||||
)
|
||||
|
||||
|
||||
def last_finished_job(db: Session) -> TopoDiscoverJob | None:
|
||||
return (
|
||||
db.query(TopoDiscoverJob)
|
||||
.filter(TopoDiscoverJob.status.in_(["done", "failed"]))
|
||||
.order_by(TopoDiscoverJob.created_at.desc())
|
||||
.first()
|
||||
)
|
||||
|
||||
|
||||
def next_due_at(db: Session, policy: LldpCollectPolicy) -> datetime | None:
|
||||
if not policy.enabled:
|
||||
return None
|
||||
days = max(1, int(policy.interval_days or 1))
|
||||
last = (
|
||||
db.query(TopoDiscoverJob)
|
||||
.filter(TopoDiscoverJob.status == "done", 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)
|
||||
|
||||
|
||||
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] = []
|
||||
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:
|
||||
raise HTTPException(status_code=400, detail="no_selected_targets")
|
||||
return FabricDiscoverRequest(
|
||||
scope="ne_ids",
|
||||
ne_ids=ne_ids,
|
||||
concurrency=concurrency,
|
||||
auto_add_unmatched=auto_add,
|
||||
)
|
||||
return FabricDiscoverRequest(
|
||||
scope="all_inventory",
|
||||
ne_ids=[],
|
||||
concurrency=concurrency,
|
||||
auto_add_unmatched=auto_add,
|
||||
)
|
||||
|
||||
|
||||
def start_collect(db: Session, *, trigger_mode: str = "manual") -> 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)
|
||||
return {"ok": True, "job": job.model_dump()}
|
||||
|
||||
|
||||
def get_dashboard(db: Session) -> LldpCollectDashboardOut:
|
||||
policy = ensure_policy(db)
|
||||
stats = db.get(TopoFabricStats, "global")
|
||||
running = has_running_job(db)
|
||||
last = last_finished_job(db)
|
||||
return LldpCollectDashboardOut(
|
||||
policy=_policy_out(policy),
|
||||
fabric_node_count=int(stats.node_count if stats else 0),
|
||||
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),
|
||||
last_discover_at=stats.last_discover_at if stats else None,
|
||||
running_job=_job_summary(running),
|
||||
last_job=_job_summary(last),
|
||||
next_due_at=next_due_at(db, policy),
|
||||
)
|
||||
|
||||
|
||||
def list_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict:
|
||||
page = max(1, int(page or 1))
|
||||
page_size = max(1, min(100, int(page_size or 20)))
|
||||
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()
|
||||
return {
|
||||
"total": total,
|
||||
"page": page,
|
||||
"page_size": page_size,
|
||||
"items": [_job_summary(r).model_dump() for r in rows if _job_summary(r)],
|
||||
}
|
||||
|
||||
|
||||
def get_job_detail(db: Session, job_id: str) -> dict:
|
||||
return get_discover_job(db, job_id).model_dump()
|
||||
|
|
@ -31,6 +31,7 @@ from .port_traffic_router import router as port_traffic_router
|
|||
from .managed_ne_router import router as managed_ne_router
|
||||
from .webcrt_router import router as webcrt_router
|
||||
from .topology_router import router as topology_router
|
||||
from .lldp_collect_router import router as lldp_collect_router
|
||||
from .importer import aggregate_alarms, import_alarm_excel, query_alarms
|
||||
from .models import (
|
||||
AiAnalyzeHistory,
|
||||
|
|
@ -139,6 +140,7 @@ app.include_router(config_sync_router)
|
|||
app.include_router(port_traffic_router)
|
||||
app.include_router(webcrt_router)
|
||||
app.include_router(topology_router)
|
||||
app.include_router(lldp_collect_router)
|
||||
parser_cfg = load_parser_config()
|
||||
_UME_CLIENT_SINGLETON = UMEClient(
|
||||
token_loader=lambda: load_shared_token(),
|
||||
|
|
@ -814,10 +816,12 @@ def on_startup() -> None:
|
|||
)
|
||||
conn.exec_driver_sql("ALTER TABLE api_token ADD COLUMN IF NOT EXISTS expires_at TIMESTAMP")
|
||||
from .port_traffic_migrate import ensure_port_traffic_series_schema
|
||||
from .topology_migrate import ensure_topology_schema
|
||||
|
||||
ensure_port_traffic_series_schema(conn)
|
||||
ensure_topology_schema(conn)
|
||||
except Exception:
|
||||
_schedule_log.exception("startup: auth/port_traffic schema migration failed")
|
||||
_schedule_log.exception("startup: auth/port_traffic/topology schema migration failed")
|
||||
_reset_runtime_pause_flags()
|
||||
_fail_stale_running_sync_jobs_on_startup()
|
||||
if _needs_startup_alarm_sync_before_ws():
|
||||
|
|
@ -847,6 +851,12 @@ def on_startup() -> None:
|
|||
cfg_resumed = recover_config_sync_on_startup(db)
|
||||
if cfg_resumed:
|
||||
_schedule_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed)
|
||||
try:
|
||||
from .lldp_collect_service import ensure_policy as ensure_lldp_collect_policy
|
||||
|
||||
ensure_lldp_collect_policy(db)
|
||||
except Exception:
|
||||
_schedule_log.exception("startup: lldp_collect policy ensure failed")
|
||||
pt_cleared = recover_port_traffic_on_startup(db)
|
||||
if pt_cleared:
|
||||
_schedule_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared)
|
||||
|
|
@ -866,6 +876,12 @@ def on_startup() -> None:
|
|||
start_config_sync_scheduler()
|
||||
except Exception:
|
||||
_schedule_log.exception("startup: config_sync scheduler init failed")
|
||||
try:
|
||||
from .lldp_collect_scheduler import start_lldp_collect_scheduler
|
||||
|
||||
start_lldp_collect_scheduler()
|
||||
except Exception:
|
||||
_schedule_log.exception("startup: lldp_collect scheduler init failed")
|
||||
try:
|
||||
from .port_traffic_scheduler import start_port_traffic_scheduler
|
||||
|
||||
|
|
|
|||
|
|
@ -393,57 +393,188 @@ class NeCollectionRun(Base):
|
|||
ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
|
||||
|
||||
class TopologyMap(Base):
|
||||
"""Named topology canvas (document-style graph)."""
|
||||
class TopoFabricNode(Base):
|
||||
"""Global topology fabric node (inventory-aligned). Target scale ~50k."""
|
||||
|
||||
__tablename__ = "topology_map"
|
||||
__tablename__ = "topo_fabric_node"
|
||||
__table_args__ = (
|
||||
UniqueConstraint("managed_ne_id", name="uq_topo_fabric_node_managed_ne_id"),
|
||||
UniqueConstraint("ume_ne_id", name="uq_topo_fabric_node_ume_ne_id"),
|
||||
)
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
managed_ne_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||||
ume_ne_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True)
|
||||
name: Mapped[str] = mapped_column(String(256), default="", index=True)
|
||||
remark: Mapped[str] = mapped_column(String(1024), default="")
|
||||
ip: Mapped[str] = mapped_column(String(128), default="", index=True)
|
||||
vendor: Mapped[str] = mapped_column(String(64), default="")
|
||||
device_type: Mapped[str] = mapped_column(String(64), default="")
|
||||
attrs: Mapped[dict] = mapped_column(_JsonType, default=dict)
|
||||
last_seen_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, index=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
|
||||
|
||||
|
||||
class TopologyNode(Base):
|
||||
"""Node on a topology map; preferably references an inventory NE."""
|
||||
class TopoFabricEdge(Base):
|
||||
"""Global fabric link. Target scale ~1M; layer reserved for future BGP/tunnel/l2vpn."""
|
||||
|
||||
__tablename__ = "topology_node"
|
||||
__tablename__ = "topo_fabric_edge"
|
||||
__table_args__ = (
|
||||
UniqueConstraint(
|
||||
"layer",
|
||||
"a_node_id",
|
||||
"b_node_id",
|
||||
"a_port",
|
||||
"b_port",
|
||||
name="uq_topo_fabric_edge_endpoints",
|
||||
),
|
||||
)
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
map_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
managed_ne_id: Mapped[str] = mapped_column(String(64), default="", index=True)
|
||||
ume_ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
|
||||
label: Mapped[str] = mapped_column(String(256), default="")
|
||||
x: Mapped[float] = mapped_column(Float, default=0.0)
|
||||
y: Mapped[float] = mapped_column(Float, default=0.0)
|
||||
# physical | bgp | tunnel | l2vpn (P1 writes physical only)
|
||||
layer: Mapped[str] = mapped_column(String(32), default="physical", index=True)
|
||||
a_node_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
b_node_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
a_port: Mapped[str] = mapped_column(String(128), default="")
|
||||
b_port: Mapped[str] = mapped_column(String(128), default="")
|
||||
# lldp | manual | stale
|
||||
source: Mapped[str] = mapped_column(String(32), default="lldp", index=True)
|
||||
# active | stale
|
||||
status: Mapped[str] = mapped_column(String(32), default="active", index=True)
|
||||
attrs: Mapped[dict] = mapped_column(_JsonType, default=dict)
|
||||
discovered_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
last_seen_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, index=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class TopologyEdge(Base):
|
||||
"""Link between two topology nodes."""
|
||||
class TopoView(Base):
|
||||
"""Named topology view (presentation); replaces legacy topology_map."""
|
||||
|
||||
__tablename__ = "topology_edge"
|
||||
__tablename__ = "topo_view"
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
map_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
source_node_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
target_node_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
source_port: Mapped[str] = mapped_column(String(128), default="")
|
||||
target_port: Mapped[str] = mapped_column(String(128), default="")
|
||||
# manual | lldp | cdp | stale
|
||||
source: Mapped[str] = mapped_column(String(32), default="manual", index=True)
|
||||
# Optional visual overrides; empty / 0 = use provenance defaults.
|
||||
name: Mapped[str] = mapped_column(String(256), default="", index=True)
|
||||
remark: Mapped[str] = mapped_column(String(1024), default="")
|
||||
# { node_ids?: [], layer?: "physical", status?: "active", keyword?: "" }
|
||||
filter: Mapped[dict] = mapped_column(_JsonType, default=dict)
|
||||
viewport: Mapped[dict] = mapped_column(_JsonType, default=dict)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
|
||||
|
||||
|
||||
class TopoViewNode(Base):
|
||||
"""Node placement on a view; no link facts."""
|
||||
|
||||
__tablename__ = "topo_view_node"
|
||||
__table_args__ = (UniqueConstraint("view_id", "fabric_node_id", name="uq_topo_view_node"),)
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
view_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
fabric_node_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
x: Mapped[float] = mapped_column(Float, default=0.0)
|
||||
y: Mapped[float] = mapped_column(Float, default=0.0)
|
||||
label: Mapped[str] = mapped_column(String(256), default="")
|
||||
locked: Mapped[bool] = mapped_column(Boolean, default=False)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class TopoViewEdgeStyle(Base):
|
||||
"""Optional per-view edge style override."""
|
||||
|
||||
__tablename__ = "topo_view_edge_style"
|
||||
__table_args__ = (UniqueConstraint("view_id", "fabric_edge_id", name="uq_topo_view_edge_style"),)
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
view_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
fabric_edge_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
stroke_color: Mapped[str] = mapped_column(String(32), default="")
|
||||
stroke_width: Mapped[int] = mapped_column(Integer, default=0)
|
||||
# "" | solid | dashed | dotted
|
||||
line_style: Mapped[str] = mapped_column(String(16), default="")
|
||||
discovered_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class LldpCollectPolicy(Base):
|
||||
"""Singleton policy for periodic fabric LLDP collect (id=1)."""
|
||||
|
||||
__tablename__ = "lldp_collect_policy"
|
||||
|
||||
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)
|
||||
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)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class TopoDiscoverJob(Base):
|
||||
"""Async LLDP discovery job over inventory / NE ids."""
|
||||
|
||||
__tablename__ = "topo_discover_job"
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
# all_inventory | ne_ids
|
||||
scope: Mapped[str] = mapped_column(String(32), default="ne_ids", index=True)
|
||||
# manual | schedule | topology (ad-hoc from canvas)
|
||||
trigger_mode: Mapped[str] = mapped_column(String(32), default="manual", index=True)
|
||||
ne_ids_json: Mapped[list] = mapped_column(_JsonType, default=list)
|
||||
status: Mapped[str] = mapped_column(String(32), default="pending", index=True)
|
||||
total: Mapped[int] = mapped_column(Integer, default=0)
|
||||
done: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edges_added: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edges_updated: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edges_stale: Mapped[int] = mapped_column(Integer, default=0)
|
||||
error: Mapped[str] = mapped_column(String(1024), default="")
|
||||
started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
|
||||
|
||||
|
||||
class TopoDiscoverJobItem(Base):
|
||||
"""Per-NE result row for a discover job."""
|
||||
|
||||
__tablename__ = "topo_discover_job_item"
|
||||
|
||||
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
|
||||
job_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
|
||||
ume_ne_id: Mapped[str] = mapped_column(String(128), default="")
|
||||
fabric_node_id: Mapped[str] = mapped_column(String(64), default="", index=True)
|
||||
ne_name: Mapped[str] = mapped_column(String(256), default="")
|
||||
ne_ip: Mapped[str] = mapped_column(String(128), default="")
|
||||
ok: Mapped[bool] = mapped_column(Boolean, default=False)
|
||||
command: Mapped[str] = mapped_column(String(256), default="")
|
||||
neighbors: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edges_added: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edges_updated: Mapped[int] = mapped_column(Integer, default=0)
|
||||
unmatched_count: Mapped[int] = mapped_column(Integer, default=0)
|
||||
unmatched_json: Mapped[list] = mapped_column(_JsonType, default=list)
|
||||
parser_key: Mapped[str] = mapped_column(String(64), default="")
|
||||
parser_stub: Mapped[bool] = mapped_column(Boolean, default=False)
|
||||
error: Mapped[str] = mapped_column(String(1024), default="")
|
||||
raw_preview: Mapped[str] = mapped_column(Text, default="")
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class TopoFabricStats(Base):
|
||||
"""Cached fabric counters for summary (avoid COUNT on 1M edges)."""
|
||||
|
||||
__tablename__ = "topo_fabric_stats"
|
||||
|
||||
id: Mapped[str] = mapped_column(String(32), primary_key=True, default="global")
|
||||
node_count: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edge_count: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edge_active: Mapped[int] = mapped_column(Integer, default=0)
|
||||
edge_stale: Mapped[int] = mapped_column(Integer, default=0)
|
||||
last_discover_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
|
||||
class AppUser(Base):
|
||||
"""Local netx application user (login account)."""
|
||||
|
||||
|
|
|
|||
|
|
@ -147,25 +147,19 @@ def lldp_command_for_vendor(vendor: str = "", device_type: str = "") -> str:
|
|||
|
||||
|
||||
def cdp_command_for_vendor(vendor: str = "", device_type: str = "") -> str:
|
||||
"""Deprecated: CDP is not used for fabric discovery (LLDP only)."""
|
||||
return get_vendor_profile(vendor, device_type).cdp_command
|
||||
|
||||
|
||||
def pick_neighbor_command(
|
||||
*,
|
||||
protocol: str = "auto",
|
||||
protocol: str = "lldp",
|
||||
vendor: str = "",
|
||||
device_type: str = "",
|
||||
) -> tuple[str, str]:
|
||||
"""Return (command, protocol_tag).
|
||||
|
||||
Default/auto uses LLDP for all vendors (multi-vendor fabrics).
|
||||
Pass protocol=\"cdp\" only when explicitly requesting CDP (Cisco).
|
||||
"""
|
||||
proto = str(protocol or "auto").strip().lower()
|
||||
"""Return (lldp_command, \"lldp\"). Physical discovery is LLDP-only (CDP ignored)."""
|
||||
_ = protocol # accepted for call-site compat; always LLDP
|
||||
profile = get_vendor_profile(vendor, device_type)
|
||||
if proto == "cdp":
|
||||
cmd = profile.cdp_command or "show cdp neighbors detail"
|
||||
return cmd, "cdp"
|
||||
return profile.lldp_command, "lldp"
|
||||
|
||||
|
||||
|
|
@ -274,12 +268,9 @@ def parse_neighbor_output(
|
|||
raw = str(text or "")
|
||||
if not raw.strip():
|
||||
return []
|
||||
proto = str(protocol or "lldp").strip().lower()
|
||||
_ = protocol # CDP discovery removed; always parse as LLDP
|
||||
key = resolve_vendor_key(vendor, device_type)
|
||||
|
||||
if proto == "cdp":
|
||||
return parse_cisco_cdp(raw)
|
||||
|
||||
parser = _VENDOR_PARSERS.get(key) or parse_generic_lldp
|
||||
hits = parser(raw)
|
||||
if hits:
|
||||
|
|
|
|||
64
netx_api/topology_migrate.py
Normal file
64
netx_api/topology_migrate.py
Normal file
|
|
@ -0,0 +1,64 @@
|
|||
"""Startup helpers for final topology schema (drop legacy map tables)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.engine import Connection
|
||||
|
||||
_log = logging.getLogger("netx.topology.migrate")
|
||||
|
||||
_LEGACY_TABLES = ("topology_edge", "topology_node", "topology_map")
|
||||
|
||||
|
||||
def drop_legacy_topology_tables(conn: Connection) -> None:
|
||||
"""Remove document-style topology_* tables after cutover to fabric/view."""
|
||||
dialect = str(getattr(conn.dialect, "name", "") or "").lower()
|
||||
for table in _LEGACY_TABLES:
|
||||
try:
|
||||
if dialect.startswith("postgres"):
|
||||
conn.execute(text(f'DROP TABLE IF EXISTS "{table}" CASCADE'))
|
||||
else:
|
||||
conn.execute(text(f"DROP TABLE IF EXISTS {table}"))
|
||||
except Exception:
|
||||
_log.exception("drop legacy topology table failed: %s", table)
|
||||
|
||||
|
||||
def ensure_topology_schema(conn: Connection) -> None:
|
||||
drop_legacy_topology_tables(conn)
|
||||
dialect = str(getattr(conn.dialect, "name", "") or "").lower()
|
||||
# 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'"
|
||||
)
|
||||
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'"
|
||||
)
|
||||
for sql in alter_stmts:
|
||||
try:
|
||||
conn.execute(text(sql))
|
||||
except Exception:
|
||||
_log.debug("topology alter skipped/failed: %s", sql[:80], exc_info=True)
|
||||
|
||||
if not dialect.startswith("postgres"):
|
||||
return
|
||||
stmts = [
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_edge_layer_a ON topo_fabric_edge (layer, a_node_id)",
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_edge_layer_b ON topo_fabric_edge (layer, b_node_id)",
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_edge_layer_seen ON topo_fabric_edge (layer, last_seen_at)",
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_edge_active ON topo_fabric_edge (layer) WHERE status = 'active'",
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_view_node_view ON topo_view_node (view_id)",
|
||||
"CREATE UNIQUE INDEX IF NOT EXISTS uq_topo_fabric_node_managed_nn ON topo_fabric_node (managed_ne_id) WHERE managed_ne_id IS NOT NULL",
|
||||
"CREATE UNIQUE INDEX IF NOT EXISTS uq_topo_fabric_node_ume_nn ON topo_fabric_node (ume_ne_id) WHERE ume_ne_id IS NOT NULL",
|
||||
"CREATE INDEX IF NOT EXISTS ix_topo_discover_job_trigger ON topo_discover_job (trigger_mode)",
|
||||
]
|
||||
for sql in stmts:
|
||||
try:
|
||||
conn.execute(text(sql))
|
||||
except Exception:
|
||||
_log.exception("ensure topology index failed: %s", sql[:80])
|
||||
|
|
@ -1,106 +1,203 @@
|
|||
"""Topology HTTP routes."""
|
||||
"""Topology HTTP routes — fabric + views (final model)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from typing import Any, Iterator
|
||||
from typing import Any
|
||||
|
||||
from fastapi import APIRouter, Depends
|
||||
from fastapi.responses import StreamingResponse
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from .db import get_db
|
||||
from .topology_schemas import (
|
||||
TopologyDiscoverRequest,
|
||||
TopologyGraphPut,
|
||||
TopologyMapCreate,
|
||||
TopologyMapUpdate,
|
||||
FabricDiscoverRequest,
|
||||
FabricManualEdgeIn,
|
||||
TopologyViewCreate,
|
||||
TopologyViewUpdate,
|
||||
ViewEdgeStylePatch,
|
||||
ViewNodesAdd,
|
||||
ViewPositionsPatch,
|
||||
)
|
||||
from .topology_service import (
|
||||
create_map,
|
||||
delete_map,
|
||||
discover_neighbors,
|
||||
get_graph,
|
||||
iter_discover_neighbors,
|
||||
list_maps,
|
||||
put_graph,
|
||||
update_map,
|
||||
add_nodes_to_view,
|
||||
create_view,
|
||||
delete_view,
|
||||
get_discover_job,
|
||||
get_fabric_neighborhood,
|
||||
get_fabric_summary,
|
||||
get_view_graph,
|
||||
list_fabric_edges,
|
||||
list_fabric_nodes,
|
||||
list_views,
|
||||
merge_duplicate_fabric_nodes,
|
||||
patch_view_edge_style,
|
||||
patch_view_positions,
|
||||
project_fabric_neighbors_to_view,
|
||||
refresh_fabric_stats,
|
||||
remove_view_nodes,
|
||||
start_discover_job,
|
||||
update_view,
|
||||
upsert_fabric_edge,
|
||||
)
|
||||
|
||||
router = APIRouter(prefix="/v1/topology", tags=["topology"])
|
||||
|
||||
|
||||
def _sse_pack(event: dict[str, Any]) -> str:
|
||||
etype = str(event.get("type") or "message")
|
||||
return f"event: {etype}\ndata: {json.dumps(event, ensure_ascii=False, default=str)}\n\n"
|
||||
# --- Fabric -----------------------------------------------------------------
|
||||
|
||||
|
||||
@router.get("/maps")
|
||||
def api_list_maps(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return list_maps(db)
|
||||
@router.get("/fabric/summary")
|
||||
def api_fabric_summary(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_fabric_summary(db).model_dump()
|
||||
|
||||
|
||||
@router.post("/maps")
|
||||
def api_create_map(body: TopologyMapCreate, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return create_map(db, body).model_dump()
|
||||
|
||||
|
||||
@router.get("/maps/{map_id}")
|
||||
def api_get_map(map_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_graph(db, map_id).model_dump()
|
||||
|
||||
|
||||
@router.patch("/maps/{map_id}")
|
||||
def api_patch_map(
|
||||
map_id: str, body: TopologyMapUpdate, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return update_map(db, map_id, body).model_dump()
|
||||
|
||||
|
||||
@router.delete("/maps/{map_id}")
|
||||
def api_delete_map(map_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return delete_map(db, map_id)
|
||||
|
||||
|
||||
@router.put("/maps/{map_id}/graph")
|
||||
def api_put_graph(
|
||||
map_id: str, body: TopologyGraphPut, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return put_graph(db, map_id, body).model_dump()
|
||||
|
||||
|
||||
@router.post("/maps/{map_id}/discover")
|
||||
def api_discover(
|
||||
map_id: str,
|
||||
body: TopologyDiscoverRequest | None = None,
|
||||
@router.get("/fabric/nodes")
|
||||
def api_fabric_nodes(
|
||||
keyword: str = "",
|
||||
page: int = Query(default=1, ge=1),
|
||||
page_size: int = Query(default=100, ge=1, le=2000),
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
req = body or TopologyDiscoverRequest()
|
||||
return discover_neighbors(db, map_id, req).model_dump()
|
||||
return list_fabric_nodes(db, keyword=keyword, page=page, page_size=page_size)
|
||||
|
||||
|
||||
@router.post("/maps/{map_id}/discover/stream")
|
||||
def api_discover_stream(
|
||||
map_id: str,
|
||||
body: TopologyDiscoverRequest | None = None,
|
||||
@router.get("/fabric/edges")
|
||||
def api_fabric_edges(
|
||||
node_id: str = "",
|
||||
layer: str = "physical",
|
||||
status: str = "",
|
||||
source: str = "",
|
||||
page: int = Query(default=1, ge=1),
|
||||
page_size: int = Query(default=100, ge=1, le=2000),
|
||||
db: Session = Depends(get_db),
|
||||
) -> StreamingResponse:
|
||||
"""SSE stream: start → ne_start/ne_result (per NE) → done."""
|
||||
req = body or TopologyDiscoverRequest()
|
||||
|
||||
def generate() -> Iterator[str]:
|
||||
try:
|
||||
for event in iter_discover_neighbors(db, map_id, req):
|
||||
yield _sse_pack(event)
|
||||
except Exception as exc: # noqa: BLE001 — surface to client then end stream
|
||||
yield _sse_pack({"type": "error", "detail": str(exc)[:800]})
|
||||
|
||||
return StreamingResponse(
|
||||
generate(),
|
||||
media_type="text/event-stream",
|
||||
headers={
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive",
|
||||
"X-Accel-Buffering": "no",
|
||||
},
|
||||
) -> dict[str, Any]:
|
||||
return list_fabric_edges(
|
||||
db,
|
||||
node_id=node_id,
|
||||
layer=layer,
|
||||
status=status,
|
||||
source=source,
|
||||
page=page,
|
||||
page_size=page_size,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/fabric/neighborhood")
|
||||
def api_fabric_neighborhood(
|
||||
node_id: str,
|
||||
depth: int = Query(default=1, ge=1, le=3),
|
||||
layer: str = "physical",
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
return get_fabric_neighborhood(db, node_id, depth=depth, layer=layer).model_dump()
|
||||
|
||||
|
||||
@router.post("/fabric/edges")
|
||||
def api_fabric_manual_edge(
|
||||
body: FabricManualEdgeIn,
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
edge, action = upsert_fabric_edge(
|
||||
db,
|
||||
a_node_id=body.a_node_id,
|
||||
b_node_id=body.b_node_id,
|
||||
a_port=body.a_port,
|
||||
b_port=body.b_port,
|
||||
source="manual",
|
||||
)
|
||||
db.commit()
|
||||
refresh_fabric_stats(db)
|
||||
return {"ok": True, "action": action, "edge": {
|
||||
"id": edge.id,
|
||||
"a_node_id": edge.a_node_id,
|
||||
"b_node_id": edge.b_node_id,
|
||||
"a_port": edge.a_port,
|
||||
"b_port": edge.b_port,
|
||||
"source": edge.source,
|
||||
"status": edge.status,
|
||||
}}
|
||||
|
||||
|
||||
@router.post("/fabric/discover")
|
||||
def api_fabric_discover(
|
||||
body: FabricDiscoverRequest | None = None,
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
return start_discover_job(db, body or FabricDiscoverRequest()).model_dump()
|
||||
|
||||
|
||||
@router.post("/fabric/cleanup-duplicates")
|
||||
def api_fabric_cleanup_duplicates(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
"""Merge duplicate fabric nodes (same managed/ume/name/ip) and retarget edges."""
|
||||
result = merge_duplicate_fabric_nodes(db)
|
||||
return {"ok": True, **result}
|
||||
|
||||
|
||||
@router.get("/fabric/discover/{job_id}")
|
||||
def api_fabric_discover_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_discover_job(db, job_id).model_dump()
|
||||
|
||||
|
||||
# --- Views ------------------------------------------------------------------
|
||||
|
||||
|
||||
@router.get("/views")
|
||||
def api_list_views(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return list_views(db)
|
||||
|
||||
|
||||
@router.post("/views")
|
||||
def api_create_view(body: TopologyViewCreate, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return create_view(db, body).model_dump()
|
||||
|
||||
|
||||
@router.get("/views/{view_id}")
|
||||
def api_get_view(view_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return get_view_graph(db, view_id).model_dump()
|
||||
|
||||
|
||||
@router.patch("/views/{view_id}")
|
||||
def api_patch_view(
|
||||
view_id: str, body: TopologyViewUpdate, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return update_view(db, view_id, body).model_dump()
|
||||
|
||||
|
||||
@router.delete("/views/{view_id}")
|
||||
def api_delete_view(view_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return delete_view(db, view_id)
|
||||
|
||||
|
||||
@router.patch("/views/{view_id}/positions")
|
||||
def api_patch_positions(
|
||||
view_id: str, body: ViewPositionsPatch, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return patch_view_positions(db, view_id, body).model_dump()
|
||||
|
||||
|
||||
@router.post("/views/{view_id}/nodes")
|
||||
def api_add_nodes(
|
||||
view_id: str, body: ViewNodesAdd, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return add_nodes_to_view(db, view_id, body).model_dump()
|
||||
|
||||
|
||||
@router.post("/views/{view_id}/project-neighbors")
|
||||
def api_project_neighbors(view_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||
return project_fabric_neighbors_to_view(db, view_id).model_dump()
|
||||
|
||||
|
||||
@router.post("/views/{view_id}/nodes/remove")
|
||||
def api_remove_nodes(
|
||||
view_id: str,
|
||||
body: dict[str, Any],
|
||||
db: Session = Depends(get_db),
|
||||
) -> dict[str, Any]:
|
||||
ids = body.get("fabric_node_ids") if isinstance(body, dict) else None
|
||||
return remove_view_nodes(db, view_id, list(ids or [])).model_dump()
|
||||
|
||||
|
||||
@router.patch("/views/{view_id}/edge-style")
|
||||
def api_edge_style(
|
||||
view_id: str, body: ViewEdgeStylePatch, db: Session = Depends(get_db)
|
||||
) -> dict[str, Any]:
|
||||
return patch_view_edge_style(db, view_id, body).model_dump()
|
||||
|
|
|
|||
|
|
@ -1,123 +1,86 @@
|
|||
"""Pydantic schemas for topology maps / nodes / edges."""
|
||||
"""Pydantic schemas for fabric topology + views (final model)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class TopologyMapCreate(BaseModel):
|
||||
name: str = Field(min_length=1, max_length=256)
|
||||
remark: str = Field(default="", max_length=1024)
|
||||
# ---------------------------------------------------------------------------
|
||||
# Fabric
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TopologyMapUpdate(BaseModel):
|
||||
name: str | None = Field(default=None, min_length=1, max_length=256)
|
||||
remark: str | None = Field(default=None, max_length=1024)
|
||||
|
||||
|
||||
class TopologyMapOut(BaseModel):
|
||||
class FabricNodeOut(BaseModel):
|
||||
id: str
|
||||
name: str
|
||||
remark: str
|
||||
managed_ne_id: str = ""
|
||||
ume_ne_id: str = ""
|
||||
name: str = ""
|
||||
ip: str = ""
|
||||
vendor: str = ""
|
||||
device_type: str = ""
|
||||
attrs: dict[str, Any] = Field(default_factory=dict)
|
||||
last_seen_at: datetime | None = None
|
||||
|
||||
|
||||
class FabricEdgeOut(BaseModel):
|
||||
id: str
|
||||
layer: str = "physical"
|
||||
a_node_id: str
|
||||
b_node_id: str
|
||||
a_port: str = ""
|
||||
b_port: 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
|
||||
|
||||
|
||||
class FabricSummaryOut(BaseModel):
|
||||
node_count: int = 0
|
||||
edge_count: int = 0
|
||||
created_at: datetime | None = None
|
||||
edge_active: int = 0
|
||||
edge_stale: int = 0
|
||||
last_discover_at: datetime | None = None
|
||||
updated_at: datetime | None = None
|
||||
|
||||
|
||||
class TopologyNodeIn(BaseModel):
|
||||
id: str = Field(min_length=1, max_length=64)
|
||||
managed_ne_id: str = ""
|
||||
ume_ne_id: str = ""
|
||||
label: str = ""
|
||||
x: float = 0.0
|
||||
y: float = 0.0
|
||||
created_at: datetime | None = None
|
||||
class FabricNeighborhoodOut(BaseModel):
|
||||
center_node_id: str
|
||||
depth: int = 1
|
||||
nodes: list[FabricNodeOut] = Field(default_factory=list)
|
||||
edges: list[FabricEdgeOut] = Field(default_factory=list)
|
||||
|
||||
|
||||
class TopologyEdgeIn(BaseModel):
|
||||
id: str = Field(min_length=1, max_length=64)
|
||||
source_node_id: str = Field(min_length=1, max_length=64)
|
||||
target_node_id: str = Field(min_length=1, max_length=64)
|
||||
source_port: str = ""
|
||||
target_port: str = ""
|
||||
source: str = "manual"
|
||||
stroke_color: str = Field(default="", max_length=32)
|
||||
stroke_width: int = Field(default=0, ge=0, le=12)
|
||||
line_style: str = Field(default="", max_length=16)
|
||||
discovered_at: datetime | None = None
|
||||
created_at: datetime | None = None
|
||||
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)
|
||||
auto_add_unmatched: bool = Field(
|
||||
default=True,
|
||||
description="Create SSH placeholder ManagedNEs for LLDP neighbors not in inventory",
|
||||
)
|
||||
concurrency: int = Field(default=4, ge=1, le=32)
|
||||
trigger_mode: str = Field(default="manual", description="manual | schedule | topology")
|
||||
|
||||
|
||||
class TopologyNodeOut(BaseModel):
|
||||
id: str
|
||||
map_id: str
|
||||
managed_ne_id: str = ""
|
||||
ume_ne_id: str = ""
|
||||
label: str = ""
|
||||
x: float = 0.0
|
||||
y: float = 0.0
|
||||
ne_name: str = ""
|
||||
ne_ip: str = ""
|
||||
vendor: str = ""
|
||||
protocol: str = ""
|
||||
connect_status: str = ""
|
||||
|
||||
|
||||
class TopologyEdgeOut(BaseModel):
|
||||
id: str
|
||||
map_id: str
|
||||
source_node_id: str
|
||||
target_node_id: str
|
||||
source_port: str = ""
|
||||
target_port: str = ""
|
||||
source: str = "manual"
|
||||
stroke_color: str = ""
|
||||
stroke_width: int = 0
|
||||
line_style: str = ""
|
||||
discovered_at: datetime | None = None
|
||||
|
||||
|
||||
class TopologyGraphOut(BaseModel):
|
||||
map: TopologyMapOut
|
||||
nodes: list[TopologyNodeOut]
|
||||
edges: list[TopologyEdgeOut]
|
||||
|
||||
|
||||
class TopologyGraphPut(BaseModel):
|
||||
nodes: list[TopologyNodeIn] = Field(default_factory=list)
|
||||
edges: list[TopologyEdgeIn] = Field(default_factory=list)
|
||||
|
||||
|
||||
class TopologyDiscoverRequest(BaseModel):
|
||||
"""Run LLDP/CDP discovery for managed NEs currently on the map."""
|
||||
|
||||
protocol: str = Field(default="auto", description="auto | lldp | cdp")
|
||||
ne_ids: list[str] | None = None
|
||||
|
||||
|
||||
class TopologyDiscoverUnmatched(BaseModel):
|
||||
class FabricDiscoverUnmatched(BaseModel):
|
||||
remote_name: str = ""
|
||||
remote_ip: str = ""
|
||||
local_port: str = ""
|
||||
remote_port: str = ""
|
||||
|
||||
|
||||
class TopologyDiscoverLink(BaseModel):
|
||||
peer_node_id: str = ""
|
||||
peer_ne_id: str = ""
|
||||
peer_name: str = ""
|
||||
peer_ip: str = ""
|
||||
local_port: str = ""
|
||||
remote_port: str = ""
|
||||
protocol: str = ""
|
||||
action: str = "" # added | updated | kept_manual
|
||||
|
||||
|
||||
class TopologyDiscoverNeResult(BaseModel):
|
||||
ne_id: str
|
||||
class FabricDiscoverJobItemOut(BaseModel):
|
||||
id: str
|
||||
job_id: str
|
||||
ne_id: str = ""
|
||||
ume_ne_id: str = ""
|
||||
fabric_node_id: str = ""
|
||||
ne_name: str = ""
|
||||
ne_ip: str = ""
|
||||
ok: bool = False
|
||||
|
|
@ -126,20 +89,130 @@ class TopologyDiscoverNeResult(BaseModel):
|
|||
edges_added: int = 0
|
||||
edges_updated: int = 0
|
||||
unmatched_count: int = 0
|
||||
unmatched: list[TopologyDiscoverUnmatched] = Field(default_factory=list)
|
||||
links: list[TopologyDiscoverLink] = Field(default_factory=list)
|
||||
unmatched: list[FabricDiscoverUnmatched] = Field(default_factory=list)
|
||||
parser_key: str = ""
|
||||
parser_stub: bool = False
|
||||
error: str = ""
|
||||
raw_preview: str = ""
|
||||
|
||||
|
||||
class TopologyDiscoverOut(BaseModel):
|
||||
map_id: str
|
||||
protocol: str
|
||||
scanned: int = 0
|
||||
class FabricDiscoverJobOut(BaseModel):
|
||||
id: str
|
||||
scope: str
|
||||
trigger_mode: str = "manual"
|
||||
status: str
|
||||
total: int = 0
|
||||
done: int = 0
|
||||
edges_added: int = 0
|
||||
edges_updated: int = 0
|
||||
edges_stale: int = 0
|
||||
results: list[TopologyDiscoverNeResult] = Field(default_factory=list)
|
||||
graph: TopologyGraphOut | None = None
|
||||
error: str = ""
|
||||
started_at: datetime | None = None
|
||||
ended_at: datetime | None = None
|
||||
items: list[FabricDiscoverJobItemOut] = Field(default_factory=list)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Views
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TopologyViewCreate(BaseModel):
|
||||
name: str = Field(min_length=1, max_length=256)
|
||||
remark: str = Field(default="", max_length=1024)
|
||||
filter: dict[str, Any] = Field(default_factory=dict)
|
||||
|
||||
|
||||
class TopologyViewUpdate(BaseModel):
|
||||
name: str | None = Field(default=None, min_length=1, max_length=256)
|
||||
remark: str | None = Field(default=None, max_length=1024)
|
||||
filter: dict[str, Any] | None = None
|
||||
viewport: dict[str, Any] | None = None
|
||||
|
||||
|
||||
class TopologyViewOut(BaseModel):
|
||||
id: str
|
||||
name: str
|
||||
remark: str = ""
|
||||
filter: dict[str, Any] = Field(default_factory=dict)
|
||||
viewport: dict[str, Any] = Field(default_factory=dict)
|
||||
node_count: int = 0
|
||||
created_at: datetime | None = None
|
||||
updated_at: datetime | None = None
|
||||
|
||||
|
||||
class ViewNodeIn(BaseModel):
|
||||
fabric_node_id: str = Field(min_length=1, max_length=64)
|
||||
x: float = 0.0
|
||||
y: float = 0.0
|
||||
label: str = ""
|
||||
locked: bool = False
|
||||
|
||||
|
||||
class ViewNodeOut(BaseModel):
|
||||
fabric_node_id: str
|
||||
managed_ne_id: str = ""
|
||||
ume_ne_id: str = ""
|
||||
label: str = ""
|
||||
x: float = 0.0
|
||||
y: float = 0.0
|
||||
locked: bool = False
|
||||
name: str = ""
|
||||
ip: str = ""
|
||||
vendor: str = ""
|
||||
device_type: str = ""
|
||||
connect_status: str = ""
|
||||
|
||||
|
||||
class ViewEdgeOut(BaseModel):
|
||||
id: str
|
||||
a_node_id: str
|
||||
b_node_id: str
|
||||
a_port: str = ""
|
||||
b_port: str = ""
|
||||
source: str = "lldp"
|
||||
status: str = "active"
|
||||
layer: str = "physical"
|
||||
stroke_color: str = ""
|
||||
stroke_width: int = 0
|
||||
line_style: str = ""
|
||||
discovered_at: datetime | None = None
|
||||
|
||||
|
||||
class TopologyViewGraphOut(BaseModel):
|
||||
view: TopologyViewOut
|
||||
nodes: list[ViewNodeOut]
|
||||
edges: list[ViewEdgeOut]
|
||||
truncated: bool = False
|
||||
truncate_reason: str = ""
|
||||
|
||||
|
||||
class ViewPositionsPatch(BaseModel):
|
||||
positions: list[ViewNodeIn] = Field(default_factory=list)
|
||||
|
||||
|
||||
class ViewNodesAdd(BaseModel):
|
||||
"""Add inventory NEs onto a view (creates fabric nodes as needed)."""
|
||||
|
||||
managed_ne_ids: list[str] = Field(default_factory=list)
|
||||
ume_ne_ids: list[str] = Field(default_factory=list)
|
||||
fabric_node_ids: list[str] = Field(
|
||||
default_factory=list,
|
||||
description="Place existing fabric nodes onto the view",
|
||||
)
|
||||
# Optional initial positions keyed by managed/ume id
|
||||
layout: str = Field(default="grid", description="grid | keep")
|
||||
|
||||
|
||||
class ViewEdgeStylePatch(BaseModel):
|
||||
fabric_edge_id: str
|
||||
stroke_color: str = ""
|
||||
stroke_width: int = Field(default=0, ge=0, le=12)
|
||||
line_style: str = Field(default="", max_length=16)
|
||||
|
||||
|
||||
class FabricManualEdgeIn(BaseModel):
|
||||
a_node_id: str = Field(min_length=1, max_length=64)
|
||||
b_node_id: str = Field(min_length=1, max_length=64)
|
||||
a_port: str = ""
|
||||
b_port: str = ""
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
Loading…
Add table
Add a link
Reference in a new issue