feat: refine cutover operations board and monitoring performance

Protect previous wave coverage, normalize batch identities, cache shared metric evidence, and batch diff inserts. Add live continuity KPIs, server paging and filters, scoped business selection, and preserve legacy monitor rule metadata.

Validation: 198 backend tests passed with 1 opt-in benchmark skipped; 14 frontend unit tests and 24 browser checks passed; TypeScript and production build passed.
This commit is contained in:
oliver 2026-10-11 07:55:04 +08:00
parent fe4c1db763
commit 32d7d969f6
12 changed files with 731 additions and 151 deletions

View file

@ -108,11 +108,11 @@ def parse_expect_set(raw: dict[str, Any] | None) -> dict[str, set[str]]:
if it.get("key") is not None:
k = it.get("key")
if isinstance(k, (list, tuple)):
joined = KEY_SEP.join(str(p).strip() for p in k if str(p).strip())
if joined:
joined = KEY_SEP.join(scalar_text(p).strip() for p in k)
if joined.strip(KEY_SEP):
bucket.add(joined)
else:
s = normalize_expect_key(str(k or ""))
s = normalize_expect_key(scalar_text(k))
if s:
bucket.add(s)
keys = it.get("keys")
@ -120,14 +120,14 @@ def parse_expect_set(raw: dict[str, Any] | None) -> dict[str, set[str]]:
# Flat list of segments → one composite key; else each entry is a key
# (string or nested list/tuple of segments).
if keys and all(not isinstance(x, (list, tuple, dict)) for x in keys):
joined = KEY_SEP.join(str(x).strip() for x in keys if str(x).strip())
if joined:
joined = KEY_SEP.join(scalar_text(x).strip() for x in keys)
if joined.strip(KEY_SEP):
bucket.add(joined)
else:
for x in keys:
if isinstance(x, (list, tuple)):
joined = KEY_SEP.join(str(p).strip() for p in x if str(p).strip())
if joined:
joined = KEY_SEP.join(scalar_text(p).strip() for p in x)
if joined.strip(KEY_SEP):
bucket.add(joined)
else:
s = normalize_expect_key(str(x or ""))
@ -272,6 +272,10 @@ def _eval_field_op(
raw = "" if not row else scalar_text(row.get(f)).strip()
val = raw.lower()
o = str(op or "eq").strip().lower()
if row and f not in row:
# A missing parser field is not evidence for a negative condition such as
# state != down. Only an explicit empty check may match that absence.
return o == "empty"
if o in ("changed", "diff"):
base = "" if not base_row else scalar_text(base_row.get(f)).strip().lower()
if not row:
@ -1119,6 +1123,7 @@ def evaluate_metric_dual(
from ..biz_state.iface_normalize import (
apply_iface_normalize_rows,
normalize_iface_rules,
normalize_iface_name,
)
sid = str(sheet_id or "").strip() or str(metric_id or "").strip()
@ -1155,6 +1160,15 @@ def evaluate_metric_dual(
)
norm_rules = normalize_iface_rules(iface_normalize_rules)
def normalized_key(key: str) -> str:
parts = key.split(KEY_SEP) if KEY_SEP in key else ([key] if len(key_fields) == 1 else key.split("|"))
if len(parts) != len(key_fields):
return key
return KEY_SEP.join(normalize_iface_name(p, norm_rules) if f in iface_fields else p
for p, f in zip(parts, key_fields))
expect_keys = {normalized_key(key) for key in expect_keys}
previous_keys = {normalized_key(key) for key in previous_keys}
old_baseline_rows = apply_iface_normalize_rows(
[strip_netx(r) for r in raw_old_base], iface_fields=iface_fields, rules=norm_rules
)
@ -1238,7 +1252,7 @@ def evaluate_metric_dual(
# Mapping describes identity; an explicit batch selection defines scope.
# Port selection also covers business rows attached to those interfaces.
for scope, keys in ((expect, expect_keys), (previous_expect or {}, previous_keys)):
ports = scope.get("_ports") or set()
ports = {normalize_iface_name(p, norm_rules) for p in scope.get("_ports") or set()}
if ports and not keys and iface_fields:
for ks, diff in old_idx.items():
row = diff.get("before") or diff.get("after") or {}
@ -1293,6 +1307,7 @@ def evaluate_metric_dual(
progress_total = 0
anomaly = 0
anomaly_in_expect = 0
steady_outside = 0
rules_by_field = field_rule_map(field_rules)
integrity_fields = effective_compare_fields(compare_fields, field_rules)
# Status transitions are interpreted by the monitor's correction rules.
@ -1390,9 +1405,10 @@ def evaluate_metric_dual(
verdict, rule_hit = "migrated_previous", "previous_batch_received"
elif color != "red":
verdict, color, rule_hit = "unfinished_previous", "red", "previous_batch_not_received"
# Retain all drift outside scope, but avoid filling the board with stable
# unmapped inventory when a port map defines the legacy migration scope.
if map_scoped and not in_exp and verdict == "not_involved":
# Stable inventory is covered by the comparison and counted, but does
# not need a full persisted evidence card on every high-frequency run.
if not in_exp and verdict == "not_involved":
steady_outside += 1
continue
if in_exp and verdict == "migrated" and color == "green":
progress_ok += 1
@ -1404,9 +1420,9 @@ def evaluate_metric_dual(
# Prefer live raw rows; fall back to baseline raw. Never use port-mapped
# synthetic "before" as the new-side display row.
raw_old = old_raw_cur_idx.get(old_ks) or old_raw_base_idx.get(old_ks)
raw_new = new_raw_cur_idx.get(new_ks) or new_raw_cur_idx.get(old_ks)
raw_new = new_raw_cur_idx.get(new_ks)
if not raw_new:
raw_new = new_raw_base_idx.get(new_ks) or new_raw_base_idx.get(old_ks)
raw_new = new_raw_base_idx.get(new_ks)
display_old = display_key_from_raw(raw_old, key_fields) or old_ks
display_new = display_key_from_raw(raw_new, key_fields) or (
@ -1420,8 +1436,8 @@ def evaluate_metric_dual(
new_disp = strip_netx(raw_new) if raw_new else {}
old_base_raw = strip_netx(old_raw_base_idx.get(old_ks)) if old_raw_base_idx.get(old_ks) else {}
new_base_raw = (
strip_netx(new_raw_base_idx.get(new_ks) or new_raw_base_idx.get(old_ks))
if (new_raw_base_idx.get(new_ks) or new_raw_base_idx.get(old_ks))
strip_netx(new_raw_base_idx.get(new_ks))
if new_raw_base_idx.get(new_ks)
else {}
)
@ -1512,6 +1528,7 @@ def evaluate_metric_dual(
"progress_total": progress_total,
"anomaly": anomaly,
"anomaly_in_expect": anomaly_in_expect,
"steady_outside": steady_outside,
"rows": rows_out,
"new_baseline_mode": new_baseline_mode,
"new_baseline_missing": new_baseline_missing,

View file

@ -7,9 +7,10 @@ from typing import Any
from uuid import uuid4
from fastapi import HTTPException
from sqlalchemy import insert
from sqlalchemy.orm import Session
from ..biz_state.compare_rules import apply_row_filters
from ..biz_state.compare_rules import apply_row_filters, scalar_text
from ..biz_state.compare_service import (
_load_metric_rows,
_port_map_dict,
@ -949,6 +950,13 @@ def run_evaluate(
from ..biz_state.compare_validation import validate_template_body
validate_template_body({"metrics": sheets, "iface_normalize_rules": iface_norm})
from ..biz_state.iface_normalize import normalize_iface_name
normalized_map = {normalize_iface_name(k, iface_norm): normalize_iface_name(v, iface_norm)
for k, v in port_map.items()}
if len(normalized_map) != len(port_map) or len(set(normalized_map.values())) != len(normalized_map):
raise HTTPException(status_code=400, detail="mapping_target_ambiguous")
port_map = normalized_map
known_sheets = {sheet_key(s) for s in sheets}
known_metrics = {str(s.get("metric_id") or "") for s in sheets}
for override in sheet_overrides:
@ -967,7 +975,8 @@ def run_evaluate(
rows_cache: dict[tuple[str, str], list[dict[str, Any]]] = {}
command_cache: dict[str, dict[str, Any]] = {}
scope_ids = {str(s.get("metric_id") or "") for s in sheets} | {sheet_key(s) for s in sheets} | {"_ports"}
uncovered_scope = sorted(k for k, values in expect.items() if values and k not in scope_ids)
uncovered_scope = sorted({k for scope in (expect, previous_expect) for k, values in scope.items()
if values and k not in scope_ids})
def metric_rows(bid: str, mid: str) -> list[dict[str, Any]]:
key = (bid, mid)
@ -1116,6 +1125,7 @@ def run_evaluate(
"progress_total": one["progress_total"],
"anomaly": one["anomaly"],
"anomaly_in_expect": int(one.get("anomaly_in_expect") or 0),
"steady_outside": int(one.get("steady_outside") or 0),
"old_summary": one["old_summary"],
"new_summary": one["new_summary"],
"new_baseline_mode": one.get("new_baseline_mode") or "provided",
@ -1189,12 +1199,14 @@ def run_evaluate(
],
"verdict_counts": verdict_counts,
"anomaly_outside_expect": sum(int(c.get("anomaly") or 0) - int(c.get("anomaly_in_expect") or 0) for c in active_cards),
"steady_outside": sum(int(c.get("steady_outside") or 0) for c in active_cards),
"duplicate_keys": sum(int(c.get("duplicate_keys_before") or 0) + int(c.get("duplicate_keys_after") or 0) for c in active_cards),
"uncovered_scope": uncovered_scope,
"coverage_complete": bool(sheet_cards) and not uncovered_scope and not any(c.get("collect_skipped") or c.get("current_missing") or c.get("awaiting_peer") or c.get("duplicate_keys_before") or c.get("duplicate_keys_after") for c in sheet_cards),
"config_snapshot": {"sheets": sheets, "sheet_overrides": sheet_overrides,
"defaults": defaults, "iface_normalize": iface_norm,
"port_map": port_map, "expect_set": mb.expect_set_json},
"port_map": port_map, "expect_set": mb.expect_set_json,
"previous_expect": {k: sorted(v) for k, v in previous_expect.items()}},
"window_active": window_active,
"expect_ports": sorted(expect.get("_ports") or ()),
"old_current_task_id": current_task_id(proj, "old"),
@ -1204,6 +1216,7 @@ def run_evaluate(
)
db.add(run)
db.flush()
pending_diffs: list[dict[str, Any]] = []
for r in all_rows:
key_list = r.get("key") or []
ev = r.get("evidence") if isinstance(r.get("evidence"), dict) else {}
@ -1222,8 +1235,8 @@ def run_evaluate(
str(r.get("metric_id") or ""),
]
)
db.add(
BizMigrationDiff(
pending_diffs.append(
dict(
id=uuid4().hex,
run_id=run.id,
metric_id=str(r.get("metric_id") or ""),
@ -1254,6 +1267,11 @@ def run_evaluate(
search_text=search[:2000],
)
)
if len(pending_diffs) >= 500:
db.execute(insert(BizMigrationDiff), pending_diffs)
pending_diffs.clear()
if pending_diffs:
db.execute(insert(BizMigrationDiff), pending_diffs)
db.flush()
if _commit:
db.commit()
@ -1436,6 +1454,11 @@ def list_baseline_expect_objects(db: Session, project_id: str) -> dict[str, Any]
collect_ids = set(resolve_collect_metric_ids(db, p))
port_map = _port_map_dict(db, p.mapping_id)
out_sheets: list[dict[str, Any]] = []
from ..biz_state.iface_normalize import apply_iface_normalize, resolve_mapped_iface, normalize_iface_name
from .evaluate import KEY_SEP
rows_cache: dict[str, list[dict[str, Any]]] = {}
port_map = {normalize_iface_name(k, _norm): normalize_iface_name(v, _norm) for k, v in port_map.items()}
for sheet in sheets:
mid = str(sheet.get("metric_id") or "").strip()
sid = sheet_key(sheet)
@ -1449,25 +1472,32 @@ def list_baseline_expect_objects(db: Session, project_id: str) -> dict[str, Any]
if sheet_ov.get("skip_dual"):
continue
iface_fields = [str(k) for k in (sheet.get("iface_fields") or []) if str(k).strip()]
if mid not in rows_cache:
rows_cache[mid] = _load_metric_rows(db, batch_id=p.old_baseline_batch_id, metric_id=mid)
rows_raw = apply_row_filters(
_load_metric_rows(db, batch_id=p.old_baseline_batch_id, metric_id=mid),
rows_cache[mid],
list(sheet.get("row_filters") or []),
)
items: list[dict[str, Any]] = []
seen_keys: set[str] = set()
for r in rows_raw:
keys = [str(r.get(f) or "").strip() for f in key_fields]
normalized = apply_iface_normalize(dict(r), iface_fields=iface_fields, rules=_norm)
keys = [scalar_text(normalized.get(f)).strip() for f in key_fields]
if not any(keys):
continue
key_str = "|".join(keys)
key_str = KEY_SEP.join(keys)
if key_str in seen_keys:
continue
seen_keys.add(key_str)
mapped = ""
if iface_fields and key_fields and key_fields[0] in iface_fields:
mapped = port_map.get(keys[0]) or ""
mapped = resolve_mapped_iface(keys[0], port_map) if port_map else ""
items.append(
{
"key": key_str,
"keys": keys,
"mapped_to": mapped,
"label": key_str,
"label": "|".join(scalar_text(r.get(f)).strip() for f in key_fields),
"row": {f: r.get(f) for f in list(dict.fromkeys([*key_fields, *iface_fields, "description", "admin", "phy", "prot"])) if f in r},
}
)