diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py index d6d1553..57bcbe5 100644 --- a/netx_api/biz_state/compare_service.py +++ b/netx_api/biz_state/compare_service.py @@ -36,7 +36,6 @@ from .compare_rules import ( arp_dynamic_row_filters, effective_compare_fields, effective_display_fields, - row_matches_filter, ) from .iface_normalize import ( apply_iface_normalize_rows, @@ -773,35 +772,7 @@ def _builtin_source_splits() -> dict[str, list[dict[str, Any]]]: } -def _packaged_zte_status_template_path(): - from pathlib import Path - - return Path(__file__).resolve().parent / "data" / "default_zte_status_template.json" - - -def _load_packaged_zte_status_template() -> dict[str, Any]: - """IOH CN migration sheet set shipped as the built-in status default.""" - path = _packaged_zte_status_template_path() - if not path.is_file(): - return {} - try: - return dict(json.loads(path.read_text(encoding="utf-8")) or {}) - except Exception: - _log.exception("failed to load packaged ZTE status template %s", path) - return {} - - def _default_zte_status_sheets() -> list[dict[str, Any]]: - """Built-in status sheets — prefer packaged IOH CN migration rules.""" - raw = _load_packaged_zte_status_template() - out: list[dict[str, Any]] = [] - for item in list(raw.get("metrics") or []): - sheet = _normalize_sheet(item) - if sheet: - out.append(sheet) - if out: - return out - # Fallback if package missing (tests / incomplete install) return [ *_isis_af_sheets(), _default_sheet_for_metric("interface_brief", compare_roles=("state",)), @@ -818,23 +789,6 @@ def _default_zte_status_sheets() -> list[dict[str, Any]]: ] -def _builtin_status_needs_packaged_upgrade(existing: list[dict[str, Any]]) -> bool: - """True when built-in template still lacks filtered BGP route sheets.""" - mids = {str(s.get("metric_id") or "") for s in existing} - if "bgp_route" not in mids and "l2vpn_mac" not in mids: - return True - has_filtered_route = any( - str(s.get("metric_id") or "") == "bgp_route" and list(s.get("row_filters") or []) - for s in existing - ) - if not has_filtered_route: - return True - packaged_keys = {sheet_key(s) for s in _default_zte_status_sheets()} - have_keys = {sheet_key(s) for s in existing} - # Missing several packaged sheet ids → sync to packaged default - return len(packaged_keys - have_keys) >= 3 - - def _default_zte_config_sheets() -> list[dict[str, Any]]: """Config-intent metrics for cutover / intent-vs-intent compare.""" return [ @@ -1101,63 +1055,52 @@ def ensure_default_lldp_template(db: Session) -> BizCompareTemplate: def ensure_default_zte_status_template(db: Session) -> BizCompareTemplate: name = "ZTE status default" row = db.query(BizCompareTemplate).filter(BizCompareTemplate.name == name).one_or_none() - packaged = _load_packaged_zte_status_template() sheets = _default_zte_status_sheets() - note = str( - packaged.get("note") - or "Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP route AF sheets)" - )[:512] - iface_rules = list(packaged.get("iface_normalize_rules") or []) + # Volatile counters must not stay in compare_fields on upgraded installs. + _STRIP_COMPARE: dict[str, frozenset[str]] = { + "interface_detail": frozenset({"input_bps", "output_bps", "in_util", "out_util"}), + "optical_brief": frozenset({"rx_power", "tx_power"}), + "bgp_peer": frozenset({"pfx_rcd"}), + } if row: existing = template_metrics(row) - if _builtin_status_needs_packaged_upgrade(existing) and sheets: - _apply_sheets_to_row(row, sheets) - row.note = note - row.updated_at = _utcnow() - if iface_rules: - _set_template_iface_normalize(row, iface_rules) - elif not template_iface_normalize(row): - # Packaged IOH rules use empty normalize; leave empty when explicit - _set_template_iface_normalize(row, []) - db.commit() - db.refresh(row) - return row - # Incremental patches for already-upgraded installs - changed = False + want = {s["metric_id"] for s in sheets} + have = {s["metric_id"] for s in existing} + changed = bool(want - have) upgraded: list[dict[str, Any]] = [] - by_sid = {sheet_key(s): s for s in sheets} + by_want = {s["metric_id"]: s for s in sheets} for s in existing: cur = dict(s) mid = str(cur.get("metric_id") or "") - sid = sheet_key(cur) if mid == "arp" and not cur.get("row_filters"): - src = by_sid.get(sid) or next( - (x for x in sheets if x.get("metric_id") == "arp"), None - ) - cur["row_filters"] = list( - (src or {}).get("row_filters") or arp_dynamic_row_filters() - ) - if not cur.get("field_rules") and src and src.get("field_rules"): - cur["field_rules"] = list(src["field_rules"]) + cur["row_filters"] = list(by_want.get("arp", {}).get("row_filters") or arp_dynamic_row_filters()) + if not cur.get("field_rules") and by_want.get("arp", {}).get("field_rules"): + cur["field_rules"] = list(by_want["arp"]["field_rules"]) changed = True + strip = _STRIP_COMPARE.get(mid) + if strip: + old_cmp = list(cur.get("compare_fields") or []) + new_cmp = [f for f in old_cmp if f not in strip] + if new_cmp != old_cmp: + cur["compare_fields"] = new_cmp + want_disp = list((by_want.get(mid) or {}).get("display_fields") or []) + if want_disp: + cur["display_fields"] = want_disp + changed = True upgraded.append(_normalize_sheet(cur) or cur) - have_mids = {str(s.get("metric_id") or "") for s in upgraded} - for s in sheets: - if str(s.get("metric_id") or "") not in have_mids: - # Only append wholly missing metrics (e.g. bgp_route family) - if str(s.get("metric_id") or "") == "bgp_route" and "bgp_route" not in have_mids: - upgraded.extend( - [x for x in sheets if x.get("metric_id") == "bgp_route"] - ) - have_mids.add("bgp_route") - changed = True - elif str(s.get("metric_id") or "") not in have_mids: + if want - have: + for s in sheets: + if s["metric_id"] not in have: upgraded.append(s) - have_mids.add(str(s.get("metric_id") or "")) - changed = True + changed = True if changed: _apply_sheets_to_row(row, upgraded if upgraded else sheets) - row.note = note + row.note = "Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP/LLDP)" + row.updated_at = _utcnow() + db.commit() + db.refresh(row) + if not template_iface_normalize(row): + _set_template_iface_normalize(row, default_zte_iface_normalize_rules()) row.updated_at = _utcnow() db.commit() db.refresh(row) @@ -1165,12 +1108,12 @@ def ensure_default_zte_status_template(db: Session) -> BizCompareTemplate: row = BizCompareTemplate( id=uuid4().hex, name=name, - note=note, + note="Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP/LLDP)", created_at=_utcnow(), updated_at=_utcnow(), ) _apply_sheets_to_row(row, sheets) - _set_template_iface_normalize(row, iface_rules) + _set_template_iface_normalize(row, default_zte_iface_normalize_rules()) db.add(row) db.commit() db.refresh(row) @@ -1226,10 +1169,10 @@ def ensure_default_templates(db: Session) -> None: def upgrade_builtin_split_sheets(db: Session) -> None: - """Upgrade built-in ZTE status template to packaged AF / BGP route sheets. + """Split unfiltered whole-table sheets on the built-in ZTE template only. - Custom templates are left alone. Built-in is replaced wholesale when it - still lacks filtered ``bgp_route`` sheets (IOH CN migration default). + A sheet is replaced when its id is still the source metric and it has no + row filters. Custom templates and already-split sheets are left alone. """ row = ( db.query(BizCompareTemplate) @@ -1239,18 +1182,6 @@ def upgrade_builtin_split_sheets(db: Session) -> None: if not row: return existing = template_metrics(row) - packaged_sheets = _default_zte_status_sheets() - if _builtin_status_needs_packaged_upgrade(existing) and packaged_sheets: - packaged = _load_packaged_zte_status_template() - _apply_sheets_to_row(row, packaged_sheets) - row.note = str( - packaged.get("note") - or "Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP route AF sheets)" - )[:512] - _set_template_iface_normalize(row, list(packaged.get("iface_normalize_rules") or [])) - row.updated_at = _utcnow() - db.commit() - return splits = _builtin_source_splits() out: list[dict[str, Any]] = [] changed = False @@ -1502,28 +1433,12 @@ def _load_metric_rows( batch_id: str, metric_id: str, on_chunk: Callable[[int], None] | None = None, - row_filters: list[dict[str, Any]] | None = None, ) -> list[dict[str, Any]]: """Load metric rows in keyset chunks (stable on million-row sheets). - When ``row_filters`` are SQL-pushdown-safe on PostgreSQL, they are applied in - the SELECT (critical for BGP afi/vrf sheet splits — avoids loading 1M+ then - discarding). Otherwise filters are applied in Python after each chunk. + Avoids ORM ``yield_per``/``unique()`` clash and OFFSET degradation on BGP-sized tables. + ``on_chunk(loaded_count)`` is called after each chunk for progress UI. """ - from .compare_sql import ( - _dialect_is_postgres, - _filters_sql_compatible, - compile_row_filters_sql, - ) - - filters = [f for f in (row_filters or []) if isinstance(f, dict)] - pushdown = bool( - filters and _dialect_is_postgres(db) and _filters_sql_compatible(filters) - ) - filter_sql, filter_params = ("TRUE", {}) - if pushdown: - filter_sql, filter_params = compile_row_filters_sql(filters) - if metric_id == "lldp_neighbor": out: list[dict[str, Any]] = [] last_id = "" @@ -1535,104 +1450,39 @@ def _load_metric_rows( if not chunk: break for n in chunk: - row = { - "local_if": n.local_if, - "remote_sys": n.remote_sys, - "remote_if": n.remote_if, - "remote_ip": n.remote_ip, - "protocol": n.protocol, - "_netx": { - "batch_id": batch_id, - "batch_command_id": n.batch_command_id or "", - "task_id": n.task_id or "", - "ne_id": n.ne_id or "", - "collected_at": n.collected_at.isoformat() + "Z" - if n.collected_at - else None, - "row_id": n.id, - }, - } - if filters and not pushdown and not all( - row_matches_filter(row, f) for f in filters - ): - db.expunge(n) - continue - out.append(row) + out.append( + { + "local_if": n.local_if, + "remote_sys": n.remote_sys, + "remote_if": n.remote_if, + "remote_ip": n.remote_ip, + "protocol": n.protocol, + "_netx": { + "batch_id": batch_id, + "batch_command_id": n.batch_command_id or "", + "task_id": n.task_id or "", + "ne_id": n.ne_id or "", + "collected_at": n.collected_at.isoformat() + "Z" + if n.collected_at + else None, + "row_id": n.id, + }, + } + ) db.expunge(n) last_id = str(chunk[-1].id) if on_chunk: on_chunk(len(out)) if len(chunk) < _LOAD_YIELD_PER: break - if filters and not pushdown: - return apply_row_filters(out, filters) return out - - # Generic tabular metrics — PG + pushdown uses SQL keyset with JSON filters + # Generic tabular metrics (ISIS / interface / ARP / ND6 / BGP …) from ..models import BizStateMetricRow - from sqlalchemy import text as sql_text out: list[dict[str, Any]] = [] last_seq = -1 last_id = "" while True: - if pushdown: - params = { - "bid": batch_id, - "mid": metric_id, - "last_seq": last_seq, - "last_id": last_id, - "lim": int(_LOAD_YIELD_PER), - **filter_params, - } - keyset = ( - "(seq > :last_seq OR (seq = :last_seq AND id > :last_id))" - if last_id - else "TRUE" - ) - 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 ({keyset}) - ORDER BY seq ASC, id ASC - LIMIT :lim - """ - ), - params, - ).mappings().all() - if not rows: - break - for r in rows: - data = dict(r["data_json"] or {}) - collected = r["collected_at"] - out.append( - { - **data, - "_netx": { - "batch_id": batch_id, - "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"]), - }, - } - ) - last_seq = int(rows[-1]["seq"] or 0) - last_id = str(rows[-1]["id"]) - if on_chunk: - on_chunk(len(out)) - if len(rows) < _LOAD_YIELD_PER: - break - continue - q = db.query(BizStateMetricRow).filter( BizStateMetricRow.batch_id == batch_id, BizStateMetricRow.metric_id == metric_id, @@ -1655,23 +1505,24 @@ def _load_metric_rows( if not chunk: break for r in chunk: - row = { - **dict(r.data_json or {}), - "_netx": { - "batch_id": batch_id, - "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 - out.append(row) + # Raw rows only — filtering belongs to the compare sheet template + # (``row_filters``), not metric-specific branches here. + # ``_netx`` is collector provenance (stripped before field compare). + out.append( + { + **dict(r.data_json or {}), + "_netx": { + "batch_id": batch_id, + "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, + }, + } + ) db.expunge(r) last_seq = int(chunk[-1].seq or 0) last_id = str(chunk[-1].id) @@ -1943,7 +1794,7 @@ def _run_sheet( on_load_progress(side, n) # PostgreSQL path: pushdown-safe sheets join in-DB (BGP-scale). - from .compare_sql import SqlCompareSkip, run_sql_sheet_compare, sql_compare_skip_reason + from .compare_sql import run_sql_sheet_compare, sql_compare_skip_reason skip_reason = sql_compare_skip_reason( db, @@ -1983,18 +1834,6 @@ def _run_sheet( diffs=list(result.get("diffs") or []), mapping_stats=dict(result.get("mapping_stats") or {}), ) - except SqlCompareSkip as skip: - skip_reason = skip.reason or "skip" - _log.info( - "python compare sheet=%s metric=%s skip_sql=%s", - sheet_key(sheet), - mid, - skip_reason, - ) - try: - db.rollback() - except Exception: - pass except Exception: _log.exception( "sql compare fallback sheet=%s metric=%s — using Python engine", @@ -2002,11 +1841,6 @@ def _run_sheet( mid, ) skip_reason = "sql_error_fallback" - # Roll back aborted SQL transaction so Python path can use the session - try: - db.rollback() - except Exception: - pass else: _log.info( "python compare sheet=%s metric=%s skip_sql=%s", @@ -2024,22 +1858,13 @@ def _run_sheet( _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, - row_filters=row_filters, + db, batch_id=before_batch_id, metric_id=mid, on_chunk=_before_chunk ) after_raw = _load_metric_rows( - db, - batch_id=after_batch_id, - metric_id=mid, - on_chunk=_after_chunk, - row_filters=row_filters, + db, batch_id=after_batch_id, metric_id=mid, on_chunk=_after_chunk ) - # Filters already applied in load when pushdown-safe; keep apply for safety - before_rows = apply_row_filters(before_raw, row_filters) if row_filters else before_raw - after_rows = apply_row_filters(after_raw, row_filters) if row_filters else after_raw + before_rows = apply_row_filters(before_raw, row_filters) + after_rows = apply_row_filters(after_raw, row_filters) policy = resolve_unchanged_policy( store_unchanged, before_n=len(before_rows), after_n=len(after_rows) ) @@ -2177,14 +2002,9 @@ def _set_run_progress( sheet: dict[str, Any] | None, started_mono: float, extra: dict[str, Any] | None = None, - detach: bool = False, ) -> None: - """Update run progress. - - ``detach=True`` writes via a fresh session so SQL compare can keep an open - transaction (TEMP CTAS) without mid-flight commits on the worker ``db``. - """ elapsed_ms = int((time.monotonic() - started_mono) * 1000) + prev = dict(run.summary_json or {}) progress = { "phase": phase, "sheet_index": sheet_index, @@ -2195,39 +2015,6 @@ def _set_run_progress( } if extra: progress.update(extra) - title = progress["sheet_title"] or progress["sheet_id"] or "" - message = ( - f"{phase} {sheet_index}/{sheet_total}" - + (f" · {title}" if title else "") - + f" · {elapsed_ms // 1000}s" - )[:1024] - - if detach: - from ..db import SessionLocal - - s = SessionLocal() - try: - r = s.get(BizCompareRun, str(run.id)) - if not r: - return - prev = dict(r.summary_json or {}) - prev["progress"] = progress - r.summary_json = prev - if str(r.status or "") != "cancelled": - r.status = "running" - r.message = message - s.commit() - # Mirror into worker instance for later in-memory reads (do not commit db) - prev_w = dict(run.summary_json or {}) - prev_w["progress"] = progress - run.summary_json = prev_w - if str(run.status or "") != "cancelled": - run.message = message - finally: - s.close() - return - - prev = dict(run.summary_json or {}) prev["progress"] = progress run.summary_json = prev # Re-read status from DB — cancel may have been committed by another session @@ -2235,7 +2022,12 @@ def _set_run_progress( db.expire(run, ["status", "message"]) if str(run.status or "") != "cancelled": run.status = "running" - run.message = message + title = progress["sheet_title"] or progress["sheet_id"] or "" + run.message = ( + f"{phase} {sheet_index}/{sheet_total}" + + (f" · {title}" if title else "") + + f" · {elapsed_ms // 1000}s" + )[:1024] db.commit() @@ -2354,7 +2146,6 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: } if note: extra["engine_note"] = str(note)[:128] - # SQL path: detach progress commits so TEMP CTAS stays in one txn _set_run_progress( db, run, @@ -2364,7 +2155,6 @@ def _execute_compare_into_run(db: Session, run_id: str) -> dict[str, Any]: sheet=_sheet, started_mono=started_mono, extra=extra, - detach=(eng == "sql"), ) _set_run_progress( diff --git a/netx_api/biz_state/compare_sql.py b/netx_api/biz_state/compare_sql.py index ae42ed8..7df648b 100644 --- a/netx_api/biz_state/compare_sql.py +++ b/netx_api/biz_state/compare_sql.py @@ -18,42 +18,13 @@ from .compare_rules import effective_compare_fields, field_rule_map _log = logging.getLogger("netx.biz_state.compare_sql") -# Below this, Python hash-join is faster and avoids TEMP CTAS / progress races. -_SQL_MIN_ROWS = 20_000 -# Cap a single SQL statement so a stuck planner/lock surfaces as fallback. -_SQL_STATEMENT_TIMEOUT_MS = 180_000 - _FIELD_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$") - - -class SqlCompareSkip(Exception): - """Soft skip — caller should use the Python engine (not an error).""" - - def __init__(self, reason: str) -> None: - super().__init__(reason) - self.reason = str(reason or "skip") _SQL_FILTER_OPS = frozenset( {"eq", "==", "ne", "!=", "in", "not_in", "nin", "contains", "empty", "not_empty", "nonempty", "ci_eq"} ) _SQL_NORMALIZE = frozenset({"", "none", "strip", "lower", "upper", "empty_as_blank"}) -_SQL_COMPARE_MODES = frozenset( - { - "", - "eq", - "ignore", - "skip", - "off", - "numeric", - "number", - "int", - "float", - "percent", - "pct", - "rel", - } -) +_SQL_COMPARE_MODES = frozenset({"", "eq", "ignore", "skip", "off"}) _EMPTY_AS_BLANK = ("n/a", "na", "-", "--", "none", "null") -_NUM_RE_SQL = r"^-?[0-9]+(\.[0-9]+)?([eE][-+]?[0-9]+)?$" def _dialect_is_postgres(db: Session) -> bool: @@ -307,36 +278,6 @@ def _rk_sql(key_fields: list[str], *, json_col: str = "data_json") -> str: return "concat_ws('|', " + ", ".join(parts) + ")" -def _field_differs_sql(bv: str, av: str, rule: Mapping[str, Any]) -> str: - """SQL boolean: True when before/after values differ under the field rule.""" - mode = str(rule.get("compare") or "eq").strip().lower() or "eq" - if mode in ("ignore", "skip", "off") or rule.get("ignore") is True: - return "FALSE" - try: - tol = float(rule.get("tolerance") or 0) - except (TypeError, ValueError): - tol = 0.0 - if mode in ("numeric", "number", "int", "float", "percent", "pct", "rel"): - # Match Python values_equal: parseable → numeric compare; else string eq - num_b = f"(({bv}) ~ '{_NUM_RE_SQL}')" - num_a = f"(({av}) ~ '{_NUM_RE_SQL}')" - bn = f"({bv})::double precision" - an = f"({av})::double precision" - if mode in ("percent", "pct", "rel"): - # differ when relative % > tol (before==0 → after must be 0) - num_diff = ( - f"(CASE WHEN {bn} = 0 THEN {an} IS DISTINCT FROM 0 " - f"ELSE (abs({an} - {bn}) / abs({bn}) * 100.0) > {tol} END)" - ) - else: - num_diff = f"(abs({an} - {bn}) > {tol})" - return ( - f"(CASE WHEN {num_b} AND {num_a} THEN {num_diff} " - f"ELSE ({bv} IS DISTINCT FROM {av}) END)" - ) - return f"({bv} IS DISTINCT FROM {av})" - - def _changed_predicate( compare_fields: list[str], rules: dict[str, dict[str, Any]], @@ -359,7 +300,7 @@ def _changed_predicate( norm = str(rule.get("normalize") or "strip").strip().lower() or "strip" bv = _norm_expr(f"{before_alias}.data", name, norm) av = _norm_expr(f"{after_alias}.data", name, norm) - clauses.append(_field_differs_sql(bv, av, rule)) + clauses.append(f"({bv} IS DISTINCT FROM {av})") if not clauses: return "FALSE" return "(" + " OR ".join(clauses) + ")" @@ -440,9 +381,6 @@ def run_sql_sheet_compare( tb = f"_netx_cmp_b_{tag}" ta = f"_netx_cmp_a_{tag}" - # Keep worker transaction open across CTAS — progress must NOT commit this session. - db.execute(text(f"SET LOCAL statement_timeout = '{int(_SQL_STATEMENT_TIMEOUT_MS)}'")) - _prog("before", 0, phase="sql_count", note="count") # Raw counts (no row_filters) @@ -469,20 +407,13 @@ def run_sql_sheet_compare( ) _prog("after", raw_a, phase="sql_count") - if max(raw_b, raw_a) < int(_SQL_MIN_ROWS): - raise SqlCompareSkip("small_sheet") - base_params = {"bid": bid_b, "mid": mid, **filter_params} - # Build TEMP sides (one transaction — no mid-flight commits on ``db``) - _prog("before", raw_b, phase="sql_project", note="temp_before") - for tname, batch_id, side, note in ( - (tb, bid_b, "before", "temp_before"), - (ta, bid_a, "after", "temp_after"), - ): + # Build TEMP sides + 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} - if side == "after": - _prog(side, raw_a, phase="sql_project", note=note) + _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( f""" @@ -503,9 +434,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 ON {tname} (rk) WHERE dup_rn = 1")) before_n = int( db.execute(text(f"SELECT count(*) FROM {tb}")).scalar() or 0 diff --git a/netx_api/biz_state/data/default_zte_status_template.json b/netx_api/biz_state/data/default_zte_status_template.json deleted file mode 100644 index d643cbc..0000000 --- a/netx_api/biz_state/data/default_zte_status_template.json +++ /dev/null @@ -1,1051 +0,0 @@ -{ - "format": "netx.biz_compare_template", - "version": 1, - "name": "ZTE status default", - "note": "Built-in ZTE status cutover (ISIS/IF/ARP/ND6/BGP, address-family sheets)", - "iface_normalize_rules": [], - "metrics": [ - { - "sheet_id": "isis_adjacency.ipv4", - "title": "ISIS IPv4", - "metric_id": "isis_adjacency", - "key_fields": [ - "process_id", - "interface", - "system_id" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "state", - "lev", - "af" - ], - "display_fields": [ - "process_id", - "interface", - "system_id", - "state", - "lev", - "af" - ], - "row_filters": [ - { - "op": "contains", - "field": "af", - "value": "IPv4" - } - ], - "field_rules": [] - }, - { - "sheet_id": "isis_adjacency.ipv6", - "title": "ISIS IPv6", - "metric_id": "isis_adjacency", - "key_fields": [ - "process_id", - "interface", - "system_id" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "state", - "lev", - "af" - ], - "display_fields": [ - "process_id", - "interface", - "system_id", - "state", - "lev", - "af" - ], - "row_filters": [ - { - "op": "contains", - "field": "af", - "value": "IPv6" - } - ], - "field_rules": [] - }, - { - "sheet_id": "interface_brief", - "title": "interface_brief", - "metric_id": "interface_brief", - "key_fields": [ - "interface" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "prot" - ], - "display_fields": [ - "interface", - "prot", - "admin", - "phy" - ], - "row_filters": [ - { - "op": "regex", - "field": "interface", - "value": "^[^.]+$" - }, - { - "op": "contains", - "field": "interface", - "value": "gei" - } - ], - "field_rules": [] - }, - { - "sheet_id": "arp", - "title": "arp", - "metric_id": "arp", - "key_fields": [ - "ip", - "interface", - "vrf" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "mac" - ], - "display_fields": [ - "ip", - "interface", - "vrf", - "mac", - "age", - "entry_type" - ], - "row_filters": [ - { - "op": "age_timer", - "field": "age", - "value": "" - } - ], - "field_rules": [ - { - "field": "mac", - "normalize": "mac" - } - ] - }, - { - "sheet_id": "nd6_cache", - "title": "nd6_cache", - "metric_id": "nd6_cache", - "key_fields": [ - "address", - "interface" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "link_address", - "status" - ], - "display_fields": [ - "address", - "interface", - "link_address", - "status", - "type", - "age" - ], - "row_filters": [ - { - "op": "ne", - "field": "status", - "value": " Incomplete" - } - ], - "field_rules": [] - }, - { - "sheet_id": "bgp_peer.ipv4", - "title": "BGP IPv4", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "ipv4" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_peer.ipv6", - "title": "BGP IPv6", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "ipv6" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_peer.vpnv4", - "title": "BGP VPNv4", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "vpnv4" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_peer.vpnv6", - "title": "BGP VPNv6", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "vpnv6" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_peer.evpn", - "title": "BGP EVPN", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "evpn" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_peer.vpls", - "title": "BGP VPLS", - "metric_id": "bgp_peer", - "key_fields": [ - "afi", - "vrf", - "neighbor", - "as_num" - ], - "iface_fields": [], - "compare_fields": [ - "state", - "pfx_rcd" - ], - "display_fields": [ - "afi", - "vrf", - "neighbor", - "as_num", - "state", - "pfx_rcd" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "vpls" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "interface_detail", - "title": "interface_detail", - "metric_id": "interface_detail", - "key_fields": [ - "interface" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "admin", - "input_bps", - "output_bps", - "port_status", - "in_util", - "out_util", - "ip_mtu", - "mtu", - "mpls_mtu", - "ipv6_mtu", - "port_media", - "negotiation" - ], - "display_fields": [ - "interface", - "admin", - "input_bps", - "output_bps", - "port_status", - "in_util", - "out_util", - "ip_mtu", - "mtu", - "mpls_mtu", - "ipv6_mtu", - "port_media", - "negotiation", - "bw", - "description", - "rate_period" - ], - "row_filters": [], - "field_rules": [ - { - "field": "input_bps", - "compare": "percent", - "tolerance": 5 - }, - { - "field": "output_bps", - "compare": "numeric", - "tolerance": 5 - }, - { - "field": "in_util", - "compare": "numeric", - "tolerance": 5 - }, - { - "field": "out_util", - "compare": "numeric", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "ospf_neighbor", - "title": "ospf_neighbor", - "metric_id": "ospf_neighbor", - "key_fields": [ - "neighbor_id", - "interface" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "state" - ], - "display_fields": [ - "neighbor_id", - "interface", - "state", - "process_id", - "address" - ], - "row_filters": [], - "field_rules": [] - }, - { - "sheet_id": "vrrp.ipv4", - "title": "VRRP IPv4", - "metric_id": "vrrp", - "key_fields": [ - "af", - "interface", - "vr_id" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "state", - "priority", - "master_addr", - "vrouter_addr" - ], - "display_fields": [ - "af", - "interface", - "vr_id", - "state", - "priority", - "master_addr", - "vrouter_addr" - ], - "row_filters": [ - { - "op": "eq", - "field": "af", - "value": "ipv4" - } - ], - "field_rules": [] - }, - { - "sheet_id": "vrrp.ipv6", - "title": "VRRP IPv6", - "metric_id": "vrrp", - "key_fields": [ - "af", - "interface", - "vr_id" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "state", - "priority", - "master_addr", - "vrouter_addr" - ], - "display_fields": [ - "af", - "interface", - "vr_id", - "state", - "priority", - "master_addr", - "vrouter_addr" - ], - "row_filters": [ - { - "op": "eq", - "field": "af", - "value": "ipv6" - } - ], - "field_rules": [] - }, - { - "sheet_id": "optical_brief", - "title": "optical_brief", - "metric_id": "optical_brief", - "key_fields": [ - "interface" - ], - "iface_fields": [ - "interface" - ], - "compare_fields": [ - "status", - "rx_power" - ], - "display_fields": [ - "interface", - "status", - "rx_power", - "tx_power", - "optic_type", - "wavelength", - "rx_threshold", - "tx_threshold" - ], - "row_filters": [], - "field_rules": [ - { - "field": "rx_power", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "lldp_neighbor", - "title": "lldp_neighbor", - "metric_id": "lldp_neighbor", - "key_fields": [ - "local_if", - "remote_sys", - "remote_if" - ], - "iface_fields": [ - "local_if" - ], - "compare_fields": [], - "display_fields": [ - "local_if", - "remote_sys", - "remote_if", - "remote_ip", - "protocol" - ], - "row_filters": [], - "field_rules": [] - }, - { - "sheet_id": "l2vpn_pw_detail", - "title": "l2vpn_pw_detail", - "metric_id": "l2vpn_pw_detail", - "key_fields": [ - "peer", - "vcid", - "service_instance" - ], - "iface_fields": [], - "compare_fields": [ - "vc_status", - "remote_status", - "vc_type", - "control_word" - ], - "display_fields": [ - "peer", - "vcid", - "service_instance", - "vc_status", - "remote_status", - "vc_type", - "control_word", - "pw_name", - "activation_status", - "service_instance_type", - "conn_mode", - "signaling", - "tunnel_dest" - ], - "row_filters": [], - "field_rules": [] - }, - { - "sheet_id": "l2vpn_mac", - "title": "l2vpn_mac", - "metric_id": "l2vpn_mac", - "key_fields": [ - "mac", - "vpn", - "neighbor", - "ac_port", - "exter_vlan", - "vpn_sid", - "neighbor_sid" - ], - "iface_fields": [], - "compare_fields": [], - "display_fields": [ - "mac", - "vpn", - "neighbor", - "ac_port", - "exter_vlan", - "vpn_sid", - "neighbor_sid", - "vlan", - "pw", - "attribute" - ], - "row_filters": [], - "field_rules": [] - }, - { - "sheet_id": "bgp_route", - "title": "bgp_route_ipv4", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "empty", - "field": "vrf", - "value": "" - }, - { - "op": "eq", - "field": "afi", - "value": "ipv4" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_route.vpnv4", - "title": "bgp_route_vpnv4", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "vpnv4" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_route.ipv6", - "title": "bgp_route_ipv6", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "empty", - "field": "vrf", - "value": "" - }, - { - "op": "eq", - "field": "afi", - "value": "ipv6" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_route.vpnv6", - "title": "bgp_route_vpnv6", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "eq", - "field": "afi", - "value": "vpnv6" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_route.vpnv6_vrf", - "title": "bgp_route_vpnv6_vrf", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "not_empty", - "field": "vrf", - "value": "" - }, - { - "op": "eq", - "field": "afi", - "value": "ipv6" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - }, - { - "sheet_id": "bgp_route.vpnv4_vrf", - "title": "bgp_route_vpnv4_vrf", - "metric_id": "bgp_route", - "key_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf" - ], - "iface_fields": [], - "compare_fields": [ - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as" - ], - "display_fields": [ - "local_as", - "afi", - "neighbor", - "direction", - "rd", - "network", - "vrf", - "metric", - "loc_prf", - "path", - "as_num", - "state", - "pfx_rcd", - "remote_as", - "next_hop", - "tag", - "status_codes" - ], - "row_filters": [ - { - "op": "not_empty", - "field": "vrf", - "value": "" - }, - { - "op": "eq", - "field": "afi", - "value": "ipv4" - } - ], - "field_rules": [ - { - "field": "pfx_rcd", - "compare": "percent", - "tolerance": 5 - } - ] - } - ] -} diff --git a/tests/test_biz_state_compare.py b/tests/test_biz_state_compare.py index 79a5f86..ce8636e 100644 --- a/tests/test_biz_state_compare.py +++ b/tests/test_biz_state_compare.py @@ -455,31 +455,21 @@ class CompareSheetDefaultsTests(unittest.TestCase): self.assertIn("isis_adjacency.ipv6", ids) # Same source metric may appear multiple times (afi splits) self.assertEqual(sum(1 for s in sheets if s["metric_id"] == "bgp_peer"), 6) - self.assertGreaterEqual(sum(1 for s in sheets if s["metric_id"] == "bgp_route"), 4) self.assertIn("bgp_peer.evpn", ids) self.assertIn("bgp_peer.vpls", ids) self.assertIn("vrrp.ipv4", ids) self.assertIn("lldp_neighbor", ids) - # Packaged IOH default: filtered bgp_route ipv4 + percent pfx_rcd - route4 = next(s for s in sheets if sheet_key(s) == "bgp_route") - self.assertEqual(route4["metric_id"], "bgp_route") - self.assertTrue( - any( - f.get("field") == "afi" and f.get("value") == "ipv4" - for f in (route4.get("row_filters") or []) - ) - ) - self.assertTrue( - any( - r.get("field") == "pfx_rcd" and r.get("compare") == "percent" - for r in (route4.get("field_rules") or []) - ) - ) + detail = next(s for s in sheets if sheet_key(s) == "interface_detail") + self.assertEqual(detail["compare_fields"], ["port_status"]) + self.assertIn("input_bps", detail["display_fields"]) + optical = next(s for s in sheets if sheet_key(s) == "optical_brief") + self.assertEqual(optical["compare_fields"], ["status"]) + self.assertIn("rx_power", optical["display_fields"]) + bgp4 = next(s for s in sheets if sheet_key(s) == "bgp_peer.ipv4") + self.assertEqual(bgp4["compare_fields"], ["as_num", "state"]) + self.assertIn("pfx_rcd", bgp4["display_fields"]) vpnv4 = next(s for s in sheets if sheet_key(s) == "bgp_peer.vpnv4") - self.assertEqual( - vpnv4["row_filters"], - [{"field": "afi", "op": "eq", "value": "vpnv4"}], - ) + self.assertEqual(vpnv4["row_filters"], [{"field": "afi", "op": "eq", "value": "vpnv4"}]) isis4 = next(s for s in sheets if sheet_key(s) == "isis_adjacency.ipv4") self.assertEqual(isis4["row_filters"][0]["op"], "contains") diff --git a/tests/test_biz_state_compare_sql.py b/tests/test_biz_state_compare_sql.py index 44b2891..e61f83c 100644 --- a/tests/test_biz_state_compare_sql.py +++ b/tests/test_biz_state_compare_sql.py @@ -112,18 +112,13 @@ class CompareSqlGateTests(unittest.TestCase): "network", "next_hop", ], - "compare_fields": ["path", "as_num", "pfx_rcd"], + "compare_fields": ["path", "as_num"], "iface_fields": [], - "field_rules": [ - {"field": "pfx_rcd", "compare": "percent", "tolerance": 5}, - ], - "row_filters": [ - {"field": "vrf", "op": "empty"}, - {"field": "afi", "op": "eq", "value": "ipv4"}, - ], + "field_rules": [], + "row_filters": [{"field": "afi", "op": "eq", "value": "ipv4"}], } self.assertTrue(can_sql_compare(_db(), sheet, port_map={})) - # BGP afi/vrf sheet splits + percent rules must not force Python + # 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: @@ -209,16 +204,11 @@ class CompareSqlFilterCompileTests(unittest.TestCase): def test_field_rules_matrix(self) -> None: self.assertTrue(_field_rules_sql_compatible([{"field": "mac", "normalize": "lower"}])) self.assertFalse(_field_rules_sql_compatible([{"field": "mac", "normalize": "mac"}])) - self.assertTrue( + self.assertFalse( _field_rules_sql_compatible( [{"field": "rx", "compare": "numeric", "tolerance": 1}] ) ) - self.assertTrue( - _field_rules_sql_compatible( - [{"field": "pfx_rcd", "compare": "percent", "tolerance": 5}] - ) - ) self.assertTrue( _field_rules_sql_compatible([{"field": "x", "compare": "ignore"}]) )