diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index 2636826..ad04d77 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -1049,27 +1049,50 @@ def _run_collect_lane( # Resolve expand_all → concrete per-VRF commands via discover profile. flat_work: list[WorkItem] = [] persisted = aux_persisted if aux_persisted is not None else set() + + def _queue_expand_fail( + *, + profile_id: str, + item_id: str, + message: str, + command: str = "", + ) -> None: + nonlocal any_fail + any_fail = True + msg = str(message or "").strip() + _emit_task_event(task_id=task_id, message=msg, level="error") + _queue( + SpooledCommand( + id=uuid4().hex, + batch_id=batch_id, + task_item_id=item_id, + profile_id=profile_id, + raw_command=(command or profile_id or "expand_all")[:512], + parse_status="failed", + message=msg[:1020], + ) + ) + for concrete, params, profile_id, item_id, mode in work: if mode != "expand_all": flat_work.append((concrete, params, profile_id, item_id, mode)) continue profile = get_profile(profile_id) if profile is None or not profile.placeholders: - any_fail = True - _emit_task_event( - task_id=task_id, + _queue_expand_fail( + profile_id=profile_id, + item_id=item_id, message=f"expand_all missing profile {profile_id}", - level="error", ) continue ph = profile.placeholders[0] disc = get_profile(str(ph.discover_profile_id or "").strip()) if disc is None: - any_fail = True - _emit_task_event( - task_id=task_id, + _queue_expand_fail( + profile_id=profile_id, + item_id=item_id, message=f"expand_all discover profile missing for {profile_id}", - level="error", + command=str(profile.command_template or "")[:512], ) continue disc_cmd = normalize_command(disc.command_template) @@ -1080,11 +1103,11 @@ def _run_collect_lane( params={}, ) if not entry.ok: - any_fail = True - _emit_task_event( - task_id=task_id, + _queue_expand_fail( + profile_id=profile_id, + item_id=item_id, message=f"expand_all discover failed: {entry.error}", - level="error", + command=disc_cmd[:512], ) continue try: @@ -1093,8 +1116,12 @@ def _run_collect_lane( records=entry.records, ) except ValueError as exc: - any_fail = True - _emit_task_event(task_id=task_id, message=str(exc), level="error") + _queue_expand_fail( + profile_id=profile_id, + item_id=item_id, + message=str(exc), + command=str(profile.command_template or "")[:512], + ) continue for cmd, p in pairs: flat_work.append((cmd, p, profile_id, item_id, "normal")) @@ -1502,17 +1529,51 @@ def _is_skip_parse_status(status: str | None) -> bool: return st.startswith("skipped") +def _is_issue_parse_status(status: str | None) -> bool: + """Statuses that explain partial/failed batches (not ok / successful aux).""" + st = str(status or "").strip().lower() + if not st or st in ("ok", "aux", "aux_cached"): + return False + if st.startswith("aux") and "fail" not in st: + return False + return ( + _is_fail_parse_status(st) + or _is_skip_parse_status(st) + or st in ("unmatched", "error") + or "fail" in st + ) + + +def _cmd_issue_line(c: Any) -> str: + """One line: [status] profile|metric | command | reason.""" + st = str(getattr(c, "parse_status", "") or "").strip() or "unknown" + profile = str(getattr(c, "profile_id", "") or "").strip() + metric = str(getattr(c, "metric_id", "") or "").strip() + item = profile or metric or "-" + if profile and metric and profile != metric: + item = f"{profile}/{metric}" + cmd = str(getattr(c, "raw_command", "") or "").strip() or "(no command)" + if len(cmd) > 120: + cmd = cmd[:117] + "..." + reason = str(getattr(c, "message", "") or "").strip() or st + if len(reason) > 180: + reason = reason[:177] + "..." + return f"[{st}] {item} | {cmd} | {reason}" + + def _batch_issue_summary( db, batch_id: str, lane_errors: list[str] | None = None, + *, + task_id: str = "", ) -> str: - """Human-readable reason for partial/failed batches (never leave message empty).""" - parts: list[str] = [] + """List which profile/command failed or was skipped (fits batch.message 1024).""" + lines: list[str] = [] for err in lane_errors or []: e = str(err or "").strip() - if e and e not in parts: - parts.append(e) + if e and e not in lines: + lines.append(e) cmds = ( db.query(BizStateBatchCommand) @@ -1523,36 +1584,79 @@ def _batch_issue_summary( n_fail = 0 n_skip = 0 n_aux_fail = 0 - samples: list[str] = [] + n_unmatched = 0 + detail: 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): + if not _is_issue_parse_status(st): + continue + if st == "unmatched": + n_unmatched += 1 + elif 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]) + else: + n_fail += 1 + detail.append(_cmd_issue_line(c)) + counts: list[str] = [] if n_fail: - parts.append(f"{n_fail} command(s) failed") + counts.append(f"failed={n_fail}") if n_aux_fail: - parts.append(f"{n_aux_fail} aux command(s) failed") + counts.append(f"aux_failed={n_aux_fail}") + if n_unmatched: + counts.append(f"unmatched={n_unmatched}") 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) + counts.append(f"skipped={n_skip}") + if counts: + lines.append("issues: " + " ".join(counts)) - text = "; ".join(parts).strip() + # Prefer as many concrete command lines as fit in the remaining budget. + budget = 1020 + used = sum(len(x) + 1 for x in lines) + shown = 0 + omitted = 0 + for d in detail: + # +1 for newline + if used + len(d) + 1 > budget - 40: + omitted = len(detail) - shown + break + lines.append(d) + used += len(d) + 1 + shown += 1 + if omitted > 0: + lines.append(f"...and {omitted} more (see Commands sheet)") + + if not detail and task_id: + # expand_all / barrier failures may only leave task events. + evs = ( + db.query(BizStateEvent) + .filter( + BizStateEvent.task_id == task_id, + BizStateEvent.level.in_(("error", "warn", "warning")), + ) + .order_by(BizStateEvent.created_at.desc()) + .limit(5) + .all() + ) + for ev in reversed(evs): + msg = str(ev.message or "").strip() + if not msg: + continue + line = f"[event] {msg}"[:200] + if used + len(line) + 1 > budget: + break + if line not in lines: + lines.append(line) + used += len(line) + 1 + + text = "\n".join(lines).strip() if not text: - text = "partial success (some steps failed or were skipped)" + text = ( + "partial: some steps failed, but no per-command details were persisted " + "(check task events / Commands sheet)" + ) return text[:1020] @@ -1577,6 +1681,7 @@ def _finalize_batch_status( batch.command_count = max(int(batch.command_count or 0), int(cmd_count or 0)) batch.row_count = max(int(batch.row_count or 0), int(total_rows or 0)) batch.ended_at = _utcnow() + tid = str(task_id or batch.task_id or "") if stopped: if any_ok or int(batch.command_count or 0) > 0 or int(batch.row_count or 0) > 0: batch.status = "partial" @@ -1586,10 +1691,14 @@ def _finalize_batch_status( batch.message = STOP_USER_MESSAGE elif any_fail and any_ok: batch.status = "partial" - batch.message = _batch_issue_summary(db, batch_id, lane_errors) + batch.message = _batch_issue_summary( + db, batch_id, lane_errors, task_id=tid + ) elif any_fail and not any_ok: batch.status = "failed" - batch.message = _batch_issue_summary(db, batch_id, lane_errors) or ( + batch.message = _batch_issue_summary( + db, batch_id, lane_errors, task_id=tid + ) or ( "; ".join(lane_errors)[:1020] if lane_errors else "all commands failed" ) else: @@ -1603,7 +1712,9 @@ def _finalize_batch_status( 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) + batch.message = _batch_issue_summary( + db, batch_id, lane_errors, task_id=tid + ) else: batch.status = "success" batch.message = "" diff --git a/netx_api/biz_state/parsers/zte/bgp_route.py b/netx_api/biz_state/parsers/zte/bgp_route.py index cdb04d2..f3fc985 100644 --- a/netx_api/biz_state/parsers/zte/bgp_route.py +++ b/netx_api/biz_state/parsers/zte/bgp_route.py @@ -193,7 +193,9 @@ def _skip_noise_line(line: str) -> bool: return True if low.startswith("valid ") or low.startswith("invalid "): return True - if re.match(r"^\d{1,2}:\d{2}:\d{2}\b", low): + # Clock banners like "09:30:01" — require whitespace/EOL so IPv6 + # prefixes such as "56:16:10::/64" are not treated as HH:MM:SS. + if re.match(r"^\d{1,2}:\d{2}:\d{2}(?:\s|$)", low): return True if low.endswith("#") or "#'" in low: return True diff --git a/tests/test_biz_state_collect_finalize.py b/tests/test_biz_state_collect_finalize.py index 9012ab9..66d63f1 100644 --- a/tests/test_biz_state_collect_finalize.py +++ b/tests/test_biz_state_collect_finalize.py @@ -127,7 +127,9 @@ class BizStateCollectFinalizeTests(unittest.TestCase): assert batch is not None self.assertEqual(batch.status, "partial") self.assertTrue(str(batch.message or "").strip()) - self.assertIn("aux", (batch.message or "").lower()) + self.assertIn("aux_failed", (batch.message or "").lower()) + self.assertIn("show interface brief", batch.message or "") + self.assertIn("parse boom", batch.message or "") def test_finalize_partial_when_skipped_bindings_with_success(self) -> None: from netx_api.models import BizStateBatchCommand @@ -172,6 +174,71 @@ class BizStateCollectFinalizeTests(unittest.TestCase): batch = self.db.get(BizStateBatch, "b-finalize") assert batch is not None self.assertIn("skipped", (batch.message or "").lower()) + self.assertIn("zte.bgp_vpnv4_neighbor_in", batch.message or "") + self.assertIn("bindings", (batch.message or "").lower()) + + def test_finalize_partial_lists_unmatched_and_failed_commands(self) -> None: + from netx_api.models import BizStateBatchCommand + + self.db.add( + BizStateBatchCommand( + id="c-ok3", + batch_id="b-finalize", + profile_id="zte.arp", + metric_id="arp", + raw_command="show arp", + parse_status="ok", + row_count=1, + ) + ) + self.db.add( + BizStateBatchCommand( + id="c-unmatched", + batch_id="b-finalize", + profile_id="zte.legacy_x", + raw_command="show weird-legacy", + parse_status="unmatched", + message="no profile matched concrete command", + row_count=0, + ) + ) + self.db.add( + BizStateBatchCommand( + id="c-fail", + batch_id="b-finalize", + profile_id="zte.bgp_route", + metric_id="bgp_route", + raw_command="show bgp vpnv4 unicast neighbor in 1.1.1.1", + parse_status="failed", + message="parse: ValueError: 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=3, + total_rows=1, + 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 + msg = batch.message or "" + self.assertIn("unmatched=1", msg) + self.assertIn("failed=1", msg) + self.assertIn("show weird-legacy", msg) + self.assertIn("show bgp vpnv4", msg) + self.assertIn("zte.bgp_route", msg) + self.assertNotIn("partial success (some steps failed", msg) def test_finalize_retries_once_on_operational_error(self) -> None: calls = {"n": 0} diff --git a/tests/test_zte_extended_parsers.py b/tests/test_zte_extended_parsers.py index 88c14e1..51e485e 100644 --- a/tests/test_zte_extended_parsers.py +++ b/tests/test_zte_extended_parsers.py @@ -1242,6 +1242,43 @@ Route Distinguisher:10.0.0.1:100 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)) + def test_bgp_route_vpnv6_prefix_not_skipped_as_clock(self) -> None: + """Prefixes like 56:16:10::/64 must not match HH:MM:SS noise filter.""" + from netx_api.biz_state.parsers.zte.bgp_route import _skip_noise_line + + self.assertTrue(_skip_noise_line("09:30:01")) + self.assertTrue(_skip_noise_line("9:05:00 system ready")) + self.assertFalse(_skip_noise_line("56:16:10::/64")) + self.assertFalse(_skip_noise_line("* i 56:16:10::/64")) + + raw = """ +Routes Advertised to This Neighbor: +Status codes: * valid, i - internal +Total number of routes: 2 + Dest Next Hop Metric LocPrf InTag Path +Route Distinguisher:65525:30001 (default for vrf CUST_V6) +* i 56:16:10::/64 ::FFFF:10.1.1.1 100 0 ? +* i 56:16:32::/64 ::FFFF:10.1.1.1 100 0 ? +""" + routes = normalize_bgp_route( + raw_text=raw, + command="show bgp vpnv6 unicast neighbor out 24.11.0.8 as 24208 | one-line", + vendor="ZTE", + device_type="zte_zxros", + params={ + "neighbor": "24.11.0.8", + "direction": "out", + "afi": "vpnv6", + "local_as": "24208", + }, + ) + self.assertEqual(len(routes), 2) + self.assertEqual( + {r["network"] for r in routes}, + {"56:16:10::/64", "56:16:32::/64"}, + ) + self.assertTrue(all(r.get("rd") == "65525:30001" for r in routes)) + def test_bgp_route_prod_snippets_rd_vrf_and_wraps(self) -> None: """Live RR snippets: OUT without status, RD+VRF, vpnv6 *i + NH wrap.""" from pathlib import Path