Add ZTE status metrics (ISIS/IF/ARP/ND6/BGP) for monitor and compare.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-17 21:11:40 +08:00
parent fe34319d65
commit 2f0c9eb1da
10 changed files with 823 additions and 1 deletions

View file

@ -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:

View file

@ -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}")

View file

@ -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,
}

View file

@ -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<iface>\S+)\s+(?P<attr>\S+)\s+(?P<mode>\S+)"
r"(?:\s+(?P<bw>\S+))?\s+(?P<admin>up|down)\s+(?P<phy>up|down)\s+(?P<prot>up|down)"
r"(?:\s+(?P<desc>.*))?$",
re.I,
)
_ISIS_ROW_RE = re.compile(
r"^(?P<iface>\S+)\s+(?P<sys>\S+)\s+(?P<state>\S+)\s+(?P<lev>\S+)\s+"
r"(?P<holds>\S+)\s+(?P<snpa>\S+)\s+(?P<pri>\S+)\s+(?P<mt>\S+)\s+"
r"(?P<nsf>\S+)\s+(?P<af>\S+)\s*$",
re.I,
)
_ARP_ROW_RE = re.compile(
r"^(?P<ip>\d{1,3}(?:\.\d{1,3}){3})\s+(?P<age>\S+)\s+(?P<mac>\S+)\s+"
r"(?P<iface>\S+)\s+(?P<ext_vlan>\S+)\s+(?P<int_vlan>\S+)\s+(?P<sub>\S+)\s*$",
re.I,
)
_ND6_ROW_RE = re.compile(
r"^(?P<addr>\S+)\s+(?P<link>\S+)\s+(?P<age>\S+)\s+(?P<status>\S+)\s+"
r"(?P<iface>\S+)\s+(?P<type>\S+)\s*$",
re.I,
)
_BGP_PEER_RE = re.compile(
r"^(?P<nei>\d{1,3}(?:\.\d{1,3}){3})\s+(?P<ver>\d+)\s+(?P<asn>\d+)\s+"
r"(?P<rx>\d+)\s+(?P<tx>\d+)\s+(?P<up>\S+)\s+(?P<state>\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

View file

@ -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)

View file

@ -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)",

View file

@ -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()

View file

@ -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",

View file

@ -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"

View file

@ -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()