From 4fa1b9b4dc48a3a621ac5042d528adfa55b6ba39 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 17 Sep 2026 16:17:04 +0800 Subject: [PATCH] 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 --- netx_api/app_shutdown.py | 7 + netx_api/app_startup.py | 22 ++ netx_api/biz_state/__init__.py | 3 + netx_api/biz_state/collect_runner.py | 486 +++++++++++++++++++++++++ netx_api/biz_state/command_match.py | 162 +++++++++ netx_api/biz_state/parsers/__init__.py | 47 +++ netx_api/biz_state/profiles.py | 213 +++++++++++ netx_api/biz_state/schema_ensure.py | 27 ++ netx_api/biz_state/service.py | 474 ++++++++++++++++++++++++ netx_api/biz_state_router.py | 165 +++++++++ netx_api/biz_state_scheduler.py | 138 +++++++ netx_api/config.py | 4 + netx_api/lldp_shared.py | 35 ++ netx_api/main.py | 2 + netx_api/models/__init__.py | 18 + netx_api/models/biz_state.py | 151 ++++++++ netx_api/scheduler_heartbeat.py | 6 + netx_api/topology_discover_scan.py | 2 +- netx_api/ume_runtime.py | 2 + tests/test_biz_state_profiles.py | 63 ++++ web/WEB.md | 11 + web/src/App.tsx | 4 + web/src/config/networkNav.ts | 6 + web/src/i18n/en.ts | 1 + web/src/i18n/zh.ts | 1 + web/src/pages/network/BizStatePage.tsx | 382 +++++++++++++++++++ web/src/services/api.ts | 56 +++ 27 files changed, 2487 insertions(+), 1 deletion(-) create mode 100644 netx_api/biz_state/__init__.py create mode 100644 netx_api/biz_state/collect_runner.py create mode 100644 netx_api/biz_state/command_match.py create mode 100644 netx_api/biz_state/parsers/__init__.py create mode 100644 netx_api/biz_state/profiles.py create mode 100644 netx_api/biz_state/schema_ensure.py create mode 100644 netx_api/biz_state/service.py create mode 100644 netx_api/biz_state_router.py create mode 100644 netx_api/biz_state_scheduler.py create mode 100644 netx_api/lldp_shared.py create mode 100644 netx_api/models/biz_state.py create mode 100644 tests/test_biz_state_profiles.py create mode 100644 web/src/pages/network/BizStatePage.tsx diff --git a/netx_api/app_shutdown.py b/netx_api/app_shutdown.py index 3570103..842ef1e 100644 --- a/netx_api/app_shutdown.py +++ b/netx_api/app_shutdown.py @@ -45,6 +45,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None: except Exception: # noqa: BLE001 _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: from .fabric_reconcile_scheduler import stop_fabric_reconcile_scheduler diff --git a/netx_api/app_startup.py b/netx_api/app_startup.py index 51cc941..c4e563e 100644 --- a/netx_api/app_startup.py +++ b/netx_api/app_startup.py @@ -62,6 +62,12 @@ def run_api_startup() -> None: apply_topology_schema_safety_net(conn) apply_collection_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: _log.exception("startup: auth/topology/collection/hop schema safety patches failed") if skip_ddl and alembic_ok: @@ -133,6 +139,22 @@ def run_api_startup() -> None: 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) + 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: from .port_traffic_migrate import backfill_port_traffic_series diff --git a/netx_api/biz_state/__init__.py b/netx_api/biz_state/__init__.py new file mode 100644 index 0000000..6e53038 --- /dev/null +++ b/netx_api/biz_state/__init__.py @@ -0,0 +1,3 @@ +"""Business state monitoring: ParseProfile collect → batch → (Phase2) compare.""" + +from __future__ import annotations diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py new file mode 100644 index 0000000..4d01dca --- /dev/null +++ b/netx_api/biz_state/collect_runner.py @@ -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} diff --git a/netx_api/biz_state/command_match.py b/netx_api/biz_state/command_match.py new file mode 100644 index 0000000..3fa668c --- /dev/null +++ b/netx_api/biz_state/command_match.py @@ -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": "", + } diff --git a/netx_api/biz_state/parsers/__init__.py b/netx_api/biz_state/parsers/__init__.py new file mode 100644 index 0000000..95a3597 --- /dev/null +++ b/netx_api/biz_state/parsers/__init__.py @@ -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()) diff --git a/netx_api/biz_state/profiles.py b/netx_api/biz_state/profiles.py new file mode 100644 index 0000000..9065c45 --- /dev/null +++ b/netx_api/biz_state/profiles.py @@ -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, + } diff --git a/netx_api/biz_state/schema_ensure.py b/netx_api/biz_state/schema_ensure.py new file mode 100644 index 0000000..8155098 --- /dev/null +++ b/netx_api/biz_state/schema_ensure.py @@ -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) diff --git a/netx_api/biz_state/service.py b/netx_api/biz_state/service.py new file mode 100644 index 0000000..426a0c4 --- /dev/null +++ b/netx_api/biz_state/service.py @@ -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 diff --git a/netx_api/biz_state_router.py b/netx_api/biz_state_router.py new file mode 100644 index 0000000..6a004c4 --- /dev/null +++ b/netx_api/biz_state_router.py @@ -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"'}, + ) diff --git a/netx_api/biz_state_scheduler.py b/netx_api/biz_state_scheduler.py new file mode 100644 index 0000000..2dd9243 --- /dev/null +++ b/netx_api/biz_state_scheduler.py @@ -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, + } diff --git a/netx_api/config.py b/netx_api/config.py index 4d67ab7..684e929 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -103,6 +103,10 @@ class Settings(BaseSettings): # Port traffic monitoring (CLI rate bit/s samples) port_traffic_scheduler_enabled: bool = True 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). ne_exec_max_commands: int = 5 # WebCRT interactive terminal sessions (multi-operator concurrent terminals). diff --git a/netx_api/lldp_shared.py b/netx_api/lldp_shared.py new file mode 100644 index 0000000..79c22ec --- /dev/null +++ b/netx_api/lldp_shared.py @@ -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", +] diff --git a/netx_api/main.py b/netx_api/main.py index b365699..6999391 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -23,6 +23,7 @@ from .managed_ne_router import router as managed_ne_router from .ops_router import router as ops_router from .parser_config import load_parser_config 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 sql_query, sql_ume_query # noqa: F401 — tests import from main 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(config_sync_router) app.include_router(port_traffic_router) +app.include_router(biz_state_router) app.include_router(webcrt_router) app.include_router(topology_router) app.include_router(lldp_collect_router) diff --git a/netx_api/models/__init__.py b/netx_api/models/__init__.py index cf03ace..932406a 100644 --- a/netx_api/models/__init__.py +++ b/netx_api/models/__init__.py @@ -24,6 +24,16 @@ from .managed_ne import ( NeCollectionRun, UmeCliOverride, ) +from .biz_state import ( + BizStateBatch, + BizStateBatchCommand, + BizStateCommandOverride, + BizStateEvent, + BizStateLldpNeighbor, + BizStateTask, + BizStateTaskItem, + BizStateTaskItemBinding, +) from .port_traffic import ( PortTrafficBoard, PortTrafficDevice, @@ -112,4 +122,12 @@ __all__ = [ "PortTrafficEvent", "PortTrafficBoard", "PortTrafficPanel", + "BizStateTask", + "BizStateTaskItem", + "BizStateTaskItemBinding", + "BizStateBatch", + "BizStateBatchCommand", + "BizStateLldpNeighbor", + "BizStateEvent", + "BizStateCommandOverride", ] diff --git a/netx_api/models/biz_state.py b/netx_api/models/biz_state.py new file mode 100644 index 0000000..6210d70 --- /dev/null +++ b/netx_api/models/biz_state.py @@ -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) diff --git a/netx_api/scheduler_heartbeat.py b/netx_api/scheduler_heartbeat.py index 5683ca0..4b0fcc8 100644 --- a/netx_api/scheduler_heartbeat.py +++ b/netx_api/scheduler_heartbeat.py @@ -67,6 +67,12 @@ def local_device_scheduler_status(*, role: str = "unknown") -> dict[str, Any]: out["port_traffic"] = port_traffic_scheduler_status() except Exception: # noqa: BLE001 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: from .fabric_reconcile_scheduler import fabric_reconcile_scheduler_status diff --git a/netx_api/topology_discover_scan.py b/netx_api/topology_discover_scan.py index ae8f428..76febea 100644 --- a/netx_api/topology_discover_scan.py +++ b/netx_api/topology_discover_scan.py @@ -25,7 +25,7 @@ from .topology_fabric import ( ensure_fabric_node_for_ume, upsert_fabric_edge, ) -from .topology_lldp import ( +from .lldp_shared import ( NeighborHit, can_discover_lldp, parse_neighbor_output, diff --git a/netx_api/ume_runtime.py b/netx_api/ume_runtime.py index 66ffe14..bc39676 100644 --- a/netx_api/ume_runtime.py +++ b/netx_api/ume_runtime.py @@ -48,12 +48,14 @@ def start_device_schedulers() -> None: 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 .biz_state_scheduler import start_biz_state_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_biz_state_scheduler() start_fabric_reconcile_scheduler() # 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" diff --git a/tests/test_biz_state_profiles.py b/tests/test_biz_state_profiles.py new file mode 100644 index 0000000..7cef8d0 --- /dev/null +++ b/tests/test_biz_state_profiles.py @@ -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() diff --git a/web/WEB.md b/web/WEB.md index 1662e18..d25ff6b 100644 --- a/web/WEB.md +++ b/web/WEB.md @@ -44,6 +44,7 @@ src/ | `/network/tasks/config-sync` | 配置同步 | `network` | | `/network/tasks/port-traffic` | 端口流量监控(设备管理) | `network` | | `/network/tasks/port-traffic/wall` | 流量大屏列表(打开独立页签) | `network` | +| `/network/tasks/biz-state` | 业务状态监控(LLDP 快照 Phase1) | `network` | | `/port-traffic/wall/:boardId` | 流量大屏专有页签(无网络侧栏) | `port-traffic-wall` | | `/topology` | 拓扑管理(模式化编辑器:选择/平移/拖动/连线、框选、自动布局、拖放添加) | `topology` | | `/webcrt` | WebCRT 终端 | `webcrt` | @@ -124,6 +125,16 @@ src/ - 采集日志:`GET /v1/port-traffic/devices/{id}/events`;失败写入 `port_traffic_event`,列表操作可查看 - 支持拓扑深链:`?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 + 站点目录,对齐厂商) - 三层库存:**运维** `managed_ne`(凭据/采集/WebCRT)· **EMS** `ume_inventory_ne` · **拓扑投影** `topo_fabric_node`(上图/LLDP/分类)。Fabric 由 ensure 按需创建,与运维表不是同一张表。 diff --git a/web/src/App.tsx b/web/src/App.tsx index f0cc8e7..e949005 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -43,6 +43,9 @@ const PortTrafficBoardListPage = lazy(() => const PortTrafficWallPage = lazy(() => 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 AuditLayout = lazy(() => import("./pages/audit/AuditLayout").then((m) => ({ default: m.AuditLayout })), @@ -120,6 +123,7 @@ function ProtectedApp() { } /> } /> } /> + } /> } /> } /> diff --git a/web/src/config/networkNav.ts b/web/src/config/networkNav.ts index cb6e696..2964c29 100644 --- a/web/src/config/networkNav.ts +++ b/web/src/config/networkNav.ts @@ -74,6 +74,12 @@ export const NETWORK_NAV: readonly NetworkNavGroup[] = [ labelKey: "network.nav.portTrafficWall", group: "tasks", }, + { + id: "biz-state", + path: "/network/tasks/biz-state", + labelKey: "network.nav.bizState", + group: "tasks", + }, ], }, ] as const; diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index c1e54b8..d062ef3 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -156,6 +156,7 @@ const en = { configSync: "Config sync", portTraffic: "Port traffic", portTrafficWall: "Traffic wall", + bizState: "Business state", }, collapseNav: "Collapse sidebar", expandNav: "Expand sidebar", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index d792e3c..c1bc95c 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -156,6 +156,7 @@ const zh = { configSync: "配置同步", portTraffic: "流量监控", portTrafficWall: "流量大屏", + bizState: "业务状态监控", }, collapseNav: "折叠侧栏", expandNav: "展开侧栏", diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx new file mode 100644 index 0000000..a77d001 --- /dev/null +++ b/web/src/pages/network/BizStatePage.tsx @@ -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([]); + const [selectedId, setSelectedId] = useState(""); + const [detail, setDetail] = useState(null); + const [batches, setBatches] = useState([]); + const [batchDetail, setBatchDetail] = useState(null); + const [profiles, setProfiles] = useState([]); + const [nes, setNes] = useState([]); + 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 ( +
+
+

