diff --git a/netx_api/biz_state/compare_engine.py b/netx_api/biz_state/compare_engine.py index b5e1251..3871cc4 100644 --- a/netx_api/biz_state/compare_engine.py +++ b/netx_api/biz_state/compare_engine.py @@ -2,6 +2,7 @@ from __future__ import annotations +from collections import defaultdict, deque from typing import Any, Mapping, Sequence from .compare_rules import field_rule_map, values_equal, explain_diff @@ -11,6 +12,53 @@ from .iface_normalize import ( resolve_mapped_iface, ) +# Prefer these fields when stratifying success samples (BGP multipath / multi-cmd). +_STRATUM_FIELDS = ("neighbor", "direction", "afi", "vrf") + + +def stratum_key(row: Mapping[str, Any] | None) -> str: + """Bucket key for stratified unchanged sampling.""" + if not isinstance(row, Mapping): + return "_" + parts: list[str] = [] + for f in _STRATUM_FIELDS: + v = str(row.get(f) or "").strip() + if v: + parts.append(f"{f}={v}") + return "|".join(parts) if parts else "_" + + +def stratify_take(items: list[Any], limit: int, *, key_fn) -> list[Any]: + """Round-robin across strata so one neighbor/direction cannot consume the whole sample.""" + lim = max(0, int(limit)) + if lim <= 0 or not items: + return [] + if len(items) <= lim: + return list(items) + buckets: dict[str, deque[Any]] = defaultdict(deque) + order: list[str] = [] + for it in items: + sk = str(key_fn(it) or "_") + if sk not in buckets: + order.append(sk) + buckets[sk].append(it) + out: list[Any] = [] + while len(out) < lim and buckets: + drained: list[str] = [] + for sk in order: + q = buckets.get(sk) + if not q: + drained.append(sk) + continue + out.append(q.popleft()) + if len(out) >= lim: + break + for sk in drained: + buckets.pop(sk, None) + if sk in order: + order = [x for x in order if x != sk] + return out + def apply_port_map( row: dict[str, Any], @@ -111,9 +159,14 @@ def compare_rows( ) -> dict[str, Any]: """Return summary + diffs list. - Diff kinds: added | removed | changed | unchanged | duplicate + Diff kinds: added | removed | changed | unchanged - Pipeline: iface normalize (both sides) → port map (before) → match. + Pipeline: iface normalize (both sides) → port map (before) → ordered + same-key pairing (load / list order within each match key). + + Same match key with N before and M after rows: zip by index + ``0..min(N,M)-1`` for field compare; extras become ``removed`` (before) + or ``added`` (after). Example: 5 vs 2 → 2 compared + 3 removed. ``include_unchanged``: when False, matching rows still increment ``summary.unchanged`` but are omitted from ``diffs``. @@ -129,8 +182,8 @@ def compare_rows( - ``True``: force drop iface from match key when a non-empty candidate exists. - ``False``: never drop iface from match key. - Duplicate match keys are not silently discarded: extras become ``duplicate`` - diffs and ``summary.duplicate_key_list`` lists the colliding keys. + ``summary.duplicate`` is always 0. ``duplicate_keys_*`` count match keys + that appear more than once on a side (diagnostic only). ``field_rules`` drives normalize / numeric tolerance / per-field compare mode (template-driven; no metric-specific branches here). @@ -177,26 +230,22 @@ def compare_rows( apply_port_map(r, iface_fields=iface_list, port_map=pmap) for r in before_norm ] - after_index: dict[tuple[str, ...], dict[str, Any]] = {} - after_dup = 0 - after_dup_keys: list[tuple[str, ...]] = [] - after_dup_rows: list[tuple[tuple[str, ...], dict[str, Any]]] = [] + before_groups: dict[tuple[str, ...], list[tuple[dict[str, Any], dict[str, Any]]]] = ( + defaultdict(list) + ) + after_groups: dict[tuple[str, ...], list[dict[str, Any]]] = defaultdict(list) + for orig, mapped in zip(before_rows, before_mapped): + before_groups[row_key(mapped, match_keys)].append((orig, mapped)) for r in after_norm: - k = row_key(r, match_keys) - if k in after_index: - after_dup += 1 - after_dup_keys.append(k) - after_dup_rows.append((k, r)) - continue # first wins — do not overwrite - after_index[k] = r + after_groups[row_key(r, match_keys)].append(r) - before_keys: set[tuple[str, ...]] = set() - before_dup = 0 - before_dup_keys: list[tuple[str, ...]] = [] diffs: list[dict[str, Any]] = [] - added = removed = changed = unchanged = duplicate = 0 + added = removed = changed = unchanged = 0 unchanged_listed = 0 limit_n = None if unchanged_limit is None else max(0, int(unchanged_limit)) + multi_before_keys: list[tuple[str, ...]] = [] + multi_after_keys: list[tuple[str, ...]] = [] + unchanged_candidates: list[tuple[dict[str, Any], dict[str, Any], dict[str, Any]]] = [] def _key_obj(row: dict[str, Any]) -> dict[str, Any]: return {f: row.get(f, "") for f in key_fields} @@ -209,12 +258,10 @@ def compare_rows( return str(netx.get("row_id") or "") return "" - def _emit_unchanged(orig: dict[str, Any], mapped: dict[str, Any], after_row: dict[str, Any]) -> None: + def _append_unchanged_diff( + orig: dict[str, Any], mapped: dict[str, Any], after_row: dict[str, Any] + ) -> None: nonlocal unchanged_listed - if not include_unchanged: - return - if limit_n is not None and unchanged_listed >= limit_n: - return unchanged_listed += 1 before_rid = _row_id(orig) after_rid = _row_id(after_row) @@ -246,96 +293,103 @@ def compare_rows( } ) + # Stable key order: before encounter order, then after-only keys + seen_keys: set[tuple[str, ...]] = set() + ordered_keys: list[tuple[str, ...]] = [] for orig, mapped in zip(before_rows, before_mapped): k = row_key(mapped, match_keys) - if k in before_keys: - before_dup += 1 - before_dup_keys.append(k) - duplicate += 1 - diffs.append( - { - "kind": "duplicate", - "side": "before", - "key": _key_obj(mapped), - "before": orig, - "after": after_index.get(k), - "mapped_before": mapped, - "changes": {}, - } - ) - continue # only first before row participates in match - before_keys.add(k) - after = after_index.get(k) - if after is None: - removed += 1 - diffs.append( - { - "kind": "removed", - "key": _key_obj(mapped), - "before": orig, - "after": None, - "mapped_before": mapped, - "changes": {}, - } - ) - continue - # Empty compare_fields = presence-only: keyed rows that exist on both - # sides are unchanged (no value checks). - field_changes: dict[str, dict[str, Any]] = {} - for f in compare_fields: - bv = mapped.get(f, "") - av = after.get(f, "") - rule = rules.get(f) - if not values_equal(bv, av, rule=rule): - entry: dict[str, Any] = {"before": bv, "after": av} - reason = explain_diff(bv, av, rule=rule) - if reason: - entry["reason"] = reason - field_changes[f] = entry - if field_changes: - changed += 1 - diffs.append( - { - "kind": "changed", - "key": _key_obj(mapped), - "before": orig, - "after": after, - "mapped_before": mapped, - "changes": field_changes, - } - ) + if k not in seen_keys: + seen_keys.add(k) + ordered_keys.append(k) + for r in after_norm: + k = row_key(r, match_keys) + if k not in seen_keys: + seen_keys.add(k) + ordered_keys.append(k) + for k in ordered_keys: + b_list = before_groups.get(k) or [] + a_list = after_groups.get(k) or [] + if len(b_list) > 1: + multi_before_keys.append(k) + if len(a_list) > 1: + multi_after_keys.append(k) + n = max(len(b_list), len(a_list)) + for i in range(n): + if i >= len(b_list): + after = a_list[i] + added += 1 + diffs.append( + { + "kind": "added", + "key": _key_obj(after), + "before": None, + "after": after, + "mapped_before": None, + "changes": {}, + "before_row_id": "", + "after_row_id": _row_id(after), + } + ) + continue + if i >= len(a_list): + orig, mapped = b_list[i] + removed += 1 + diffs.append( + { + "kind": "removed", + "key": _key_obj(mapped), + "before": orig, + "after": None, + "mapped_before": mapped, + "changes": {}, + "before_row_id": _row_id(orig), + "after_row_id": "", + } + ) + continue + orig, mapped = b_list[i] + after = a_list[i] + field_changes: dict[str, dict[str, Any]] = {} + for f in compare_fields: + bv = mapped.get(f, "") + av = after.get(f, "") + rule = rules.get(f) + if not values_equal(bv, av, rule=rule): + entry: dict[str, Any] = {"before": bv, "after": av} + reason = explain_diff(bv, av, rule=rule) + if reason: + entry["reason"] = reason + field_changes[f] = entry + if field_changes: + changed += 1 + diffs.append( + { + "kind": "changed", + "key": _key_obj(mapped), + "before": orig, + "after": after, + "mapped_before": mapped, + "changes": field_changes, + "before_row_id": _row_id(orig), + "after_row_id": _row_id(after), + } + ) + else: + unchanged += 1 + if include_unchanged: + unchanged_candidates.append((orig, mapped, after)) + + if include_unchanged and unchanged_candidates: + if limit_n is None: + picked = unchanged_candidates else: - unchanged += 1 - _emit_unchanged(orig, mapped, after) - - for k, after in after_index.items(): - if k in before_keys: - continue - added += 1 - diffs.append( - { - "kind": "added", - "key": _key_obj(after), - "before": None, - "after": after, - "mapped_before": None, - "changes": {}, - } - ) - - for k, after in after_dup_rows: - duplicate += 1 - diffs.append( - { - "kind": "duplicate", - "side": "after", - "key": _key_obj(after), - "before": None, - "after": after, - "mapped_before": None, - "changes": {}, - } - ) + picked = stratify_take( + unchanged_candidates, + limit_n, + key_fn=lambda t: stratum_key(t[1]), + ) + for orig, mapped, after_row in picked: + _append_unchanged_diff(orig, mapped, after_row) def _fmt_keys(keys: list[tuple[str, ...]]) -> list[str]: seen: set[str] = set() @@ -345,7 +399,7 @@ def compare_rows( if s not in seen: seen.add(s) out.append(s) - return out + return out[:64] stats = mapping_stats( before_rows=before_norm, @@ -363,16 +417,21 @@ def compare_rows( "removed": removed, "changed": changed, "unchanged": unchanged, - "duplicate": duplicate, + "duplicate": 0, "match_key_fields": match_keys, - "duplicate_keys_before": before_dup, - "duplicate_keys_after": after_dup, - "duplicate_key_list": _fmt_keys(before_dup_keys + after_dup_keys), + "duplicate_keys_before": len(multi_before_keys), + "duplicate_keys_after": len(multi_after_keys), + "duplicate_key_list": _fmt_keys(multi_before_keys + multi_after_keys), "unchanged_listed": unchanged_listed, "unchanged_truncated": bool( include_unchanged and limit_n is not None and unchanged > unchanged_listed ), "unchanged_compact": bool(compact_unchanged and unchanged_listed > 0), + "unchanged_sample_mode": ( + "stratified" + if include_unchanged and limit_n is not None and unchanged > unchanged_listed + else ("full" if include_unchanged else "none") + ), }, "diffs": diffs, "mapping_stats": stats, diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py index d6d1553..b1fa66c 100644 --- a/netx_api/biz_state/compare_service.py +++ b/netx_api/biz_state/compare_service.py @@ -186,6 +186,13 @@ def _compare_side( _DIFF_CHUNK = 2000 _LOAD_YIELD_PER = 5000 _SEARCH_TEXT_MAX = 4000 +# Heartbeat while bulk-inserting large fail/ok diff sets (vpnv4-scale). +_PERSIST_PROGRESS_EVERY = 10_000 +# Above this, store fail diffs as key + row_id + changes (hydrate sides on read). +_FAIL_COMPACT_MIN = 50_000 +# Live search (kw): load at most this many matching rows per side, return ≤ this many pairs. +_LIVE_SEARCH_LOAD_CAP = 2_000 +_LIVE_SEARCH_RESULT_CAP = 200 # Success-row persist policy (see resolve_unchanged_policy) _STORE_UNCHANGED_MODES = frozenset({"auto", "always", "never", "sample", "keys"}) _UNCHANGED_FULL_MAX = 20_000 @@ -273,15 +280,41 @@ def _persist_sheet_diffs( metric_id: str, diffs: list[dict[str, Any]], seq_start: int = 0, + on_progress: Callable[[int, int], None] | None = None, ) -> int: """Bulk-insert diff rows already selected by the engine policy. Compact success rows carry key + before/after_row_id; JSON sides stay empty and are hydrated from metric tables on read. + + ``on_progress(written, total)`` fires periodically so UI elapsed time moves + during multi-minute inserts (e.g. large vpnv4 fail sets). """ buf: list[dict[str, Any]] = [] seq = int(seq_start or 0) written = 0 + total = len(diffs) + last_prog = 0 + last_prog_t = time.monotonic() + + def _maybe_prog(force: bool = False) -> None: + nonlocal last_prog, last_prog_t + if not on_progress: + return + now = time.monotonic() + if ( + not force + and written - last_prog < _PERSIST_PROGRESS_EVERY + and now - last_prog_t < 2.0 + ): + return + last_prog = written + last_prog_t = now + try: + on_progress(written, total) + except Exception: + _log.exception("persist progress callback failed run=%s metric=%s", run_id, metric_id) + for d in diffs: kind = str(d.get("kind") or "") before = _strip_netx(d.get("before")) @@ -295,10 +328,14 @@ def _persist_sheet_diffs( "mapped_before": mapped, "changes": dict(d.get("changes") or {}), } - # Compact success: search_text = kind + key only (no fat sides) + # Compact rows: search_text = kind + key (+ changes) only — no fat sides search_src = ( - {"kind": kind, "key": payload["key"]} - if kind == "unchanged" and bool(d.get("compact")) + { + "kind": kind, + "key": payload["key"], + **({"changes": payload["changes"]} if payload["changes"] else {}), + } + if bool(d.get("compact")) else payload ) buf.append( @@ -323,8 +360,14 @@ def _persist_sheet_diffs( if len(buf) >= _DIFF_CHUNK: db.bulk_insert_mappings(BizCompareDiff, buf) buf.clear() + # Commit chunks so progress/UI can see mid-write fail rows and + # elapsed_ms advances (otherwise persisting_* looks frozen). + db.commit() + _maybe_prog() if buf: db.bulk_insert_mappings(BizCompareDiff, buf) + db.commit() + _maybe_prog(force=True) return written @@ -514,6 +557,7 @@ def _pending_sheet_meta(sheet: dict[str, Any]) -> dict[str, Any]: "compare_fields": compare_fields, "display_fields": list(sheet.get("display_fields") or []), "field_rules": list(sheet.get("field_rules") or []), + "row_filters": list(sheet.get("row_filters") or []), "ignore_port_changes": sheet.get("ignore_port_changes"), "mode": "presence" if not compare_fields else "fields", "status": "pending", @@ -2423,6 +2467,39 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: fail_diffs = [d for d in diffs if str(d.get("kind") or "") != "unchanged"] ok_diffs = [d for d in diffs if str(d.get("kind") or "") == "unchanged"] mid = sheet_key(one) + # Million-row vpnv4 with many diffs: keep key/row_id/changes only. + if len(fail_diffs) >= _FAIL_COMPACT_MIN: + for d in fail_diffs: + d["before"] = {} + d["after"] = {} + d["mapped_before"] = {} + d["compact"] = True + + def _on_persist( + written: int, + total_n: int, + *, + phase: str, + kind_key: str, + ) -> None: + _set_run_progress( + db, + run, + phase=phase, + sheet_index=idx, + sheet_total=total, + sheet=sheet, + started_mono=started_mono, + extra={ + kind_key: total_n, + "persisted": written, + "persist_total": total_n, + }, + # Separate session so chunk commits in persist do not race + # with progress JSON writes on the worker session. + detach=True, + ) + _set_run_progress( db, run, @@ -2431,12 +2508,22 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: sheet_total=total, sheet=sheet, started_mono=started_mono, - extra={"fail_rows": len(fail_diffs)}, + extra={ + "fail_rows": len(fail_diffs), + "persisted": 0, + "persist_total": len(fail_diffs), + }, ) n_fail = _persist_sheet_diffs( - db, run_id=run.id, metric_id=mid, diffs=fail_diffs, seq_start=0 + db, + run_id=run.id, + metric_id=mid, + diffs=fail_diffs, + seq_start=0, + on_progress=lambda w, n: _on_persist( + w, n, phase="persisting_fail", kind_key="fail_rows" + ), ) - db.commit() if ok_diffs: _set_run_progress( db, @@ -2446,7 +2533,11 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: sheet_total=total, sheet=sheet, started_mono=started_mono, - extra={"ok_rows": len(ok_diffs)}, + extra={ + "ok_rows": len(ok_diffs), + "persisted": 0, + "persist_total": len(ok_diffs), + }, ) _persist_sheet_diffs( db, @@ -2454,6 +2545,9 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: metric_id=mid, diffs=ok_diffs, seq_start=n_fail, + on_progress=lambda w, n: _on_persist( + w, n, phase="persisting_ok", kind_key="ok_rows" + ), ) for k in agg: agg[k] += int(s.get(k) or 0) @@ -2467,6 +2561,7 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: "compare_fields": one["compare_fields"], "display_fields": one.get("display_fields") or [], "field_rules": one.get("field_rules") or [], + "row_filters": one.get("row_filters") or [], "ignore_port_changes": one.get("ignore_port_changes"), "mode": one["mode"], "status": "done", @@ -2948,6 +3043,270 @@ def _lookup_sheet(sheets: list[dict[str, Any]], key: str) -> dict[str, Any] | No return None +def _kind_allows(kind_n: str, diff_kind: str) -> bool: + dk = str(diff_kind or "") + if kind_n == "all": + return True + if kind_n == "diff": + return dk in ("removed", "changed") + return dk == kind_n + + +def _kw_match_sql(key_fields: list[str], *, param: str = "kw") -> str: + """OR of ILIKE on key fields (and data_json::text fallback).""" + from .compare_sql import _FIELD_RE, _safe_field + + parts: list[str] = [] + for f in key_fields: + name = str(f or "").strip() + if not name or not _FIELD_RE.match(name): + continue + sf = _safe_field(name) + parts.append(f"lower(trim(both from coalesce(data_json->>'{sf}', ''))) LIKE :{param}") + # Broad fallback so free-text still hits non-key columns (path, next_hop, …) + parts.append(f"lower(data_json::text) LIKE :{param}") + return "(" + " OR ".join(parts) + ")" if parts else f"(lower(data_json::text) LIKE :{param})" + + +def _load_metric_rows_for_search( + db: Session, + *, + batch_id: str, + metric_id: str, + row_filters: list[dict[str, Any]] | None, + key_fields: list[str], + kw: str, + cap: int = _LIVE_SEARCH_LOAD_CAP, +) -> tuple[list[dict[str, Any]], bool]: + """Load rows matching sheet filters + kw. Returns (rows, truncated).""" + from .compare_sql import ( + _dialect_is_postgres, + _filters_sql_compatible, + compile_row_filters_sql, + ) + from ..models import BizStateMetricRow + from sqlalchemy import text as sql_text + + bid = str(batch_id or "").strip() + mid = str(metric_id or "").strip() + needle = str(kw or "").strip() + if not bid or not mid or not needle: + return [], False + lim = max(1, min(int(cap), _LIVE_SEARCH_LOAD_CAP)) + filters = [f for f in (row_filters or []) if isinstance(f, dict)] + like = f"%{needle.lower()}%" + + if _dialect_is_postgres(db): + filter_sql, filter_params = ("TRUE", {}) + if filters and _filters_sql_compatible(filters): + filter_sql, filter_params = compile_row_filters_sql(filters) + kw_sql = _kw_match_sql(key_fields) + # Fetch lim+1 to detect truncation + params = { + "bid": bid, + "mid": mid, + "kw": like, + "lim": lim + 1, + **filter_params, + } + rows = db.execute( + sql_text( + f""" + SELECT id, batch_command_id, task_id, ne_id, seq, data_json, collected_at + FROM biz_state_metric_row + WHERE batch_id = :bid + AND metric_id = :mid + AND ({filter_sql}) + AND ({kw_sql}) + ORDER BY seq ASC, id ASC + LIMIT :lim + """ + ), + params, + ).mappings().all() + truncated = len(rows) > lim + rows = rows[:lim] + out: list[dict[str, Any]] = [] + for r in rows: + data = dict(r["data_json"] or {}) + collected = r["collected_at"] + out.append( + { + **data, + "_netx": { + "batch_id": bid, + "batch_command_id": str(r["batch_command_id"] or ""), + "task_id": str(r["task_id"] or ""), + "ne_id": str(r["ne_id"] or ""), + "collected_at": collected.isoformat() + "Z" + if collected is not None + else None, + "row_id": str(r["id"]), + }, + } + ) + return out, truncated + + # Non-PG / fallback: scan with early stop (OK for tests / small sheets) + q = ( + db.query(BizStateMetricRow) + .filter( + BizStateMetricRow.batch_id == bid, + BizStateMetricRow.metric_id == mid, + ) + .order_by(BizStateMetricRow.seq.asc(), BizStateMetricRow.id.asc()) + ) + out = [] + truncated = False + needle_l = needle.lower() + key_set = [str(k).strip() for k in key_fields if str(k).strip()] + for r in q.yield_per(500): + data = dict(r.data_json or {}) + row = { + **data, + "_netx": { + "batch_id": bid, + "batch_command_id": r.batch_command_id or "", + "task_id": r.task_id or "", + "ne_id": r.ne_id or "", + "collected_at": r.collected_at.isoformat() + "Z" + if r.collected_at + else None, + "row_id": r.id, + }, + } + if filters and not all(row_matches_filter(row, f) for f in filters): + db.expunge(r) + continue + hit = False + for kf in key_set: + if needle_l in str(row.get(kf) or "").lower(): + hit = True + break + if not hit: + blob = json.dumps(data, ensure_ascii=False, default=str).lower() + hit = needle_l in blob + if not hit: + db.expunge(r) + continue + out.append(row) + db.expunge(r) + if len(out) >= lim: + # peek one more? + truncated = True + break + return out, truncated + + +def _live_search_sheet_diffs( + db: Session, + run: BizCompareRun, + sheet: dict[str, Any], + *, + kind: str, + kw: str, + page: int, + page_size: int, +) -> dict[str, Any]: + """Search before/after metric tables, zip-compare, filter by kind tab.""" + mid_src = str(sheet.get("metric_id") or "").strip() + sid = sheet_key(sheet) + key_fields = list(sheet.get("key_fields") or []) + iface_fields = list(sheet.get("iface_fields") or []) + field_rules = list(sheet.get("field_rules") or []) + compare_fields = effective_compare_fields( + list(sheet.get("compare_fields") or []), + field_rules, + ) + row_filters = list(sheet.get("row_filters") or []) + # Older runs may lack row_filters on sheet meta — fall back to template + if not row_filters and run.template_id: + tpl = db.get(BizCompareTemplate, run.template_id) + if tpl: + for s in template_metrics(tpl): + if sheet_key(s) == sid: + row_filters = list(s.get("row_filters") or []) + if not key_fields: + key_fields = list(s.get("key_fields") or []) + if not field_rules: + field_rules = list(s.get("field_rules") or []) + compare_fields = effective_compare_fields( + list(s.get("compare_fields") or compare_fields), + field_rules, + ) + break + + if not key_fields or not mid_src: + return { + "total": 0, + "page": page, + "page_size": page_size, + "metric_id": sid, + "items": [], + "source": "live", + "truncated": False, + } + + ignore_ports = sheet.get("ignore_port_changes") + if ignore_ports is not None: + ignore_ports = bool(ignore_ports) + pmap = _port_map_dict(db, str(run.mapping_id or "")) + tpl = db.get(BizCompareTemplate, run.template_id) if run.template_id else None + norm_rules = template_iface_normalize(tpl) + + before_rows, trunc_b = _load_metric_rows_for_search( + db, + batch_id=str(run.before_batch_id or ""), + metric_id=mid_src, + row_filters=row_filters, + key_fields=key_fields, + kw=kw, + ) + after_rows, trunc_a = _load_metric_rows_for_search( + db, + batch_id=str(run.after_batch_id or ""), + metric_id=mid_src, + row_filters=row_filters, + key_fields=key_fields, + kw=kw, + ) + result = compare_rows( + before_rows=before_rows, + after_rows=after_rows, + key_fields=key_fields, + iface_fields=iface_fields, + compare_fields=compare_fields, + port_map=pmap, + field_rules=field_rules, + iface_normalize_rules=norm_rules, + ignore_port_changes=ignore_ports, + include_unchanged=True, + unchanged_limit=None, + compact_unchanged=False, + ) + kind_n = (kind or "diff").strip().lower() + filtered = [ + d for d in list(result.get("diffs") or []) if _kind_allows(kind_n, str(d.get("kind") or "")) + ] + # Cap pairs returned to keep UI snappy + truncated = bool(trunc_b or trunc_a or len(filtered) > _LIVE_SEARCH_RESULT_CAP) + filtered = filtered[:_LIVE_SEARCH_RESULT_CAP] + total = len(filtered) + start = (page - 1) * page_size + page_items = filtered[start : start + page_size] + return { + "total": total, + "page": page, + "page_size": page_size, + "metric_id": sid, + "items": page_items, + "source": "live", + "truncated": truncated, + "live_before_matched": len(before_rows), + "live_after_matched": len(after_rows), + } + + def list_run_diffs( db: Session, run_id: str, @@ -2973,6 +3332,18 @@ def list_run_diffs( sheet = _lookup_sheet(sheets, asked) if asked else (sheets[0] if sheets else None) mid = sheet_key(sheet) if sheet else (asked or str(r.metric_id or "")) + # Unified search: any kw → live source lookup; kind tab only filters. + if kw_n and sheet: + return _live_search_sheet_diffs( + db, + r, + sheet, + kind=kind_n, + kw=kw_n, + page=page_n, + page_size=size_n, + ) + if _run_has_diff_rows(db, run_id): q = db.query(BizCompareDiff).filter( BizCompareDiff.run_id == run_id, @@ -2980,10 +3351,12 @@ def list_run_diffs( ) if kind_n == "diff": q = q.filter(BizCompareDiff.kind.in_(("removed", "changed"))) + elif kind_n == "removed": + q = q.filter(BizCompareDiff.kind == "removed") + elif kind_n == "changed": + q = q.filter(BizCompareDiff.kind == "changed") elif kind_n != "all": q = q.filter(BizCompareDiff.kind == kind_n) - if kw_n: - q = q.filter(BizCompareDiff.search_text.ilike(f"%{kw_n}%")) total = q.count() rows = ( q.order_by(BizCompareDiff.seq.asc(), BizCompareDiff.id.asc()) @@ -2998,6 +3371,8 @@ def list_run_diffs( "page_size": size_n, "metric_id": mid, "items": items, + "source": "stored", + "truncated": False, } # Legacy: diffs embedded in summary_json / diffs_json @@ -3007,7 +3382,7 @@ def list_run_diffs( inline = list((sheet or {}).get("diffs") or []) if not inline and mid == r.metric_id: inline = list(r.diffs_json or []) - filtered = _filter_inline_diffs(inline, kind=kind_n, kw=kw_n) + filtered = _filter_inline_diffs(inline, kind=kind_n, kw="") total = len(filtered) start = (page_n - 1) * size_n page_items = filtered[start : start + size_n] @@ -3017,6 +3392,8 @@ def list_run_diffs( "page_size": size_n, "metric_id": mid, "items": page_items, + "source": "stored", + "truncated": False, } diff --git a/netx_api/biz_state/compare_sql.py b/netx_api/biz_state/compare_sql.py index ae42ed8..b7ba9d3 100644 --- a/netx_api/biz_state/compare_sql.py +++ b/netx_api/biz_state/compare_sql.py @@ -406,6 +406,9 @@ def run_sql_sheet_compare( ) -> dict[str, Any]: """Compare one sheet via PostgreSQL TEMP tables + FULL OUTER JOIN. + Same-key multipath: rows get ``dup_rn`` by ``seq,id`` and join on + ``(rk, dup_rn)`` (ordered zip). Extras are added/removed, not duplicate. + Returns the same ``{summary, diffs, mapping_stats}`` shape as ``compare_rows``. ``on_progress(side, n, *, engine=\"sql\", note=..., phase=...)``. """ @@ -503,9 +506,7 @@ def run_sql_sheet_compare( ), params, ) - db.execute( - text(f"CREATE INDEX IF NOT EXISTS {tname}_rk ON {tname} (rk) WHERE dup_rn = 1") - ) + db.execute(text(f"CREATE INDEX IF NOT EXISTS {tname}_rk ON {tname} (rk, dup_rn)")) before_n = int( db.execute(text(f"SELECT count(*) FROM {tb}")).scalar() or 0 @@ -525,14 +526,14 @@ def run_sql_sheet_compare( END """ - # Aggregate primary match kinds (first-wins keys only) + # Ordered same-key pairing: join on (rk, dup_rn) so 5 vs 2 → 2 compared + 3 removed agg_rows = db.execute( text( f""" SELECT {kind_expr} AS kind, count(*)::bigint AS n - FROM (SELECT * FROM {tb} WHERE dup_rn = 1) b - FULL OUTER JOIN (SELECT * FROM {ta} WHERE dup_rn = 1) a - ON b.rk = a.rk + FROM {tb} b + FULL OUTER JOIN {ta} a + ON b.rk = a.rk AND b.dup_rn = a.dup_rn GROUP BY 1 """ ) @@ -543,13 +544,19 @@ def run_sql_sheet_compare( changed = counts.get("changed", 0) unchanged = counts.get("unchanged", 0) - dup_b = int( - db.execute(text(f"SELECT count(*) FROM {tb} WHERE dup_rn > 1")).scalar() or 0 + # Diagnostic: count of match keys with >1 row (not fail rows) + multi_b = int( + db.execute( + text(f"SELECT count(*) FROM (SELECT rk FROM {tb} GROUP BY rk HAVING count(*) > 1) t") + ).scalar() + or 0 ) - dup_a = int( - db.execute(text(f"SELECT count(*) FROM {ta} WHERE dup_rn > 1")).scalar() or 0 + multi_a = int( + db.execute( + text(f"SELECT count(*) FROM (SELECT rk FROM {ta} GROUP BY rk HAVING count(*) > 1) t") + ).scalar() + or 0 ) - duplicate = dup_b + dup_a # Lazy import — avoid circular import with compare_service from .compare_service import resolve_unchanged_policy @@ -559,7 +566,7 @@ def run_sql_sheet_compare( ) diffs: list[dict[str, Any]] = [] - # Fail + duplicate rows (stream into Python — should be << million) + # Fail rows (stream into Python — should be << million) fail_sql = text( f""" SELECT @@ -569,12 +576,14 @@ def run_sql_sheet_compare( b.data AS before_data, a.data AS after_data, COALESCE(b.rk, a.rk) AS rk - FROM (SELECT * FROM {tb} WHERE dup_rn = 1) b - FULL OUTER JOIN (SELECT * FROM {ta} WHERE dup_rn = 1) a - ON b.rk = a.rk + FROM {tb} b + FULL OUTER JOIN {ta} a + ON b.rk = a.rk AND b.dup_rn = a.dup_rn WHERE {kind_expr} IN ('added', 'removed', 'changed') """ ) + fail_n = 0 + _prog("after", after_n, phase="sql_fail_fetch", note="fetch_fails") for row in db.execute(fail_sql).mappings(): kind = str(row["kind"] or "") before = dict(row["before_data"] or {}) if row["before_data"] is not None else None @@ -593,54 +602,77 @@ def run_sql_sheet_compare( if kind == "changed": item["changes"] = _changes_from_rows(before, after, compare_fields, rules) diffs.append(item) - - # Duplicate extras - for side, tname in (("before", tb), ("after", ta)): - q = text( - f""" - SELECT id, data, rk FROM {tname} WHERE dup_rn > 1 - """ - ) - for row in db.execute(q).mappings(): - data = dict(row["data"] or {}) - diffs.append( - { - "kind": "duplicate", - "side": side, - "key": _key_obj_from_data(data, key_fields), - "before": data if side == "before" else None, - "after": data if side == "after" else None, - "mapped_before": data if side == "before" else None, - "changes": {}, - "before_row_id": str(row["id"] or "") if side == "before" else "", - "after_row_id": str(row["id"] or "") if side == "after" else "", - } - ) + fail_n += 1 + if fail_n == 1 or fail_n % 25_000 == 0: + _prog("fail", fail_n, phase="sql_fail_fetch", note="fetch_fails") + if fail_n: + _prog("fail", fail_n, phase="sql_fail_fetch", note="fetch_fails") unchanged_listed = 0 include_u = bool(policy.get("include")) compact = bool(policy.get("compact")) limit_n = policy.get("limit") if include_u and unchanged > 0: - lim_sql = "" params_u: dict[str, Any] = {} + # Stratified sample: round-robin by neighbor|direction|afi|vrf so one + # BGP command cannot consume the whole success sample. + stratum_expr = """ + concat_ws('|', + coalesce(nullif(trim(both from coalesce(b.data->>'neighbor','')), ''), '_'), + coalesce(nullif(trim(both from coalesce(b.data->>'direction','')), ''), '_'), + coalesce(nullif(trim(both from coalesce(b.data->>'afi','')), ''), '_'), + coalesce(nullif(trim(both from coalesce(b.data->>'vrf','')), ''), '_') + ) + """ if limit_n is not None: - lim_sql = " LIMIT :lim" params_u["lim"] = max(0, int(limit_n)) - u_sql = text( - f""" - SELECT - b.id AS before_row_id, - a.id AS after_row_id, - b.data AS before_data, - a.data AS after_data - FROM (SELECT * FROM {tb} WHERE dup_rn = 1) b - INNER JOIN (SELECT * FROM {ta} WHERE dup_rn = 1) a - ON b.rk = a.rk - WHERE NOT ({changed_pred}) - {lim_sql} - """ - ) + u_sql = text( + f""" + WITH u AS ( + SELECT + b.id AS before_row_id, + a.id AS after_row_id, + b.data AS before_data, + a.data AS after_data, + {stratum_expr} AS stratum + FROM {tb} b + INNER JOIN {ta} a + ON b.rk = a.rk AND b.dup_rn = a.dup_rn + WHERE NOT ({changed_pred}) + ), + ranked AS ( + SELECT *, + row_number() OVER ( + PARTITION BY stratum ORDER BY before_row_id ASC + ) AS rn_in + FROM u + ), + picked AS ( + SELECT *, + row_number() OVER ( + ORDER BY rn_in ASC, stratum ASC, before_row_id ASC + ) AS pick_ord + FROM ranked + ) + SELECT before_row_id, after_row_id, before_data, after_data + FROM picked + WHERE pick_ord <= :lim + """ + ) + else: + u_sql = text( + f""" + SELECT + b.id AS before_row_id, + a.id AS after_row_id, + b.data AS before_data, + a.data AS after_data + FROM {tb} b + INNER JOIN {ta} a + ON b.rk = a.rk AND b.dup_rn = a.dup_rn + WHERE NOT ({changed_pred}) + """ + ) for row in db.execute(u_sql, params_u).mappings(): before = dict(row["before_data"] or {}) after = dict(row["after_data"] or {}) @@ -677,7 +709,6 @@ def run_sql_sheet_compare( db.execute(text(f"DROP TABLE IF EXISTS {tb}")) db.execute(text(f"DROP TABLE IF EXISTS {ta}")) - dup_key_list: list[str] = [] summary = { "before_count": before_n, "after_count": after_n, @@ -685,16 +716,21 @@ def run_sql_sheet_compare( "removed": removed, "changed": changed, "unchanged": unchanged, - "duplicate": duplicate, + "duplicate": 0, "match_key_fields": list(key_fields), - "duplicate_keys_before": dup_b, - "duplicate_keys_after": dup_a, - "duplicate_key_list": dup_key_list, + "duplicate_keys_before": multi_b, + "duplicate_keys_after": multi_a, + "duplicate_key_list": [], "unchanged_listed": unchanged_listed, "unchanged_truncated": bool( include_u and limit_n is not None and unchanged > unchanged_listed ), "unchanged_compact": bool(compact and unchanged_listed > 0), + "unchanged_sample_mode": ( + "stratified" + if include_u and limit_n is not None and unchanged > unchanged_listed + else ("full" if include_u else "none") + ), "engine": "sql", "before_raw_count": raw_b, "after_raw_count": raw_a, diff --git a/tests/test_biz_state_compare.py b/tests/test_biz_state_compare.py index 79a5f86..dcb07e1 100644 --- a/tests/test_biz_state_compare.py +++ b/tests/test_biz_state_compare.py @@ -5,7 +5,8 @@ from __future__ import annotations import copy import unittest -from netx_api.biz_state.compare_engine import compare_rows, mapping_stats +from netx_api.biz_state.compare_engine import compare_rows, mapping_stats, stratify_take +from netx_api.biz_state.compare_service import _kind_allows from netx_api.biz_state.compare_rules import ( apply_row_filters, arp_dynamic_row_filters, @@ -198,12 +199,70 @@ class CompareEngineTests(unittest.TestCase): self.assertEqual(out["summary"]["unchanged_listed"], 2) self.assertTrue(out["summary"]["unchanged_truncated"]) self.assertTrue(out["summary"]["unchanged_compact"]) + self.assertEqual(out["summary"]["unchanged_sample_mode"], "stratified") self.assertEqual(len(out["diffs"]), 2) self.assertEqual(out["diffs"][0]["before"], {}) self.assertEqual(out["diffs"][0]["after"], {}) self.assertEqual(out["diffs"][0]["before_row_id"], "b0") self.assertEqual(out["diffs"][0]["after_row_id"], "a0") + def test_stratify_take_round_robin(self) -> None: + items = [("a", i) for i in range(5)] + [("b", i) for i in range(5)] + got = stratify_take(items, 4, key_fn=lambda x: x[0]) + self.assertEqual([x[0] for x in got], ["a", "b", "a", "b"]) + + def test_kind_allows_tabs(self) -> None: + self.assertTrue(_kind_allows("all", "unchanged")) + self.assertTrue(_kind_allows("diff", "removed")) + self.assertTrue(_kind_allows("diff", "changed")) + self.assertFalse(_kind_allows("diff", "added")) + self.assertFalse(_kind_allows("diff", "unchanged")) + self.assertTrue(_kind_allows("unchanged", "unchanged")) + self.assertFalse(_kind_allows("unchanged", "removed")) + self.assertTrue(_kind_allows("removed", "removed")) + self.assertTrue(_kind_allows("changed", "changed")) + + def test_unchanged_sample_stratified_across_neighbors(self) -> None: + """Sample must cover multiple neighbors, not only the first command's rows.""" + before = [] + after = [] + for n_i, neigh in enumerate(("1.1.1.1", "2.2.2.2", "3.3.3.3")): + for j in range(10): + net = f"10.{n_i}.{j}.0/24" + before.append( + { + "neighbor": neigh, + "direction": "in", + "network": net, + "v": "1", + "_netx": {"row_id": f"b-{neigh}-{j}"}, + } + ) + after.append( + { + "neighbor": neigh, + "direction": "in", + "network": net, + "v": "1", + "_netx": {"row_id": f"a-{neigh}-{j}"}, + } + ) + out = compare_rows( + before_rows=before, + after_rows=after, + key_fields=["neighbor", "direction", "network"], + iface_fields=[], + compare_fields=["v"], + port_map={}, + include_unchanged=True, + unchanged_limit=6, + compact_unchanged=True, + ) + self.assertEqual(out["summary"]["unchanged"], 30) + self.assertEqual(out["summary"]["unchanged_listed"], 6) + keys = [d["key"]["neighbor"] for d in out["diffs"]] + self.assertEqual(set(keys), {"1.1.1.1", "2.2.2.2", "3.3.3.3"}) + def test_mac_normalize_via_field_rules(self) -> None: before = [{"ip": "1.1.1.1", "mac": "00:11:22:33:44:55", "iface": "gei-0/1"}] after = [{"ip": "1.1.1.1", "mac": "0011.2233.4455", "iface": "gei-0/1"}] @@ -475,6 +534,13 @@ class CompareSheetDefaultsTests(unittest.TestCase): for r in (route4.get("field_rules") or []) ) ) + route_vpn = next(s for s in sheets if sheet_key(s) == "bgp_route.vpnv4") + self.assertTrue( + any( + f.get("field") == "afi" and f.get("value") == "vpnv4" + for f in (route_vpn.get("row_filters") or []) + ) + ) vpnv4 = next(s for s in sheets if sheet_key(s) == "bgp_peer.vpnv4") self.assertEqual( vpnv4["row_filters"], @@ -483,7 +549,8 @@ class CompareSheetDefaultsTests(unittest.TestCase): isis4 = next(s for s in sheets if sheet_key(s) == "isis_adjacency.ipv4") self.assertEqual(isis4["row_filters"][0]["op"], "contains") - def test_duplicate_match_keys_are_reported(self) -> None: + def test_ordered_same_key_pairing_2v2(self) -> None: + """Same key, two rows each: zip by order (not first-wins + duplicate).""" before = [ {"local_if": "a", "remote_sys": "X", "remote_if": "1", "remote_ip": "1"}, {"local_if": "a", "remote_sys": "X", "remote_if": "1", "remote_ip": "9"}, @@ -500,15 +567,47 @@ class CompareSheetDefaultsTests(unittest.TestCase): compare_fields=["remote_ip"], port_map={"a": "a"}, ) + self.assertEqual(out["summary"]["before_count"], 2) + self.assertEqual(out["summary"]["after_count"], 2) + self.assertEqual(out["summary"]["duplicate"], 0) self.assertEqual(out["summary"]["duplicate_keys_before"], 1) self.assertEqual(out["summary"]["duplicate_keys_after"], 1) - self.assertEqual(out["summary"]["duplicate"], 2) self.assertIn("a|X|1", out["summary"]["duplicate_key_list"]) kinds = [d["kind"] for d in out["diffs"]] - self.assertEqual(kinds.count("duplicate"), 2) - # First before wins → matches first after → unchanged (same remote_ip) + self.assertNotIn("duplicate", kinds) + # Pair 0: 1==1 unchanged; pair 1: 9!=2 changed self.assertEqual(out["summary"]["unchanged"], 1) + self.assertEqual(out["summary"]["changed"], 1) + self.assertEqual(out["summary"]["added"], 0) + self.assertEqual(out["summary"]["removed"], 0) + + def test_ordered_multipath_5_vs_2(self) -> None: + """5 before / 2 after same key → 2 compared + 3 removed failures.""" + key = {"local_if": "a", "remote_sys": "X", "remote_if": "1"} + before = [{**key, "remote_ip": str(i)} for i in range(5)] + after = [{**key, "remote_ip": "0"}, {**key, "remote_ip": "1"}] + out = compare_rows( + before_rows=before, + after_rows=after, + key_fields=["local_if", "remote_sys", "remote_if"], + iface_fields=["local_if"], + compare_fields=["remote_ip"], + port_map={"a": "a"}, + ) + self.assertEqual(out["summary"]["before_count"], 5) + self.assertEqual(out["summary"]["after_count"], 2) + self.assertEqual(out["summary"]["removed"], 3) + self.assertEqual(out["summary"]["added"], 0) + self.assertEqual( + out["summary"]["unchanged"] + out["summary"]["changed"], + 2, + ) + self.assertEqual(out["summary"]["unchanged"], 2) self.assertEqual(out["summary"]["changed"], 0) + self.assertEqual(out["summary"]["duplicate"], 0) + kinds = [d["kind"] for d in out["diffs"]] + self.assertEqual(kinds.count("removed"), 3) + self.assertNotIn("duplicate", kinds) def test_ignore_port_changes_false_keeps_iface(self) -> None: before = [ diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 1ee7ef0..1367d79 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -502,20 +502,23 @@ const en = { runProgress: "{{phase}} · sheet {{i}}/{{n}} · {{sheet}} · {{s}}s elapsed", runRowsLoaded: "Loaded {{side}} {{n}} rows", runRowsSqlCount: "DB count {{side}} {{n}} rows", + runPersisting: "Wrote {{done}}/{{total}} rows", runEngineSql: "engine SQL", runEnginePython: "engine Python", runEngineNote: "reason {{note}}", sheetPending: "Pending", ranWithDuration: "Compare finished ({{s}}s)", unchangedNotStored: "Success rows were counted but not stored. Set “Store success rows” to sample and re-run for spot-check.", - unchangedSampleHint: "{{total}} success rows total; browsing a sample of {{listed}} (hydrated from source tables for cutover spot-check). Pass rate uses all {{total}}.", + unchangedSampleHint: "{{total}} success rows total; browsing a stratified sample of {{listed}}. Search any route for live source lookup. Pass rate uses all {{total}}.", + liveSearchHint: "Search: live verdict from source batches (tabs only filter kind). Works for fail, success, and added.", + liveSearchTruncatedHint: "Search results truncated — narrow the query (full prefix / neighbor).", storeUnchanged: "Store success rows", storeUnchangedAuto: "Auto (full if small / sample 5k + hydrate if large)", storeUnchangedSample: "Sample only (up to 5k keys + hydrate)", storeUnchangedKeys: "All success keys (full browse; large write still slow)", storeUnchangedAlways: "Always full JSON (slow on huge sheets)", storeUnchangedNever: "Never (counts only)", - storeUnchangedHint: "Cutover focuses on fails and pass rate; success defaults to a sample. Pass rate always uses the full success count.", + storeUnchangedHint: "Cutover focuses on fails and pass rate; success defaults to a stratified sample. Type a filter to look up any route from source tables. Pass rate uses the full success count.", saveJob: "Save config", jobSaved: "Job config saved", jobDeleted: "Compare job deleted", @@ -542,7 +545,7 @@ const en = { pickRun: "Select run…", runCount: "{{n}} runs", resultEmpty: "No matching diff rows", - resultFilterPh: "Filter key / values…", + resultFilterPh: "Search prefix / neighbor / RD… (live source lookup)", kindAll: "All", kindDiff: "Fail", kindAdded: "Added", @@ -574,7 +577,7 @@ const en = { filterFailField: "Failed compare field", filterFailFieldAll: "Any failed field", resultEmpty: "No matching rows", - resultFilterPh: "Filter identity / values…", + resultFilterPh: "Search prefix / neighbor / RD… (live source lookup)", diffCount: "{{n}} failed", passRate: "Pass rate", passOk: "All passed", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index f6f364a..bd7c3e9 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -501,20 +501,23 @@ const zh = { runProgress: "{{phase}} · 表 {{i}}/{{n}} · {{sheet}} · 已用 {{s}}s", runRowsLoaded: "已加载 {{side}} {{n}} 行", runRowsSqlCount: "库内统计 {{side}} {{n}} 行", + runPersisting: "已写入 {{done}}/{{total}} 行", runEngineSql: "引擎 SQL", runEnginePython: "引擎 Python", runEngineNote: "原因 {{note}}", sheetPending: "等待中", ranWithDuration: "比对完成(耗时 {{s}} 秒)", unchangedNotStored: "成功行仅统计数量未落库。可在任务配置将「成功行保存」改为抽样后重新比对(抽查用)。", - unchangedSampleHint: "成功共 {{total}} 条,明细抽样 {{listed}} 条(从原表补全显示,割接抽查用)。通过率按全部 {{total}} 计。", + unchangedSampleHint: "成功共 {{total}} 条,明细抽样 {{listed}} 条(分层抽查)。要查任意路由请输入筛选条件——将按原表即时判定。通过率按全部 {{total}} 计。", + liveSearchHint: "搜索:原表即时判定(页签只过滤种类)。失败/成功/新增均可查到。", + liveSearchTruncatedHint: "搜索结果已截断,请收窄条件(如完整前缀 / neighbor)。", storeUnchanged: "成功行保存", storeUnchangedAuto: "自动(小表全量 / 大表抽样 5000+原表补全)", storeUnchangedSample: "仅抽样(最多 5000,身份键+原表补全)", storeUnchangedKeys: "全量身份键(可翻完全部成功,大表写入仍较久)", storeUnchangedAlways: "始终全量 JSON(大表很慢,慎用)", storeUnchangedNever: "不保存(只计数量)", - storeUnchangedHint: "割接优先看失败与通过率;成功默认抽样抽查。通过率始终按全部成功计数。", + storeUnchangedHint: "割接优先看失败与通过率;成功默认分层抽样。输入筛选可回查原表(不依赖抽样)。通过率按全部成功计数。", saveJob: "保存配置", jobSaved: "任务配置已保存", jobDeleted: "比对任务已删除", @@ -541,7 +544,7 @@ const zh = { pickRun: "选择比对记录…", runCount: "{{n}} 次", resultEmpty: "无匹配失败/结果行", - resultFilterPh: "筛选身份 / 字段值…", + resultFilterPh: "搜索前缀 / neighbor / RD…(回查原表)", kindAll: "全部", kindDiff: "失败", kindAdded: "新增", diff --git a/web/src/pages/network/BizComparePage.tsx b/web/src/pages/network/BizComparePage.tsx index fdd13d3..af63c82 100644 --- a/web/src/pages/network/BizComparePage.tsx +++ b/web/src/pages/network/BizComparePage.tsx @@ -38,6 +38,7 @@ import { jobChipColor, NmStatusChip } from "./nmChips"; type PageTab = "templates" | "jobs"; type JobDetailTab = "config" | "runs" | "result"; type KindFilter = "diff" | "all" | "added" | "removed" | "changed" | "unchanged"; +type DiffListSource = "stored" | "live" | ""; type CreateJobStep = 0 | 1 | 2 | 3; const CREATE_JOB_STEPS = 4; @@ -847,6 +848,8 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage const [resultTotal, setResultTotal] = useState(0); const [pagedDiffs, setPagedDiffs] = useState([]); const [diffsLoading, setDiffsLoading] = useState(false); + const [diffsSource, setDiffsSource] = useState(""); + const [diffsTruncated, setDiffsTruncated] = useState(false); const boardRef = useRef(null); const tableScrollRef = useRef(null); const tableScrollPosRef = useRef({ top: 0, left: 0 }); @@ -1069,6 +1072,8 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage if (!runId || !mid || jobDetailTab !== "result" || sheetStillRunning) { setPagedDiffs([]); setResultTotal(0); + setDiffsSource(""); + setDiffsTruncated(false); return; } let cancelled = false; @@ -1086,6 +1091,10 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage if (cancelled) return; setPagedDiffs((res.items || []) as DiffRow[]); setResultTotal(Number(res.total || 0)); + setDiffsSource( + String((res as any).source || "").toLowerCase() === "live" ? "live" : "stored", + ); + setDiffsTruncated(Boolean((res as any).truncated)); const pages = Math.max( 1, Math.ceil(Number(res.total || 0) / Number(res.page_size || resultPageSize)), @@ -1257,8 +1266,15 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage String(activeSheetCard?.status || ""), ); const activeAdded = Number(activeSheetCard?.added || 0); + const activeRemoved = Number(activeSheetCard?.removed || 0); + const activeChanged = Number(activeSheetCard?.changed || 0); const showFailCol = - kindFilter === "diff" || kindFilter === "all" || kindFilter === "added"; + kindFilter === "diff" || + kindFilter === "all" || + kindFilter === "added" || + kindFilter === "removed" || + kindFilter === "changed"; + const isLiveSearch = diffsSource === "live" && Boolean(debouncedResultKw.trim()); const resultEmptyColSpan = 1 + (showFailCol ? 1 : 0) + @@ -1895,6 +1911,10 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage rows_loaded?: number; engine?: string; engine_note?: string; + persisted?: number; + persist_total?: number; + fail_rows?: number; + ok_rows?: number; }; const runEngine = String(runProgress.engine || "").toLowerCase(); @@ -3292,6 +3312,15 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage })} ) : null} + {String(runProgress.phase || "").startsWith("persisting") && + Number(runProgress.persist_total || 0) > 0 ? ( + + {t("bizCompare.runPersisting", { + done: String(runProgress.persisted || 0), + total: String(runProgress.persist_total || 0), + })} + + ) : null} {runDetail.message ? ( {String(runDetail.message)} ) : null} @@ -3491,6 +3520,8 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage {( [ ["diff", activeFail, "diff"], + ["removed", activeRemoved, "removed"], + ["changed", activeChanged, "changed"], ["added", activeAdded, "added"], ["unchanged", activeSuccess, "unchanged"], ["all", null, "all"], @@ -3510,7 +3541,11 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage ? t("bizCompare.kindAll") : id === "added" ? t("bizCompare.kindAddedShort") - : t("bizCompare.kindSuccess")} + : id === "removed" + ? t("bizCompare.kindRemovedShort") + : id === "changed" + ? t("bizCompare.kindChangedShort") + : t("bizCompare.kindSuccess")} {n !== null ? ( <> {" "} @@ -3530,11 +3565,31 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage {diffsLoading ? "…" : `${pagedDiffs.length}/${resultTotal}`} - {kindFilter === "unchanged" && + {isLiveSearch ? ( +

