mirror of
https://github.com/GNS3/gns3-server.git
synced 2026-08-27 12:30:13 +03:00
refactor(marker): converge to filter single-path model, remove dual-apply endpoints
Markers now follow exactly the same apply pattern as packet filters: state lives in Link._markers, application goes through NIO (update() -> PUT /nio -> _ubridge_apply_markers). The former immediate-apply REST endpoints (/markers/start, /markers/stop on the compute side) and the per-node start_marker/stop_marker methods are removed — they were a legacy of the original capture-inspired design and have been superseded by the NIO flow. Changes: - controller/udp_link: start_marker/stop_marker/update_marker now set _markers state + call self.update() (mirrors update_filters). Removed _marker_capture_nodes runtime dict and its helpers. - controller/project: _create_link_from_topology_data restores _markers directly from persisted data (with BPF validation, like filter reload). No long calls start_marker during load. - compute: _ubridge_apply_markers swallows BPF compile errors (warn+skip), matching _ubridge_apply_filters behaviour so a single bad expression cannot break link creation / node restart. - Removed: /markers/start,stop endpoints (6 handlers across vpcs/qemu/docker route files), node start_marker/stop_marker methods (3 VM files), _ubridge_delete_marker_filter, _marker_capture_nodes, MarkerDelete schema. Net: ~280 lines of dead code removed; marker and packet filter now share a single, unified apply path via the NIO.
This commit is contained in:
parent
88c03d1431
commit
c62b9b0283
@ -29,8 +29,6 @@ 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
|
||||
|
||||
responses = {404: {"model": schemas.ErrorMessage, "description": "Could not find project or Docker node"}}
|
||||
@ -355,60 +353,6 @@ 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
|
||||
)
|
||||
nio = node.get_nio(adapter_number)
|
||||
if nio:
|
||||
nio.markers[marker_data.name] = {
|
||||
"bpf": marker_data.bpf, "tag": marker_data.tag, "link_id": marker_data.link_id
|
||||
}
|
||||
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)
|
||||
nio = node.get_nio(adapter_number)
|
||||
if nio:
|
||||
nio.markers.pop(marker_data.name, None)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream",
|
||||
dependencies=[Depends(compute_authentication)]
|
||||
|
||||
@ -30,8 +30,6 @@ 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
|
||||
|
||||
import logging
|
||||
@ -384,60 +382,6 @@ 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
|
||||
)
|
||||
nio = node.get_nio(adapter_number)
|
||||
if nio:
|
||||
nio.markers[marker_data.name] = {
|
||||
"bpf": marker_data.bpf, "tag": marker_data.tag, "link_id": marker_data.link_id
|
||||
}
|
||||
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)
|
||||
nio = node.get_nio(adapter_number)
|
||||
if nio:
|
||||
nio.markers.pop(marker_data.name, None)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream",
|
||||
dependencies=[Depends(compute_authentication)]
|
||||
|
||||
@ -29,8 +29,6 @@ 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
|
||||
|
||||
responses = {404: {"model": schemas.ErrorMessage, "description": "Could not find project or VMware node"}}
|
||||
@ -305,65 +303,6 @@ 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
|
||||
)
|
||||
# Mirror the marker spec onto the NIO so it survives a node restart (the NIO
|
||||
# persists across stop/start and _ubridge_apply_markers re-applies its markers).
|
||||
nio = node.get_nio(port_number)
|
||||
if nio:
|
||||
nio.markers[marker_data.name] = {
|
||||
"bpf": marker_data.bpf,
|
||||
"tag": marker_data.tag,
|
||||
"link_id": marker_data.link_id,
|
||||
}
|
||||
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)
|
||||
nio = node.get_nio(port_number)
|
||||
if nio:
|
||||
nio.markers.pop(marker_data.name, None)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/{node_id}/adapters/{adapter_number}/ports/{port_number}/capture/stream",
|
||||
|
||||
@ -1101,19 +1101,6 @@ class BaseNode:
|
||||
# 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 _ubridge_apply_markers(self, bridge_name, nio):
|
||||
"""
|
||||
(Re-)apply every traffic-insight marker carried by *nio* to the uBridge
|
||||
@ -1137,7 +1124,18 @@ class BaseNode:
|
||||
pcap_path = os.path.join(
|
||||
markers_dir, f"{self._id}_{link_id}_{name}.pcap"
|
||||
)
|
||||
await self._ubridge_add_marker_filter(bridge_name, name, bpf, pcap_path, tag)
|
||||
try:
|
||||
await self._ubridge_add_marker_filter(bridge_name, name, bpf, pcap_path, tag)
|
||||
except UbridgeError as e:
|
||||
# Swallow BPF compile errors (warn + skip) so a single bad
|
||||
# expression can't break link creation / node restart — mirrors
|
||||
# _ubridge_apply_filters, which does the same for packet filters.
|
||||
if "syntax error" in str(e).lower() or "compile filter" in str(e).lower():
|
||||
message = f"Warning: ignoring marker '{name}' due to BPF syntax error: {e}"
|
||||
log.warning(message)
|
||||
self.project.emit("log.warning", {"message": message})
|
||||
continue
|
||||
raise
|
||||
manager.register(
|
||||
str(self.project.id), self._id, name, link_id, tag
|
||||
)
|
||||
|
||||
@ -1440,44 +1440,6 @@ 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
|
||||
|
||||
@ -1697,44 +1697,6 @@ 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
|
||||
|
||||
@ -512,42 +512,6 @@ 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.
|
||||
|
||||
@ -41,6 +41,7 @@ from ..config import Config
|
||||
from ..utils.path import check_path_allowed, get_default_project_directory
|
||||
from ..utils.application_id import get_next_application_id
|
||||
from ..utils.asyncio.pool import Pool
|
||||
from ..utils.packet_filter_validation import validate_bpf_syntax
|
||||
from ..utils.asyncio import locking
|
||||
from ..utils.asyncio import aiozipstream
|
||||
from ..utils.asyncio import wait_run_in_executor
|
||||
@ -765,22 +766,31 @@ class Project:
|
||||
"Dropping invalid filters on link %s: %s",
|
||||
link_data.get("link_id"), e
|
||||
)
|
||||
# Restore traffic-insight markers. Each marker's capture side is resolved
|
||||
# when the link is (re)created and the marker is applied to uBridge via
|
||||
# _ubridge_apply_markers in add_ubridge_udp_connection.
|
||||
# Restore traffic-insight markers directly into link state (mirrors how
|
||||
# filters are restored via update_filters). The capture_node_id persisted
|
||||
# last time is reused for NIO routing; no side resolution is possible here
|
||||
# because the link's nodes are added later. The marker is applied to
|
||||
# uBridge by _ubridge_apply_markers when create() runs. Invalid BPF is
|
||||
# dropped (like invalid filters).
|
||||
for name, marker in (link_data.get("markers") or {}).items():
|
||||
try:
|
||||
await link.start_marker(
|
||||
name=name,
|
||||
bpf=marker["bpf"],
|
||||
tag=marker.get("tag"),
|
||||
color=marker.get("color"),
|
||||
)
|
||||
except (ControllerError, KeyError) as e:
|
||||
bpf = marker.get("bpf")
|
||||
if not bpf:
|
||||
log.warning("Dropping marker %s on link %s: missing bpf", name, link_data.get("link_id"))
|
||||
continue
|
||||
result = validate_bpf_syntax(bpf)
|
||||
if not result.get("valid"):
|
||||
log.warning(
|
||||
"Dropping marker %s on link %s: %s",
|
||||
name, link_data.get("link_id"), e
|
||||
"Dropping marker %s on link %s: invalid BPF (%s)",
|
||||
name, link_data.get("link_id"), result.get("error")
|
||||
)
|
||||
continue
|
||||
link._markers[name] = {
|
||||
"bpf": bpf,
|
||||
"tag": marker.get("tag"),
|
||||
"enabled": marker.get("enabled", True),
|
||||
"color": marker.get("color"),
|
||||
"capture_node_id": marker.get("capture_node_id"),
|
||||
}
|
||||
if "link_style" in link_data:
|
||||
await link.update_link_style(link_data["link_style"])
|
||||
if "show_filters_icon" in link_data:
|
||||
|
||||
@ -35,8 +35,6 @@ 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):
|
||||
@ -321,24 +319,16 @@ class UDPLink(Link):
|
||||
# explicitly deletes a marker via the REST API, and a marker is torn
|
||||
# down automatically only when its link is deleted.
|
||||
|
||||
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, color=None):
|
||||
"""
|
||||
Attach a traffic-insight marker to this link.
|
||||
|
||||
State-only model (mirrors ``update_filters``): record the marker in
|
||||
``_markers`` (with its capture-side node id for NIO routing), then push
|
||||
via ``self.update()`` so it rides the NIO and is applied by
|
||||
``_ubridge_apply_markers``. No dedicated uBridge round-trip — exactly
|
||||
how packet filters are applied.
|
||||
|
||||
:param name: stable filter name — echoed in MARK signals + pcap identity
|
||||
:param bpf: libpcap BPF expression
|
||||
:param tag: optional correlation id
|
||||
@ -349,28 +339,20 @@ class UDPLink(Link):
|
||||
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')}")
|
||||
|
||||
marker_side = self._choose_marker_side()
|
||||
# Record state + runtime capture-side ref unconditionally (so stop/update
|
||||
# work and the marker is persisted), but only push to uBridge when the
|
||||
# link is already live. During project load the link is not yet created
|
||||
# (self._created is False); the marker then rides the NIO via create()
|
||||
# and is applied once by _ubridge_apply_markers — mirroring exactly how
|
||||
# update_filters guards its update() call.
|
||||
self._store_capture_node_for_marker(name, marker_side)
|
||||
self._markers[name].update({"bpf": bpf, "tag": tag, "enabled": True, "color": color})
|
||||
self._markers[name] = {
|
||||
"bpf": bpf,
|
||||
"tag": tag,
|
||||
"enabled": True,
|
||||
"color": color,
|
||||
"capture_node_id": marker_side["node"].id,
|
||||
}
|
||||
if self._created:
|
||||
data = {"name": name, "bpf": bpf, "tag": tag, "link_id": self._id}
|
||||
await marker_side["node"].post(
|
||||
"/adapters/{adapter_number}/ports/{port_number}/markers/start".format(
|
||||
adapter_number=marker_side["adapter_number"], port_number=marker_side["port_number"]
|
||||
),
|
||||
data=data,
|
||||
)
|
||||
await self.update()
|
||||
self._project.emit_notification("link.updated", self.asdict())
|
||||
self._project.dump()
|
||||
|
||||
@ -378,30 +360,27 @@ class UDPLink(Link):
|
||||
"""
|
||||
Remove a traffic-insight marker from this link.
|
||||
|
||||
Drop it from ``_markers`` and push via ``self.update()``: the NIO
|
||||
reset+reapply in ``_ubridge_apply_filters``/``_ubridge_apply_markers``
|
||||
drops it from uBridge. Mirrors how deleting a packet filter works.
|
||||
|
||||
: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)
|
||||
del self._markers[name]
|
||||
if self._created:
|
||||
await self.update()
|
||||
self._project.emit_notification("link.updated", self.asdict())
|
||||
self._project.dump()
|
||||
|
||||
async def update_marker(self, name, bpf=None, tag=None, enabled=None, color=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.
|
||||
Update an existing marker's BPF/tag/enabled/color. Any change pushes via
|
||||
``self.update()``; uBridge picks up the new params on the next NIO
|
||||
reset+reapply (same as packet filters).
|
||||
|
||||
:param name: filter name to update
|
||||
:param bpf: new BPF expression (None = keep existing)
|
||||
@ -414,54 +393,19 @@ class UDPLink(Link):
|
||||
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)
|
||||
new_color = color if color is not None else marker_info.get("color")
|
||||
|
||||
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, "color": new_color,
|
||||
}
|
||||
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, "color": new_color,
|
||||
}
|
||||
elif new_color != marker_info.get("color"):
|
||||
# Color-only change: no ubridge round-trip, just update state.
|
||||
self._markers[name] = {**marker_info, "color": new_color}
|
||||
if bpf is not None and bpf != marker_info["bpf"]:
|
||||
result = validate_bpf_syntax(bpf)
|
||||
if not result.get("valid"):
|
||||
raise ControllerError(f"Invalid BPF expression: {result.get('error', 'unknown error')}")
|
||||
marker_info["bpf"] = bpf
|
||||
if tag is not None:
|
||||
marker_info["tag"] = tag
|
||||
if enabled is not None:
|
||||
marker_info["enabled"] = enabled
|
||||
if color is not None:
|
||||
marker_info["color"] = color
|
||||
|
||||
if self._created:
|
||||
await self.update()
|
||||
self._project.emit_notification("link.updated", self.asdict())
|
||||
self._project.dump()
|
||||
|
||||
@ -20,7 +20,7 @@ from .common import ErrorMessage
|
||||
from .version import Version
|
||||
|
||||
# Controller schemas
|
||||
from .controller.links import LinkCreate, LinkUpdate, Link, UDPPortInfo, EthernetPortInfo, LinkCapture, MarkerCreate, MarkerDelete
|
||||
from .controller.links import LinkCreate, LinkUpdate, Link, UDPPortInfo, EthernetPortInfo, LinkCapture, MarkerCreate
|
||||
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
|
||||
|
||||
@ -159,9 +159,3 @@ class MarkerCreate(BaseModel):
|
||||
)
|
||||
|
||||
|
||||
class MarkerDelete(BaseModel):
|
||||
"""
|
||||
Body for removing a traffic-insight marker from a link.
|
||||
"""
|
||||
|
||||
name: str
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user