Ship gated eye polish, fabric levels, and collection UI refresh.

Topology MCP adds pull/compact/bundle/suggest-hubs with a no-template skill path; API gains fabric level and NE collection policy; web list pages get paging and denser collect/network workflows.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-12 16:13:05 +08:00
parent b1aee86701
commit b81e5869a6
91 changed files with 10395 additions and 1513 deletions

View file

@ -31,6 +31,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None:
except Exception: # noqa: BLE001
_log.exception("stop_lldp_collect_scheduler failed")
try:
from .ne_collect_scheduler import stop_ne_collect_scheduler
stop_ne_collect_scheduler()
except Exception: # noqa: BLE001
_log.exception("stop_ne_collect_scheduler failed")
try:
from .port_traffic_scheduler import stop_port_traffic_scheduler

View file

@ -115,11 +115,17 @@ def run_api_startup() -> None:
if cfg_resumed:
_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed)
try:
from .collection_policy import ensure_policy as ensure_ne_collect_policy
from .collection_policy import history_keep_value, prune_collection_jobs
from .lldp_collect_service import ensure_policy as ensure_lldp_collect_policy
ensure_lldp_collect_policy(db)
pol = ensure_ne_collect_policy(db)
pruned_jobs = prune_collection_jobs(db, keep=history_keep_value(pol))
if pruned_jobs:
_log.info("startup: pruned %s ne_collection job(s) by history_keep", pruned_jobs)
except Exception:
_log.exception("startup: lldp_collect policy ensure failed")
_log.exception("startup: lldp/ne_collect policy ensure failed")
pt_cleared = recover_port_traffic_on_startup(db)
if pt_cleared:
_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared)
@ -151,7 +157,7 @@ def run_api_startup() -> None:
else:
_log.info(
"startup: inline schedulers disabled — run `python -m netx_api.worker` for "
"config_sync / lldp_collect / port_traffic"
"config_sync / lldp_collect / ne_collect / port_traffic"
)
start_api_sideband_threads()

View file

@ -48,6 +48,12 @@ def finalize_collection_job(db: Session, job_id: str) -> None:
job.ended_at = finish_at
job.last_run_at = job.ended_at or finish_at
db.commit()
try:
from .collection_policy import ensure_policy, history_keep_value, prune_collection_jobs
prune_collection_jobs(db, keep=history_keep_value(ensure_policy(db)))
except Exception: # noqa: BLE001
pass
def reconcile_stale_collection_job(db: Session, job_id: str) -> bool:

View file

