diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index c70ad96..c0d9a16 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -20,6 +20,7 @@ from ..models import ( BizStateBatchCommand, BizStateEvent, BizStateLldpNeighbor, + BizStateMetricRow, BizStateTask, BizStateTaskItem, BizStateTaskItemBinding, @@ -148,6 +149,55 @@ def _persist_vrf_route_summary( return n +_GENERIC_METRICS = { + "isis_adjacency", + "interface_brief", + "arp", + "nd6_cache", + "bgp_peer", +} +_METRIC_CHUNK = 2000 + + +def _persist_metric_rows( + db, + *, + batch: BizStateBatch, + cmd_row: BizStateBatchCommand, + metric_id: str, + records: list[dict[str, Any]], +) -> int: + """Bulk-insert generic metric rows (JSON payload per row).""" + mid = str(metric_id or "").strip() + if not mid or not records: + return 0 + buf: list[dict[str, Any]] = [] + n = 0 + for i, rec in enumerate(records): + if not isinstance(rec, dict) or not rec: + continue + buf.append( + { + "id": uuid4().hex, + "batch_id": batch.id, + "batch_command_id": cmd_row.id, + "task_id": batch.task_id, + "ne_id": batch.ne_id, + "metric_id": mid, + "seq": i, + "data_json": dict(rec), + "collected_at": _utcnow(), + } + ) + n += 1 + if len(buf) >= _METRIC_CHUNK: + db.bulk_insert_mappings(BizStateMetricRow, buf) + buf.clear() + if buf: + db.bulk_insert_mappings(BizStateMetricRow, buf) + return n + + def _finish_task(task_id: str, *, error: str = "") -> None: db = SessionLocal() try: @@ -454,6 +504,14 @@ def _run_collect_session( n = _persist_vrf_route_summary( sdb, batch=batch_row, cmd_row=cmd_row, records=records ) + elif hit.profile.metric_id in _GENERIC_METRICS: + n = _persist_metric_rows( + sdb, + batch=batch_row, + cmd_row=cmd_row, + metric_id=hit.profile.metric_id, + records=records, + ) cmd_row.row_count = n cmd_row.parse_status = "ok" total_rows += n @@ -517,6 +575,7 @@ def _purge_old_batches(db, *, task_id: str, keep: int) -> None: bid = b.id db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == bid).delete() db.query(BizStateVrfRouteSummary).filter(BizStateVrfRouteSummary.batch_id == bid).delete() + db.query(BizStateMetricRow).filter(BizStateMetricRow.batch_id == bid).delete() db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == bid).delete() db.delete(b) if drop: diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py index d579850..1092023 100644 --- a/netx_api/biz_state/compare_service.py +++ b/netx_api/biz_state/compare_service.py @@ -212,6 +212,29 @@ def _default_vrf_sheet() -> dict[str, Any]: ) +def _default_sheet_for_metric(metric_id: str, *, compare_roles: tuple[str, ...] = ("state",)) -> dict[str, Any]: + fields = metric_field_map().get(metric_id) or [] + keys = [f.name for f in fields if f.is_key] + ifaces = [f.name for f in fields if f.is_interface] + compare = [f.name for f in fields if (not f.is_key) and f.role in compare_roles] + return _sheet_def( + metric_id=metric_id, + key_fields=keys, + iface_fields=ifaces, + compare_fields=compare, + ) + + +def _default_zte_status_sheets() -> list[dict[str, Any]]: + return [ + _default_sheet_for_metric("isis_adjacency", compare_roles=("state",)), + _default_sheet_for_metric("interface_brief", compare_roles=("state",)), + _default_sheet_for_metric("arp", compare_roles=("state",)), + _default_sheet_for_metric("nd6_cache", compare_roles=("state",)), + _default_sheet_for_metric("bgp_peer", compare_roles=("state",)), + ] + + def _normalize_sheet(raw: Any) -> dict[str, Any] | None: if not isinstance(raw, dict): return None @@ -420,10 +443,40 @@ def ensure_default_vrf_template(db: Session) -> BizCompareTemplate: return row +def ensure_default_zte_status_template(db: Session) -> BizCompareTemplate: + name = "ZTE status default" + row = db.query(BizCompareTemplate).filter(BizCompareTemplate.name == name).one_or_none() + sheets = _default_zte_status_sheets() + if row: + existing = template_metrics(row) + want = {s["metric_id"] for s in sheets} + have = {s["metric_id"] for s in existing} + if want - have: + _apply_sheets_to_row(row, sheets) + row.note = "Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP)" + row.updated_at = _utcnow() + db.commit() + db.refresh(row) + return row + row = BizCompareTemplate( + id=uuid4().hex, + name=name, + note="Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP)", + created_at=_utcnow(), + updated_at=_utcnow(), + ) + _apply_sheets_to_row(row, sheets) + db.add(row) + db.commit() + db.refresh(row) + return row + + def ensure_default_templates(db: Session) -> None: ensure_default_cutover_template(db) ensure_default_lldp_template(db) ensure_default_vrf_template(db) + ensure_default_zte_status_template(db) def list_templates(db: Session) -> list[dict[str, Any]]: @@ -638,6 +691,23 @@ def _load_metric_rows(db: Session, *, batch_id: str, metric_id: str) -> list[dic {"vrf": r.vrf, "source": r.source, "networks": r.networks} for r in rows ] + # Generic tabular metrics (ISIS / interface / ARP / ND6 / BGP …) + from ..models import BizStateMetricRow + + rows = ( + db.query(BizStateMetricRow) + .filter( + BizStateMetricRow.batch_id == batch_id, + BizStateMetricRow.metric_id == metric_id, + ) + .order_by(BizStateMetricRow.seq.asc(), BizStateMetricRow.id.asc()) + .all() + ) + if rows: + return [dict(r.data_json or {}) for r in rows] + # Known metric with zero rows is OK; unknown metric still errors + if metric_id in metric_field_map(): + return [] raise HTTPException(status_code=400, detail=f"unsupported_metric:{metric_id}") diff --git a/netx_api/biz_state/parsers/__init__.py b/netx_api/biz_state/parsers/__init__.py index 7178c3e..11521b3 100644 --- a/netx_api/biz_state/parsers/__init__.py +++ b/netx_api/biz_state/parsers/__init__.py @@ -6,6 +6,13 @@ from typing import Any, Callable from ...lldp_shared import NeighborHit, parse_neighbor_output from .vrf import normalize_vrf_list, normalize_vrf_route_summary +from .zte_status import ( + normalize_arp, + normalize_bgp_peer, + normalize_interface_brief, + normalize_isis_adjacency, + normalize_nd6_cache, +) NormalizeFn = Callable[..., list[dict[str, Any]]] @@ -43,6 +50,11 @@ _REGISTRY: dict[str, NormalizeFn] = { "lldp_neighbors": normalize_lldp_neighbors, "vrf_list": normalize_vrf_list, "vrf_route_summary": normalize_vrf_route_summary, + "isis_adjacency": normalize_isis_adjacency, + "interface_brief": normalize_interface_brief, + "arp": normalize_arp, + "nd6_cache": normalize_nd6_cache, + "bgp_peer": normalize_bgp_peer, } diff --git a/netx_api/biz_state/parsers/zte_status.py b/netx_api/biz_state/parsers/zte_status.py new file mode 100644 index 0000000..42350c6 --- /dev/null +++ b/netx_api/biz_state/parsers/zte_status.py @@ -0,0 +1,321 @@ +"""ZTE ZXROS status table parsers (ISIS / interface / ARP / ND6 / BGP).""" + +from __future__ import annotations + +import re +from typing import Any + +from ...lldp_shared import resolve_vendor_key +from ...ntc_parse import parse_cli, resolve_cli_platform, row_get + +_IFACE_RE = re.compile( + r"^(?P\S+)\s+(?P\S+)\s+(?P\S+)" + r"(?:\s+(?P\S+))?\s+(?Pup|down)\s+(?Pup|down)\s+(?Pup|down)" + r"(?:\s+(?P.*))?$", + re.I, +) + +_ISIS_ROW_RE = re.compile( + r"^(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+" + r"(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+" + r"(?P\S+)\s+(?P\S+)\s*$", + re.I, +) + +_ARP_ROW_RE = re.compile( + r"^(?P\d{1,3}(?:\.\d{1,3}){3})\s+(?P\S+)\s+(?P\S+)\s+" + r"(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+(?P\S+)\s*$", + re.I, +) + +_ND6_ROW_RE = re.compile( + r"^(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+(?P\S+)\s+" + r"(?P\S+)\s+(?P\S+)\s*$", + re.I, +) + +_BGP_PEER_RE = re.compile( + r"^(?P\d{1,3}(?:\.\d{1,3}){3})\s+(?P\d+)\s+(?P\d+)\s+" + r"(?P\d+)\s+(?P\d+)\s+(?P\S+)\s+(?P\S+)\s*$", + re.I, +) + +_PROCESS_RE = re.compile(r"(?i)^\s*Process\s+ID\s*:\s*(\d+)\s*$") +_HEADER_HINTS = ( + "interface", + "system id", + "address", + "neighbor", + "ip", + "hardware", + "link-address", +) + + +def _is_header(line: str) -> bool: + low = line.lower() + return any(h in low for h in _HEADER_HINTS) and ("---" in low or " " in line) + + +def normalize_isis_adjacency( + *, + raw_text: str, + vendor: str = "", + device_type: str = "", + command: str = "", + params: dict[str, str] | None = None, +) -> list[dict[str, Any]]: + _ = (vendor, device_type, command, params) + out: list[dict[str, Any]] = [] + process_id = "" + for raw in str(raw_text or "").splitlines(): + line = raw.strip() + if not line: + continue + pm = _PROCESS_RE.match(line) + if pm: + process_id = pm.group(1) + continue + if line.lower().startswith("interface") and "system" in line.lower(): + continue + m = _ISIS_ROW_RE.match(line) + if not m: + continue + out.append( + { + "process_id": process_id, + "interface": m.group("iface")[:128], + "system_id": m.group("sys")[:128], + "state": m.group("state")[:32], + "lev": m.group("lev")[:16], + "holds": m.group("holds")[:32], + "snpa": m.group("snpa")[:64], + "pri": m.group("pri")[:16], + "mt": m.group("mt")[:16], + "nsf": m.group("nsf")[:32], + "af": m.group("af")[:64], + } + ) + return out + + +def normalize_interface_brief( + *, + raw_text: str, + vendor: str = "", + device_type: str = "", + command: str = "", + params: dict[str, str] | None = None, +) -> list[dict[str, Any]]: + _ = params + platform = resolve_cli_platform( + vendor=vendor, + device_type=device_type, + vendor_key=resolve_vendor_key(vendor, device_type), + ) + cmd = str(command or "show interface brief").strip() or "show interface brief" + rows = parse_cli(platform=platform, command=cmd, text=raw_text) if platform else [] + out: list[dict[str, Any]] = [] + seen: set[str] = set() + for r in rows: + iface = row_get(r, "INTERFACE", "interface") + if not iface or iface.lower() == "interface": + continue + if iface in seen: + continue + seen.add(iface) + out.append( + { + "interface": iface[:128], + "attribute": row_get(r, "ATTRIBUTE", "attribute")[:64], + "mode": row_get(r, "MODE", "mode")[:64], + "bw": row_get(r, "BW", "bw")[:32], + "admin": row_get(r, "ADMIN", "admin")[:16], + "phy": row_get(r, "PHY", "phy")[:16], + "prot": row_get(r, "PROT", "prot")[:16], + "description": row_get(r, "DESCRIPTION", "description")[:256], + } + ) + if out: + return out + for raw in str(raw_text or "").splitlines(): + line = raw.strip() + if not line or line.lower().startswith("interface"): + continue + m = _IFACE_RE.match(line) + if not m: + continue + iface = m.group("iface") + if iface in seen: + continue + seen.add(iface) + out.append( + { + "interface": iface[:128], + "attribute": (m.group("attr") or "")[:64], + "mode": (m.group("mode") or "")[:64], + "bw": (m.group("bw") or "")[:32], + "admin": (m.group("admin") or "")[:16], + "phy": (m.group("phy") or "")[:16], + "prot": (m.group("prot") or "")[:16], + "description": (m.group("desc") or "").strip()[:256], + } + ) + return out + + +def normalize_arp( + *, + raw_text: str, + vendor: str = "", + device_type: str = "", + command: str = "", + params: dict[str, str] | None = None, +) -> list[dict[str, Any]]: + _ = (vendor, device_type, command, params) + out: list[dict[str, Any]] = [] + seen: set[tuple[str, str]] = set() + ip_re = re.compile(r"^\d{1,3}(?:\.\d{1,3}){3}$") + for raw in str(raw_text or "").splitlines(): + line = raw.strip() + if not line or line.startswith("---") or line.lower().startswith("arp protect"): + continue + if line.lower().startswith("the count"): + continue + if "hardware" in line.lower() and "address" in line.lower(): + continue + parts = line.split() + if len(parts) < 4 or not ip_re.match(parts[0]): + continue + ip = parts[0] + age = parts[1] + mac = parts[2] + iface = parts[3] + exter = parts[4] if len(parts) > 4 else "" + inter = parts[5] if len(parts) > 5 else "" + sub = parts[6] if len(parts) > 6 else "" + key = (ip, iface) + if key in seen: + continue + seen.add(key) + out.append( + { + "ip": ip[:64], + "age": age[:32], + "mac": mac[:64], + "interface": iface[:128], + "exter_vlan": exter[:32], + "inter_vlan": inter[:32], + "sub_interface": sub[:128], + } + ) + return out + + +def normalize_nd6_cache( + *, + raw_text: str, + vendor: str = "", + device_type: str = "", + command: str = "", + params: dict[str, str] | None = None, +) -> list[dict[str, Any]]: + _ = (vendor, device_type, command, params) + out: list[dict[str, Any]] = [] + seen: set[tuple[str, str]] = set() + for raw in str(raw_text or "").splitlines(): + line = raw.strip() + if not line: + continue + low = line.lower() + if low.startswith("s-static") or low.startswith("total cache") or low.startswith("only current"): + continue + if low.startswith("address") and "link" in low: + continue + m = _ND6_ROW_RE.match(line) + if not m: + continue + addr = m.group("addr") + iface = m.group("iface") + key = (addr, iface) + if key in seen: + continue + seen.add(key) + out.append( + { + "address": addr[:128], + "link_address": m.group("link")[:64], + "age": m.group("age")[:64], + "status": m.group("status")[:32], + "interface": iface[:128], + "type": m.group("type")[:32], + } + ) + return out + + +def _detect_bgp_afi(command: str, params: dict[str, str] | None) -> str: + if params and params.get("afi"): + return str(params.get("afi") or "").strip().lower() + low = str(command or "").lower() + if "vpnv6" in low: + return "vpnv6" + if "vpnv4" in low: + return "vpnv4" + if "ipv6" in low: + return "ipv6" + if "ipv4" in low: + return "ipv4" + return "unknown" + + +def normalize_bgp_peer( + *, + raw_text: str, + vendor: str = "", + device_type: str = "", + command: str = "", + params: dict[str, str] | None = None, +) -> list[dict[str, Any]]: + _ = (vendor, device_type) + afi = _detect_bgp_afi(command, params) + out: list[dict[str, Any]] = [] + seen: set[str] = set() + for raw in str(raw_text or "").splitlines(): + line = raw.strip() + if not line: + continue + low = line.lower() + if low.startswith("bgp router") or low.startswith("local as") or low.startswith("all "): + continue + if low.startswith("neighbor") and "msg" in low: + continue + m = _BGP_PEER_RE.match(line) + if not m: + continue + nei = m.group("nei") + if nei in seen: + continue + seen.add(nei) + state_raw = m.group("state") + if state_raw.isdigit(): + state = "Established" + pfx = state_raw + else: + state = state_raw + pfx = "" + out.append( + { + "afi": afi[:32], + "neighbor": nei[:64], + "ver": m.group("ver")[:8], + "as_num": m.group("asn")[:16], + "msg_rcvd": m.group("rx")[:32], + "msg_send": m.group("tx")[:32], + "up_down": m.group("up")[:32], + "state": state[:64], + "pfx_rcd": pfx[:32], + "state_or_pfx": state_raw[:64], + } + ) + return out diff --git a/netx_api/biz_state/profiles.py b/netx_api/biz_state/profiles.py index b197b70..330d370 100644 --- a/netx_api/biz_state/profiles.py +++ b/netx_api/biz_state/profiles.py @@ -227,13 +227,191 @@ def _vrf_profiles() -> list[ParseProfile]: return out +# --- ZTE ZXROS status tables (cutover monitoring + compare) --- + +_ISIS_FIELDS: list[FieldDef] = [ + FieldDef("process_id", length=32, indexed=True, is_key=True, display_name="Process ID"), + FieldDef("interface", length=128, indexed=True, is_key=True, is_interface=True, display_name="接口"), + FieldDef("system_id", length=128, indexed=True, is_key=True, display_name="System ID"), + FieldDef("state", length=32, role="state", display_name="状态"), + FieldDef("lev", length=16, role="state", display_name="Level"), + FieldDef("holds", length=32, role="meta", display_name="Holds"), + FieldDef("snpa", length=64, role="meta", display_name="SNPA"), + FieldDef("pri", length=16, role="meta", display_name="Pri"), + FieldDef("mt", length=16, role="meta", display_name="MT"), + FieldDef("nsf", length=32, role="meta", display_name="NSF"), + FieldDef("af", length=64, role="state", display_name="AF"), +] + +_IFACE_BRIEF_FIELDS: list[FieldDef] = [ + FieldDef("interface", length=128, indexed=True, is_key=True, is_interface=True, display_name="接口"), + FieldDef("attribute", length=64, role="meta", display_name="属性"), + FieldDef("mode", length=64, role="meta", display_name="模式"), + FieldDef("bw", length=32, role="meta", display_name="带宽"), + FieldDef("admin", length=16, role="state", display_name="Admin"), + FieldDef("phy", length=16, role="state", display_name="Phy"), + FieldDef("prot", length=16, role="state", display_name="Prot"), + FieldDef("description", length=256, role="meta", display_name="描述"), +] + +_ARP_FIELDS: list[FieldDef] = [ + FieldDef("ip", length=64, indexed=True, is_key=True, display_name="IP"), + FieldDef("interface", length=128, indexed=True, is_key=True, is_interface=True, display_name="接口"), + FieldDef("mac", length=64, role="state", display_name="MAC"), + FieldDef("age", length=32, role="meta", display_name="Age"), + FieldDef("exter_vlan", length=32, role="meta", display_name="Exter VLAN"), + FieldDef("inter_vlan", length=32, role="meta", display_name="Inter VLAN"), + FieldDef("sub_interface", length=128, role="meta", display_name="Sub-IF"), +] + +_ND6_FIELDS: list[FieldDef] = [ + FieldDef("address", length=128, indexed=True, is_key=True, display_name="IPv6"), + FieldDef("interface", length=128, indexed=True, is_key=True, is_interface=True, display_name="接口"), + FieldDef("link_address", length=64, role="state", display_name="Link-Address"), + FieldDef("status", length=32, role="state", display_name="Status"), + FieldDef("type", length=32, role="meta", display_name="Type"), + FieldDef("age", length=64, role="meta", display_name="Age"), +] + +_BGP_PEER_FIELDS: list[FieldDef] = [ + FieldDef("afi", length=32, indexed=True, is_key=True, display_name="AFI", from_command_param=True), + FieldDef("neighbor", length=64, indexed=True, is_key=True, display_name="Neighbor"), + FieldDef("as_num", length=16, role="state", display_name="AS"), + FieldDef("state", length=64, role="state", display_name="State"), + FieldDef("pfx_rcd", length=32, role="state", display_name="PfxRcd"), + FieldDef("state_or_pfx", length=64, role="meta", display_name="State/PfxRcd"), + FieldDef("ver", length=8, role="meta", display_name="Ver"), + FieldDef("msg_rcvd", length=32, role="meta", display_name="MsgRcvd"), + FieldDef("msg_send", length=32, role="meta", display_name="MsgSend"), + FieldDef("up_down", length=32, role="meta", display_name="Up/Down"), +] + + +def _zte_status_profiles() -> list[ParseProfile]: + """ZTE ZXROS status snapshots from lab show commands.""" + return [ + ParseProfile( + profile_id="zte.isis_adjacency", + vendor_key="zte", + metric_id="isis_adjacency", + parser_id="isis_adjacency", + title="ISIS Adjacency", + command_template="show isis adjacency | one-line", + match=r"(?i)^\s*show\s+isis\s+adjacency(?:\s*\|\s*one-line)?\s*$", + textfsm_command="show isis adjacency", + description="ISIS adjacency table (Process ID blocks).", + fields=list(_ISIS_FIELDS), + tags=["isis", "l3", "status"], + sort_order=300, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.interface_brief", + vendor_key="zte", + metric_id="interface_brief", + parser_id="interface_brief", + title="Interface Brief", + command_template="show interface brief", + match=r"(?i)^\s*show\s+interface\s+brief\s*$", + textfsm_command="show interface brief", + description="Interface admin/phy/prot status brief.", + fields=list(_IFACE_BRIEF_FIELDS), + tags=["interface", "l2", "status"], + sort_order=310, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.arp", + vendor_key="zte", + metric_id="arp", + parser_id="arp", + title="ARP Table", + command_template="show arp | one-line", + match=r"(?i)^\s*show\s+arp(?:\s*\|\s*one-line)?\s*$", + textfsm_command="show arp", + description="ARP entries (IP/MAC/interface).", + fields=list(_ARP_FIELDS), + tags=["arp", "l3", "status"], + sort_order=320, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.nd6_cache", + vendor_key="zte", + metric_id="nd6_cache", + parser_id="nd6_cache", + title="ND6 Cache", + command_template="show nd6 cache | one-line", + match=r"(?i)^\s*show\s+nd6\s+cache(?:\s*\|\s*one-line)?\s*$", + textfsm_command="show nd6 cache", + description="IPv6 neighbor discovery cache.", + fields=list(_ND6_FIELDS), + tags=["nd6", "ipv6", "status"], + sort_order=330, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.bgp_vpnv4_summary", + vendor_key="zte", + metric_id="bgp_peer", + parser_id="bgp_peer", + title="BGP VPNv4 Summary", + command_template="show bgp vpnv4 unicast summary", + match=r"(?i)^\s*show\s+bgp\s+vpnv4\s+unicast\s+summary\s*$", + textfsm_command="show bgp vpnv4 unicast summary", + description="BGP VPNv4 peer summary (afi=vpnv4).", + fields=list(_BGP_PEER_FIELDS), + tags=["bgp", "vpnv4", "status"], + sort_order=340, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.bgp_ipv4_summary", + vendor_key="zte", + metric_id="bgp_peer", + parser_id="bgp_peer", + title="BGP IPv4 Summary", + command_template="show bgp ipv4 unicast summary", + match=r"(?i)^\s*show\s+bgp\s+ipv4\s+unicast\s+summary\s*$", + textfsm_command="show bgp ipv4 unicast summary", + description="BGP IPv4 unicast peer summary (afi=ipv4).", + fields=list(_BGP_PEER_FIELDS), + tags=["bgp", "ipv4", "status"], + sort_order=350, + enabled=True, + kind="collect", + ), + ParseProfile( + profile_id="zte.bgp_vpnv6_summary", + vendor_key="zte", + metric_id="bgp_peer", + parser_id="bgp_peer", + title="BGP VPNv6 Summary", + command_template="show bgp vpnv6 unicast summary", + match=r"(?i)^\s*show\s+bgp\s+vpnv6\s+unicast\s+summary\s*$", + textfsm_command="show bgp vpnv6 unicast summary", + description="BGP VPNv6 peer summary (afi=vpnv6).", + fields=list(_BGP_PEER_FIELDS), + tags=["bgp", "vpnv6", "status"], + sort_order=360, + enabled=True, + kind="collect", + ), + ] + + _PROFILES: list[ParseProfile] | None = None def all_profiles() -> list[ParseProfile]: global _PROFILES if _PROFILES is None: - _PROFILES = _lldp_profiles() + _vrf_profiles() + _PROFILES = _lldp_profiles() + _vrf_profiles() + _zte_status_profiles() return list(_PROFILES) diff --git a/netx_api/biz_state/schema_ensure.py b/netx_api/biz_state/schema_ensure.py index 3dd8bfd..9ceade8 100644 --- a/netx_api/biz_state/schema_ensure.py +++ b/netx_api/biz_state/schema_ensure.py @@ -24,6 +24,8 @@ def apply_biz_state_schema(conn: Connection) -> None: "CREATE INDEX IF NOT EXISTS ix_biz_compare_job_status ON biz_compare_job (status)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_run_job_id ON biz_compare_run (job_id)", "CREATE INDEX IF NOT EXISTS ix_biz_state_vrf_route_batch_id ON biz_state_vrf_route_summary (batch_id)", + "CREATE INDEX IF NOT EXISTS ix_biz_state_metric_row_batch_id ON biz_state_metric_row (batch_id)", + "CREATE INDEX IF NOT EXISTS ix_biz_state_metric_row_batch_metric ON biz_state_metric_row (batch_id, metric_id)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_id ON biz_compare_diff (run_id)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_metric_kind ON biz_compare_diff (run_id, metric_id, kind)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_metric_seq ON biz_compare_diff (run_id, metric_id, seq)", diff --git a/netx_api/biz_state/service.py b/netx_api/biz_state/service.py index 1d1f1aa..faa0ac5 100644 --- a/netx_api/biz_state/service.py +++ b/netx_api/biz_state/service.py @@ -18,6 +18,7 @@ from ..models import ( BizStateCommandOverride, BizStateEvent, BizStateLldpNeighbor, + BizStateMetricRow, BizStateTask, BizStateTaskItem, BizStateTaskItemBinding, @@ -436,6 +437,21 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: .limit(5000) .all() ) + metric_rows = ( + db.query(BizStateMetricRow) + .filter(BizStateMetricRow.batch_id == batch_id) + .order_by( + BizStateMetricRow.metric_id.asc(), + BizStateMetricRow.seq.asc(), + BizStateMetricRow.id.asc(), + ) + .limit(20000) + .all() + ) + metrics_by_id: dict[str, list[dict[str, Any]]] = {} + for r in metric_rows: + mid = str(r.metric_id or "") + metrics_by_id.setdefault(mid, []).append(dict(r.data_json or {})) return { "id": b.id, "task_id": b.task_id, @@ -473,6 +489,7 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "vrf_route_summary": [ {"vrf": r.vrf, "source": r.source, "networks": r.networks} for r in vrf_rows ], + "metrics": metrics_by_id, } @@ -527,6 +544,20 @@ def export_batch_zip(db: Session, batch_id: str) -> bytes: ",".join([_csv(r["vrf"]), _csv(r["source"]), _csv(str(r["networks"]))]) ) zf.writestr("tables/vrf_route_summary.csv", "\n".join(vrf_csv) + "\n") + + for mid, rows in sorted((detail.get("metrics") or {}).items()): + if not rows: + continue + cols: list[str] = [] + for rec in rows: + for k in rec.keys(): + if k not in cols: + cols.append(str(k)) + lines = [",".join(_csv(c) for c in cols)] + for rec in rows: + lines.append(",".join(_csv(str(rec.get(c, "") or "")) for c in cols)) + safe = "".join(ch if ch.isalnum() or ch in "-_" else "_" for ch in mid)[:80] or "metric" + zf.writestr(f"tables/{safe}.csv", "\n".join(lines) + "\n") return buf.getvalue() diff --git a/netx_api/models/__init__.py b/netx_api/models/__init__.py index 5c14ee2..d9392cf 100644 --- a/netx_api/models/__init__.py +++ b/netx_api/models/__init__.py @@ -36,6 +36,7 @@ from .biz_state import ( BizStateCommandOverride, BizStateEvent, BizStateLldpNeighbor, + BizStateMetricRow, BizStateVrfRouteSummary, BizStateTask, BizStateTaskItem, @@ -136,6 +137,7 @@ __all__ = [ "BizStateBatchCommand", "BizStateLldpNeighbor", "BizStateVrfRouteSummary", + "BizStateMetricRow", "BizStateEvent", "BizStateCommandOverride", "BizCompareTemplate", diff --git a/netx_api/models/biz_state.py b/netx_api/models/biz_state.py index f2cad0f..a9de3a1 100644 --- a/netx_api/models/biz_state.py +++ b/netx_api/models/biz_state.py @@ -145,6 +145,26 @@ class BizStateVrfRouteSummary(Base): collected_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) +class BizStateMetricRow(Base): + """Generic structured rows for tabular status metrics (ISIS/ARP/BGP/…).""" + + __tablename__ = "biz_state_metric_row" + __table_args__ = ( + Index("ix_biz_state_metric_row_batch_metric", "batch_id", "metric_id"), + Index("ix_biz_state_metric_row_batch_metric_seq", "batch_id", "metric_id", "seq"), + ) + + 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) + metric_id: Mapped[str] = mapped_column(String(64), default="", index=True) + seq: Mapped[int] = mapped_column(Integer, default=0) + data_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + collected_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + + class BizStateEvent(Base): __tablename__ = "biz_state_event" diff --git a/tests/test_zte_status_parsers.py b/tests/test_zte_status_parsers.py new file mode 100644 index 0000000..a4dd1b9 --- /dev/null +++ b/tests/test_zte_status_parsers.py @@ -0,0 +1,127 @@ +"""Unit tests for ZTE ZXROS status parsers (samples from test/log).""" + +from __future__ import annotations + +import unittest +from pathlib import Path + +from netx_api.biz_state.parsers.zte_status import ( + normalize_arp, + normalize_bgp_peer, + normalize_interface_brief, + normalize_isis_adjacency, + normalize_nd6_cache, +) +from netx_api.biz_state.profiles import get_profile, metric_field_map, profiles_for_vendor + + +def _log_text() -> str: + p = Path(__file__).resolve().parents[2] / "test" / "log" + if not p.is_file(): + # Fallback: relative to monorepo root when tests run from netx/ + p = Path(__file__).resolve().parents[3] / "test" / "log" + return p.read_text(encoding="utf-8", errors="ignore") if p.is_file() else "" + + +def _section(blob: str, start_marker: str, end_markers: tuple[str, ...]) -> str: + i = blob.find(start_marker) + if i < 0: + return "" + rest = blob[i:] + cut = len(rest) + for em in end_markers: + j = rest.find(em, len(start_marker)) + if j >= 0: + cut = min(cut, j) + return rest[:cut] + + +class ZteStatusParserTests(unittest.TestCase): + @classmethod + def setUpClass(cls) -> None: + cls.log = _log_text() + + def test_profiles_registered(self) -> None: + zte = {p.profile_id for p in profiles_for_vendor("zte", kind="collect")} + self.assertIn("zte.isis_adjacency", zte) + self.assertIn("zte.interface_brief", zte) + self.assertIn("zte.arp", zte) + self.assertIn("zte.nd6_cache", zte) + self.assertIn("zte.bgp_vpnv4_summary", zte) + self.assertIn("zte.bgp_ipv4_summary", zte) + self.assertIn("zte.bgp_vpnv6_summary", zte) + for mid in ("isis_adjacency", "interface_brief", "arp", "nd6_cache", "bgp_peer"): + self.assertIn(mid, metric_field_map()) + self.assertIsNotNone(get_profile("zte.isis_adjacency")) + + def test_isis_adjacency(self) -> None: + text = _section( + self.log, + "show isis adjacency", + ("show interface brief", "show arp", "M6000-4SE-3#show"), + ) + rows = normalize_isis_adjacency(raw_text=text) + self.assertGreaterEqual(len(rows), 5) + procs = {r["process_id"] for r in rows} + self.assertIn("1", procs) + self.assertIn("20", procs) + up = [r for r in rows if r["state"].upper() == "UP"] + self.assertEqual(len(up), len(rows)) + + def test_interface_brief(self) -> None: + text = _section( + self.log, + "show interface brief", + ("show arp", "show nd6", "M6000-4SE-3#show arp"), + ) + rows = normalize_interface_brief( + raw_text=text, vendor="zte", device_type="zte_zxros", command="show interface brief" + ) + self.assertGreaterEqual(len(rows), 10) + by_if = {r["interface"]: r for r in rows} + self.assertIn("cgei-0/1/0/1", by_if) + self.assertEqual(by_if["cgei-0/1/0/1"]["admin"].lower(), "up") + self.assertEqual(by_if["cgei-0/1/0/3"]["admin"].lower(), "down") + + def test_arp(self) -> None: + text = _section(self.log, "show arp", ("show nd6", "PAG3_", "M6000-4SE-3#show nd6")) + rows = normalize_arp(raw_text=text) + self.assertGreaterEqual(len(rows), 10) + ips = {r["ip"] for r in rows} + self.assertIn("192.166.1.65", ips) + self.assertIn("10.229.234.1", ips) + + def test_nd6(self) -> None: + text = _section(self.log, "show nd6 cache", ("PAG3_", "show bgp")) + rows = normalize_nd6_cache(raw_text=text) + self.assertGreaterEqual(len(rows), 5) + addrs = {r["address"] for r in rows} + self.assertTrue(any("fe80::" in a for a in addrs)) + + def test_bgp_peers(self) -> None: + v4 = _section(self.log, "show bgp vpnv4 unicast summary", ("show bgp ipv4",)) + rows = normalize_bgp_peer(raw_text=v4, command="show bgp vpnv4 unicast summary") + self.assertGreaterEqual(len(rows), 8) + self.assertTrue(all(r["afi"] == "vpnv4" for r in rows)) + est = [r for r in rows if r["state"] == "Established"] + conn = [r for r in rows if r["state"] == "Connect"] + self.assertGreaterEqual(len(est), 3) + self.assertGreaterEqual(len(conn), 3) + + ipv4 = _section(self.log, "show bgp ipv4 unicast summary", ("show bgp vpnv6",)) + rows2 = normalize_bgp_peer(raw_text=ipv4, command="show bgp ipv4 unicast summary") + self.assertGreaterEqual(len(rows2), 8) + self.assertTrue(all(r["afi"] == "ipv4" for r in rows2)) + + v6 = _section(self.log, "show bgp vpnv6 unicast summary", ("END",)) + # file ends after vpnv6 — take remainder + if not v6.strip(): + i = self.log.find("show bgp vpnv6 unicast summary") + v6 = self.log[i:] if i >= 0 else "" + rows3 = normalize_bgp_peer(raw_text=v6, command="show bgp vpnv6 unicast summary") + self.assertGreaterEqual(len(rows3), 5) + self.assertTrue(all(r["afi"] == "vpnv6" for r in rows3)) + + +if __name__ == "__main__": + unittest.main()