diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py index 2f8b2f5..57bcbe5 100644 --- a/netx_api/biz_state/compare_service.py +++ b/netx_api/biz_state/compare_service.py @@ -1762,7 +1762,7 @@ def _run_sheet( port_map: dict[str, str], iface_normalize_rules: list[dict[str, str]] | None = None, store_unchanged: str = "auto", - on_load_progress: Callable[[str, int], None] | None = None, + on_load_progress: Callable[..., None] | None = None, ) -> dict[str, Any]: key_fields = list(sheet.get("key_fields") or []) iface_fields = list(sheet.get("iface_fields") or []) @@ -1785,28 +1785,32 @@ def _run_sheet( ignore_ports = bool(ignore_ports) mid = sheet["metric_id"] - # PostgreSQL path: pushdown-safe sheets join in-DB (BGP-scale). - from .compare_sql import can_sql_compare, run_sql_sheet_compare + def _emit_load(side: str, n: int, **meta: Any) -> None: + if not on_load_progress: + return + try: + on_load_progress(side, n, **meta) + except TypeError: + on_load_progress(side, n) - if can_sql_compare( + # PostgreSQL path: pushdown-safe sheets join in-DB (BGP-scale). + from .compare_sql import run_sql_sheet_compare, sql_compare_skip_reason + + skip_reason = sql_compare_skip_reason( db, sheet, port_map=port_map, iface_normalize_rules=iface_normalize_rules, - ): + ) + if not skip_reason: try: - - def _sql_progress(side: str, n: int) -> None: - if on_load_progress: - on_load_progress(side, n) - result = run_sql_sheet_compare( db, sheet=sheet, before_batch_id=before_batch_id, after_batch_id=after_batch_id, store_unchanged=store_unchanged, - on_progress=_sql_progress, + on_progress=_emit_load, ) summary = dict(result["summary"]) # Engine already sets raw counts / policy; keep keys stable @@ -1836,14 +1840,22 @@ def _run_sheet( sheet_key(sheet), mid, ) + skip_reason = "sql_error_fallback" + else: + _log.info( + "python compare sheet=%s metric=%s skip_sql=%s", + sheet_key(sheet), + mid, + skip_reason, + ) + + _emit_load("before", 0, engine="python", note=skip_reason or "python", phase="loading") def _before_chunk(n: int) -> None: - if on_load_progress: - on_load_progress("before", n) + _emit_load("before", n, engine="python", note=skip_reason or "python", phase="loading") def _after_chunk(n: int) -> None: - if on_load_progress: - on_load_progress("after", n) + _emit_load("after", n, engine="python", note=skip_reason or "python", phase="loading") before_raw = _load_metric_rows( db, batch_id=before_batch_id, metric_id=mid, on_chunk=_before_chunk @@ -2106,27 +2118,43 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: break _publish_sheets() - _load_pub = {"t": 0.0, "n": -1} + _load_pub = {"t": 0.0, "n": -1, "engine": ""} - def _on_load(side: str, n: int, _idx: int = idx, _sheet: dict = sheet) -> None: + def _on_load( + side: str, + n: int, + *, + engine: str = "python", + note: str = "", + phase: str | None = None, + _idx: int = idx, + _sheet: dict = sheet, + ) -> None: now = time.monotonic() - # Avoid committing every 5k on million-row BGP loads - if n - _load_pub["n"] < 25_000 and now - _load_pub["t"] < 2.0: - return + eng = str(engine or "python") + # SQL emits sparse updates; Python still throttle chunk spam + if eng != "sql": + if n - _load_pub["n"] < 25_000 and now - _load_pub["t"] < 2.0: + return _load_pub["t"] = now _load_pub["n"] = n + _load_pub["engine"] = eng + extra: dict[str, Any] = { + "load_side": side, + "rows_loaded": n, + "engine": eng, + } + if note: + extra["engine_note"] = str(note)[:128] _set_run_progress( db, run, - phase="loading", + phase=str(phase or ("loading" if eng != "sql" else "sql_count")), sheet_index=_idx, sheet_total=total, sheet=_sheet, started_mono=started_mono, - extra={ - "load_side": side, - "rows_loaded": n, - }, + extra=extra, ) _set_run_progress( @@ -2137,6 +2165,7 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: sheet_total=total, sheet=sheet, started_mono=started_mono, + extra={"engine": "", "engine_note": ""}, ) one = _run_sheet( db, diff --git a/netx_api/biz_state/compare_sql.py b/netx_api/biz_state/compare_sql.py index e66475b..7df648b 100644 --- a/netx_api/biz_state/compare_sql.py +++ b/netx_api/biz_state/compare_sql.py @@ -119,6 +119,53 @@ def _field_rules_sql_compatible(rules: Sequence[Mapping[str, Any]] | None) -> bo return True +def sql_compare_skip_reason( + db: Session, + sheet: Mapping[str, Any], + *, + port_map: Mapping[str, str] | None = None, + iface_normalize_rules: Sequence[Mapping[str, str]] | None = None, +) -> str: + """Empty string when SQL is allowed; otherwise a short reason for progress/logs. + + Simple ``row_filters`` (eq/in/contains/…) used for BGP sheet splits are fine — + they do **not** force Python by themselves. + """ + if not _dialect_is_postgres(db): + return "not_postgres" + if port_map: + return "port_map" + key_fields = [str(k).strip() for k in (sheet.get("key_fields") or []) if str(k).strip()] + if not key_fields: + return "no_key_fields" + if any(not _FIELD_RE.match(k) for k in key_fields): + return "unsafe_key_field" + iface_fields = [str(f).strip() for f in (sheet.get("iface_fields") or []) if str(f).strip()] + iface_set = set(iface_fields) + ignore_ports = sheet.get("ignore_port_changes") + if ignore_ports is True: + return "ignore_port_changes" + if ignore_ports is None and iface_set and any(k in iface_set for k in key_fields): + return "auto_ignore_ports" + field_rules = list(sheet.get("field_rules") or []) + if not _field_rules_sql_compatible(field_rules): + return "field_rules" + compare_fields = effective_compare_fields( + list(sheet.get("compare_fields") or []), + field_rules, + ) + if any(not _FIELD_RE.match(str(f).strip()) for f in compare_fields if str(f).strip()): + return "unsafe_compare_field" + used = set(key_fields) | {str(f).strip() for f in compare_fields if str(f).strip()} + if iface_normalize_rules and (used & iface_set): + return "iface_normalize" + if not _filters_sql_compatible(list(sheet.get("row_filters") or [])): + return "row_filters" + if not str(sheet.get("metric_id") or "").strip(): + return "no_metric_id" + return "" + + def can_sql_compare( db: Session, sheet: Mapping[str, Any], @@ -127,39 +174,14 @@ def can_sql_compare( iface_normalize_rules: Sequence[Mapping[str, str]] | None = None, ) -> bool: """True when this sheet can run entirely as a PostgreSQL JOIN.""" - if not _dialect_is_postgres(db): - return False - if port_map: - return False - key_fields = [str(k).strip() for k in (sheet.get("key_fields") or []) if str(k).strip()] - if not key_fields or any(not _FIELD_RE.match(k) for k in key_fields): - return False - iface_fields = [str(f).strip() for f in (sheet.get("iface_fields") or []) if str(f).strip()] - iface_set = set(iface_fields) - # Port-rename heuristic / map rewrite cannot be expressed here - ignore_ports = sheet.get("ignore_port_changes") - if ignore_ports is True: - return False - if ignore_ports is None and iface_set and any(k in iface_set for k in key_fields): - # Auto path may drop iface from match key — stay on Python - return False - field_rules = list(sheet.get("field_rules") or []) - if not _field_rules_sql_compatible(field_rules): - return False - compare_fields = effective_compare_fields( - list(sheet.get("compare_fields") or []), - field_rules, + return not bool( + sql_compare_skip_reason( + db, + sheet, + port_map=port_map, + iface_normalize_rules=iface_normalize_rules, + ) ) - if any(not _FIELD_RE.match(str(f).strip()) for f in compare_fields if str(f).strip()): - return False - used = set(key_fields) | {str(f).strip() for f in compare_fields if str(f).strip()} - if iface_normalize_rules and (used & iface_set): - return False - if not _filters_sql_compatible(list(sheet.get("row_filters") or [])): - return False - if not str(sheet.get("metric_id") or "").strip(): - return False - return True def compile_row_filters_sql( @@ -321,15 +343,24 @@ def run_sql_sheet_compare( before_batch_id: str, after_batch_id: str, store_unchanged: str = "auto", - on_progress: Callable[[str, int], None] | None = None, + on_progress: Callable[..., None] | None = None, ) -> dict[str, Any]: """Compare one sheet via PostgreSQL TEMP tables + FULL OUTER JOIN. Returns the same ``{summary, diffs, mapping_stats}`` shape as ``compare_rows``. + ``on_progress(side, n, *, engine=\"sql\", note=..., phase=...)``. """ if not _dialect_is_postgres(db): raise RuntimeError("sql_compare_requires_postgres") + def _prog(side: str, n: int, *, phase: str = "sql_count", note: str = "") -> None: + if not on_progress: + return + try: + on_progress(side, n, engine="sql", note=note, phase=phase) + except TypeError: + on_progress(side, n) + key_fields = [str(k).strip() for k in (sheet.get("key_fields") or []) if str(k).strip()] field_rules = list(sheet.get("field_rules") or []) rules = field_rule_map(field_rules) @@ -350,6 +381,8 @@ def run_sql_sheet_compare( tb = f"_netx_cmp_b_{tag}" ta = f"_netx_cmp_a_{tag}" + _prog("before", 0, phase="sql_count", note="count") + # Raw counts (no row_filters) raw_b = int( db.execute( @@ -361,8 +394,7 @@ def run_sql_sheet_compare( ).scalar() or 0 ) - if on_progress: - on_progress("before", raw_b) + _prog("before", raw_b, phase="sql_count") raw_a = int( db.execute( text( @@ -373,14 +405,14 @@ def run_sql_sheet_compare( ).scalar() or 0 ) - if on_progress: - on_progress("after", raw_a) + _prog("after", raw_a, phase="sql_count") base_params = {"bid": bid_b, "mid": mid, **filter_params} # Build TEMP sides - for tname, batch_id in ((tb, bid_b), (ta, bid_a)): + for tname, batch_id, side in ((tb, bid_b, "before"), (ta, bid_a, "after")): db.execute(text(f"DROP TABLE IF EXISTS {tname}")) params = {**base_params, "bid": batch_id} + _prog(side, raw_b if side == "before" else raw_a, phase="sql_project", note="temp") # PRESERVE ROWS: compare progress commits must not drop temps mid-run db.execute( text( @@ -410,8 +442,7 @@ def run_sql_sheet_compare( after_n = int( db.execute(text(f"SELECT count(*) FROM {ta}")).scalar() or 0 ) - if on_progress: - on_progress("after", after_n) + _prog("after", after_n, phase="sql_join", note="join") changed_pred = _changed_predicate(compare_fields, rules) kind_expr = f""" diff --git a/tests/test_biz_state_compare_sql.py b/tests/test_biz_state_compare_sql.py index f885318..e61f83c 100644 --- a/tests/test_biz_state_compare_sql.py +++ b/tests/test_biz_state_compare_sql.py @@ -8,6 +8,7 @@ from unittest.mock import MagicMock from netx_api.biz_state.compare_sql import ( can_sql_compare, compile_row_filters_sql, + sql_compare_skip_reason, _field_rules_sql_compatible, _filters_sql_compatible, _safe_field, @@ -117,6 +118,8 @@ class CompareSqlGateTests(unittest.TestCase): "row_filters": [{"field": "afi", "op": "eq", "value": "ipv4"}], } self.assertTrue(can_sql_compare(_db(), sheet, port_map={})) + # BGP afi/vrf sheet splits via row_filters must not force Python + self.assertEqual(sql_compare_skip_reason(_db(), sheet, port_map={}), "") def test_rejects_iface_normalize_when_key_uses_iface(self) -> None: sheet = { diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 2304298..1ee7ef0 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -501,6 +501,10 @@ const en = { runStatusDoneSec: "Done {{s}}s", runProgress: "{{phase}} · sheet {{i}}/{{n}} · {{sheet}} · {{s}}s elapsed", runRowsLoaded: "Loaded {{side}} {{n}} rows", + runRowsSqlCount: "DB count {{side}} {{n}} 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.", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 2b191f9..f6f364a 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -500,6 +500,10 @@ const zh = { runStatusDoneSec: "完成 {{s}}s", runProgress: "{{phase}} · 表 {{i}}/{{n}} · {{sheet}} · 已用 {{s}}s", runRowsLoaded: "已加载 {{side}} {{n}} 行", + runRowsSqlCount: "库内统计 {{side}} {{n}} 行", + runEngineSql: "引擎 SQL", + runEnginePython: "引擎 Python", + runEngineNote: "原因 {{note}}", sheetPending: "等待中", ranWithDuration: "比对完成(耗时 {{s}} 秒)", unchangedNotStored: "成功行仅统计数量未落库。可在任务配置将「成功行保存」改为抽样后重新比对(抽查用)。", diff --git a/web/src/pages/network/BizComparePage.tsx b/web/src/pages/network/BizComparePage.tsx index ba8c9ca..fdd13d3 100644 --- a/web/src/pages/network/BizComparePage.tsx +++ b/web/src/pages/network/BizComparePage.tsx @@ -1893,7 +1893,10 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage diff_rows?: number; load_side?: string; rows_loaded?: number; + engine?: string; + engine_note?: string; }; + const runEngine = String(runProgress.engine || "").toLowerCase(); // Poll active compare runs so the modal can be closed and reopened safely. useEffect(() => { @@ -3003,7 +3006,10 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage sheet_title?: string; rows_loaded?: number; load_side?: string; + engine?: string; + engine_note?: string; }; + const eng = String(prog.engine || "").toLowerCase(); const stLabel = st === "running" || st === "queued" ? t("bizCompare.runStatusRunning") @@ -3018,15 +3024,25 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage {stLabel} - {active && prog.phase ? ( + {active && (prog.phase || eng) ? (
- {prog.phase} + {eng === "sql" + ? t("bizCompare.runEngineSql") + : eng === "python" + ? t("bizCompare.runEnginePython") + : ""} + {eng && prog.engine_note + ? ` · ${prog.engine_note}` + : ""} + {prog.phase ? `${eng ? " · " : ""}${prog.phase}` : ""} {prog.sheet_total ? ` · ${prog.sheet_index || 0}/${prog.sheet_total}` : ""} {prog.sheet_title ? ` · ${prog.sheet_title}` : ""} {Number(prog.rows_loaded || 0) > 0 - ? ` · ${prog.load_side || ""} ${prog.rows_loaded}` + ? eng === "sql" + ? ` · ${prog.load_side || ""} ${prog.rows_loaded}` + : ` · ${prog.load_side || ""} ${prog.rows_loaded}` : ""}
) : null} @@ -3240,6 +3256,18 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage {runDetail && runIsActive ? (
{t("bizCompare.runStatusRunning")} + {runEngine === "sql" || runEngine === "python" ? ( + + {runEngine === "sql" + ? t("bizCompare.runEngineSql") + : t("bizCompare.runEnginePython")} + {runProgress.engine_note + ? ` · ${t("bizCompare.runEngineNote", { + note: String(runProgress.engine_note), + })}` + : ""} + + ) : null} {t("bizCompare.runProgress", { phase: String(runProgress.phase || "…"), @@ -3253,10 +3281,15 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage {Number(runProgress.rows_loaded || 0) > 0 ? ( - {t("bizCompare.runRowsLoaded", { - side: String(runProgress.load_side || "—"), - n: String(runProgress.rows_loaded || 0), - })} + {runEngine === "sql" + ? t("bizCompare.runRowsSqlCount", { + side: String(runProgress.load_side || "—"), + n: String(runProgress.rows_loaded || 0), + }) + : t("bizCompare.runRowsLoaded", { + side: String(runProgress.load_side || "—"), + n: String(runProgress.rows_loaded || 0), + })} ) : null} {runDetail.message ? (