Wire cutover monitor to templates: project binding, multi-sheet evaluate, dual status semantics, and UI expect/HF.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-18 20:54:07 +08:00
parent 202196a197
commit 60a663e2f6
11 changed files with 808 additions and 159 deletions

View file

@ -24,6 +24,39 @@ def port_status_label(row: dict[str, Any] | None) -> str:
return "/".join(parts)
def status_label(row: dict[str, Any] | None, sheet_override: dict[str, Any] | None = None) -> str:
"""Human status for board; uses sheet_override.status_fields when set."""
ov = sheet_override or {}
fields = [str(f) for f in (ov.get("status_fields") or []) if str(f).strip()]
if not fields:
return port_status_label(row)
if not row:
return "—"
parts = [str(row.get(f) or "-").lower() for f in fields]
if all(p == "-" for p in parts):
return "—"
return "/".join(parts)
def classify_status(row: dict[str, Any] | None, sheet_override: dict[str, Any] | None) -> str:
"""Classify current row as up|down|other|none from sheet_override status semantics."""
ov = sheet_override or {}
fields = [str(f) for f in (ov.get("status_fields") or []) if str(f).strip()]
if not fields or not row:
return "none"
down_vals = {str(x).lower() for x in (ov.get("down_values") or ["down"])}
up_vals = {str(x).lower() for x in (ov.get("up_values") or ["up"])}
vals = [str(row.get(f) or "").strip().lower() for f in fields]
vals = [v for v in vals if v]
if not vals:
return "none"
if any(v in down_vals for v in vals):
return "down"
if vals and all(v in up_vals for v in vals):
return "up"
return "other"
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).
@ -86,6 +119,34 @@ def side_verdict(
return "ok", "gray"
def _side_tokens(kind: str, status: str) -> set[str]:
toks: set[str] = set()
k = str(kind or "").strip()
if k:
toks.add(k)
s = str(status or "").strip()
if s and s != "none":
toks.add(s)
return toks
def _match_success(
old_tokens: set[str],
new_tokens: set[str],
success_patterns: list[dict[str, Any]] | None,
) -> bool:
for pat in success_patterns or []:
if not isinstance(pat, dict):
continue
old_need = {str(x) for x in (pat.get("old") or []) if str(x)}
new_need = {str(x) for x in (pat.get("new") or []) if str(x)}
if not old_need or not new_need:
continue
if (old_need & old_tokens) and (new_need & new_tokens):
return True
return False
def dual_verdict(
*,
old_kind: str,
@ -93,12 +154,25 @@ def dual_verdict(
in_expect: bool,
window_active: bool,
acceptance: bool = False,
old_status: str = "none",
new_status: str = "none",
success_patterns: list[dict[str, Any]] | None = None,
) -> tuple[str, str]:
"""Synthesize old+new into migration board verdict.
``acceptance=True`` (本批完成终验): unfinished expect items become red
(``unfinished`` / ``lost``), not yellow migrating.
When ``success_patterns`` is set (from monitor sheet_overrides), tokens may
include kinds and status classes (up/down) so e.g. old down + new up → migrated.
"""
if in_expect and _match_success(
_side_tokens(old_kind, old_status),
_side_tokens(new_kind, new_status),
success_patterns,
):
return "migrated", "green"
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"):
@ -179,6 +253,17 @@ def _current_row(diff: dict[str, Any] | None) -> dict[str, Any]:
return dict(diff.get("after") or {})
def override_for_metric(
sheet_overrides: list[dict[str, Any]] | None,
metric_id: str,
) -> dict[str, Any]:
mid = str(metric_id or "").strip()
for ov in sheet_overrides or []:
if isinstance(ov, dict) and str(ov.get("metric_id") or "").strip() == mid:
return ov
return {}
def evaluate_metric_dual(
*,
metric_id: str,
@ -193,9 +278,13 @@ def evaluate_metric_dual(
expect: dict[str, set[str]],
window_active: bool,
acceptance: bool = False,
field_rules: list[dict[str, Any]] | None = None,
sheet_override: dict[str, Any] | None = None,
) -> 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)
ov = sheet_override or {}
success_patterns = list(ov.get("success") or []) if isinstance(ov.get("success"), list) else []
old_cmp = compare_rows(
before_rows=old_baseline_rows,
@ -204,6 +293,7 @@ def evaluate_metric_dual(
iface_fields=[],
compare_fields=compare_fields,
port_map=None,
field_rules=field_rules,
)
old_idx = build_diff_index_from_compare(old_cmp)
@ -224,6 +314,7 @@ def evaluate_metric_dual(
iface_fields=[],
compare_fields=compare_fields,
port_map=None,
field_rules=field_rules,
)
new_idx = build_diff_index_from_compare(new_cmp)
@ -267,21 +358,26 @@ def evaluate_metric_dual(
if nd is None:
new_kind = ""
old_cur = _current_row(od)
new_cur = _current_row(nd)
old_st = classify_status(old_cur, ov)
new_st = classify_status(new_cur, ov)
verdict, color = dual_verdict(
old_kind=old_kind or "",
new_kind=new_kind or "",
in_expect=in_exp,
window_active=window_active,
acceptance=acceptance,
old_status=old_st,
new_status=new_st,
success_patterns=success_patterns or None,
)
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,
@ -297,10 +393,10 @@ def evaluate_metric_dual(
"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)
"old_status": status_label(old_cur, ov)
if old_cur
else ("gone" if old_kind == "removed" else "—"),
"new_status": port_status_label(new_cur)
"new_status": status_label(new_cur, ov)
if new_cur
else ("gone" if new_kind == "removed" else "—"),
"old_side": side_verdict(

View file

@ -123,6 +123,27 @@ def ensure_default_monitor_templates(db: Session) -> None:
db.commit()
def default_port_monitor_template_id(db: Session) -> str:
"""Ensure seeds exist and return id of 「端口割接监控」 (or first template)."""
ensure_default_monitor_templates(db)
row = (
db.query(BizMonitorTemplate)
.filter(BizMonitorTemplate.name == "端口割接监控")
.one_or_none()
)
if row:
return row.id
first = db.query(BizMonitorTemplate).order_by(BizMonitorTemplate.created_at.asc()).first()
return first.id if first else ""
def get_monitor_template_row(db: Session, template_id: str) -> BizMonitorTemplate | None:
tid = str(template_id or "").strip()
if not tid:
return None
return db.get(BizMonitorTemplate, tid)
def list_monitor_templates(db: Session) -> list[dict[str, Any]]:
ensure_default_monitor_templates(db)
names = _compare_name_map(db)

View file

@ -9,19 +9,28 @@ 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 ..biz_state.compare_service import _load_metric_rows, _port_map_dict, template_metrics
from ..models import (
BizCompareTemplate,
BizMigrationBatch,
BizMigrationDiff,
BizMigrationProject,
BizMigrationRedTicket,
BizMigrationRun,
BizMonitorTemplate,
BizPortMapping,
BizStateBatch,
BizStateTask,
)
from ..timeutil import utcnow_naive
from .evaluate import PORT_METRIC_ID, evaluate_metric_dual, parse_expect_set, port_sheet_def
from . import monitor_templates as mon_tpl
from .evaluate import (
PORT_METRIC_ID,
evaluate_metric_dual,
override_for_metric,
parse_expect_set,
port_sheet_def,
)
def _task_brief(db: Session, task_id: str) -> dict[str, Any]:
@ -52,6 +61,72 @@ def _batch_brief(db: Session, batch_id: str) -> dict[str, Any]:
}
def _monitor_template_brief(db: Session, template_id: str) -> dict[str, Any]:
tid = str(template_id or "").strip()
if not tid:
return {"id": "", "name": "", "compare_template_id": "", "compare_template_name": ""}
row = db.get(BizMonitorTemplate, tid)
if not row:
return {"id": tid, "name": "", "compare_template_id": "", "compare_template_name": ""}
cmp_name = ""
if row.compare_template_id:
ct = db.get(BizCompareTemplate, row.compare_template_id)
cmp_name = (ct.name if ct else "") or ""
return {
"id": row.id,
"name": row.name or "",
"compare_template_id": row.compare_template_id or "",
"compare_template_name": cmp_name,
"collect_metric_ids": list(row.collect_metric_ids_json or []),
}
def resolve_project_monitor_template(db: Session, proj: BizMigrationProject) -> BizMonitorTemplate:
"""Return monitor template for project; seed default port template if unbound."""
mon_tpl.ensure_default_monitor_templates(db)
tid = str(getattr(proj, "monitor_template_id", None) or "").strip()
row = db.get(BizMonitorTemplate, tid) if tid else None
if row:
return row
default_id = mon_tpl.default_port_monitor_template_id(db)
row = db.get(BizMonitorTemplate, default_id) if default_id else None
if not row:
raise HTTPException(status_code=400, detail="monitor_template_required")
# Persist default on first use so UI shows binding
if not tid:
proj.monitor_template_id = row.id
proj.updated_at = utcnow_naive()
db.commit()
db.refresh(proj)
return row
def resolve_evaluate_sheets(
db: Session, mt: BizMonitorTemplate
) -> tuple[list[dict[str, Any]], list[dict[str, Any]], dict[str, Any]]:
"""Return (sheets, sheet_overrides, defaults) from monitor → compare template."""
overrides = list(mt.sheet_overrides_json or []) if isinstance(mt.sheet_overrides_json, list) else []
defaults = dict(mt.defaults_json or {}) if isinstance(mt.defaults_json, dict) else {}
cid = str(mt.compare_template_id or "").strip()
if cid:
ct = db.get(BizCompareTemplate, cid)
if ct:
sheets = template_metrics(ct)
if sheets:
return sheets, overrides, defaults
# Fallback: built-in port sheet (legacy)
return [port_sheet_def()], overrides, defaults
def resolve_collect_metric_ids(db: Session, proj: BizMigrationProject) -> list[str]:
mt = resolve_project_monitor_template(db, proj)
collect = [str(x).strip() for x in (mt.collect_metric_ids_json or []) if str(x).strip()]
if collect:
return collect
sheets, _, _ = resolve_evaluate_sheets(db, mt)
return [str(s.get("metric_id") or "").strip() for s in sheets if str(s.get("metric_id") or "").strip()]
def project_to_dict(db: Session, p: BizMigrationProject) -> dict[str, Any]:
return {
"id": p.id,
@ -61,6 +136,8 @@ def project_to_dict(db: Session, p: BizMigrationProject) -> dict[str, Any]:
"old_baseline_batch_id": p.old_baseline_batch_id,
"new_baseline_batch_id": p.new_baseline_batch_id,
"mapping_id": p.mapping_id,
"monitor_template_id": getattr(p, "monitor_template_id", None) or "",
"monitor_template": _monitor_template_brief(db, getattr(p, "monitor_template_id", None) or ""),
"status": p.status,
"note": p.note,
"old_task": _task_brief(db, p.old_task_id),
@ -108,6 +185,12 @@ def create_project(db: Session, body: dict[str, Any]) -> dict[str, Any]:
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")
monitor_template_id = str(body.get("monitor_template_id") or "").strip()
if monitor_template_id:
if not db.get(BizMonitorTemplate, monitor_template_id):
raise HTTPException(status_code=404, detail="monitor_template_not_found")
else:
monitor_template_id = mon_tpl.default_port_monitor_template_id(db)
p = BizMigrationProject(
id=uuid4().hex,
name=name,
@ -116,6 +199,7 @@ def create_project(db: Session, body: dict[str, Any]) -> dict[str, Any]:
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,
monitor_template_id=monitor_template_id,
status=str(body.get("status") or "draft").strip() or "draft",
note=str(body.get("note") or "")[:500],
)
@ -157,6 +241,11 @@ def patch_project(db: Session, project_id: str, body: dict[str, Any]) -> dict[st
if bid and not db.get(BizStateBatch, bid):
raise HTTPException(status_code=404, detail="batch_not_found")
p.new_baseline_batch_id = bid
if "monitor_template_id" in body and body["monitor_template_id"] is not None:
mid = str(body["monitor_template_id"] or "").strip()
if mid and not db.get(BizMonitorTemplate, mid):
raise HTTPException(status_code=404, detail="monitor_template_not_found")
p.monitor_template_id = mid
p.updated_at = utcnow_naive()
db.commit()
db.refresh(p)
@ -270,7 +359,8 @@ def run_evaluate(
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()]
mt = resolve_project_monitor_template(db, proj)
sheets, sheet_overrides, _defaults = resolve_evaluate_sheets(db, mt)
sheet_cards: list[dict[str, Any]] = []
all_rows: list[dict[str, Any]] = []
@ -278,13 +368,15 @@ def run_evaluate(
verdict_counts: dict[str, int] = {}
for sheet in sheets:
mid = sheet["metric_id"]
mid = str(sheet.get("metric_id") or "").strip()
key_fields = list(sheet.get("key_fields") or [])
if not key_fields:
if not mid or 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 [])
field_rules = list(sheet.get("field_rules") or [])
sheet_ov = override_for_metric(sheet_overrides, mid)
old_base = apply_row_filters(
_load_metric_rows(db, batch_id=proj.old_baseline_batch_id, metric_id=mid),
@ -314,15 +406,17 @@ def run_evaluate(
old_current_rows=old_now,
new_baseline_rows=new_base_rows,
new_current_rows=new_now,
port_map=port_map,
port_map=port_map if iface_fields else {},
expect=expect,
window_active=window_active,
acceptance=acceptance,
field_rules=field_rules,
sheet_override=sheet_ov,
)
sheet_cards.append(
{
"metric_id": mid,
"title": "端口状态",
"title": mid,
"progress_ok": one["progress_ok"],
"progress_total": one["progress_total"],
"anomaly": one["anomaly"],
@ -347,7 +441,9 @@ def run_evaluate(
purpose=str(purpose or ("acceptance" if acceptance else "manual"))[:32],
status="success",
summary_json={
"metric_focus": PORT_METRIC_ID,
"metric_focus": sheets[0].get("metric_id") if sheets else PORT_METRIC_ID,
"monitor_template_id": mt.id,
"compare_template_id": mt.compare_template_id or "",
"acceptance": acceptance,
"sheet_cards": sheet_cards,
"progress": {
@ -486,25 +582,86 @@ def list_baseline_ports(db: Session, project_id: str) -> dict[str, Any]:
}
def _iface_brief_item(*, vendor: str, device_type: str) -> dict[str, Any]:
def list_baseline_expect_objects(db: Session, project_id: str) -> dict[str, Any]:
"""Per-sheet baseline keys for multi-metric expect 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": "", "sheets": [], "mapped": {}}
mt = resolve_project_monitor_template(db, p)
sheets, _, _ = resolve_evaluate_sheets(db, mt)
port_map = _port_map_dict(db, p.mapping_id)
out_sheets: list[dict[str, Any]] = []
for sheet in sheets:
mid = str(sheet.get("metric_id") or "").strip()
key_fields = [str(k) for k in (sheet.get("key_fields") or []) if str(k).strip()]
if not mid or not key_fields:
continue
iface_fields = [str(k) for k in (sheet.get("iface_fields") or []) if str(k).strip()]
rows_raw = apply_row_filters(
_load_metric_rows(db, batch_id=p.old_baseline_batch_id, metric_id=mid),
list(sheet.get("row_filters") or []),
)
items: list[dict[str, Any]] = []
for r in rows_raw:
keys = [str(r.get(f) or "").strip() for f in key_fields]
if not any(keys):
continue
key_str = "|".join(keys)
mapped = ""
if iface_fields and key_fields and key_fields[0] in iface_fields:
mapped = port_map.get(keys[0]) or ""
items.append(
{
"key": key_str,
"keys": keys,
"mapped_to": mapped,
"label": key_str,
"row": {f: r.get(f) for f in list(dict.fromkeys([*key_fields, *iface_fields, "description", "admin", "phy", "prot"])) if f in r},
}
)
items.sort(key=lambda x: str(x["key"]))
out_sheets.append(
{
"metric_id": mid,
"key_fields": key_fields,
"iface_fields": iface_fields,
"items": items,
}
)
return {
"batch_id": p.old_baseline_batch_id,
"monitor_template_id": mt.id,
"sheets": out_sheets,
"mapped": port_map,
}
def _catalog_item_for_metric(*, vendor: str, device_type: str, metric_id: str) -> dict[str, Any]:
from ..biz_state.profiles import profiles_for_vendor
from ..lldp_shared import resolve_vendor_key
mid = str(metric_id or "").strip()
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":
if p.metric_id == mid and p.kind == "collect":
return {
"source_profile_id": p.profile_id,
"kind": "catalog",
"enabled": True,
"title": p.title or "interface brief",
"title": p.title or mid,
}
raise HTTPException(
status_code=400,
detail=f"no_interface_brief_profile_for_vendor:{vkey or vendor or 'unknown'}",
detail=f"no_profile_for_metric:{mid}:{vkey or vendor or 'unknown'}",
)
def _iface_brief_item(*, vendor: str, device_type: str) -> dict[str, Any]:
return _catalog_item_for_metric(vendor=vendor, device_type=device_type, metric_id=PORT_METRIC_ID)
def _enabled_metric_ids(db: Session, task_id: str) -> set[str]:
from ..biz_state.profiles import get_profile
from ..models import BizStateTaskItem
@ -525,10 +682,15 @@ def _enabled_metric_ids(db: Session, task_id: str) -> set[str]:
return out
def _is_port_highfreq_task(db: Session, task: BizStateTask) -> bool:
"""True when task is interface_brief-only with short interval (cutover HF)."""
def _is_highfreq_task(db: Session, task: BizStateTask, want_metrics: set[str]) -> bool:
"""True when task metrics match want set and interval is short (cutover HF)."""
metrics = _enabled_metric_ids(db, task.id)
return metrics == {PORT_METRIC_ID} and int(task.interval_sec or 0) <= 300
return metrics == set(want_metrics) and int(task.interval_sec or 0) <= 300
def _is_port_highfreq_task(db: Session, task: BizStateTask) -> bool:
"""Back-compat: interface_brief-only HF."""
return _is_highfreq_task(db, task, {PORT_METRIC_ID})
def _ensure_side_highfreq(
@ -538,19 +700,31 @@ def _ensure_side_highfreq(
project_name: str,
interval_sec: int,
retention_days: int,
metric_ids: list[str],
) -> tuple[BizStateTask, bool]:
"""Return (task, created). Reuse if already HF port-only; else create sibling."""
"""Return (task, created). Reuse if already HF for metric set; else create sibling."""
from ..biz_state import service as biz_svc
if _is_port_highfreq_task(db, template):
want = {str(m).strip() for m in metric_ids if str(m).strip()}
if not want:
want = {PORT_METRIC_ID}
if _is_highfreq_task(db, template, want):
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]
items = [
_catalog_item_for_metric(
vendor=template.vendor,
device_type=template.device_type,
metric_id=mid,
)
for mid in sorted(want)
]
note = f"割接高频/{'+'.join(sorted(want)[:3])}/{project_name}"[:256]
created = biz_svc.create_task(
db,
{
@ -564,7 +738,7 @@ def _ensure_side_highfreq(
"status": "running",
"interval_sec": interval_sec,
"retention_days": retention_days,
"items": [item],
"items": items,
},
)
task = db.get(BizStateTask, str(created.get("id") or ""))
@ -581,11 +755,25 @@ def ensure_port_highfreq(
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.
"""Create/bind high-freq biz_state tasks for old/new NEs using monitor collect_metric_ids."""
return ensure_highfreq(
db,
project_id,
interval_sec=interval_sec,
retention_days=retention_days,
collect_now=collect_now,
)
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.
"""
def ensure_highfreq(
db: Session,
project_id: str,
*,
interval_sec: int = 60,
retention_days: int = 7,
collect_now: bool = True,
) -> dict[str, Any]:
"""Create/bind HF collect tasks from project's monitor template collect_metric_ids."""
from ..biz_state.collect_runner import dispatch_collect
proj = db.get(BizMigrationProject, project_id)
@ -596,6 +784,7 @@ def ensure_port_highfreq(
if not old_tpl or not new_tpl:
raise HTTPException(status_code=400, detail="old_new_task_required")
metric_ids = resolve_collect_metric_ids(db, proj)
iv = max(60, int(interval_sec or 60))
ret = max(1, int(retention_days or 7))
old_task, old_created = _ensure_side_highfreq(
@ -604,8 +793,8 @@ def ensure_port_highfreq(
project_name=proj.name,
interval_sec=iv,
retention_days=ret,
metric_ids=metric_ids,
)
# 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")
@ -615,6 +804,7 @@ def ensure_port_highfreq(
project_name=proj.name,
interval_sec=iv,
retention_days=ret,
metric_ids=metric_ids,
)
proj = db.get(BizMigrationProject, project_id)
@ -641,6 +831,7 @@ def ensure_port_highfreq(
"old_created": old_created,
"new_created": new_created,
"interval_sec": iv,
"collect_metric_ids": metric_ids,
"collect": collect,
}

View file

@ -22,6 +22,7 @@ class ProjectIn(BaseModel):
old_baseline_batch_id: str = ""
new_baseline_batch_id: str = ""
mapping_id: str = ""
monitor_template_id: str = ""
status: str = "draft"
note: str = ""
@ -31,6 +32,7 @@ class ProjectPatchIn(BaseModel):
old_baseline_batch_id: str | None = None
new_baseline_batch_id: str | None = None
mapping_id: str | None = None
monitor_template_id: str | None = None
status: str | None = None
note: str | None = None
@ -103,6 +105,12 @@ def api_baseline_ports(project_id: str, db: Session = Depends(get_db)):
return svc.list_baseline_ports(db, project_id)
@router.get("/projects/{project_id}/baseline-expect")
def api_baseline_expect(project_id: str, db: Session = Depends(get_db)):
"""Multi-metric baseline keys for expect-set picking (from monitor template sheets)."""
return svc.list_baseline_expect_objects(db, project_id)
class EnsureHighfreqIn(BaseModel):
interval_sec: int = 60
retention_days: int = 7
@ -110,14 +118,15 @@ class EnsureHighfreqIn(BaseModel):
@router.post("/projects/{project_id}/ensure-port-highfreq")
@router.post("/projects/{project_id}/ensure-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."""
"""Create/bind HF biz_state tasks from monitor template collect_metric_ids."""
payload = body.model_dump() if body else {}
return svc.ensure_port_highfreq(
return svc.ensure_highfreq(
db,
project_id,
interval_sec=int(payload.get("interval_sec") or 60),

View file

@ -61,6 +61,8 @@ def apply_biz_state_schema(conn: Connection) -> None:
")",
"CREATE INDEX IF NOT EXISTS ix_biz_monitor_template_name ON biz_monitor_template (name)",
"CREATE INDEX IF NOT EXISTS ix_biz_monitor_template_compare ON biz_monitor_template (compare_template_id)",
"ALTER TABLE biz_migration_project ADD COLUMN IF NOT EXISTS monitor_template_id VARCHAR(64) DEFAULT ''",
"CREATE INDEX IF NOT EXISTS ix_biz_migration_project_monitor_tpl ON biz_migration_project (monitor_template_id)",
):
try:
_run_sql(conn, sql)

View file

@ -26,6 +26,8 @@ class BizMigrationProject(Base):
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)
# BizMonitorTemplate — HOW (via compare template) + dual/status rules
monitor_template_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)