From bf82a01a03763564db0cfd5870f4732e6a350daf Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 30 Jul 2026 00:35:16 +0800 Subject: [PATCH] Stream topology discovery progress and mark missing links stale. SSE updates the UI per NE; edges not seen in a successful scan turn red and can be cleared. Co-authored-by: Cursor --- netx_api/models.py | 2 +- netx_api/topology_router.py | 37 ++++- netx_api/topology_schemas.py | 1 + netx_api/topology_service.py | 154 ++++++++++++++----- tests/test_topology.py | 114 ++++++++++++++ web/src/i18n/en.ts | 17 ++- web/src/i18n/zh.ts | 17 ++- web/src/index.css | 82 ++++++++++ web/src/pages/TopologyPage.tsx | 272 ++++++++++++++++++++++++++++++--- web/src/services/api.ts | 98 ++++++++++++ web/src/types.ts | 21 +++ 11 files changed, 748 insertions(+), 67 deletions(-) diff --git a/netx_api/models.py b/netx_api/models.py index 1aa467e..d57d73a 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -425,7 +425,7 @@ class TopologyEdge(Base): target_node_id: Mapped[str] = mapped_column(String(64), index=True) source_port: Mapped[str] = mapped_column(String(128), default="") target_port: Mapped[str] = mapped_column(String(128), default="") - # manual | lldp | cdp + # manual | lldp | cdp | stale source: Mapped[str] = mapped_column(String(32), default="manual", index=True) discovered_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) diff --git a/netx_api/topology_router.py b/netx_api/topology_router.py index f8e48c4..b7f501e 100644 --- a/netx_api/topology_router.py +++ b/netx_api/topology_router.py @@ -2,9 +2,11 @@ from __future__ import annotations -from typing import Any +import json +from typing import Any, Iterator from fastapi import APIRouter, Depends +from fastapi.responses import StreamingResponse from sqlalchemy.orm import Session from .db import get_db @@ -19,6 +21,7 @@ from .topology_service import ( delete_map, discover_neighbors, get_graph, + iter_discover_neighbors, list_maps, put_graph, update_map, @@ -27,6 +30,11 @@ from .topology_service import ( router = APIRouter(prefix="/v1/topology", tags=["topology"]) +def _sse_pack(event: dict[str, Any]) -> str: + etype = str(event.get("type") or "message") + return f"event: {etype}\ndata: {json.dumps(event, ensure_ascii=False, default=str)}\n\n" + + @router.get("/maps") def api_list_maps(db: Session = Depends(get_db)) -> dict[str, Any]: return list_maps(db) @@ -69,3 +77,30 @@ def api_discover( ) -> dict[str, Any]: req = body or TopologyDiscoverRequest() return discover_neighbors(db, map_id, req).model_dump() + + +@router.post("/maps/{map_id}/discover/stream") +def api_discover_stream( + map_id: str, + body: TopologyDiscoverRequest | None = None, + db: Session = Depends(get_db), +) -> StreamingResponse: + """SSE stream: start → ne_start/ne_result (per NE) → done.""" + req = body or TopologyDiscoverRequest() + + def generate() -> Iterator[str]: + try: + for event in iter_discover_neighbors(db, map_id, req): + yield _sse_pack(event) + except Exception as exc: # noqa: BLE001 — surface to client then end stream + yield _sse_pack({"type": "error", "detail": str(exc)[:800]}) + + return StreamingResponse( + generate(), + media_type="text/event-stream", + headers={ + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", + }, + ) diff --git a/netx_api/topology_schemas.py b/netx_api/topology_schemas.py index 3f2eb9b..8db50fb 100644 --- a/netx_api/topology_schemas.py +++ b/netx_api/topology_schemas.py @@ -108,5 +108,6 @@ class TopologyDiscoverOut(BaseModel): scanned: int = 0 edges_added: int = 0 edges_updated: int = 0 + edges_stale: int = 0 results: list[TopologyDiscoverNeResult] = Field(default_factory=list) graph: TopologyGraphOut | None = None diff --git a/netx_api/topology_service.py b/netx_api/topology_service.py index 2fc4248..daf23be 100644 --- a/netx_api/topology_service.py +++ b/netx_api/topology_service.py @@ -221,7 +221,7 @@ def put_graph(db: Session, map_id: str, body: TopologyGraphPut) -> TopologyGraph if sid not in node_ids or tid not in node_ids: raise HTTPException(status_code=400, detail="edge_endpoint_not_in_nodes") src = str(e.source or "manual").strip().lower() or "manual" - if src not in {"manual", "lldp", "cdp"}: + if src not in {"manual", "lldp", "cdp", "stale"}: raise HTTPException(status_code=400, detail="invalid_edge_source") now = _utcnow() @@ -243,6 +243,7 @@ def put_graph(db: Session, map_id: str, body: TopologyGraphPut) -> TopologyGraph ) ) for e in edges_in: + src = str(e.source or "manual").strip().lower() or "manual" db.add( TopologyEdge( id=str(e.id).strip() or uuid4().hex, @@ -251,8 +252,8 @@ def put_graph(db: Session, map_id: str, body: TopologyGraphPut) -> TopologyGraph target_node_id=str(e.target_node_id).strip(), source_port=str(e.source_port or "").strip()[:128], target_port=str(e.target_port or "").strip()[:128], - source=str(e.source or "manual").strip().lower() or "manual", - discovered_at=now if str(e.source or "").lower() in {"lldp", "cdp"} else None, + source=src, + discovered_at=now if src in {"lldp", "cdp", "stale"} else None, created_at=now, updated_at=now, ) @@ -306,11 +307,12 @@ def _edge_pair_key(a: str, b: str, local_port: str, remote_port: str) -> tuple[s return (b, a, remote_port, local_port) -def discover_neighbors( +def iter_discover_neighbors( db: Session, map_id: str, body: TopologyDiscoverRequest, -) -> TopologyDiscoverOut: +): + """Yield discovery progress events: start / ne_start / ne_result / done / error.""" row = _get_map_or_404(db, map_id) nodes = db.query(TopologyNode).filter(TopologyNode.map_id == row.id).all() edges = db.query(TopologyEdge).filter(TopologyEdge.map_id == row.id).all() @@ -325,7 +327,6 @@ def discover_neighbors( and (not filter_ids or n.managed_ne_id in filter_ids) ] - # Index existing edges for upsert (undirected + ports). existing: dict[tuple[str, str, str, str], TopologyEdge] = {} for e in edges: key = _edge_pair_key( @@ -339,28 +340,55 @@ def discover_neighbors( results: list[TopologyDiscoverNeResult] = [] added = 0 updated = 0 + stale_count = 0 now = _utcnow() proto_req = str(body.protocol or "auto").strip().lower() or "auto" + total = len(scan_nodes) + touched_edge_ids: set[str] = set() + scanned_ok_node_ids: set[str] = set() - for n in scan_nodes: + yield { + "type": "start", + "map_id": row.id, + "protocol": proto_req, + "total": total, + } + + for index, n in enumerate(scan_nodes, start=1): ne = nes.get(n.managed_ne_id) if ne is None: continue + yield { + "type": "ne_start", + "index": index, + "total": total, + "ne_id": ne.id, + "ne_name": ne.name or "", + "ne_ip": ne.ip_address or "", + } cmd, proto_tag = pick_neighbor_command( protocol=proto_req, vendor=ne.vendor or "", device_type=ne.device_type or "", ) if not cmd: - results.append( - TopologyDiscoverNeResult( - ne_id=ne.id, - ne_name=ne.name or "", - ne_ip=ne.ip_address or "", - ok=False, - error="no_command_for_vendor", - ) + result = TopologyDiscoverNeResult( + ne_id=ne.id, + ne_name=ne.name or "", + ne_ip=ne.ip_address or "", + ok=False, + error="no_command_for_vendor", ) + results.append(result) + yield { + "type": "ne_result", + "index": index, + "total": total, + "result": result.model_dump(mode="json"), + "edges_added": added, + "edges_updated": updated, + "edges_stale": stale_count, + } continue exec_out = execute_managed_ne_commands( db, @@ -369,18 +397,27 @@ def discover_neighbors( read_timeout_sec=60, ) if not exec_out.get("ok"): - results.append( - TopologyDiscoverNeResult( - ne_id=ne.id, - ne_name=ne.name or "", - ne_ip=ne.ip_address or "", - ok=False, - command=cmd, - error=str(exec_out.get("detail") or exec_out.get("error") or "exec_failed")[:500], - ) + result = TopologyDiscoverNeResult( + ne_id=ne.id, + ne_name=ne.name or "", + ne_ip=ne.ip_address or "", + ok=False, + command=cmd, + error=str(exec_out.get("detail") or exec_out.get("error") or "exec_failed")[:500], ) + results.append(result) + yield { + "type": "ne_result", + "index": index, + "total": total, + "result": result.model_dump(mode="json"), + "edges_added": added, + "edges_updated": updated, + "edges_stale": stale_count, + } continue + scanned_ok_node_ids.add(n.id) raw = str(exec_out.get("output") or "") hits = parse_neighbor_output( raw, @@ -402,7 +439,6 @@ def discover_neighbors( edge_proto = hit.protocol if hit.protocol in {"lldp", "cdp"} else proto_tag cur = existing.get(key) if cur is not None: - # Never downgrade manual edges; refresh discovery metadata only for discovered. if (cur.source or "manual") == "manual": continue cur.source = edge_proto @@ -410,10 +446,10 @@ def discover_neighbors( cur.target_port = remote_port if cur.source_node_id == n.id else local_port cur.discovered_at = now cur.updated_at = now + touched_edge_ids.add(cur.id) ne_updated += 1 updated += 1 continue - # Prefer orientation: scanning node as source. new_edge = TopologyEdge( id=uuid4().hex, map_id=row.id, @@ -428,32 +464,72 @@ def discover_neighbors( ) db.add(new_edge) existing[key] = new_edge + touched_edge_ids.add(new_edge.id) ne_added += 1 added += 1 - results.append( - TopologyDiscoverNeResult( - ne_id=ne.id, - ne_name=ne.name or "", - ne_ip=ne.ip_address or "", - ok=True, - command=cmd, - neighbors=len(hits), - edges_added=ne_added, - edges_updated=ne_updated, - raw_preview=raw[:800], - ) + result = TopologyDiscoverNeResult( + ne_id=ne.id, + ne_name=ne.name or "", + ne_ip=ne.ip_address or "", + ok=True, + command=cmd, + neighbors=len(hits), + edges_added=ne_added, + edges_updated=ne_updated, + raw_preview=raw[:800], ) + results.append(result) + yield { + "type": "ne_result", + "index": index, + "total": total, + "result": result.model_dump(mode="json"), + "edges_added": added, + "edges_updated": updated, + "edges_stale": stale_count, + } + + # Mark previously discovered edges not refreshed by this run as stale + # (only when at least one endpoint was successfully scanned). + if scanned_ok_node_ids: + for e in edges: + src = (e.source or "manual").strip().lower() + if src not in {"lldp", "cdp", "stale"}: + continue + if e.id in touched_edge_ids: + continue + if e.source_node_id not in scanned_ok_node_ids and e.target_node_id not in scanned_ok_node_ids: + continue + e.source = "stale" + e.updated_at = now + stale_count += 1 row.updated_at = now db.commit() graph = get_graph(db, row.id) - return TopologyDiscoverOut( + report = TopologyDiscoverOut( map_id=row.id, protocol=proto_req, scanned=len(results), edges_added=added, edges_updated=updated, + edges_stale=stale_count, results=results, graph=graph, ) + yield {"type": "done", "report": report.model_dump(mode="json")} + + +def discover_neighbors( + db: Session, + map_id: str, + body: TopologyDiscoverRequest, +) -> TopologyDiscoverOut: + report: TopologyDiscoverOut | None = None + for event in iter_discover_neighbors(db, map_id, body): + if event.get("type") == "done": + report = TopologyDiscoverOut.model_validate(event.get("report") or {}) + if report is None: + raise HTTPException(status_code=500, detail="discover_failed") + return report diff --git a/tests/test_topology.py b/tests/test_topology.py index 7186409..9d61487 100644 --- a/tests/test_topology.py +++ b/tests/test_topology.py @@ -397,6 +397,120 @@ class TopologyServiceTests(unittest.TestCase): self.db.delete(ne_b) self.db.commit() + def test_iter_discover_emits_live_progress_events(self) -> None: + suffix = uuid4().hex[:8] + ne_a_id = f"nea-{suffix}" + ne_b_id = f"neb-{suffix}" + ip_a = f"198.51.100.{(int(suffix[:2], 16) % 100) + 1}" + ip_b = f"198.51.100.{(int(suffix[2:4], 16) % 100) + 101}" + ne_a = ManagedNE( + id=ne_a_id, + name="R2", + vendor="Cisco", + device_type="cisco_ios", + ip_address=ip_a, + ) + ne_b = ManagedNE( + id=ne_b_id, + name="R1", + vendor="Cisco", + device_type="cisco_ios", + ip_address=ip_b, + ) + self.db.add(ne_a) + self.db.add(ne_b) + self.db.commit() + created = svc.create_map(self.db, TopologyMapCreate(name=f"Stream-{suffix}")) + mid = created.id + svc.put_graph( + self.db, + mid, + TopologyGraphPut( + nodes=[ + TopologyNodeIn(id="n1", managed_ne_id=ne_a_id, label="R2", x=0, y=0), + TopologyNodeIn(id="n2", managed_ne_id=ne_b_id, label="R1", x=100, y=0), + ], + edges=[], + ), + ) + fake_exec = { + "ok": True, + "output": CISCO_CDP_DETAIL, + "commands": ["show cdp neighbors detail"], + } + events: list[str] = [] + with patch.object(svc, "execute_managed_ne_commands", return_value=fake_exec): + for ev in svc.iter_discover_neighbors( + self.db, mid, TopologyDiscoverRequest(protocol="cdp", ne_ids=[ne_a_id]) + ): + events.append(str(ev.get("type") or "")) + self.assertEqual(events[0], "start") + self.assertIn("ne_start", events) + self.assertIn("ne_result", events) + self.assertEqual(events[-1], "done") + svc.delete_map(self.db, mid) + self.db.delete(ne_a) + self.db.delete(ne_b) + self.db.commit() + + def test_discover_marks_missing_edges_stale(self) -> None: + suffix = uuid4().hex[:8] + ne_a_id = f"nea-{suffix}" + ne_b_id = f"neb-{suffix}" + ip_a = f"203.0.113.{(int(suffix[:2], 16) % 80) + 10}" + ip_b = f"203.0.113.{(int(suffix[2:4], 16) % 80) + 100}" + ne_a = ManagedNE( + id=ne_a_id, + name="R2", + vendor="Cisco", + device_type="cisco_ios", + ip_address=ip_a, + ) + ne_b = ManagedNE( + id=ne_b_id, + name="R1", + vendor="Cisco", + device_type="cisco_ios", + ip_address=ip_b, + ) + self.db.add(ne_a) + self.db.add(ne_b) + self.db.commit() + mid = svc.create_map(self.db, TopologyMapCreate(name=f"Stale-{suffix}")).id + svc.put_graph( + self.db, + mid, + TopologyGraphPut( + nodes=[ + TopologyNodeIn(id="n1", managed_ne_id=ne_a_id, label="R2", x=0, y=0), + TopologyNodeIn(id="n2", managed_ne_id=ne_b_id, label="R1", x=100, y=0), + ], + edges=[], + ), + ) + with_neighbors = { + "ok": True, + "output": CISCO_CDP_DETAIL, + "commands": ["show cdp neighbors detail"], + } + empty = {"ok": True, "output": "Total entries displayed: 0\n", "commands": ["show cdp neighbors detail"]} + with patch.object(svc, "execute_managed_ne_commands", return_value=with_neighbors): + out1 = svc.discover_neighbors( + self.db, mid, TopologyDiscoverRequest(protocol="cdp", ne_ids=[ne_a_id]) + ) + self.assertGreaterEqual(out1.edges_added, 1) + with patch.object(svc, "execute_managed_ne_commands", return_value=empty): + out2 = svc.discover_neighbors( + self.db, mid, TopologyDiscoverRequest(protocol="cdp", ne_ids=[ne_a_id]) + ) + self.assertGreaterEqual(out2.edges_stale, 1) + graph = svc.get_graph(self.db, mid) + self.assertTrue(any(e.source == "stale" for e in graph.edges)) + svc.delete_map(self.db, mid) + self.db.delete(ne_a) + self.db.delete(ne_b) + self.db.commit() + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index a94b83a..8f63e35 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -611,17 +611,20 @@ const en = { saved: "Topology saved", discover: "Discover links", discovering: "Discovering…", - discovered: "Discovery done: +{{added}} updated {{updated}}", + discovered: "Discovery done: +{{added}} updated {{updated}} missing {{stale}}", discoverFail: "Discovery failed: {{detail}}", selectMap: "Select or create a topology map", - canvasHint: "Add NEs from the left, drag and connect; dashed edges are LLDP/CDP", + canvasHint: "Add NEs from the left; blue dashed=found, red dashed=missing (deletable)", selected: "Selected", + selectedEdge: "Selected edge", openWebcrt: "Open terminal", openNe: "NE details", removeNode: "Remove from canvas", + removeEdge: "Delete edge", noNeLink: "Not linked to a managed NE", edgeManual: "Manual", edgeDiscovered: "Discovered", + edgeStale: "Missing", fit: "Fit view", display: "Display", hideIp: "Hide IP", @@ -637,6 +640,16 @@ const en = { paletteManaged: "Managed", paletteUme: "UME", paletteLoading: "Loading…", + discoverReport: "Discovery progress", + discoverClose: "Close", + discoverProgress: "Scanning {{count}} managed NE(s) on the map (serial LLDP; please wait)…", + discoverProgressLive: "Scanning {{index}}/{{total}}: {{name}}", + discoverSummary: "Scanned {{scanned}} · added {{added}} · updated {{updated}} · missing {{stale}} · failed {{failed}}", + discoverNeOk: "neighbors {{neighbors}} · added {{added}} · updated {{updated}}", + discoverNeFail: "Collect failed", + removeStale: "Clear missing ({{count}})", + removeStaleHint: "Delete red-marked missing edges", + staleRemoved: "Removed {{count}} missing edge(s)", }, }; diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 455fc83..f2f7389 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -609,17 +609,20 @@ const zh = { saved: "拓扑已保存", discover: "发现链路", discovering: "发现中…", - discovered: "发现完成:新增 {{added}},更新 {{updated}}", + discovered: "发现完成:新增 {{added}},更新 {{updated}},未发现 {{stale}}", discoverFail: "发现失败:{{detail}}", selectMap: "请选择或新建一张拓扑图", - canvasHint: "从左侧添加网元,拖动节点并连线;虚线为 LLDP/CDP 发现链路", + canvasHint: "从左侧添加网元,拖动节点并连线;蓝虚线=发现链路,红虚线=本次未发现(可删)", selected: "已选节点", + selectedEdge: "已选链路", openWebcrt: "打开终端", openNe: "网元详情", removeNode: "从画布移除", + removeEdge: "删除链路", noNeLink: "未关联托管网元", edgeManual: "人工", edgeDiscovered: "发现", + edgeStale: "未发现", fit: "适应画布", display: "显示选项", hideIp: "隐藏 IP", @@ -635,6 +638,16 @@ const zh = { paletteManaged: "托管网元", paletteUme: "UME 网元", paletteLoading: "加载中…", + discoverReport: "发现进度", + discoverClose: "关闭", + discoverProgress: "正在扫描画布上 {{count}} 台托管网元(串行登录采集 LLDP,请稍候)…", + discoverProgressLive: "正在扫描 {{index}}/{{total}}:{{name}}", + discoverSummary: "扫描 {{scanned}} 台 · 新增 {{added}} · 更新 {{updated}} · 未发现 {{stale}} · 失败 {{failed}}", + discoverNeOk: "邻居 {{neighbors}} · 新增 {{added}} · 更新 {{updated}}", + discoverNeFail: "采集失败", + removeStale: "清除未发现 ({{count}})", + removeStaleHint: "删除红色标记的未发现链路", + staleRemoved: "已删除 {{count}} 条未发现链路", }, }; diff --git a/web/src/index.css b/web/src/index.css index 1927806..9e7697c 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -2089,6 +2089,88 @@ pre { gap: 8px; } +.topo-discover { + display: flex; + flex-direction: column; + gap: 8px; + max-height: 220px; + padding: 10px 12px; + background: #f8fafc; + border-bottom: 1px solid #e2e8f0; + overflow: hidden; +} + +.topo-discover__head { + display: flex; + align-items: center; + justify-content: space-between; + gap: 8px; +} + +.topo-discover__summary { + margin: 0; + font-size: 13px; + color: #334155; +} + +.topo-discover__error { + margin: 0; + font-size: 13px; + color: #b91c1c; +} + +.topo-discover__list { + list-style: none; + margin: 0; + padding: 0; + overflow: auto; + flex: 1; + min-height: 0; +} + +.topo-discover__list li { + padding: 8px; + border-radius: 8px; + margin-bottom: 4px; + background: #fff; + border: 1px solid #e2e8f0; +} + +.topo-discover__list li.is-fail { + border-color: #fecaca; + background: #fff1f2; +} + +.topo-discover__item-title { + display: flex; + align-items: center; + justify-content: space-between; + gap: 8px; + font-weight: 600; + color: #0f172a; +} + +.topo-discover__badge { + font-size: 11px; + font-weight: 700; + color: #166534; +} + +.topo-discover__list li.is-fail .topo-discover__badge { + color: #b91c1c; +} + +.topo-discover__item-meta, +.topo-discover__item-cmd { + margin-top: 2px; + font-size: 12px; + color: #64748b; +} + +.topo-discover__item-cmd { + font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; +} + .topo-canvas { position: relative; flex: 1; diff --git a/web/src/pages/TopologyPage.tsx b/web/src/pages/TopologyPage.tsx index fb1b84e..9e7f62f 100644 --- a/web/src/pages/TopologyPage.tsx +++ b/web/src/pages/TopologyPage.tsx @@ -22,7 +22,7 @@ import "@xyflow/react/dist/style.css"; import { createTopologyMap, deleteTopologyMap, - discoverTopologyNeighbors, + discoverTopologyNeighborsStream, fetchManagedNe, fetchTopologyGraph, fetchTopologyMaps, @@ -34,7 +34,14 @@ import { queryKeys } from "../constants/queryKeys"; import { useI18n } from "../i18n"; import { useToast } from "../hooks/useToast"; import { openNewModuleWindow, openOrFocusModule } from "../utils/moduleWindows"; -import type { ManagedNeItem, TopologyEdgeItem, TopologyNodeItem, UmeNeItem } from "../types"; +import type { + ManagedNeItem, + TopologyDiscoverNeResult, + TopologyDiscoverOut, + TopologyEdgeItem, + TopologyNodeItem, + UmeNeItem, +} from "../types"; type NeNodeData = { label: string; @@ -140,6 +147,17 @@ function NeNode({ data, selected }: NodeProps>) { const nodeTypes = { neNode: NeNode }; +function edgeStyle(source: string): { stroke: string; strokeDasharray?: string; strokeWidth?: number } { + const src = (source || "manual").toLowerCase(); + if (src === "stale") { + return { stroke: "#dc2626", strokeDasharray: "4 4", strokeWidth: 2 }; + } + if (src === "lldp" || src === "cdp") { + return { stroke: "#0ea5e9", strokeDasharray: "6 4" }; + } + return { stroke: "#64748b" }; +} + function graphToFlow(nodes: TopologyNodeItem[], edges: TopologyEdgeItem[]) { const rfNodes: Node[] = nodes.map((n) => ({ id: n.id, @@ -155,7 +173,7 @@ function graphToFlow(nodes: TopologyNodeItem[], edges: TopologyEdgeItem[]) { }, })); const rfEdges: Edge[] = edges.map((e) => { - const discovered = e.source === "lldp" || e.source === "cdp"; + const src = e.source || "manual"; const label = [e.source_port, e.target_port].filter(Boolean).join(" ↔ "); return { id: e.id, @@ -164,12 +182,10 @@ function graphToFlow(nodes: TopologyNodeItem[], edges: TopologyEdgeItem[]) { type: "straight", label: label || undefined, animated: false, - style: discovered - ? { stroke: "#0ea5e9", strokeDasharray: "6 4" } - : { stroke: "#64748b" }, + style: edgeStyle(src), markerEnd: { type: MarkerType.ArrowClosed, width: 16, height: 16 }, data: { - source: e.source || "manual", + source: src, source_port: e.source_port || "", target_port: e.target_port || "", }, @@ -206,6 +222,7 @@ export function TopologyPage() { const [mapId, setMapId] = useState(""); const [keyword, setKeyword] = useState(""); const [selectedNodeId, setSelectedNodeId] = useState(null); + const [selectedEdgeId, setSelectedEdgeId] = useState(null); const [hideIp, setHideIp] = useState(true); const [hideVendor, setHideVendor] = useState(true); const [hidePorts, setHidePorts] = useState(true); @@ -213,6 +230,20 @@ export function TopologyPage() { const [sidebarCollapsed, setSidebarCollapsed] = useState(false); const [hideAddedNes, setHideAddedNes] = useState(true); const [paletteSource, setPaletteSource] = useState("managed"); + const [discoverOpen, setDiscoverOpen] = useState(false); + const [discovering, setDiscovering] = useState(false); + const [discoverReport, setDiscoverReport] = useState(null); + const [discoverLiveResults, setDiscoverLiveResults] = useState([]); + const [discoverProgress, setDiscoverProgress] = useState({ + index: 0, + total: 0, + neName: "", + neIp: "", + edgesAdded: 0, + edgesUpdated: 0, + }); + const [discoverError, setDiscoverError] = useState(""); + const discoverAbortRef = useRef(null); const [nodes, setNodes, onNodesChange] = useNodesState>([]); const [edges, setEdges, onEdgesChange] = useEdgesState([]); const rfRef = useRef, Edge> | null>(null); @@ -338,14 +369,60 @@ export function TopologyPage() { onError: (err) => showError(String(err)), }); - const discoverMut = useMutation({ - mutationFn: async () => { + const runDiscover = useCallback(async () => { + if (!mapId || discovering) return; + discoverAbortRef.current?.abort(); + const ac = new AbortController(); + discoverAbortRef.current = ac; + setDiscoverOpen(true); + setDiscovering(true); + setDiscoverReport(null); + setDiscoverLiveResults([]); + setDiscoverError(""); + setDiscoverProgress({ + index: 0, + total: nodes.filter((n) => Boolean(n.data.managed_ne_id)).length, + neName: "", + neIp: "", + edgesAdded: 0, + edgesUpdated: 0, + }); + try { if (dirtyRef.current) { await putTopologyGraph(mapId, flowToGraphPayload(nodes, edges)); + dirtyRef.current = false; } - return discoverTopologyNeighbors(mapId, { protocol: "auto" }); - }, - onSuccess: async (out) => { + const out = await discoverTopologyNeighborsStream( + mapId, + { protocol: "auto" }, + { + onStart: (ev) => { + setDiscoverProgress((p) => ({ ...p, total: ev.total, index: 0 })); + }, + onNeStart: (ev) => { + setDiscoverProgress((p) => ({ + ...p, + index: ev.index, + total: ev.total, + neName: ev.ne_name || ev.ne_ip || ev.ne_id, + neIp: ev.ne_ip, + })); + }, + onNeResult: (ev) => { + if (ev.result) { + setDiscoverLiveResults((prev) => [...prev, ev.result]); + } + setDiscoverProgress((p) => ({ + ...p, + index: ev.index, + total: ev.total, + edgesAdded: ev.edges_added, + edgesUpdated: ev.edges_updated, + })); + }, + }, + ac.signal, + ); if (out.graph) { queryClient.setQueryData(queryKeys.topologyGraph(mapId), out.graph); const { rfNodes, rfEdges } = graphToFlow(out.graph.nodes, out.graph.edges); @@ -353,16 +430,42 @@ export function TopologyPage() { setEdges(rfEdges); dirtyRef.current = false; } + setDiscoverReport(out); await queryClient.invalidateQueries({ queryKey: queryKeys.topologyMaps }); showOk( t("topology.discovered") .replace("{{added}}", String(out.edges_added)) - .replace("{{updated}}", String(out.edges_updated)), + .replace("{{updated}}", String(out.edges_updated)) + .replace("{{stale}}", String(out.edges_stale || 0)), ); - }, - onError: (err) => - showError(t("topology.discoverFail").replace("{{detail}}", String(err))), - }); + } catch (err) { + if (ac.signal.aborted) return; + setDiscoverError(String(err)); + showError(t("topology.discoverFail").replace("{{detail}}", String(err))); + } finally { + if (discoverAbortRef.current === ac) discoverAbortRef.current = null; + setDiscovering(false); + } + }, [mapId, discovering, nodes, edges, queryClient, setNodes, setEdges, showOk, showError, t]); + + const discoverResults = discoverReport?.results?.length + ? discoverReport.results + : discoverLiveResults; + const discoverSummary = discoverReport + ? { + scanned: discoverReport.scanned, + added: discoverReport.edges_added, + updated: discoverReport.edges_updated, + stale: discoverReport.edges_stale || 0, + failed: discoverReport.results.filter((r) => !r.ok).length, + } + : { + scanned: discoverLiveResults.length, + added: discoverProgress.edgesAdded, + updated: discoverProgress.edgesUpdated, + stale: 0, + failed: discoverLiveResults.filter((r) => !r.ok).length, + }; const onConnect = useCallback( (connection: Connection) => { @@ -455,8 +558,24 @@ export function TopologyPage() { () => nodes.find((n) => n.id === selectedNodeId) || null, [nodes, selectedNodeId], ); + const selectedEdge = useMemo( + () => edges.find((e) => e.id === selectedEdgeId) || null, + [edges, selectedEdgeId], + ); + const staleEdgeCount = useMemo( + () => + edges.filter((e) => String((e.data as { source?: string } | undefined)?.source || "") === "stale") + .length, + [edges], + ); const removeSelected = () => { + if (selectedEdgeId) { + dirtyRef.current = true; + setEdges((es) => es.filter((e) => e.id !== selectedEdgeId)); + setSelectedEdgeId(null); + return; + } if (!selectedNodeId) return; dirtyRef.current = true; setNodes((ns) => ns.filter((n) => n.id !== selectedNodeId)); @@ -464,6 +583,17 @@ export function TopologyPage() { setSelectedNodeId(null); }; + const removeStaleEdges = () => { + const n = staleEdgeCount; + if (n <= 0) return; + dirtyRef.current = true; + setEdges((es) => + es.filter((e) => String((e.data as { source?: string } | undefined)?.source || "") !== "stale"), + ); + setSelectedEdgeId(null); + showOk(t("topology.staleRemoved").replace("{{count}}", String(n))); + }; + const openWebcrt = () => { const managedId = selectedNode?.data.managed_ne_id; const umeId = selectedNode?.data.ume_ne_id; @@ -790,10 +920,19 @@ export function TopologyPage() { + + + + ) : null} + + {discoverOpen ? ( +
+
+ {t("topology.discoverReport")} + +
+ {discovering ? ( +

+ {t("topology.discoverProgressLive") + .replace("{{index}}", String(discoverProgress.index)) + .replace("{{total}}", String(discoverProgress.total)) + .replace( + "{{name}}", + discoverProgress.neName || discoverProgress.neIp || "…", + )} +

+ ) : null} + {discoverError ?

{discoverError}

: null} + {discoverResults.length > 0 || discoverReport ? ( + <> +

+ {t("topology.discoverSummary") + .replace("{{scanned}}", String(discoverSummary.scanned)) + .replace("{{added}}", String(discoverSummary.added)) + .replace("{{updated}}", String(discoverSummary.updated)) + .replace("{{stale}}", String(discoverSummary.stale)) + .replace("{{failed}}", String(discoverSummary.failed))} +

+
    + {discoverResults.map((r) => ( +
  • +
    + {r.ne_name || r.ne_ip || r.ne_id} + {r.ok ? "OK" : "FAIL"} +
    +
    + {r.ne_ip ? `${r.ne_ip} · ` : ""} + {r.ok + ? t("topology.discoverNeOk") + .replace("{{neighbors}}", String(r.neighbors)) + .replace("{{added}}", String(r.edges_added)) + .replace("{{updated}}", String(r.edges_updated)) + : r.error || t("topology.discoverNeFail")} +
    + {r.command ? ( +
    {r.command}
    + ) : null} +
  • + ))} +
+ + ) : null} +
) : null}
@@ -844,13 +1060,25 @@ export function TopologyPage() { onEdgesChange(changes); }} onConnect={onConnect} - onNodeClick={(_e, node) => setSelectedNodeId(node.id)} - onPaneClick={() => setSelectedNodeId(null)} + onNodeClick={(_e, node) => { + setSelectedEdgeId(null); + setSelectedNodeId(node.id); + }} + onEdgeClick={(_e, edge) => { + setSelectedNodeId(null); + setSelectedEdgeId(edge.id); + }} + onPaneClick={() => { + setSelectedNodeId(null); + setSelectedEdgeId(null); + }} onInit={(inst) => { rfRef.current = inst; }} fitView deleteKeyCode={["Backspace", "Delete"]} + edgesFocusable + elementsSelectable > diff --git a/web/src/services/api.ts b/web/src/services/api.ts index aaabcd6..703f3cf 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -19,6 +19,8 @@ import type { CliTargetListResponse, UmeCliOverrideItem, TopologyDiscoverOut, + TopologyDiscoverNeResult, + TopologyDiscoverStreamHandlers, TopologyGraph, TopologyMapItem, } from "../types"; @@ -473,3 +475,99 @@ export const discoverTopologyNeighbors = ( `/v1/topology/maps/${encodeURIComponent(mapId)}/discover`, body || {}, ); + +function parseSseChunks(buffer: string): { events: Array<{ event: string; data: string }>; rest: string } { + const events: Array<{ event: string; data: string }> = []; + let rest = buffer; + while (true) { + const sep = rest.indexOf("\n\n"); + if (sep < 0) break; + const raw = rest.slice(0, sep); + rest = rest.slice(sep + 2); + let event = "message"; + const dataLines: string[] = []; + for (const line of raw.split("\n")) { + if (line.startsWith("event:")) event = line.slice(6).trim(); + else if (line.startsWith("data:")) dataLines.push(line.slice(5).trim()); + } + if (dataLines.length) events.push({ event, data: dataLines.join("\n") }); + } + return { events, rest }; +} + +export async function discoverTopologyNeighborsStream( + mapId: string, + body: { protocol?: string; ne_ids?: string[] } | undefined, + handlers: TopologyDiscoverStreamHandlers, + signal?: AbortSignal, +): Promise { + const res = await fetch(`/v1/topology/maps/${encodeURIComponent(mapId)}/discover/stream`, { + method: "POST", + headers: { + accept: "text/event-stream", + "content-type": "application/json", + }, + body: JSON.stringify(body || {}), + signal, + }); + if (!res.ok) { + const text = await res.text().catch(() => ""); + throw new Error(text || `discover_stream_failed:${res.status}`); + } + if (!res.body) throw new Error("discover_stream_empty_body"); + + const reader = res.body.getReader(); + const decoder = new TextDecoder(); + let buf = ""; + let finalReport: TopologyDiscoverOut | null = null; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + buf += decoder.decode(value, { stream: true }); + const parsed = parseSseChunks(buf); + buf = parsed.rest; + for (const item of parsed.events) { + let payload: Record = {}; + try { + payload = JSON.parse(item.data) as Record; + } catch { + continue; + } + const type = String(payload.type || item.event || ""); + if (type === "start") { + handlers.onStart?.({ + map_id: String(payload.map_id || ""), + protocol: String(payload.protocol || ""), + total: Number(payload.total || 0), + }); + } else if (type === "ne_start") { + handlers.onNeStart?.({ + index: Number(payload.index || 0), + total: Number(payload.total || 0), + ne_id: String(payload.ne_id || ""), + ne_name: String(payload.ne_name || ""), + ne_ip: String(payload.ne_ip || ""), + }); + } else if (type === "ne_result") { + handlers.onNeResult?.({ + index: Number(payload.index || 0), + total: Number(payload.total || 0), + result: payload.result as TopologyDiscoverNeResult, + edges_added: Number(payload.edges_added || 0), + edges_updated: Number(payload.edges_updated || 0), + }); + } else if (type === "done") { + finalReport = payload.report as TopologyDiscoverOut; + if (finalReport) handlers.onDone?.(finalReport); + } else if (type === "error") { + const detail = String(payload.detail || "discover_stream_error"); + handlers.onError?.(detail); + throw new Error(detail); + } + } + } + + if (!finalReport) throw new Error("discover_stream_incomplete"); + return finalReport; +} diff --git a/web/src/types.ts b/web/src/types.ts index c6f3120..88ea69e 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -405,6 +405,27 @@ export type TopologyDiscoverOut = { scanned: number; edges_added: number; edges_updated: number; + edges_stale?: number; results: TopologyDiscoverNeResult[]; graph: TopologyGraph | null; }; + +export type TopologyDiscoverStreamHandlers = { + onStart?: (ev: { map_id: string; protocol: string; total: number }) => void; + onNeStart?: (ev: { + index: number; + total: number; + ne_id: string; + ne_name: string; + ne_ip: string; + }) => void; + onNeResult?: (ev: { + index: number; + total: number; + result: TopologyDiscoverNeResult; + edges_added: number; + edges_updated: number; + }) => void; + onDone?: (report: TopologyDiscoverOut) => void; + onError?: (detail: string) => void; +};