diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index a0605ff..9375d79 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -45,6 +45,13 @@ from .command_match import ( ) from .parsers import get_parser from .profiles import get_profile +from .collect_stop import ( + STOP_USER_MESSAGE, + clear_stop_requested, + is_stop_requested, + register_lane_holder, + unregister_lane_holder, +) _log = logging.getLogger("netx.biz_state.runner") @@ -457,17 +464,32 @@ def execute_claimed_batch( error = "" try: - _run_collect_session( - task_id=task_id, - batch_id=batch_id, - source=source, - ne_id=ne_id, - vendor=vendor, - device_type=device_type, - ) + if is_stop_requested(batch_id): + _finalize_batch_status( + batch_id=batch_id, + task_id=task_id, + cmd_count=0, + total_rows=0, + any_fail=False, + any_ok=False, + lane_errors=[], + stopped=True, + ) + error = STOP_USER_MESSAGE + else: + _run_collect_session( + task_id=task_id, + batch_id=batch_id, + source=source, + ne_id=ne_id, + vendor=vendor, + device_type=device_type, + ) except Exception as exc: _log.exception("biz_state collect failed task=%s batch=%s", task_id, batch_id) error = _format_error(exc) + if "_stopped" in error or is_stop_requested(batch_id): + error = STOP_USER_MESSAGE try: _fail_batch_status(batch_id, error) except Exception: @@ -509,10 +531,14 @@ def _run_collect_lane( budget = min(int(cap), int(per_cmd) * max(1, len(work)) + 90) holder: dict[str, Any] = {} flush_every = persist_every_cmds() + register_lane_holder(batch_id, holder) def _session() -> tuple[int, int, bool, bool]: from ..ne_netmiko import drain_read_channel + if is_stop_requested(batch_id): + raise TimeoutError(f"{label}_stopped") + conn = open_netmiko_connection(creds, session_timeout=budget) holder["conn"] = conn total_rows = 0 @@ -636,8 +662,10 @@ def _run_collect_lane( flat_work.append((cmd, p, profile_id, item_id, "normal")) for concrete, params, profile_id, item_id, mode in flat_work: - if holder.get("timed_out"): - raise TimeoutError(f"{label}_aborted") + if holder.get("timed_out") or holder.get("stop_requested") or is_stop_requested( + batch_id + ): + raise TimeoutError(f"{label}_stopped") cmd_id = uuid4().hex raw_text = "" @@ -923,6 +951,8 @@ def _run_collect_lane( ) except TimeoutError as exc: raise RuntimeError(str(exc)[:1020]) from exc + finally: + unregister_lane_holder(batch_id, holder) def _absorb_lane_result( @@ -1052,6 +1082,7 @@ def _finalize_batch_status( any_fail: bool, any_ok: bool, lane_errors: list[str], + stopped: bool = False, ) -> str: """Write terminal batch status on a fresh Session (retry once on disconnect).""" @@ -1063,7 +1094,14 @@ def _finalize_batch_status( batch.command_count = max(int(batch.command_count or 0), int(cmd_count or 0)) batch.row_count = max(int(batch.row_count or 0), int(total_rows or 0)) batch.ended_at = _utcnow() - if any_fail and any_ok: + if stopped: + if any_ok or int(batch.command_count or 0) > 0 or int(batch.row_count or 0) > 0: + batch.status = "partial" + batch.message = STOP_USER_MESSAGE + else: + batch.status = "cancelled" + batch.message = STOP_USER_MESSAGE + elif any_fail and any_ok: batch.status = "partial" if lane_errors: batch.message = "; ".join(lane_errors)[:1020] @@ -1077,7 +1115,7 @@ def _finalize_batch_status( batch.message = "" status = str(batch.status or "") db.commit() - if status in ("success", "partial"): + if status in ("success", "partial") and not stopped: try: from .compare_service import try_auto_compare_for_task @@ -1086,28 +1124,40 @@ def _finalize_batch_status( _log.exception("biz_state auto compare hook failed task=%s", task_id) return status - return str( - _run_db_with_reconnect(_write, label="biz_state_finalize") or "" - ) + try: + return str( + _run_db_with_reconnect(_write, label="biz_state_finalize") or "" + ) + finally: + clear_stop_requested(batch_id) def _fail_batch_status(batch_id: str, error: str) -> None: """Mark batch failed on a fresh Session (retry once on disconnect).""" msg = str(error or "")[:1020] + stopped = "_stopped" in msg or is_stop_requested(batch_id) def _write(db) -> None: batch = db.get(BizStateBatch, batch_id) if not batch: return - if str(batch.status or "") != "running": + st = str(batch.status or "") + if st not in ("running", "queued"): return - batch.status = "failed" - batch.message = msg + if stopped: + has_progress = int(batch.command_count or 0) > 0 or int(batch.row_count or 0) > 0 + batch.status = "partial" if has_progress else "cancelled" + batch.message = STOP_USER_MESSAGE + else: + batch.status = "failed" + batch.message = msg batch.ended_at = _utcnow() db.commit() - _run_db_with_reconnect(_write, label="biz_state_fail_batch") - + try: + _run_db_with_reconnect(_write, label="biz_state_fail_batch") + finally: + clear_stop_requested(batch_id) def _run_collect_session( *, @@ -1340,6 +1390,18 @@ def _run_collect_session( if lane_errors and not any_ok and cmd_count == 0: # Progressive bumps may already have cmds; only hard-fail if nothing landed. if not _batch_has_progress(batch_id): + if is_stop_requested(batch_id) or any("_stopped" in e for e in lane_errors): + _finalize_batch_status( + batch_id=batch_id, + task_id=task_id, + cmd_count=cmd_count, + total_rows=total_rows, + any_fail=False, + any_ok=False, + lane_errors=lane_errors, + stopped=True, + ) + return raise RuntimeError("; ".join(lane_errors)[:1020]) any_fail = True @@ -1358,6 +1420,7 @@ def _run_collect_session( except Exception: _log.exception("biz_state persist barrier before finalize failed batch=%s", batch_id) + stopped = is_stop_requested(batch_id) or any("_stopped" in e for e in lane_errors) _finalize_batch_status( batch_id=batch_id, task_id=task_id, @@ -1366,6 +1429,7 @@ def _run_collect_session( any_fail=any_fail, any_ok=any_ok, lane_errors=lane_errors, + stopped=stopped, ) def _purge(db) -> None: diff --git a/netx_api/biz_state/collect_stop.py b/netx_api/biz_state/collect_stop.py new file mode 100644 index 0000000..8c92db0 --- /dev/null +++ b/netx_api/biz_state/collect_stop.py @@ -0,0 +1,221 @@ +"""Stop an in-flight / queued biz_state collect batch. + +Queued batches are cancelled in DB immediately. Running lanes poll +``is_stop_requested`` between commands (DB-backed for dedicated workers) +and process-local holders are force-closed when stop is requested in the +same process. +""" + +from __future__ import annotations + +import logging +import threading +import time +from typing import Any + +from ..db import SessionLocal +from ..models import BizStateBatch, BizStateTask +from ..ne_session_factory import close_netmiko_connection +from ..timeutil import utcnow_naive + +_log = logging.getLogger("netx.biz_state.collect_stop") + +STOP_REQUEST_TOKEN = "stop_requested" +STOP_USER_MESSAGE = "stopped_by_user" + +_lock = threading.Lock() +_stop_batches: set[str] = set() +_holders: dict[str, list[dict[str, Any]]] = {} +# batch_id → (monotonic_ts, stop_requested) +_db_cache: dict[str, tuple[float, bool]] = {} +_DB_CACHE_TTL_SEC = 1.5 + + +def register_lane_holder(batch_id: str, holder: dict[str, Any]) -> None: + bid = str(batch_id or "").strip() + if not bid or holder is None: + return + with _lock: + _holders.setdefault(bid, []).append(holder) + + +def unregister_lane_holder(batch_id: str, holder: dict[str, Any]) -> None: + bid = str(batch_id or "").strip() + if not bid: + return + with _lock: + lst = _holders.get(bid) or [] + try: + lst.remove(holder) + except ValueError: + pass + if not lst: + _holders.pop(bid, None) + + +def mark_stop_requested(batch_id: str) -> None: + bid = str(batch_id or "").strip() + if not bid: + return + with _lock: + _stop_batches.add(bid) + _db_cache[bid] = (time.monotonic(), True) + + +def clear_stop_requested(batch_id: str) -> None: + bid = str(batch_id or "").strip() + if not bid: + return + with _lock: + _stop_batches.discard(bid) + _db_cache.pop(bid, None) + _holders.pop(bid, None) + + +def is_stop_requested(batch_id: str) -> bool: + """True if user asked to stop this batch (memory or DB).""" + bid = str(batch_id or "").strip() + if not bid: + return False + with _lock: + if bid in _stop_batches: + return True + hit = _db_cache.get(bid) + now = time.monotonic() + if hit and (now - hit[0]) < _DB_CACHE_TTL_SEC: + return bool(hit[1]) + + stopped = _read_stop_from_db(bid) + with _lock: + _db_cache[bid] = (time.monotonic(), stopped) + if stopped: + _stop_batches.add(bid) + return stopped + + +def _read_stop_from_db(batch_id: str) -> bool: + db = SessionLocal() + try: + batch = db.get(BizStateBatch, batch_id) + if not batch: + return False + st = str(batch.status or "") + if st in ("cancelled", "success", "partial", "failed"): + # Already terminal — treat as stop so lanes exit quickly. + return st == "cancelled" or STOP_REQUEST_TOKEN in str(batch.message or "") + msg = str(batch.message or "") + return msg.startswith(STOP_REQUEST_TOKEN) or msg == STOP_USER_MESSAGE + except Exception: + _log.exception("read stop flag failed batch=%s", batch_id) + return False + finally: + db.close() + + +def _force_close_holders(batch_id: str) -> int: + bid = str(batch_id or "").strip() + with _lock: + holders = list(_holders.get(bid) or []) + n = 0 + for holder in holders: + try: + holder["stop_requested"] = True + holder["timed_out"] = True + conn = holder.get("conn") + if conn is not None: + close_netmiko_connection(conn) + n += 1 + except Exception: + _log.exception("force-close on stop failed batch=%s", bid) + return n + + +def request_stop_collect(task_id: str) -> dict[str, Any]: + """Cancel queued batches and signal running collect for this task to abort.""" + tid = str(task_id or "").strip() + if not tid: + return {"ok": False, "reason": "missing_task_id"} + + db = SessionLocal() + cancelled_ids: list[str] = [] + signaled_ids: list[str] = [] + closed = 0 + try: + task = db.get(BizStateTask, tid) + if not task: + return {"ok": False, "reason": "task_not_found", "task_id": tid} + + active = ( + db.query(BizStateBatch) + .filter( + BizStateBatch.task_id == tid, + BizStateBatch.status.in_(("queued", "running")), + ) + .all() + ) + if not active: + # Nothing to stop — clear sticky collect_running if orphaned. + if bool(task.collect_running): + task.collect_running = False + if hasattr(task, "collect_queued_at"): + task.collect_queued_at = None + task.updated_at = utcnow_naive() + db.commit() + return { + "ok": True, + "task_id": tid, + "stopped": False, + "reason": "not_collecting", + "cancelled_batches": [], + "signaled_batches": [], + } + + now = utcnow_naive() + still_running = False + for batch in active: + bid = str(batch.id) + mark_stop_requested(bid) + if str(batch.status or "") == "queued": + batch.status = "cancelled" + batch.message = STOP_USER_MESSAGE + batch.ended_at = now + cancelled_ids.append(bid) + else: + batch.message = STOP_REQUEST_TOKEN + signaled_ids.append(bid) + still_running = True + closed += _force_close_holders(bid) + + if not still_running: + task.collect_running = False + if hasattr(task, "collect_queued_at"): + task.collect_queued_at = None + task.last_collect_ended_at = now + task.last_error = STOP_USER_MESSAGE + task.updated_at = now + + db.commit() + _log.info( + "stop collect task=%s cancelled=%s signaled=%s closed_conns=%s", + tid, + cancelled_ids, + signaled_ids, + closed, + ) + return { + "ok": True, + "task_id": tid, + "stopped": True, + "cancelled_batches": cancelled_ids, + "signaled_batches": signaled_ids, + "closed_connections": closed, + } + except Exception: + _log.exception("request_stop_collect failed task=%s", tid) + try: + db.rollback() + except Exception: + pass + return {"ok": False, "reason": "stop_error", "task_id": tid} + finally: + db.close() diff --git a/netx_api/biz_state_router.py b/netx_api/biz_state_router.py index 83dabdf..14ed384 100644 --- a/netx_api/biz_state_router.py +++ b/netx_api/biz_state_router.py @@ -252,6 +252,19 @@ def api_collect_now( return {"ok": True, "started": True, "queued": True, "task_id": task_id} +@router.post("/tasks/{task_id}/collect/stop") +def api_collect_stop(task_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + """Stop queued or running collect for this task (best-effort mid-command).""" + from .biz_state.collect_stop import request_stop_collect + + task = db.get(BizStateTask, task_id) + if not task: + from fastapi import HTTPException + + raise HTTPException(status_code=404, detail="task not found") + return request_stop_collect(task_id) + + @router.get("/tasks/{task_id}/batches") def api_list_batches( task_id: str, limit: int = 50, db: Session = Depends(get_db) diff --git a/netx_api/models/biz_state.py b/netx_api/models/biz_state.py index ee6e2a4..480e88a 100644 --- a/netx_api/models/biz_state.py +++ b/netx_api/models/biz_state.py @@ -90,7 +90,7 @@ class BizStateBatch(Base): ne_id: Mapped[str] = mapped_column(String(128), default="", index=True) ne_name: Mapped[str] = mapped_column(String(256), default="") vendor: Mapped[str] = mapped_column(String(64), default="") - status: Mapped[str] = mapped_column(String(32), default="running", index=True) # queued|running|success|partial|failed + status: Mapped[str] = mapped_column(String(32), default="running", index=True) # queued|running|success|partial|failed|cancelled command_count: Mapped[int] = mapped_column(Integer, default=0) row_count: Mapped[int] = mapped_column(Integer, default=0) message: Mapped[str] = mapped_column(String(1024), default="") diff --git a/tests/test_biz_state_collect_stop.py b/tests/test_biz_state_collect_stop.py new file mode 100644 index 0000000..546b271 --- /dev/null +++ b/tests/test_biz_state_collect_stop.py @@ -0,0 +1,111 @@ +"""biz_state collect stop: cancel queued / signal running.""" + +from __future__ import annotations + +import unittest +from unittest.mock import patch + +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker +from sqlalchemy.pool import StaticPool + +from netx_api.biz_state import collect_stop as stop_mod +from netx_api.db import Base +from netx_api.models import BizStateBatch, BizStateTask + + +class BizStateCollectStopTests(unittest.TestCase): + def setUp(self) -> None: + engine = create_engine( + "sqlite+pysqlite:///:memory:", + future=True, + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + TestingSession = sessionmaker( + bind=engine, autoflush=False, autocommit=False, expire_on_commit=False + ) + Base.metadata.create_all(bind=engine) + self.Session = TestingSession + self.db = TestingSession() + self._session_patch = patch.object(stop_mod, "SessionLocal", TestingSession) + self._session_patch.start() + stop_mod._stop_batches.clear() + stop_mod._holders.clear() + stop_mod._db_cache.clear() + + self.task = BizStateTask( + id="t-stop", + source="managed", + ne_id="ne-1", + ne_name="NE1", + status="running", + collect_running=True, + ) + self.db.add(self.task) + self.db.commit() + + def tearDown(self) -> None: + self._session_patch.stop() + self.db.close() + + def test_stop_cancels_queued_batch(self) -> None: + batch = BizStateBatch( + id="b-q", + task_id=self.task.id, + status="queued", + message="queued_manual", + ) + self.db.add(batch) + self.db.commit() + + out = stop_mod.request_stop_collect(self.task.id) + self.assertTrue(out.get("ok")) + self.assertTrue(out.get("stopped")) + self.assertEqual(out.get("cancelled_batches"), ["b-q"]) + self.assertEqual(out.get("signaled_batches"), []) + + self.db.expire_all() + b = self.db.get(BizStateBatch, "b-q") + t = self.db.get(BizStateTask, self.task.id) + assert b is not None and t is not None + self.assertEqual(b.status, "cancelled") + self.assertEqual(b.message, stop_mod.STOP_USER_MESSAGE) + self.assertFalse(t.collect_running) + + def test_stop_signals_running_batch(self) -> None: + batch = BizStateBatch( + id="b-r", + task_id=self.task.id, + status="running", + message="", + ) + self.db.add(batch) + self.db.commit() + + out = stop_mod.request_stop_collect(self.task.id) + self.assertTrue(out.get("ok")) + self.assertTrue(out.get("stopped")) + self.assertEqual(out.get("signaled_batches"), ["b-r"]) + self.assertTrue(stop_mod.is_stop_requested("b-r")) + + self.db.expire_all() + b = self.db.get(BizStateBatch, "b-r") + t = self.db.get(BizStateTask, self.task.id) + assert b is not None and t is not None + self.assertEqual(b.status, "running") + self.assertEqual(b.message, stop_mod.STOP_REQUEST_TOKEN) + # Still running until worker finishes. + self.assertTrue(t.collect_running) + + def test_stop_idle(self) -> None: + self.task.collect_running = False + self.db.commit() + out = stop_mod.request_stop_collect(self.task.id) + self.assertTrue(out.get("ok")) + self.assertFalse(out.get("stopped")) + self.assertEqual(out.get("reason"), "not_collecting") + + +if __name__ == "__main__": + unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 83bbe5e..6ed4ff3 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -195,6 +195,9 @@ const en = { scheduleOff: "Manual", scheduleHint: "When checked, collect on interval; otherwise only Collect now.", collectNow: "Collect now", + stopCollect: "Stop collect", + stopCollectOk: "Stop collect requested", + stopCollectIdle: "No collect in progress", delete: "Delete", confirmDelete: "Delete this task and all batches?", deleted: "Deleted", @@ -235,6 +238,11 @@ const en = { scheduleSaved: "Schedule saved", colInterval: "Interval", profiles: "Monitor items", + colLane: "Collect lane", + laneLight: "Light", + laneHeavy: "Heavy", + laneLightHint: "Shared SSH with most commands; ~10 min lane cap by default", + laneHeavyHint: "Dedicated SSH + longer timeouts; ~40 min lane cap by default", enable: "On", command: "Command template", colAux: "Aux commands", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 5014268..a2364d4 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -195,6 +195,9 @@ const zh = { scheduleOff: "手动", scheduleHint: "勾选后按周期自动采集;不勾选则仅支持「立即采集」。", collectNow: "立即采集", + stopCollect: "停止采集", + stopCollectOk: "已请求停止采集", + stopCollectIdle: "当前没有进行中的采集", delete: "删除", confirmDelete: "删除该业务监控任务及所有批次?", deleted: "已删除", @@ -235,6 +238,11 @@ const zh = { scheduleSaved: "周期配置已保存", colInterval: "周期", profiles: "监控项", + colLane: "采集车道", + laneLight: "轻车道", + laneHeavy: "重车道", + laneLightHint: "与多数命令共用 SSH;整轮默认约 10 分钟上限", + laneHeavyHint: "独立 SSH + 更长超时;整轮默认约 40 分钟上限", enable: "启用", command: "命令模板", colAux: "辅命令", diff --git a/web/src/pages/network/BizComparePage.tsx b/web/src/pages/network/BizComparePage.tsx index 5e16516..8ffddee 100644 --- a/web/src/pages/network/BizComparePage.tsx +++ b/web/src/pages/network/BizComparePage.tsx @@ -31,6 +31,7 @@ import { formatErr, } from "../../services/api"; import { formatSystemTime } from "../../utils/time"; +import { cutoverCachedGet, cutoverCachedGetSWR, invalidateCutoverCache } from "./cutoverDataCache"; import { jobChipColor, NmStatusChip } from "./nmChips"; type PageTab = "templates" | "jobs"; @@ -846,19 +847,36 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage const [navCollapsed, setNavCollapsed] = useState(false); const tplImportRef = useRef(null); - const refresh = useCallback(async () => { - const [taskRes, tpl, maps, j, met] = await Promise.all([ - bizStateListTasks(), - bizCompareListTemplates(), - bizCompareListMappings(), - bizCompareListJobs(), - bizCompareListMetrics(), - ]); - setTasks((taskRes.items || []) as TaskOpt[]); - setTemplates((tpl.items || []) as Template[]); - setMappings((maps.items || []) as Mapping[]); - setJobs((j.items || []) as Job[]); - setMetrics((met.items || []) as MetricSchema[]); + const refresh = useCallback(async (opts?: { force?: boolean }) => { + type Bundle = { + taskRes: Awaited>; + tpl: Awaited>; + maps: Awaited>; + j: Awaited>; + met: Awaited>; + }; + const fetchBundle = async (): Promise => { + const [taskRes, tpl, maps, j, met] = await Promise.all([ + bizStateListTasks(), + bizCompareListTemplates(), + bizCompareListMappings(), + bizCompareListJobs(), + bizCompareListMetrics(), + ]); + return { taskRes, tpl, maps, j, met }; + }; + const apply = (b: Bundle) => { + setTasks((b.taskRes.items || []) as TaskOpt[]); + setTemplates((b.tpl.items || []) as Template[]); + setMappings((b.maps.items || []) as Mapping[]); + setJobs((b.j.items || []) as Job[]); + setMetrics((b.met.items || []) as MetricSchema[]); + }; + if (opts?.force) { + apply(await cutoverCachedGet("bizCompare:lists", fetchBundle, { force: true })); + return; + } + apply(await cutoverCachedGetSWR("bizCompare:lists", fetchBundle, apply)); }, []); useEffect(() => { @@ -1457,7 +1475,9 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage else await bizCompareCreateTemplate(body); showOk(t("bizCompare.templateSaved")); setTplOpen(false); - await refresh(); + invalidateCutoverCache("bizMonitor:"); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1471,7 +1491,9 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage try { await bizCompareDeleteTemplate(id); showOk(t("bizCompare.templateDeleted")); - await refresh(); + invalidateCutoverCache("bizMonitor:"); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1519,7 +1541,9 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage })), }); showOk(t("bizCompare.templateImported")); - await refresh(); + invalidateCutoverCache("bizMonitor:"); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); if (pageMode !== "jobs") setPageTab("templates"); } catch (e) { showError(formatErr(e)); @@ -1536,7 +1560,7 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage await bizCompareDeleteJob(id); if (jobId === id) closeJob(); showOk(t("bizCompare.jobDeleted")); - await refresh(); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1577,7 +1601,7 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage setMappingId(String(m.id)); } showOk(t("bizCompare.mappingSaved")); - await refresh(); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1712,7 +1736,7 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage const j = await bizCompareCreateJob(jobConfigBody()); showOk(t("bizCompare.created")); closeCreateJob(); - await refresh(); + await refresh({ force: true }); await openJob(String(j.id)); } catch (e) { showError(formatErr(e)); @@ -1764,7 +1788,7 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage try { await bizCompareUpdateJob(jobId, jobConfigBody()); showOk(t("bizCompare.jobSaved")); - await refresh(); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1787,7 +1811,7 @@ export function BizComparePage({ pageMode = "all" }: { pageMode?: BizComparePage showOk(t("bizCompare.ran")); const r = await bizCompareListRuns(jobId); setRuns(r.items || []); - await refresh(); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { diff --git a/web/src/pages/network/BizMigrationPage.tsx b/web/src/pages/network/BizMigrationPage.tsx index 002ad3a..2ab2cd9 100644 --- a/web/src/pages/network/BizMigrationPage.tsx +++ b/web/src/pages/network/BizMigrationPage.tsx @@ -40,6 +40,7 @@ import type { CliTargetItem } from "../../types"; import { writeClipboardText } from "../../utils/clipboard"; import { pageCount } from "../../utils/display"; import { formatSystemTime, localDatetimeInputToUtcIso, utcIsoToLocalDatetimeInput } from "../../utils/time"; +import { cutoverCachedGet, cutoverCachedGetSWR, invalidateCutoverCache } from "./cutoverDataCache"; import { jobChipColor, NmStatusChip, sourceChipColor } from "./nmChips"; type NeSourceFilter = "all" | "managed" | "ume"; @@ -520,12 +521,17 @@ export function BizMigrationPage() { }, [expectSheets, expectMetricId, portFilter]); const reloadProjects = useCallback(async () => { - const res = await bizMigrationListProjects(); + const res = await cutoverCachedGet("bizMigration:projects", () => bizMigrationListProjects(), { + force: true, + }); setProjects((res.items || []) as Project[]); + invalidateCutoverCache("bizMigration:bootstrap"); }, []); const reloadMappings = useCallback(async () => { - const mp = await bizCompareListMappings(); + const mp = await cutoverCachedGet("bizMigration:mappings", () => bizCompareListMappings(), { + force: true, + }); setMappings( ((mp.items || []) as Record[]).map((x) => ({ id: String(x.id || ""), @@ -538,6 +544,7 @@ export function BizMigrationPage() { : [], })), ); + invalidateCutoverCache("bizMigration:bootstrap"); }, []); const parseMapRows = () => { @@ -626,45 +633,56 @@ export function BizMigrationPage() { useEffect(() => { void (async () => { try { - const [pr, mp, mt] = await Promise.all([ - bizMigrationListProjects(), - bizCompareListMappings(), - bizMonitorListTemplates(), - ]); - setProjects((pr.items || []) as Project[]); - setMappings( - ((mp.items || []) as Record[]).map((x) => ({ + type Bundle = { + pr: Awaited>; + mp: Awaited>; + mt: Awaited>; + }; + const fetchBundle = async (): Promise => { + const [pr, mp, mt] = await Promise.all([ + bizMigrationListProjects(), + bizCompareListMappings(), + bizMonitorListTemplates(), + ]); + return { pr, mp, mt }; + }; + const apply = (b: Bundle) => { + setProjects((b.pr.items || []) as Project[]); + setMappings( + ((b.mp.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + name: String(x.name || x.id || ""), + rows: Array.isArray(x.rows) + ? (x.rows as Array<{ before_if?: string; after_if?: string }>).map((r) => ({ + before_if: String(r.before_if || ""), + after_if: String(r.after_if || ""), + })) + : [], + })), + ); + const mts = ((b.mt.items || []) as Record[]).map((x) => ({ id: String(x.id || ""), name: String(x.name || x.id || ""), - rows: Array.isArray(x.rows) - ? (x.rows as Array<{ before_if?: string; after_if?: string }>).map((r) => ({ - before_if: String(r.before_if || ""), - after_if: String(r.after_if || ""), - })) + compare_template_id: String(x.compare_template_id || ""), + compare_template_name: String(x.compare_template_name || ""), + collect_metric_ids: Array.isArray(x.collect_metric_ids) + ? (x.collect_metric_ids as string[]) : [], - })), - ); - const mts = ((mt.items || []) as Record[]).map((x) => ({ - id: String(x.id || ""), - name: String(x.name || x.id || ""), - compare_template_id: String(x.compare_template_id || ""), - compare_template_name: String(x.compare_template_name || ""), - collect_metric_ids: Array.isArray(x.collect_metric_ids) - ? (x.collect_metric_ids as string[]) - : [], - collect_metric_ids_effective: Array.isArray(x.collect_metric_ids_effective) - ? (x.collect_metric_ids_effective as string[]) - : [], - })); - setMonitorTpls(mts); - if (!createMonitorTplId && mts.length) { - const preferred = - mts.find((x) => x.name === "默认割接监控") || - mts.find((x) => /默认|default|状态|status/i.test(x.name)) || - mts[0]; - setCreateMonitorTplId(preferred.id); - setCreateCollectMetricIds(monitorTplMetrics(preferred)); - } + collect_metric_ids_effective: Array.isArray(x.collect_metric_ids_effective) + ? (x.collect_metric_ids_effective as string[]) + : [], + })); + setMonitorTpls(mts); + if (!createMonitorTplId && mts.length) { + const preferred = + mts.find((x) => x.name === "默认割接监控") || + mts.find((x) => /默认|default|状态|status/i.test(x.name)) || + mts[0]; + setCreateMonitorTplId(preferred.id); + setCreateCollectMetricIds(monitorTplMetrics(preferred)); + } + }; + apply(await cutoverCachedGetSWR("bizMigration:bootstrap", fetchBundle, apply)); } catch (e) { showError(formatErr(e)); } diff --git a/web/src/pages/network/BizMonitorTemplatesPage.tsx b/web/src/pages/network/BizMonitorTemplatesPage.tsx index ee69735..3308205 100644 --- a/web/src/pages/network/BizMonitorTemplatesPage.tsx +++ b/web/src/pages/network/BizMonitorTemplatesPage.tsx @@ -15,6 +15,7 @@ import { bizMonitorUpdateTemplate, formatErr, } from "../../services/api"; +import { cutoverCachedGet, cutoverCachedGetSWR, invalidateCutoverCache } from "./cutoverDataCache"; type MetricField = { name: string; display_name?: string }; type MetricSchema = { metric_id: string; fields: MetricField[] }; @@ -828,29 +829,44 @@ export function BizMonitorTemplatesPage() { const [overridesText, setOverridesText] = useState("[]"); const tplImportRef = useRef(null); - const refresh = useCallback(async () => { - const [mon, cmp, metrics] = await Promise.all([ - bizMonitorListTemplates(), - bizCompareListTemplates(), - bizCompareListMetrics(), - ]); - setItems((mon.items || []) as MonitorTpl[]); - setCompareTpls( - ((cmp.items || []) as Record[]).map((x) => ({ - id: String(x.id || ""), - name: String(x.name || x.id || ""), - metrics: Array.isArray(x.metrics) ? (x.metrics as CompareSheet[]) : [], - })), - ); - setMetricSchemas( - ((metrics.items || []) as MetricSchema[]).map((m) => ({ - metric_id: m.metric_id, - fields: (m.fields || []).map((f) => ({ - name: f.name, - display_name: f.display_name, + const refresh = useCallback(async (opts?: { force?: boolean }) => { + type Bundle = { + mon: Awaited>; + cmp: Awaited>; + metrics: Awaited>; + }; + const fetchBundle = async (): Promise => { + const [mon, cmp, metrics] = await Promise.all([ + bizMonitorListTemplates(), + bizCompareListTemplates(), + bizCompareListMetrics(), + ]); + return { mon, cmp, metrics }; + }; + const apply = (b: Bundle) => { + setItems((b.mon.items || []) as MonitorTpl[]); + setCompareTpls( + ((b.cmp.items || []) as Record[]).map((x) => ({ + id: String(x.id || ""), + name: String(x.name || x.id || ""), + metrics: Array.isArray(x.metrics) ? (x.metrics as CompareSheet[]) : [], })), - })), - ); + ); + setMetricSchemas( + ((b.metrics.items || []) as MetricSchema[]).map((m) => ({ + metric_id: m.metric_id, + fields: (m.fields || []).map((f) => ({ + name: f.name, + display_name: f.display_name, + })), + })), + ); + }; + if (opts?.force) { + apply(await cutoverCachedGet("bizMonitor:lists", fetchBundle, { force: true })); + return; + } + apply(await cutoverCachedGetSWR("bizMonitor:lists", fetchBundle, apply)); }, []); useEffect(() => { @@ -1016,7 +1032,8 @@ export function BizMonitorTemplatesPage() { showOk(t("bizMonitorTpl.created")); } closeEdit(); - await refresh(); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1030,7 +1047,8 @@ export function BizMonitorTemplatesPage() { try { await bizMonitorDeleteTemplate(id); showOk(t("bizMonitorTpl.deleted")); - await refresh(); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { @@ -1083,7 +1101,8 @@ export function BizMonitorTemplatesPage() { sheet_overrides: body.sheet_overrides, }); showOk(t("bizMonitorTpl.templateImported")); - await refresh(); + invalidateCutoverCache("bizMigration:"); + await refresh({ force: true }); } catch (e) { showError(formatErr(e)); } finally { diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index f4235a9..d4cc7a1 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -9,6 +9,7 @@ import { useI18n } from "../../i18n"; import { bizStateBulkDeleteBatches, bizStateCollectNow, + bizStateCollectStop, bizStateCreateTask, bizStateDeleteBatch, bizStateDeleteTask, @@ -37,6 +38,7 @@ import type { CliTargetItem } from "../../types"; import { pageCount } from "../../utils/display"; import { writeClipboardText } from "../../utils/clipboard"; import { formatSystemTime } from "../../utils/time"; +import { cutoverCachedGet, cutoverCachedGetSWR, invalidateCutoverCache } from "./cutoverDataCache"; import { jobChipColor, NmStatusChip, sourceChipColor } from "./nmChips"; type TaskRow = { @@ -68,6 +70,8 @@ type Profile = { description: string; metric_id: string; kind?: string; + /** light | heavy — collect dual-lane */ + collect_lane?: string; placeholders?: Placeholder[]; aux_commands?: Array<{ key: string; @@ -305,7 +309,8 @@ export function BizStatePage() { const refreshTasks = useCallback(async () => { const purpose = purposeFilter === "all" ? "" : purposeFilter === "portrait" ? "portrait" : "cutover_hf"; - const res = await bizStateListTasks(purpose); + const key = `bizState:tasks:${purpose || "all"}`; + const res = await cutoverCachedGet(key, () => bizStateListTasks(purpose), { force: true }); const items = (res.items || []) as TaskRow[]; setTasks(items); return items; @@ -333,13 +338,19 @@ export function BizStatePage() { useEffect(() => { void (async () => { try { - await refreshTasks(); + const purpose = + purposeFilter === "all" ? "" : purposeFilter === "portrait" ? "portrait" : "cutover_hf"; + const key = `bizState:tasks:${purpose || "all"}`; + const apply = (res: Awaited>) => { + setTasks((res.items || []) as TaskRow[]); + }; + apply(await cutoverCachedGetSWR(key, () => bizStateListTasks(purpose), apply)); } catch (e) { showError(formatErr(e)); } })(); - // eslint-disable-next-line react-hooks/exhaustive-deps -- refresh when purpose filter / refreshTasks changes - }, [refreshTasks]); + // eslint-disable-next-line react-hooks/exhaustive-deps -- refresh when purpose filter changes + }, [purposeFilter]); // Progress poll while any collect is running (list chips and/or open task). // Does NOT block navigation; cleans up on unmount / when nothing is collecting. @@ -625,6 +636,7 @@ export function BizStatePage() { }); showOk(t("bizState.created")); closeCreate(); + invalidateCutoverCache("bizCompare:"); await refreshTasks(); await openTask(String(task.id), "profiles"); } catch (e) { @@ -815,6 +827,35 @@ export function BizStatePage() { await collectNowForTask(taskId, true); }; + const stopCollectForTask = async (id: string) => { + setBusy(true); + try { + const out = await bizStateCollectStop(id); + if (out.stopped) { + showOk(t("bizState.stopCollectOk")); + } else { + showOk(t("bizState.stopCollectIdle")); + } + setTaskCollecting(id, false); + try { + await refreshTasks(); + } catch { + /* ignore */ + } + if (taskId === id) { + try { + await refreshTaskProgress(id); + } catch { + /* ignore */ + } + } + } catch (e) { + showError(formatErr(e)); + } finally { + setBusy(false); + } + }; + const exportTaskCommands = async () => { if (!taskId) return; setBusy(true); @@ -835,6 +876,7 @@ export function BizStatePage() { await bizStateDeleteTask(id); showOk(t("bizState.deleted")); if (taskId === id) closeTask(); + invalidateCutoverCache("bizCompare:"); await refreshTasks(); } catch (e) { showError(formatErr(e)); @@ -1225,6 +1267,16 @@ export function BizStatePage() { > {t("bizState.collectNow")} + {row.collect_running || collectingIds[row.id] ? ( + + ) : null} + {detail?.collect_running || (taskId && collectingIds[taskId]) ? ( + + ) : null}