业务状态监控

+

Phase1:LLDP 邻居快照采集(共享解析,独立批次;比对 Phase2)

+
+ {err ?
{err}
: null} + +
+

新建任务

+
+ + +
+
+ +
+
+

任务列表

+ + + + + + + + + + {tasks.map((t) => ( + void openTask(t.id)} + style={{ cursor: "pointer", background: selectedId === t.id ? "var(--wb-row-active, #eef)" : undefined }} + > + + + + + ))} + {!tasks.length ? ( + + + + ) : null} + +
网元状态最近采集
+ {t.ne_name || t.ne_ip} +
{t.vendor}
+
+ {t.status} + {t.collect_running ? " · 采集中" : ""} + {t.last_collect_ended_at || "-"}
+ 暂无任务 +
+
+ +
+ {!detail ? ( +

选择左侧任务查看详情

+ ) : ( + <> +

+ {detail.ne_name} · {detail.status} +

+
+ + + + +
+ {detail.last_error ?
{detail.last_error}
: null} + +

启用项(勾选表)

+ + + + + + + + + + + {(detail.items || []).map((it: any) => { + const prof = profiles.find((p) => p.profile_id === it.source_profile_id); + return ( + + + + + + + ); + })} + +
启用项命令类型
{it.enabled ? "✓" : ""} + {it.title || prof?.title || it.source_profile_id} + {prof?.description ? ( +
+ {prof.description} +
+ ) : null} +
+ {it.command_override || prof?.command_template || "-"} + {it.kind}
+ +

采集批次

+ + + + + + + + + + + {batches.map((b) => ( + + + + + + + ))} + {!batches.length ? ( + + + + ) : null} + +
时间状态行数操作
{b.started_at || "-"}{b.status}{b.row_count} + + +
+ 尚无批次 +
+ + {batchDetail ? ( +
+

批次详情 {batchDetail.id}

+

+ 命令 {batchDetail.command_count} · 行 {batchDetail.row_count} · {batchDetail.status} +

+ + + + + + + + + + + {(batchDetail.lldp_neighbors || []).slice(0, 200).map((n: any, i: number) => ( + + + + + + + ))} + +
本端口对端系统对端口管理IP
{n.local_if}{n.remote_sys}{n.remote_if}{n.remote_ip}
+
+ ) : null} + + )} +
+
+
+ ); +} diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 91ab328..7c66e57 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -1753,3 +1753,59 @@ export const deletePortTrafficBoard = (boardId: string) => apiDelete<{ ok: boolean; id: string }>( `/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[] }>(`/v1/biz-state/profiles${q ? `?${q}` : ""}`); +}; + +export const bizStateListTasks = () => + apiGet<{ items: Record[] }>("/v1/biz-state/tasks"); + +export const bizStateCreateTask = (body: Record) => + apiPost>("/v1/biz-state/tasks", body); + +export const bizStateGetTask = (taskId: string) => + apiGet>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}`); + +export const bizStatePatchTask = (taskId: string, body: Record) => + apiPatch>(`/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[] }>( + `/v1/biz-state/tasks/${encodeURIComponent(taskId)}/batches?limit=${limit}`, + ); + +export const bizStateGetBatch = (batchId: string) => + apiGet>(`/v1/biz-state/batches/${encodeURIComponent(batchId)}`); + +export const bizStateDownloadExport = async (batchId: string): Promise => { + 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); + } +};