diff --git a/gns3server/api/routes/compute/docker_nodes.py b/gns3server/api/routes/compute/docker_nodes.py index 3e30628a2..f8e8b8bb6 100644 --- a/gns3server/api/routes/compute/docker_nodes.py +++ b/gns3server/api/routes/compute/docker_nodes.py @@ -29,6 +29,7 @@ from typing import Union from gns3server import schemas from gns3server.compute.docker import Docker from gns3server.compute.docker.docker_vm import DockerVM +from gns3server.compute.marker.marker_manager import MarkerManager from .dependencies.authentication import compute_authentication, ws_compute_authentication @@ -353,6 +354,52 @@ async def stop_docker_node_capture( await node.stop_capture(adapter_number) +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/start", + dependencies=[Depends(compute_authentication)] +) +async def start_docker_node_marker( + *, + project_id: UUID, + adapter_number: int, + port_number: int, + marker_data: schemas.MarkerCreate, + node: DockerVM = Depends(dep_node) +) -> dict: + """ + Attach a traffic-insight ``mark`` filter to the Docker node's uBridge bridge. + """ + + pcap_path = os.path.join( + node.project.markers_working_directory(), + f"{node.id}_{marker_data.link_id}_{marker_data.name}.pcap" + ) + await node.start_marker(adapter_number, marker_data.name, marker_data.bpf, pcap_path, marker_data.tag) + MarkerManager.instance().register( + str(project_id), node.id, marker_data.name, marker_data.link_id, marker_data.tag + ) + return {"pcap_file_path": str(pcap_path)} + + +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/stop", + status_code=status.HTTP_204_NO_CONTENT, + dependencies=[Depends(compute_authentication)] +) +async def stop_docker_node_marker( + adapter_number: int, + port_number: int, + marker_data: schemas.MarkerDelete, + node: DockerVM = Depends(dep_node) +) -> None: + """ + Remove a traffic-insight ``mark`` filter from the Docker node's uBridge bridge. + """ + + await node.stop_marker(adapter_number, marker_data.name) + MarkerManager.instance().unregister(node.id, marker_data.name) + + @router.get( "/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream", dependencies=[Depends(compute_authentication)] diff --git a/gns3server/api/routes/compute/qemu_nodes.py b/gns3server/api/routes/compute/qemu_nodes.py index 623ae8dae..6c8456eeb 100644 --- a/gns3server/api/routes/compute/qemu_nodes.py +++ b/gns3server/api/routes/compute/qemu_nodes.py @@ -30,6 +30,7 @@ from gns3server import schemas from gns3server.compute import qemu from gns3server.compute.qemu import Qemu from gns3server.compute.qemu.qemu_vm import QemuVM +from gns3server.compute.marker.marker_manager import MarkerManager from .dependencies.authentication import compute_authentication, ws_compute_authentication @@ -382,6 +383,52 @@ async def stop_qemu_node_capture( await node.stop_capture(adapter_number) +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/start", + dependencies=[Depends(compute_authentication)] +) +async def start_qemu_node_marker( + *, + project_id: UUID, + adapter_number: int, + marker_data: schemas.MarkerCreate, + port_number: int = Path(..., ge=0, le=0), + node: QemuVM = Depends(dep_node) +) -> dict: + """ + Attach a traffic-insight ``mark`` filter to the QEMU node's uBridge bridge. + """ + + pcap_path = os.path.join( + node.project.markers_working_directory(), + f"{node.id}_{marker_data.link_id}_{marker_data.name}.pcap" + ) + await node.start_marker(adapter_number, marker_data.name, marker_data.bpf, pcap_path, marker_data.tag) + MarkerManager.instance().register( + str(project_id), node.id, marker_data.name, marker_data.link_id, marker_data.tag + ) + return {"pcap_file_path": str(pcap_path)} + + +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/stop", + status_code=status.HTTP_204_NO_CONTENT, + dependencies=[Depends(compute_authentication)] +) +async def stop_qemu_node_marker( + adapter_number: int, + marker_data: schemas.MarkerDelete, + port_number: int = Path(..., ge=0, le=0), + node: QemuVM = Depends(dep_node) +) -> None: + """ + Remove a traffic-insight ``mark`` filter from the QEMU node's uBridge bridge. + """ + + await node.stop_marker(adapter_number, marker_data.name) + MarkerManager.instance().unregister(node.id, marker_data.name) + + @router.get( "/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream", dependencies=[Depends(compute_authentication)] diff --git a/gns3server/api/routes/compute/vpcs_nodes.py b/gns3server/api/routes/compute/vpcs_nodes.py index c51f0254d..60360f579 100644 --- a/gns3server/api/routes/compute/vpcs_nodes.py +++ b/gns3server/api/routes/compute/vpcs_nodes.py @@ -29,6 +29,7 @@ from uuid import UUID from gns3server import schemas from gns3server.compute.vpcs import VPCS from gns3server.compute.vpcs.vpcs_vm import VPCSVM +from gns3server.compute.marker.marker_manager import MarkerManager from .dependencies.authentication import compute_authentication, ws_compute_authentication @@ -303,6 +304,54 @@ async def stop_vpcs_node_capture( await node.stop_capture(port_number) +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/start", + dependencies=[Depends(compute_authentication)] +) +async def start_vpcs_node_marker( + *, + project_id: UUID, + port_number: int, + marker_data: schemas.MarkerCreate, + adapter_number: int = Path(..., ge=0, le=0), + node: VPCSVM = Depends(dep_node) +) -> dict: + """ + Attach a traffic-insight ``mark`` filter to the VPCS node's uBridge bridge. + On BPF match uBridge emits a MARK signal and appends the packet to the pcap. + """ + + pcap_path = os.path.join( + node.project.markers_working_directory(), + f"{node.id}_{marker_data.link_id}_{marker_data.name}.pcap" + ) + await node.start_marker(port_number, marker_data.name, marker_data.bpf, pcap_path, marker_data.tag) + MarkerManager.instance().register( + str(project_id), node.id, marker_data.name, marker_data.link_id, marker_data.tag + ) + return {"pcap_file_path": pcap_path} + + +@router.post( + "/{node_id}/adapters/{adapter_number}/ports/{port_number}/markers/stop", + status_code=status.HTTP_204_NO_CONTENT, + dependencies=[Depends(compute_authentication)] +) +async def stop_vpcs_node_marker( + *, + port_number: int, + marker_data: schemas.MarkerDelete, + adapter_number: int = Path(..., ge=0, le=0), + node: VPCSVM = Depends(dep_node) +) -> None: + """ + Remove a traffic-insight ``mark`` filter from the VPCS node's uBridge bridge. + """ + + await node.stop_marker(port_number, marker_data.name) + MarkerManager.instance().unregister(node.id, marker_data.name) + + @router.get( "/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream", dependencies=[Depends(compute_authentication)] diff --git a/gns3server/api/routes/controller/links.py b/gns3server/api/routes/controller/links.py index ce9caf450..bb723927a 100644 --- a/gns3server/api/routes/controller/links.py +++ b/gns3server/api/routes/controller/links.py @@ -424,6 +424,85 @@ async def web_wireshark_websocket( pass +@router.get( + "/{link_id}/markers", + dependencies=[Depends(has_privilege("Link.Audit"))] +) +async def get_markers(link: Link = Depends(dep_link)) -> dict: + """ + Return all traffic-insight markers configured on this link. + + Required privilege: Link.Audit + """ + + return link.markers + + +@router.post( + "/{link_id}/markers", + status_code=status.HTTP_201_CREATED, + dependencies=[Depends(has_privilege("Link.Modify"))] +) +async def create_marker( + marker_data: schemas.MarkerCreate, + link: Link = Depends(dep_link) +) -> dict: + """ + Attach a traffic-insight marker to the link. + On BPF match uBridge emits MARK signals and appends packets to a pcap. + + Required privilege: Link.Modify + """ + + await link.start_marker( + name=marker_data.name or f"marker-{link.id[:8]}", + bpf=marker_data.bpf, + tag=marker_data.tag, + ) + return link.markers.get(marker_data.name, {}) + + +@router.delete( + "/{link_id}/markers/{marker_name}", + status_code=status.HTTP_204_NO_CONTENT, + dependencies=[Depends(has_privilege("Link.Modify"))] +) +async def delete_marker( + marker_name: str, + link: Link = Depends(dep_link) +) -> None: + """ + Remove a traffic-insight marker from the link. + + Required privilege: Link.Modify + """ + + await link.stop_marker(marker_name) + + +@router.put( + "/{link_id}/markers/{marker_name}", + dependencies=[Depends(has_privilege("Link.Modify"))] +) +async def update_marker( + marker_name: str, + marker_data: schemas.MarkerCreate, + link: Link = Depends(dep_link) +) -> dict: + """ + Update a traffic-insight marker (change BPF, tag, or enabled). + + Required privilege: Link.Modify + """ + + await link.update_marker( + name=marker_name, + bpf=marker_data.bpf if marker_data.bpf else None, + tag=marker_data.tag, + ) + return link.markers.get(marker_name, {}) + + @router.get( "/{link_id}/iface", response_model=Union[schemas.UDPPortInfo, schemas.EthernetPortInfo], diff --git a/gns3server/compute/base_node.py b/gns3server/compute/base_node.py index 8e7681d9e..fc2701ab0 100644 --- a/gns3server/compute/base_node.py +++ b/gns3server/compute/base_node.py @@ -935,9 +935,38 @@ class BaseNode: f"Hypervisor {self._ubridge_hypervisor.host}:{self._ubridge_hypervisor.port} has successfully started" ) await self._ubridge_hypervisor.connect() + # Tell this uBridge where to send MARK signals and which node id to + # tag them with. Marker is opt-in and inert until a `mark` filter is + # added, so this never disturbs the data plane. + await self._ubridge_configure_marker_sink() # save if privileged are required in case uBridge needs to be restarted in self._ubridge_send() self._ubridge_require_privileged_access = require_privileged_access + async def _ubridge_configure_marker_sink(self): + """ + Point this node's uBridge at the compute's marker UDP sink and tag its + signals with this node's id. Safe to call before any marker filter + exists — uBridge stays inert until a ``mark`` filter is configured. + + Old uBridge builds without the marker module are tolerated: the failure + is downgraded to a warning so node start is not blocked by an opt-in + observability feature. + """ + + from gns3server.compute.marker.marker_manager import MarkerManager + + manager = MarkerManager.instance() + if not manager.running or not manager.host or not manager.port: + return + try: + await self._ubridge_send(f"marker sink {manager.host} {manager.port}") + await self._ubridge_send(f"marker node {self._id}") + except UbridgeError: + log.warning( + "uBridge does not support the marker module; traffic insight disabled for node %r", + self.name, + ) + async def _stop_ubridge(self): """ Stops uBridge. @@ -1042,6 +1071,47 @@ class BaseNode: ) i += 1 + async def _ubridge_add_marker_filter(self, bridge_name, name, bpf, pcap_path, tag=None): + """ + Attach a `mark` packet filter to a uBridge bridge for traffic insight. + + On BPF match uBridge (a) emits a UDP MARK signal to the configured sink + and (b) appends the packet to ``pcap_path``. Unlike the impairment + filters, this is an observability tap: it never drops or alters traffic, + and it is added/removed on its own (not via reset_packet_filters) so the + pcap is not closed/reopened on unrelated filter changes. + + :param bridge_name: uBridge bridge carrying the link's traffic + :param name: stable, gns3server-chosen filter name (pcap identity + echoed in signals) + :param bpf: libpcap BPF expression + :param pcap_path: absolute path ubridge appends matched packets to + :param tag: optional correlation id echoed in MARK signals + """ + + # mark [tag ] [pcap ] — tag/pcap keyword pairs, any order. + cmd = 'bridge add_packet_filter {bridge} {name} mark "{bpf}"'.format( + bridge=bridge_name, name=name, bpf=bpf + ) + if tag is not None: + cmd += f" tag {tag}" + cmd += ' pcap "{path}"'.format(path=pcap_path) + # Let BPF compile errors propagate — the marker is the user's intent, so a + # bad expression must surface instead of being silently dropped. + await self._ubridge_send(cmd) + + async def _ubridge_delete_marker_filter(self, bridge_name, name): + """ + Remove a `mark` filter from a uBridge bridge. + + uBridge closes and flushes the filter's pcap on delete; the file itself + persists on disk for later replay. + + :param bridge_name: uBridge bridge the filter is attached to + :param name: filter name previously passed to _ubridge_add_marker_filter + """ + + await self._ubridge_send(f"bridge delete_packet_filter {bridge_name} {name}") + async def _add_ubridge_ethernet_connection(self, bridge_name, ethernet_interface, block_host_traffic=False): """ Creates a connection with an Ethernet interface in uBridge. diff --git a/gns3server/compute/docker/docker_vm.py b/gns3server/compute/docker/docker_vm.py index 8c139cdfd..5b787a39f 100644 --- a/gns3server/compute/docker/docker_vm.py +++ b/gns3server/compute/docker/docker_vm.py @@ -1440,6 +1440,44 @@ class DockerVM(BaseNode): ) ) + async def start_marker(self, adapter_number, name, bpf, pcap_path, tag=None): + """ + Attach a traffic-insight ``mark`` filter to this adapter's uBridge bridge. + On BPF match uBridge emits a MARK signal and appends the packet to the pcap. + + :param adapter_number: adapter number + :param name: stable filter name — pcap identity and echoed in MARK signals + :param bpf: libpcap BPF expression + :param pcap_path: absolute path ubridge appends matched packets to + :param tag: optional correlation id echoed in MARK signals + """ + + if self.status == "started" and self.ubridge: + adapter = f"bridge{adapter_number}" + await self._ubridge_add_marker_filter(adapter, name, bpf, pcap_path, tag) + log.info( + "Docker VM '{name}' [{id}]: starting marker '{marker}' on adapter {adapter_number}".format( + name=self.name, id=self.id, marker=name, adapter_number=adapter_number + ) + ) + + async def stop_marker(self, adapter_number, name): + """ + Remove a traffic-insight ``mark`` filter from this adapter's uBridge bridge. + + :param adapter_number: adapter number + :param name: filter name previously passed to start_marker + """ + + if self.status == "started" and self.ubridge: + adapter = f"bridge{adapter_number}" + await self._ubridge_delete_marker_filter(adapter, name) + log.info( + "Docker VM '{name}' [{id}]: stopping marker '{marker}' on adapter {adapter_number}".format( + name=self.name, id=self.id, marker=name, adapter_number=adapter_number + ) + ) + async def _get_log(self): """ Returns the log from the container diff --git a/gns3server/compute/marker/__init__.py b/gns3server/compute/marker/__init__.py new file mode 100644 index 000000000..8fdb2b775 --- /dev/null +++ b/gns3server/compute/marker/__init__.py @@ -0,0 +1,24 @@ +#!/usr/bin/env python +# +# Copyright (C) 2024 GNS3 Technologies Inc. +# +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . +# +# +# Traffic-insight marker subsystem (compute side). +# +# ubridge's ``marker`` module is a passive tap: on a BPF match it emits a UDP +# ``MARK`` signal to a configured sink and/or appends the packet to a pcap. +# This package owns the compute-side UDP sink: one listener per compute process +# serves every ubridge on that host, disambiguated by ``node=``. diff --git a/gns3server/compute/marker/marker_listener.py b/gns3server/compute/marker/marker_listener.py new file mode 100644 index 000000000..d8bfa3303 --- /dev/null +++ b/gns3server/compute/marker/marker_listener.py @@ -0,0 +1,102 @@ +#!/usr/bin/env python +# +# Copyright (C) 2024 GNS3 Technologies Inc. +# +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . + +import asyncio +import logging + +log = logging.getLogger(__name__) + + +class MarkerListener(asyncio.DatagramProtocol): + """ + Receives ubridge ``MARK`` signal datagrams and turns each into a + ``marker.match`` notification. + + Signal format (one datagram per match, ASCII):: + + MARK node= filter= tag= len=\\n + + The signal carries metadata only (no packet bytes). The compute-side + :class:`~gns3server.compute.marker.marker_manager.MarkerManager` registry + resolves ``(node_id, filter_name)`` to ``(project_id, link_id, tag)`` so the + event can be emitted on the right project-scoped notification stream. + """ + + def __init__(self, manager): + # MarkerManager owns this listener and the registry. + self._manager = manager + self.transport = None + + def connection_made(self, transport): + self.transport = transport + + def datagram_received(self, data, addr): + try: + self._handle(data) + except Exception: + # Never let a malformed datagram kill the listener. + log.exception("Failed to process MARK datagram from %s: %r", addr, data) + + def _handle(self, data): + line = data.decode("utf-8", errors="replace").strip() + if not line.startswith("MARK"): + return + + parts = line.split() + # parts[0] == "MARK"; parts[1] == "" + if len(parts) < 2: + return + + try: + ts = float(parts[1]) + except ValueError: + log.warning("Ignoring MARK signal with bad timestamp: %r", line) + return + + kv = {} + for token in parts[2:]: + if "=" in token: + key, value = token.split("=", 1) + kv[key] = value + + node_id = kv.get("node") + filter_name = kv.get("filter") + if not node_id or not filter_name: + return + + # "-" means the field was unset on the ubridge side (see contract §3.3). + tag = kv.get("tag") + length = kv.get("len") + + project_id, link_id, registered_tag = self._manager.lookup(node_id, filter_name) + if project_id is None: + log.warning( + "MARK signal for unregistered node=%s filter=%s, dropping", node_id, filter_name + ) + return + + event = { + "project_id": project_id, + "node_id": node_id, + "link_id": link_id, + "filter": filter_name, + # Prefer the value carried in the signal; fall back to the one we registered. + "tag": tag if tag and tag != "-" else registered_tag, + "ts": ts, + "len": int(length) if length and length.isdigit() else 0, + } + self._manager.emit_match(project_id, event) diff --git a/gns3server/compute/marker/marker_manager.py b/gns3server/compute/marker/marker_manager.py new file mode 100644 index 000000000..7a4168762 --- /dev/null +++ b/gns3server/compute/marker/marker_manager.py @@ -0,0 +1,163 @@ +#!/usr/bin/env python +# +# Copyright (C) 2024 GNS3 Technologies Inc. +# +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . + +import asyncio +import logging + +from gns3server.compute.marker.marker_listener import MarkerListener +from gns3server.compute.notification_manager import NotificationManager + +log = logging.getLogger(__name__) + + +class MarkerManager: + """ + Singleton owning the compute-side UDP sink for ubridge ``MARK`` signals and + the registry that maps each ``(node_id, filter_name)`` back to its + ``(project_id, link_id, tag)``. + + The registry is populated when a marker is created on a link (the compute + endpoint has project_id + node_id from its route path and link_id/name/tag + from the request body) and cleared when the marker is deleted or the project + closed. At signal time it is an O(1) lookup — no node-table scan, and the + signal payload is untouched. + + One listener per compute process serves every ubridge on that host; source + ubridges are disambiguated by ``node=`` (UUID, globally unique). + """ + + def __init__(self): + + self._listener = None + self._transport = None + self._host = None + self._port = None + # Flat lookup: (node_id, filter_name) -> {"project_id", "link_id", "tag"} + self._entries = {} + # Reverse index for O(1) per-project teardown: project_id -> set of keys + self._by_project = {} + + @property + def host(self): + """The host the UDP sink is reachable on (for ``marker sink``).""" + return self._host + + @property + def port(self): + """The UDP port the sink is bound on (for ``marker sink``).""" + return self._port + + @property + def running(self): + return self._transport is not None + + async def start(self, host="127.0.0.1", port=0): + """ + Bind the UDP sink. ``port=0`` lets the OS choose a free port, which is + then read back and exposed via :attr:`port` for ``marker sink`` commands. + """ + + if self.running: + return + loop = asyncio.get_running_loop() + self._listener = MarkerListener(self) + self._transport, _ = await loop.create_datagram_endpoint( + lambda: self._listener, local_addr=(host, port) + ) + sock = self._transport.get_extra_info("socket") + self._host = host + self._port = sock.getsockname()[1] if sock else port + log.info("Marker signal sink listening on %s:%s", self._host, self._port) + + async def stop(self): + """Close the UDP sink and drop the whole registry.""" + + if self._transport: + self._transport.close() + self._transport = None + self._listener = None + self._entries.clear() + self._by_project.clear() + self._host = None + self._port = None + + def register(self, project_id, node_id, filter_name, link_id, tag=None): + """ + Record that ``filter_name`` on ``node_id`` belongs to ``project_id`` / + ``link_id``. Called from the compute marker-start endpoint. + + Re-registering the same key updates the stored tag (e.g. on re-add). + """ + + key = (node_id, filter_name) + self._entries[key] = {"project_id": project_id, "link_id": link_id, "tag": tag} + self._by_project.setdefault(project_id, set()).add(key) + + def unregister(self, node_id, filter_name): + """Forget a single marker. Returns True if something was removed.""" + + key = (node_id, filter_name) + entry = self._entries.pop(key, None) + if entry is None: + return False + project_entries = self._by_project.get(entry["project_id"]) + if project_entries is not None: + project_entries.discard(key) + if not project_entries: + self._by_project.pop(entry["project_id"], None) + return True + + def unregister_project(self, project_id): + """Drop every marker belonging to ``project_id`` (project close).""" + + keys = self._by_project.pop(project_id, None) + if not keys: + return + for key in keys: + self._entries.pop(key, None) + + def lookup(self, node_id, filter_name): + """ + O(1) resolution of an incoming signal to its project/link/tag. + + :returns: (project_id, link_id, tag) or (None, None, None) on miss. + """ + + entry = self._entries.get((node_id, filter_name)) + if entry is None: + return None, None, None + return entry["project_id"], entry["link_id"], entry["tag"] + + def emit_match(self, project_id, event): + """ + Forward a parsed match as a project-scoped ``marker.match`` notification. + Flows compute -> controller dispatch -> project_emit -> web UI WS. + """ + + NotificationManager.instance().emit("marker.match", event, project_id=project_id) + + _instance = None + + @staticmethod + def instance(): + if MarkerManager._instance is None: + MarkerManager._instance = MarkerManager() + return MarkerManager._instance + + @staticmethod + def reset(): + MarkerManager._instance = None diff --git a/gns3server/compute/project.py b/gns3server/compute/project.py index eebb65ca1..ae3b5c8d1 100644 --- a/gns3server/compute/project.py +++ b/gns3server/compute/project.py @@ -246,6 +246,22 @@ class Project: raise ComputeError(f"Could not create the capture working directory: {e}") return workdir + def markers_working_directory(self): + """ + Returns the working directory where uBridge writes per-link marker pcaps + (matched packets, kept for later replay). + + :returns: path to the directory + """ + + workdir = os.path.join(self._path, "project-files", "markers") + if not self._deleted: + try: + os.makedirs(workdir, exist_ok=True) + except OSError as e: + raise ComputeError(f"Could not create the markers working directory: {e}") + return workdir + def add_node(self, node): """ Adds a node to the project. diff --git a/gns3server/compute/qemu/qemu_vm.py b/gns3server/compute/qemu/qemu_vm.py index f91177111..b659d9cb4 100644 --- a/gns3server/compute/qemu/qemu_vm.py +++ b/gns3server/compute/qemu/qemu_vm.py @@ -1697,6 +1697,44 @@ class QemuVM(BaseNode): ) ) + async def start_marker(self, adapter_number, name, bpf, pcap_path, tag=None): + """ + Attach a traffic-insight ``mark`` filter to this adapter's uBridge bridge. + On BPF match uBridge emits a MARK signal and appends the packet to the pcap. + + :param adapter_number: adapter number + :param name: stable filter name — pcap identity and echoed in MARK signals + :param bpf: libpcap BPF expression + :param pcap_path: absolute path ubridge appends matched packets to + :param tag: optional correlation id echoed in MARK signals + """ + + if self.ubridge: + await self._ubridge_add_marker_filter( + f"QEMU-{self._id}-{adapter_number}", name, bpf, pcap_path, tag + ) + log.info( + "QEMU VM '{name}' [{id}]: starting marker '{marker}' on adapter {adapter_number}".format( + name=self.name, id=self.id, marker=name, adapter_number=adapter_number + ) + ) + + async def stop_marker(self, adapter_number, name): + """ + Remove a traffic-insight ``mark`` filter from this adapter's uBridge bridge. + + :param adapter_number: adapter number + :param name: filter name previously passed to start_marker + """ + + if self.ubridge: + await self._ubridge_delete_marker_filter(f"QEMU-{self._id}-{adapter_number}", name) + log.info( + "QEMU VM '{name}' [{id}]: stopping marker '{marker}' on adapter {adapter_number}".format( + name=self.name, id=self.id, marker=name, adapter_number=adapter_number + ) + ) + async def create_disk_image(self, disk_name, options): """ Create a Qemu disk diff --git a/gns3server/compute/vpcs/vpcs_vm.py b/gns3server/compute/vpcs/vpcs_vm.py index b42c4095a..0d3353a9a 100644 --- a/gns3server/compute/vpcs/vpcs_vm.py +++ b/gns3server/compute/vpcs/vpcs_vm.py @@ -512,6 +512,42 @@ class VPCSVM(BaseNode): ) ) + async def start_marker(self, port_number, name, bpf, pcap_path, tag=None): + """ + Attach a traffic-insight ``mark`` filter to this node's uBridge bridge. + On BPF match uBridge emits a MARK signal and appends the packet to the pcap. + + :param port_number: port number (kept for API symmetry; VPCS has a single bridge) + :param name: stable filter name — pcap identity and echoed in MARK signals + :param bpf: libpcap BPF expression + :param pcap_path: absolute path ubridge appends matched packets to + :param tag: optional correlation id echoed in MARK signals + """ + + if self.ubridge: + await self._ubridge_add_marker_filter(f"VPCS-{self._id}", name, bpf, pcap_path, tag) + log.info( + "VPCS '{name}' [{id}]: starting marker '{marker}' on port {port_number}".format( + name=self.name, id=self.id, marker=name, port_number=port_number + ) + ) + + async def stop_marker(self, port_number, name): + """ + Remove a traffic-insight ``mark`` filter from this node's uBridge bridge. + + :param port_number: port number (kept for API symmetry) + :param name: filter name previously passed to start_marker + """ + + if self.ubridge: + await self._ubridge_delete_marker_filter(f"VPCS-{self._id}", name) + log.info( + "VPCS '{name}' [{id}]: stopping marker '{marker}' on port {port_number}".format( + name=self.name, id=self.id, marker=name, port_number=port_number + ) + ) + def _build_command(self): """ Command to start the VPCS process. diff --git a/gns3server/controller/link.py b/gns3server/controller/link.py index 556f5ba5f..0e959cef8 100644 --- a/gns3server/controller/link.py +++ b/gns3server/controller/link.py @@ -88,6 +88,7 @@ class Link: self._link_type = "ethernet" self._suspended = False self._filters = {} + self._markers = {} self._link_style = {} self._wireshark = False self._show_filters_icon = True @@ -99,6 +100,13 @@ class Link: """ return self._filters + @property + def markers(self): + """ + Get the traffic insight markers dict: name → {bpf, tag, enabled} + """ + return self._markers + @property def show_filters_icon(self): """ @@ -298,6 +306,27 @@ class Link: raise NotImplementedError + async def start_marker(self, name, bpf, tag=None): + """ + Attach a traffic-insight marker to this link (base — UDPLink overrides). + """ + raise NotImplementedError + + async def stop_marker(self, name): + """ + Remove a traffic-insight marker from this link (base — UDPLink overrides). + """ + raise NotImplementedError + + async def update_marker(self, name, bpf=None, tag=None, enabled=None): + """ + Update an existing marker's BPF, tag, or enabled flag. + + A BPF change is a delete+re-add on the ubridge side so the pcap is + flushed and the new filter takes effect. + """ + raise NotImplementedError + async def start_capture(self, data_link_type="DLT_EN10MB", capture_file_name=None, wireshark=False, jwt_token=None): """ Start capture on the link @@ -571,6 +600,7 @@ class Link: "nodes": res, "link_id": self._id, "filters": self._filters, + "markers": self._markers, "link_style": self._link_style, "suspend": self._suspended, "show_filters_icon": getattr(self, '_show_filters_icon', True), @@ -585,6 +615,7 @@ class Link: "capture_compute_id": self.capture_compute_id, "link_type": self._link_type, "filters": self._filters, + "markers": self._markers, "suspend": self._suspended, "link_style": self._link_style, "wireshark": self._wireshark, diff --git a/gns3server/controller/udp_link.py b/gns3server/controller/udp_link.py index 7c2e6f187..42e34ae47 100644 --- a/gns3server/controller/udp_link.py +++ b/gns3server/controller/udp_link.py @@ -19,6 +19,7 @@ from .controller_error import ControllerError, ControllerNotFoundError from .link import Link from .node_types import BUILTIN_NODE_TYPES +from gns3server.utils.packet_filter_validation import validate_bpf_syntax, FilterValidationError class UDPLink(Link): @@ -26,6 +27,8 @@ class UDPLink(Link): super().__init__(project, link_id=link_id) self._created = False self._link_data = [] + # Runtime-only Node references for marker commands (not serialized). + self._marker_capture_nodes = {} @property def debug_link_data(self): @@ -251,3 +254,132 @@ class UDPLink(Link): """ if self._capture_node and node == self._capture_node["node"] and node.status != "started": await self.stop_capture() + # Tear down any marker whose capture-side node just stopped. + for name, marker_info in list(self._markers.items()): + if marker_info.get("capture_node_id") == node.id and node.status != "started": + await self.stop_marker(name) + + def _capture_node_for_marker(self, name): + """Return the stored (node, adapter_number, port_number) for a marker's capture side.""" + return self._marker_capture_nodes.get(name) + + def _store_capture_node_for_marker(self, name, capture_side): + """Persist the capture-side identity (serializable refs) + runtime Node.""" + self._markers[name] = { + **self._markers.get(name, {}), + "capture_node_id": capture_side["node"].id, + "capture_adapter": capture_side["adapter_number"], + "capture_port": capture_side["port_number"], + } + self._marker_capture_nodes[name] = capture_side + + async def start_marker(self, name, bpf, tag=None): + """ + Attach a traffic-insight marker to this link. + + :param name: stable filter name — echoed in MARK signals + pcap identity + :param bpf: libpcap BPF expression + :param tag: optional correlation id + """ + + if name in self._markers: + raise ControllerError(f"Marker '{name}' already exists on link {self._id}") + + # Pre-validate BPF on the controller side before reaching ubridge. + result = validate_bpf_syntax(bpf) + if not result.get("valid"): + raise ControllerError(f"Invalid BPF expression: {result.get('error', 'unknown error')}") + + capture_side = self._choose_capture_side() + data = {"name": name, "bpf": bpf, "tag": tag, "link_id": self._id} + await capture_side["node"].post( + "/adapters/{adapter_number}/ports/{port_number}/markers/start".format( + adapter_number=capture_side["adapter_number"], port_number=capture_side["port_number"] + ), + data=data, + ) + self._store_capture_node_for_marker(name, capture_side) + self._markers[name].update({"bpf": bpf, "tag": tag, "enabled": True}) + self._project.emit_notification("link.updated", self.asdict()) + self._project.dump() + + async def stop_marker(self, name): + """ + Remove a traffic-insight marker from this link. + + :param name: filter name to remove + """ + + if name not in self._markers: + raise ControllerNotFoundError(f"Marker '{name}' not found on link {self._id}") + + capture_side = self._marker_capture_nodes.get(name) + if capture_side: + await capture_side["node"].post( + "/adapters/{adapter_number}/ports/{port_number}/markers/stop".format( + adapter_number=capture_side["adapter_number"], + port_number=capture_side["port_number"], + ), + data={"name": name}, + ) + self._markers.pop(name, None) + self._marker_capture_nodes.pop(name, None) + self._project.emit_notification("link.updated", self.asdict()) + self._project.dump() + + async def update_marker(self, name, bpf=None, tag=None, enabled=None): + """ + Update an existing marker. A BPF change requires delete+re-add so the + ubridge side flushes the pcap and the new filter takes effect. + + :param name: filter name to update + :param bpf: new BPF expression (None = keep existing) + :param tag: new tag id (None = keep existing) + :param enabled: toggle (None = keep existing) + """ + + marker_info = self._markers.get(name) + if not marker_info: + raise ControllerNotFoundError(f"Marker '{name}' not found on link {self._id}") + + new_bpf = bpf if bpf is not None else marker_info["bpf"] + new_tag = tag if tag is not None else marker_info.get("tag") + new_enabled = enabled if enabled is not None else marker_info.get("enabled", True) + + if not new_enabled and marker_info.get("enabled", True): + # Toggle off: remove from ubridge but keep state. + await self.stop_marker(name) + self._markers[name] = {**marker_info, "bpf": new_bpf, "tag": new_tag, "enabled": False} + self._project.emit_notification("link.updated", self.asdict()) + self._project.dump() + return + + capture_side = self._marker_capture_nodes.get(name) + if new_bpf != marker_info.get("bpf") or new_tag != marker_info.get("tag"): + # BPF or tag changed — re-validate, delete, re-add. + if new_bpf != marker_info.get("bpf"): + result = validate_bpf_syntax(new_bpf) + if not result.get("valid"): + raise ControllerError(f"Invalid BPF expression: {result.get('error', 'unknown error')}") + if capture_side: + # Delete old filter from ubridge. + await capture_side["node"].post( + "/adapters/{adapter_number}/ports/{port_number}/markers/stop".format( + adapter_number=capture_side["adapter_number"], + port_number=capture_side["port_number"], + ), + data={"name": name}, + ) + # Re-add with new params. + data = {"name": name, "bpf": new_bpf, "tag": new_tag, "link_id": self._id} + await capture_side["node"].post( + "/adapters/{adapter_number}/ports/{port_number}/markers/start".format( + adapter_number=capture_side["adapter_number"], + port_number=capture_side["port_number"], + ), + data=data, + ) + self._markers[name] = {**marker_info, "bpf": new_bpf, "tag": new_tag, "enabled": True} + + self._project.emit_notification("link.updated", self.asdict()) + self._project.dump() diff --git a/gns3server/core/tasks.py b/gns3server/core/tasks.py index d4968d8f5..9fc911d55 100644 --- a/gns3server/core/tasks.py +++ b/gns3server/core/tasks.py @@ -24,6 +24,7 @@ from gns3server.controller import Controller from gns3server.config import Config from gns3server.compute import MODULES from gns3server.compute.port_manager import PortManager +from gns3server.compute.marker.marker_manager import MarkerManager from gns3server.utils.http_client import HTTPClient from gns3server.db.tasks import connect_to_db, get_computes, disconnect_from_db, discover_images_on_filesystem @@ -84,6 +85,14 @@ async def startup(app: FastAPI) -> None: m = module.instance() m.port_manager = PortManager.instance() + # Start the marker (traffic-insight) UDP sink. One listener per compute + # process receives ubridge MARK signals; ubridges are told its host/port at + # startup (see BaseNode._start_ubridge). + server_settings = Config.instance().settings.Server + await MarkerManager.instance().start( + host=server_settings.marker_listen_host, port=server_settings.marker_listen_port + ) + # Mark MCP server as ready to accept connections (if MCP is available) from gns3server.agent import MCP_AVAILABLE @@ -101,6 +110,7 @@ async def shutdown(app: FastAPI) -> None: if auto_discover_images_task_handle is not None and not auto_discover_images_task_handle.cancelled(): auto_discover_images_task_handle.cancel() await HTTPClient.close_session() + await MarkerManager.instance().stop() await Controller.instance().stop() for module in MODULES: diff --git a/gns3server/schemas/__init__.py b/gns3server/schemas/__init__.py index d228f4b84..3c8d0ac63 100644 --- a/gns3server/schemas/__init__.py +++ b/gns3server/schemas/__init__.py @@ -20,7 +20,7 @@ from .common import ErrorMessage from .version import Version # Controller schemas -from .controller.links import LinkCreate, LinkUpdate, Link, UDPPortInfo, EthernetPortInfo, LinkCapture +from .controller.links import LinkCreate, LinkUpdate, Link, UDPPortInfo, EthernetPortInfo, LinkCapture, MarkerCreate, MarkerDelete from .controller.computes import ComputeCreate, ComputeUpdate, ComputeVirtualBoxVM, ComputeVMwareVM, ComputeDockerImage, AutoIdlePC, Compute from .controller.templates import TemplateCreate, TemplateUpdate, TemplateUsage, Template from .controller.images import Image, ImageType diff --git a/gns3server/schemas/config.py b/gns3server/schemas/config.py index 294ea4cf6..6b00bff75 100644 --- a/gns3server/schemas/config.py +++ b/gns3server/schemas/config.py @@ -153,6 +153,12 @@ class ServerSettings(BaseModel): udp_start_port_range: int = Field(10000, gt=0, le=65535) udp_end_port_range: int = Field(30000, gt=0, le=65535) ubridge_path: str = "ubridge" + # Marker (traffic-insight) UDP sink: one listener per compute process that + # receives ubridge MARK signals from every ubridge on this host. The host + # defaults to loopback because ubridge runs on the same host as the compute. + # port=0 lets the OS choose a free port (read back and handed to ubridge). + marker_listen_host: str = "127.0.0.1" + marker_listen_port: int = Field(0, ge=0, le=65535) compute_username: str = "gns3" compute_password: SecretStr = SecretStr("") allowed_interfaces: List[str] = Field(default_factory=list) diff --git a/gns3server/schemas/controller/links.py b/gns3server/schemas/controller/links.py index be6846bb8..7a0676e04 100644 --- a/gns3server/schemas/controller/links.py +++ b/gns3server/schemas/controller/links.py @@ -62,6 +62,10 @@ class LinkBase(BaseModel): suspend: Optional[bool] = None link_style: Optional[LinkStyle] = None filters: Optional[dict] = None + markers: Optional[dict] = Field( + None, + description="Traffic-insight markers on this link: name → {bpf, tag, enabled}" + ) show_filters_icon: Optional[bool] = Field( True, description="Show filters icon in Web UI" @@ -135,3 +139,25 @@ class LinkCapture(BaseModel): data_link_type: str = "DLT_EN10MB" capture_file_name: Optional[str] = None wireshark: bool = False + + +class MarkerCreate(BaseModel): + """ + Body for attaching a traffic-insight marker to a link. + + ``name`` is optional at the controller REST layer (auto-generated when + absent) but always set when the controller forwards to the compute. + """ + + name: Optional[str] = None + bpf: str + tag: Optional[int] = None + link_id: Optional[str] = None + + +class MarkerDelete(BaseModel): + """ + Body for removing a traffic-insight marker from a link. + """ + + name: str diff --git a/tests/compute/marker/__init__.py b/tests/compute/marker/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/tests/compute/marker/test_marker_manager.py b/tests/compute/marker/test_marker_manager.py new file mode 100644 index 000000000..dc8e826ff --- /dev/null +++ b/tests/compute/marker/test_marker_manager.py @@ -0,0 +1,218 @@ +#!/usr/bin/env python +# +# Copyright (C) 2024 GNS3 Technologies Inc. +# +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . + +import asyncio +import pytest + +from gns3server.compute.marker.marker_manager import MarkerManager +from gns3server.compute.marker.marker_listener import MarkerListener + + +# --------------------------------------------------------------------------- +# Registry +# --------------------------------------------------------------------------- + +class TestMarkerRegistry: + + def test_register_and_lookup(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + mgr.register("proj1", "node1", "filter1", "link1", tag=5) + pid, lid, tag = mgr.lookup("node1", "filter1") + assert pid == "proj1" + assert lid == "link1" + assert tag == 5 + + def test_miss_returns_none(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + pid, lid, tag = mgr.lookup("no-such-node", "no-such-filter") + assert pid is None + + def test_reregister_updates(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + mgr.register("p", "n", "f", "l", tag=1) + mgr.register("p", "n", "f", "l", tag=99) + _, _, tag = mgr.lookup("n", "f") + assert tag == 99 + + def test_unregister(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + mgr.register("p", "n", "f", "l") + assert mgr.unregister("n", "f") is True + pid, _, _ = mgr.lookup("n", "f") + assert pid is None + assert mgr.unregister("n", "f") is False + + def test_unregister_project(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + mgr.register("p1", "n1", "f1", "l1") + mgr.register("p1", "n2", "f2", "l2") + mgr.register("p2", "n3", "f3", "l3") + mgr.unregister_project("p1") + assert mgr.lookup("n1", "f1") == (None, None, None) + assert mgr.lookup("n2", "f2") == (None, None, None) + assert mgr.lookup("n3", "f3")[0] == "p2" + + def test_re_add_after_project_clear(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + mgr.register("p", "n", "f", "l") + mgr.unregister_project("p") + mgr.register("p", "n", "f", "l2", tag=42) + pid, lid, tag = mgr.lookup("n", "f") + assert pid == "p" and lid == "l2" and tag == 42 + + +# --------------------------------------------------------------------------- +# MarkerListener parsing +# --------------------------------------------------------------------------- + +class FakeMarkerManager: + def __init__(self): + self.events = [] + self._entries = {} + + def lookup(self, node_id, filter_name): + e = self._entries.get((node_id, filter_name)) + if e is None: + return None, None, None + return e["project_id"], e["link_id"], e["tag"] + + def emit_match(self, project_id, event): + self.events.append((project_id, event)) + + def register(self, project_id, node_id, filter_name, link_id, tag): + self._entries[(node_id, filter_name)] = { + "project_id": project_id, "link_id": link_id, "tag": tag + } + + +class TestMarkerListener: + + def test_parses_valid_mark_datagram(self): + fmgr = FakeMarkerManager() + fmgr.register("p1", "n1", "f1", "l1", tag=7) + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received( + b"MARK 1700000000.123456 node=n1 filter=f1 tag=7 len=98\n", + ("127.0.0.1", 9999), + ) + assert len(fmgr.events) == 1 + _, ev = fmgr.events[0] + assert ev["node_id"] == "n1" + assert ev["link_id"] == "l1" + assert ev["filter"] == "f1" + assert ev["tag"] == "7" + assert ev["ts"] == pytest.approx(1700000000.123456) + assert ev["len"] == 98 + + def test_unknown_node_dropped(self): + fmgr = FakeMarkerManager() + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received(b"MARK 1.0 node=bad filter=bad len=10\n", None) + assert fmgr.events == [] + + def test_bad_timestamp_ignored(self): + fmgr = FakeMarkerManager() + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received(b"MARK badts node=n filter=f len=1\n", None) + assert fmgr.events == [] + + def test_not_mark_line_ignored(self): + fmgr = FakeMarkerManager() + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received(b"HELLO world\n", None) + assert fmgr.events == [] + + def test_missing_node_ignored(self): + fmgr = FakeMarkerManager() + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received(b"MARK 1.0 filter=f len=1\n", None) + assert fmgr.events == [] + + def test_tag_dash_falls_back_to_registered(self): + fmgr = FakeMarkerManager() + fmgr.register("p", "n", "f", "l", tag=42) + lis = MarkerListener(fmgr) + lis.connection_made(None) + lis.datagram_received(b"MARK 2.0 node=n filter=f tag=- len=20\n", None) + assert fmgr.events[0][1]["tag"] == 42 + + def test_exception_does_not_kill_listener(self): + fmgr = FakeMarkerManager() + lis = MarkerListener(fmgr) + lis.connection_made(None) + # Non-decodable bytes + lis.datagram_received(b"\xff\xfe\xfd", None) + # The listener swallows exceptions; reaching here proves it survived. + assert True + + +# --------------------------------------------------------------------------- +# UDP round-trip +# --------------------------------------------------------------------------- + +@pytest.mark.asyncio +class TestMarkerManagerUDP: + + async def test_listener_receives_and_dispatches(self): + MarkerManager.reset() + mgr = MarkerManager.instance() + + captured = [] + original_emit = mgr.emit_match + mgr.emit_match = lambda pid, ev: captured.append((pid, ev)) + + await mgr.start("127.0.0.1", 0) + assert mgr.running + assert mgr.port is not None + + mgr.register("proj-rt", "node-rt", "filt-rt", "link-rt", tag=10) + + loop = asyncio.get_running_loop() + + class SendProto(asyncio.DatagramProtocol): + def connection_made(self, transport): + self.transport = transport + + sp = SendProto() + transport, _ = await loop.create_datagram_endpoint( + lambda: sp, remote_addr=("127.0.0.1", mgr.port) + ) + transport.sendto( + b"MARK 123.456 node=node-rt filter=filt-rt tag=10 len=88\n" + ) + await asyncio.sleep(0.15) + transport.close() + + mgr.emit_match = original_emit + await mgr.stop() + + assert len(captured) == 1 + pid, ev = captured[0] + assert pid == "proj-rt" + assert ev["link_id"] == "link-rt" + assert ev["len"] == 88 diff --git a/tests/controller/test_link.py b/tests/controller/test_link.py index 239682948..54d1db5ba 100644 --- a/tests/controller/test_link.py +++ b/tests/controller/test_link.py @@ -221,6 +221,7 @@ async def test_json(project, compute): } ], "filters": {}, + "markers": {}, "show_filters_icon": True, "link_style": {}, "suspend": False, @@ -255,6 +256,7 @@ async def test_json(project, compute): ], "link_style": {}, "filters": {}, + "markers": {}, "show_filters_icon": True, "suspend": False }