diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index 211633b..a3ba046 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -367,6 +367,10 @@ def _run_collect_lane( per_cmd: int, cap: int, label: str, + shared_cache: dict[str, Any] | None = None, + cache_lock: Any | None = None, + cmd_locks: dict[str, Any] | None = None, + aux_persisted: set[tuple[str, str]] | None = None, ) -> tuple[int, int, bool, bool]: """Run one SSH lane (own connection + CollectSession + timeout budget).""" if not work: @@ -405,6 +409,9 @@ def _run_collect_lane( device_type=device_type_eff, vendor_key=vendor_key, read_timeout=per_cmd, + shared_cache=shared_cache, + cache_lock=cache_lock, + cmd_locks=cmd_locks, ) try: batch_row = sdb.get(BizStateBatch, batch_id) @@ -413,7 +420,7 @@ def _run_collect_lane( # Resolve expand_all → concrete per-VRF commands via discover profile. flat_work: list[WorkItem] = [] - aux_persisted: set[tuple[str, str]] = set() + persisted = aux_persisted if aux_persisted is not None else set() for concrete, params, profile_id, item_id, mode in work: if mode != "expand_all": flat_work.append((concrete, params, profile_id, item_id, mode)) @@ -627,7 +634,16 @@ def _run_collect_lane( and aux_mid in _GENERIC_METRICS ): persist_key = (normalize_command(ra.command), aux_mid) - if persist_key not in aux_persisted: + do_persist = False + if cache_lock is not None: + with cache_lock: + if persist_key not in persisted: + persisted.add(persist_key) + do_persist = True + elif persist_key not in persisted: + persisted.add(persist_key) + do_persist = True + if do_persist: n_aux = _persist_metric_rows( sdb, batch=batch_row, @@ -637,7 +653,6 @@ def _run_collect_lane( ) aux_row.row_count = n_aux total_rows += n_aux - aux_persisted.add(persist_key) if n_aux: _bump_batch_progress(batch_id, add_rows=n_aux) sdb.add(aux_row) @@ -900,12 +915,20 @@ def _run_collect_session( raise RuntimeError("no commands to run") light_work, heavy_work = partition_work(work) + shared_cache: dict[str, Any] = {} + cache_lock = threading.RLock() + cmd_locks: dict[str, Any] = {} + aux_persisted: set[tuple[str, str]] = set() lane_kwargs = dict( batch_id=batch_id, creds=creds, vendor_eff=vendor_eff, device_type_eff=device_type_eff, vendor_key=vendor_key, + shared_cache=shared_cache, + cache_lock=cache_lock, + cmd_locks=cmd_locks, + aux_persisted=aux_persisted, ) def _run_light() -> tuple[int, int, bool, bool]: diff --git a/netx_api/biz_state/collect_session.py b/netx_api/biz_state/collect_session.py index 4fe88ed..5d7da43 100644 --- a/netx_api/biz_state/collect_session.py +++ b/netx_api/biz_state/collect_session.py @@ -2,9 +2,11 @@ from __future__ import annotations +from contextlib import nullcontext from dataclasses import dataclass, field from typing import Any, Callable import re +import threading from .command_match import normalize_command from .enrich import EnrichJoin, apply_enrich_joins @@ -88,7 +90,11 @@ class ParseBundle: class CollectSession: - """SSH session-scoped command cache + aux fetch/parse.""" + """SSH session-scoped command cache + aux fetch/parse. + + Optional ``shared_cache`` / ``cache_lock`` / ``cmd_locks`` let light+heavy + lanes reuse the same CLI results within one batch (config_vrf / FIB aux). + """ def __init__( self, @@ -99,6 +105,9 @@ class CollectSession: vendor_key: str = "", read_timeout: int = 120, send_fn: SendFn | None = None, + shared_cache: dict[str, CachedCommand] | None = None, + cache_lock: threading.RLock | None = None, + cmd_locks: dict[str, threading.Lock] | None = None, ) -> None: self.conn = conn self.vendor = vendor @@ -106,7 +115,25 @@ class CollectSession: self.vendor_key = vendor_key self.read_timeout = int(read_timeout or 120) self._send = send_fn - self.cache: dict[str, CachedCommand] = {} + self.cache: dict[str, CachedCommand] = ( + shared_cache if shared_cache is not None else {} + ) + self._cache_lock = cache_lock + self._cmd_locks = cmd_locks if cmd_locks is not None else {} + + def _meta_lock(self): + return self._cache_lock if self._cache_lock is not None else nullcontext() + + def _command_lock(self, command: str): + ck = normalize_command(command) + if self._cache_lock is None: + return nullcontext() + with self._cache_lock: + lock = self._cmd_locks.get(ck) + if lock is None: + lock = threading.Lock() + self._cmd_locks[ck] = lock + return lock def _send_show(self, command: str) -> str: if self._send is None: @@ -135,14 +162,16 @@ class CollectSession: error=str(error or ""), cmd_row_id=str(cmd_row_id or ""), ) - self.cache[ck] = entry + with self._meta_lock(): + self.cache[ck] = entry return entry def get_cached(self, command: str) -> CachedCommand | None: ck = normalize_command(command) - hit = self.cache.get(ck) - if hit and hit.ok: - return hit + with self._meta_lock(): + hit = self.cache.get(ck) + if hit and hit.ok: + return hit return None def fetch_and_parse( @@ -154,52 +183,57 @@ class CollectSession: params: dict[str, str] | None = None, cmd_row_id: str = "", ) -> tuple[CachedCommand, bool]: - """Return ``(entry, cache_hit)``. On miss: CLI + optional parser.""" - cached = self.get_cached(command) - if cached is not None: - return cached, True - try: - raw = self._send_show(command) - except Exception as exc: - entry = self.remember( - command, - raw="", - ok=False, - error=f"{type(exc).__name__}: {exc}", - cmd_row_id=cmd_row_id, - ) - return entry, False - records: list[dict[str, Any]] = [] - fsm_tables: dict[str, list[dict[str, Any]]] = {} - if parser_id and get_parser(parser_id): + """Return ``(entry, cache_hit)``. On miss: CLI + optional parser. + + Same concrete CLI is serialized across shared-cache lanes so aux of + one monitor item can be reused by the next without re-collecting. + """ + with self._command_lock(command): + cached = self.get_cached(command) + if cached is not None: + return cached, True try: - records, fsm_tables, _keys = run_parser( - parser_id, - raw_text=raw, - vendor=self.vendor, - device_type=self.device_type, - command=textfsm_command or command, - textfsm_command=textfsm_command or "", - params=params or {}, - ) + raw = self._send_show(command) except Exception as exc: entry = self.remember( command, - raw=raw, + raw="", ok=False, - error=f"parse: {type(exc).__name__}: {exc}", + error=f"{type(exc).__name__}: {exc}", cmd_row_id=cmd_row_id, ) return entry, False - entry = self.remember( - command, - raw=raw, - fsm_tables=fsm_tables, - records=records, - ok=True, - cmd_row_id=cmd_row_id, - ) - return entry, False + records: list[dict[str, Any]] = [] + fsm_tables: dict[str, list[dict[str, Any]]] = {} + if parser_id and get_parser(parser_id): + try: + records, fsm_tables, _keys = run_parser( + parser_id, + raw_text=raw, + vendor=self.vendor, + device_type=self.device_type, + command=textfsm_command or command, + textfsm_command=textfsm_command or "", + params=params or {}, + ) + except Exception as exc: + entry = self.remember( + command, + raw=raw, + ok=False, + error=f"parse: {type(exc).__name__}: {exc}", + cmd_row_id=cmd_row_id, + ) + return entry, False + entry = self.remember( + command, + raw=raw, + fsm_tables=fsm_tables, + records=records, + ok=True, + cmd_row_id=cmd_row_id, + ) + return entry, False def primary_rule_keys(parser_id: str) -> list[str]: diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index ee144cc..956b19b 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -238,6 +238,12 @@ const en = { bindTitle: "Select VRF bindings", unbound: "Not bound", allVrfsDefault: "All VRFs (default)", + bindHintRequired: "Bind VRF params before collect", + bindHintOptional: "Bind VRFs, or leave empty for all", + discoverLoading: "Discovering VRFs…", + discoverEmpty: "No VRFs discovered", + selectAllVrfs: "Select all", + clearVrfs: "Clear", batches: "Batches", viewBatch: "Open", export: "Export", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 0963198..2fe3ac3 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -238,6 +238,12 @@ const zh = { bindTitle: "选择 VRF 绑定", unbound: "未关联", allVrfsDefault: "默认全部 VRF", + bindHintRequired: "需关联 VRF 参数后才可采集", + bindHintOptional: "可关联指定 VRF,未关联则采集全部", + discoverLoading: "正在发现 VRF…", + discoverEmpty: "未发现可用 VRF", + selectAllVrfs: "全选", + clearVrfs: "清空", batches: "采集批次", viewBatch: "查看", export: "导出", diff --git a/web/src/index.css b/web/src/index.css index 8efa9e7..9a1a9eb 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -11261,11 +11261,52 @@ html.login-page--paused .login-page__flare { overflow: auto; } +.bs-bind-list--modal { + max-height: min(48vh, 360px); + padding: 4px 2px; + border: 1px solid rgba(148, 163, 184, 0.18); + border-radius: 8px; +} + .bs-bind-item { display: block; margin-bottom: 4px; } +.bs-profiles-table-wrap { + overflow-x: auto !important; + overflow-y: auto; + max-height: min(58vh, 560px); +} + +.bs-profiles-table { + min-width: 960px; +} + +.bs-profiles-table .bs-cmd-cell { + white-space: nowrap; + font-size: 12px; +} + +.bs-profiles-table .bs-params-cell { + min-width: 140px; + max-width: 280px; + word-break: break-word; +} + +.bs-bind-hint { + margin-top: 4px; + font-size: 12px; +} + +.bs-params-warn { + color: #fbbf24; +} + +.bs-bind-modal.app-heroui-modal { + z-index: 80; +} + /* biz_state compare result workbook / board — dark surfaces (match nm modal) */ .bs-cmp-job-body { min-height: min(64vh, 680px); diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index ddd819f..21a0e2c 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -234,11 +234,13 @@ export function BizStatePage() { const [dailyKeepCount, setDailyKeepCount] = useState(10); const [selectedBatchIds, setSelectedBatchIds] = useState([]); - // VRF bind (inside task modal) + // VRF bind modal (blocking) const [bindItemId, setBindItemId] = useState(""); const [candidates, setCandidates] = useState([]); const [selectedVrfs, setSelectedVrfs] = useState([]); const [discoverCmd, setDiscoverCmd] = useState(""); + const [discoverLoading, setDiscoverLoading] = useState(false); + const [discoverError, setDiscoverError] = useState(""); // batch workbook modal (summary + lazy-paged metric sheets) const [batchDetail, setBatchDetail] = useState(null); @@ -756,6 +758,15 @@ export function BizStatePage() { } }; + const closeBindModal = () => { + setBindItemId(""); + setCandidates([]); + setSelectedVrfs([]); + setDiscoverCmd(""); + setDiscoverError(""); + setDiscoverLoading(false); + }; + const startDiscover = async (item: any) => { if (!taskId) return; const prof = profiles.find((p) => p.profile_id === item.source_profile_id); @@ -764,8 +775,13 @@ export function BizStatePage() { showError(t("bizState.noNeedBind")); return; } - setBusy(true); setBindItemId(item.id); + setCandidates([]); + setSelectedVrfs([]); + setDiscoverCmd(""); + setDiscoverError(""); + setDiscoverLoading(true); + setBusy(true); try { const res = await bizStateDiscover({ task_id: taskId, @@ -773,8 +789,9 @@ export function BizStatePage() { placeholder: ph.name, }); if (!res.ok) { - showError(res.error || t("bizState.discoverFailed")); - setCandidates([]); + const err = res.error || t("bizState.discoverFailed"); + setDiscoverError(err); + showError(err); return; } setDiscoverCmd(res.command || ""); @@ -784,9 +801,15 @@ export function BizStatePage() { .filter((b: any) => b.placeholder === ph.name) .map((b: any) => String(b.value)); setSelectedVrfs(existing.length ? existing : cand.map((c) => c.value)); + if (!cand.length) { + setDiscoverError(t("bizState.discoverEmpty")); + } } catch (e) { - showError(formatErr(e)); + const err = formatErr(e); + setDiscoverError(err); + showError(err); } finally { + setDiscoverLoading(false); setBusy(false); } }; @@ -808,8 +831,7 @@ export function BizStatePage() { ); showOk(t("bizState.bindingsSaved")); await loadTask(taskId); - setBindItemId(""); - setCandidates([]); + closeBindModal(); } catch (e) { showError(formatErr(e)); } finally { @@ -1311,14 +1333,14 @@ export function BizStatePage() { {taskTab === "profiles" ? ( <> -
- +
+
- + @@ -1333,6 +1355,13 @@ export function BizStatePage() { const bindOptional = (prof.placeholders || []).every( (ph) => ph.required === false, ); + const bindHint = needsBind + ? binds.length + ? binds.map((b: any) => b.value).join(", ") + : bindOptional + ? t("bizState.allVrfsDefault") + : t("bizState.unbound") + : "—"; return ( + -
{t("bizState.enable")} {t("bizState.profiles")}{t("bizState.command")} {t("bizState.params")}{t("bizState.command")}
@@ -1346,25 +1375,38 @@ export function BizStatePage() {
{prof.title}
{prof.description ?
{prof.description}
: null} + {needsBind ? ( +
+ {bindOptional + ? t("bizState.bindHintOptional") + : t("bizState.bindHintRequired")} +
+ ) : null} +
+ {needsBind ? ( + + {bindHint} + + ) : ( + "—" + )} - {prof.command_template} - - {needsBind - ? binds.length - ? binds.map((b: any) => b.value).join(", ") - : bindOptional - ? t("bizState.allVrfsDefault") - : t("bizState.unbound") - : "—"} + {prof.command_template} {needsBind && enabled && it ? (
- - {bindItemId && candidates.length ? ( -
-

{t("bizState.bindTitle")}

-

- {discoverCmd} · {selectedVrfs.length} -

-
- {candidates.map((c) => ( - - ))} -
-
- - - -
-
- ) : null} ) : (
@@ -1575,6 +1547,113 @@ export function BizStatePage() { + {/* VRF discover / bind — blocking modal above task dialog */} + + + {t("bizState.bindTitle")} + + + + {discoverLoading ? ( +

{t("bizState.discoverLoading")}

+ ) : null} + {discoverCmd ? ( +

+ {discoverCmd} + {candidates.length ? ` · ${selectedVrfs.length}/${candidates.length}` : null} +

+ ) : null} + {discoverError ?

{discoverError}

: null} + {!discoverLoading && candidates.length ? ( + <> +
+ + +
+
+ {candidates.map((c) => ( + + ))} +
+ + ) : null} +
+ + {(() => { + const item = (detail?.items || []).find((it: any) => it.id === bindItemId); + const prof = profiles.find((p) => p.profile_id === item?.source_profile_id); + const bindOptional = (prof?.placeholders || []).every((ph) => ph.required === false); + return ( + <> + + {bindOptional ? ( + + ) : null} + + + ); + })()} + +
+ {/* Batch workbook: summary + lazy-paged metric sheets */}