mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 04:20:45 +08:00
Improve VRF bind UX with modal and share CLI cache across lanes.
Open a blocking discover dialog for VRF selection, reorder profile columns with scroll/hints, and reuse shared collect cache between light/heavy monitor items. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
3576e4a69c
commit
771f7e21c2
6 changed files with 327 additions and 138 deletions
|
|
@ -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]:
|
||||
|
|
|
|||
|
|
@ -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]:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue