Add biz_state monitoring with shared LLDP parse and batch collect.

Introduce ParseProfile-based business state tasks that reuse LLDP parsing for snapshot batches (separate from topology Fabric writes), with scheduler, export, and a Phase1 UI under network tasks.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-17 16:17:04 +08:00
parent 066158d2f7
commit 4fa1b9b4dc
27 changed files with 2487 additions and 1 deletions

View file

@ -45,6 +45,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None:
except Exception: # noqa: BLE001 except Exception: # noqa: BLE001
_log.exception("stop_port_traffic_scheduler failed") _log.exception("stop_port_traffic_scheduler failed")
try:
from .biz_state_scheduler import stop_biz_state_scheduler
stop_biz_state_scheduler()
except Exception: # noqa: BLE001
_log.exception("stop_biz_state_scheduler failed")
try: try:
from .fabric_reconcile_scheduler import stop_fabric_reconcile_scheduler from .fabric_reconcile_scheduler import stop_fabric_reconcile_scheduler

View file

@ -62,6 +62,12 @@ def run_api_startup() -> None:
apply_topology_schema_safety_net(conn) apply_topology_schema_safety_net(conn)
apply_collection_schema_safety_net(conn) apply_collection_schema_safety_net(conn)
apply_hop_schema_safety_net(conn) apply_hop_schema_safety_net(conn)
try:
from .biz_state.schema_ensure import apply_biz_state_schema
apply_biz_state_schema(conn)
except Exception:
_log.exception("startup: biz_state schema safety patch failed")
except Exception: except Exception:
_log.exception("startup: auth/topology/collection/hop schema safety patches failed") _log.exception("startup: auth/topology/collection/hop schema safety patches failed")
if skip_ddl and alembic_ok: if skip_ddl and alembic_ok:
@ -133,6 +139,22 @@ def run_api_startup() -> None:
pt_cleared = recover_port_traffic_on_startup(db) pt_cleared = recover_port_traffic_on_startup(db)
if pt_cleared: if pt_cleared:
_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared) _log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared)
try:
from .models import BizStateTask
stuck = (
db.query(BizStateTask)
.filter(BizStateTask.collect_running.is_(True))
.all()
)
for t in stuck:
t.collect_running = False
t.last_error = (t.last_error or "")[:900] + " | reset_on_startup"
if stuck:
db.commit()
_log.info("startup: cleared %s biz_state stuck collect_running flag(s)", len(stuck))
except Exception:
_log.exception("startup: biz_state collect_running recovery failed")
try: try:
from .port_traffic_migrate import backfill_port_traffic_series from .port_traffic_migrate import backfill_port_traffic_series

View file

@ -0,0 +1,3 @@
"""Business state monitoring: ParseProfile collect → batch → (Phase2) compare."""
from __future__ import annotations

View file

@ -0,0 +1,486 @@
"""Collect runner: expand task items → CLI → match → parse → batch rows."""
from __future__ import annotations
import logging
from datetime import datetime
from typing import Any
from uuid import uuid4
from fastapi import HTTPException
from ..cli_creds import cli_creds_skip_reason
from ..cli_resolve import resolve_cli_target
from ..cli_timeout import run_cli_with_timeout
from ..config import settings
from ..db import SessionLocal
from ..lldp_shared import resolve_vendor_key
from ..models import (
BizStateBatch,
BizStateBatchCommand,
BizStateEvent,
BizStateLldpNeighbor,
BizStateTask,
BizStateTaskItem,
BizStateTaskItemBinding,
)
from ..ne_netmiko import disable_target_paging, send_show_command
from ..ne_session_factory import close_netmiko_connection, open_netmiko_connection
from .command_match import expand_from_bindings, match_command, normalize_command
from .parsers import get_parser
from .profiles import get_profile
_log = logging.getLogger("netx.biz_state.runner")
def _utcnow() -> datetime:
return datetime.utcnow()
def _format_error(exc: BaseException) -> str:
return f"{type(exc).__name__}: {exc}"[:1020]
def _append_event(db, *, task_id: str, message: str, level: str = "error") -> None:
msg = str(message or "").strip()
if not msg or not task_id:
return
db.add(
BizStateEvent(
id=uuid4().hex,
task_id=task_id,
level=str(level or "error")[:16],
message=msg[:4000],
created_at=_utcnow(),
)
)
def _bindings_for_item(db, item_id: str) -> list[dict[str, str]]:
rows = (
db.query(BizStateTaskItemBinding)
.filter(BizStateTaskItemBinding.item_id == item_id)
.all()
)
params: dict[str, str] = {}
for r in rows:
ph = str(r.placeholder or "").strip()
val = str(r.value or "").strip()
if ph and val:
params[ph] = val
return [params] if params else []
def _persist_lldp_rows(
db,
*,
batch: BizStateBatch,
cmd_row: BizStateBatchCommand,
records: list[dict[str, Any]],
) -> int:
n = 0
seen: set[tuple[str, str, str]] = set()
for rec in records:
local_if = str(rec.get("local_if") or "").strip()[:128]
remote_sys = str(rec.get("remote_sys") or "").strip()[:256]
remote_if = str(rec.get("remote_if") or "").strip()[:128]
if not local_if and not remote_sys and not remote_if:
continue
key = (local_if, remote_sys, remote_if)
if key in seen:
continue
seen.add(key)
db.add(
BizStateLldpNeighbor(
id=uuid4().hex,
batch_id=batch.id,
batch_command_id=cmd_row.id,
task_id=batch.task_id,
ne_id=batch.ne_id,
local_if=local_if,
remote_sys=remote_sys,
remote_if=remote_if,
remote_ip=str(rec.get("remote_ip") or "")[:128],
protocol=str(rec.get("protocol") or "lldp")[:32],
collected_at=_utcnow(),
)
)
n += 1
return n
def _finish_task(task_id: str, *, error: str = "") -> None:
db = SessionLocal()
try:
task = db.get(BizStateTask, task_id)
if not task:
return
task.collect_running = False
task.last_collect_ended_at = _utcnow()
task.last_error = str(error or "")[:1020]
task.updated_at = _utcnow()
if error:
_append_event(db, task_id=task_id, message=error, level="error")
db.commit()
finally:
db.close()
def dispatch_collect(task_id: str) -> None:
"""Claim and run one collect round."""
db = SessionLocal()
batch_id = ""
try:
task = db.get(BizStateTask, task_id)
if not task:
return
if task.collect_running:
return
if str(task.status or "") not in ("running", "draft", "paused"):
return
items = (
db.query(BizStateTaskItem)
.filter(
BizStateTaskItem.task_id == task_id,
BizStateTaskItem.enabled.is_(True),
)
.order_by(BizStateTaskItem.sort_order.asc())
.all()
)
if not items:
task.last_error = "no enabled task items"
task.updated_at = _utcnow()
db.commit()
return
task.collect_running = True
task.last_collect_started_at = _utcnow()
task.last_error = ""
task.updated_at = _utcnow()
batch = BizStateBatch(
id=uuid4().hex,
task_id=task.id,
source=task.source,
ne_id=task.ne_id,
ne_name=task.ne_name,
vendor=task.vendor,
status="running",
started_at=_utcnow(),
)
db.add(batch)
db.commit()
batch_id = batch.id
vendor = str(task.vendor or "")
device_type = str(task.device_type or "")
source = str(task.source or "managed").strip().lower()
ne_id = str(task.ne_id or "").strip()
retention = int(task.retention_batches or 30)
finally:
db.close()
if not batch_id:
return
error = ""
try:
_run_collect_session(
task_id=task_id,
batch_id=batch_id,
source=source,
ne_id=ne_id,
vendor=vendor,
device_type=device_type,
retention=retention,
)
except Exception as exc:
_log.exception("biz_state collect failed task=%s", task_id)
error = _format_error(exc)
db = SessionLocal()
try:
batch = db.get(BizStateBatch, batch_id)
if batch:
batch.status = "failed"
batch.message = error
batch.ended_at = _utcnow()
db.commit()
finally:
db.close()
finally:
_finish_task(task_id, error=error)
def _run_collect_session(
*,
task_id: str,
batch_id: str,
source: str,
ne_id: str,
vendor: str,
device_type: str,
retention: int,
) -> None:
per_cmd = int(settings.ne_collect_read_timeout_sec or 120)
cap = int(settings.ne_collect_run_timeout_cap_sec or 600)
db = SessionLocal()
try:
task = db.get(BizStateTask, task_id)
batch = db.get(BizStateBatch, batch_id)
if not task or not batch:
return
try:
if source == "managed":
creds, info = resolve_cli_target(db, managed_ne_id=ne_id)
elif source == "ume":
creds, info = resolve_cli_target(db, ume_ne_id=ne_id)
else:
raise RuntimeError("invalid_source")
except HTTPException as exc:
raise RuntimeError(str(exc.detail or "resolve_failed")) from exc
skip = cli_creds_skip_reason(creds, interactive=False)
if skip:
raise RuntimeError(skip)
vendor_eff = str(info.get("vendor") or vendor or "")
device_type_eff = str(info.get("device_type") or device_type or "")
if vendor_eff and vendor_eff != task.vendor:
task.vendor = vendor_eff
if device_type_eff and device_type_eff != task.device_type:
task.device_type = device_type_eff
db.commit()
vendor_key = resolve_vendor_key(vendor_eff, device_type_eff)
items = (
db.query(BizStateTaskItem)
.filter(
BizStateTaskItem.task_id == task_id,
BizStateTaskItem.enabled.is_(True),
)
.order_by(BizStateTaskItem.sort_order.asc())
.all()
)
# Build work list before opening session
work: list[tuple[str, dict[str, str], str, str, str]] = []
# concrete, params, profile_id, item_id, mode
for item in items:
if item.kind == "custom_raw":
cmd = normalize_command(item.command_override)
if cmd:
work.append((cmd, {}, "", item.id, "custom"))
continue
profile = get_profile(item.source_profile_id)
if profile is None:
_append_event(
db,
task_id=task_id,
message=f"unknown profile {item.source_profile_id}",
level="error",
)
continue
binds = _bindings_for_item(db, item.id)
try:
pairs = expand_from_bindings(
profile=profile,
bindings=binds,
command_override=item.command_override,
)
except ValueError as exc:
_append_event(db, task_id=task_id, message=str(exc), level="error")
continue
for concrete, params in pairs:
work.append((concrete, params, profile.profile_id, item.id, "normal"))
if not work:
batch.status = "failed"
batch.message = "no commands to run"
batch.ended_at = _utcnow()
db.commit()
raise RuntimeError("no commands to run")
budget = min(cap, per_cmd * max(1, len(work)) + 90)
holder: dict[str, Any] = {}
def _session() -> tuple[int, int, bool, bool]:
from ..ne_netmiko import drain_read_channel
conn = open_netmiko_connection(creds, session_timeout=budget)
holder["conn"] = conn
total_rows = 0
cmd_count = 0
any_fail = False
any_ok = False
try:
try:
disable_target_paging(
conn,
vendor=str(creds.get("vendor") or vendor_eff or ""),
device_type=str(creds.get("device_type") or device_type_eff or ""),
)
except Exception:
pass
try:
drain_read_channel(conn)
except Exception:
pass
sdb = SessionLocal()
try:
batch_row = sdb.get(BizStateBatch, batch_id)
if not batch_row:
return 0, 0, True, True
for concrete, params, profile_id, item_id, mode in work:
if holder.get("timed_out"):
raise TimeoutError("biz_state_aborted")
cmd_count += 1
cmd_row = BizStateBatchCommand(
id=uuid4().hex,
batch_id=batch_id,
task_item_id=item_id,
profile_id=profile_id,
raw_command=concrete[:512],
params_json=dict(params or {}),
created_at=_utcnow(),
)
try:
raw = send_show_command(conn, concrete, read_timeout=per_cmd)
cmd_row.raw_text = str(raw or "")
except Exception as exc:
any_fail = True
cmd_row.parse_status = "failed"
cmd_row.message = _format_error(exc)
sdb.add(cmd_row)
sdb.commit()
continue
if mode == "custom":
cmd_row.parse_status = "skipped_custom"
cmd_row.message = "custom_raw"
sdb.add(cmd_row)
sdb.commit()
any_ok = True
continue
hit = match_command(vendor_key=vendor_key, command=concrete)
if not hit:
any_fail = True
cmd_row.parse_status = "unmatched"
cmd_row.message = "no profile matched concrete command"
sdb.add(cmd_row)
sdb.commit()
continue
cmd_row.profile_id = hit.profile.profile_id
cmd_row.parser_id = hit.profile.parser_id
cmd_row.metric_id = hit.profile.metric_id
merged = {**params, **hit.params}
cmd_row.params_json = merged
parser = get_parser(hit.profile.parser_id)
if not parser:
any_fail = True
cmd_row.parse_status = "failed"
cmd_row.message = f"unknown parser {hit.profile.parser_id}"
sdb.add(cmd_row)
sdb.commit()
continue
try:
records = parser(
raw_text=cmd_row.raw_text,
vendor=vendor_eff,
device_type=device_type_eff,
command=hit.profile.textfsm_command or concrete,
params=merged,
)
except Exception as exc:
any_fail = True
cmd_row.parse_status = "failed"
cmd_row.message = f"parse: {_format_error(exc)}"
sdb.add(cmd_row)
sdb.commit()
continue
n = 0
if hit.profile.metric_id == "lldp_neighbor":
n = _persist_lldp_rows(
sdb, batch=batch_row, cmd_row=cmd_row, records=records
)
cmd_row.row_count = n
cmd_row.parse_status = "ok"
total_rows += n
any_ok = True
sdb.add(cmd_row)
sdb.commit()
finally:
sdb.close()
return total_rows, cmd_count, any_fail, any_ok
finally:
holder.pop("conn", None)
close_netmiko_connection(conn)
try:
total_rows, cmd_count, any_fail, any_ok = run_cli_with_timeout(
_session,
timeout_sec=budget,
conn_holder=holder,
label="biz_state",
acquire_budget=True,
)
except TimeoutError as exc:
raise RuntimeError(str(exc)[:1020]) from exc
batch = db.get(BizStateBatch, batch_id)
if batch:
batch.command_count = cmd_count
batch.row_count = total_rows
batch.ended_at = _utcnow()
if any_fail and any_ok:
batch.status = "partial"
elif any_fail and not any_ok:
batch.status = "failed"
batch.message = "all commands failed"
else:
batch.status = "success"
db.commit()
_purge_old_batches(db, task_id=task_id, keep=retention)
finally:
db.close()
def _purge_old_batches(db, *, task_id: str, keep: int) -> None:
keep_n = max(1, int(keep or 30))
rows = (
db.query(BizStateBatch)
.filter(BizStateBatch.task_id == task_id)
.order_by(BizStateBatch.started_at.desc())
.all()
)
drop = rows[keep_n:]
for b in drop:
bid = b.id
db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == bid).delete()
db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == bid).delete()
db.delete(b)
if drop:
db.commit()
def trigger_collect_now(task_id: str) -> dict[str, Any]:
db = SessionLocal()
try:
task = db.get(BizStateTask, task_id)
if not task:
raise HTTPException(status_code=404, detail="task not found")
if task.collect_running:
raise HTTPException(status_code=409, detail="collect already running")
finally:
db.close()
dispatch_collect(task_id)
return {"ok": True, "task_id": task_id}

