diff --git a/netx_api/biz_state/compare_engine.py b/netx_api/biz_state/compare_engine.py
index b5e1251..ddf0051 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
from typing import Any, Mapping, Sequence
from .compare_rules import field_rule_map, values_equal, explain_diff
@@ -111,9 +112,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 +135,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 +183,21 @@ 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, ...]] = []
def _key_obj(row: dict[str, Any]) -> dict[str, Any]:
return {f: row.get(f, "") for f in key_fields}
@@ -246,96 +247,90 @@ 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,
- }
- )
- 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": {},
- }
- )
+ 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
+ _emit_unchanged(orig, mapped, after)
def _fmt_keys(keys: list[tuple[str, ...]]) -> list[str]:
seen: set[str] = set()
@@ -345,7 +340,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,11 +358,11 @@ 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
diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py
index d6d1553..39ad3af 100644
--- a/netx_api/biz_state/compare_service.py
+++ b/netx_api/biz_state/compare_service.py
@@ -186,6 +186,10 @@ 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
# Success-row persist policy (see resolve_unchanged_policy)
_STORE_UNCHANGED_MODES = frozenset({"auto", "always", "never", "sample", "keys"})
_UNCHANGED_FULL_MAX = 20_000
@@ -273,15 +277,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 +325,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 +357,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
@@ -2423,6 +2463,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 +2504,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 +2529,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 +2541,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)
diff --git a/netx_api/biz_state/compare_sql.py b/netx_api/biz_state/compare_sql.py
index ae42ed8..5c0ae18 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,29 +602,11 @@ 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"))
@@ -634,9 +625,9 @@ def run_sql_sheet_compare(
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
+ FROM {tb} b
+ INNER JOIN {ta} a
+ ON b.rk = a.rk AND b.dup_rn = a.dup_rn
WHERE NOT ({changed_pred})
{lim_sql}
"""
@@ -677,7 +668,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,11 +675,11 @@ 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
diff --git a/tests/test_biz_state_compare.py b/tests/test_biz_state_compare.py
index 79a5f86..ed684af 100644
--- a/tests/test_biz_state_compare.py
+++ b/tests/test_biz_state_compare.py
@@ -475,6 +475,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 +490,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 +508,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..32b302d 100644
--- a/web/src/i18n/en.ts
+++ b/web/src/i18n/en.ts
@@ -502,6 +502,7 @@ 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}}",
diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts
index f6f364a..930fb5a 100644
--- a/web/src/i18n/zh.ts
+++ b/web/src/i18n/zh.ts
@@ -501,6 +501,7 @@ const zh = {
runProgress: "{{phase}} · 表 {{i}}/{{n}} · {{sheet}} · 已用 {{s}}s",
runRowsLoaded: "已加载 {{side}} {{n}} 行",
runRowsSqlCount: "库内统计 {{side}} {{n}} 行",
+ runPersisting: "已写入 {{done}}/{{total}} 行",
runEngineSql: "引擎 SQL",
runEnginePython: "引擎 Python",
runEngineNote: "原因 {{note}}",
diff --git a/web/src/pages/network/BizComparePage.tsx b/web/src/pages/network/BizComparePage.tsx
index fdd13d3..a93d310 100644
--- a/web/src/pages/network/BizComparePage.tsx
+++ b/web/src/pages/network/BizComparePage.tsx
@@ -1895,6 +1895,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 +3296,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}