+ {diffsTruncated + ? t("bizCompare.liveSearchTruncatedHint") + : t("bizCompare.liveSearchHint")} +

+ ) : null} + {!isLiveSearch && + kindFilter === "unchanged" && (runDetail?.summary?.unchanged_truncated || - (Number(runDetail?.summary?.unchanged || 0) > - Number(runDetail?.summary?.unchanged_listed || 0) && - Number(runDetail?.summary?.unchanged_listed || 0) > 0)) ? ( + (Number( + (activeRunSheet?.summary as any)?.unchanged || + runDetail?.summary?.unchanged || + 0, + ) > + Number( + (activeRunSheet?.summary as any)?.unchanged_listed ?? + runDetail?.summary?.unchanged_listed ?? + 0, + ) && + Number( + (activeRunSheet?.summary as any)?.unchanged_listed ?? + runDetail?.summary?.unchanged_listed ?? + 0, + ) > 0)) ? (

{t("bizCompare.unchangedSampleHint", { listed: String( diff --git a/web/src/services/api.ts b/web/src/services/api.ts index d3b22f8..6facb74 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -2185,6 +2185,10 @@ export const bizCompareListRunDiffs = (params: { page_size: number; metric_id: string; items: Record[]; + source?: string; + truncated?: boolean; + live_before_matched?: number; + live_after_matched?: number; }>(`/v1/biz-state/compare/runs/${encodeURIComponent(params.runId)}/diffs?${p.toString()}`); };