View file

@ -0,0 +1,162 @@
"""Match concrete CLI commands to ParseProfiles (longest / most specific wins)."""
from __future__ import annotations
import re
from dataclasses import dataclass
from typing import Any
from .profiles import ParseProfile, all_profiles, get_profile
@dataclass(frozen=True)
class MatchResult:
profile: ParseProfile
params: dict[str, str]
matched_length: int
def normalize_command(command: str) -> str:
text = str(command or "").strip()
text = re.sub(r"\s+", " ", text)
return text
def match_command(
*,
vendor_key: str,
command: str,
profiles: list[ParseProfile] | None = None,
) -> MatchResult | None:
cmd = normalize_command(command)
if not cmd:
return None
key = str(vendor_key or "").strip().lower()
cands = profiles if profiles is not None else [
p for p in all_profiles() if p.enabled and (p.vendor_key == key or p.vendor_key == "*")
]
best: MatchResult | None = None
for p in cands:
try:
m = re.match(p.match, cmd)
except re.error:
continue
if not m:
continue
params = {k: str(v or "").strip() for k, v in (m.groupdict() or {}).items()}
scored = MatchResult(profile=p, params=params, matched_length=len(p.match))
if best is None or scored.matched_length > best.matched_length:
best = scored
return best
def expand_from_bindings(
*,
profile: ParseProfile,
bindings: list[dict[str, str]] | None = None,
command_override: str = "",
) -> list[tuple[str, dict[str, str]]]:
"""Return list of (concrete_command, params). Reject leftover placeholders."""
override = normalize_command(command_override)
if override:
if "<" in override and ">" in override:
raise ValueError(f"command still has placeholders: {override}")
return [(override, {})]
tmpl = profile.command_template
if not profile.placeholders:
concrete = normalize_command(tmpl)
if "<" in concrete and ">" in concrete:
raise ValueError(f"template has placeholders but profile defines none: {concrete}")
return [(concrete, {})]
binds = list(bindings or [])
if not binds:
raise ValueError(f"profile {profile.profile_id} requires parameter bindings")
out: list[tuple[str, dict[str, str]]] = []
for b in binds:
params = {str(k): str(v).strip() for k, v in dict(b or {}).items() if str(v).strip()}
rendered = tmpl
for ph in profile.placeholders:
val = params.get(ph.name) or params.get(ph.schema_field) or ""
if ph.required and not val:
raise ValueError(f"missing placeholder {ph.name} for {profile.profile_id}")
rendered = rendered.replace(f"<{ph.name}>", val)
concrete = normalize_command(rendered)
if re.search(r"<[^>]+>", concrete):
raise ValueError(f"unresolved placeholders in: {concrete}")
out.append((concrete, params))
return out
def preview_task_item(
*,
vendor_key: str,
profile_id: str = "",
command: str = "",
bindings: list[dict[str, str]] | None = None,
kind: str = "catalog",
) -> dict[str, Any]:
"""Dry-run expand + match for UI preview."""
if kind == "custom_raw":
cmd = normalize_command(command)
return {
"ok": bool(cmd),
"kind": "custom_raw",
"commands": [cmd] if cmd else [],
"parse": "skipped_custom",
"message": "custom row: collect only, no parse",
}
profile = get_profile(profile_id) if profile_id else None
if profile is None and command:
hit = match_command(vendor_key=vendor_key, command=command)
if hit:
profile = hit.profile
if profile is None:
return {
"ok": False,
"kind": kind,
"commands": [],
"parse": "unmatched",
"message": "no profile matched",
}
try:
pairs = expand_from_bindings(
profile=profile,
bindings=bindings,
command_override=command if command and command != profile.command_template else "",
)
except ValueError as exc:
return {
"ok": False,
"kind": kind,
"profile_id": profile.profile_id,
"commands": [],
"parse": "invalid",
"message": str(exc),
}
previews = []
for concrete, params in pairs:
hit = match_command(vendor_key=vendor_key, command=concrete)
previews.append(
{
"command": concrete,
"params": params,
"matched_profile_id": hit.profile.profile_id if hit else "",
"metric_id": hit.profile.metric_id if hit else "",
"parser_id": hit.profile.parser_id if hit else "",
"match_params": hit.params if hit else {},
}
)
return {
"ok": all(bool(x["matched_profile_id"]) for x in previews),
"kind": kind,
"profile_id": profile.profile_id,
"commands": previews,
"parse": "ok" if previews and all(x["matched_profile_id"] for x in previews) else "partial",
"message": "",
}

View file

@ -0,0 +1,47 @@
"""Parser callbacks keyed by parser_id."""
from __future__ import annotations
from typing import Any, Callable
from ...lldp_shared import NeighborHit, parse_neighbor_output
NormalizeFn = Callable[..., list[dict[str, Any]]]
def normalize_lldp_neighbors(
*,
raw_text: str,
vendor: str = "",
device_type: str = "",
command: str = "",
params: dict[str, str] | None = None,
) -> list[dict[str, Any]]:
hits: list[NeighborHit] = parse_neighbor_output(
raw_text,
vendor=vendor,
device_type=device_type,
command=command,
)
_ = params
rows: list[dict[str, Any]] = []
for h in hits:
rows.append(
{
"local_if": str(h.local_port or "").strip(),
"remote_sys": str(h.remote_name or "").strip(),
"remote_if": str(h.remote_port or "").strip(),
"remote_ip": str(h.remote_ip or "").strip(),
"protocol": str(h.protocol or "lldp").strip() or "lldp",
}
)
return rows
_REGISTRY: dict[str, NormalizeFn] = {
"lldp_neighbors": normalize_lldp_neighbors,
}
def get_parser(parser_id: str) -> NormalizeFn | None:
return _REGISTRY.get(str(parser_id or "").strip())

View file

