From 8631f58113eb6e8b9987792d059b111b58fcb7fd Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 18 Sep 2026 16:43:26 +0800 Subject: [PATCH] Add cutover migration monitor and day-based batch retention. Protect baselines and compare/migration-referenced batches from purge or delete, and let operators mark baselines and clean high-freq data manually. Co-authored-by: Cursor --- netx_api/biz_migration/__init__.py | 1 + netx_api/biz_migration/evaluate.py | 343 ++++++++ netx_api/biz_migration/service.py | 973 +++++++++++++++++++++ netx_api/biz_migration_router.py | 211 +++++ netx_api/biz_state/collect_runner.py | 41 +- netx_api/biz_state/retention.py | 168 ++++ netx_api/biz_state/schema_ensure.py | 19 + netx_api/biz_state/service.py | 169 +++- netx_api/biz_state_router.py | 42 +- netx_api/biz_state_scheduler.py | 32 + netx_api/main.py | 2 + netx_api/models/__init__.py | 12 + netx_api/models/biz_migration.py | 129 +++ netx_api/models/biz_state.py | 16 +- tests/test_biz_migration_evaluate.py | 136 +++ tests/test_biz_state_retention.py | 97 ++ web/src/App.tsx | 4 + web/src/config/networkNav.ts | 6 + web/src/i18n/en.ts | 104 ++- web/src/i18n/zh.ts | 103 ++- web/src/pages/network/BizMigrationPage.tsx | 938 ++++++++++++++++++++ web/src/pages/network/BizStatePage.tsx | 245 +++++- web/src/services/api.ts | 137 +++ 23 files changed, 3837 insertions(+), 91 deletions(-) create mode 100644 netx_api/biz_migration/__init__.py create mode 100644 netx_api/biz_migration/evaluate.py create mode 100644 netx_api/biz_migration/service.py create mode 100644 netx_api/biz_migration_router.py create mode 100644 netx_api/biz_state/retention.py create mode 100644 netx_api/models/biz_migration.py create mode 100644 tests/test_biz_migration_evaluate.py create mode 100644 tests/test_biz_state_retention.py create mode 100644 web/src/pages/network/BizMigrationPage.tsx diff --git a/netx_api/biz_migration/__init__.py b/netx_api/biz_migration/__init__.py new file mode 100644 index 0000000..29155b0 --- /dev/null +++ b/netx_api/biz_migration/__init__.py @@ -0,0 +1 @@ +"""biz_migration package — cutover monitor (separate from CompareJob).""" diff --git a/netx_api/biz_migration/evaluate.py b/netx_api/biz_migration/evaluate.py new file mode 100644 index 0000000..223eaf5 --- /dev/null +++ b/netx_api/biz_migration/evaluate.py @@ -0,0 +1,343 @@ +"""Cutover monitor: relative-to-baseline drift + dual-device migration verdict.""" + +from __future__ import annotations + +from typing import Any + +from ..biz_state.compare_engine import apply_port_map, compare_rows + +PORT_METRIC_ID = "interface_brief" +PORT_STATUS_FIELDS = ("admin", "phy", "prot") + + +def _key_str(key: tuple[str, ...] | list[str]) -> str: + return "|".join(str(x) for x in key) + + +def port_status_label(row: dict[str, Any] | None) -> str: + """Compact admin/phy/prot for board (e.g. up/up/up).""" + if not row: + return "—" + parts = [str(row.get(f) or "-").lower() for f in PORT_STATUS_FIELDS] + if all(p == "-" for p in parts): + return "—" + return "/".join(parts) + + +def parse_expect_set(raw: dict[str, Any] | None) -> dict[str, set[str]]: + """Return metric_id → set of key strings (old-side / before-map identity). + + Supports: + {"ports": ["gei-0/1/0/1", ...]} → interface_brief + {"items": [{"metric_id": "...", "key": "a|b"}, {"metric_id":"...", "keys":["a","b"]}]} + """ + out: dict[str, set[str]] = {} + data = raw or {} + ports = data.get("ports") or [] + if isinstance(ports, list): + for p in ports: + s = str(p or "").strip() + if s: + out.setdefault(PORT_METRIC_ID, set()).add(s) + out.setdefault("_ports", set()).add(s) + items = data.get("items") or [] + if isinstance(items, list): + for it in items: + if not isinstance(it, dict): + continue + mid = str(it.get("metric_id") or "").strip() + if not mid: + continue + bucket = out.setdefault(mid, set()) + if it.get("key") is not None: + bucket.add(str(it.get("key") or "").strip()) + keys = it.get("keys") + if isinstance(keys, list): + if all(not isinstance(x, (list, tuple, dict)) for x in keys): + bucket.add("|".join(str(x).strip() for x in keys)) + else: + for x in keys: + bucket.add(str(x).strip()) + return {k: {x for x in v if x} for k, v in out.items() if v} + + +def side_verdict( + *, + kind: str, + in_expect: bool, + window_active: bool, +) -> tuple[str, str]: + """Map compare kind → (verdict, color) for one device vs its baseline.""" + k = str(kind or "") + if k == "unchanged": + return "ok", "green" + if k == "removed": + if in_expect: + return "expected_gone", "yellow" + return "anomaly_gone", "red" + if k == "added": + if in_expect: + return "expected_new", "green" + return "unexpected_new", "yellow" if window_active else "yellow" + if k == "changed": + if in_expect: + return "expected_change", "yellow" + return "anomaly_change", "red" + return "ok", "gray" + + +def dual_verdict( + *, + old_kind: str, + new_kind: str, + in_expect: bool, + window_active: bool, + acceptance: bool = False, +) -> tuple[str, str]: + """Synthesize old+new into migration board verdict. + + ``acceptance=True`` (本批完成终验): unfinished expect items become red + (``unfinished`` / ``lost``), not yellow migrating. + """ + if not in_expect: + if old_kind in ("removed", "changed") or new_kind in ("removed", "changed"): + if old_kind == "removed" and new_kind in ("", "removed"): + return "lost", "red" + if old_kind in ("removed", "changed") or new_kind in ("changed",): + return "anomaly", "red" + if old_kind == "added" or new_kind == "added": + return "unexpected_new", "yellow" + return "not_involved", "gray" + + if old_kind in ("removed", "changed") and new_kind in ("added", "unchanged", "changed"): + return "migrated", "green" + if old_kind == "removed" and new_kind in ("", "removed"): + return "lost", "red" + if old_kind in ("unchanged", "") and new_kind in ("unchanged", "added", "changed"): + if acceptance: + return "unfinished", "red" + return "migrating", "yellow" + if old_kind == "unchanged" and new_kind in ("", "removed"): + if acceptance: + return "unfinished", "red" + return "migrating", "yellow" + ov, oc = side_verdict(kind=old_kind or "unchanged", in_expect=True, window_active=window_active) + if oc == "red" or (new_kind in ("removed",) and old_kind != "removed"): + return "anomaly", "red" + if acceptance: + return "unfinished", "red" + if window_active or ov.startswith("expected"): + return "migrating", "yellow" + return "migrating", "yellow" + + +def build_diff_index_from_compare(result: dict[str, Any]) -> dict[str, dict[str, Any]]: + """Normalize compare_rows result diffs → key_str → {kind, before, after, key}.""" + out: dict[str, dict[str, Any]] = {} + for d in result.get("diffs") or []: + key = d.get("key") + if isinstance(key, dict): + key_list = [str(v) for v in key.values()] + ks = _key_str(key_list) + elif isinstance(key, (list, tuple)): + key_list = [str(x) for x in key] + ks = _key_str(key_list) + else: + ks = str(key or "") + key_list = [ks] if ks else [] + before = dict(d.get("before") or {}) if d.get("before") else {} + after = dict(d.get("after") or {}) if d.get("after") else {} + out[ks] = { + "kind": str(d.get("kind") or ""), + "key": key_list, + "before": before, + "after": after, + "mapped_before": dict(d.get("mapped_before") or {}) if d.get("mapped_before") else {}, + } + return out + + +def expect_keys_for_metric( + expect: dict[str, set[str]], + *, + metric_id: str, + iface_fields: list[str], +) -> set[str]: + keys = set(expect.get(metric_id) or ()) + if iface_fields and expect.get("_ports"): + keys |= set(expect["_ports"]) + return keys + + +def _current_row(diff: dict[str, Any] | None) -> dict[str, Any]: + """Current-side snapshot from a compare diff (empty if removed / missing).""" + if not diff: + return {} + kind = str(diff.get("kind") or "") + if kind == "removed": + return {} + return dict(diff.get("after") or {}) + + +def evaluate_metric_dual( + *, + metric_id: str, + key_fields: list[str], + iface_fields: list[str], + compare_fields: list[str], + old_baseline_rows: list[dict[str, Any]], + old_current_rows: list[dict[str, Any]], + new_baseline_rows: list[dict[str, Any]] | None, + new_current_rows: list[dict[str, Any]], + port_map: dict[str, str], + expect: dict[str, set[str]], + window_active: bool, + acceptance: bool = False, +) -> dict[str, Any]: + """Run old vs old-baseline, new vs new-baseline (or mapped old baseline), dual merge.""" + expect_keys = expect_keys_for_metric(expect, metric_id=metric_id, iface_fields=iface_fields) + + old_cmp = compare_rows( + before_rows=old_baseline_rows, + after_rows=old_current_rows, + key_fields=key_fields, + iface_fields=[], + compare_fields=compare_fields, + port_map=None, + ) + old_idx = build_diff_index_from_compare(old_cmp) + + if new_baseline_rows is not None: + new_before = new_baseline_rows + elif port_map and iface_fields: + new_before = [ + apply_port_map(r, iface_fields=iface_fields, port_map=port_map) + for r in old_baseline_rows + ] + else: + new_before = [] + + new_cmp = compare_rows( + before_rows=new_before, + after_rows=new_current_rows, + key_fields=key_fields, + iface_fields=[], + compare_fields=compare_fields, + port_map=None, + ) + new_idx = build_diff_index_from_compare(new_cmp) + + rev_map = {v: k for k, v in port_map.items()} + + canon_keys: list[str] = [] + seen: set[str] = set() + + def _add(canon: str) -> None: + if canon and canon not in seen: + seen.add(canon) + canon_keys.append(canon) + + for ok in sorted(expect_keys): + _add(ok) + for ks in sorted(old_idx): + _add(ks) + for ks in sorted(new_idx): + _add(rev_map.get(ks, ks)) + + rows_out: list[dict[str, Any]] = [] + progress_ok = 0 + progress_total = 0 + anomaly = 0 + + for old_ks in canon_keys: + new_ks = port_map.get(old_ks, old_ks) + in_exp = old_ks in expect_keys + if in_exp: + progress_total += 1 + + od = old_idx.get(old_ks) + nd = new_idx.get(new_ks) or new_idx.get(old_ks) + + old_kind = str((od or {}).get("kind") or "") + new_kind = str((nd or {}).get("kind") or "") + if od is None and nd is None and not in_exp: + continue + if od is None: + old_kind = "" + if nd is None: + new_kind = "" + + verdict, color = dual_verdict( + old_kind=old_kind or "", + new_kind=new_kind or "", + in_expect=in_exp, + window_active=window_active, + acceptance=acceptance, + ) + if in_exp and verdict == "migrated" and color == "green": + progress_ok += 1 + if color == "red": + anomaly += 1 + + old_cur = _current_row(od) + new_cur = _current_row(nd) + + rows_out.append( + { + "metric_id": metric_id, + "key": (od or nd or {}).get("key") or [old_ks], + "key_str": old_ks, + "new_key_str": new_ks, + "verdict": verdict, + "color": color, + "in_expect": in_exp, + "old_kind": old_kind, + "new_kind": new_kind, + "old": old_cur, + "new": new_cur, + "old_baseline": dict((od or {}).get("before") or {}), + "new_baseline": dict((nd or {}).get("before") or {}), + "old_status": port_status_label(old_cur) + if old_cur + else ("gone" if old_kind == "removed" else "—"), + "new_status": port_status_label(new_cur) + if new_cur + else ("gone" if new_kind == "removed" else "—"), + "old_side": side_verdict( + kind=old_kind or "unchanged", + in_expect=in_exp, + window_active=window_active, + ) + if old_kind + else ("", "gray"), + "new_side": side_verdict( + kind=new_kind or "unchanged", + in_expect=in_exp, + window_active=window_active, + ) + if new_kind + else ("", "gray"), + } + ) + + return { + "metric_id": metric_id, + "old_summary": old_cmp.get("summary") or {}, + "new_summary": new_cmp.get("summary") or {}, + "progress_ok": progress_ok, + "progress_total": progress_total if progress_total else len(expect_keys), + "anomaly": anomaly, + "rows": rows_out, + } + + +def port_sheet_def() -> dict[str, Any]: + """Default sheet for port-status cutover monitor (interface_brief only).""" + return { + "metric_id": PORT_METRIC_ID, + "key_fields": ["interface"], + "iface_fields": ["interface"], + "compare_fields": ["admin", "phy", "prot"], + "row_filters": [], + "field_rules": [], + } diff --git a/netx_api/biz_migration/service.py b/netx_api/biz_migration/service.py new file mode 100644 index 0000000..00bd895 --- /dev/null +++ b/netx_api/biz_migration/service.py @@ -0,0 +1,973 @@ +"""CRUD + evaluate for cutover migration monitor.""" + +from __future__ import annotations + +from typing import Any +from uuid import uuid4 + +from fastapi import HTTPException +from sqlalchemy.orm import Session + +from ..biz_state.compare_rules import apply_row_filters +from ..biz_state.compare_service import _load_metric_rows, _port_map_dict +from ..models import ( + BizMigrationBatch, + BizMigrationDiff, + BizMigrationProject, + BizMigrationRedTicket, + BizMigrationRun, + BizPortMapping, + BizStateBatch, + BizStateTask, +) +from ..timeutil import utcnow_naive +from .evaluate import PORT_METRIC_ID, evaluate_metric_dual, parse_expect_set, port_sheet_def + + +def _task_brief(db: Session, task_id: str) -> dict[str, Any]: + t = db.get(BizStateTask, task_id) if task_id else None + if not t: + return {"id": task_id or "", "ne_name": "", "ne_ip": "", "vendor": ""} + return { + "id": t.id, + "ne_name": t.ne_name, + "ne_ip": t.ne_ip, + "vendor": t.vendor, + "status": t.status, + "note": t.note, + "interval_sec": t.interval_sec, + "collect_running": bool(t.collect_running), + } + + +def _batch_brief(db: Session, batch_id: str) -> dict[str, Any]: + b = db.get(BizStateBatch, batch_id) if batch_id else None + if not b: + return {"id": batch_id or "", "status": "", "started_at": None} + return { + "id": b.id, + "status": b.status, + "started_at": b.started_at.isoformat() if b.started_at else None, + "row_count": b.row_count, + } + + +def project_to_dict(db: Session, p: BizMigrationProject) -> dict[str, Any]: + return { + "id": p.id, + "name": p.name, + "old_task_id": p.old_task_id, + "new_task_id": p.new_task_id, + "old_baseline_batch_id": p.old_baseline_batch_id, + "new_baseline_batch_id": p.new_baseline_batch_id, + "mapping_id": p.mapping_id, + "status": p.status, + "note": p.note, + "old_task": _task_brief(db, p.old_task_id), + "new_task": _task_brief(db, p.new_task_id), + "old_baseline": _batch_brief(db, p.old_baseline_batch_id), + "new_baseline": _batch_brief(db, p.new_baseline_batch_id), + "created_at": p.created_at.isoformat() if p.created_at else None, + "updated_at": p.updated_at.isoformat() if p.updated_at else None, + } + + +def batch_to_dict(b: BizMigrationBatch) -> dict[str, Any]: + return { + "id": b.id, + "project_id": b.project_id, + "batch_label": b.batch_label, + "status": b.status, + "expect_set": dict(b.expect_set_json or {}), + "started_at": b.started_at.isoformat() if b.started_at else None, + "ended_at": b.ended_at.isoformat() if b.ended_at else None, + "accept_status": getattr(b, "accept_status", None) or "none", + "accept_run_id": getattr(b, "accept_run_id", None) or "", + "accept_summary": dict(getattr(b, "accept_summary_json", None) or {}), + "note": b.note, + "created_at": b.created_at.isoformat() if b.created_at else None, + "updated_at": b.updated_at.isoformat() if b.updated_at else None, + } + + +def list_projects(db: Session) -> list[dict[str, Any]]: + rows = db.query(BizMigrationProject).order_by(BizMigrationProject.created_at.desc()).all() + return [project_to_dict(db, p) for p in rows] + + +def create_project(db: Session, body: dict[str, Any]) -> dict[str, Any]: + name = str(body.get("name") or "").strip() + if not name: + raise HTTPException(status_code=400, detail="name_required") + old_task_id = str(body.get("old_task_id") or "").strip() + new_task_id = str(body.get("new_task_id") or "").strip() + if not old_task_id or not new_task_id: + raise HTTPException(status_code=400, detail="old_new_task_required") + if not db.get(BizStateTask, old_task_id) or not db.get(BizStateTask, new_task_id): + raise HTTPException(status_code=404, detail="task_not_found") + mapping_id = str(body.get("mapping_id") or "").strip() + if mapping_id and not db.get(BizPortMapping, mapping_id): + raise HTTPException(status_code=404, detail="mapping_not_found") + p = BizMigrationProject( + id=uuid4().hex, + name=name, + old_task_id=old_task_id, + new_task_id=new_task_id, + old_baseline_batch_id=str(body.get("old_baseline_batch_id") or "").strip(), + new_baseline_batch_id=str(body.get("new_baseline_batch_id") or "").strip(), + mapping_id=mapping_id, + status=str(body.get("status") or "draft").strip() or "draft", + note=str(body.get("note") or "")[:500], + ) + db.add(p) + db.commit() + db.refresh(p) + return project_to_dict(db, p) + + +def get_project(db: Session, project_id: str) -> dict[str, Any]: + p = db.get(BizMigrationProject, project_id) + if not p: + raise HTTPException(status_code=404, detail="project_not_found") + return project_to_dict(db, p) + + +def patch_project(db: Session, project_id: str, body: dict[str, Any]) -> dict[str, Any]: + p = db.get(BizMigrationProject, project_id) + if not p: + raise HTTPException(status_code=404, detail="project_not_found") + if "name" in body and body["name"] is not None: + p.name = str(body["name"]).strip() or p.name + if "note" in body and body["note"] is not None: + p.note = str(body["note"])[:500] + if "status" in body and body["status"] is not None: + p.status = str(body["status"]).strip() or p.status + if "mapping_id" in body and body["mapping_id"] is not None: + mid = str(body["mapping_id"] or "").strip() + if mid and not db.get(BizPortMapping, mid): + raise HTTPException(status_code=404, detail="mapping_not_found") + p.mapping_id = mid + if "old_baseline_batch_id" in body and body["old_baseline_batch_id"] is not None: + bid = str(body["old_baseline_batch_id"] or "").strip() + if bid and not db.get(BizStateBatch, bid): + raise HTTPException(status_code=404, detail="batch_not_found") + p.old_baseline_batch_id = bid + if "new_baseline_batch_id" in body and body["new_baseline_batch_id"] is not None: + bid = str(body["new_baseline_batch_id"] or "").strip() + if bid and not db.get(BizStateBatch, bid): + raise HTTPException(status_code=404, detail="batch_not_found") + p.new_baseline_batch_id = bid + p.updated_at = utcnow_naive() + db.commit() + db.refresh(p) + return project_to_dict(db, p) + + +def delete_project(db: Session, project_id: str) -> dict[str, Any]: + p = db.get(BizMigrationProject, project_id) + if not p: + raise HTTPException(status_code=404, detail="project_not_found") + batches = db.query(BizMigrationBatch).filter(BizMigrationBatch.project_id == project_id).all() + for b in batches: + runs = db.query(BizMigrationRun).filter(BizMigrationRun.batch_id == b.id).all() + for r in runs: + db.query(BizMigrationDiff).filter(BizMigrationDiff.run_id == r.id).delete() + db.delete(r) + db.delete(b) + db.delete(p) + db.commit() + return {"ok": True} + + +def list_batches(db: Session, project_id: str) -> list[dict[str, Any]]: + if not db.get(BizMigrationProject, project_id): + raise HTTPException(status_code=404, detail="project_not_found") + rows = ( + db.query(BizMigrationBatch) + .filter(BizMigrationBatch.project_id == project_id) + .order_by(BizMigrationBatch.created_at.desc()) + .all() + ) + return [batch_to_dict(b) for b in rows] + + +def create_batch(db: Session, project_id: str, body: dict[str, Any]) -> dict[str, Any]: + if not db.get(BizMigrationProject, project_id): + raise HTTPException(status_code=404, detail="project_not_found") + label = str(body.get("batch_label") or "").strip() or "batch" + expect = body.get("expect_set") if isinstance(body.get("expect_set"), dict) else {} + b = BizMigrationBatch( + id=uuid4().hex, + project_id=project_id, + batch_label=label, + status=str(body.get("status") or "pending").strip() or "pending", + expect_set_json=dict(expect or {}), + note=str(body.get("note") or "")[:500], + ) + db.add(b) + db.commit() + db.refresh(b) + out = batch_to_dict(b) + out["open_red_count"] = count_open_red_tickets(db, project_id) + return out + + +def _latest_success_batch(db: Session, task_id: str) -> BizStateBatch | None: + return ( + db.query(BizStateBatch) + .filter( + BizStateBatch.task_id == task_id, + BizStateBatch.status.in_(("success", "partial")), + ) + .order_by(BizStateBatch.started_at.desc()) + .first() + ) + + +def pinned_baseline_batch_ids(db: Session) -> set[str]: + """Batch IDs that must survive retention purge.""" + ids: set[str] = set() + for p in db.query(BizMigrationProject).all(): + if p.old_baseline_batch_id: + ids.add(p.old_baseline_batch_id) + if p.new_baseline_batch_id: + ids.add(p.new_baseline_batch_id) + return ids + + +def run_evaluate( + db: Session, + *, + batch_id: str, + old_batch_id: str = "", + new_batch_id: str = "", + acceptance: bool = False, + purpose: str = "manual", +) -> dict[str, Any]: + mb = db.get(BizMigrationBatch, batch_id) + if not mb: + raise HTTPException(status_code=404, detail="batch_not_found") + proj = db.get(BizMigrationProject, mb.project_id) + if not proj: + raise HTTPException(status_code=404, detail="project_not_found") + if not proj.old_baseline_batch_id: + raise HTTPException(status_code=400, detail="old_baseline_required") + + old_batch = db.get(BizStateBatch, old_batch_id.strip()) if old_batch_id.strip() else None + if not old_batch: + old_batch = _latest_success_batch(db, proj.old_task_id) + new_batch = db.get(BizStateBatch, new_batch_id.strip()) if new_batch_id.strip() else None + if not new_batch: + new_batch = _latest_success_batch(db, proj.new_task_id) + if not old_batch: + raise HTTPException(status_code=400, detail="old_current_batch_required") + if not new_batch: + raise HTTPException(status_code=400, detail="new_current_batch_required") + old_cur = old_batch.id + new_cur = new_batch.id + + port_map = _port_map_dict(db, proj.mapping_id) + expect = parse_expect_set(mb.expect_set_json if isinstance(mb.expect_set_json, dict) else {}) + # Final acceptance: window closed → unfinished expect = red + window_active = (mb.status == "active") and (not acceptance) + sheets = [port_sheet_def()] + + sheet_cards: list[dict[str, Any]] = [] + all_rows: list[dict[str, Any]] = [] + seq = 0 + verdict_counts: dict[str, int] = {} + + for sheet in sheets: + mid = sheet["metric_id"] + key_fields = list(sheet.get("key_fields") or []) + if not key_fields: + continue + iface_fields = list(sheet.get("iface_fields") or []) + compare_fields = list(sheet.get("compare_fields") or []) + row_filters = list(sheet.get("row_filters") or []) + + old_base = apply_row_filters( + _load_metric_rows(db, batch_id=proj.old_baseline_batch_id, metric_id=mid), + row_filters, + ) + old_now = apply_row_filters( + _load_metric_rows(db, batch_id=old_cur, metric_id=mid), + row_filters, + ) + new_base_rows = None + if proj.new_baseline_batch_id: + new_base_rows = apply_row_filters( + _load_metric_rows(db, batch_id=proj.new_baseline_batch_id, metric_id=mid), + row_filters, + ) + new_now = apply_row_filters( + _load_metric_rows(db, batch_id=new_cur, metric_id=mid), + row_filters, + ) + + one = evaluate_metric_dual( + metric_id=mid, + key_fields=key_fields, + iface_fields=iface_fields, + compare_fields=compare_fields, + old_baseline_rows=old_base, + old_current_rows=old_now, + new_baseline_rows=new_base_rows, + new_current_rows=new_now, + port_map=port_map, + expect=expect, + window_active=window_active, + acceptance=acceptance, + ) + sheet_cards.append( + { + "metric_id": mid, + "title": "端口状态", + "progress_ok": one["progress_ok"], + "progress_total": one["progress_total"], + "anomaly": one["anomaly"], + "old_summary": one["old_summary"], + "new_summary": one["new_summary"], + } + ) + for r in one["rows"]: + r["seq"] = seq + seq += 1 + all_rows.append(r) + v = str(r.get("verdict") or "") + if v: + verdict_counts[v] = verdict_counts.get(v, 0) + 1 + + run = BizMigrationRun( + id=uuid4().hex, + project_id=proj.id, + batch_id=mb.id, + old_batch_id=old_cur, + new_batch_id=new_cur, + purpose=str(purpose or ("acceptance" if acceptance else "manual"))[:32], + status="success", + summary_json={ + "metric_focus": PORT_METRIC_ID, + "acceptance": acceptance, + "sheet_cards": sheet_cards, + "progress": { + "ok": sum(c["progress_ok"] for c in sheet_cards), + "total": sum(c["progress_total"] for c in sheet_cards), + }, + "anomaly": sum(c["anomaly"] for c in sheet_cards), + "verdict_counts": verdict_counts, + "window_active": window_active, + "expect_ports": sorted(expect.get("_ports") or ()), + }, + message="", + ) + db.add(run) + db.flush() + for r in all_rows: + key_list = r.get("key") or [] + search = " ".join( + [ + str(r.get("key_str") or ""), + str(r.get("new_key_str") or ""), + str(r.get("verdict") or ""), + str(r.get("old_status") or ""), + str(r.get("new_status") or ""), + str(r.get("metric_id") or ""), + ] + ) + db.add( + BizMigrationDiff( + id=uuid4().hex, + run_id=run.id, + metric_id=str(r.get("metric_id") or ""), + seq=int(r.get("seq") or 0), + verdict=str(r.get("verdict") or ""), + color=str(r.get("color") or ""), + key_json={ + "key": key_list, + "key_str": r.get("key_str"), + "new_key_str": r.get("new_key_str"), + "old_status": r.get("old_status"), + "new_status": r.get("new_status"), + }, + old_kind=str(r.get("old_kind") or ""), + new_kind=str(r.get("new_kind") or ""), + old_json=dict(r.get("old") or {}), + new_json=dict(r.get("new") or {}), + in_expect=bool(r.get("in_expect")), + search_text=search[:2000], + ) + ) + db.commit() + db.refresh(run) + return run_to_dict(db, run) + + +def run_to_dict(db: Session, run: BizMigrationRun, *, include_diffs: bool = False) -> dict[str, Any]: + out: dict[str, Any] = { + "id": run.id, + "project_id": run.project_id, + "batch_id": run.batch_id, + "old_batch_id": run.old_batch_id, + "new_batch_id": run.new_batch_id, + "purpose": getattr(run, "purpose", None) or "manual", + "status": run.status, + "summary": dict(run.summary_json or {}), + "message": run.message, + "created_at": run.created_at.isoformat() if run.created_at else None, + "old_batch": _batch_brief(db, run.old_batch_id), + "new_batch": _batch_brief(db, run.new_batch_id), + } + if include_diffs: + diffs = ( + db.query(BizMigrationDiff) + .filter(BizMigrationDiff.run_id == run.id) + .order_by(BizMigrationDiff.seq.asc()) + .limit(5000) + .all() + ) + out["diffs"] = [diff_to_dict(d) for d in diffs] + return out + + +def diff_to_dict(d: BizMigrationDiff) -> dict[str, Any]: + kj = d.key_json if isinstance(d.key_json, dict) else {} + return { + "id": d.id, + "metric_id": d.metric_id, + "seq": d.seq, + "verdict": d.verdict, + "color": d.color, + "key": kj, + "key_str": kj.get("key_str") or "", + "new_key_str": kj.get("new_key_str") or "", + "old_status": kj.get("old_status") or "", + "new_status": kj.get("new_status") or "", + "old_kind": d.old_kind, + "new_kind": d.new_kind, + "old": d.old_json, + "new": d.new_json, + "in_expect": d.in_expect, + } + + +def list_baseline_ports(db: Session, project_id: str) -> dict[str, Any]: + """List interface names from project old baseline for expect-set picking.""" + p = db.get(BizMigrationProject, project_id) + if not p: + raise HTTPException(status_code=404, detail="project_not_found") + if not p.old_baseline_batch_id: + return {"batch_id": "", "ports": [], "mapped": {}} + rows = _load_metric_rows( + db, batch_id=p.old_baseline_batch_id, metric_id=PORT_METRIC_ID + ) + port_map = _port_map_dict(db, p.mapping_id) + ports: list[dict[str, Any]] = [] + for r in rows: + name = str(r.get("interface") or "").strip() + if not name: + continue + ports.append( + { + "interface": name, + "admin": r.get("admin") or "", + "phy": r.get("phy") or "", + "prot": r.get("prot") or "", + "description": r.get("description") or "", + "mapped_to": port_map.get(name) or "", + } + ) + ports.sort(key=lambda x: str(x["interface"])) + return { + "batch_id": p.old_baseline_batch_id, + "metric_id": PORT_METRIC_ID, + "ports": ports, + "mapped": port_map, + } + + +def _iface_brief_item(*, vendor: str, device_type: str) -> dict[str, Any]: + from ..biz_state.profiles import profiles_for_vendor + from ..lldp_shared import resolve_vendor_key + + vkey = resolve_vendor_key(vendor, device_type) + for p in profiles_for_vendor(vkey): + if p.metric_id == PORT_METRIC_ID and p.kind == "collect": + return { + "source_profile_id": p.profile_id, + "kind": "catalog", + "enabled": True, + "title": p.title or "interface brief", + } + raise HTTPException( + status_code=400, + detail=f"no_interface_brief_profile_for_vendor:{vkey or vendor or 'unknown'}", + ) + + +def _enabled_metric_ids(db: Session, task_id: str) -> set[str]: + from ..biz_state.profiles import get_profile + from ..models import BizStateTaskItem + + out: set[str] = set() + items = ( + db.query(BizStateTaskItem) + .filter(BizStateTaskItem.task_id == task_id, BizStateTaskItem.enabled.is_(True)) + .all() + ) + for it in items: + if it.kind == "custom_raw": + out.add("__custom__") + continue + p = get_profile(str(it.source_profile_id or "")) + if p and p.metric_id: + out.add(p.metric_id) + return out + + +def _is_port_highfreq_task(db: Session, task: BizStateTask) -> bool: + """True when task is interface_brief-only with short interval (cutover HF).""" + metrics = _enabled_metric_ids(db, task.id) + return metrics == {PORT_METRIC_ID} and int(task.interval_sec or 0) <= 300 + + +def _ensure_side_highfreq( + db: Session, + *, + template: BizStateTask, + project_name: str, + interval_sec: int, + retention_days: int, +) -> tuple[BizStateTask, bool]: + """Return (task, created). Reuse if already HF port-only; else create sibling.""" + from ..biz_state import service as biz_svc + + if _is_port_highfreq_task(db, template): + if template.status != "running": + biz_svc.update_task(db, template.id, {"status": "running"}) + refreshed = db.get(BizStateTask, template.id) + return refreshed or template, False + return template, False + + item = _iface_brief_item(vendor=template.vendor, device_type=template.device_type) + note = f"割接高频-端口/{project_name}"[:256] + created = biz_svc.create_task( + db, + { + "source": template.source, + "ne_id": template.ne_id, + "ne_name": template.ne_name, + "ne_ip": template.ne_ip, + "vendor": template.vendor, + "device_type": template.device_type, + "note": note, + "status": "running", + "interval_sec": interval_sec, + "retention_days": retention_days, + "items": [item], + }, + ) + task = db.get(BizStateTask, str(created.get("id") or "")) + if not task: + raise HTTPException(status_code=500, detail="highfreq_task_create_failed") + return task, True + + +def ensure_port_highfreq( + db: Session, + project_id: str, + *, + interval_sec: int = 60, + retention_days: int = 7, + collect_now: bool = True, +) -> dict[str, Any]: + """Create/bind interface_brief-only high-freq biz_state tasks for old/new NEs. + + Collection stays in biz_state — migration only points at the tasks. + Later metrics can be added on the same tasks via the biz-state UI. + """ + from ..biz_state.collect_runner import dispatch_collect + + proj = db.get(BizMigrationProject, project_id) + if not proj: + raise HTTPException(status_code=404, detail="project_not_found") + old_tpl = db.get(BizStateTask, proj.old_task_id) if proj.old_task_id else None + new_tpl = db.get(BizStateTask, proj.new_task_id) if proj.new_task_id else None + if not old_tpl or not new_tpl: + raise HTTPException(status_code=400, detail="old_new_task_required") + + iv = max(60, int(interval_sec or 60)) + ret = max(1, int(retention_days or 7)) + old_task, old_created = _ensure_side_highfreq( + db, + template=old_tpl, + project_name=proj.name, + interval_sec=iv, + retention_days=ret, + ) + # re-load templates after possible commits inside ensure + new_tpl = db.get(BizStateTask, proj.new_task_id) + if not new_tpl: + raise HTTPException(status_code=400, detail="new_task_required") + new_task, new_created = _ensure_side_highfreq( + db, + template=new_tpl, + project_name=proj.name, + interval_sec=iv, + retention_days=ret, + ) + + proj = db.get(BizMigrationProject, project_id) + if not proj: + raise HTTPException(status_code=404, detail="project_not_found") + proj.old_task_id = old_task.id + proj.new_task_id = new_task.id + proj.updated_at = utcnow_naive() + db.commit() + + collect: dict[str, Any] = {"old": None, "new": None} + if collect_now: + for side, tid in (("old", old_task.id), ("new", new_task.id)): + try: + dispatch_collect(tid) + collect[side] = {"ok": True, "task_id": tid} + except Exception as exc: # noqa: BLE001 + collect[side] = {"ok": False, "task_id": tid, "error": str(exc)[:200]} + + return { + "project": project_to_dict(db, proj), + "old_task": _task_brief(db, old_task.id), + "new_task": _task_brief(db, new_task.id), + "old_created": old_created, + "new_created": new_created, + "interval_sec": iv, + "collect": collect, + } + + +def collect_project_now(db: Session, project_id: str) -> dict[str, Any]: + """Trigger immediate collect on project's old/new biz_state tasks.""" + from ..biz_state.collect_runner import dispatch_collect + + proj = db.get(BizMigrationProject, project_id) + if not proj: + raise HTTPException(status_code=404, detail="project_not_found") + out: dict[str, Any] = {"old": None, "new": None} + for side, tid in (("old", proj.old_task_id), ("new", proj.new_task_id)): + if not tid: + out[side] = {"ok": False, "error": "task_missing"} + continue + task = db.get(BizStateTask, tid) + if not task: + out[side] = {"ok": False, "error": "task_not_found"} + continue + if bool(task.collect_running): + out[side] = {"ok": False, "error": "collect_already_running", "task_id": tid} + continue + try: + dispatch_collect(tid) + out[side] = {"ok": True, "task_id": tid} + except Exception as exc: # noqa: BLE001 + out[side] = {"ok": False, "task_id": tid, "error": str(exc)[:200]} + return out + + +def _red_ticket_to_dict(t: BizMigrationRedTicket) -> dict[str, Any]: + return { + "id": t.id, + "project_id": t.project_id, + "batch_id": t.batch_id, + "run_id": t.run_id, + "metric_id": t.metric_id, + "key_str": t.key_str, + "new_key_str": t.new_key_str, + "verdict": t.verdict, + "color": t.color, + "old_status": t.old_status, + "new_status": t.new_status, + "detail": dict(t.detail_json or {}), + "status": t.status, + "carried_to_batch_id": t.carried_to_batch_id, + "note": t.note, + "created_at": t.created_at.isoformat() if t.created_at else None, + "resolved_at": t.resolved_at.isoformat() if t.resolved_at else None, + } + + +def count_open_red_tickets(db: Session, project_id: str) -> int: + return ( + db.query(BizMigrationRedTicket) + .filter( + BizMigrationRedTicket.project_id == project_id, + BizMigrationRedTicket.status.in_(("open", "carried")), + ) + .count() + ) + + +def list_red_tickets( + db: Session, + project_id: str, + *, + status: str = "", + limit: int = 200, +) -> dict[str, Any]: + if not db.get(BizMigrationProject, project_id): + raise HTTPException(status_code=404, detail="project_not_found") + q = db.query(BizMigrationRedTicket).filter(BizMigrationRedTicket.project_id == project_id) + if status.strip(): + q = q.filter(BizMigrationRedTicket.status == status.strip()) + rows = ( + q.order_by(BizMigrationRedTicket.created_at.desc()) + .limit(min(500, max(1, limit))) + .all() + ) + return { + "open_count": count_open_red_tickets(db, project_id), + "items": [_red_ticket_to_dict(t) for t in rows], + } + + +def resolve_red_ticket(db: Session, ticket_id: str, *, note: str = "") -> dict[str, Any]: + t = db.get(BizMigrationRedTicket, ticket_id) + if not t: + raise HTTPException(status_code=404, detail="red_ticket_not_found") + t.status = "resolved" + t.resolved_at = utcnow_naive() + if note: + t.note = str(note)[:500] + db.commit() + db.refresh(t) + return _red_ticket_to_dict(t) + + +def _persist_red_tickets_from_run( + db: Session, + *, + project_id: str, + batch_id: str, + run_id: str, +) -> list[BizMigrationRedTicket]: + """Create open red tickets from acceptance run red diffs (expect + anomaly).""" + diffs = ( + db.query(BizMigrationDiff) + .filter(BizMigrationDiff.run_id == run_id, BizMigrationDiff.color == "red") + .order_by(BizMigrationDiff.seq.asc()) + .all() + ) + created: list[BizMigrationRedTicket] = [] + for d in diffs: + kj = d.key_json if isinstance(d.key_json, dict) else {} + t = BizMigrationRedTicket( + id=uuid4().hex, + project_id=project_id, + batch_id=batch_id, + run_id=run_id, + metric_id=d.metric_id, + key_str=str(kj.get("key_str") or "")[:256], + new_key_str=str(kj.get("new_key_str") or "")[:256], + verdict=d.verdict, + color=d.color or "red", + old_status=str(kj.get("old_status") or "")[:64], + new_status=str(kj.get("new_status") or "")[:64], + detail_json={ + "old_kind": d.old_kind, + "new_kind": d.new_kind, + "in_expect": d.in_expect, + "old": d.old_json, + "new": d.new_json, + }, + status="open", + ) + db.add(t) + created.append(t) + return created + + +def finish_batch( + db: Session, + batch_id: str, + *, + mark_done: bool = False, +) -> dict[str, Any]: + """本批完成:关窗 → 终验 → 验收小结 → 红单留痕(不硬卡下一批).""" + mb = db.get(BizMigrationBatch, batch_id) + if not mb: + raise HTTPException(status_code=404, detail="batch_not_found") + if mb.status not in ("active", "review"): + raise HTTPException( + status_code=400, + detail=f"batch_not_finishable:{mb.status}", + ) + + # Close window first so acceptance uses non-active rules + if mb.status == "active": + mb.ended_at = utcnow_naive() + mb.status = "done" if mark_done else "review" + mb.updated_at = utcnow_naive() + db.commit() + + run = run_evaluate(db, batch_id=batch_id, acceptance=True, purpose="acceptance") + summary = dict(run.get("summary") or {}) + progress = dict(summary.get("progress") or {}) + anomaly = int(summary.get("anomaly") or 0) + ok = int(progress.get("ok") or 0) + total = int(progress.get("total") or 0) + passed = anomaly == 0 and (total == 0 or ok >= total) + + accept_summary = { + "passed": passed, + "progress_ok": ok, + "progress_total": total, + "anomaly": anomaly, + "verdict_counts": dict(summary.get("verdict_counts") or {}), + "expect_ports": list(summary.get("expect_ports") or []), + "run_id": run.get("id"), + "old_batch_id": run.get("old_batch_id"), + "new_batch_id": run.get("new_batch_id"), + "finished_at": utcnow_naive().isoformat(), + } + + mb = db.get(BizMigrationBatch, batch_id) + if not mb: + raise HTTPException(status_code=404, detail="batch_not_found") + mb.accept_run_id = str(run.get("id") or "") + mb.accept_status = "passed" if passed else "failed" + mb.accept_summary_json = accept_summary + mb.updated_at = utcnow_naive() + + # Replace prior open tickets from this batch's previous acceptance (re-finish) + db.query(BizMigrationRedTicket).filter( + BizMigrationRedTicket.batch_id == batch_id, + BizMigrationRedTicket.status == "open", + ).delete() + reds = _persist_red_tickets_from_run( + db, + project_id=mb.project_id, + batch_id=batch_id, + run_id=str(run.get("id") or ""), + ) + db.commit() + db.refresh(mb) + + return { + "batch": batch_to_dict(mb), + "run": run, + "accept_summary": accept_summary, + "red_tickets": [_red_ticket_to_dict(t) for t in reds], + "open_red_count": count_open_red_tickets(db, mb.project_id), + "can_continue": True, # 带红继续:永不硬卡 + } + + +def patch_batch(db: Session, batch_id: str, body: dict[str, Any]) -> dict[str, Any]: + b = db.get(BizMigrationBatch, batch_id) + if not b: + raise HTTPException(status_code=404, detail="batch_not_found") + if "batch_label" in body and body["batch_label"] is not None: + b.batch_label = str(body["batch_label"]).strip() or b.batch_label + if "note" in body and body["note"] is not None: + b.note = str(body["note"])[:500] + if "expect_set" in body and isinstance(body["expect_set"], dict): + b.expect_set_json = dict(body["expect_set"]) + if "status" in body and body["status"] is not None: + st = str(body["status"]).strip() + # Prefer finish_batch for review/done from active (auto acceptance) + if st in ("review", "done") and b.status == "active": + return finish_batch(db, batch_id, mark_done=(st == "done"))["batch"] + prev = b.status + b.status = st or b.status + if st == "active" and prev != "active": + b.started_at = utcnow_naive() + # 带红:标记既有 open 红单为 carried(不关闭) + open_reds = ( + db.query(BizMigrationRedTicket) + .filter( + BizMigrationRedTicket.project_id == b.project_id, + BizMigrationRedTicket.status == "open", + ) + .all() + ) + for t in open_reds: + t.status = "carried" + t.carried_to_batch_id = b.id + b.updated_at = utcnow_naive() + db.commit() + db.refresh(b) + return batch_to_dict(b) + + +def get_run(db: Session, run_id: str) -> dict[str, Any]: + run = db.get(BizMigrationRun, run_id) + if not run: + raise HTTPException(status_code=404, detail="run_not_found") + return run_to_dict(db, run, include_diffs=False) + + +def list_run_diffs( + db: Session, + run_id: str, + *, + metric_id: str = "", + verdict: str = "", + color: str = "", + kw: str = "", + offset: int = 0, + limit: int = 100, +) -> dict[str, Any]: + if not db.get(BizMigrationRun, run_id): + raise HTTPException(status_code=404, detail="run_not_found") + q = db.query(BizMigrationDiff).filter(BizMigrationDiff.run_id == run_id) + if metric_id: + q = q.filter(BizMigrationDiff.metric_id == metric_id) + if verdict: + q = q.filter(BizMigrationDiff.verdict == verdict) + if color: + q = q.filter(BizMigrationDiff.color == color) + if kw.strip(): + like = f"%{kw.strip()}%" + q = q.filter(BizMigrationDiff.search_text.ilike(like)) + total = q.count() + rows = q.order_by(BizMigrationDiff.seq.asc()).offset(max(0, offset)).limit(min(500, max(1, limit))).all() + return { + "total": total, + "offset": offset, + "limit": limit, + "items": [diff_to_dict(d) for d in rows], + } + + +def list_runs(db: Session, batch_id: str, limit: int = 20) -> list[dict[str, Any]]: + rows = ( + db.query(BizMigrationRun) + .filter(BizMigrationRun.batch_id == batch_id) + .order_by(BizMigrationRun.created_at.desc()) + .limit(min(100, max(1, limit))) + .all() + ) + return [run_to_dict(db, r) for r in rows] + + +def board(db: Session, batch_id: str, run_id: str = "") -> dict[str, Any]: + """Progress / anomaly board for a migration batch (latest or given run).""" + mb = db.get(BizMigrationBatch, batch_id) + if not mb: + raise HTTPException(status_code=404, detail="batch_not_found") + proj = get_project(db, mb.project_id) + run = None + if run_id: + run = db.get(BizMigrationRun, run_id) + if not run or run.batch_id != batch_id: + raise HTTPException(status_code=404, detail="run_not_found") + else: + run = ( + db.query(BizMigrationRun) + .filter(BizMigrationRun.batch_id == batch_id) + .order_by(BizMigrationRun.created_at.desc()) + .first() + ) + return { + "project": proj, + "batch": batch_to_dict(mb), + "run": run_to_dict(db, run) if run else None, + } diff --git a/netx_api/biz_migration_router.py b/netx_api/biz_migration_router.py new file mode 100644 index 0000000..a12a7f1 --- /dev/null +++ b/netx_api/biz_migration_router.py @@ -0,0 +1,211 @@ +"""HTTP API for cutover migration monitor (/v1/biz-migration).""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter, Depends +from pydantic import BaseModel, Field +from sqlalchemy.orm import Session + +from .biz_migration import service as svc +from .db import get_db + +router = APIRouter(prefix="/v1/biz-migration", tags=["biz-migration"]) + + +class ProjectIn(BaseModel): + name: str + old_task_id: str + new_task_id: str + old_baseline_batch_id: str = "" + new_baseline_batch_id: str = "" + mapping_id: str = "" + status: str = "draft" + note: str = "" + + +class ProjectPatchIn(BaseModel): + name: str | None = None + old_baseline_batch_id: str | None = None + new_baseline_batch_id: str | None = None + mapping_id: str | None = None + status: str | None = None + note: str | None = None + + +class BatchIn(BaseModel): + batch_label: str = "batch" + expect_set: dict[str, Any] = Field(default_factory=dict) + status: str = "pending" + note: str = "" + + +class BatchPatchIn(BaseModel): + batch_label: str | None = None + expect_set: dict[str, Any] | None = None + status: str | None = None + note: str | None = None + + +class EvaluateIn(BaseModel): + old_batch_id: str = "" + new_batch_id: str = "" + + +@router.get("/projects") +def api_list_projects(db: Session = Depends(get_db)): + return {"items": svc.list_projects(db)} + + +@router.post("/projects") +def api_create_project(body: ProjectIn, db: Session = Depends(get_db)): + return svc.create_project(db, body.model_dump()) + + +@router.get("/projects/{project_id}") +def api_get_project(project_id: str, db: Session = Depends(get_db)): + return svc.get_project(db, project_id) + + +@router.patch("/projects/{project_id}") +def api_patch_project(project_id: str, body: ProjectPatchIn, db: Session = Depends(get_db)): + return svc.patch_project(db, project_id, body.model_dump(exclude_unset=True)) + + +@router.delete("/projects/{project_id}") +def api_delete_project(project_id: str, db: Session = Depends(get_db)): + return svc.delete_project(db, project_id) + + +@router.get("/projects/{project_id}/baseline-ports") +def api_baseline_ports(project_id: str, db: Session = Depends(get_db)): + """List old-baseline interface_brief ports for expect-set selection.""" + return svc.list_baseline_ports(db, project_id) + + +class EnsureHighfreqIn(BaseModel): + interval_sec: int = 60 + retention_days: int = 7 + collect_now: bool = True + + +@router.post("/projects/{project_id}/ensure-port-highfreq") +def api_ensure_port_highfreq( + project_id: str, + body: EnsureHighfreqIn | None = None, + db: Session = Depends(get_db), +): + """Create/bind biz_state high-freq interface_brief tasks for old/new NEs.""" + payload = body.model_dump() if body else {} + return svc.ensure_port_highfreq( + db, + project_id, + interval_sec=int(payload.get("interval_sec") or 60), + retention_days=int(payload.get("retention_days") or 7), + collect_now=bool(payload.get("collect_now", True)), + ) + + +@router.post("/projects/{project_id}/collect-now") +def api_collect_now(project_id: str, db: Session = Depends(get_db)): + """Trigger immediate collect on project's bound biz_state tasks.""" + return svc.collect_project_now(db, project_id) + + +@router.get("/projects/{project_id}/batches") +def api_list_batches(project_id: str, db: Session = Depends(get_db)): + return {"items": svc.list_batches(db, project_id)} + + +@router.post("/projects/{project_id}/batches") +def api_create_batch(project_id: str, body: BatchIn, db: Session = Depends(get_db)): + return svc.create_batch(db, project_id, body.model_dump()) + + +@router.post("/batches/{batch_id}/finish") +def api_finish_batch( + batch_id: str, + mark_done: bool = False, + db: Session = Depends(get_db), +): + """本批完成:终验 + 小结 + 红单(允许带红继续).""" + return svc.finish_batch(db, batch_id, mark_done=mark_done) + + +@router.get("/projects/{project_id}/red-tickets") +def api_list_red_tickets( + project_id: str, + status: str = "", + limit: int = 200, + db: Session = Depends(get_db), +): + return svc.list_red_tickets(db, project_id, status=status, limit=limit) + + +class RedTicketPatchIn(BaseModel): + note: str = "" + + +@router.post("/red-tickets/{ticket_id}/resolve") +def api_resolve_red_ticket( + ticket_id: str, + body: RedTicketPatchIn | None = None, + db: Session = Depends(get_db), +): + note = body.note if body else "" + return svc.resolve_red_ticket(db, ticket_id, note=note) + + +@router.patch("/batches/{batch_id}") +def api_patch_batch(batch_id: str, body: BatchPatchIn, db: Session = Depends(get_db)): + return svc.patch_batch(db, batch_id, body.model_dump(exclude_unset=True)) + + +@router.post("/batches/{batch_id}/evaluate") +def api_evaluate(batch_id: str, body: EvaluateIn | None = None, db: Session = Depends(get_db)): + payload = body.model_dump() if body else {} + return svc.run_evaluate( + db, + batch_id=batch_id, + old_batch_id=str(payload.get("old_batch_id") or ""), + new_batch_id=str(payload.get("new_batch_id") or ""), + ) + + +@router.get("/batches/{batch_id}/runs") +def api_list_runs(batch_id: str, limit: int = 20, db: Session = Depends(get_db)): + return {"items": svc.list_runs(db, batch_id, limit=limit)} + + +@router.get("/batches/{batch_id}/board") +def api_board(batch_id: str, run_id: str = "", db: Session = Depends(get_db)): + return svc.board(db, batch_id, run_id=run_id) + + +@router.get("/runs/{run_id}") +def api_get_run(run_id: str, db: Session = Depends(get_db)): + return svc.get_run(db, run_id) + + +@router.get("/runs/{run_id}/diffs") +def api_list_diffs( + run_id: str, + metric_id: str = "", + verdict: str = "", + color: str = "", + kw: str = "", + offset: int = 0, + limit: int = 100, + db: Session = Depends(get_db), +): + return svc.list_run_diffs( + db, + run_id, + metric_id=metric_id, + verdict=verdict, + color=color, + kw=kw, + offset=offset, + limit=limit, + ) diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index 37a1b9d..eba8379 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -272,7 +272,6 @@ def dispatch_collect(task_id: str) -> None: device_type = str(task.device_type or "") source = str(task.source or "managed").strip().lower() ne_id = str(task.ne_id or "").strip() - retention = int(task.retention_batches or 30) finally: db.close() @@ -288,7 +287,6 @@ def dispatch_collect(task_id: str) -> None: ne_id=ne_id, vendor=vendor, device_type=device_type, - retention=retention, ) except Exception as exc: _log.exception("biz_state collect failed task=%s", task_id) @@ -315,7 +313,6 @@ def _run_collect_session( ne_id: str, vendor: str, device_type: str, - retention: int, ) -> None: per_cmd = int(settings.ne_collect_read_timeout_sec or 120) cap = int(settings.ne_collect_run_timeout_cap_sec or 600) @@ -678,29 +675,29 @@ def _run_collect_session( except Exception: _log.exception("biz_state auto compare hook failed task=%s", task_id) - _purge_old_batches(db, task_id=task_id, keep=retention) + _purge_task_retention(db, task_id=task_id) finally: db.close() -def _purge_old_batches(db, *, task_id: str, keep: int) -> None: - keep_n = max(1, int(keep or 30)) - rows = ( - db.query(BizStateBatch) - .filter(BizStateBatch.task_id == task_id) - .order_by(BizStateBatch.started_at.desc()) - .all() - ) - drop = rows[keep_n:] - for b in drop: - bid = b.id - db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == bid).delete() - db.query(BizStateVrfRouteSummary).filter(BizStateVrfRouteSummary.batch_id == bid).delete() - db.query(BizStateMetricRow).filter(BizStateMetricRow.batch_id == bid).delete() - db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == bid).delete() - db.delete(b) - if drop: - db.commit() +def _purge_task_retention(db, *, task_id: str) -> None: + from .retention import purge_task_batches + + task = db.get(BizStateTask, task_id) + if not task: + return + try: + info = purge_task_batches(db, task) + if info.get("dropped"): + _log.info( + "biz_state retention purged task=%s dropped=%s days=%s daily=%s", + task_id, + info.get("dropped"), + info.get("retention_days"), + info.get("daily_keep_enabled"), + ) + except Exception: + _log.exception("biz_state retention purge failed task=%s", task_id) def trigger_collect_now(task_id: str) -> dict[str, Any]: diff --git a/netx_api/biz_state/retention.py b/netx_api/biz_state/retention.py new file mode 100644 index 0000000..d14c554 --- /dev/null +++ b/netx_api/biz_state/retention.py @@ -0,0 +1,168 @@ +"""Biz-state batch retention: age/day policies + reference protection.""" + +from __future__ import annotations + +from collections import defaultdict +from datetime import datetime, timedelta +from typing import Any + +from sqlalchemy.orm import Session + +from ..models import ( + BizCompareJob, + BizCompareRun, + BizMigrationProject, + BizMigrationRun, + BizStateBatch, + BizStateBatchCommand, + BizStateLldpNeighbor, + BizStateMetricRow, + BizStateTask, + BizStateVrfRouteSummary, +) + + +def _utcnow() -> datetime: + return datetime.utcnow() + + +def protected_batch_map(db: Session, *, task_id: str = "") -> dict[str, list[str]]: + """batch_id → reason codes. Empty task_id = all tasks.""" + out: dict[str, list[str]] = defaultdict(list) + + def add(bid: str, reason: str) -> None: + b = str(bid or "").strip() + if not b: + return + if reason not in out[b]: + out[b].append(reason) + + q = db.query(BizStateBatch) + if task_id: + q = q.filter(BizStateBatch.task_id == task_id) + for b in q.filter(BizStateBatch.is_baseline.is_(True)).all(): + add(b.id, "manual_baseline") + + for j in db.query(BizCompareJob).all(): + add(j.before_batch_id, "compare_job_before") + add(j.after_batch_id, "compare_job_after") + + for r in db.query(BizCompareRun).all(): + add(r.before_batch_id, "compare_run_before") + add(r.after_batch_id, "compare_run_after") + + for p in db.query(BizMigrationProject).all(): + add(p.old_baseline_batch_id, "migration_old_baseline") + add(p.new_baseline_batch_id, "migration_new_baseline") + + for r in db.query(BizMigrationRun).all(): + add(r.old_batch_id, "migration_run_old") + add(r.new_batch_id, "migration_run_new") + + return dict(out) + + +def protected_batch_ids(db: Session, *, task_id: str = "") -> set[str]: + return set(protected_batch_map(db, task_id=task_id).keys()) + + +def batch_protect_info(db: Session, batch_id: str) -> dict[str, Any]: + reasons = protected_batch_map(db).get(batch_id, []) + b = db.get(BizStateBatch, batch_id) + if b and bool(getattr(b, "is_baseline", False)) and "manual_baseline" not in reasons: + reasons = ["manual_baseline", *reasons] + return { + "protected": bool(reasons), + "reasons": reasons, + "is_baseline": bool(b and getattr(b, "is_baseline", False)), + } + + +def delete_batch_data(db: Session, batch_id: str) -> None: + """Hard-delete one batch and child rows (caller must check protection).""" + bid = str(batch_id or "").strip() + if not bid: + return + db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == bid).delete() + db.query(BizStateVrfRouteSummary).filter(BizStateVrfRouteSummary.batch_id == bid).delete() + db.query(BizStateMetricRow).filter(BizStateMetricRow.batch_id == bid).delete() + db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == bid).delete() + b = db.get(BizStateBatch, bid) + if b: + db.delete(b) + + +def purge_task_batches(db: Session, task: BizStateTask) -> dict[str, Any]: + """Apply retention_days + optional daily_keep policy. Never deletes protected batches. + + Returns counts for logging/UI. + """ + task_id = str(task.id) + retention_days = max(1, int(getattr(task, "retention_days", None) or 30)) + daily_on = bool(getattr(task, "daily_keep_enabled", False)) + daily_n = max(1, int(getattr(task, "daily_keep_count", None) or 10)) + + protected = protected_batch_ids(db, task_id=task_id) + rows = ( + db.query(BizStateBatch) + .filter(BizStateBatch.task_id == task_id) + .order_by(BizStateBatch.started_at.desc()) + .all() + ) + now = _utcnow() + cutoff = now - timedelta(days=retention_days) + today = now.date() + + to_drop: set[str] = set() + + # 1) Age: older than retention_days + for b in rows: + if b.id in protected: + continue + started = b.started_at or now + if started < cutoff: + to_drop.add(b.id) + + # 2) Daily keep (default off): for each past calendar day, keep N newest + if daily_on: + by_day: dict[Any, list[BizStateBatch]] = defaultdict(list) + for b in rows: + started = b.started_at or now + d = started.date() + if d >= today: + continue # 当天的多余留到第二天再清 + by_day[d].append(b) + for _day, day_rows in by_day.items(): + # already ordered desc globally; re-sort + day_rows.sort(key=lambda x: x.started_at or now, reverse=True) + for b in day_rows[daily_n:]: + if b.id not in protected: + to_drop.add(b.id) + + dropped = 0 + for bid in to_drop: + delete_batch_data(db, bid) + dropped += 1 + if dropped: + db.commit() + return { + "task_id": task_id, + "dropped": dropped, + "protected": len(protected), + "retention_days": retention_days, + "daily_keep_enabled": daily_on, + "daily_keep_count": daily_n, + } + + +def purge_all_tasks(db: Session) -> dict[str, Any]: + """Periodic sweep for all tasks (HF idle tasks included).""" + tasks = db.query(BizStateTask).all() + total = 0 + details: list[dict[str, Any]] = [] + for t in tasks: + info = purge_task_batches(db, t) + total += int(info.get("dropped") or 0) + if info.get("dropped"): + details.append(info) + return {"dropped": total, "tasks": details} diff --git a/netx_api/biz_state/schema_ensure.py b/netx_api/biz_state/schema_ensure.py index 9ceade8..99e2a37 100644 --- a/netx_api/biz_state/schema_ensure.py +++ b/netx_api/biz_state/schema_ensure.py @@ -29,6 +29,25 @@ def apply_biz_state_schema(conn: Connection) -> None: "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_id ON biz_compare_diff (run_id)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_metric_kind ON biz_compare_diff (run_id, metric_id, kind)", "CREATE INDEX IF NOT EXISTS ix_biz_compare_diff_run_metric_seq ON biz_compare_diff (run_id, metric_id, seq)", + "CREATE INDEX IF NOT EXISTS ix_biz_migration_project_status ON biz_migration_project (status)", + "CREATE INDEX IF NOT EXISTS ix_biz_migration_batch_project_id ON biz_migration_batch (project_id)", + "CREATE INDEX IF NOT EXISTS ix_biz_migration_run_batch_id ON biz_migration_run (batch_id)", + "CREATE INDEX IF NOT EXISTS ix_biz_migration_diff_run_id ON biz_migration_diff (run_id)", + # Allow multiple biz_state tasks per NE (cutover high-freq + full) + "ALTER TABLE biz_state_task DROP CONSTRAINT IF EXISTS uq_biz_state_task_ne", + "DROP INDEX IF EXISTS uq_biz_state_task_ne", + "CREATE INDEX IF NOT EXISTS ix_biz_state_task_source_ne ON biz_state_task (source, ne_id)", + "ALTER TABLE biz_state_task ADD COLUMN IF NOT EXISTS retention_days INTEGER DEFAULT 30", + "ALTER TABLE biz_state_task ADD COLUMN IF NOT EXISTS daily_keep_enabled BOOLEAN DEFAULT FALSE", + "ALTER TABLE biz_state_task ADD COLUMN IF NOT EXISTS daily_keep_count INTEGER DEFAULT 10", + "ALTER TABLE biz_state_batch ADD COLUMN IF NOT EXISTS is_baseline BOOLEAN DEFAULT FALSE", + "ALTER TABLE biz_state_batch ADD COLUMN IF NOT EXISTS baseline_marked_at TIMESTAMP", + "CREATE INDEX IF NOT EXISTS ix_biz_state_batch_is_baseline ON biz_state_batch (is_baseline)", + "ALTER TABLE biz_migration_batch ADD COLUMN IF NOT EXISTS accept_status VARCHAR(32) DEFAULT 'none'", + "ALTER TABLE biz_migration_batch ADD COLUMN IF NOT EXISTS accept_run_id VARCHAR(64) DEFAULT ''", + "ALTER TABLE biz_migration_batch ADD COLUMN IF NOT EXISTS accept_summary_json JSON DEFAULT '{}'", + "ALTER TABLE biz_migration_run ADD COLUMN IF NOT EXISTS purpose VARCHAR(32) DEFAULT 'manual'", + "CREATE INDEX IF NOT EXISTS ix_biz_migration_red_project_status ON biz_migration_red_ticket (project_id, status)", ): try: _run_sql(conn, sql) diff --git a/netx_api/biz_state/service.py b/netx_api/biz_state/service.py index 84d5aee..63a0d22 100644 --- a/netx_api/biz_state/service.py +++ b/netx_api/biz_state/service.py @@ -28,6 +28,12 @@ from ..models import ( from ..timeutil import utcnow_naive from .command_match import preview_task_item from .profiles import all_profiles, get_profile, profile_to_public_dict, profiles_for_vendor +from .retention import ( + batch_protect_info, + delete_batch_data, + protected_batch_map, + purge_task_batches, +) def _utcnow() -> datetime: @@ -134,17 +140,13 @@ def create_task(db: Session, body: dict[str, Any]) -> dict[str, Any]: ne_id = str(body.get("ne_id") or "").strip() if not ne_id: raise HTTPException(status_code=400, detail="ne_id_required") - existing = ( - db.query(BizStateTask) - .filter(BizStateTask.source == source, BizStateTask.ne_id == ne_id) - .one_or_none() - ) - if existing: - raise HTTPException(status_code=409, detail="task_already_exists_for_ne") meta = _ne_meta(db, source=source, ne_id=ne_id) vendor = str(body.get("vendor") or meta["vendor"] or "") device_type = str(body.get("device_type") or meta["device_type"] or "") + status = str(body.get("status") or "draft").strip() or "draft" + if status not in ("draft", "running", "paused", "stopped"): + status = "draft" task = BizStateTask( id=uuid4().hex, source=source, @@ -154,9 +156,13 @@ def create_task(db: Session, body: dict[str, Any]) -> dict[str, Any]: vendor=vendor, device_type=device_type, note=str(body.get("note") or "")[:256], - status="draft", + status=status, interval_sec=max(60, int(body.get("interval_sec") or 3600)), - retention_batches=max(1, int(body.get("retention_batches") or 30)), + retention_days=max(1, min(3650, int(body.get("retention_days") or 30))), + daily_keep_enabled=bool(body.get("daily_keep_enabled") or False), + daily_keep_count=max(1, min(1000, int(body.get("daily_keep_count") or 10))), + # keep legacy column in sync for brownfield readers + retention_batches=max(1, int(body.get("retention_days") or body.get("retention_batches") or 30)), created_at=_utcnow(), updated_at=_utcnow(), ) @@ -230,8 +236,17 @@ def update_task(db: Session, task_id: str, body: dict[str, Any]) -> dict[str, An task.note = str(body.get("note") or "")[:256] if "interval_sec" in body: task.interval_sec = max(60, int(body.get("interval_sec") or 3600)) - if "retention_batches" in body: - task.retention_batches = max(1, int(body.get("retention_batches") or 30)) + if "retention_days" in body: + task.retention_days = max(1, min(3650, int(body.get("retention_days") or 30))) + task.retention_batches = task.retention_days # legacy mirror + elif "retention_batches" in body: + # backward compat: treat as days if old clients still send it + task.retention_days = max(1, min(3650, int(body.get("retention_batches") or 30))) + task.retention_batches = task.retention_days + if "daily_keep_enabled" in body: + task.daily_keep_enabled = bool(body.get("daily_keep_enabled")) + if "daily_keep_count" in body: + task.daily_keep_count = max(1, min(1000, int(body.get("daily_keep_count") or 10))) if "items" in body: _replace_items(db, task.id, list(body.get("items") or [])) if "status" in body: @@ -337,7 +352,9 @@ def get_task(db: Session, task_id: str) -> dict[str, Any]: "note": task.note, "status": task.status, "interval_sec": task.interval_sec, - "retention_batches": task.retention_batches, + "retention_days": int(getattr(task, "retention_days", None) or 30), + "daily_keep_enabled": bool(getattr(task, "daily_keep_enabled", False)), + "daily_keep_count": int(getattr(task, "daily_keep_count", None) or 10), "collect_running": bool(task.collect_running), "last_collect_started_at": task.last_collect_started_at.isoformat() + "Z" if task.last_collect_started_at @@ -360,6 +377,7 @@ def list_tasks(db: Session) -> list[dict[str, Any]]: "ne_name": t.ne_name, "ne_ip": t.ne_ip, "vendor": t.vendor, + "note": t.note, "status": t.status, "interval_sec": t.interval_sec, "collect_running": bool(t.collect_running), @@ -377,11 +395,20 @@ def delete_task(db: Session, task_id: str) -> None: if not task: raise HTTPException(status_code=404, detail="task_not_found") batches = db.query(BizStateBatch).filter(BizStateBatch.task_id == task_id).all() + # Refuse if any batch is still referenced by compare/migration (manual baseline is OK to drop with task) + pmap = protected_batch_map(db, task_id=task_id) + blocked: list[dict[str, Any]] = [] for b in batches: - db.query(BizStateLldpNeighbor).filter(BizStateLldpNeighbor.batch_id == b.id).delete() - db.query(BizStateVrfRouteSummary).filter(BizStateVrfRouteSummary.batch_id == b.id).delete() - db.query(BizStateBatchCommand).filter(BizStateBatchCommand.batch_id == b.id).delete() - db.delete(b) + reasons = [r for r in pmap.get(b.id, []) if r != "manual_baseline"] + if reasons: + blocked.append({"batch_id": b.id, "reasons": reasons}) + if blocked: + raise HTTPException( + status_code=409, + detail={"error": "batches_referenced", "items": blocked[:20]}, + ) + for b in batches: + delete_batch_data(db, b.id) items = db.query(BizStateTaskItem).filter(BizStateTaskItem.task_id == task_id).all() for it in items: db.query(BizStateTaskItemBinding).filter(BizStateTaskItemBinding.item_id == it.id).delete() @@ -391,28 +418,101 @@ def delete_task(db: Session, task_id: str) -> None: db.commit() +def _batch_list_item(b: BizStateBatch, protect: dict[str, Any]) -> dict[str, Any]: + return { + "id": b.id, + "status": b.status, + "command_count": b.command_count, + "row_count": b.row_count, + "message": b.message, + "ne_name": b.ne_name or "", + "ne_id": b.ne_id or "", + "started_at": b.started_at.isoformat() + "Z" if b.started_at else None, + "ended_at": b.ended_at.isoformat() + "Z" if b.ended_at else None, + "is_baseline": bool(getattr(b, "is_baseline", False)), + "baseline_marked_at": b.baseline_marked_at.isoformat() + "Z" + if getattr(b, "baseline_marked_at", None) + else None, + "protected": bool(protect.get("protected")), + "protect_reasons": list(protect.get("reasons") or []), + } + + def list_batches(db: Session, task_id: str, *, limit: int = 50) -> list[dict[str, Any]]: rows = ( db.query(BizStateBatch) .filter(BizStateBatch.task_id == task_id) .order_by(BizStateBatch.started_at.desc()) - .limit(max(1, min(200, int(limit)))) + .limit(max(1, min(500, int(limit)))) .all() ) - return [ - { - "id": b.id, - "status": b.status, - "command_count": b.command_count, - "row_count": b.row_count, - "message": b.message, - "ne_name": b.ne_name or "", - "ne_id": b.ne_id or "", - "started_at": b.started_at.isoformat() + "Z" if b.started_at else None, - "ended_at": b.ended_at.isoformat() + "Z" if b.ended_at else None, - } - for b in rows - ] + pmap = protected_batch_map(db, task_id=task_id) + out: list[dict[str, Any]] = [] + for b in rows: + reasons = list(pmap.get(b.id, [])) + if bool(getattr(b, "is_baseline", False)) and "manual_baseline" not in reasons: + reasons = ["manual_baseline", *reasons] + out.append( + _batch_list_item( + b, + {"protected": bool(reasons), "reasons": reasons}, + ) + ) + return out + + +def set_batch_baseline(db: Session, batch_id: str, *, marked: bool) -> dict[str, Any]: + b = db.get(BizStateBatch, batch_id) + if not b: + raise HTTPException(status_code=404, detail="batch_not_found") + b.is_baseline = bool(marked) + b.baseline_marked_at = _utcnow() if marked else None + db.commit() + return _batch_list_item(b, batch_protect_info(db, batch_id)) + + +def delete_batch(db: Session, batch_id: str) -> dict[str, Any]: + b = db.get(BizStateBatch, batch_id) + if not b: + raise HTTPException(status_code=404, detail="batch_not_found") + info = batch_protect_info(db, batch_id) + if info.get("protected"): + raise HTTPException( + status_code=409, + detail={"error": "batch_protected", "reasons": info.get("reasons") or []}, + ) + delete_batch_data(db, batch_id) + db.commit() + return {"ok": True, "batch_id": batch_id} + + +def delete_batches_bulk(db: Session, batch_ids: list[str]) -> dict[str, Any]: + deleted: list[str] = [] + skipped: list[dict[str, Any]] = [] + for raw in batch_ids: + bid = str(raw or "").strip() + if not bid: + continue + b = db.get(BizStateBatch, bid) + if not b: + skipped.append({"batch_id": bid, "reasons": ["not_found"]}) + continue + info = batch_protect_info(db, bid) + if info.get("protected"): + skipped.append({"batch_id": bid, "reasons": info.get("reasons") or []}) + continue + delete_batch_data(db, bid) + deleted.append(bid) + if deleted: + db.commit() + return {"ok": True, "deleted": deleted, "skipped": skipped, "deleted_count": len(deleted)} + + +def run_purge_for_task(db: Session, task_id: str) -> dict[str, Any]: + task = db.get(BizStateTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + return purge_task_batches(db, task) def get_batch(db: Session, batch_id: str) -> dict[str, Any]: @@ -454,6 +554,7 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: for r in metric_rows: mid = str(r.metric_id or "") metrics_by_id.setdefault(mid, []).append(dict(r.data_json or {})) + protect = batch_protect_info(db, batch_id) return { "id": b.id, "task_id": b.task_id, @@ -463,6 +564,12 @@ def get_batch(db: Session, batch_id: str) -> dict[str, Any]: "message": b.message, "started_at": b.started_at.isoformat() + "Z" if b.started_at else None, "ended_at": b.ended_at.isoformat() + "Z" if b.ended_at else None, + "is_baseline": bool(getattr(b, "is_baseline", False)), + "baseline_marked_at": b.baseline_marked_at.isoformat() + "Z" + if getattr(b, "baseline_marked_at", None) + else None, + "protected": bool(protect.get("protected")), + "protect_reasons": list(protect.get("reasons") or []), "commands": [ { "id": c.id, diff --git a/netx_api/biz_state_router.py b/netx_api/biz_state_router.py index 42c0c8d..0568c58 100644 --- a/netx_api/biz_state_router.py +++ b/netx_api/biz_state_router.py @@ -44,19 +44,32 @@ class TaskCreateIn(BaseModel): vendor: str = "" device_type: str = "" note: str = "" + status: str = "draft" interval_sec: int = 3600 - retention_batches: int = 30 + retention_days: int = 30 + daily_keep_enabled: bool = False + daily_keep_count: int = 10 items: list[TaskItemIn] = Field(default_factory=list) class TaskPatchIn(BaseModel): note: str | None = None interval_sec: int | None = None - retention_batches: int | None = None + retention_days: int | None = None + daily_keep_enabled: bool | None = None + daily_keep_count: int | None = None status: str | None = None items: list[TaskItemIn] | None = None +class BatchBaselineIn(BaseModel): + marked: bool = True + + +class BatchBulkDeleteIn(BaseModel): + batch_ids: list[str] = Field(default_factory=list) + + class PreviewIn(BaseModel): vendor: str = "" device_type: str = "" @@ -192,6 +205,31 @@ def api_list_batches( return {"items": svc.list_batches(db, task_id, limit=limit)} +@router.post("/tasks/{task_id}/purge") +def api_purge_task(task_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + """Apply retention policy now (age + optional daily keep). Protected batches skipped.""" + return svc.run_purge_for_task(db, task_id) + + +@router.post("/batches/bulk-delete") +def api_bulk_delete_batches( + body: BatchBulkDeleteIn, db: Session = Depends(get_db) +) -> dict[str, Any]: + return svc.delete_batches_bulk(db, list(body.batch_ids or [])) + + +@router.post("/batches/{batch_id}/baseline") +def api_set_batch_baseline( + batch_id: str, body: BatchBaselineIn, db: Session = Depends(get_db) +) -> dict[str, Any]: + return svc.set_batch_baseline(db, batch_id, marked=bool(body.marked)) + + +@router.delete("/batches/{batch_id}") +def api_delete_batch(batch_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + return svc.delete_batch(db, batch_id) + + @router.get("/batches/{batch_id}") def api_get_batch(batch_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: return svc.get_batch(db, batch_id) diff --git a/netx_api/biz_state_scheduler.py b/netx_api/biz_state_scheduler.py index 2dd9243..e0b7fdb 100644 --- a/netx_api/biz_state_scheduler.py +++ b/netx_api/biz_state_scheduler.py @@ -19,12 +19,40 @@ _thread: threading.Thread | None = None _dispatch_pool: ThreadPoolExecutor | None = None _pool_lock = threading.Lock() _last_tick_mono: float = 0.0 +_last_purge_mono: float = 0.0 +_PURGE_INTERVAL_SEC = 3600.0 def _utcnow() -> datetime: return datetime.utcnow() +def _maybe_purge_retention() -> None: + """Hourly sweep so paused/HF tasks still get retention cleanup.""" + import time as _time + + global _last_purge_mono + now = _time.monotonic() + if _last_purge_mono and (now - _last_purge_mono) < _PURGE_INTERVAL_SEC: + return + _last_purge_mono = now + db = SessionLocal() + try: + from .biz_state.retention import purge_all_tasks + + info = purge_all_tasks(db) + if info.get("dropped"): + _log.info( + "biz_state periodic retention dropped=%s tasks=%s", + info.get("dropped"), + len(info.get("tasks") or []), + ) + except Exception: + _log.exception("biz_state periodic retention failed") + finally: + db.close() + + def _dispatch_pool_get() -> ThreadPoolExecutor: global _dispatch_pool with _pool_lock: @@ -98,6 +126,10 @@ def _loop() -> None: try_dispatch_due_tasks() except Exception: _log.exception("biz_state scheduler tick failed") + try: + _maybe_purge_retention() + except Exception: + _log.exception("biz_state retention tick failed") def start_biz_state_scheduler() -> None: diff --git a/netx_api/main.py b/netx_api/main.py index 6999391..98eb06f 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -24,6 +24,7 @@ from .ops_router import router as ops_router from .parser_config import load_parser_config from .port_traffic_router import router as port_traffic_router from .biz_state_router import router as biz_state_router +from .biz_migration_router import router as biz_migration_router from .sql_router import router as sql_router from .sql_router import sql_query, sql_ume_query # noqa: F401 — tests import from main from .topology_router import router as topology_router @@ -76,6 +77,7 @@ app.include_router(collection_router) app.include_router(config_sync_router) app.include_router(port_traffic_router) app.include_router(biz_state_router) +app.include_router(biz_migration_router) app.include_router(webcrt_router) app.include_router(topology_router) app.include_router(lldp_collect_router) diff --git a/netx_api/models/__init__.py b/netx_api/models/__init__.py index d9392cf..7700855 100644 --- a/netx_api/models/__init__.py +++ b/netx_api/models/__init__.py @@ -42,6 +42,13 @@ from .biz_state import ( BizStateTaskItem, BizStateTaskItemBinding, ) +from .biz_migration import ( + BizMigrationBatch, + BizMigrationDiff, + BizMigrationProject, + BizMigrationRedTicket, + BizMigrationRun, +) from .port_traffic import ( PortTrafficBoard, PortTrafficDevice, @@ -146,4 +153,9 @@ __all__ = [ "BizCompareJob", "BizCompareRun", "BizCompareDiff", + "BizMigrationProject", + "BizMigrationBatch", + "BizMigrationRun", + "BizMigrationDiff", + "BizMigrationRedTicket", ] diff --git a/netx_api/models/biz_migration.py b/netx_api/models/biz_migration.py new file mode 100644 index 0000000..5e17d9d --- /dev/null +++ b/netx_api/models/biz_migration.py @@ -0,0 +1,129 @@ +"""Cutover / migration monitor ORM (separate from biz_state CompareJob).""" + +from __future__ import annotations + +from datetime import datetime +from uuid import uuid4 + +from sqlalchemy import Boolean, DateTime, Index, Integer, String, Text +from sqlalchemy.orm import Mapped, mapped_column + +from ..db import Base +from ..timeutil import utcnow_naive +from ._types import JsonType as _JsonType + + +class BizMigrationProject(Base): + """Cutover project linking old/new collect tasks + baselines.""" + + __tablename__ = "biz_migration_project" + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + name: Mapped[str] = mapped_column(String(256), default="", index=True) + old_task_id: Mapped[str] = mapped_column(String(64), default="", index=True) + new_task_id: Mapped[str] = mapped_column(String(64), default="", index=True) + old_baseline_batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + new_baseline_batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + # Reuse biz_port_mapping (before_if=old, after_if=new) + mapping_id: Mapped[str] = mapped_column(String(64), default="", index=True) + status: Mapped[str] = mapped_column(String(32), default="draft", index=True) # draft|active|done + note: Mapped[str] = mapped_column(String(512), default="") + created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive) + + +class BizMigrationBatch(Base): + """One cutover night / wave with an expect set.""" + + __tablename__ = "biz_migration_batch" + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + project_id: Mapped[str] = mapped_column(String(64), default="", index=True) + batch_label: Mapped[str] = mapped_column(String(128), default="") + # pending | active | review | done + status: Mapped[str] = mapped_column(String(32), default="pending", index=True) + # {"ports": ["gei-..."], "items": [{"metric_id":"bgp_peer","key":"..."}]} + expect_set_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + # Acceptance (set by finish_batch): none | passed | failed + accept_status: Mapped[str] = mapped_column(String(32), default="none", index=True) + accept_run_id: Mapped[str] = mapped_column(String(64), default="", index=True) + accept_summary_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + note: Mapped[str] = mapped_column(String(512), default="") + created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive) + updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive) + + +class BizMigrationRun(Base): + """One user-triggered evaluation against baselines + current batches.""" + + __tablename__ = "biz_migration_run" + __table_args__ = (Index("ix_biz_migration_run_batch_created", "batch_id", "created_at"),) + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + project_id: Mapped[str] = mapped_column(String(64), default="", index=True) + batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + old_batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + new_batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + # manual | acceptance + purpose: Mapped[str] = mapped_column(String(32), default="manual", index=True) + status: Mapped[str] = mapped_column(String(32), default="success", index=True) + summary_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + message: Mapped[str] = mapped_column(String(1024), default="") + created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + + +class BizMigrationDiff(Base): + """Per-row / dual verdict row for a migration run (paged board detail).""" + + __tablename__ = "biz_migration_diff" + __table_args__ = ( + Index("ix_biz_migration_diff_run_metric", "run_id", "metric_id"), + Index("ix_biz_migration_diff_run_verdict", "run_id", "verdict"), + ) + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + run_id: Mapped[str] = mapped_column(String(64), default="", index=True) + metric_id: Mapped[str] = mapped_column(String(64), default="", index=True) + seq: Mapped[int] = mapped_column(Integer, default=0) + # ok | expected | anomaly | migrating | migrated | lost | not_involved | unfinished + verdict: Mapped[str] = mapped_column(String(32), default="", index=True) + # green | yellow | red | gray + color: Mapped[str] = mapped_column(String(16), default="") + key_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + old_kind: Mapped[str] = mapped_column(String(16), default="") + new_kind: Mapped[str] = mapped_column(String(16), default="") + old_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + new_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + in_expect: Mapped[bool] = mapped_column(Boolean, default=False) + search_text: Mapped[str] = mapped_column(Text, default="") + + +class BizMigrationRedTicket(Base): + """Persisted anomaly / unfinished expect items — survive across batches (带红继续).""" + + __tablename__ = "biz_migration_red_ticket" + __table_args__ = ( + Index("ix_biz_migration_red_project_status", "project_id", "status"), + Index("ix_biz_migration_red_batch", "batch_id"), + ) + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + project_id: Mapped[str] = mapped_column(String(64), default="", index=True) + batch_id: Mapped[str] = mapped_column(String(64), default="", index=True) + run_id: Mapped[str] = mapped_column(String(64), default="", index=True) + metric_id: Mapped[str] = mapped_column(String(64), default="", index=True) + key_str: Mapped[str] = mapped_column(String(256), default="") + new_key_str: Mapped[str] = mapped_column(String(256), default="") + verdict: Mapped[str] = mapped_column(String(32), default="") + color: Mapped[str] = mapped_column(String(16), default="red") + old_status: Mapped[str] = mapped_column(String(64), default="") + new_status: Mapped[str] = mapped_column(String(64), default="") + detail_json: Mapped[dict] = mapped_column(_JsonType, default=dict) + # open | carried | resolved + status: Mapped[str] = mapped_column(String(32), default="open", index=True) + carried_to_batch_id: Mapped[str] = mapped_column(String(64), default="") + note: Mapped[str] = mapped_column(String(512), default="") + created_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + resolved_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) diff --git a/netx_api/models/biz_state.py b/netx_api/models/biz_state.py index a9de3a1..df0b296 100644 --- a/netx_api/models/biz_state.py +++ b/netx_api/models/biz_state.py @@ -14,10 +14,13 @@ from ._types import JsonType as _JsonType class BizStateTask(Base): - """Per-NE business state monitoring config.""" + """Per-NE business state monitoring config. + + Multiple tasks per ``(source, ne_id)`` are allowed (e.g. hourly full + + cutover high-freq port-only). + """ __tablename__ = "biz_state_task" - __table_args__ = (UniqueConstraint("source", "ne_id", name="uq_biz_state_task_ne"),) id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) source: Mapped[str] = mapped_column(String(32), default="managed", index=True) @@ -29,6 +32,12 @@ class BizStateTask(Base): note: Mapped[str] = mapped_column(String(256), default="") status: Mapped[str] = mapped_column(String(32), default="draft", index=True) # draft|running|paused|stopped interval_sec: Mapped[int] = mapped_column(Integer, default=300) + # Keep snapshots for N calendar days (protected baselines/refs never auto-deleted) + retention_days: Mapped[int] = mapped_column(Integer, default=30) + # Optional: keep only N batches per past calendar day (purge extras next day). Default off. + daily_keep_enabled: Mapped[bool] = mapped_column(Boolean, default=False) + daily_keep_count: Mapped[int] = mapped_column(Integer, default=10) + # Legacy column kept for brownfield reads; unused by new purge path retention_batches: Mapped[int] = mapped_column(Integer, default=30) collect_running: Mapped[bool] = mapped_column(Boolean, default=False) last_collect_started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) @@ -83,6 +92,9 @@ class BizStateBatch(Base): message: Mapped[str] = mapped_column(String(1024), default="") started_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + # Manual pin: excluded from auto-purge until revoked + is_baseline: Mapped[bool] = mapped_column(Boolean, default=False, index=True) + baseline_marked_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) class BizStateBatchCommand(Base): diff --git a/tests/test_biz_migration_evaluate.py b/tests/test_biz_migration_evaluate.py new file mode 100644 index 0000000..2142a33 --- /dev/null +++ b/tests/test_biz_migration_evaluate.py @@ -0,0 +1,136 @@ +"""Unit tests for cutover migration evaluate helpers (port status focus).""" + +from __future__ import annotations + +import unittest + +from netx_api.biz_migration.evaluate import ( + dual_verdict, + evaluate_metric_dual, + parse_expect_set, + port_sheet_def, + port_status_label, + side_verdict, +) + + +class ParseExpectSetTests(unittest.TestCase): + def test_ports_go_to_interface_brief(self): + got = parse_expect_set({"ports": ["gei-1", "gei-2", ""]}) + self.assertEqual(got["interface_brief"], {"gei-1", "gei-2"}) + self.assertEqual(got["_ports"], {"gei-1", "gei-2"}) + + +class PortSheetTests(unittest.TestCase): + def test_port_sheet_fields(self): + s = port_sheet_def() + self.assertEqual(s["metric_id"], "interface_brief") + self.assertEqual(s["key_fields"], ["interface"]) + self.assertEqual(s["compare_fields"], ["admin", "phy", "prot"]) + + def test_status_label(self): + self.assertEqual( + port_status_label({"admin": "up", "phy": "up", "prot": "up"}), + "up/up/up", + ) + self.assertEqual(port_status_label({}), "—") + + +class VerdictTests(unittest.TestCase): + def test_side_anomaly_gone(self): + self.assertEqual( + side_verdict(kind="removed", in_expect=False, window_active=True), + ("anomaly_gone", "red"), + ) + self.assertEqual( + side_verdict(kind="removed", in_expect=True, window_active=True), + ("expected_gone", "yellow"), + ) + + def test_dual_migrated(self): + self.assertEqual( + dual_verdict(old_kind="removed", new_kind="added", in_expect=True, window_active=True), + ("migrated", "green"), + ) + + def test_dual_lost(self): + self.assertEqual( + dual_verdict(old_kind="removed", new_kind="", in_expect=True, window_active=False), + ("lost", "red"), + ) + + def test_acceptance_unfinished_is_red(self): + self.assertEqual( + dual_verdict( + old_kind="unchanged", + new_kind="", + in_expect=True, + window_active=False, + acceptance=True, + ), + ("unfinished", "red"), + ) + # window mode still yellow + self.assertEqual( + dual_verdict( + old_kind="unchanged", + new_kind="", + in_expect=True, + window_active=True, + acceptance=False, + ), + ("migrating", "yellow"), + ) + + +class EvaluateMetricDualTests(unittest.TestCase): + def test_port_migration_happy_path(self): + old_base = [{"interface": "gei-old", "admin": "up", "phy": "up", "prot": "up"}] + old_cur: list[dict] = [] + new_cur = [{"interface": "gei-new", "admin": "up", "phy": "up", "prot": "up"}] + expect = parse_expect_set({"ports": ["gei-old"]}) + out = evaluate_metric_dual( + metric_id="interface_brief", + key_fields=["interface"], + iface_fields=["interface"], + compare_fields=["admin", "phy", "prot"], + old_baseline_rows=old_base, + old_current_rows=old_cur, + new_baseline_rows=None, + new_current_rows=new_cur, + port_map={"gei-old": "gei-new"}, + expect=expect, + window_active=True, + ) + self.assertEqual(out["progress_total"], 1) + migrated = [r for r in out["rows"] if r["verdict"] == "migrated"] + self.assertTrue(migrated, out["rows"]) + self.assertEqual(migrated[0]["color"], "green") + self.assertEqual(migrated[0]["old_status"], "gone") + self.assertEqual(migrated[0]["new_status"], "up/up/up") + self.assertEqual(out["progress_ok"], 1) + + def test_unexpected_loss_is_anomaly(self): + old_base = [{"interface": "gei-keep", "admin": "up", "phy": "up", "prot": "up"}] + old_cur: list[dict] = [] + new_cur = [{"interface": "gei-keep", "admin": "up", "phy": "up", "prot": "up"}] + expect = parse_expect_set({"ports": []}) + out = evaluate_metric_dual( + metric_id="interface_brief", + key_fields=["interface"], + iface_fields=["interface"], + compare_fields=["admin", "phy", "prot"], + old_baseline_rows=old_base, + old_current_rows=old_cur, + new_baseline_rows=None, + new_current_rows=new_cur, + port_map={}, + expect=expect, + window_active=False, + ) + reds = [r for r in out["rows"] if r["color"] == "red"] + self.assertTrue(reds, out["rows"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_biz_state_retention.py b/tests/test_biz_state_retention.py new file mode 100644 index 0000000..7e4837b --- /dev/null +++ b/tests/test_biz_state_retention.py @@ -0,0 +1,97 @@ +"""Retention purge: age days + daily keep + baseline protection.""" + +from __future__ import annotations + +import unittest +from datetime import datetime, timedelta + +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker + +from netx_api.db import Base +from netx_api.models import BizStateBatch, BizStateTask +from netx_api.biz_state.retention import purge_task_batches, protected_batch_ids + + +class RetentionPurgeTests(unittest.TestCase): + def setUp(self) -> None: + engine = create_engine("sqlite+pysqlite:///:memory:", future=True) + TestingSession = sessionmaker( + bind=engine, autoflush=False, autocommit=False, expire_on_commit=False + ) + Base.metadata.create_all(bind=engine) + self.db = TestingSession() + self.task = BizStateTask( + id="t1", + source="managed", + ne_id="ne1", + status="running", + interval_sec=60, + retention_days=7, + daily_keep_enabled=False, + daily_keep_count=10, + ) + self.db.add(self.task) + self.db.commit() + + def tearDown(self) -> None: + self.db.close() + + def _add_batch(self, bid: str, *, days_ago: float, baseline: bool = False) -> BizStateBatch: + started = datetime.utcnow() - timedelta(days=days_ago) + b = BizStateBatch( + id=bid, + task_id=self.task.id, + status="success", + started_at=started, + ended_at=started, + is_baseline=baseline, + baseline_marked_at=started if baseline else None, + ) + self.db.add(b) + self.db.commit() + return b + + def test_age_purge_keeps_baseline(self) -> None: + self._add_batch("old", days_ago=10) + self._add_batch("pin", days_ago=20, baseline=True) + self._add_batch("fresh", days_ago=1) + info = purge_task_batches(self.db, self.task) + self.assertEqual(info["dropped"], 1) + ids = {b.id for b in self.db.query(BizStateBatch).all()} + self.assertEqual(ids, {"pin", "fresh"}) + self.assertIn("pin", protected_batch_ids(self.db, task_id=self.task.id)) + + def test_daily_keep_skips_today(self) -> None: + self.task.daily_keep_enabled = True + self.task.daily_keep_count = 1 + self.db.commit() + # same past calendar day: 3 batches → keep 1 newest + yesterday_noon = datetime.utcnow().replace(hour=12, minute=0, second=0, microsecond=0) - timedelta( + days=1 + ) + for i, bid in enumerate(("y1", "y2", "y3")): + b = BizStateBatch( + id=bid, + task_id=self.task.id, + status="success", + started_at=yesterday_noon - timedelta(hours=i), + ended_at=yesterday_noon - timedelta(hours=i), + ) + self.db.add(b) + # today: keep all even if > N + self._add_batch("tod1", days_ago=0.01) + self._add_batch("tod2", days_ago=0.02) + self.db.commit() + info = purge_task_batches(self.db, self.task) + self.assertEqual(info["dropped"], 2) + ids = {b.id for b in self.db.query(BizStateBatch).all()} + self.assertIn("y1", ids) # newest yesterday + self.assertIn("tod1", ids) + self.assertIn("tod2", ids) + self.assertNotIn("y2", ids) + self.assertNotIn("y3", ids) + + +if __name__ == "__main__": + unittest.main() diff --git a/web/src/App.tsx b/web/src/App.tsx index 943f404..59e189a 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -49,6 +49,9 @@ const BizStatePage = lazy(() => const BizComparePage = lazy(() => import("./pages/network/BizComparePage").then((m) => ({ default: m.BizComparePage })), ); +const BizMigrationPage = lazy(() => + import("./pages/network/BizMigrationPage").then((m) => ({ default: m.BizMigrationPage })), +); const UsersPage = lazy(() => import("./pages/UsersPage").then((m) => ({ default: m.UsersPage }))); const AuditLayout = lazy(() => import("./pages/audit/AuditLayout").then((m) => ({ default: m.AuditLayout })), @@ -128,6 +131,7 @@ function ProtectedApp() { } /> } /> } /> + } /> } /> } /> diff --git a/web/src/config/networkNav.ts b/web/src/config/networkNav.ts index 0892c88..18b465d 100644 --- a/web/src/config/networkNav.ts +++ b/web/src/config/networkNav.ts @@ -86,6 +86,12 @@ export const NETWORK_NAV: readonly NetworkNavGroup[] = [ labelKey: "network.nav.bizCompare", group: "tasks", }, + { + id: "biz-migration", + path: "/network/tasks/biz-migration", + labelKey: "network.nav.bizMigration", + group: "tasks", + }, ], }, ] as const; diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 736163e..ee9bd27 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -158,6 +158,7 @@ const en = { portTrafficWall: "Traffic wall", bizState: "Business state", bizCompare: "Business compare", + bizMigration: "Cutover monitor", }, collapseNav: "Collapse sidebar", expandNav: "Expand sidebar", @@ -199,7 +200,25 @@ const en = { intervalUnitDays: "days", intervalUnitHours: "hours", intervalUnitSeconds: "seconds", - retention: "Keep batches", + retention: "Retention days", + retentionHint: "Batches older than this are purged; baselines and compare/cutover refs are kept.", + dailyKeepEnabled: "Daily keep policy", + dailyKeepCount: "Keep per day", + dailyKeepHint: "When on: keep only the newest N batches per past day; extras purge next day (off by default).", + markBaseline: "Mark baseline", + unmarkBaseline: "Unmark baseline", + baseline: "Baseline", + protected: "Protected", + deleteBatch: "Delete", + bulkDelete: "Bulk delete", + confirmDeleteBatch: "Delete this batch? Protected batches cannot be deleted.", + confirmBulkDelete: "Delete {{count}} selected batches? Protected ones are skipped.", + batchDeleted: "Batch deleted", + bulkDeleted: "Deleted {{deleted}}, skipped {{skipped}}", + purgeNow: "Purge now", + purgeOk: "Purge done: dropped {{dropped}}", + selectAll: "Select deletable", + colProtect: "Protect", saveSchedule: "Save schedule", scheduleSaved: "Schedule saved", colInterval: "Interval", @@ -428,6 +447,89 @@ const en = { colKey: "Identity", colChange: "Changes", }, + bizMigration: { + title: "Cutover monitor", + hint: "Separate from business compare: baseline drift, expect set, dual-device verdict. Reuses collect tasks and port mappings.", + hintPort: + "Metric focus: port status (interface_brief / admin·phy·prot). Baseline → pick ports → evaluate board.", + createProject: "New cutover project", + projectName: "Project name", + oldTask: "Old-device collect task", + newTask: "New-device collect task", + portMapping: "Port mapping", + oldBaseline: "Old baseline batch", + newBaseline: "New baseline batch (optional)", + needProjectFields: "Enter a name and pick old/new collect tasks", + needExpectPorts: "Select at least one expected port for this batch", + needBaselineFirst: "Save an old baseline batch first to list ports", + projectCreated: "Project created", + projects: "Projects", + emptyProjects: "No cutover projects yet", + saveBaseline: "Save baseline / mapping", + baselineSaved: "Baseline saved", + pickExpectPorts: "Expected ports (from old baseline)", + portFilterPh: "Filter interface / description", + emptyBaselinePorts: "No interface_brief rows in baseline (enable that profile on the task)", + batchLabel: "Batch label", + defaultBatchLabel: "Night N", + expectPorts: "Expected ports (old side)", + expectPortsPh: "Comma or newline separated", + createBatch: "New batch", + batchCreated: "Batch created", + batches: "Batches", + startBatch: "Start batch", + finishBatch: "Finish batch", + evaluate: "Evaluate now", + evaluated: "Evaluation done", + statusUpdated: "Status updated", + board: "Progress board", + metricPort: "Port status", + progress: "Progress", + anomaly: "Anomaly", + onlyExpect: "Expected ports only", + colMetric: "Metric", + colKey: "Object", + colPort: "Old port", + colMapped: "Mapped new", + colDesc: "Description", + colOldPort: "Old port", + colNewPort: "New port", + colOldStatus: "Old admin/phy/prot", + colNewStatus: "New admin/phy/prot", + colOld: "Old", + colNew: "New", + colVerdict: "Verdict", + emptyDiffs: "No rows yet — run Evaluate", + verdictMigrated: "Migrated", + verdictMigrating: "Migrating", + verdictLost: "Lost", + verdictAnomaly: "Anomaly", + verdictNotInvolved: "Not involved", + verdictUnexpected: "Unexpected new", + verdictOk: "OK", + enableHighfreq: "Enable high-freq port collect", + highfreqHint: + "Creates/reuses biz-state tasks with interface_brief only on old/new NEs. Add BGP etc. later on the same tasks in Business state.", + highfreqReady: "High-freq port collect tasks ready and bound to this project", + collectNow: "Collect old+new now", + collectTriggered: "Collect triggered", + openBizState: "Open business state", + boundTasks: "Bound collect tasks", + acceptTitle: "Batch acceptance summary", + acceptPassed: "Batch acceptance passed", + acceptFailed: "Batch acceptance failed ({{n}} red tickets; you may continue)", + acceptPassedShort: "Passed", + acceptFailedShort: "Failed", + acceptCarryHint: "Failure does not block the next batch; red tickets stay until resolved.", + openRedHint: "{{n}} open red ticket(s) on this project (carry allowed)", + redTitle: "Red tickets", + redOpen: "open", + resolveRed: "Resolve", + redResolved: "Red ticket resolved", + batchCreatedWithRed: "Batch created ({{n}} open red tickets; carry allowed)", + colStatus: "Status", + verdictUnfinished: "Unfinished", + }, portTraffic: { title: "Port traffic", create: "Monitor device", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 80f1aa4..3c5e3be 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -158,6 +158,7 @@ const zh = { portTrafficWall: "流量大屏", bizState: "业务状态监控", bizCompare: "业务状态比对", + bizMigration: "割接监控", }, collapseNav: "折叠侧栏", expandNav: "展开侧栏", @@ -199,7 +200,25 @@ const zh = { intervalUnitDays: "天", intervalUnitHours: "小时", intervalUnitSeconds: "秒", - retention: "保留批次数", + retention: "保留天数", + retentionHint: "超过天数的批次会自动清理;基线与比对/割接引用的批次不会删除。", + dailyKeepEnabled: "每日保留策略", + dailyKeepCount: "每天保留次数", + dailyKeepHint: "开启后:过去每一天只保留最近 N 次,多余在次日清理(默认关闭)。", + markBaseline: "标为基线", + unmarkBaseline: "撤销基线", + baseline: "基线", + protected: "已保护", + deleteBatch: "删除", + bulkDelete: "批量删除", + confirmDeleteBatch: "删除该采集批次?受保护的批次无法删除。", + confirmBulkDelete: "删除选中的 {{count}} 个批次?受保护的会自动跳过。", + batchDeleted: "已删除批次", + bulkDeleted: "已删除 {{deleted}} 个,跳过 {{skipped}} 个", + purgeNow: "立即清理", + purgeOk: "清理完成:删除 {{dropped}} 个批次", + selectAll: "全选可删", + colProtect: "保护", saveSchedule: "保存周期", scheduleSaved: "周期配置已保存", colInterval: "周期", @@ -426,6 +445,88 @@ const zh = { colKey: "身份", colChange: "变更", }, + bizMigration: { + title: "割接监控", + hint: "与业务状态比对独立:相对基线漂移 + 本批预期 + 老/新双端判定。复用采集任务与端口映射。", + hintPort: "当前监控项:端口状态(interface_brief / admin·phy·prot)。选基线 → 勾选本批端口 → 判定看板。", + createProject: "新建割接项目", + projectName: "项目名称", + oldTask: "老设备采集任务", + newTask: "新设备采集任务", + portMapping: "端口映射", + oldBaseline: "老设备基线批次", + newBaseline: "新设备基线批次(可选)", + needProjectFields: "请填写项目名并选择老/新采集任务", + needExpectPorts: "请先勾选本批预期迁移的端口", + needBaselineFirst: "请先保存老设备基线批次,才能列出端口", + projectCreated: "项目已创建", + projects: "项目", + emptyProjects: "暂无割接项目", + saveBaseline: "保存基线/映射", + baselineSaved: "基线已保存", + pickExpectPorts: "本批预期端口(来自老基线)", + portFilterPh: "筛选接口名 / 描述", + emptyBaselinePorts: "基线中无端口状态数据(确认任务已采 interface_brief)", + batchLabel: "批次标签", + defaultBatchLabel: "第N晚", + expectPorts: "本批预期端口(老侧)", + expectPortsPh: "逗号或换行分隔", + createBatch: "新建批次", + batchCreated: "批次已创建", + batches: "批次", + startBatch: "本批开始", + finishBatch: "本批完成", + evaluate: "立即判定", + evaluated: "判定完成", + statusUpdated: "状态已更新", + board: "进度看板", + metricPort: "端口状态", + progress: "进度", + anomaly: "异常", + onlyExpect: "仅看预期端口", + colMetric: "监控项", + colKey: "对象", + colPort: "老侧端口", + colMapped: "映射新口", + colDesc: "描述", + colOldPort: "老端口", + colNewPort: "新端口", + colOldStatus: "老状态 admin/phy/prot", + colNewStatus: "新状态 admin/phy/prot", + colOld: "老侧", + colNew: "新侧", + colVerdict: "判定", + emptyDiffs: "暂无明细,请先「立即判定」", + verdictMigrated: "已迁移", + verdictMigrating: "迁移中", + verdictLost: "丢失", + verdictAnomaly: "异常", + verdictNotInvolved: "未涉及", + verdictUnexpected: "意外新增", + verdictOk: "正常", + enableHighfreq: "开启高频端口采集", + highfreqHint: + "在业务监控下为老/新设备各建(或复用)仅含 interface_brief 的任务并周期采集;后续要加 BGP 等直接去业务监控勾选即可。", + highfreqReady: "高频端口采集任务已就绪(已绑定到本项目)", + collectNow: "立即采集老+新", + collectTriggered: "已触发采集", + openBizState: "打开业务监控", + boundTasks: "当前绑定采集任务", + acceptTitle: "本批验收小结", + acceptPassed: "本批验收通过", + acceptFailed: "本批验收未通过(红单 {{n}},可带红开下一批)", + acceptPassedShort: "通过", + acceptFailedShort: "未通过", + acceptCarryHint: "验收不过不拦截:可新建下一批继续;红单保留留痕,可手动关闭。", + openRedHint: "项目仍有 {{n}} 条未关闭红单(带红继续)", + redTitle: "红单留痕", + redOpen: "未关闭", + resolveRed: "关闭", + redResolved: "红单已关闭", + batchCreatedWithRed: "批次已创建(当前有 {{n}} 条未关闭红单,允许带红继续)", + colStatus: "状态", + verdictUnfinished: "未完成", + }, portTraffic: { title: "端口流量监控", create: "添加设备监控", diff --git a/web/src/pages/network/BizMigrationPage.tsx b/web/src/pages/network/BizMigrationPage.tsx new file mode 100644 index 0000000..2b0a333 --- /dev/null +++ b/web/src/pages/network/BizMigrationPage.tsx @@ -0,0 +1,938 @@ +import { Button, Input } from "@heroui/react"; +import { useCallback, useEffect, useMemo, useState } from "react"; +import { FieldSelect } from "../../components/ui/FieldSelect"; +import { useToast } from "../../hooks/useToast"; +import { useI18n } from "../../i18n"; +import { + bizCompareListMappings, + bizMigrationCollectNow, + bizMigrationCreateBatch, + bizMigrationCreateProject, + bizMigrationEnsurePortHighfreq, + bizMigrationEvaluate, + bizMigrationFinishBatch, + bizMigrationGetBoard, + bizMigrationListBaselinePorts, + bizMigrationListBatches, + bizMigrationListDiffs, + bizMigrationListProjects, + bizMigrationListRedTickets, + bizMigrationPatchBatch, + bizMigrationPatchProject, + bizMigrationResolveRedTicket, + bizStateListBatches, + bizStateListTasks, + formatErr, +} from "../../services/api"; +import { Link } from "react-router-dom"; + +type TaskOpt = { id: string; ne_name: string; ne_ip: string; note?: string; interval_sec?: number }; +type Project = { + id: string; + name: string; + old_task_id: string; + new_task_id: string; + old_baseline_batch_id: string; + new_baseline_batch_id: string; + mapping_id: string; + status: string; + old_task?: TaskOpt; + new_task?: TaskOpt; +}; +type MigBatch = { + id: string; + batch_label: string; + status: string; + expect_set?: { ports?: string[] }; + accept_status?: string; + accept_run_id?: string; + accept_summary?: { + passed?: boolean; + progress_ok?: number; + progress_total?: number; + anomaly?: number; + verdict_counts?: Record; + }; +}; +type RedTicket = { + id: string; + batch_id: string; + key_str: string; + new_key_str?: string; + verdict: string; + old_status?: string; + new_status?: string; + status: string; +}; +type SheetCard = { + metric_id: string; + title?: string; + progress_ok: number; + progress_total: number; + anomaly: number; +}; +type DiffRow = { + id: string; + metric_id: string; + verdict: string; + color: string; + old_kind: string; + new_kind: string; + in_expect: boolean; + key_str?: string; + new_key_str?: string; + old_status?: string; + new_status?: string; +}; +type BaselinePort = { + interface: string; + admin?: string; + phy?: string; + prot?: string; + description?: string; + mapped_to?: string; +}; + +const COLOR: Record = { + green: "#16a34a", + yellow: "#ca8a04", + red: "#dc2626", + gray: "#6b7280", +}; + + const VERDICT_I18N: Record = { + migrated: "bizMigration.verdictMigrated", + migrating: "bizMigration.verdictMigrating", + lost: "bizMigration.verdictLost", + anomaly: "bizMigration.verdictAnomaly", + not_involved: "bizMigration.verdictNotInvolved", + unexpected_new: "bizMigration.verdictUnexpected", + unfinished: "bizMigration.verdictUnfinished", + ok: "bizMigration.verdictOk", +}; + +export function BizMigrationPage() { + const { t } = useI18n(); + const { showOk, showError } = useToast(); + + const [projects, setProjects] = useState([]); + const [tasks, setTasks] = useState([]); + const [mappings, setMappings] = useState<{ id: string; name: string }[]>([]); + const [projectId, setProjectId] = useState(""); + const [batches, setBatches] = useState([]); + const [batchId, setBatchId] = useState(""); + const [board, setBoard] = useState<{ + run?: { + id: string; + summary?: { + sheet_cards?: SheetCard[]; + progress?: { ok: number; total: number }; + anomaly?: number; + expect_ports?: string[]; + }; + } | null; + } | null>(null); + const [diffs, setDiffs] = useState([]); + const [busy, setBusy] = useState(false); + const [onlyExpect, setOnlyExpect] = useState(true); + + const [newName, setNewName] = useState(""); + const [oldTaskId, setOldTaskId] = useState(""); + const [newTaskId, setNewTaskId] = useState(""); + const [mappingId, setMappingId] = useState(""); + const [oldBaselineId, setOldBaselineId] = useState(""); + const [newBaselineId, setNewBaselineId] = useState(""); + const [oldBatches, setOldBatches] = useState<{ id: string; started_at?: string | null }[]>([]); + const [newBatches, setNewBatches] = useState<{ id: string; started_at?: string | null }[]>([]); + const [batchLabel, setBatchLabel] = useState(""); + const [baselinePorts, setBaselinePorts] = useState([]); + const [selectedPorts, setSelectedPorts] = useState>(new Set()); + const [portFilter, setPortFilter] = useState(""); + const [redTickets, setRedTickets] = useState([]); + const [openRedCount, setOpenRedCount] = useState(0); + const [acceptInfo, setAcceptInfo] = useState(null); + + const project = useMemo( + () => projects.find((p) => p.id === projectId) || null, + [projects, projectId], + ); + const batch = useMemo(() => batches.find((b) => b.id === batchId) || null, [batches, batchId]); + const sheetCards = board?.run?.summary?.sheet_cards || []; + const visibleDiffs = useMemo( + () => (onlyExpect ? diffs.filter((d) => d.in_expect) : diffs), + [diffs, onlyExpect], + ); + + useEffect(() => { + if (batch?.accept_summary && Object.keys(batch.accept_summary).length) { + setAcceptInfo(batch.accept_summary); + } else if (batch && batch.accept_status === "none") { + setAcceptInfo(null); + } + }, [batch]); + + const filteredBaselinePorts = useMemo(() => { + const kw = portFilter.trim().toLowerCase(); + if (!kw) return baselinePorts; + return baselinePorts.filter( + (p) => + p.interface.toLowerCase().includes(kw) || + String(p.description || "") + .toLowerCase() + .includes(kw), + ); + }, [baselinePorts, portFilter]); + + const reloadProjects = useCallback(async () => { + const res = await bizMigrationListProjects(); + setProjects((res.items || []) as Project[]); + }, []); + + const loadBaselinePorts = useCallback(async (pid: string) => { + if (!pid) { + setBaselinePorts([]); + return; + } + try { + const res = await bizMigrationListBaselinePorts(pid); + setBaselinePorts(res.ports || []); + } catch { + setBaselinePorts([]); + } + }, []); + + const loadRedTickets = useCallback(async (pid: string) => { + if (!pid) { + setRedTickets([]); + setOpenRedCount(0); + return; + } + try { + const res = await bizMigrationListRedTickets(pid); + setRedTickets((res.items || []) as RedTicket[]); + setOpenRedCount(Number(res.open_count || 0)); + } catch { + setRedTickets([]); + setOpenRedCount(0); + } + }, []); + + useEffect(() => { + void (async () => { + try { + const [pr, tk, mp] = await Promise.all([ + bizMigrationListProjects(), + bizStateListTasks(), + bizCompareListMappings(), + ]); + setProjects((pr.items || []) as Project[]); + setTasks( + ((tk.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + ne_name: String(x.ne_name || ""), + ne_ip: String(x.ne_ip || ""), + note: String(x.note || ""), + interval_sec: Number(x.interval_sec || 0) || undefined, + })), + ); + setMappings( + ((mp.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + name: String(x.name || x.id || ""), + })), + ); + } catch (e) { + showError(formatErr(e)); + } + })(); + }, [showError]); + + useEffect(() => { + if (!projectId) { + setBatches([]); + setBatchId(""); + setBaselinePorts([]); + return; + } + void (async () => { + try { + const res = await bizMigrationListBatches(projectId); + const items = (res.items || []) as MigBatch[]; + setBatches(items); + if (!items.find((b) => b.id === batchId)) { + setBatchId(items[0]?.id || ""); + } + await loadBaselinePorts(projectId); + await loadRedTickets(projectId); + } catch (e) { + showError(formatErr(e)); + } + })(); + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [projectId, showError, loadBaselinePorts, loadRedTickets]); + + useEffect(() => { + if (!batchId) { + setBoard(null); + setDiffs([]); + return; + } + void loadBoard(batchId); + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [batchId]); + + useEffect(() => { + if (!oldTaskId) { + setOldBatches([]); + return; + } + void bizStateListBatches(oldTaskId, 30).then((r) => { + setOldBatches( + ((r.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + started_at: (x.started_at as string) || null, + })), + ); + }); + }, [oldTaskId]); + + useEffect(() => { + if (!newTaskId) { + setNewBatches([]); + return; + } + void bizStateListBatches(newTaskId, 30).then((r) => { + setNewBatches( + ((r.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + started_at: (x.started_at as string) || null, + })), + ); + }); + }, [newTaskId]); + + async function loadBoard(bid: string, runId = "") { + try { + const b = await bizMigrationGetBoard(bid, runId); + setBoard(b as typeof board); + const rid = String((b as { run?: { id?: string } })?.run?.id || ""); + if (rid) { + const d = await bizMigrationListDiffs({ runId: rid, limit: 500 }); + setDiffs((d.items || []) as DiffRow[]); + } else { + setDiffs([]); + } + } catch (e) { + showError(formatErr(e)); + } + } + + async function onCreateProject() { + if (!newName.trim() || !oldTaskId || !newTaskId) { + showError(t("bizMigration.needProjectFields")); + return; + } + setBusy(true); + try { + const p = (await bizMigrationCreateProject({ + name: newName.trim(), + old_task_id: oldTaskId, + new_task_id: newTaskId, + mapping_id: mappingId || "", + old_baseline_batch_id: oldBaselineId || "", + new_baseline_batch_id: newBaselineId || "", + status: "active", + })) as Project; + await reloadProjects(); + setProjectId(p.id); + setNewName(""); + showOk(t("bizMigration.projectCreated")); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function onSaveBaseline() { + if (!projectId) return; + setBusy(true); + try { + await bizMigrationPatchProject(projectId, { + old_baseline_batch_id: oldBaselineId || project?.old_baseline_batch_id || "", + new_baseline_batch_id: newBaselineId || project?.new_baseline_batch_id || "", + mapping_id: mappingId || project?.mapping_id || "", + }); + await reloadProjects(); + await loadBaselinePorts(projectId); + showOk(t("bizMigration.baselineSaved")); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function onCreateBatch() { + if (!projectId) return; + const ports = [...selectedPorts]; + if (!ports.length) { + showError(t("bizMigration.needExpectPorts")); + return; + } + setBusy(true); + try { + const b = (await bizMigrationCreateBatch(projectId, { + batch_label: batchLabel.trim() || t("bizMigration.defaultBatchLabel"), + expect_set: { ports }, + status: "pending", + })) as MigBatch; + const res = await bizMigrationListBatches(projectId); + setBatches((res.items || []) as MigBatch[]); + setBatchId(b.id); + setBatchLabel(""); + setSelectedPorts(new Set()); + if (Number((b as { open_red_count?: number }).open_red_count || 0) > 0) { + showOk(t("bizMigration.batchCreatedWithRed", { n: String((b as { open_red_count?: number }).open_red_count) })); + } else { + showOk(t("bizMigration.batchCreated")); + } + await loadRedTickets(projectId); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function setBatchStatus(status: string) { + if (!batchId) return; + setBusy(true); + try { + if (status === "review" || status === "done") { + const res = (await bizMigrationFinishBatch(batchId, status === "done")) as { + batch?: MigBatch; + accept_summary?: MigBatch["accept_summary"]; + open_red_count?: number; + }; + const list = await bizMigrationListBatches(projectId); + setBatches((list.items || []) as MigBatch[]); + setAcceptInfo(res.accept_summary || res.batch?.accept_summary || null); + await loadRedTickets(projectId); + if (res.accept_summary?.passed) { + showOk(t("bizMigration.acceptPassed")); + } else { + showOk(t("bizMigration.acceptFailed", { n: String(res.open_red_count ?? openRedCount) })); + } + if (res.batch?.accept_run_id) { + await loadBoard(batchId, res.batch.accept_run_id); + } + } else { + await bizMigrationPatchBatch(batchId, { status }); + const res = await bizMigrationListBatches(projectId); + setBatches((res.items || []) as MigBatch[]); + showOk(t("bizMigration.statusUpdated")); + if (status === "active") await loadRedTickets(projectId); + } + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function onResolveTicket(id: string) { + try { + await bizMigrationResolveRedTicket(id); + await loadRedTickets(projectId); + showOk(t("bizMigration.redResolved")); + } catch (e) { + showError(formatErr(e)); + } + } + + async function onEvaluate() { + if (!batchId) return; + setBusy(true); + try { + const run = await bizMigrationEvaluate(batchId, {}); + await loadBoard(batchId, String((run as { id?: string }).id || "")); + showOk(t("bizMigration.evaluated")); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function onEnsureHighfreq() { + if (!projectId) return; + setBusy(true); + try { + const res = await bizMigrationEnsurePortHighfreq(projectId, { + interval_sec: 60, + collect_now: true, + }); + await reloadProjects(); + const pr = await bizMigrationListProjects(); + setProjects((pr.items || []) as Project[]); + showOk(t("bizMigration.highfreqReady")); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + async function onCollectNow() { + if (!projectId) return; + setBusy(true); + try { + await bizMigrationCollectNow(projectId); + showOk(t("bizMigration.collectTriggered")); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + } + + useEffect(() => { + if (project) { + setOldTaskId(project.old_task_id); + setNewTaskId(project.new_task_id); + setMappingId(project.mapping_id || ""); + setOldBaselineId(project.old_baseline_batch_id || ""); + setNewBaselineId(project.new_baseline_batch_id || ""); + } + }, [project]); + + function togglePort(name: string) { + setSelectedPorts((prev) => { + const next = new Set(prev); + if (next.has(name)) next.delete(name); + else next.add(name); + return next; + }); + } + + function verdictLabel(v: string) { + const key = VERDICT_I18N[v]; + return key ? t(key) : v; + } + + return ( +
+
+