@ -0,0 +1,265 @@
"""NE batch-collect policy, prune-by-count, and schedule due helpers."""
from __future__ import annotations
import logging
import shutil
from datetime import datetime, timedelta
from typing import Any
from fastapi import HTTPException
from sqlalchemy.orm import Session
from .collection_schemas import (
CollectionPolicyOut,
CollectionPolicyUpdate,
CollectionTargetRef,
)
from .models import ManagedNE, NeCollectionJob, NeCollectionPolicy, NeCollectionRun, UmeInventoryNE
from .ne_collection_paths import collection_data_root
from .timeutil import utcnow_naive # used by ensure_policy.updated_at
_log = logging.getLogger("netx.collection.policy")
POLICY_ID = 1
DEFAULT_HISTORY_KEEP = 3
MAX_INTERVAL_HOURS = 8760 # 365d
def _utcnow() -> datetime:
return datetime.now()
def _normalize_interval_hours(row: NeCollectionPolicy) -> 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) -> NeCollectionPolicy:
row = db.get(NeCollectionPolicy, POLICY_ID)
if row is None:
row = NeCollectionPolicy(
id=POLICY_ID,
enabled=False,
interval_days=1,
interval_hours=24,
scope_mode="all",
selected_targets=[],
title="",
commands="",
history_keep=DEFAULT_HISTORY_KEEP,
updated_at=_utcnow(),
)
db.add(row)
db.commit()
db.refresh(row)
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
def _targets_from_json(raw: Any) -> list[CollectionTargetRef]:
items: list[CollectionTargetRef] = []
if not isinstance(raw, list):
return items
for x in raw:
if not isinstance(x, dict):
continue
tid = str(x.get("id") or "").strip()
if not tid:
continue
src = str(x.get("source") or "managed").strip().lower() or "managed"
if src not in {"managed", "ume"}:
src = "managed"
items.append(CollectionTargetRef(source=src, id=tid))
return items
def policy_to_out(row: NeCollectionPolicy) -> CollectionPolicyOut:
hours = _normalize_interval_hours(row)
days = max(1, min(365, (hours + 23) // 24))
keep = getattr(row, "history_keep", None)
if keep is None:
keep = DEFAULT_HISTORY_KEEP
return CollectionPolicyOut(
enabled=bool(row.enabled),
interval_days=days,
interval_hours=hours,
scope_mode="selected" if str(row.scope_mode or "") == "selected" else "all",
selected_targets=_targets_from_json(row.selected_targets),
title=str(row.title or ""),
commands=str(row.commands or ""),
history_keep=max(0, min(200, int(keep))),
updated_at=row.updated_at,
)
def get_policy(db: Session) -> CollectionPolicyOut:
return policy_to_out(ensure_policy(db))
def history_keep_value(row: NeCollectionPolicy | None = None) -> int:
if row is None:
return DEFAULT_HISTORY_KEEP
keep = getattr(row, "history_keep", None)
if keep is None:
keep = DEFAULT_HISTORY_KEEP
return max(0, min(200, int(keep)))
def prune_collection_jobs(db: Session, *, keep: int = DEFAULT_HISTORY_KEEP) -> int:
"""Delete finished jobs beyond ``keep`` (newest kept). Active jobs always retained."""
keep = max(0, min(200, int(keep)))
finished = (
db.query(NeCollectionJob)
.filter(NeCollectionJob.status.in_(("done", "failed")))
.order_by(NeCollectionJob.created_at.desc())
.all()
)
to_drop = finished if keep == 0 else finished[keep:]
if not to_drop:
return 0
root = collection_data_root().resolve()
dropped = 0
for job in to_drop:
jid = str(job.id)
db.query(NeCollectionRun).filter(NeCollectionRun.job_id == jid).delete(
synchronize_session=False
)
db.delete(job)
dropped += 1
job_dir = (root / jid).resolve()
if str(job_dir).startswith(str(root)) and job_dir.is_dir():
shutil.rmtree(job_dir, ignore_errors=True)
if dropped:
db.commit()
_log.info("pruned %s ne_collection job(s); keep=%s", dropped, keep)
return dropped
def update_policy(db: Session, body: CollectionPolicyUpdate) -> CollectionPolicyOut:
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_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 "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, CollectionTargetRef):
tid = ref.id.strip()
src = (ref.source or "managed").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 "title" in data and data["title"] is not None:
row.title = str(data["title"] or "").strip()[:256]
if "commands" in data and data["commands"] is not None:
row.commands = str(data["commands"] or "")
if "history_keep" in data and data["history_keep"] is not None:
row.history_keep = max(0, min(200, int(data["history_keep"])))
if bool(row.enabled):
from .collection_service import _parse_commands
if not _parse_commands(str(row.commands or "")):
raise HTTPException(status_code=400, detail="commands_required_for_schedule")
if str(row.scope_mode or "") == "selected" and not (row.selected_targets or []):
raise HTTPException(status_code=400, detail="no_selected_targets")
row.updated_at = _utcnow()
db.commit()
db.refresh(row)
prune_collection_jobs(db, keep=history_keep_value(row))
return policy_to_out(row)
def next_due_at(db: Session, policy: NeCollectionPolicy | None = None) -> datetime | None:
"""Due time based on last *scheduled* successful collect only (manual must not reset)."""
pol = policy or ensure_policy(db)
if not pol.enabled:
return None
hours = _normalize_interval_hours(pol)
last = (
db.query(NeCollectionJob)
.filter(
NeCollectionJob.status == "done",
NeCollectionJob.trigger_mode == "schedule",
NeCollectionJob.ended_at.isnot(None),
)
.order_by(NeCollectionJob.ended_at.desc())
.first()
)
if last is None or last.ended_at is None:
return _utcnow()
return last.ended_at + timedelta(hours=hours)
def expand_policy_targets(db: Session, policy: NeCollectionPolicy) -> list[tuple[str, str, str, str]]:
"""Return list of (source, id, name, ip) for a policy."""
from .device_types import WEBCRT_NE_SOURCE
mode = str(policy.scope_mode or "all").strip().lower()
out: list[tuple[str, str, str, str]] = []
seen: set[tuple[str, str]] = set()
def _add(source: str, tid: str, name: str, ip: str) -> None:
key = (source, tid)
if key in seen:
return
seen.add(key)
out.append((source, tid, name, ip))
if mode == "selected":
for ref in _targets_from_json(policy.selected_targets):
if ref.source == "managed":
ne = db.get(ManagedNE, ref.id)
if ne:
_add(
"managed",
str(ne.id),
str(ne.name or ne.ip_address or ""),
str(ne.ip_address or ""),
)
else:
inv = db.get(UmeInventoryNE, ref.id)
if inv:
name = str(inv.host_name or inv.user_label or inv.ne_name or inv.ip_address or inv.ne_id)
_add("ume", str(inv.ne_id), name, str(inv.ip_address or ""))
return out
for ne in (
db.query(ManagedNE)
.filter(ManagedNE.source != WEBCRT_NE_SOURCE)
.order_by(ManagedNE.name.asc())
.all()
):
_add("managed", str(ne.id), str(ne.name or ne.ip_address or ""), str(ne.ip_address or ""))
for inv in db.query(UmeInventoryNE).order_by(UmeInventoryNE.host_name.asc()).all():
if not str(inv.ip_address or "").strip():
continue
name = str(inv.host_name or inv.user_label or inv.ne_name or inv.ip_address or inv.ne_id)
_add("ume", str(inv.ne_id), name, str(inv.ip_address or ""))
return out

View file

@ -4,8 +4,11 @@ from fastapi import APIRouter, BackgroundTasks, Depends, Query
from fastapi.responses import FileResponse, Response
from sqlalchemy.orm import Session
from .collection_policy import get_policy, update_policy
from .collection_schemas import CollectionJobCreate, CollectionPolicyUpdate
from .collection_service import (
build_collection_job_zip,
create_and_start_from_policy,
create_collection,
delete_collection_job,
get_collection_dashboard,
@ -19,7 +22,6 @@ from .collection_service import (
start_collection_job,
retry_failed_collection_job,
)
from .collection_schemas import CollectionJobCreate
from .db import get_db
from .models import NeCollectionRun
from .ne_collect_runner import dispatch_collection_runs
@ -42,6 +44,32 @@ def api_collection_dashboard(db: Session = Depends(get_db)):
return get_collection_dashboard(db).model_dump()
@router.get("/policy")
def api_get_collection_policy(db: Session = Depends(get_db)):
return get_policy(db).model_dump()
@router.put("/policy")
def api_put_collection_policy(body: CollectionPolicyUpdate, db: Session = Depends(get_db)):
return update_policy(db, body).model_dump()
@router.post("/start-from-policy")
def api_start_from_policy(
background_tasks: BackgroundTasks,
db: Session = Depends(get_db),
):
"""Create + start a collect job from the saved policy (manual trigger)."""
out, payload = create_and_start_from_policy(db, trigger_mode="manual")
background_tasks.add_task(
dispatch_collection_runs,
payload["job_id"],
payload["run_ids"],
payload["commands"],
)
return out.model_dump()
@router.post("")
def api_create_collection(body: CollectionJobCreate, db: Session = Depends(get_db)):
return create_collection(db, body).model_dump()
@ -51,9 +79,13 @@ def api_create_collection(body: CollectionJobCreate, db: Session = Depends(get_d
def api_list_collections(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
):
return list_collection_jobs(db, page=page, page_size=page_size)
return list_collection_jobs(
db, page=page, page_size=page_size, status=status, keyword=keyword
)
@router.get("/runs/{run_id}/download")

View file

@ -5,6 +5,34 @@ from datetime import datetime
from pydantic import BaseModel, Field
class CollectionTargetRef(BaseModel):
source: str = "managed" # managed | ume
id: str
class CollectionPolicyOut(BaseModel):
enabled: bool = False
interval_days: int = 1
interval_hours: int = 24
scope_mode: str = "all"
selected_targets: list[CollectionTargetRef] = Field(default_factory=list)
title: str = ""
commands: str = ""
history_keep: int = 3
updated_at: datetime | None = None
class CollectionPolicyUpdate(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)
scope_mode: str | None = None
selected_targets: list[CollectionTargetRef] | None = None
title: str | None = None
commands: str | None = None
history_keep: int | None = Field(default=None, ge=0, le=200)
class CollectionJobCreate(BaseModel):
title: str = ""
commands: str = Field(min_length=1)
@ -31,6 +59,7 @@ class CollectionJobOut(BaseModel):
id: str
title: str
commands: str
trigger_mode: str = "manual"
status: str
ne_count: int
success_count: int
@ -47,6 +76,7 @@ class CollectionJobSummary(BaseModel):
id: str
title: str
status: str
trigger_mode: str = "manual"
ne_count: int
success_count: int
fail_count: int
@ -61,3 +91,5 @@ class CollectionDashboardOut(BaseModel):
active_count: int = 0
running_job: CollectionJobSummary | None = None
last_job: CollectionJobSummary | None = None
next_due_at: datetime | None = None
policy: CollectionPolicyOut | None = None

View file

@ -20,6 +20,14 @@ from .collection_job_state import (
sync_job_progress,
_sync_job_counts,
)
from .collection_policy import (
ensure_policy,
expand_policy_targets,
history_keep_value,
next_due_at,
policy_to_out,
prune_collection_jobs,
)
from .collection_schemas import (
CollectionDashboardOut,
CollectionJobCreate,
@ -87,6 +95,7 @@ def job_to_out(row: NeCollectionJob, *, output_count: int | None = None) -> Coll
id=str(row.id),
title=str(row.title or ""),
commands=str(row.commands or ""),
trigger_mode=str(getattr(row, "trigger_mode", None) or "manual"),
status=str(row.status or "pending"),
ne_count=int(row.ne_count or 0),
success_count=int(row.success_count or 0),
@ -107,6 +116,7 @@ def job_to_summary(row: NeCollectionJob | None) -> CollectionJobSummary | None:
id=str(row.id),
title=str(row.title or "").strip() or str(row.id)[:8],
status=str(row.status or "pending"),
trigger_mode=str(getattr(row, "trigger_mode", None) or "manual"),
ne_count=int(row.ne_count or 0),
success_count=int(row.success_count or 0),
fail_count=int(row.fail_count or 0),
@ -162,11 +172,14 @@ def get_collection_dashboard(db: Session) -> CollectionDashboardOut:
or 0
)
last = last_finished_collection_job(db)
policy = ensure_policy(db)
return CollectionDashboardOut(
job_count=job_count,
active_count=active_count,
running_job=job_to_summary(running),
last_job=job_to_summary(last),
next_due_at=next_due_at(db, policy),
policy=policy_to_out(policy),
)
@ -318,6 +331,7 @@ def create_collection(db: Session, body: CollectionJobCreate) -> CollectionJobOu
job = NeCollectionJob(
title=str(body.title or "").strip() or f"collect-{now.strftime('%Y%m%d-%H%M%S')}",
commands="\n".join(commands),
trigger_mode="manual",
status="pending",
ne_count=len(targets),
created_at=now,
@ -342,8 +356,81 @@ def create_collection(db: Session, body: CollectionJobCreate) -> CollectionJobOu
return job_to_out(job, output_count=0)
def list_collection_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict[str, Any]:
def create_and_start_from_policy(
db: Session,
*,
trigger_mode: str = "manual",
) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
"""Create a job from the singleton policy and start it immediately."""
if has_active_collection_job(db) is not None:
raise HTTPException(status_code=409, detail="collection_job_running")
policy = ensure_policy(db)
commands = _parse_commands(str(policy.commands or ""))
if not commands:
raise HTTPException(status_code=400, detail="commands_empty")
targets = expand_policy_targets(db, policy)
if not targets:
raise HTTPException(status_code=400, detail="no_eligible_ne")
mode = str(trigger_mode or "manual").strip().lower() or "manual"
if mode not in {"manual", "schedule"}:
mode = "manual"
now = _now()
title = str(policy.title or "").strip() or f"collect-{now.strftime('%Y%m%d-%H%M%S')}"
job = NeCollectionJob(
title=title,
commands="\n".join(commands),
trigger_mode=mode,
status="pending",
ne_count=len(targets),
created_at=now,
started_at=None,
last_run_at=None,
)
db.add(job)
db.flush()
for source, tid, name, ip in targets:
db.add(
NeCollectionRun(
job_id=str(job.id),
ne_id=tid,
ne_source=source,
ne_name=name,
ne_ip=ip,
status="pending",
)
)
db.commit()
db.refresh(job)
out, payload = start_collection_job(db, str(job.id))
try:
prune_collection_jobs(db, keep=history_keep_value(policy))
except Exception: # noqa: BLE001
_log.exception("prune_collection_jobs after start failed")
return out, payload
def list_collection_jobs(
db: Session,
*,
page: int = 1,
page_size: int = 20,
status: str = "",
keyword: str = "",
) -> dict[str, Any]:
stmt = db.query(NeCollectionJob)
st = str(status or "").strip()
if st:
stmt = stmt.filter(NeCollectionJob.status == st)
kw = str(keyword or "").strip()
if kw:
like = f"%{kw}%"
stmt = stmt.filter(
or_(
NeCollectionJob.title.ilike(like),
NeCollectionJob.id.ilike(like),
NeCollectionJob.error_message.ilike(like),
)
)
total = int(stmt.count())
rows = stmt.order_by(NeCollectionJob.created_at.desc()).offset((page - 1) * page_size).limit(page_size).all()
for row in rows:

View file

@ -81,6 +81,10 @@ class Settings(BaseSettings):
ne_collect_pending_stale_sec: int = 180
ne_collect_run_timeout_cap_sec: int = 600
ne_collection_data_dir: str = "data/ne_collections"
# Batch CLI collect scheduler (policy.enabled defaults False — manual until turned on).
ne_collect_scheduler_enabled: bool = True
ne_collect_scheduler_tick_sec: int = 60
ne_collect_startup_grace_sec: int = 3600
# Config sync (periodic running-config backup into DB)
config_sync_scheduler_enabled: bool = True
config_sync_scheduler_tick_sec: int = 60

View file

@ -6,7 +6,7 @@ from typing import Any
from uuid import uuid4
from fastapi import HTTPException
from sqlalchemy import func
from sqlalchemy import func, or_
from sqlalchemy.orm import Session
from .cli_resolve import cli_profile_ready
@ -252,8 +252,29 @@ def create_cycle(db: Session, body: ConfigSyncCycleCreate) -> ConfigSyncCycleOut
return cycle_to_out(cycle)
def list_cycles(db: Session, *, page: int, page_size: int) -> dict[str, Any]:
q = db.query(ConfigSyncCycle).order_by(ConfigSyncCycle.created_at.desc())
def list_cycles(
db: Session,
*,
page: int,
page_size: int,
status: str = "",
keyword: str = "",
) -> dict[str, Any]:
q = db.query(ConfigSyncCycle)
st = str(status or "").strip()
if st:
q = q.filter(ConfigSyncCycle.status == st)
kw = str(keyword or "").strip()
if kw:
like = f"%{kw}%"
q = q.filter(
or_(
ConfigSyncCycle.id.ilike(like),
ConfigSyncCycle.trigger_mode.ilike(like),
ConfigSyncCycle.error_message.ilike(like),
)
)
q = q.order_by(ConfigSyncCycle.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": [cycle_to_out(r) for r in rows]}

View file

@ -74,9 +74,11 @@ def api_dashboard(db: Session = Depends(get_db)):
def api_list_cycles(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
):
return list_cycles(db, page=page, page_size=page_size)
return list_cycles(db, page=page, page_size=page_size, status=status, keyword=keyword)
@router.post("/cycles")

View file

@ -67,9 +67,11 @@ def api_stop_job(job_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
def api_list_jobs(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
) -> dict[str, Any]:
return list_jobs(db, page=page, page_size=page_size)
return list_jobs(db, page=page, page_size=page_size, status=status, keyword=keyword)
@router.get("/jobs/{job_id}")
@ -77,6 +79,15 @@ def api_get_job(
job_id: str,
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
) -> dict[str, Any]:
return get_job_detail(db, job_id, page=page, page_size=page_size)
return get_job_detail(
db,
job_id,
page=page,
page_size=page_size,
item_status=status,
item_keyword=keyword,
)

View file

@ -481,10 +481,32 @@ def get_dashboard(db: Session) -> LldpCollectDashboardOut:
)
def list_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict:
def list_jobs(
db: Session,
*,
page: int = 1,
page_size: int = 20,
status: str = "",
keyword: str = "",
) -> 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())
q = db.query(TopoDiscoverJob)
st = str(status or "").strip()
if st:
q = q.filter(TopoDiscoverJob.status == st)
kw = str(keyword or "").strip()
if kw:
like = f"%{kw}%"
q = q.filter(
or_(
TopoDiscoverJob.id.ilike(like),
TopoDiscoverJob.error.ilike(like),
TopoDiscoverJob.scope.ilike(like),
TopoDiscoverJob.trigger_mode.ilike(like),
)
)
q = q.order_by(TopoDiscoverJob.created_at.desc())
total = int(q.count())
rows = q.offset((page - 1) * page_size).limit(page_size).all()
outcomes = _job_outcome_map(db, [r.id for r in rows])
@ -503,6 +525,19 @@ def list_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict:
def get_job_detail(
db: Session, job_id: str, *, page: int | None = None, page_size: int | None = None
db: Session,
job_id: str,
*,
page: int | None = None,
page_size: int | None = None,
item_status: str = "",
item_keyword: str = "",
) -> dict:
return get_discover_job(db, job_id, page=page, page_size=page_size).model_dump()
return get_discover_job(
db,
job_id,
page=page,
page_size=page_size,
item_status=item_status,
item_keyword=item_keyword,
).model_dump()

View file

@ -111,6 +111,7 @@ def _prom_lines(metrics: dict[str, Any]) -> str:
for name, key in (
("config_sync", "netx_config_sync_scheduler_running"),
("lldp_collect", "netx_lldp_collect_scheduler_running"),
("ne_collect", "netx_ne_collect_scheduler_running"),
("port_traffic", "netx_port_traffic_scheduler_running"),
("fabric_reconcile", "netx_fabric_reconcile_scheduler_running"),
):

View file

@ -20,6 +20,7 @@ from .managed_ne import (
CliConnectProfile,
ManagedNE,
NeCollectionJob,
NeCollectionPolicy,
NeCollectionRun,
UmeCliOverride,
)
@ -82,6 +83,7 @@ __all__ = [
"CliConnectProfile",
"UmeCliOverride",
"NeCollectionJob",
"NeCollectionPolicy",
"NeCollectionRun",
"TopoFabricNode",
"TopoClassifyRule",

View file

@ -98,6 +98,28 @@ class UmeCliOverride(Base):
updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True)
class NeCollectionPolicy(Base):
"""Singleton policy for periodic batch CLI collect (id=1).
Default ``enabled=False``: one-shot / manual only until the operator turns on schedule.
``history_keep`` defaults to 3 finished jobs.
"""
__tablename__ = "ne_collection_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) # legacy mirror of hours
interval_hours: Mapped[int] = mapped_column(Integer, default=24)
scope_mode: Mapped[str] = mapped_column(String(32), default="all") # all | selected
selected_targets: Mapped[list] = mapped_column(_JsonType, default=list)
title: Mapped[str] = mapped_column(String(256), default="")
commands: Mapped[str] = mapped_column(Text, default="")
# Finished jobs to retain (newest kept); active jobs always kept.
history_keep: Mapped[int] = mapped_column(Integer, default=3)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)
class NeCollectionJob(Base):
"""Batch CLI collection job over managed NEs."""
@ -106,6 +128,8 @@ class NeCollectionJob(Base):
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
title: Mapped[str] = mapped_column(String(256), default="")
commands: Mapped[str] = mapped_column(Text, default="")
# manual | schedule
trigger_mode: Mapped[str] = mapped_column(String(32), default="manual", index=True)
status: Mapped[str] = mapped_column(String(32), default="pending", index=True)
ne_count: Mapped[int] = mapped_column(Integer, default=0)
success_count: Mapped[int] = mapped_column(Integer, default=0)

View file

@ -26,10 +26,12 @@ class TopoFabricNode(Base):
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="")
# Classify tags (regex rules / manual). role: core|aggregation|access|unknown|""
# Layout rank: major.minor (0=external … 1=core … 2=agg … 3=access). None = unclassified.
level: Mapped[float | None] = mapped_column(Float, nullable=True, index=True)
# Synced alias from floor(level): external|core|aggregation|access|edge|""
role: Mapped[str] = mapped_column(String(32), default="", index=True)
region_folder_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
# rule | manual | ""
# rule | manual | "" (applies to level; role is derived)
role_source: Mapped[str] = mapped_column(String(16), default="")
region_source: Mapped[str] = mapped_column(String(16), default="")
# Composed flat-world coordinates (packed per-SBN local layouts). Not raw UME xPos/yPos.
@ -47,15 +49,15 @@ class TopoClassifyRule(Base):
__tablename__ = "topo_classify_rule"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
# role | region
scope: Mapped[str] = mapped_column(String(16), default="role", index=True)
# level | region (legacy scope "role" accepted as alias of level)
scope: Mapped[str] = mapped_column(String(16), default="level", index=True)
name: Mapped[str] = mapped_column(String(256), default="")
pattern: Mapped[str] = mapped_column(String(512), default="")
# name | ip | name_ip
match_field: Mapped[str] = mapped_column(String(32), default="name")
priority: Mapped[int] = mapped_column(Integer, default=100, index=True)
enabled: Mapped[bool] = mapped_column(Boolean, default=True, index=True)
# role: {role}; region: {folder_id} or {region_name_from_group}
# role: {level} or legacy {role}; region: {folder_id} or {region_name_from_group}
payload: Mapped[dict] = mapped_column(_JsonType, default=dict)
remark: Mapped[str] = mapped_column(String(512), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)

View file

@ -0,0 +1,113 @@
"""Background scheduler for periodic NE batch collect."""
from __future__ import annotations
import logging
import threading
import time
from datetime import datetime
from .collection_policy import ensure_policy, next_due_at
from .collection_service import create_and_start_from_policy, has_active_collection_job
from .config import settings
from .db import SessionLocal
from .ne_collect_runner import dispatch_collection_runs
_log = logging.getLogger("netx.ne_collect.scheduler")
_stop = threading.Event()
_thread: threading.Thread | None = None
_BOOT_MONO = time.monotonic()
_last_tick_mono: float = 0.0
def _utcnow() -> datetime:
# Match NeCollectionJob timestamps (collection_service uses datetime.now()).
return datetime.now()
def startup_grace_remaining_sec() -> float:
grace = max(0, int(getattr(settings, "ne_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, "ne_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_active_collection_job(db) is not None:
return None
due = next_due_at(db, policy)
if due is not None and due > _utcnow():
return None
out, payload = create_and_start_from_policy(db, trigger_mode="schedule")
job_id = str(out.id)
dispatch_collection_runs(payload["job_id"], payload["run_ids"], payload["commands"])
_log.info("ne_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 {
"collection_job_running",
"commands_empty",
"no_eligible_ne",
"commands_required_for_schedule",
"no_selected_targets",
}:
_log.info("ne_collect schedule skip: %s", detail)
return None
_log.exception("ne_collect schedule start failed")
return None
finally:
db.close()
def _loop() -> None:
global _last_tick_mono
tick = max(15, int(getattr(settings, "ne_collect_scheduler_tick_sec", 60) or 60))
grace = max(0, int(getattr(settings, "ne_collect_startup_grace_sec", 3600) or 0))
_log.info("ne_collect scheduler started tick=%ss startup_grace=%ss", tick, grace)
while not _stop.is_set():
try:
_last_tick_mono = time.monotonic()
try_start_scheduled_collect()
except Exception:
_log.exception("ne_collect scheduler tick failed")
_stop.wait(tick)
_log.info("ne_collect scheduler stopped")
def start_ne_collect_scheduler() -> None:
global _thread
if not bool(getattr(settings, "ne_collect_scheduler_enabled", True)):
_log.info("ne_collect scheduler disabled by settings")
return
if _thread is not None and _thread.is_alive():
return
_stop.clear()
_thread = threading.Thread(target=_loop, name="ne-collect-scheduler", daemon=True)
_thread.start()
def stop_ne_collect_scheduler() -> None:
_stop.set()
def ne_collect_scheduler_status() -> dict:
now = time.monotonic()
return {
"running": bool(_thread is not None and _thread.is_alive()),
"last_tick_age_sec": (now - _last_tick_mono) if _last_tick_mono else None,
"startup_grace_remaining_sec": round(startup_grace_remaining_sec(), 1),
}

View file

@ -7,6 +7,7 @@ from typing import Any
from uuid import uuid4
from fastapi import HTTPException
from sqlalchemy import or_
from sqlalchemy.orm import Session
from .cli_creds import require_cli_creds_ready
@ -51,8 +52,30 @@ from .port_traffic_schemas import (
_log = logging.getLogger("netx.port_traffic.service")
def list_devices(db: Session, *, page: int = 1, page_size: int = 20) -> dict[str, Any]:
q = db.query(PortTrafficDevice).order_by(PortTrafficDevice.ne_name.asc(), PortTrafficDevice.ne_ip.asc())
def list_devices(
db: Session,
*,
page: int = 1,
page_size: int = 20,
status: str = "",
keyword: str = "",
) -> dict[str, Any]:
q = db.query(PortTrafficDevice)
st = str(status or "").strip()
if st:
q = q.filter(PortTrafficDevice.status == st)
kw = str(keyword or "").strip()
if kw:
like = f"%{kw}%"
q = q.filter(
or_(
PortTrafficDevice.ne_name.ilike(like),
PortTrafficDevice.ne_ip.ilike(like),
PortTrafficDevice.ne_id.ilike(like),
PortTrafficDevice.vendor.ilike(like),
)
)
q = q.order_by(PortTrafficDevice.ne_name.asc(), PortTrafficDevice.ne_ip.asc())
total = q.count()
rows = q.offset((page - 1) * page_size).limit(page_size).all()
return {

View file

@ -163,9 +163,11 @@ def api_delete_board(
def api_list_devices(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
):
return list_devices(db, page=page, page_size=page_size)
return list_devices(db, page=page, page_size=page_size, status=status, keyword=keyword)
@router.post("/devices")

View file

@ -55,6 +55,12 @@ def local_device_scheduler_status(*, role: str = "unknown") -> dict[str, Any]:
out["lldp_collect"] = lldp_collect_scheduler_status()
except Exception: # noqa: BLE001
out["lldp_collect"] = {"running": False, "error": "unavailable"}
try:
from .ne_collect_scheduler import ne_collect_scheduler_status
out["ne_collect"] = ne_collect_scheduler_status()
except Exception: # noqa: BLE001
out["ne_collect"] = {"running": False, "error": "unavailable"}
try:
from .port_traffic_scheduler import port_traffic_scheduler_status

View file

@ -154,6 +154,7 @@ def apply_topology_schema_safety_net(conn: Connection) -> None:
)
_run_sql(conn, "ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS world_x DOUBLE PRECISION")
_run_sql(conn, "ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS world_y DOUBLE PRECISION")
_run_sql(conn, "ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS level DOUBLE PRECISION")
_run_sql(
conn,
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_node_world_x ON topo_fabric_node (world_x)",
@ -162,6 +163,10 @@ def apply_topology_schema_safety_net(conn: Connection) -> None:
conn,
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_node_world_y ON topo_fabric_node (world_y)",
)
_run_sql(
conn,
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_node_level ON topo_fabric_node (level)",
)
# UME link ifnames — required by topology sync/apply; missing when DB predates the columns
# or Alembic was stamped head without running domain patches.
_run_sql(
@ -329,6 +334,7 @@ def apply_domain_schema_patches(conn: Connection) -> None:
"ALTER TABLE ne_collection_job ADD COLUMN IF NOT EXISTS last_run_at TIMESTAMP",
"UPDATE ne_collection_job SET last_run_at = COALESCE(ended_at, started_at, created_at) "
"WHERE last_run_at IS NULL",
"ALTER TABLE ne_collection_job ADD COLUMN IF NOT EXISTS trigger_mode VARCHAR(32) DEFAULT 'manual'",
"ALTER TABLE topology_edge ADD COLUMN IF NOT EXISTS stroke_color VARCHAR(32) DEFAULT ''",
"ALTER TABLE topology_edge ADD COLUMN IF NOT EXISTS stroke_width INTEGER DEFAULT 0",
"ALTER TABLE topology_edge ADD COLUMN IF NOT EXISTS line_style VARCHAR(16) DEFAULT ''",

View file

@ -1,25 +1,23 @@
"""Topology classify preview/apply and fabric node tagging."""
from __future__ import annotations
import re
from typing import Any
from fastapi import HTTPException
from sqlalchemy.orm import Session
from .models import TopoClassifyRule, TopoFabricNode, TopoFolder
from .models import TopoFabricNode, TopoFolder
from .topology_classify_common import (
_MATCH_FIELDS,
_ROLE_VALUES,
_compile_pattern,
_enabled_rules,
_ensure_region_by_name,
_match_text,
_resolve_level_hit,
_resolve_region_hit,
_resolve_role_hit,
_utcnow,
apply_level_fields,
)
from .topology_membership import normalize_view_role
from .topology_level import level_to_role, normalize_level, role_to_level
from .topology_schemas import (
ClassifyApplyOut,
ClassifyPreviewOut,
@ -31,31 +29,33 @@ from .topology_schemas import (
FabricNodeTagPatch,
)
def preview_classify(db: Session, *, sample_limit: int = 20) -> ClassifyPreviewOut:
role_rules = _enabled_rules(db, "role")
level_rules = _enabled_rules(db, "level")
region_rules = _enabled_rules(db, "region")
nodes = db.query(TopoFabricNode).order_by(TopoFabricNode.name.asc()).all()
role_matched = role_unmatched = role_conflict = 0
level_matched = level_unmatched = level_conflict = 0
region_matched = region_unmatched = region_conflict = 0
role_samples: list[dict[str, Any]] = []
level_samples: list[dict[str, Any]] = []
region_samples: list[dict[str, Any]] = []
unmatched_samples: list[dict[str, Any]] = []
for n in nodes:
role, _rid, multi_r = _resolve_role_hit(n, role_rules)
if role is None:
role_unmatched += 1
level, _rid, multi_r = _resolve_level_hit(n, level_rules)
if level is None:
level_unmatched += 1
else:
role_matched += 1
level_matched += 1
if multi_r:
role_conflict += 1
if len(role_samples) < sample_limit:
role_samples.append(
level_conflict += 1
if len(level_samples) < sample_limit:
level_samples.append(
{
"fabric_node_id": n.id,
"name": n.name,
"ip": n.ip,
"role": role,
"level": level,
"role": level_to_role(level),
"multi_hit": multi_r,
}
)
@ -80,20 +80,24 @@ def preview_classify(db: Session, *, sample_limit: int = 20) -> ClassifyPreviewO
}
)
if role is None and region_id is None and len(unmatched_samples) < sample_limit:
if level is None and region_id is None and len(unmatched_samples) < sample_limit:
unmatched_samples.append(
{"fabric_node_id": n.id, "name": n.name, "ip": n.ip, "vendor": n.vendor}
)
return ClassifyPreviewOut(
total_nodes=len(nodes),
role_matched=role_matched,
role_unmatched=role_unmatched,
role_conflicts=role_conflict,
level_matched=level_matched,
level_unmatched=level_unmatched,
level_conflicts=level_conflict,
role_matched=level_matched,
role_unmatched=level_unmatched,
role_conflicts=level_conflict,
region_matched=region_matched,
region_unmatched=region_unmatched,
region_conflicts=region_conflict,
role_samples=role_samples,
level_samples=level_samples,
role_samples=level_samples,
region_samples=region_samples,
unmatched_samples=unmatched_samples,
)
@ -105,27 +109,22 @@ def apply_classify(
skip_manual: bool = True,
fill_empty_only: bool = False,
) -> ClassifyApplyOut:
role_rules = _enabled_rules(db, "role")
level_rules = _enabled_rules(db, "level")
region_rules = _enabled_rules(db, "region")
nodes = db.query(TopoFabricNode).all()
role_updated = region_updated = skipped_manual = 0
level_updated = region_updated = skipped_manual = 0
for n in nodes:
role, _, _ = _resolve_role_hit(n, role_rules)
if role is not None:
level, _, _ = _resolve_level_hit(n, level_rules)
if level is not None:
if skip_manual and str(n.role_source or "") == "manual":
skipped_manual += 1
elif fill_empty_only and str(n.role or "").strip():
elif fill_empty_only and n.level is not None:
pass
else:
n.role = role
n.role_source = "rule"
apply_level_fields(n, level, source="rule")
n.updated_at = _utcnow()
role_updated += 1
elif not str(n.role or "").strip() and str(n.role_source or "") != "manual":
n.role = "unknown"
n.role_source = "rule"
n.updated_at = _utcnow()
level_updated += 1
region_id, _, _ = _resolve_region_hit(db, n, region_rules, create_missing=True)
if region_id is not None and not str(region_id).startswith("new:"):
@ -141,7 +140,8 @@ def apply_classify(
db.commit()
return ClassifyApplyOut(
role_updated=role_updated,
level_updated=level_updated,
role_updated=level_updated,
region_updated=region_updated,
skipped_manual=skipped_manual,
total_nodes=len(nodes),
@ -159,17 +159,19 @@ def list_unmatched(
q = db.query(TopoFabricNode)
k = str(kind or "any").strip().lower()
role_miss = or_(TopoFabricNode.role == "", TopoFabricNode.role == "unknown")
if k == "role":
k = "level"
level_miss = TopoFabricNode.level.is_(None)
region_miss = or_(
TopoFabricNode.region_folder_id.is_(None),
TopoFabricNode.region_folder_id == "",
)
if k == "role":
q = q.filter(role_miss)
if k == "level":
q = q.filter(level_miss)
elif k == "region":
q = q.filter(region_miss)
else:
q = q.filter(or_(role_miss, region_miss))
q = q.filter(or_(level_miss, region_miss))
total = q.count()
rows = (
q.order_by(TopoFabricNode.name.asc())
@ -195,14 +197,18 @@ def patch_fabric_node_tags(
n = db.get(TopoFabricNode, fabric_node_id)
if n is None:
raise HTTPException(status_code=404, detail="fabric_node_not_found")
if body.role is not None:
role = str(body.role or "").strip().lower()
if role and role not in _ROLE_VALUES:
raise HTTPException(status_code=400, detail="role_invalid")
n.role = role
n.role_source = "manual"
if body.region_folder_id is not None:
fid = str(body.region_folder_id or "").strip()
data = body.model_dump(exclude_unset=True)
if "level" in data or "role" in data:
try:
if "level" in data:
lv = normalize_level(data.get("level"))
else:
lv = role_to_level(data.get("role"))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
apply_level_fields(n, lv, source="manual")
if "region_folder_id" in data:
fid = str(data.get("region_folder_id") or "").strip()
if fid:
folder = db.get(TopoFolder, fid)
if folder is None or str(folder.kind or "") != "region":
@ -245,6 +251,7 @@ def match_fabric_nodes(db: Session, body: FabricNodesMatchRequest) -> FabricNode
"fabric_node_id": n.id,
"name": n.name,
"ip": n.ip,
"level": n.level,
"role": n.role or "",
"region_folder_id": n.region_folder_id or "",
"link_status": fabric_link_status(n),
@ -263,19 +270,23 @@ def match_fabric_nodes(db: Session, body: FabricNodesMatchRequest) -> FabricNode
def bulk_tag_fabric_nodes(
db: Session, body: FabricNodesBulkTagRequest
) -> FabricNodesBulkTagOut:
"""Assign role/region after user confirms a regex or explicit selection."""
if body.role is None and body.region_folder_id is None:
raise HTTPException(status_code=400, detail="role_or_region_required")
"""Assign level/region after user confirms a regex or explicit selection."""
data = body.model_dump(exclude_unset=True)
level_v: float | None | object = Ellipsis
if "level" in data:
try:
level_v = normalize_level(data.get("level"))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
elif "role" in data:
try:
level_v = role_to_level(data.get("role"))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
role_v: str | None = None
if body.role is not None:
role_v = str(body.role or "").strip().lower()
if role_v and role_v not in _ROLE_VALUES:
raise HTTPException(status_code=400, detail="role_invalid")
region_v: str | None = None
if body.region_folder_id is not None:
region_v = str(body.region_folder_id or "").strip()
region_v: str | None | object = Ellipsis
if "region_folder_id" in data:
region_v = str(data.get("region_folder_id") or "").strip()
if region_v:
folder = db.get(TopoFolder, region_v)
if folder is None or str(folder.kind or "") != "region":
@ -283,6 +294,9 @@ def bulk_tag_fabric_nodes(
else:
region_v = ""
if level_v is Ellipsis and region_v is Ellipsis:
raise HTTPException(status_code=400, detail="level_or_region_required")
ids = [str(x).strip() for x in (body.fabric_node_ids or []) if str(x).strip()]
if str(body.pattern or "").strip():
matched_nodes = _iter_regex_matches(
@ -298,11 +312,18 @@ def bulk_tag_fabric_nodes(
else:
raise HTTPException(status_code=400, detail="ids_or_pattern_required")
role_alias = level_to_role(level_v) if isinstance(level_v, float) else (
"" if level_v is None else None
)
if level_v is Ellipsis:
role_alias = None
samples = [
{
"fabric_node_id": n.id,
"name": n.name,
"ip": n.ip,
"level": n.level,
"role": n.role or "",
"region_folder_id": n.region_folder_id or "",
}
@ -313,19 +334,19 @@ def bulk_tag_fabric_nodes(
dry_run=True,
matched=len(matched_nodes),
updated=0,
role=role_v,
region_folder_id=region_v,
level=None if level_v is Ellipsis else level_v, # type: ignore[arg-type]
role=role_alias,
region_folder_id=None if region_v is Ellipsis else (region_v or None), # type: ignore[arg-type]
samples=samples,
)
now = _utcnow()
updated = 0
for n in matched_nodes:
if role_v is not None:
n.role = role_v
n.role_source = "manual"
if region_v is not None:
n.region_folder_id = region_v or None
if level_v is not Ellipsis:
apply_level_fields(n, level_v, source="manual") # type: ignore[arg-type]
if region_v is not Ellipsis:
n.region_folder_id = region_v or None # type: ignore[operator]
n.region_source = "manual"
n.updated_at = now
updated += 1
@ -334,8 +355,9 @@ def bulk_tag_fabric_nodes(
dry_run=False,
matched=len(matched_nodes),
updated=updated,
role=role_v,
region_folder_id=region_v,
level=None if level_v is Ellipsis else level_v, # type: ignore[arg-type]
role=role_alias,
region_folder_id=None if region_v is Ellipsis else (region_v or None), # type: ignore[arg-type]
samples=samples,
)
@ -343,5 +365,3 @@ def bulk_tag_fabric_nodes(
def apply_classify_empty_only(db: Session) -> ClassifyApplyOut:
"""Incremental classify for newly synced nodes (fill empty tags only)."""
return apply_classify(db, skip_manual=True, fill_empty_only=True)

View file

@ -10,23 +10,30 @@ from sqlalchemy.orm import Session
from .models import TopoClassifyRule, TopoFabricNode, TopoFolder
from .timeutil import utcnow_naive
from .topology_membership import VIEW_ROLES, normalize_view_role
from .topology_level import LEVEL_PRESETS, level_to_role, normalize_level, role_to_level
from .topology_schemas import (
ClassifyRuleOut,
TopologyFolderCreate,
)
_MAX_PATTERN_LEN = 512
_ROLE_VALUES = VIEW_ROLES | {"unknown"}
_MATCH_FIELDS = frozenset({"name", "ip", "name_ip"})
_SCOPES = frozenset({"role", "region"})
_SCOPES = frozenset({"level", "region", "role"}) # role = legacy alias of level
_SLICE_TEMPLATES = frozenset({"core_only", "core_agg", "agg_access"})
_ROLE_VALUES = frozenset(LEVEL_PRESETS) | {"unknown", "edge", ""}
def _utcnow() -> datetime:
return utcnow_naive()
def _normalize_scope(scope: str) -> str:
s = str(scope or "level").strip().lower()
if s == "role":
return "level"
return s
def _compile_pattern(pattern: str) -> re.Pattern[str]:
p = str(pattern or "").strip()
if not p:
@ -51,9 +58,10 @@ def _match_text(node: TopoFabricNode, match_field: str) -> str:
def _rule_out(row: TopoClassifyRule) -> ClassifyRuleOut:
scope = _normalize_scope(str(row.scope or "level"))
return ClassifyRuleOut(
id=row.id,
scope=str(row.scope or "role"),
scope=scope,
name=str(row.name or ""),
pattern=str(row.pattern or ""),
match_field=str(row.match_field or "name"),
@ -68,11 +76,25 @@ def _rule_out(row: TopoClassifyRule) -> ClassifyRuleOut:
def _validate_payload(scope: str, payload: dict[str, Any]) -> dict[str, Any]:
out = dict(payload or {})
if scope == "role":
role = normalize_view_role(str(out.get("role") or ""))
if str(out.get("role") or "").strip().lower() not in VIEW_ROLES:
raise HTTPException(status_code=400, detail="role_payload_invalid")
return {"role": role}
scope_n = _normalize_scope(scope)
if scope_n == "level":
if "level" in out and out.get("level") is not None:
try:
lv = normalize_level(out.get("level"))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
if lv is None:
raise HTTPException(status_code=400, detail="level_payload_invalid")
return {"level": lv}
if "role" in out:
try:
lv = role_to_level(str(out.get("role") or ""))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
if lv is None:
raise HTTPException(status_code=400, detail="level_payload_invalid")
return {"level": lv, "role": level_to_role(lv)}
raise HTTPException(status_code=400, detail="level_payload_invalid")
if "folder_id" in out and str(out.get("folder_id") or "").strip():
return {"folder_id": str(out["folder_id"]).strip()}
if "region_name_from_group" in out:
@ -87,9 +109,11 @@ def _validate_payload(scope: str, payload: dict[str, Any]) -> dict[str, Any]:
def _enabled_rules(db: Session, scope: str) -> list[tuple[TopoClassifyRule, re.Pattern[str]]]:
scope_n = _normalize_scope(scope)
scopes = ("level", "role") if scope_n == "level" else (scope_n,)
rows = (
db.query(TopoClassifyRule)
.filter(TopoClassifyRule.scope == scope, TopoClassifyRule.enabled.is_(True))
.filter(TopoClassifyRule.scope.in_(scopes), TopoClassifyRule.enabled.is_(True))
.order_by(TopoClassifyRule.priority.asc(), TopoClassifyRule.name.asc())
.all()
)
@ -122,23 +146,50 @@ def _ensure_region_by_name(db: Session, name: str) -> TopoFolder:
return folder
def _resolve_role_hit(
def _payload_level(payload: dict[str, Any] | None) -> float | None:
p = dict(payload or {})
if "level" in p and p.get("level") is not None:
try:
return normalize_level(p.get("level"))
except ValueError:
return None
if "role" in p:
try:
return role_to_level(str(p.get("role") or ""))
except ValueError:
return None
return None
def _resolve_level_hit(
node: TopoFabricNode, rules: list[tuple[TopoClassifyRule, re.Pattern[str]]]
) -> tuple[str | None, str | None, bool]:
"""Return (role, rule_id, multi_hit)."""
hits: list[tuple[str, str]] = []
) -> tuple[float | None, str | None, bool]:
"""Return (level, rule_id, multi_hit)."""
hits: list[tuple[float, str]] = []
for rule, cre in rules:
text = _match_text(node, rule.match_field)
if not text:
continue
if cre.search(text):
role = normalize_view_role(str((rule.payload or {}).get("role") or ""))
hits.append((role, rule.id))
lv = _payload_level(rule.payload)
if lv is None:
continue
hits.append((lv, rule.id))
if not hits:
return None, None, False
return hits[0][0], hits[0][1], len(hits) > 1
# Back-compat name used by older imports
def _resolve_role_hit(
node: TopoFabricNode, rules: list[tuple[TopoClassifyRule, re.Pattern[str]]]
) -> tuple[str | None, str | None, bool]:
lv, rid, multi = _resolve_level_hit(node, rules)
if lv is None:
return None, None, False
return level_to_role(lv), rid, multi
def _resolve_region_hit(
db: Session,
node: TopoFabricNode,
@ -182,3 +233,8 @@ def _resolve_region_hit(
return hits[0][0], hits[0][1], len(hits) > 1
def apply_level_fields(node: TopoFabricNode, level: float | None, *, source: str) -> None:
"""Write level + synced role alias."""
node.level = level
node.role = level_to_role(level)
node.role_source = source

View file

@ -13,6 +13,7 @@ from .topology_classify_common import (
_MAX_PATTERN_LEN,
_SCOPES,
_compile_pattern,
_normalize_scope,
_rule_out,
_utcnow,
_validate_payload,
@ -21,8 +22,13 @@ from .topology_schemas import ClassifyRuleCreate, ClassifyRuleOut, ClassifyRuleU
def list_rules(db: Session, *, scope: str = "") -> list[ClassifyRuleOut]:
q = db.query(TopoClassifyRule)
if scope.strip():
q = q.filter(TopoClassifyRule.scope == scope.strip().lower())
scope_f = str(scope or "").strip().lower()
if scope_f:
scope_n = _normalize_scope(scope_f)
if scope_n == "level":
q = q.filter(TopoClassifyRule.scope.in_(("level", "role")))
else:
q = q.filter(TopoClassifyRule.scope == scope_n)
rows = q.order_by(
TopoClassifyRule.scope.asc(),
TopoClassifyRule.priority.asc(),
@ -32,9 +38,10 @@ def list_rules(db: Session, *, scope: str = "") -> list[ClassifyRuleOut]:
def create_rule(db: Session, body: ClassifyRuleCreate) -> ClassifyRuleOut:
scope = str(body.scope or "").strip().lower()
if scope not in _SCOPES:
scope_raw = str(body.scope or "").strip().lower()
if scope_raw not in _SCOPES:
raise HTTPException(status_code=400, detail="scope_invalid")
scope = _normalize_scope(scope_raw)
match_field = str(body.match_field or "name").strip().lower()
if match_field not in _MATCH_FIELDS:
raise HTTPException(status_code=400, detail="match_field_invalid")

View file

@ -4,6 +4,7 @@ from __future__ import annotations
from typing import Any
from fastapi import HTTPException
from sqlalchemy import or_
from sqlalchemy.orm import Session
from .cli_resolve import get_default_profile, infer_device_type_vendor
@ -48,6 +49,8 @@ def _job_out(
include_items: bool = True,
page: int | None = None,
page_size: int | None = None,
item_status: str = "",
item_keyword: str = "",
) -> FabricDiscoverJobOut:
items_out: list[FabricDiscoverJobItemOut] = []
items_total = 0
@ -59,6 +62,37 @@ def _job_out(
.filter(TopoDiscoverJobItem.job_id == job.id)
.order_by(TopoDiscoverJobItem.created_at.asc())
)
st = str(item_status or "").strip().lower()
if st in ("ok", "success", "pass"):
q = q.filter(
TopoDiscoverJobItem.ok.is_(True),
TopoDiscoverJobItem.parser_stub.is_(False),
TopoDiscoverJobItem.unmatched_count == 0,
)
elif st in ("fail", "failed", "error"):
q = q.filter(
or_(
TopoDiscoverJobItem.ok.is_(False),
TopoDiscoverJobItem.parser_stub.is_(True),
)
)
elif st in ("warn", "warning"):
q = q.filter(
TopoDiscoverJobItem.ok.is_(True),
TopoDiscoverJobItem.parser_stub.is_(False),
TopoDiscoverJobItem.unmatched_count > 0,
)
kw = str(item_keyword or "").strip()
if kw:
like = f"%{kw}%"
q = q.filter(
or_(
TopoDiscoverJobItem.ne_name.ilike(like),
TopoDiscoverJobItem.ne_ip.ilike(like),
TopoDiscoverJobItem.ne_id.ilike(like),
TopoDiscoverJobItem.error.ilike(like),
)
)
items_total = int(q.count())
if page is not None and page_size is not None:
items_page = max(1, int(page or 1))
@ -121,11 +155,20 @@ def get_discover_job(
*,
page: int | None = None,
page_size: int | None = None,
item_status: str = "",
item_keyword: str = "",
) -> FabricDiscoverJobOut:
job = db.get(TopoDiscoverJob, str(job_id or "").strip())
if job is None:
raise HTTPException(status_code=404, detail="discover_job_not_found")
return _job_out(db, job, page=page, page_size=page_size)
return _job_out(
db,
job,
page=page,
page_size=page_size,
item_status=item_status,
item_keyword=item_keyword,
)
def _ume_target_dict(db: Session, uid: str, default_profile: Any) -> dict[str, str] | None:

View file

@ -55,7 +55,13 @@ from .topology_schemas import (
def _node_out(n: TopoFabricNode) -> FabricNodeOut:
lv = getattr(n, "level", None)
try:
level_v = float(lv) if lv is not None else None
except (TypeError, ValueError):
level_v = None
return FabricNodeOut(
level=level_v,
role=str(getattr(n, "role", "") or ""),
region_folder_id=str(getattr(n, "region_folder_id", None) or "") or None,
role_source=str(getattr(n, "role_source", "") or ""),
@ -264,6 +270,8 @@ def list_fabric_nodes(
*,
keyword: str = "",
role: str = "",
level: str = "",
level_major: str = "",
region_folder_id: str = "",
unmatched: str = "",
link_status: str = "",
@ -271,6 +279,7 @@ def list_fabric_nodes(
page_size: int = PAGE_DEFAULT,
) -> dict[str, Any]:
from .topology_inventory_lifecycle import enrich_fabric_node_dicts
from .topology_level import LEVEL_PRESETS, normalize_level
page = max(1, int(page or 1))
page_size = max(1, min(PAGE_MAX, int(page_size or PAGE_DEFAULT)))
@ -286,15 +295,46 @@ def list_fabric_nodes(
TopoFabricNode.ume_ne_id.ilike(like),
)
)
level_raw = str(level or "").strip()
if level_raw:
try:
lv = normalize_level(level_raw)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
if lv is None:
q = q.filter(TopoFabricNode.level.is_(None))
else:
q = q.filter(TopoFabricNode.level == lv)
maj_raw = str(level_major or "").strip()
if maj_raw:
try:
maj = int(float(maj_raw))
except ValueError as exc:
raise HTTPException(status_code=400, detail="level_major_invalid") from exc
q = q.filter(
TopoFabricNode.level.isnot(None),
TopoFabricNode.level >= float(maj),
TopoFabricNode.level < float(maj) + 1.0,
)
role_v = str(role or "").strip().lower()
if role_v:
q = q.filter(TopoFabricNode.role == role_v)
if role_v in LEVEL_PRESETS and not level_raw and not maj_raw:
# Prefer major-band filter so 1.1 still matches role=core
preset_lv = LEVEL_PRESETS[role_v]
maj = int(preset_lv)
q = q.filter(
TopoFabricNode.level.isnot(None),
TopoFabricNode.level >= float(maj),
TopoFabricNode.level < float(maj) + 1.0,
)
else:
q = q.filter(TopoFabricNode.role == role_v)
region_v = str(region_folder_id or "").strip()
if region_v:
q = q.filter(TopoFabricNode.region_folder_id == region_v)
um = str(unmatched or "").strip().lower()
if um == "role":
q = q.filter(or_(TopoFabricNode.role == "", TopoFabricNode.role == "unknown"))
if um in {"role", "level"}:
q = q.filter(TopoFabricNode.level.is_(None))
elif um == "region":
q = q.filter(
or_(TopoFabricNode.region_folder_id.is_(None), TopoFabricNode.region_folder_id == "")
@ -302,8 +342,7 @@ def list_fabric_nodes(
elif um == "any":
q = q.filter(
or_(
TopoFabricNode.role == "",
TopoFabricNode.role == "unknown",
TopoFabricNode.level.is_(None),
TopoFabricNode.region_folder_id.is_(None),
TopoFabricNode.region_folder_id == "",
)

138
netx_api/topology_level.py Normal file
View file

@ -0,0 +1,138 @@
"""Fabric topology level (layout rank): major.minor, smaller = closer to external/WAN."""
from __future__ import annotations
import math
from typing import Any
# Preset alias → default level (major.0). Sub-tiers use 1.1, 2.1, …
LEVEL_PRESETS: dict[str, float] = {
"external": 0.0,
"core": 1.0,
"aggregation": 2.0,
"aggregate": 2.0,
"agg": 2.0,
"access": 3.0,
"edge": 4.0,
"cpe": 4.0,
}
# floor(level) → synced role alias (filters / UI chips)
_MAJOR_TO_ROLE: dict[int, str] = {
0: "external",
1: "core",
2: "aggregation",
3: "access",
}
_ROLE_VALUES = frozenset(LEVEL_PRESETS) | {"unknown", ""}
def normalize_level(value: Any) -> float | None:
"""Parse level; empty/None → None. Snap to one decimal in [0, 99.9]."""
if value is None:
return None
if isinstance(value, str):
s = value.strip().lower()
if not s or s in {"unknown", "null", "none"}:
return None
if s in LEVEL_PRESETS:
return LEVEL_PRESETS[s]
try:
value = float(s)
except ValueError as exc:
raise ValueError("level_invalid") from exc
try:
lv = float(value)
except (TypeError, ValueError) as exc:
raise ValueError("level_invalid") from exc
if not math.isfinite(lv):
raise ValueError("level_invalid")
if lv < 0 or lv > 99.9:
raise ValueError("level_out_of_range")
return round(lv + 1e-9, 1)
def level_major(level: float | None) -> int | None:
if level is None:
return None
return int(math.floor(float(level)))
def level_to_role(level: float | None) -> str:
"""Denormalized role alias for filters; empty when unclassified."""
maj = level_major(level)
if maj is None:
return ""
if maj in _MAJOR_TO_ROLE:
return _MAJOR_TO_ROLE[maj]
if maj >= 4:
return "edge"
return ""
def role_to_level(role: str | None) -> float | None:
r = str(role or "").strip().lower()
if not r or r == "unknown":
return None
if r in LEVEL_PRESETS:
return LEVEL_PRESETS[r]
raise ValueError("role_invalid")
def coerce_level_input(
*,
level: Any = None,
role: Any = None,
level_provided: bool = False,
role_provided: bool = False,
) -> float | None | object:
"""Resolve patch/bulk input.
Returns:
- float | None: concrete level (None clears)
- Ellipsis: neither field provided
"""
if level_provided:
return normalize_level(level)
if role_provided:
return role_to_level(None if role is None else str(role))
return Ellipsis
def format_level(level: float | None) -> str:
if level is None:
return ""
lv = float(level)
if abs(lv - round(lv)) < 1e-9:
return str(int(round(lv)))
return f"{lv:.1f}".rstrip("0").rstrip(".") if "." in f"{lv:.1f}" else f"{lv:.1f}"
def infer_layer_from_level(
level: float | None,
*,
name: str = "",
role: str | None = None,
) -> str:
"""Map to layout layer key: external|core|agg|access|other."""
maj = level_major(level)
if maj is not None:
if maj <= 0:
return "external"
if maj == 1:
return "core"
if maj == 2:
return "agg"
if maj >= 3:
return "access"
# Fallbacks when unclassified
r = str(role or "").strip().lower()
if r in LEVEL_PRESETS:
return infer_layer_from_level(LEVEL_PRESETS[r])
import re
m = re.search(r"-(CN|AN|EN)(\d*)-", name or "", re.I)
if not m:
return "other"
return {"CN": "core", "AN": "agg", "EN": "access"}[m.group(1).upper()]

View file

@ -40,9 +40,11 @@ def ensure_topology_schema(conn: Connection) -> None:
"ALTER TABLE topo_view ADD COLUMN IF NOT EXISTS role VARCHAR(32) DEFAULT 'core'",
"ALTER TABLE topo_view ADD COLUMN IF NOT EXISTS sort_order INTEGER DEFAULT 0",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS role VARCHAR(32) DEFAULT ''",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS level DOUBLE PRECISION",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS region_folder_id VARCHAR(64)",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS role_source VARCHAR(16) DEFAULT ''",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS region_source VARCHAR(16) DEFAULT ''",
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_node_level ON topo_fabric_node (level)",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS world_x DOUBLE PRECISION",
"ALTER TABLE topo_fabric_node ADD COLUMN IF NOT EXISTS world_y DOUBLE PRECISION",
"CREATE INDEX IF NOT EXISTS ix_topo_fabric_node_world_x ON topo_fabric_node (world_x)",
@ -89,6 +91,7 @@ def ensure_topology_schema(conn: Connection) -> None:
"ALTER TABLE topo_view ADD COLUMN role VARCHAR(32) DEFAULT 'core'",
"ALTER TABLE topo_view ADD COLUMN sort_order INTEGER DEFAULT 0",
"ALTER TABLE topo_fabric_node ADD COLUMN role VARCHAR(32) DEFAULT ''",
"ALTER TABLE topo_fabric_node ADD COLUMN level REAL",
"ALTER TABLE topo_fabric_node ADD COLUMN region_folder_id VARCHAR(64)",
"ALTER TABLE topo_fabric_node ADD COLUMN role_source VARCHAR(16) DEFAULT ''",
"ALTER TABLE topo_fabric_node ADD COLUMN region_source VARCHAR(16) DEFAULT ''",

View file

@ -93,8 +93,10 @@ def api_fabric_summary(db: Session = Depends(get_db)) -> dict[str, Any]:
def api_fabric_nodes(
keyword: str = "",
role: str = "",
level: str = Query(default="", description="Exact level e.g. 1.1"),
level_major: str = Query(default="", description="Major band e.g. 2 → [2.0, 3.0)"),
region_folder_id: str = "",
unmatched: str = Query(default="", description="any | role | region"),
unmatched: str = Query(default="", description="any | level | role | region"),
link_status: str = Query(
default="",
description="linked | orphaned | managed | ume | both",
@ -107,6 +109,8 @@ def api_fabric_nodes(
db,
keyword=keyword,
role=role,
level=level,
level_major=level_major,
region_folder_id=region_folder_id,
unmatched=unmatched,
link_status=link_status,
@ -255,9 +259,18 @@ def api_fabric_discover_job(
job_id: str,
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
status: str = Query(default=""),
keyword: str = Query(default=""),
db: Session = Depends(get_db),
) -> dict[str, Any]:
return get_discover_job(db, job_id, page=page, page_size=page_size).model_dump()
return get_discover_job(
db,
job_id,
page=page,
page_size=page_size,
item_status=status,
item_keyword=keyword,
).model_dump()
# --- Tree / folders ---------------------------------------------------------
@ -520,7 +533,7 @@ def api_classify_apply(
@router.get("/classify/unmatched")
def api_classify_unmatched(
kind: str = Query(default="any", description="any | role | region"),
kind: str = Query(default="any", description="any | level | role | region"),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=50, ge=1, le=500),
db: Session = Depends(get_db),
@ -591,6 +604,10 @@ def api_generate_slices(
body: SliceGenerateRequest,
db: Session = Depends(get_db),
) -> dict[str, Any]:
"""Deprecated for product UI: slice maps freeze custom views and bypass region canvases.
Kept for tests / legacy clients. Prefer region folders + classify tags + MCP layout.
"""
return generate_slices(db, body).model_dump()

View file

@ -21,6 +21,9 @@ class FabricNodeOut(BaseModel):
ip: str = ""
vendor: str = ""
device_type: str = ""
# Layout rank (major.minor). None = unclassified.
level: float | None = None
# Synced alias from floor(level): external|core|aggregation|access|edge|""
role: str = ""
region_folder_id: str | None = None
role_source: str = ""
@ -483,7 +486,7 @@ class ClassifyRuleOut(BaseModel):
class ClassifyRuleCreate(BaseModel):
scope: str = Field(default="role", description="role | region")
scope: str = Field(default="level", description="level | region (role accepted as level alias)")
name: str = ""
pattern: str
match_field: str = "name"
@ -505,26 +508,33 @@ class ClassifyRuleUpdate(BaseModel):
class ClassifyPreviewOut(BaseModel):
total_nodes: int = 0
level_matched: int = 0
level_unmatched: int = 0
level_conflicts: int = 0
# Legacy aliases
role_matched: int = 0
role_unmatched: int = 0
role_conflicts: int = 0
region_matched: int = 0
region_unmatched: int = 0
region_conflicts: int = 0
level_samples: list[dict[str, Any]] = Field(default_factory=list)
role_samples: list[dict[str, Any]] = Field(default_factory=list)
region_samples: list[dict[str, Any]] = Field(default_factory=list)
unmatched_samples: list[dict[str, Any]] = Field(default_factory=list)
class ClassifyApplyOut(BaseModel):
role_updated: int = 0
level_updated: int = 0
role_updated: int = 0 # alias of level_updated
region_updated: int = 0
skipped_manual: int = 0
total_nodes: int = 0
class FabricNodeTagPatch(BaseModel):
role: str | None = None
level: float | None = None
role: str | None = None # preset alias → level
region_folder_id: str | None = None
@ -545,12 +555,13 @@ class FabricNodesMatchOut(BaseModel):
class FabricNodesBulkTagRequest(BaseModel):
"""Assign role/region to explicit ids or to an ephemeral regex match."""
"""Assign level/region to explicit ids or to an ephemeral regex match."""
fabric_node_ids: list[str] = Field(default_factory=list)
pattern: str = ""
match_field: str = "name"
role: str | None = None
level: float | None = None
role: str | None = None # preset → level when level omitted
region_folder_id: str | None = None
dry_run: bool = False
@ -559,6 +570,7 @@ class FabricNodesBulkTagOut(BaseModel):
dry_run: bool = False
matched: int = 0
updated: int = 0
level: float | None = None
role: str | None = None
region_folder_id: str | None = None
samples: list[dict[str, Any]] = Field(default_factory=list)

View file

@ -467,7 +467,17 @@ def _apply_fabric_filters(
)
role_v = str(role or "").strip().lower()
if role_v:
q = q.filter(TopoFabricNode.role == role_v)
from .topology_level import LEVEL_PRESETS
if role_v in LEVEL_PRESETS:
maj = int(LEVEL_PRESETS[role_v])
q = q.filter(
TopoFabricNode.level.isnot(None),
TopoFabricNode.level >= float(maj),
TopoFabricNode.level < float(maj) + 1.0,
)
else:
q = q.filter(TopoFabricNode.role == role_v)
vendor_v = str(vendor or "").strip()
if vendor_v:
q = q.filter(TopoFabricNode.vendor.ilike(f"%{vendor_v}%"))

View file

@ -1,6 +1,6 @@
"""UME / long-task runtime helpers shared by API and optional worker process.
Device collectors (config_sync / LLDP / port_traffic) run via ``start_device_schedulers``
Device collectors (config_sync / LLDP / ne_collect / port_traffic) run via ``start_device_schedulers``
(API inline by default; set ``NETX_RUN_INLINE_SCHEDULERS=false`` and run
``python -m netx_api.worker`` for a split process).
API process also owns UME keepalive, alarm WSS, current-alarm/inventory sync loops,
@ -53,11 +53,13 @@ def start_device_schedulers() -> None:
from .config_sync_scheduler import start_config_sync_scheduler
from .fabric_reconcile_scheduler import start_fabric_reconcile_scheduler
from .lldp_collect_scheduler import start_lldp_collect_scheduler
from .ne_collect_scheduler import start_ne_collect_scheduler
from .port_traffic_scheduler import start_port_traffic_scheduler
from .scheduler_heartbeat import start_scheduler_heartbeat_publisher
start_config_sync_scheduler()
start_lldp_collect_scheduler()
start_ne_collect_scheduler()
start_port_traffic_scheduler()
start_fabric_reconcile_scheduler()
# Publish status so API /metrics can see collectors when run in a split worker.