@ -0,0 +1,213 @@
"""ParseProfile registry: command template + parser + schema (code as source of truth).
Display fields may be overridden at runtime via biz_state_command_override (hot edit).
"""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any
@dataclass(frozen=True)
class FieldDef:
name: str
dtype: str = "str" # str|int|float|bool
nullable: bool = True
indexed: bool = False
is_key: bool = False
is_interface: bool = False
role: str = "identity" # identity|state|counter|meta
display_name: str = ""
from_command_param: bool = False
length: int = 256
@dataclass(frozen=True)
class PlaceholderDef:
name: str
schema_field: str
required: bool = True
bind_mode: str = "manual_text" # discover_select | manual_text
discover_profile_id: str = ""
discover_value_field: str = ""
discover_label_field: str = ""
@dataclass
class ParseProfile:
profile_id: str
vendor_key: str # zte|huawei|cisco|... or "*" for all
metric_id: str
parser_id: str
title: str
command_template: str
match: str # regex with optional named groups
textfsm_command: str = ""
description: str = ""
sample_output: str = ""
placeholders: list[PlaceholderDef] = field(default_factory=list)
fields: list[FieldDef] = field(default_factory=list)
tags: list[str] = field(default_factory=list)
sort_order: int = 100
enabled: bool = True
kind: str = "collect" # collect | discover
_LLDP_FIELDS: list[FieldDef] = [
FieldDef("local_if", length=128, indexed=True, is_key=True, is_interface=True, display_name="本端接口"),
FieldDef("remote_sys", length=256, indexed=True, is_key=True, display_name="对端系统名"),
FieldDef("remote_if", length=128, is_key=True, display_name="对端接口"),
FieldDef("remote_ip", length=128, role="meta", display_name="对端管理IP"),
FieldDef("protocol", length=32, role="meta", display_name="协议"),
]
def _lldp_profiles() -> list[ParseProfile]:
"""One logical LLDP collect profile per vendor_key (commands differ)."""
from ..lldp_shared import VENDOR_LLDP_PROFILES, STUB_PARSER_KEYS
out: list[ParseProfile] = []
order = 10
for key, vp in VENDOR_LLDP_PROFILES.items():
if key in STUB_PARSER_KEYS and key != "nokia":
# Still register so UI can show; parse may return empty until templates exist.
pass
cmd = vp.lldp_command
# Escape for regex: match exact command ignoring extra whitespace flexibility
escaped = r"\s+".join(
__import__("re").escape(p) for p in cmd.split() if p
)
out.append(
ParseProfile(
profile_id=f"{key}.lldp_neighbors",
vendor_key=key,
metric_id="lldp_neighbor",
parser_id="lldp_neighbors",
title="LLDP Neighbors",
command_template=cmd,
match=rf"(?i)^\s*{escaped}\s*$",
textfsm_command=cmd,
description=vp.notes or "LLDP neighbor table snapshot for cutover compare.",
sample_output="",
placeholders=[],
fields=list(_LLDP_FIELDS),
tags=["lldp", "l2"],
sort_order=order,
enabled=key not in ("ericsson", "generic"),
kind="collect",
)
)
order += 10
# AOS uses a different command than generic nokia profile
out.append(
ParseProfile(
profile_id="nokia_aos.lldp_neighbors",
vendor_key="nokia",
metric_id="lldp_neighbor",
parser_id="lldp_neighbors",
title="LLDP Neighbors (AOS)",
command_template="show lldp remote-system",
match=r"(?i)^\s*show\s+lldp\s+remote-system\s*$",
textfsm_command="show lldp remote-system",
description="Alcatel AOS LLDP remote-system.",
fields=list(_LLDP_FIELDS),
tags=["lldp", "l2", "aos"],
sort_order=95,
enabled=True,
kind="collect",
)
)
return out
_PROFILES: list[ParseProfile] | None = None
def all_profiles() -> list[ParseProfile]:
global _PROFILES
if _PROFILES is None:
_PROFILES = _lldp_profiles()
return list(_PROFILES)
def reload_profiles() -> None:
global _PROFILES
_PROFILES = None
def profiles_for_vendor(vendor_key: str) -> list[ParseProfile]:
key = str(vendor_key or "").strip().lower()
return [
p
for p in all_profiles()
if p.enabled and (p.vendor_key == key or p.vendor_key == "*")
]
def get_profile(profile_id: str) -> ParseProfile | None:
pid = str(profile_id or "").strip()
for p in all_profiles():
if p.profile_id == pid:
return p
return None
def metric_field_map() -> dict[str, list[FieldDef]]:
"""Merge fields by metric_id (first-seen wins on name conflict)."""
out: dict[str, list[FieldDef]] = {}
seen: dict[str, set[str]] = {}
for p in all_profiles():
mid = p.metric_id
bucket = out.setdefault(mid, [])
names = seen.setdefault(mid, set())
for f in p.fields:
if f.name in names:
continue
names.add(f.name)
bucket.append(f)
return out
def profile_to_public_dict(p: ParseProfile, *, overrides: dict[str, Any] | None = None) -> dict[str, Any]:
ov = overrides or {}
return {
"profile_id": p.profile_id,
"vendor_key": p.vendor_key,
"metric_id": p.metric_id,
"parser_id": p.parser_id,
"title": str(ov.get("title") or p.title),
"command_template": str(ov.get("command_template") or p.command_template),
"description": str(ov.get("description") if ov.get("description") is not None else p.description),
"sample_output": str(ov.get("sample_output") if ov.get("sample_output") is not None else p.sample_output),
"placeholders": [
{
"name": ph.name,
"schema_field": ph.schema_field,
"required": ph.required,
"bind_mode": ph.bind_mode,
"discover_profile_id": ph.discover_profile_id,
"discover_value_field": ph.discover_value_field,
"discover_label_field": ph.discover_label_field,
}
for ph in p.placeholders
],
"fields": [
{
"name": f.name,
"dtype": f.dtype,
"is_key": f.is_key,
"is_interface": f.is_interface,
"role": f.role,
"display_name": f.display_name or f.name,
"from_command_param": f.from_command_param,
}
for f in p.fields
],
"tags": list(p.tags),
"sort_order": p.sort_order,
"enabled": bool(ov.get("enabled")) if "enabled" in ov else p.enabled,
"kind": p.kind,
"match": p.match,
"textfsm_command": p.textfsm_command or p.command_template,
}

View file

@ -0,0 +1,27 @@
"""Idempotent DDL safety-net for biz_state tables (create_all + ADD COLUMN)."""
from __future__ import annotations
import logging
from sqlalchemy.engine import Connection
from ..schema_patches import _run_sql
_log = logging.getLogger("netx.biz_state.schema")
def apply_biz_state_schema(conn: Connection) -> None:
"""Ensure core tables exist; create_all usually handles this — safety net for brownfield."""
# Tables are defined on ORM Base; create_all covers new installs.
# Keep lightweight indexes that older DBs might miss.
for sql in (
"CREATE INDEX IF NOT EXISTS ix_biz_state_task_status ON biz_state_task (status)",
"CREATE INDEX IF NOT EXISTS ix_biz_state_batch_task_id ON biz_state_batch (task_id)",
"CREATE INDEX IF NOT EXISTS ix_biz_state_batch_command_batch_id ON biz_state_batch_command (batch_id)",
"CREATE INDEX IF NOT EXISTS ix_biz_state_lldp_neighbor_batch_id ON biz_state_lldp_neighbor (batch_id)",
):
try:
_run_sql(conn, sql)
except Exception:
_log.debug("biz_state schema patch skipped: %s", sql[:80], exc_info=True)

View file

