diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index 5b03569..c35d487 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -3,6 +3,7 @@ from __future__ import annotations import logging +import re import threading from concurrent.futures import ThreadPoolExecutor from datetime import datetime @@ -222,7 +223,7 @@ def _run_primary_parse_job(job: Any) -> tuple[bool, bool]: """Parse primary+aux raws, write spool, submit persist. Returns (any_ok, any_fail).""" from .parse_pool import AuxRawCapture, PrimaryParseJob from .persist_pool import get_persist_pool - from .spool import SpooledCommand, write_meta, write_records + from .spool import SpooledCommand, count_text_lines, write_meta, write_records if not isinstance(job, PrimaryParseJob): return False, True @@ -233,6 +234,10 @@ def _run_primary_parse_job(job: Any) -> tuple[bool, bool]: any_ok = False any_fail = False + def _declared_total(raw: str) -> int: + m = re.search(r"(?i)total\s+number\s+of\s+routes\s*:\s*(\d+)", raw or "") + return int(m.group(1)) if m else 0 + def _flush_item(item: SpooledCommand, *, records: list[dict[str, Any]] | None = None) -> None: if records is not None and item.persist_kind: item.records_rel_path = write_records(batch_id, item.id, records) @@ -301,6 +306,7 @@ def _run_primary_parse_job(job: Any) -> tuple[bool, bool]: raw_command=cap.command[:512], params_json={}, raw_rel_path=cap.raw_rel_path, + raw_line_count=count_text_lines(cap.raw), ) if cap.cache_hit and entry.ok and entry.records: aux_sp.parse_status = "aux_cached" @@ -364,6 +370,9 @@ def _run_primary_parse_job(job: Any) -> tuple[bool, bool]: raw_command=job.concrete[:512], params_json=dict(job.merged_params or {}), raw_rel_path=job.raw_rel_path, + raw_line_count=int(getattr(job, "raw_line_count", 0) or 0) + or count_text_lines(job.raw_text), + declared_total=_declared_total(job.raw_text), ) try: records, fsm_tables, rule_keys = run_primary_with_bundle( @@ -389,6 +398,10 @@ def _run_primary_parse_job(job: Any) -> tuple[bool, bool]: hints.append( "enrich=" + ",".join(getattr(j, "from_aux", "") for j in job.enrich_joins) ) + declared = int(primary.declared_total or 0) + nrec = len(records or []) + if declared > 0: + hints.append(f"declared={declared};parsed={nrec}") if hints: primary.message = ";".join(hints)[:1020] primary.parse_status = "ok" @@ -451,8 +464,28 @@ def _flush_spooled_commands( max_raw = raw_max_bytes() for item in items: raw = "" + truncated = False + line_count = int(getattr(item, "raw_line_count", 0) or 0) if item.raw_rel_path: + from .spool import count_file_lines + + if line_count <= 0: + try: + line_count = count_file_lines(item.raw_rel_path) + except Exception: + line_count = 0 raw = read_raw_text(item.raw_rel_path, max_bytes=max_raw) + if max_raw > 0 and "[truncated" in raw: + truncated = True + # Prefer full-file line count; fall back to stored text. + if line_count <= 0 and raw: + from .spool import count_text_lines + + line_count = count_text_lines(raw) + msg = str(item.message or "").strip() + if truncated: + note = f"raw_truncated@{max_raw}B" + msg = f"{msg}; {note}" if msg else note cmd_row = BizStateBatchCommand( id=item.id, batch_id=batch_id, @@ -463,9 +496,11 @@ def _flush_spooled_commands( raw_command=str(item.raw_command or "")[:512], params_json=dict(item.params_json or {}), parse_status=item.parse_status, - message=str(item.message or "")[:1020], + message=msg[:1020], raw_text=raw, row_count=int(item.row_count or 0), + raw_line_count=int(line_count or 0), + declared_total=int(getattr(item, "declared_total", 0) or 0), created_at=_utcnow(), ) db.add(cmd_row) @@ -925,8 +960,12 @@ def _run_collect_lane( continue raw_rel = "" + raw_lines = 0 try: + from .spool import count_text_lines, write_raw_text + raw_rel = write_raw_text(batch_id, cmd_id, raw_text) + raw_lines = count_text_lines(raw_text) except Exception: _log.exception("biz_state spool raw failed cmd=%s", cmd_id) @@ -943,6 +982,7 @@ def _run_collect_lane( parse_status="skipped_custom", message="custom_raw", raw_rel_path=raw_rel, + raw_line_count=raw_lines, ) ) continue @@ -961,6 +1001,7 @@ def _run_collect_lane( parse_status="unmatched", message="no profile matched concrete command", raw_rel_path=raw_rel, + raw_line_count=raw_lines, ) ) continue @@ -981,6 +1022,7 @@ def _run_collect_lane( parse_status="failed", message=f"unknown parser {hit.profile.parser_id}", raw_rel_path=raw_rel, + raw_line_count=raw_lines, ) ) continue @@ -1070,6 +1112,7 @@ def _run_collect_lane( merged_params=dict(merged or {}), raw_text=raw_text, raw_rel_path=raw_rel, + raw_line_count=raw_lines, textfsm_command=hit.profile.textfsm_command or concrete, vendor=vendor_eff, device_type=device_type_eff, @@ -1105,6 +1148,11 @@ def _run_collect_lane( label, ) any_fail = True + _emit_task_event( + task_id=task_id, + message=f"{label}: parse_barrier_timeout", + level="error", + ) with parse_stats_lock: if parse_stats["ok"]: any_ok = True @@ -1116,6 +1164,12 @@ def _run_collect_lane( batch_id, label, ) + any_fail = True + _emit_task_event( + task_id=task_id, + message=f"{label}: persist_barrier_timeout", + level="error", + ) # Do NOT read batch.row_count here — dual light+heavy lanes would each # see the cumulative DB total and _absorb would double-count. return total_rows, cmd_count, any_fail, any_ok @@ -1265,6 +1319,70 @@ def _batch_has_progress(batch_id: str) -> bool: return bool(_run_db_with_reconnect(_read, label="biz_state_batch_progress")) +def _is_fail_parse_status(status: str | None) -> bool: + st = str(status or "").strip().lower() + return st in ("failed", "error", "fail", "aux_failed") or st.endswith("_failed") + + +def _is_skip_parse_status(status: str | None) -> bool: + st = str(status or "").strip().lower() + return st.startswith("skipped") + + +def _batch_issue_summary( + db, + batch_id: str, + lane_errors: list[str] | None = None, +) -> str: + """Human-readable reason for partial/failed batches (never leave message empty).""" + parts: list[str] = [] + for err in lane_errors or []: + e = str(err or "").strip() + if e and e not in parts: + parts.append(e) + + cmds = ( + db.query(BizStateBatchCommand) + .filter(BizStateBatchCommand.batch_id == batch_id) + .order_by(BizStateBatchCommand.created_at.asc()) + .all() + ) + n_fail = 0 + n_skip = 0 + n_aux_fail = 0 + samples: list[str] = [] + for c in cmds: + st = str(c.parse_status or "").strip().lower() + msg = str(c.message or "").strip() + if st == "aux_failed" or (st.startswith("aux") and "fail" in st): + n_aux_fail += 1 + if msg and len(samples) < 3: + samples.append(f"aux:{c.raw_command}: {msg}"[:160]) + elif _is_fail_parse_status(st): + n_fail += 1 + if msg and len(samples) < 3: + samples.append(f"{c.raw_command}: {msg}"[:160]) + elif _is_skip_parse_status(st): + n_skip += 1 + if msg and len(samples) < 3: + samples.append(f"skip:{c.profile_id or c.raw_command}: {msg}"[:160]) + + if n_fail: + parts.append(f"{n_fail} command(s) failed") + if n_aux_fail: + parts.append(f"{n_aux_fail} aux command(s) failed") + if n_skip: + parts.append(f"{n_skip} item(s) skipped (e.g. missing bindings)") + for s in samples: + if s not in parts: + parts.append(s) + + text = "; ".join(parts).strip() + if not text: + text = "partial success (some steps failed or were skipped)" + return text[:1020] + + def _finalize_batch_status( *, batch_id: str, @@ -1295,16 +1413,27 @@ def _finalize_batch_status( batch.message = STOP_USER_MESSAGE elif any_fail and any_ok: batch.status = "partial" - if lane_errors: - batch.message = "; ".join(lane_errors)[:1020] + batch.message = _batch_issue_summary(db, batch_id, lane_errors) elif any_fail and not any_ok: batch.status = "failed" - batch.message = ( + batch.message = _batch_issue_summary(db, batch_id, lane_errors) or ( "; ".join(lane_errors)[:1020] if lane_errors else "all commands failed" ) else: - batch.status = "success" - batch.message = "" + # Still surface skipped-only rows as partial when some cmds ran ok. + statuses = [ + str(c.parse_status or "") + for c in db.query(BizStateBatchCommand) + .filter(BizStateBatchCommand.batch_id == batch_id) + .all() + ] + n_skip = sum(1 for st in statuses if _is_skip_parse_status(st)) + if n_skip > 0 and any_ok: + batch.status = "partial" + batch.message = _batch_issue_summary(db, batch_id, lane_errors) + else: + batch.status = "success" + batch.message = "" status = str(batch.status or "") db.commit() # Only full success triggers auto compare; never block the collect thread. @@ -1448,7 +1577,26 @@ def _run_collect_session( command_override=item.command_override, ) except ValueError as exc: - _append_event(db, task_id=task_id, message=str(exc), level="error") + msg = str(exc) + _append_event(db, task_id=task_id, message=msg, level="error") + # Persist a visible skip row so UI / partial status can explain it. + db.add( + BizStateBatchCommand( + id=uuid4().hex, + batch_id=batch_id, + task_item_id=item.id, + profile_id=profile.profile_id, + parser_id=str(profile.parser_id or ""), + metric_id=str(profile.metric_id or ""), + raw_command=str(profile.command_template or "")[:512], + params_json={}, + parse_status="skipped", + message=msg[:1020], + row_count=0, + raw_text="", + ) + ) + batch.command_count = int(batch.command_count or 0) + 1 continue for concrete, params in pairs: if concrete == EXPAND_ALL_COMMAND: diff --git a/netx_api/biz_state/parse_pool.py b/netx_api/biz_state/parse_pool.py index 8bd06ed..58893a5 100644 --- a/netx_api/biz_state/parse_pool.py +++ b/netx_api/biz_state/parse_pool.py @@ -52,6 +52,7 @@ class PrimaryParseJob: merged_params: dict[str, Any] raw_text: str raw_rel_path: str + raw_line_count: int = 0 textfsm_command: str vendor: str device_type: str diff --git a/netx_api/biz_state/parsers/zte/bgp_route.py b/netx_api/biz_state/parsers/zte/bgp_route.py index f47026c..f753266 100644 --- a/netx_api/biz_state/parsers/zte/bgp_route.py +++ b/netx_api/biz_state/parsers/zte/bgp_route.py @@ -314,14 +314,22 @@ def _hand_parse( if _skip_noise_line(line): continue - # Continuation: indented next-hop after network-only line + # Continuation: indented next-hop after network-only line. + # RR vpnv6 often puts NH + LocPrf/Path on the same indented line + # ("24.11.0.8 100 0 ?") — take first token + # as NH or ECMP legs collapse under empty next_hop dedupe (~half count). if pending_net and not pending_nh and line[:1].isspace(): tok = line.strip() - if _looks_like_ip_or_prefix(tok) and "/" not in tok: - pending_nh = tok + parts = tok.split() + first = parts[0] if parts else "" + if _looks_like_ip_or_prefix(first) and "/" not in first: + pending_nh = first + rest = " ".join(parts[1:]) + if rest: + _flush_pending(rest=rest, path_continuation=True) continue # Metrics/path without explicit next-hop (rare) - if tok and not _looks_like_prefix(tok.split()[0] if tok.split() else ""): + if tok and not _looks_like_prefix(first): _flush_pending(rest=tok, path_continuation=True) continue @@ -417,9 +425,9 @@ def normalize_bgp_route( ) rows = prefer_fsm(tables, RULE_KEYS, _map, _hand, raw_text=raw_text) - # If device declared a large table but FSM/hand returned a tiny subset, prefer hand. + # Prefer hand when FSM under-parses vs device-declared total (common on vpnv6 wraps). declared = _declared_total(raw_text) - if declared is not None and declared > 0 and len(rows) < max(1, declared // 2): + if declared is not None and declared > 0 and len(rows) < int(declared * 0.9): hand_rows = _hand(raw_text=raw_text) if len(hand_rows) > len(rows): return hand_rows diff --git a/netx_api/biz_state/schema_ensure.py b/netx_api/biz_state/schema_ensure.py index 990e6e2..6b1013f 100644 --- a/netx_api/biz_state/schema_ensure.py +++ b/netx_api/biz_state/schema_ensure.py @@ -79,6 +79,8 @@ def apply_biz_state_schema(conn: Connection) -> None: "ALTER TABLE biz_migration_project ADD COLUMN IF NOT EXISTS new_hf_bindings_json JSON DEFAULT '[]'", "ALTER TABLE biz_state_task ADD COLUMN IF NOT EXISTS purpose VARCHAR(32) DEFAULT ''", "CREATE INDEX IF NOT EXISTS ix_biz_state_task_purpose ON biz_state_task (purpose)", + "ALTER TABLE biz_state_batch_command ADD COLUMN IF NOT EXISTS raw_line_count INTEGER DEFAULT 0", + "ALTER TABLE biz_state_batch_command ADD COLUMN IF NOT EXISTS declared_total INTEGER DEFAULT 0", "ALTER TABLE biz_migration_red_ticket ADD COLUMN IF NOT EXISTS match_key_str VARCHAR(256) DEFAULT ''", "CREATE INDEX IF NOT EXISTS ix_biz_migration_red_match ON biz_migration_red_ticket (project_id, metric_id, match_key_str)", ): diff --git a/netx_api/biz_state/service.py b/netx_api/biz_state/service.py index cf8fc35..a41192c 100644 --- a/netx_api/biz_state/service.py +++ b/netx_api/biz_state/service.py @@ -601,6 +601,18 @@ def _raw_line_count(raw: str | None) -> int: return s.count("\n") + (0 if s.endswith("\n") else 1) +def _cmd_raw_line_count(cmd: Any) -> int: + """Prefer full-file line count persisted before DB raw_text truncate.""" + stored = int(getattr(cmd, "raw_line_count", 0) or 0) + if stored > 0: + return stored + return _raw_line_count(getattr(cmd, "raw_text", None)) + + +def _cmd_declared_total(cmd: Any) -> int: + return int(getattr(cmd, "declared_total", 0) or 0) + + def get_batch(db: Session, batch_id: str) -> dict[str, Any]: """Batch workbook summary: meta + commands + sheet catalog (no metric row payload).""" b = db.get(BizStateBatch, batch_id) @@ -666,7 +678,8 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "params": c.params_json or {}, "parse_status": c.parse_status, "row_count": c.row_count, - "raw_line_count": _raw_line_count(raw), + "raw_line_count": _cmd_raw_line_count(c), + "declared_total": _cmd_declared_total(c), "message": c.message, "has_raw": bool(str(raw).strip()), "is_aux": is_aux, @@ -674,8 +687,12 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: cmd_n = normalize_command(str(c.raw_command or "")) if not is_aux and cmd_n: primary_cmds_by_cli.setdefault(cmd_n, info) - # Commands sheet: hide aux when the same CLI already has a primary row + # Commands sheet: hide successful aux when the same CLI already has a primary row; + # keep failed/skipped aux visible so partial reasons are not hidden. if is_aux and cmd_n and cmd_n in primary_cmds_by_cli: + st_l = status + if st_l in ("aux_failed",) or "fail" in st_l or st_l.startswith("skipped"): + cmd_payload.append(info) continue if is_aux and cmd_n: # aux may appear before primary in list — defer; second pass below @@ -699,6 +716,7 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "parse_status": c.parse_status, "row_count": c.row_count, "raw_line_count": info["raw_line_count"], + "declared_total": info["declared_total"], "message": c.message, "has_raw": info["has_raw"], "profile_id": c.profile_id, @@ -724,7 +742,8 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "params": c.params_json or {}, "parse_status": c.parse_status, "row_count": c.row_count, - "raw_line_count": _raw_line_count(c.raw_text), + "raw_line_count": _cmd_raw_line_count(c), + "declared_total": _cmd_declared_total(c), "message": c.message, "has_raw": bool(str(c.raw_text or "").strip()), "is_aux": True, @@ -756,6 +775,23 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: sheet_count = len(sheets) sheets_with_data = sum(1 for s in sheets if int(s.get("row_count") or 0) > 0) + # Full-batch command stats (includes hidden successful aux). + stats = {"total": 0, "ok": 0, "failed": 0, "aux_failed": 0, "skipped": 0, "other": 0} + for c in cmds: + stats["total"] += 1 + st = str(c.parse_status or "").strip().lower() + if st in ("ok", "success", "aux", "aux_ok", "aux_cached"): + stats["ok"] += 1 + elif st == "aux_failed" or (st.startswith("aux") and "fail" in st): + stats["aux_failed"] += 1 + stats["failed"] += 1 + elif st in ("failed", "error", "fail") or st.endswith("_failed"): + stats["failed"] += 1 + elif st.startswith("skipped"): + stats["skipped"] += 1 + else: + stats["other"] += 1 + return { "id": b.id, "task_id": b.task_id, @@ -777,6 +813,7 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "sheets_with_data": sheets_with_data, "sheet_count": sheet_count, "commands": cmd_payload, + "command_stats": stats, "sheets": sheets, } @@ -935,7 +972,8 @@ def get_batch_command(db: Session, batch_id: str, command_id: str) -> dict[str, "params": c.params_json or {}, "parse_status": c.parse_status, "row_count": c.row_count, - "raw_line_count": _raw_line_count(raw), + "raw_line_count": _cmd_raw_line_count(c), + "declared_total": _cmd_declared_total(c), "message": c.message, "raw_text": raw, "collected_at": c.created_at.isoformat() + "Z" if c.created_at else None, diff --git a/netx_api/biz_state/spool.py b/netx_api/biz_state/spool.py index 999117f..b0d00e9 100644 --- a/netx_api/biz_state/spool.py +++ b/netx_api/biz_state/spool.py @@ -51,11 +51,33 @@ def _cmd_paths(batch_id: str, cmd_id: str) -> tuple[Path, Path, Path]: def write_raw_text(batch_id: str, cmd_id: str, text: str) -> str: """Write CLI output; return path relative to spool root (posix).""" raw_path, _, _ = _cmd_paths(batch_id, cmd_id) - raw_path.write_bytes(str(text or "").encode("utf-8", errors="replace")) + data = str(text or "").encode("utf-8", errors="replace") + raw_path.write_bytes(data) rel = raw_path.resolve().relative_to(spool_root()) return str(rel).replace("\\", "/") +def count_file_lines(rel_path: str) -> int: + """Count lines in a spool raw file (full file, not DB-truncated).""" + if not rel_path: + return 0 + path = (spool_root() / str(rel_path)).resolve() + if not str(path).startswith(str(spool_root())) or not path.is_file(): + return 0 + n = 0 + with path.open("rb") as fh: + for _ in fh: + n += 1 + return n + + +def count_text_lines(text: str | None) -> int: + s = text or "" + if not s: + return 0 + return s.count("\n") + (0 if s.endswith("\n") else 1) + + def write_records(batch_id: str, cmd_id: str, records: list[dict[str, Any]]) -> str: """Write parsed records as JSONL; return relative path.""" _, _, rec_path = _cmd_paths(batch_id, cmd_id) @@ -130,6 +152,12 @@ class SpooledCommand: raw_rel_path: str = "" records_rel_path: str = "" row_count: int = 0 + # Full CLI line count (before DB raw_text truncate). + raw_line_count: int = 0 + # Device-declared total when present (e.g. BGP "Total number of routes"). + declared_total: int = 0 + # True when raw_text stored in DB was truncated by raw_max_bytes. + raw_truncated: bool = False # "" | "metric" | "lldp" persist_kind: str = "" @@ -148,6 +176,9 @@ class SpooledCommand: "raw_rel_path": self.raw_rel_path, "records_rel_path": self.records_rel_path, "row_count": self.row_count, + "raw_line_count": self.raw_line_count, + "declared_total": self.declared_total, + "raw_truncated": self.raw_truncated, "persist_kind": self.persist_kind, } diff --git a/netx_api/cli_templates/zte/zte_zxros_show_bgp_neighbor_routes.textfsm b/netx_api/cli_templates/zte/zte_zxros_show_bgp_neighbor_routes.textfsm index dc2dd7c..1b6f2a7 100644 --- a/netx_api/cli_templates/zte/zte_zxros_show_bgp_neighbor_routes.textfsm +++ b/netx_api/cli_templates/zte/zte_zxros_show_bgp_neighbor_routes.textfsm @@ -33,8 +33,11 @@ Routes # One-line IPv4 / RR "* i prefix nh …" ^\s*${STATUS}\s+${NETWORK}\s+${NEXT_HOP}\s+${METRIC}\s+${LOC_PRF}\s+${TAG}\s+${PATH}\s*$$ -> Record ^\s*${STATUS}\s+${NETWORK}\s+${NEXT_HOP}\s+${PATH}\s*$$ -> Record - # IPv6 wrap: network / next-hop / metrics+path on separate lines + # IPv6 wrap: network alone, then next-hop (+ optional metric/path) lines ^\s*${STATUS}\s+${NETWORK}\s*$$ + ^\s+${NEXT_HOP}\s+${METRIC}\s+${LOC_PRF}\s+${TAG}\s+${PATH}\s*$$ -> Record + ^\s+${NEXT_HOP}\s+${METRIC}\s+${LOC_PRF}\s+${PATH}\s*$$ -> Record + ^\s+${NEXT_HOP}\s+${PATH}\s*$$ -> Record ^\s+${NEXT_HOP}\s*$$ ^\s+${PATH}\s*$$ -> Record ^\s*$$ diff --git a/netx_api/models/biz_state.py b/netx_api/models/biz_state.py index 480e88a..bd8c311 100644 --- a/netx_api/models/biz_state.py +++ b/netx_api/models/biz_state.py @@ -118,6 +118,10 @@ class BizStateBatchCommand(Base): 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) + # Full CLI lines counted before raw_text truncate into Postgres. + raw_line_count: Mapped[int] = mapped_column(Integer, default=0) + # Optional device-declared total (BGP Total number of routes, etc.). + declared_total: 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) diff --git a/tests/test_batch_workbook_api.py b/tests/test_batch_workbook_api.py index e156545..2cb354f 100644 --- a/tests/test_batch_workbook_api.py +++ b/tests/test_batch_workbook_api.py @@ -123,6 +123,64 @@ class BatchWorkbookApiTests(unittest.TestCase): self.assertEqual(out["message"], "stopped") self.assertEqual(out["commands"][0]["raw_line_count"], 3) self.assertEqual(out["commands"][0]["row_count"], 3) + self.assertEqual(out["commands"][0]["declared_total"], 0) + + def test_get_batch_prefers_stored_raw_line_count(self) -> None: + """DB raw_text may be truncated; API must use persisted full-file line count.""" + batch = BizStateBatch( + id="b1", + task_id="t1", + status="ok", + command_count=1, + row_count=100, + ) + cmd = BizStateBatchCommand( + id="c1", + batch_id="b1", + profile_id="zte.bgp_vpnv4_neighbor_in", + parser_id="bgp_route", + metric_id="bgp_route", + raw_command="show bgp vpnv4 unicast neighbor in 1.1.1.1", + parse_status="ok", + row_count=100, + raw_text="a\nb\nc", # truncated stub (3 lines) + raw_line_count=119303, + declared_total=1669101, + message="declared=1669101;parsed=100", + ) + db = MagicMock() + db.get.side_effect = lambda model, pk: batch if pk == "b1" else None + cmd_q = MagicMock() + cmd_q.filter.return_value.order_by.return_value.all.return_value = [cmd] + metric_count_q = MagicMock() + metric_count_q.filter.return_value.group_by.return_value.all.return_value = [ + ("bgp_route", 100) + ] + lldp_count_q = MagicMock() + lldp_count_q.filter.return_value.scalar.return_value = 0 + + def query(*_args, **_kwargs): + n = query.n + query.n += 1 + if n == 0: + return cmd_q + if n == 1: + return metric_count_q + return lldp_count_q + + query.n = 0 + db.query.side_effect = query + + with patch( + "netx_api.biz_state.service.batch_protect_info", + return_value={"protected": False, "reasons": []}, + ): + out = get_batch(db, "b1") + + self.assertEqual(out["commands"][0]["raw_line_count"], 119303) + self.assertEqual(out["commands"][0]["declared_total"], 1669101) + self.assertEqual(out["sheets"][0]["commands"][0]["raw_line_count"], 119303) + self.assertEqual(out["sheets"][0]["commands"][0]["declared_total"], 1669101) def test_get_batch_command_and_raw_download(self) -> None: from netx_api.biz_state.service import get_batch_command diff --git a/tests/test_biz_state_collect_finalize.py b/tests/test_biz_state_collect_finalize.py index 7de469d..9012ab9 100644 --- a/tests/test_biz_state_collect_finalize.py +++ b/tests/test_biz_state_collect_finalize.py @@ -79,6 +79,100 @@ class BizStateCollectFinalizeTests(unittest.TestCase): self.assertIn("biz_state_heavy_timeout", batch.message or "") self.assertIsNotNone(batch.ended_at) + def test_finalize_partial_message_from_failed_cmds_when_lane_errors_empty(self) -> None: + from netx_api.models import BizStateBatchCommand + + self.db.add( + BizStateBatchCommand( + id="c-ok", + batch_id="b-finalize", + profile_id="zte.arp", + parser_id="arp", + metric_id="arp", + raw_command="show arp", + parse_status="ok", + row_count=10, + ) + ) + self.db.add( + BizStateBatchCommand( + id="c-aux-fail", + batch_id="b-finalize", + profile_id="zte.if_intf", + parser_id="if_intf", + metric_id="if_intf", + raw_command="show interface brief", + parse_status="aux_failed", + message="aux_for=x;parse boom", + row_count=0, + ) + ) + self.db.commit() + with patch.object(runner, "SessionLocal", self.Session): + with patch( + "netx_api.biz_state.compare_service.schedule_auto_compare_for_task", + ): + status = runner._finalize_batch_status( + batch_id="b-finalize", + task_id="t-finalize", + cmd_count=2, + total_rows=10, + any_fail=True, + any_ok=True, + lane_errors=[], + ) + self.assertEqual(status, "partial") + self.db.expire_all() + batch = self.db.get(BizStateBatch, "b-finalize") + assert batch is not None + self.assertEqual(batch.status, "partial") + self.assertTrue(str(batch.message or "").strip()) + self.assertIn("aux", (batch.message or "").lower()) + + def test_finalize_partial_when_skipped_bindings_with_success(self) -> None: + from netx_api.models import BizStateBatchCommand + + self.db.add( + BizStateBatchCommand( + id="c-ok2", + batch_id="b-finalize", + profile_id="zte.arp", + raw_command="show arp", + parse_status="ok", + row_count=1, + ) + ) + self.db.add( + BizStateBatchCommand( + id="c-skip", + batch_id="b-finalize", + profile_id="zte.bgp_vpnv4_neighbor_in", + raw_command="show bgp vpnv4 ...", + parse_status="skipped", + message="profile requires parameter bindings", + row_count=0, + ) + ) + self.db.commit() + with patch.object(runner, "SessionLocal", self.Session): + with patch( + "netx_api.biz_state.compare_service.schedule_auto_compare_for_task", + ): + status = runner._finalize_batch_status( + batch_id="b-finalize", + task_id="t-finalize", + cmd_count=2, + total_rows=1, + any_fail=False, + any_ok=True, + lane_errors=[], + ) + self.assertEqual(status, "partial") + self.db.expire_all() + batch = self.db.get(BizStateBatch, "b-finalize") + assert batch is not None + self.assertIn("skipped", (batch.message or "").lower()) + def test_finalize_retries_once_on_operational_error(self) -> None: calls = {"n": 0} real_session = self.Session @@ -93,6 +187,12 @@ class BizStateCollectFinalizeTests(unittest.TestCase): def get(self, *args, **kwargs): return self._inner.get(*args, **kwargs) + def query(self, *args, **kwargs): + return self._inner.query(*args, **kwargs) + + def add(self, *args, **kwargs): + return self._inner.add(*args, **kwargs) + def commit(self) -> None: calls["n"] += 1 if calls["n"] == 1: diff --git a/tests/test_zte_extended_parsers.py b/tests/test_zte_extended_parsers.py index ef534e3..bde8530 100644 --- a/tests/test_zte_extended_parsers.py +++ b/tests/test_zte_extended_parsers.py @@ -1198,5 +1198,42 @@ Total number of routes: 2 by_key = {(r["rd"], r["network"], r["next_hop"]) for r in routes} self.assertEqual(len(by_key), 2) + def test_bgp_route_vpnv6_ecmp_nh_with_metrics_same_line(self) -> None: + """RR vpnv6 wrap: indented 'NH LocPrf Path' must keep both ECMP legs.""" + raw = """ +Routes Learned From This Neighbor: +Status codes: * valid, i - internal +Total number of routes: 4 + Dest Next Hop Metric LocPrf InTag Path +Route Distinguisher:10.0.0.1:100 +* i 2407:1::/48 + 10.1.1.1 100 0 ? +* i 2407:1::/48 + 10.1.1.2 100 0 ? +* i 2407:2::/48 + 10.1.1.1 100 0 ? +* i 2407:2::/48 + 10.1.1.2 100 0 ? +""" + routes = normalize_bgp_route( + raw_text=raw, + command="show bgp vpnv6 unicast neighbor in 24.11.0.8 | one-line", + vendor="ZTE", + device_type="zte_zxros", + params={"neighbor": "24.11.0.8", "direction": "in", "afi": "vpnv6"}, + ) + self.assertEqual(len(routes), 4, "empty next_hop must not collapse ECMP") + self.assertEqual( + {(r["network"], r["next_hop"]) for r in routes}, + { + ("2407:1::/48", "10.1.1.1"), + ("2407:1::/48", "10.1.1.2"), + ("2407:2::/48", "10.1.1.1"), + ("2407:2::/48", "10.1.1.2"), + }, + ) + self.assertTrue(all(r.get("status_codes") == "*i" for r in routes)) + self.assertTrue(all(r.get("rd") == "10.0.0.1:100" for r in routes)) + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index db6445d..8a2aa45 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -285,7 +285,8 @@ const en = { batchId: "Batch ID", copyBatchId: "Copy batch ID", batchMessageEmpty: "No message or error", - batchCmdSummary: "{{ok}} ok / {{fail}} failed / {{other}} other · {{total}} commands", + batchCmdSummary: + "{{ok}} ok / {{fail}} failed / {{skipped}} skipped / {{other}} other · {{total}} counted (batch {{batchCmds}} cmds)", sheetCommands: "Commands", sheetLldp: "LLDP neighbors", sheetVrfRoute: "VRF route summary", @@ -305,6 +306,7 @@ const en = { rawLogTitle: "Raw collect log", rawLogEmpty: "No raw output", rawLogStats: "{{lines}} lines collected · {{rows}} rows stored", + rawLogDeclared: "declared {{declared}}", openWorkbook: "Open workbook", filterColumn: "Column", filterAllCols: "All columns", @@ -322,6 +324,7 @@ const en = { colTime: "Time", colRows: "DB rows", colRawLines: "CLI lines", + colDeclared: "Declared", colActions: "Actions", colCommand: "Command", colMessage: "Message", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index b13dc66..e0dadaf 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -285,7 +285,8 @@ const zh = { batchId: "批次 ID", copyBatchId: "复制批次 ID", batchMessageEmpty: "无额外说明或报错", - batchCmdSummary: "命令 {{ok}} 成功 / {{fail}} 失败 / {{other}} 其他 · 共 {{total}} 条", + batchCmdSummary: + "命令 {{ok}} 成功 / {{fail}} 失败 / {{skipped}} 跳过 / {{other}} 其他 · 统计 {{total}} 条(批次计数 {{batchCmds}})", sheetCommands: "命令", sheetLldp: "LLDP 邻居", sheetVrfRoute: "VRF 路由摘要", @@ -305,6 +306,7 @@ const zh = { rawLogTitle: "原始采集日志", rawLogEmpty: "无原始输出", rawLogStats: "采集 {{lines}} 行 · 入库 {{rows}} 条", + rawLogDeclared: "声明 {{declared}}", openWorkbook: "打开工作簿", filterColumn: "筛选列", filterAllCols: "全部列", @@ -322,6 +324,7 @@ const zh = { colTime: "时间", colRows: "入库条数", colRawLines: "采集行数", + colDeclared: "设备声明", colActions: "操作", colCommand: "命令", colMessage: "说明", diff --git a/web/src/index.css b/web/src/index.css index aec5011..3327d46 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -10181,6 +10181,20 @@ html.login-page--paused .login-page__flare { .bs-collect-detail-modal .bs-sheet-table { max-height: min(48vh, 480px); + overflow: auto; +} + +.bs-collect-detail-modal .bs-collect-detail-table table { + border-collapse: separate; + border-spacing: 0; +} + +.bs-collect-detail-modal .bs-collect-detail-table thead th { + position: sticky; + top: 0; + z-index: 3; + background: rgba(15, 23, 42, 0.97); + box-shadow: 0 1px 0 rgba(148, 163, 184, 0.25); } .bs-rawlog-pre { diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index 4ce238c..7aacd42 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -111,6 +111,7 @@ type SheetCmd = { parse_status?: string; row_count?: number; raw_line_count?: number; + declared_total?: number; message?: string; has_raw?: boolean; profile_id?: string; @@ -162,6 +163,20 @@ function intervalUnitMax(unit: "days" | "hours" | "seconds") { return 604800; // up to 7 days in seconds } +function formatCmdCollectStats( + t: (key: string, vars?: Record) => string, + lines: number, + rows: number, + declared?: number, +) { + const base = t("bizState.rawLogStats", { lines, rows }); + const d = Number(declared || 0); + if (d > 0) { + return `${base} · ${t("bizState.rawLogDeclared", { declared: d })}`; + } + return base; +} + function convertIntervalValue( value: number, from: "days" | "hours" | "seconds", @@ -311,6 +326,7 @@ export function BizStatePage() { const [rawLogCommandId, setRawLogCommandId] = useState(""); const [rawLogLines, setRawLogLines] = useState(0); const [rawLogRows, setRawLogRows] = useState(0); + const [rawLogDeclared, setRawLogDeclared] = useState(0); const [rawLogMessage, setRawLogMessage] = useState(""); /** Collect status / errors detail (lighter than workbook). */ const [collectDetail, setCollectDetail] = useState(null); @@ -550,17 +566,33 @@ export function BizStatePage() { }, [activeSheet, batchDetail, debouncedSheetKw, sheetColumn, sheetTotal]); const collectDetailCmdSummary = useMemo(() => { + const stats = collectDetail?.command_stats as + | { ok?: number; failed?: number; skipped?: number; other?: number; total?: number; aux_failed?: number } + | undefined; + if (stats && typeof stats.total === "number") { + return { + ok: Number(stats.ok || 0), + fail: Number(stats.failed || 0), + skipped: Number(stats.skipped || 0), + other: Number(stats.other || 0), + total: Number(stats.total || 0), + auxFailed: Number(stats.aux_failed || 0), + }; + } const cmds = (collectDetail?.commands || []) as SheetCmd[]; let ok = 0; let fail = 0; + let skipped = 0; let other = 0; for (const c of cmds) { const st = String(c.parse_status || "").toLowerCase(); if (st === "ok" || st === "success" || st === "aux" || st === "aux_ok" || st === "aux_cached") ok += 1; - else if (st === "failed" || st === "error" || st === "fail") fail += 1; + else if (st === "failed" || st === "error" || st === "fail" || st === "aux_failed" || st.endsWith("_failed")) + fail += 1; + else if (st.startsWith("skipped")) skipped += 1; else other += 1; } - return { ok, fail, other, total: cmds.length }; + return { ok, fail, skipped, other, total: cmds.length, auxFailed: 0 }; }, [collectDetail]); const loadSheetPage = useCallback( @@ -1192,6 +1224,7 @@ export function BizStatePage() { setRawLogCommandId(commandId); setRawLogLines(0); setRawLogRows(0); + setRawLogDeclared(0); setRawLogMessage(""); try { const d = await bizStateGetBatchCommand(bid, commandId); @@ -1199,13 +1232,15 @@ export function BizStatePage() { setRawLogText(String(d.raw_text || "")); const lines = Number(d.raw_line_count ?? 0); const rows = Number(d.row_count ?? 0); + const declared = Number(d.declared_total ?? 0); setRawLogLines(lines); setRawLogRows(rows); + setRawLogDeclared(declared); setRawLogMessage(String(d.message || "")); const bits = [ d.parse_status, d.metric_id, - t("bizState.rawLogStats", { lines, rows }), + formatCmdCollectStats(t, lines, rows, declared), d.collected_at ? fmtTime(d.collected_at) : "", ].filter(Boolean); setRawLogMeta(bits.join(" · ")); @@ -2336,10 +2371,12 @@ export function BizStatePage() { {t("bizState.exportRawLog")} - {t("bizState.rawLogStats", { - lines: Number(c.raw_line_count ?? 0), - rows: Number(c.row_count ?? 0), - })} + {formatCmdCollectStats( + t, + Number(c.raw_line_count ?? 0), + Number(c.row_count ?? 0), + Number(c.declared_total ?? 0), + )} {c.parse_status ? ( @@ -2383,10 +2420,12 @@ export function BizStatePage() { {t("bizState.exportRawLog")} - {t("bizState.rawLogStats", { - lines: Number(row.raw_line_count ?? 0), - rows: Number(row.row_count ?? 0), - })} + {formatCmdCollectStats( + t, + Number(row.raw_line_count ?? 0), + Number(row.row_count ?? 0), + Number(row.declared_total ?? 0), + )} {cellText(row.parse_status) ? ( @@ -2573,7 +2612,7 @@ export function BizStatePage() { {rawLogCmd ? {rawLogCmd} : null}
- {t("bizState.rawLogStats", { lines: rawLogLines, rows: rawLogRows })} + {formatCmdCollectStats(t, rawLogLines, rawLogRows, rawLogDeclared)} {rawLogMeta ? {rawLogMeta} : null}
@@ -2655,11 +2694,13 @@ export function BizStatePage() { {t("bizState.batchCmdSummary", { ok: collectDetailCmdSummary.ok, fail: collectDetailCmdSummary.fail, + skipped: collectDetailCmdSummary.skipped, other: collectDetailCmdSummary.other, total: collectDetailCmdSummary.total, + batchCmds: Number(collectDetail.command_count ?? collectDetailCmdSummary.total), })}

-
+
@@ -2667,6 +2708,7 @@ export function BizStatePage() { + @@ -2684,6 +2726,9 @@ export function BizStatePage() { + diff --git a/web/src/services/api.ts b/web/src/services/api.ts index ae46d94..7ec7e41 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -1915,6 +1915,7 @@ export const bizStateGetBatchCommand = (batchId: string, commandId: string) => parse_status: string; row_count: number; raw_line_count?: number; + declared_total?: number; message: string; raw_text: string; collected_at?: string | null;
{t("bizState.colStatus")} {t("bizState.colRawLines")} {t("bizState.colRows")}{t("bizState.colDeclared")} {t("bizState.colMessage")} {t("bizState.colActions")}
{Number(c.raw_line_count ?? 0)} {Number(c.row_count ?? 0)} + {Number(c.declared_total ?? 0) > 0 ? Number(c.declared_total) : "—"} + {c.message?.trim() ? c.message : "—"}