{t("bizMigration.title")}

+

{t("bizMigration.hintPort")}

+
+ +
+ {t("bizMigration.createProject")} +
+ + setOldTaskId(e.target.value)} + fullWidth + > + + {tasks.map((x) => ( + + ))} + + setNewTaskId(e.target.value)} + fullWidth + > + + {tasks.map((x) => ( + + ))} + + setMappingId(e.target.value)} + fullWidth + > + + {mappings.map((m) => ( + + ))} + + setOldBaselineId(e.target.value)} + fullWidth + > + + {oldBatches.map((x) => ( + + ))} + + setNewBaselineId(e.target.value)} + fullWidth + > + + {newBatches.map((x) => ( + + ))} + +
+ +
+ +
+
+ {t("bizMigration.projects")} + {projects.length === 0 && ( + {t("bizMigration.emptyProjects")} + )} + {projects.map((p) => ( + + ))} +
+ +
+ {project && ( +
+
+ {project.name} +
+ + + + + {t("bizMigration.openBizState")} + +
+
+ {(project.old_task || project.new_task) && ( +
+ {t("bizMigration.boundTasks")}:{" "} + {project.old_task?.ne_name || project.old_task_id.slice(0, 8)} + {project.old_task?.note ? ` (${project.old_task.note})` : ""} / + {project.new_task?.ne_name || project.new_task_id.slice(0, 8)} + {project.new_task?.note ? ` (${project.new_task.note})` : ""} + {project.old_task?.interval_sec + ? ` · ${project.old_task.interval_sec}s` + : ""} +
+ )} + {openRedCount > 0 && ( +
+ {t("bizMigration.openRedHint", { n: String(openRedCount) })} +
+ )} +

{t("bizMigration.highfreqHint")}

+ +
+
+ {t("bizMigration.pickExpectPorts")} ({selectedPorts.size}) +
+ + {!project.old_baseline_batch_id && ( +
+ {t("bizMigration.needBaselineFirst")} +
+ )} +
+ + + + + + + + + + + {filteredBaselinePorts.map((p) => ( + togglePort(p.interface)} style={{ cursor: "pointer" }}> + + + + + + + ))} + {filteredBaselinePorts.length === 0 && ( + + + + )} + +
+ {t("bizMigration.colPort")}admin/phy/prot{t("bizMigration.colMapped")}{t("bizMigration.colDesc")}
+ togglePort(p.interface)} + /> + {p.interface} + {p.admin || "-"}/{p.phy || "-"}/{p.prot || "-"} + {p.mapped_to || "—"} + {p.description || ""} +
+ {t("bizMigration.emptyBaselinePorts")} +
+
+
+ +
+ + + setBatchId(e.target.value)} + > + {batches.map((b) => ( + + ))} + + {batch && ( + <> + + + + + )} +
+
+ )} + + {acceptInfo && ( +
+ {t("bizMigration.acceptTitle")} +
+
+ {acceptInfo.passed ? t("bizMigration.acceptPassedShort") : t("bizMigration.acceptFailedShort")} +
+
+ {t("bizMigration.progress")}: {acceptInfo.progress_ok ?? 0}/{acceptInfo.progress_total ?? 0} +
+
+ {t("bizMigration.anomaly")}: {acceptInfo.anomaly ?? 0} +
+
+
{t("bizMigration.acceptCarryHint")}
+
+ )} + + {redTickets.length > 0 && ( +
+ + {t("bizMigration.redTitle")} ({openRedCount} {t("bizMigration.redOpen")}) + +
+ + + + + + + + + + + + + {redTickets.map((r) => ( + + + + + + + + + + ))} + +
{t("bizMigration.colOldPort")}{t("bizMigration.colNewPort")}{t("bizMigration.colOldStatus")}{t("bizMigration.colNewStatus")}{t("bizMigration.colVerdict")}{t("bizMigration.colStatus")} +
{r.key_str}{r.new_key_str || "—"}{r.old_status || "—"}{r.new_status || "—"}{verdictLabel(r.verdict)}{r.status} + {r.status !== "resolved" && ( + + )} +
+
+
+ )} + + {board?.run && ( +
+ + {t("bizMigration.board")} · {t("bizMigration.metricPort")} + +
+
+ {t("bizMigration.progress")}:{" "} + + {board.run.summary?.progress?.ok ?? 0}/{board.run.summary?.progress?.total ?? 0} + +
+
+ {t("bizMigration.anomaly")}: {board.run.summary?.anomaly ?? 0} +
+ {sheetCards[0] && ( +
+ {sheetCards[0].title || sheetCards[0].metric_id} +
+ )} + +
+ +
+ + + + + + + + + + + + {visibleDiffs.map((d) => ( + + + + + + + + ))} + {visibleDiffs.length === 0 && ( + + + + )} + +
{t("bizMigration.colOldPort")}{t("bizMigration.colNewPort")}{t("bizMigration.colOldStatus")}{t("bizMigration.colNewStatus")}{t("bizMigration.colVerdict")}
{d.key_str || "—"}{d.new_key_str || "—"}{d.old_status || "—"}{d.new_status || "—"} + {verdictLabel(d.verdict)} +
+ {t("bizMigration.emptyDiffs")} +
+
+
+ )} +
+
+
+ ); +} diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index df7a3b6..a3edb8a 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -7,8 +7,10 @@ import { useDebouncedValue } from "../../hooks/useDebouncedValue"; import { useToast } from "../../hooks/useToast"; import { useI18n } from "../../i18n"; import { + bizStateBulkDeleteBatches, bizStateCollectNow, bizStateCreateTask, + bizStateDeleteBatch, bizStateDeleteTask, bizStateDiscover, bizStateDownloadExport, @@ -18,6 +20,8 @@ import { bizStateListProfiles, bizStateListTasks, bizStatePatchTask, + bizStatePurgeTask, + bizStateSetBatchBaseline, bizStateSetBindings, fetchCliTargets, formatErr, @@ -33,6 +37,7 @@ type TaskRow = { ne_name: string; ne_ip: string; vendor: string; + note?: string; status: string; collect_running: boolean; last_error: string; @@ -57,6 +62,9 @@ type BatchRow = { row_count: number; command_count: number; started_at?: string | null; + is_baseline?: boolean; + protected?: boolean; + protect_reasons?: string[]; }; type Candidate = { value: string; label: string; rd?: string }; @@ -269,7 +277,10 @@ export function BizStatePage() { const [taskTab, setTaskTab] = useState("profiles"); const [intervalValue, setIntervalValue] = useState(1); const [intervalUnit, setIntervalUnit] = useState<"days" | "hours" | "seconds">("hours"); - const [retentionBatches, setRetentionBatches] = useState(30); + const [retentionDays, setRetentionDays] = useState(30); + const [dailyKeepEnabled, setDailyKeepEnabled] = useState(false); + const [dailyKeepCount, setDailyKeepCount] = useState(10); + const [selectedBatchIds, setSelectedBatchIds] = useState([]); // VRF bind (inside task modal) const [bindItemId, setBindItemId] = useState(""); @@ -333,7 +344,7 @@ export function BizStatePage() { if (!kw) return tasks; return tasks.filter((row) => { const blob = - `${row.ne_name} ${row.ne_ip} ${row.vendor} ${row.status} ${row.source || ""} ${row.last_error}`.toLowerCase(); + `${row.ne_name} ${row.ne_ip} ${row.vendor} ${row.note || ""} ${row.status} ${row.source || ""} ${row.last_error}`.toLowerCase(); return blob.includes(kw); }); }, [tasks, debouncedListKw]); @@ -382,9 +393,12 @@ export function BizStatePage() { const ui = secToIntervalUi(Number(task.interval_sec || 3600)); setIntervalValue(ui.value); setIntervalUnit(ui.unit); - setRetentionBatches(Math.max(1, Number(task.retention_batches || 30))); - const b = await bizStateListBatches(id); + setRetentionDays(Math.max(1, Number(task.retention_days || 30))); + setDailyKeepEnabled(Boolean(task.daily_keep_enabled)); + setDailyKeepCount(Math.max(1, Number(task.daily_keep_count || 10))); + const b = await bizStateListBatches(id, 200); setBatches((b.items || []) as BatchRow[]); + setSelectedBatchIds([]); const p = await bizStateListProfiles({ vendor: task.vendor || "", device_type: task.device_type || "", @@ -457,7 +471,9 @@ export function BizStatePage() { try { await bizStatePatchTask(taskId, { interval_sec: intervalUiToSec(intervalValue, intervalUnit), - retention_batches: Math.max(1, Number(retentionBatches) || 30), + retention_days: Math.max(1, Number(retentionDays) || 30), + daily_keep_enabled: dailyKeepEnabled, + daily_keep_count: Math.max(1, Number(dailyKeepCount) || 10), }); showOk(t("bizState.scheduleSaved")); await loadTask(taskId); @@ -469,6 +485,82 @@ export function BizStatePage() { } }; + const toggleBatchSelect = (id: string, protectedBatch: boolean) => { + if (protectedBatch) return; + setSelectedBatchIds((prev) => + prev.includes(id) ? prev.filter((x) => x !== id) : [...prev, id], + ); + }; + + const selectDeletableBatches = () => { + setSelectedBatchIds(batches.filter((b) => !b.protected).map((b) => b.id)); + }; + + const markBaseline = async (batchId: string, marked: boolean) => { + setBusy(true); + try { + await bizStateSetBatchBaseline(batchId, marked); + if (taskId) await loadTask(taskId); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + }; + + const removeBatch = async (batchId: string) => { + if (!window.confirm(t("bizState.confirmDeleteBatch"))) return; + setBusy(true); + try { + await bizStateDeleteBatch(batchId); + showOk(t("bizState.batchDeleted")); + if (taskId) await loadTask(taskId); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + }; + + const bulkRemoveBatches = async () => { + if (!selectedBatchIds.length) return; + if ( + !window.confirm( + t("bizState.confirmBulkDelete").replace("{{count}}", String(selectedBatchIds.length)), + ) + ) { + return; + } + setBusy(true); + try { + const res = await bizStateBulkDeleteBatches(selectedBatchIds); + showOk( + t("bizState.bulkDeleted") + .replace("{{deleted}}", String(res.deleted_count || 0)) + .replace("{{skipped}}", String((res.skipped || []).length)), + ); + if (taskId) await loadTask(taskId); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + }; + + const purgeNow = async () => { + if (!taskId) return; + setBusy(true); + try { + const res = await bizStatePurgeTask(taskId); + showOk(t("bizState.purgeOk").replace("{{dropped}}", String(res.dropped || 0))); + await loadTask(taskId); + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + }; + const collectNow = async () => { if (!taskId) return; setBusy(true); @@ -670,6 +762,7 @@ export function BizStatePage() {
{row.ne_name || row.ne_ip || "—"}
{row.vendor || "—"} · {row.ne_ip || "—"} + {row.note ? ` · ${row.note}` : ""}
{row.last_error ?
{row.last_error}
: null} @@ -939,16 +1032,41 @@ export function BizStatePage() { setRetentionBatches(Math.max(1, Number(e.target.value) || 1))} + max={3650} + value={String(retentionDays)} + onChange={(e) => setRetentionDays(Math.max(1, Number(e.target.value) || 1))} /> + + {dailyKeepEnabled ? ( + + ) : null} + - {detail.vendor || "—"} · {detail.ne_ip || "—"} · {t("bizState.scheduleHint")} + {detail.vendor || "—"} · {detail.ne_ip || "—"} · {t("bizState.retentionHint")} + {dailyKeepEnabled ? ` · ${t("bizState.dailyKeepHint")}` : ""} ) : null} @@ -1086,41 +1204,104 @@ export function BizStatePage() { ) : (
+
+ + +
+ + - {batches.map((b) => ( - - - - - - - ))} + {batches.map((b) => { + const locked = Boolean(b.protected); + return ( + + + + + + + + + ); + })} {!batches.length ? ( - diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 537c92c..3e00979 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -1813,6 +1813,29 @@ export const bizStateListBatches = (taskId: string, limit = 50) => `/v1/biz-state/tasks/${encodeURIComponent(taskId)}/batches?limit=${limit}`, ); +export const bizStateSetBatchBaseline = (batchId: string, marked: boolean) => + apiPost>( + `/v1/biz-state/batches/${encodeURIComponent(batchId)}/baseline`, + { marked }, + ); + +export const bizStateDeleteBatch = (batchId: string) => + apiDelete<{ ok: boolean }>(`/v1/biz-state/batches/${encodeURIComponent(batchId)}`); + +export const bizStateBulkDeleteBatches = (batchIds: string[]) => + apiPost<{ + ok: boolean; + deleted: string[]; + skipped: { batch_id: string; reasons: string[] }[]; + deleted_count: number; + }>("/v1/biz-state/batches/bulk-delete", { batch_ids: batchIds }); + +export const bizStatePurgeTask = (taskId: string) => + apiPost>( + `/v1/biz-state/tasks/${encodeURIComponent(taskId)}/purge`, + {}, + ); + export const bizStateGetBatch = (batchId: string) => apiGet>(`/v1/biz-state/batches/${encodeURIComponent(batchId)}`); @@ -1959,3 +1982,117 @@ export const bizCompareDownloadRun = async (runId: string): Promise => { URL.revokeObjectURL(url); } }; + +// --- biz-migration (cutover monitor, separate from compare) --- + +export const bizMigrationListProjects = () => + apiGet<{ items: Record[] }>("/v1/biz-migration/projects"); + +export const bizMigrationCreateProject = (body: Record) => + apiPost>("/v1/biz-migration/projects", body); + +export const bizMigrationPatchProject = (projectId: string, body: Record) => + apiPatch>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}`, + body, + ); + +export const bizMigrationListBatches = (projectId: string) => + apiGet<{ items: Record[] }>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}/batches`, + ); + +export const bizMigrationListBaselinePorts = (projectId: string) => + apiGet<{ + batch_id: string; + metric_id: string; + ports: { + interface: string; + admin?: string; + phy?: string; + prot?: string; + description?: string; + mapped_to?: string; + }[]; + }>(`/v1/biz-migration/projects/${encodeURIComponent(projectId)}/baseline-ports`); + +export const bizMigrationEnsurePortHighfreq = ( + projectId: string, + body: { interval_sec?: number; retention_days?: number; collect_now?: boolean } = {}, +) => + apiPost>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}/ensure-port-highfreq`, + body, + ); + +export const bizMigrationCollectNow = (projectId: string) => + apiPost>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}/collect-now`, + {}, + ); + +export const bizMigrationFinishBatch = (batchId: string, markDone = false) => + apiPost>( + `/v1/biz-migration/batches/${encodeURIComponent(batchId)}/finish?mark_done=${markDone ? "true" : "false"}`, + {}, + ); + +export const bizMigrationListRedTickets = (projectId: string, status = "") => { + const q = status ? `?status=${encodeURIComponent(status)}` : ""; + return apiGet<{ open_count: number; items: Record[] }>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}/red-tickets${q}`, + ); +}; + +export const bizMigrationResolveRedTicket = (ticketId: string, note = "") => + apiPost>( + `/v1/biz-migration/red-tickets/${encodeURIComponent(ticketId)}/resolve`, + { note }, + ); + +export const bizMigrationCreateBatch = (projectId: string, body: Record) => + apiPost>( + `/v1/biz-migration/projects/${encodeURIComponent(projectId)}/batches`, + body, + ); + +export const bizMigrationPatchBatch = (batchId: string, body: Record) => + apiPatch>( + `/v1/biz-migration/batches/${encodeURIComponent(batchId)}`, + body, + ); + +export const bizMigrationEvaluate = (batchId: string, body: Record = {}) => + apiPost>( + `/v1/biz-migration/batches/${encodeURIComponent(batchId)}/evaluate`, + body, + ); + +export const bizMigrationGetBoard = (batchId: string, runId = "") => { + const q = runId ? `?run_id=${encodeURIComponent(runId)}` : ""; + return apiGet>( + `/v1/biz-migration/batches/${encodeURIComponent(batchId)}/board${q}`, + ); +}; + +export const bizMigrationListDiffs = (params: { + runId: string; + metricId?: string; + verdict?: string; + color?: string; + kw?: string; + offset?: number; + limit?: number; +}) => { + const p = new URLSearchParams(); + if (params.metricId) p.set("metric_id", params.metricId); + if (params.verdict) p.set("verdict", params.verdict); + if (params.color) p.set("color", params.color); + if (params.kw) p.set("kw", params.kw); + if (params.offset != null) p.set("offset", String(params.offset)); + if (params.limit != null) p.set("limit", String(params.limit)); + const qs = p.toString(); + return apiGet<{ items: Record[]; total: number }>( + `/v1/biz-migration/runs/${encodeURIComponent(params.runId)}/diffs${qs ? `?${qs}` : ""}`, + ); +};
{t("bizState.colTime")} {t("bizState.colStatus")}{t("bizState.colProtect")} {t("bizState.colRows")} {t("bizState.colActions")}
{fmtTime(b.started_at)} - {b.status} - - {b.row_count} - / {b.command_count} cmd - -
- - -
-
+ toggleBatchSelect(b.id, locked)} + /> + {fmtTime(b.started_at)} + {b.status} + + {b.is_baseline ? ( + {t("bizState.baseline")} + ) : locked ? ( + {t("bizState.protected")} + ) : ( + "—" + )} + + {b.row_count} + / {b.command_count} cmd + +
+ + + {b.is_baseline ? ( + + ) : ( + + )} + +
+
+
{t("bizState.noBatches")}