@ -0,0 +1,474 @@
"""biz_state service: tasks, profiles, batches, export."""
from __future__ import annotations
import io
import zipfile
from datetime import datetime
from typing import Any
from uuid import uuid4
from fastapi import HTTPException
from sqlalchemy.orm import Session
from ..lldp_shared import resolve_vendor_key
from ..models import (
BizStateBatch,
BizStateBatchCommand,
BizStateCommandOverride,
BizStateEvent,
BizStateLldpNeighbor,
BizStateTask,
BizStateTaskItem,
BizStateTaskItemBinding,
ManagedNE,
)
from ..timeutil import utcnow_naive
from .command_match import preview_task_item
from .profiles import all_profiles, get_profile, profile_to_public_dict, profiles_for_vendor
def _utcnow() -> datetime:
return utcnow_naive()
def _override_map(db: Session) -> dict[str, dict[str, Any]]:
rows = db.query(BizStateCommandOverride).all()
out: dict[str, dict[str, Any]] = {}
for r in rows:
ov: dict[str, Any] = {}
if r.title:
ov["title"] = r.title
if r.command_template:
ov["command_template"] = r.command_template
if r.description:
ov["description"] = r.description
if r.sample_output:
ov["sample_output"] = r.sample_output
if r.enabled is not None:
ov["enabled"] = bool(r.enabled)
out[str(r.profile_id)] = ov
return out
def list_profiles_public(db: Session, *, vendor_key: str = "") -> list[dict[str, Any]]:
ov = _override_map(db)
key = str(vendor_key or "").strip().lower()
profiles = profiles_for_vendor(key) if key else [p for p in all_profiles() if p.enabled]
# Also include disabled-in-code but enabled via override? keep simple: filter enabled after merge
result = []
for p in sorted(profiles, key=lambda x: (x.sort_order, x.profile_id)):
d = profile_to_public_dict(p, overrides=ov.get(p.profile_id))
if d.get("enabled"):
result.append(d)
return result
def upsert_profile_override(db: Session, profile_id: str, body: dict[str, Any]) -> dict[str, Any]:
p = get_profile(profile_id)
if not p:
raise HTTPException(status_code=404, detail="profile_not_found")
row = (
db.query(BizStateCommandOverride)
.filter(BizStateCommandOverride.profile_id == profile_id)
.one_or_none()
)
if row is None:
row = BizStateCommandOverride(id=uuid4().hex, profile_id=profile_id)
db.add(row)
if "title" in body:
row.title = str(body.get("title") or "")
if "command_template" in body:
row.command_template = str(body.get("command_template") or "")
if "description" in body:
row.description = str(body.get("description") or "")
if "sample_output" in body:
row.sample_output = str(body.get("sample_output") or "")
if "enabled" in body and body.get("enabled") is not None:
row.enabled = bool(body.get("enabled"))
row.updated_at = _utcnow()
db.commit()
return profile_to_public_dict(p, overrides=_override_map(db).get(profile_id))
def _ne_meta(db: Session, *, source: str, ne_id: str) -> dict[str, str]:
src = str(source or "managed").strip().lower()
nid = str(ne_id or "").strip()
if src == "managed":
row = db.get(ManagedNE, nid)
if not row:
raise HTTPException(status_code=404, detail="managed_ne_not_found")
return {
"ne_name": str(getattr(row, "name", "") or ""),
"ne_ip": str(getattr(row, "ip_address", "") or ""),
"vendor": str(getattr(row, "vendor", "") or ""),
"device_type": str(getattr(row, "device_type", "") or ""),
}
# ume: best-effort from inventory if present
try:
from ..models import UmeInventoryNE
row = db.get(UmeInventoryNE, nid)
if row:
return {
"ne_name": str(getattr(row, "name", "") or ""),
"ne_ip": str(getattr(row, "ip", "") or ""),
"vendor": str(getattr(row, "vendor", "") or ""),
"device_type": str(getattr(row, "device_type", "") or ""),
}
except Exception:
pass
return {"ne_name": "", "ne_ip": "", "vendor": "", "device_type": ""}
def create_task(db: Session, body: dict[str, Any]) -> dict[str, Any]:
source = str(body.get("source") or "managed").strip().lower() or "managed"
ne_id = str(body.get("ne_id") or "").strip()
if not ne_id:
raise HTTPException(status_code=400, detail="ne_id_required")
existing = (
db.query(BizStateTask)
.filter(BizStateTask.source == source, BizStateTask.ne_id == ne_id)
.one_or_none()
)
if existing:
raise HTTPException(status_code=409, detail="task_already_exists_for_ne")
meta = _ne_meta(db, source=source, ne_id=ne_id)
vendor = str(body.get("vendor") or meta["vendor"] or "")
device_type = str(body.get("device_type") or meta["device_type"] or "")
task = BizStateTask(
id=uuid4().hex,
source=source,
ne_id=ne_id,
ne_name=str(body.get("ne_name") or meta["ne_name"] or ""),
ne_ip=str(body.get("ne_ip") or meta["ne_ip"] or ""),
vendor=vendor,
device_type=device_type,
note=str(body.get("note") or "")[:256],
status="draft",
interval_sec=max(60, int(body.get("interval_sec") or 300)),
retention_batches=max(1, int(body.get("retention_batches") or 30)),
created_at=_utcnow(),
updated_at=_utcnow(),
)
db.add(task)
db.flush()
items_in = list(body.get("items") or [])
if not items_in:
# Default: enable LLDP profile for this vendor
vkey = resolve_vendor_key(vendor, device_type)
for p in profiles_for_vendor(vkey):
if p.metric_id == "lldp_neighbor" and p.kind == "collect":
items_in.append(
{
"source_profile_id": p.profile_id,
"kind": "catalog",
"enabled": True,
"title": p.title,
}
)
break
_replace_items(db, task.id, items_in)
db.commit()
return get_task(db, task.id)
def _replace_items(db: Session, task_id: str, items_in: list[dict[str, Any]]) -> None:
old_items = db.query(BizStateTaskItem).filter(BizStateTaskItem.task_id == task_id).all()
for it in old_items:
db.query(BizStateTaskItemBinding).filter(BizStateTaskItemBinding.item_id == it.id).delete()
db.delete(it)
db.flush()
for idx, raw in enumerate(items_in):
kind = str(raw.get("kind") or "catalog").strip() or "catalog"
item = BizStateTaskItem(
id=uuid4().hex,
task_id=task_id,
source_profile_id=str(raw.get("source_profile_id") or "")[:128],
kind=kind,
enabled=bool(raw.get("enabled", True)),
title=str(raw.get("title") or "")[:256],
command_override=str(raw.get("command_override") or raw.get("command") or "")[:512],
sort_order=int(raw.get("sort_order") if raw.get("sort_order") is not None else idx),
created_at=_utcnow(),
)
db.add(item)
db.flush()
for b in list(raw.get("bindings") or []):
ph = str(b.get("placeholder") or b.get("name") or "").strip()
val = str(b.get("value") or "").strip()
if not ph or not val:
continue
db.add(
BizStateTaskItemBinding(
id=uuid4().hex,
item_id=item.id,
placeholder=ph[:64],
value=val[:256],
created_at=_utcnow(),
)
)
def update_task(db: Session, task_id: str, body: dict[str, Any]) -> dict[str, Any]:
task = db.get(BizStateTask, task_id)
if not task:
raise HTTPException(status_code=404, detail="task_not_found")
if "note" in body:
task.note = str(body.get("note") or "")[:256]
if "interval_sec" in body:
task.interval_sec = max(60, int(body.get("interval_sec") or 300))
if "retention_batches" in body:
task.retention_batches = max(1, int(body.get("retention_batches") or 30))
if "status" in body:
st = str(body.get("status") or "").strip()
if st in ("draft", "running", "paused", "stopped"):
task.status = st
if "items" in body:
_replace_items(db, task.id, list(body.get("items") or []))
task.updated_at = _utcnow()
db.commit()
return get_task(db, task_id)
def get_task(db: Session, task_id: str) -> dict[str, Any]:
task = db.get(BizStateTask, task_id)
if not task:
raise HTTPException(status_code=404, detail="task_not_found")
items = (
db.query(BizStateTaskItem)
.filter(BizStateTaskItem.task_id == task_id)
.order_by(BizStateTaskItem.sort_order.asc())
.all()
)
item_out = []
for it in items:
binds = (
db.query(BizStateTaskItemBinding)
.filter(BizStateTaskItemBinding.item_id == it.id)
.all()
)
item_out.append(
{
"id": it.id,
"source_profile_id": it.source_profile_id,
"kind": it.kind,
"enabled": bool(it.enabled),
"title": it.title,
"command_override": it.command_override,
"sort_order": it.sort_order,
"bindings": [{"placeholder": b.placeholder, "value": b.value} for b in binds],
}
)
return {
"id": task.id,
"source": task.source,
"ne_id": task.ne_id,
"ne_name": task.ne_name,
"ne_ip": task.ne_ip,
"vendor": task.vendor,
"device_type": task.device_type,
"note": task.note,
"status": task.status,
"interval_sec": task.interval_sec,
"retention_batches": task.retention_batches,
"collect_running": bool(task.collect_running),
"last_collect_started_at": task.last_collect_started_at.isoformat() + "Z"
if task.last_collect_started_at
else None,
"last_collect_ended_at": task.last_collect_ended_at.isoformat() + "Z"
if task.last_collect_ended_at
else None,
"last_error": task.last_error,
"items": item_out,
}
def list_tasks(db: Session) -> list[dict[str, Any]]:
rows = db.query(BizStateTask).order_by(BizStateTask.updated_at.desc()).all()
return [
{
"id": t.id,
"source": t.source,
"ne_id": t.ne_id,
"ne_name": t.ne_name,
"ne_ip": t.ne_ip,
"vendor": t.vendor,
"status": t.status,
"interval_sec": t.interval_sec,
"collect_running": bool(t.collect_running),
"last_error": t.last_error,
"last_collect_ended_at": t.last_collect_ended_at.isoformat() + "Z"
if t.last_collect_ended_at
else None,
}
for t in rows
]
def delete_task(db: Session, task_id: str) -> None:
task = db.get(BizStateTask, task_id)
if not task:
raise HTTPException(status_code=404, detail="task_not_found")
batches = db.query(BizStateBatch).filter(BizStateBatch.task_id == task_id).all()
for b in batches:
db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == b.id).delete()
db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == b.id).delete()
db.delete(b)
items = db.query(BizStateTaskItem).filter(BizStateTaskItem.task_id == task_id).all()
for it in items:
db.query(BizStateTaskItemBinding).filter(BizStateTaskItemBinding.item_id == it.id).delete()
db.delete(it)
db.query(BizStateEvent).filter(BizStateEvent.task_id == task_id).delete()
db.delete(task)
db.commit()
def list_batches(db: Session, task_id: str, *, limit: int = 50) -> list[dict[str, Any]]:
rows = (
db.query(BizStateBatch)
.filter(BizStateBatch.task_id == task_id)
.order_by(BizStateBatch.started_at.desc())
.limit(max(1, min(200, int(limit))))
.all()
)
return [
{
"id": b.id,
"status": b.status,
"command_count": b.command_count,
"row_count": b.row_count,
"message": b.message,
"started_at": b.started_at.isoformat() + "Z" if b.started_at else None,
"ended_at": b.ended_at.isoformat() + "Z" if b.ended_at else None,
}
for b in rows
]
def get_batch(db: Session, batch_id: str) -> dict[str, Any]:
b = db.get(BizStateBatch, batch_id)
if not b:
raise HTTPException(status_code=404, detail="batch_not_found")
cmds = (
db.query(BizStateBatchCommand)
.filter(BizStateBatchCommand.batch_id == batch_id)
.order_by(BizStateBatchCommand.created_at.asc())
.all()
)
neighbors = (
db.query(BizStateLldpNeighbor)
.filter(BizStateLldpNeighbor.batch_id == batch_id)
.order_by(BizStateLldpNeighbor.local_if.asc())
.limit(5000)
.all()
)
return {
"id": b.id,
"task_id": b.task_id,
"status": b.status,
"command_count": b.command_count,
"row_count": b.row_count,
"message": b.message,
"started_at": b.started_at.isoformat() + "Z" if b.started_at else None,
"ended_at": b.ended_at.isoformat() + "Z" if b.ended_at else None,
"commands": [
{
"id": c.id,
"profile_id": c.profile_id,
"parser_id": c.parser_id,
"metric_id": c.metric_id,
"raw_command": c.raw_command,
"params": c.params_json or {},
"parse_status": c.parse_status,
"row_count": c.row_count,
"message": c.message,
"raw_text_preview": (c.raw_text or "")[:2000],
}
for c in cmds
],
"lldp_neighbors": [
{
"local_if": n.local_if,
"remote_sys": n.remote_sys,
"remote_if": n.remote_if,
"remote_ip": n.remote_ip,
"protocol": n.protocol,
}
for n in neighbors
],
}
def export_batch_zip(db: Session, batch_id: str) -> bytes:
detail = get_batch(db, batch_id)
buf = io.BytesIO()
with zipfile.ZipFile(buf, "w", compression=zipfile.ZIP_DEFLATED) as zf:
# manifest
lines = [
f"batch_id={detail['id']}",
f"task_id={detail['task_id']}",
f"status={detail['status']}",
f"commands={detail['command_count']}",
f"rows={detail['row_count']}",
"",
"commands:",
]
for c in detail["commands"]:
lines.append(
f"- {c['raw_command']} | parse={c['parse_status']} | "
f"profile={c['profile_id']} | rows={c['row_count']}"
)
zf.writestr("manifest.txt", "\n".join(lines) + "\n")
for c in detail["commands"]:
safe = "".join(ch if ch.isalnum() or ch in "-_" else "_" for ch in c["raw_command"])[:80]
zf.writestr(f"raw/{c['id']}_{safe}.txt", c.get("raw_text_preview") or "")
# full raw from DB
row = db.get(BizStateBatchCommand, c["id"])
if row and row.raw_text:
zf.writestr(f"raw/{c['id']}_{safe}.full.txt", row.raw_text)
# CSV
csv_lines = ["local_if,remote_sys,remote_if,remote_ip,protocol"]
for n in detail["lldp_neighbors"]:
csv_lines.append(
",".join(
[
_csv(n["local_if"]),
_csv(n["remote_sys"]),
_csv(n["remote_if"]),
_csv(n["remote_ip"]),
_csv(n["protocol"]),
]
)
)
zf.writestr("tables/lldp_neighbor.csv", "\n".join(csv_lines) + "\n")
return buf.getvalue()
def _csv(v: str) -> str:
s = str(v or "")
if any(ch in s for ch in ",\"\n"):
return '"' + s.replace('"', '""') + '"'
return s
def preview_items(db: Session, *, vendor: str, device_type: str, items: list[dict[str, Any]]) -> list[dict[str, Any]]:
vkey = resolve_vendor_key(vendor, device_type)
out = []
for raw in items:
out.append(
preview_task_item(
vendor_key=vkey,
profile_id=str(raw.get("source_profile_id") or ""),
command=str(raw.get("command_override") or raw.get("command") or ""),
bindings=list(raw.get("bindings") or []),
kind=str(raw.get("kind") or "catalog"),
)
)
return out

View file

@ -0,0 +1,165 @@
"""HTTP API for business state monitoring (/v1/biz-state)."""
from __future__ import annotations
from typing import Any
from fastapi import APIRouter, BackgroundTasks, Depends
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
from .db import get_db
from .biz_state import service as svc
from .biz_state.collect_runner import dispatch_collect
from .lldp_shared import resolve_vendor_key
from .models import BizStateTask
router = APIRouter(prefix="/v1/biz-state", tags=["biz-state"])
class ProfileOverrideIn(BaseModel):
title: str | None = None
command_template: str | None = None
description: str | None = None
sample_output: str | None = None
enabled: bool | None = None
class TaskItemIn(BaseModel):
source_profile_id: str = ""
kind: str = "catalog"
enabled: bool = True
title: str = ""
command_override: str = ""
sort_order: int | None = None
bindings: list[dict[str, str]] = Field(default_factory=list)
class TaskCreateIn(BaseModel):
source: str = "managed"
ne_id: str
ne_name: str = ""
ne_ip: str = ""
vendor: str = ""
device_type: str = ""
note: str = ""
interval_sec: int = 300
retention_batches: int = 30
items: list[TaskItemIn] = Field(default_factory=list)
class TaskPatchIn(BaseModel):
note: str | None = None
interval_sec: int | None = None
retention_batches: int | None = None
status: str | None = None
items: list[TaskItemIn] | None = None
class PreviewIn(BaseModel):
vendor: str = ""
device_type: str = ""
items: list[TaskItemIn] = Field(default_factory=list)
@router.get("/profiles")
def api_list_profiles(
vendor: str = "",
device_type: str = "",
db: Session = Depends(get_db),
) -> dict[str, Any]:
vkey = resolve_vendor_key(vendor, device_type) if (vendor or device_type) else ""
return {"items": svc.list_profiles_public(db, vendor_key=vkey)}
@router.patch("/profiles/{profile_id}")
def api_patch_profile(
profile_id: str,
body: ProfileOverrideIn,
db: Session = Depends(get_db),
) -> dict[str, Any]:
return svc.upsert_profile_override(db, profile_id, body.model_dump(exclude_unset=True))
@router.post("/tasks/preview-items")
def api_preview_items(
body: PreviewIn,
db: Session = Depends(get_db),
) -> dict[str, Any]:
items = [i.model_dump() for i in body.items]
return {
"items": svc.preview_items(
db, vendor=body.vendor, device_type=body.device_type, items=items
)
}
@router.get("/tasks")
def api_list_tasks(db: Session = Depends(get_db)) -> dict[str, Any]:
return {"items": svc.list_tasks(db)}
@router.post("/tasks")
def api_create_task(body: TaskCreateIn, db: Session = Depends(get_db)) -> dict[str, Any]:
return svc.create_task(db, body.model_dump())
@router.get("/tasks/{task_id}")
def api_get_task(task_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
return svc.get_task(db, task_id)
@router.patch("/tasks/{task_id}")
def api_patch_task(
task_id: str, body: TaskPatchIn, db: Session = Depends(get_db)
) -> dict[str, Any]:
return svc.update_task(db, task_id, body.model_dump(exclude_unset=True))
@router.delete("/tasks/{task_id}")
def api_delete_task(task_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
svc.delete_task(db, task_id)
return {"ok": True}
@router.post("/tasks/{task_id}/collect")
def api_collect_now(
task_id: str,
background_tasks: BackgroundTasks,
db: Session = Depends(get_db),
) -> dict[str, Any]:
task = db.get(BizStateTask, task_id)
if not task:
from fastapi import HTTPException
raise HTTPException(status_code=404, detail="task not found")
if bool(task.collect_running):
return {"ok": True, "started": False, "reason": "already_collecting", "task_id": task_id}
# Allow one-shot from draft/paused
task.last_collect_ended_at = None
db.commit()
background_tasks.add_task(dispatch_collect, task_id)
return {"ok": True, "started": True, "task_id": task_id}
@router.get("/tasks/{task_id}/batches")
def api_list_batches(
task_id: str, limit: int = 50, db: Session = Depends(get_db)
) -> dict[str, Any]:
return {"items": svc.list_batches(db, task_id, limit=limit)}
@router.get("/batches/{batch_id}")
def api_get_batch(batch_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
return svc.get_batch(db, batch_id)
@router.get("/batches/{batch_id}/export")
def api_export_batch(batch_id: str, db: Session = Depends(get_db)) -> StreamingResponse:
data = svc.export_batch_zip(db, batch_id)
return StreamingResponse(
iter([data]),
media_type="application/zip",
headers={"Content-Disposition": f'attachment; filename="biz_state_{batch_id}.zip"'},
)

View file

@ -0,0 +1,138 @@
"""Background scheduler for biz_state collection."""
from __future__ import annotations
import logging
import threading
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from .cli_budget import clamp_cli_workers
from .config import settings
from .db import SessionLocal
from .models import BizStateTask
from .biz_state.collect_runner import dispatch_collect
_log = logging.getLogger("netx.biz_state.scheduler")
_stop = threading.Event()
_thread: threading.Thread | None = None
_dispatch_pool: ThreadPoolExecutor | None = None
_pool_lock = threading.Lock()
_last_tick_mono: float = 0.0
def _utcnow() -> datetime:
return datetime.utcnow()
def _dispatch_pool_get() -> ThreadPoolExecutor:
global _dispatch_pool
with _pool_lock:
if _dispatch_pool is None:
workers = clamp_cli_workers(
int(getattr(settings, "biz_state_dispatch_workers", 2) or 2),
)
_dispatch_pool = ThreadPoolExecutor(
max_workers=workers, thread_name_prefix="biz-dispatch"
)
return _dispatch_pool
def shutdown_biz_state_dispatch_pool(*, wait: bool = False) -> None:
global _dispatch_pool
with _pool_lock:
if _dispatch_pool is not None:
try:
_dispatch_pool.shutdown(wait=wait, cancel_futures=True)
except TypeError:
_dispatch_pool.shutdown(wait=wait)
_dispatch_pool = None
def try_dispatch_due_tasks() -> int:
import time as _time
global _last_tick_mono
_last_tick_mono = _time.monotonic()
db = SessionLocal()
try:
tasks = (
db.query(BizStateTask)
.filter(
BizStateTask.status == "running",
BizStateTask.collect_running.is_(False),
)
.all()
)
due_ids: list[str] = []
now = _utcnow()
for task in tasks:
interval = max(60, int(task.interval_sec or 300))
ended = task.last_collect_ended_at
if ended is None:
due_ids.append(str(task.id))
continue
if (now - ended).total_seconds() >= interval:
due_ids.append(str(task.id))
finally:
db.close()
if not due_ids:
return 0
pool = _dispatch_pool_get()
for tid in due_ids:
try:
pool.submit(dispatch_collect, tid)
except Exception:
_log.exception("biz_state submit failed task=%s", tid)
return len(due_ids)
def _loop() -> None:
tick = max(5, int(getattr(settings, "biz_state_scheduler_tick_sec", 15) or 15))
while not _stop.wait(tick):
if not bool(getattr(settings, "biz_state_scheduler_enabled", True)):
continue
try:
try_dispatch_due_tasks()
except Exception:
_log.exception("biz_state scheduler tick failed")
def start_biz_state_scheduler() -> None:
global _thread
if not bool(getattr(settings, "biz_state_scheduler_enabled", True)):
_log.info("biz_state scheduler disabled")
return
if _thread and _thread.is_alive():
return
_stop.clear()
_thread = threading.Thread(target=_loop, name="biz-state-scheduler", daemon=True)
_thread.start()
_log.info("biz_state scheduler started")
def stop_biz_state_scheduler() -> None:
_stop.set()
shutdown_biz_state_dispatch_pool(wait=False)
global _thread
t = _thread
_thread = None
if t and t.is_alive():
t.join(timeout=2.0)
_log.info("biz_state scheduler stopped")
def biz_state_scheduler_status() -> dict:
import time as _time
alive = bool(_thread and _thread.is_alive())
age = None
if _last_tick_mono:
age = max(0.0, _time.monotonic() - _last_tick_mono)
return {
"running": alive,
"enabled": bool(getattr(settings, "biz_state_scheduler_enabled", True)),
"last_tick_age_sec": age,
}

View file

@ -103,6 +103,10 @@ class Settings(BaseSettings):
# Port traffic monitoring (CLI rate bit/s samples) # Port traffic monitoring (CLI rate bit/s samples)
port_traffic_scheduler_enabled: bool = True port_traffic_scheduler_enabled: bool = True
port_traffic_scheduler_tick_sec: int = 15 port_traffic_scheduler_tick_sec: int = 15
# Business state monitoring (LLDP snapshot batches; Phase1)
biz_state_scheduler_enabled: bool = True
biz_state_scheduler_tick_sec: int = 15
biz_state_dispatch_workers: int = 2
# Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec). # Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec).
ne_exec_max_commands: int = 5 ne_exec_max_commands: int = 5
# WebCRT interactive terminal sessions (multi-operator concurrent terminals). # WebCRT interactive terminal sessions (multi-operator concurrent terminals).

35
netx_api/lldp_shared.py Normal file
View file

@ -0,0 +1,35 @@
"""Shared LLDP parse API (template + TextFSM + normalize).
Topology discovery and biz_state monitoring both call this module.
Downstream flows diverge: Fabric edges vs biz_state batches.
"""
from __future__ import annotations
from .topology_lldp import (
STUB_PARSER_KEYS,
VENDOR_LLDP_PROFILES,
NeighborHit,
VendorLldpProfile,
can_discover_lldp,
get_vendor_profile,
lldp_command_for_vendor,
parse_neighbor_output,
parser_meta,
pick_neighbor_command,
resolve_vendor_key,
)
__all__ = [
"NeighborHit",
"VendorLldpProfile",
"VENDOR_LLDP_PROFILES",
"STUB_PARSER_KEYS",
"resolve_vendor_key",
"get_vendor_profile",
"lldp_command_for_vendor",
"pick_neighbor_command",
"can_discover_lldp",
"parser_meta",
"parse_neighbor_output",
]

View file

@ -23,6 +23,7 @@ from .managed_ne_router import router as managed_ne_router
from .ops_router import router as ops_router from .ops_router import router as ops_router
from .parser_config import load_parser_config from .parser_config import load_parser_config
from .port_traffic_router import router as port_traffic_router from .port_traffic_router import router as port_traffic_router
from .biz_state_router import router as biz_state_router
from .sql_router import router as sql_router from .sql_router import router as sql_router
from .sql_router import sql_query, sql_ume_query # noqa: F401 — tests import from main from .sql_router import sql_query, sql_ume_query # noqa: F401 — tests import from main
from .topology_router import router as topology_router from .topology_router import router as topology_router
@ -74,6 +75,7 @@ app.include_router(cli_router)
app.include_router(collection_router) app.include_router(collection_router)
app.include_router(config_sync_router) app.include_router(config_sync_router)
app.include_router(port_traffic_router) app.include_router(port_traffic_router)
app.include_router(biz_state_router)
app.include_router(webcrt_router) app.include_router(webcrt_router)
app.include_router(topology_router) app.include_router(topology_router)
app.include_router(lldp_collect_router) app.include_router(lldp_collect_router)

View file

@ -24,6 +24,16 @@ from .managed_ne import (
NeCollectionRun, NeCollectionRun,
UmeCliOverride, UmeCliOverride,
) )
from .biz_state import (
BizStateBatch,
BizStateBatchCommand,
BizStateCommandOverride,
BizStateEvent,
BizStateLldpNeighbor,
BizStateTask,
BizStateTaskItem,
BizStateTaskItemBinding,
)
from .port_traffic import ( from .port_traffic import (
PortTrafficBoard, PortTrafficBoard,
PortTrafficDevice, PortTrafficDevice,
@ -112,4 +122,12 @@ __all__ = [
"PortTrafficEvent", "PortTrafficEvent",
"PortTrafficBoard", "PortTrafficBoard",
"PortTrafficPanel", "PortTrafficPanel",
"BizStateTask",
"BizStateTaskItem",
"BizStateTaskItemBinding",
"BizStateBatch",
"BizStateBatchCommand",
"BizStateLldpNeighbor",
"BizStateEvent",
"BizStateCommandOverride",
] ]

View file

@ -0,0 +1,151 @@
"""biz_state ORM: tasks, items, bindings, batches, LLDP neighbor rows."""
from __future__ import annotations
from datetime import datetime
from uuid import uuid4
from sqlalchemy import Boolean, DateTime, Integer, String, Text, UniqueConstraint
from sqlalchemy.orm import Mapped, mapped_column
from ..db import Base
from ..timeutil import utcnow_naive
from ._types import JsonType as _JsonType
class BizStateTask(Base):
"""Per-NE business state monitoring config."""
__tablename__ = "biz_state_task"
__table_args__ = (UniqueConstraint("source", "ne_id", name="uq_biz_state_task_ne"),)
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
source: Mapped[str] = mapped_column(String(32), default="managed", index=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
ne_name: Mapped[str] = mapped_column(String(256), default="")
ne_ip: Mapped[str] = mapped_column(String(128), default="")
vendor: Mapped[str] = mapped_column(String(64), default="")
device_type: Mapped[str] = mapped_column(String(64), default="")
note: Mapped[str] = mapped_column(String(256), default="")
status: Mapped[str] = mapped_column(String(32), default="draft", index=True) # draft|running|paused|stopped
interval_sec: Mapped[int] = mapped_column(Integer, default=300)
retention_batches: Mapped[int] = mapped_column(Integer, default=30)
collect_running: Mapped[bool] = mapped_column(Boolean, default=False)
last_collect_started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
last_collect_ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
last_error: Mapped[str] = mapped_column(String(1024), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)
class BizStateTaskItem(Base):
"""Enabled collect row on a task (catalog profile or custom_raw)."""
__tablename__ = "biz_state_task_item"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
task_id: Mapped[str] = mapped_column(String(64), default="", index=True)
source_profile_id: Mapped[str] = mapped_column(String(128), default="", index=True)
kind: Mapped[str] = mapped_column(String(32), default="catalog") # catalog|custom_raw
enabled: Mapped[bool] = mapped_column(Boolean, default=True)
title: Mapped[str] = mapped_column(String(256), default="")
command_override: Mapped[str] = mapped_column(String(512), default="")
sort_order: Mapped[int] = mapped_column(Integer, default=0)
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)
class BizStateTaskItemBinding(Base):
"""Placeholder binding for parameterized profiles (Phase3 UX; model ready)."""
__tablename__ = "biz_state_task_item_binding"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
item_id: Mapped[str] = mapped_column(String(64), default="", index=True)
placeholder: Mapped[str] = mapped_column(String(64), default="")
value: Mapped[str] = mapped_column(String(256), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)
class BizStateBatch(Base):
"""One collect round snapshot header."""
__tablename__ = "biz_state_batch"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
task_id: Mapped[str] = mapped_column(String(64), default="", index=True)
source: Mapped[str] = mapped_column(String(32), default="managed", index=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
ne_name: Mapped[str] = mapped_column(String(256), default="")
vendor: Mapped[str] = mapped_column(String(64), default="")
status: Mapped[str] = mapped_column(String(32), default="running", index=True) # running|success|partial|failed
command_count: Mapped[int] = mapped_column(Integer, default=0)
row_count: Mapped[int] = mapped_column(Integer, default=0)
message: Mapped[str] = mapped_column(String(1024), default="")
started_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True)
ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
class BizStateBatchCommand(Base):
"""Per-command audit inside a batch."""
__tablename__ = "biz_state_batch_command"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
batch_id: Mapped[str] = mapped_column(String(64), default="", index=True)
task_item_id: Mapped[str] = mapped_column(String(64), default="", index=True)
profile_id: Mapped[str] = mapped_column(String(128), default="")
parser_id: Mapped[str] = mapped_column(String(128), default="")
metric_id: Mapped[str] = mapped_column(String(64), default="", index=True)
raw_command: Mapped[str] = mapped_column(String(512), default="")
params_json: Mapped[dict] = mapped_column(_JsonType, default=dict)
parse_status: Mapped[str] = mapped_column(String(32), default="") # ok|unmatched|failed|skipped_custom
row_count: Mapped[int] = mapped_column(Integer, default=0)
raw_text: Mapped[str] = mapped_column(Text, default="")
message: Mapped[str] = mapped_column(String(1024), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)
class BizStateLldpNeighbor(Base):
"""Structured LLDP neighbor rows for a batch (cutover-compare ready)."""
__tablename__ = "biz_state_lldp_neighbor"
__table_args__ = (
UniqueConstraint("batch_id", "local_if", "remote_sys", "remote_if", name="uq_biz_lldp_row"),
)
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
batch_id: Mapped[str] = mapped_column(String(64), default="", index=True)
batch_command_id: Mapped[str] = mapped_column(String(64), default="", index=True)
task_id: Mapped[str] = mapped_column(String(64), default="", index=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
local_if: Mapped[str] = mapped_column(String(128), default="", index=True)
remote_sys: Mapped[str] = mapped_column(String(256), default="", index=True)
remote_if: Mapped[str] = mapped_column(String(128), default="")
remote_ip: Mapped[str] = mapped_column(String(128), default="")
protocol: Mapped[str] = mapped_column(String(32), default="lldp")
collected_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True)
class BizStateEvent(Base):
__tablename__ = "biz_state_event"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
task_id: Mapped[str] = mapped_column(String(64), default="", index=True)
level: Mapped[str] = mapped_column(String(16), default="error", index=True)
message: Mapped[str] = mapped_column(Text, default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True)
class BizStateCommandOverride(Base):
"""Hot-edit overlay for profile display / template / enabled."""
__tablename__ = "biz_state_command_override"
id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex)
profile_id: Mapped[str] = mapped_column(String(128), unique=True, index=True)
title: Mapped[str] = mapped_column(String(256), default="")
command_template: Mapped[str] = mapped_column(String(512), default="")
description: Mapped[str] = mapped_column(Text, default="")
sample_output: Mapped[str] = mapped_column(Text, default="")
enabled: Mapped[bool | None] = mapped_column(Boolean, nullable=True)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive)

View file

@ -67,6 +67,12 @@ def local_device_scheduler_status(*, role: str = "unknown") -> dict[str, Any]:
out["port_traffic"] = port_traffic_scheduler_status() out["port_traffic"] = port_traffic_scheduler_status()
except Exception: # noqa: BLE001 except Exception: # noqa: BLE001
out["port_traffic"] = {"running": False, "error": "unavailable"} out["port_traffic"] = {"running": False, "error": "unavailable"}
try:
from .biz_state_scheduler import biz_state_scheduler_status
out["biz_state"] = biz_state_scheduler_status()
except Exception: # noqa: BLE001
out["biz_state"] = {"running": False, "error": "unavailable"}
try: try:
from .fabric_reconcile_scheduler import fabric_reconcile_scheduler_status from .fabric_reconcile_scheduler import fabric_reconcile_scheduler_status

View file

@ -25,7 +25,7 @@ from .topology_fabric import (
ensure_fabric_node_for_ume, ensure_fabric_node_for_ume,
upsert_fabric_edge, upsert_fabric_edge,
) )
from .topology_lldp import ( from .lldp_shared import (
NeighborHit, NeighborHit,
can_discover_lldp, can_discover_lldp,
parse_neighbor_output, parse_neighbor_output,

View file

@ -48,12 +48,14 @@ def start_device_schedulers() -> None:
from .lldp_collect_scheduler import start_lldp_collect_scheduler from .lldp_collect_scheduler import start_lldp_collect_scheduler
from .ne_collect_scheduler import start_ne_collect_scheduler from .ne_collect_scheduler import start_ne_collect_scheduler
from .port_traffic_scheduler import start_port_traffic_scheduler from .port_traffic_scheduler import start_port_traffic_scheduler
from .biz_state_scheduler import start_biz_state_scheduler
from .scheduler_heartbeat import start_scheduler_heartbeat_publisher from .scheduler_heartbeat import start_scheduler_heartbeat_publisher
start_config_sync_scheduler() start_config_sync_scheduler()
start_lldp_collect_scheduler() start_lldp_collect_scheduler()
start_ne_collect_scheduler() start_ne_collect_scheduler()
start_port_traffic_scheduler() start_port_traffic_scheduler()
start_biz_state_scheduler()
start_fabric_reconcile_scheduler() start_fabric_reconcile_scheduler()
# Publish status so API /metrics can see collectors when run in a split worker. # Publish status so API /metrics can see collectors when run in a split worker.
role = "api_inline" if bool(getattr(settings, "run_inline_schedulers", True)) else "worker" role = "api_inline" if bool(getattr(settings, "run_inline_schedulers", True)) else "worker"

View file

@ -0,0 +1,63 @@
"""Unit tests for biz_state ParseProfile match / expand (no device)."""
from __future__ import annotations
import unittest
from netx_api.biz_state.command_match import expand_from_bindings, match_command, preview_task_item
from netx_api.biz_state.profiles import all_profiles, get_profile, profiles_for_vendor
from netx_api.biz_state.parsers import normalize_lldp_neighbors
from netx_api.lldp_shared import parse_neighbor_output
class BizStateProfileTests(unittest.TestCase):
def test_lldp_profiles_registered(self) -> None:
profiles = all_profiles()
self.assertTrue(any(p.metric_id == "lldp_neighbor" for p in profiles))
zte = profiles_for_vendor("zte")
self.assertTrue(any(p.profile_id == "zte.lldp_neighbors" for p in zte))
def test_match_zte_lldp_command(self) -> None:
hit = match_command(vendor_key="zte", command="show lldp neighbor brief")
self.assertIsNotNone(hit)
assert hit is not None
self.assertEqual(hit.profile.parser_id, "lldp_neighbors")
self.assertEqual(hit.profile.metric_id, "lldp_neighbor")
def test_match_huawei_lldp_command(self) -> None:
hit = match_command(vendor_key="huawei", command="display lldp neighbor")
self.assertIsNotNone(hit)
assert hit is not None
self.assertEqual(hit.profile.profile_id, "huawei.lldp_neighbors")
def test_expand_no_placeholder(self) -> None:
p = get_profile("zte.lldp_neighbors")
self.assertIsNotNone(p)
assert p is not None
pairs = expand_from_bindings(profile=p)
self.assertEqual(len(pairs), 1)
self.assertEqual(pairs[0][0], "show lldp neighbor brief")
def test_preview_custom_raw(self) -> None:
prev = preview_task_item(
vendor_key="zte",
kind="custom_raw",
command="show version",
)
self.assertTrue(prev["ok"])
self.assertEqual(prev["parse"], "skipped_custom")
def test_shared_parse_empty(self) -> None:
rows = normalize_lldp_neighbors(
raw_text="",
vendor="zte",
device_type="zte_zxros",
command="show lldp neighbor brief",
)
self.assertEqual(rows, [])
hits = parse_neighbor_output("", vendor="zte", device_type="zte_zxros")
self.assertEqual(hits, [])
if __name__ == "__main__":
unittest.main()

View file

@ -44,6 +44,7 @@ src/
| `/network/tasks/config-sync` | 配置同步 | `network` | | `/network/tasks/config-sync` | 配置同步 | `network` |
| `/network/tasks/port-traffic` | 端口流量监控(设备管理) | `network` | | `/network/tasks/port-traffic` | 端口流量监控(设备管理) | `network` |
| `/network/tasks/port-traffic/wall` | 流量大屏列表(打开独立页签) | `network` | | `/network/tasks/port-traffic/wall` | 流量大屏列表(打开独立页签) | `network` |
| `/network/tasks/biz-state` | 业务状态监控(LLDP 快照 Phase1) | `network` |
| `/port-traffic/wall/:boardId` | 流量大屏专有页签(无网络侧栏) | `port-traffic-wall` | | `/port-traffic/wall/:boardId` | 流量大屏专有页签(无网络侧栏) | `port-traffic-wall` |
| `/topology` | 拓扑管理(模式化编辑器:选择/平移/拖动/连线、框选、自动布局、拖放添加) | `topology` | | `/topology` | 拓扑管理(模式化编辑器:选择/平移/拖动/连线、框选、自动布局、拖放添加) | `topology` |
| `/webcrt` | WebCRT 终端 | `webcrt` | | `/webcrt` | WebCRT 终端 | `webcrt` |
@ -124,6 +125,16 @@ src/
- 采集日志:`GET /v1/port-traffic/devices/{id}/events`;失败写入 `port_traffic_event`,列表操作可查看 - 采集日志:`GET /v1/port-traffic/devices/{id}/events`;失败写入 `port_traffic_event`,列表操作可查看
- 支持拓扑深链:`?ne_id=&source=managed|ume&ifname=` 打开向导并预填网元 - 支持拓扑深链:`?ne_id=&source=managed|ume&ifname=` 打开向导并预填网元
## 业务状态监控(biz_state)
- API:`/v1/biz-state/profiles`、`/tasks*`、`/batches*`、`/batches/{id}/export`
- **ParseProfile**:命令模板 + TextFSM + 回调 + schema;LLDP 与拓扑 **共享解析**(`lldp_shared`),业务流程写批次表,拓扑写 Fabric
- Phase1 样板:LLDP 邻居快照;建任务默认启用对应厂商 LLDP profile;支持自定义只采不解析行
- 调度:`NETX_BIZ_STATE_SCHEDULER_ENABLED`(默认开),tick `NETX_BIZ_STATE_SCHEDULER_TICK_SEC`
- 前端:`/network/tasks/biz-state`(任务列表 / 启用项 / 批次 / 导出 zip)
- Phase2(未做):比对模板、端口映射、前后批次 diff、大屏
- Phase3(未做):占位符发现→人选关联(如 VRF)
## 拓扑管理(Fabric + 站点目录,对齐厂商) ## 拓扑管理(Fabric + 站点目录,对齐厂商)
- 三层库存:**运维** `managed_ne`(凭据/采集/WebCRT)· **EMS** `ume_inventory_ne` · **拓扑投影** `topo_fabric_node`(上图/LLDP/分类)。Fabric 由 ensure 按需创建,与运维表不是同一张表。 - 三层库存:**运维** `managed_ne`(凭据/采集/WebCRT)· **EMS** `ume_inventory_ne` · **拓扑投影** `topo_fabric_node`(上图/LLDP/分类)。Fabric 由 ensure 按需创建,与运维表不是同一张表。

View file

@ -43,6 +43,9 @@ const PortTrafficBoardListPage = lazy(() =>
const PortTrafficWallPage = lazy(() => const PortTrafficWallPage = lazy(() =>
import("./pages/network/PortTrafficWallPage").then((m) => ({ default: m.PortTrafficWallPage })), import("./pages/network/PortTrafficWallPage").then((m) => ({ default: m.PortTrafficWallPage })),
); );
const BizStatePage = lazy(() =>
import("./pages/network/BizStatePage").then((m) => ({ default: m.BizStatePage })),
);
const UsersPage = lazy(() => import("./pages/UsersPage").then((m) => ({ default: m.UsersPage }))); const UsersPage = lazy(() => import("./pages/UsersPage").then((m) => ({ default: m.UsersPage })));
const AuditLayout = lazy(() => const AuditLayout = lazy(() =>
import("./pages/audit/AuditLayout").then((m) => ({ default: m.AuditLayout })), import("./pages/audit/AuditLayout").then((m) => ({ default: m.AuditLayout })),
@ -120,6 +123,7 @@ function ProtectedApp() {
<Route path="tasks/config-sync" element={<ConfigSyncPage />} /> <Route path="tasks/config-sync" element={<ConfigSyncPage />} />
<Route path="tasks/port-traffic/wall" element={<LegacyPortTrafficWallRedirect />} /> <Route path="tasks/port-traffic/wall" element={<LegacyPortTrafficWallRedirect />} />
<Route path="tasks/port-traffic" element={<PortTrafficPage />} /> <Route path="tasks/port-traffic" element={<PortTrafficPage />} />
<Route path="tasks/biz-state" element={<BizStatePage />} />
</Route> </Route>
<Route path="/collect" element={<Navigate to="/network/tasks/collect" replace />} /> <Route path="/collect" element={<Navigate to="/network/tasks/collect" replace />} />
<Route path="/users" element={<UsersPage />} /> <Route path="/users" element={<UsersPage />} />

View file

@ -74,6 +74,12 @@ export const NETWORK_NAV: readonly NetworkNavGroup[] = [
labelKey: "network.nav.portTrafficWall", labelKey: "network.nav.portTrafficWall",
group: "tasks", group: "tasks",
}, },
{
id: "biz-state",
path: "/network/tasks/biz-state",
labelKey: "network.nav.bizState",
group: "tasks",
},
], ],
}, },
] as const; ] as const;

View file

@ -156,6 +156,7 @@ const en = {
configSync: "Config sync", configSync: "Config sync",
portTraffic: "Port traffic", portTraffic: "Port traffic",
portTrafficWall: "Traffic wall", portTrafficWall: "Traffic wall",
bizState: "Business state",
}, },
collapseNav: "Collapse sidebar", collapseNav: "Collapse sidebar",
expandNav: "Expand sidebar", expandNav: "Expand sidebar",

View file

@ -156,6 +156,7 @@ const zh = {
configSync: "配置同步", configSync: "配置同步",
portTraffic: "流量监控", portTraffic: "流量监控",
portTrafficWall: "流量大屏", portTrafficWall: "流量大屏",
bizState: "业务状态监控",
}, },
collapseNav: "折叠侧栏", collapseNav: "折叠侧栏",
expandNav: "展开侧栏", expandNav: "展开侧栏",

View file

@ -0,0 +1,382 @@
import { useCallback, useEffect, useState } from "react";
import {
bizStateCollectNow,
bizStateCreateTask,
bizStateDeleteTask,
bizStateDownloadExport,
bizStateGetBatch,
bizStateGetTask,
bizStateListBatches,
bizStateListProfiles,
bizStateListTasks,
bizStatePatchTask,
fetchManagedNe,
} from "../../services/api";
import type { ManagedNeItem } from "../../types";
type TaskRow = {
id: string;
ne_name: string;
ne_ip: string;
vendor: string;
status: string;
collect_running: boolean;
last_error: string;
last_collect_ended_at?: string | null;
};
type Profile = {
profile_id: string;
title: string;
command_template: string;
description: string;
metric_id: string;
};
type BatchRow = {
id: string;
status: string;
row_count: number;
command_count: number;
started_at?: string | null;
};
export function BizStatePage() {
const [tasks, setTasks] = useState<TaskRow[]>([]);
const [selectedId, setSelectedId] = useState("");
const [detail, setDetail] = useState<any>(null);
const [batches, setBatches] = useState<BatchRow[]>([]);
const [batchDetail, setBatchDetail] = useState<any>(null);
const [profiles, setProfiles] = useState<Profile[]>([]);
const [nes, setNes] = useState<ManagedNeItem[]>([]);
const [neId, setNeId] = useState("");
const [busy, setBusy] = useState(false);
const [err, setErr] = useState("");
const refreshTasks = useCallback(async () => {
const res = await bizStateListTasks();
setTasks((res.items || []) as TaskRow[]);
}, []);
useEffect(() => {
void (async () => {
try {
await refreshTasks();
const neRes = await fetchManagedNe({
keyword: "",
vendor: "",
connectStatus: "",
page: 1,
pageSize: 200,
});
setNes(neRes.items || []);
} catch (e: any) {
setErr(String(e?.message || e));
}
})();
}, [refreshTasks]);
const openTask = async (id: string) => {
setSelectedId(id);
setBatchDetail(null);
setErr("");
try {
const t = await bizStateGetTask(id);
setDetail(t);
const b = await bizStateListBatches(id);
setBatches((b.items || []) as BatchRow[]);
const p = await bizStateListProfiles({
vendor: t.vendor || "",
device_type: t.device_type || "",
});
setProfiles((p.items || []) as Profile[]);
} catch (e: any) {
setErr(String(e?.message || e));
}
};
const createTask = async () => {
if (!neId) return;
const ne = nes.find((n) => n.id === neId);
if (!ne) return;
setBusy(true);
setErr("");
try {
const t = await bizStateCreateTask({
source: "managed",
ne_id: ne.id,
ne_name: ne.name,
ne_ip: ne.ip_address,
vendor: ne.vendor,
device_type: ne.device_type,
});
await refreshTasks();
await openTask(t.id);
} catch (e: any) {
setErr(String(e?.message || e));
} finally {
setBusy(false);
}
};
const setStatus = async (status: string) => {
if (!selectedId) return;
setBusy(true);
try {
await bizStatePatchTask(selectedId, { status });
await openTask(selectedId);
await refreshTasks();
} catch (e: any) {
setErr(String(e?.message || e));
} finally {
setBusy(false);
}
};
const collectNow = async () => {
if (!selectedId) return;
setBusy(true);
setErr("");
try {
await bizStateCollectNow(selectedId);
await openTask(selectedId);
await refreshTasks();
// Poll a few times while collect_running
for (let i = 0; i < 20; i++) {
await new Promise((r) => setTimeout(r, 1500));
const t = await bizStateGetTask(selectedId);
setDetail(t);
if (!t.collect_running) {
const b = await bizStateListBatches(selectedId);
setBatches((b.items || []) as BatchRow[]);
break;
}
}
await refreshTasks();
} catch (e: any) {
setErr(String(e?.message || e));
} finally {
setBusy(false);
}
};
const openBatch = async (batchId: string) => {
try {
const d = await bizStateGetBatch(batchId);
setBatchDetail(d);
} catch (e: any) {
setErr(String(e?.message || e));
}
};
const removeTask = async () => {
if (!selectedId) return;
if (!window.confirm("删除该业务监控任务及所有批次?")) return;
setBusy(true);
try {
await bizStateDeleteTask(selectedId);
setSelectedId("");
setDetail(null);
setBatches([]);
setBatchDetail(null);
await refreshTasks();
} catch (e: any) {
setErr(String(e?.message || e));
} finally {
setBusy(false);
}
};
return (
<div className="wb-page">
<div className="wb-page__header">
<h1>业务状态监控</h1>
<p className="wb-muted">Phase1:LLDP 邻居快照采集(共享解析,独立批次;比对 Phase2)</p>
</div>
{err ? <div className="wb-alert wb-alert--error">{err}</div> : null}
<div className="wb-card" style={{ marginBottom: 16 }}>
<h3>新建任务</h3>
<div style={{ display: "flex", gap: 8, alignItems: "center", flexWrap: "wrap" }}>
<select value={neId} onChange={(e) => setNeId(e.target.value)}>
<option value="">选择托管网元…</option>
{nes.map((n) => (
<option key={n.id} value={n.id}>
{n.name || n.ip_address} ({n.vendor})
</option>
))}
</select>
<button type="button" disabled={busy || !neId} onClick={() => void createTask()}>
创建(默认启用 LLDP)
</button>
</div>
</div>
<div style={{ display: "grid", gridTemplateColumns: "1fr 1.4fr", gap: 16 }}>
<div className="wb-card">
<h3>任务列表</h3>
<table className="wb-table">
<thead>
<tr>
<th>网元</th>
<th>状态</th>
<th>最近采集</th>
</tr>
</thead>
<tbody>
{tasks.map((t) => (
<tr
key={t.id}
onClick={() => void openTask(t.id)}
style={{ cursor: "pointer", background: selectedId === t.id ? "var(--wb-row-active, #eef)" : undefined }}
>
<td>
{t.ne_name || t.ne_ip}
<div className="wb-muted">{t.vendor}</div>
</td>
<td>
{t.status}
{t.collect_running ? " · 采集中" : ""}
</td>
<td className="wb-muted">{t.last_collect_ended_at || "-"}</td>
</tr>
))}
{!tasks.length ? (
<tr>
<td colSpan={3} className="wb-muted">
暂无任务
</td>
</tr>
) : null}
</tbody>
</table>
</div>
<div className="wb-card">
{!detail ? (
<p className="wb-muted">选择左侧任务查看详情</p>
) : (
<>
<h3>
{detail.ne_name} · {detail.status}
</h3>
<div style={{ display: "flex", gap: 8, flexWrap: "wrap", marginBottom: 12 }}>
<button type="button" disabled={busy} onClick={() => void setStatus("running")}>
启动周期
</button>
<button type="button" disabled={busy} onClick={() => void setStatus("paused")}>
暂停
</button>
<button type="button" disabled={busy} onClick={() => void collectNow()}>
立即采集
</button>
<button type="button" disabled={busy} onClick={() => void removeTask()}>
删除
</button>
</div>
{detail.last_error ? <div className="wb-alert wb-alert--error">{detail.last_error}</div> : null}
<h4>启用项(勾选表)</h4>
<table className="wb-table">
<thead>
<tr>
<th>启用</th>
<th>项</th>
<th>命令</th>
<th>类型</th>
</tr>
</thead>
<tbody>
{(detail.items || []).map((it: any) => {
const prof = profiles.find((p) => p.profile_id === it.source_profile_id);
return (
<tr key={it.id}>
<td>{it.enabled ? "✓" : ""}</td>
<td>
{it.title || prof?.title || it.source_profile_id}
{prof?.description ? (
<div className="wb-muted" style={{ fontSize: 12 }}>
{prof.description}
</div>
) : null}
</td>
<td>
<code>{it.command_override || prof?.command_template || "-"}</code>
</td>
<td>{it.kind}</td>
</tr>
);
})}
</tbody>
</table>
<h4 style={{ marginTop: 16 }}>采集批次</h4>
<table className="wb-table">
<thead>
<tr>
<th>时间</th>
<th>状态</th>
<th>行数</th>
<th>操作</th>
</tr>
</thead>
<tbody>
{batches.map((b) => (
<tr key={b.id}>
<td>{b.started_at || "-"}</td>
<td>{b.status}</td>
<td>{b.row_count}</td>
<td style={{ display: "flex", gap: 8 }}>
<button type="button" onClick={() => void openBatch(b.id)}>
查看
</button>
<button type="button" onClick={() => void bizStateDownloadExport(b.id)}>
导出
</button>
</td>
</tr>
))}
{!batches.length ? (
<tr>
<td colSpan={4} className="wb-muted">
尚无批次
</td>
</tr>
) : null}
</tbody>
</table>
{batchDetail ? (
<div style={{ marginTop: 16 }}>
<h4>批次详情 {batchDetail.id}</h4>
<p className="wb-muted">
命令 {batchDetail.command_count} · 行 {batchDetail.row_count} · {batchDetail.status}
</p>
<table className="wb-table">
<thead>
<tr>
<th>本端口</th>
<th>对端系统</th>
<th>对端口</th>
<th>管理IP</th>
</tr>
</thead>
<tbody>
{(batchDetail.lldp_neighbors || []).slice(0, 200).map((n: any, i: number) => (
<tr key={i}>
<td>{n.local_if}</td>
<td>{n.remote_sys}</td>
<td>{n.remote_if}</td>
<td>{n.remote_ip}</td>
</tr>
))}
</tbody>
</table>
</div>
) : null}
</>
)}
</div>
</div>
</div>
);
}

View file

@ -1753,3 +1753,59 @@ export const deletePortTrafficBoard = (boardId: string) =>
apiDelete<{ ok: boolean; id: string }>( apiDelete<{ ok: boolean; id: string }>(
`/v1/port-traffic/boards/${encodeURIComponent(boardId)}`, `/v1/port-traffic/boards/${encodeURIComponent(boardId)}`,
); );
/* ---- biz_state (business state monitoring) ---- */
export const bizStateListProfiles = (params?: { vendor?: string; device_type?: string }) => {
const p = new URLSearchParams();
if (params?.vendor?.trim()) p.set("vendor", params.vendor.trim());
if (params?.device_type?.trim()) p.set("device_type", params.device_type.trim());
const q = p.toString();
return apiGet<{ items: Record<string, unknown>[] }>(`/v1/biz-state/profiles${q ? `?${q}` : ""}`);
};
export const bizStateListTasks = () =>
apiGet<{ items: Record<string, unknown>[] }>("/v1/biz-state/tasks");
export const bizStateCreateTask = (body: Record<string, unknown>) =>
apiPost<Record<string, unknown>>("/v1/biz-state/tasks", body);
export const bizStateGetTask = (taskId: string) =>
apiGet<Record<string, unknown>>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}`);
export const bizStatePatchTask = (taskId: string, body: Record<string, unknown>) =>
apiPatch<Record<string, unknown>>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}`, body);
export const bizStateDeleteTask = (taskId: string) =>
apiDelete<{ ok: boolean }>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}`);
export const bizStateCollectNow = (taskId: string) =>
apiPost<{ ok: boolean }>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}/collect`, {});
export const bizStateListBatches = (taskId: string, limit = 50) =>
apiGet<{ items: Record<string, unknown>[] }>(
`/v1/biz-state/tasks/${encodeURIComponent(taskId)}/batches?limit=${limit}`,
);
export const bizStateGetBatch = (batchId: string) =>
apiGet<Record<string, unknown>>(`/v1/biz-state/batches/${encodeURIComponent(batchId)}`);
export const bizStateDownloadExport = async (batchId: string): Promise<void> => {
const path = `/v1/biz-state/batches/${encodeURIComponent(batchId)}/export`;
const res = await fetch(path, { method: "GET", credentials: fetchCreds, headers: authHeaders() });
if (res.status === 401) {
handleUnauthorized(path);
throw new Error("unauthorized");
}
if (!res.ok) throw new Error(`${res.status} export`);
const blob = await res.blob();
const url = URL.createObjectURL(blob);
try {
const a = document.createElement("a");
a.href = url;
a.download = `biz_state_${batchId}.zip`;
a.click();
} finally {
URL.revokeObjectURL(url);
}
};