marker: fine-grained filter ops, clean pcap on remove

Deleting or updating a marker no longer triggers a full NIO reapply
(reset_packet_filters + re-add), which closed/reopened every sibling
marker's pcap via uBridge. Instead operate on single filters:

- stop_marker: bridge delete_packet_filter + unlink the pcap (works with
  the node stopped; filter removal is skipped, the file is still deleted).
- update_marker: bpf/tag/direction → rebuild just that filter (delete + add);
  enabled → instant toggle; color/highlight_duration → stored only.
- compute delete_marker_capture / rebuild_marker_filter + per-node routes
  (DELETE /markers/{name}, PUT /markers/{name}/rebuild) + MarkerRebuild schema.

IOU overrides _ubridge_delete_marker_filter for iol_bridge; rebuild reuses
the already-overridden add/delete/set, so IOU needs no rebuild override.
This commit is contained in:
YueGuobin 2026-08-04 21:35:31 +08:00
parent 19815f7a37
commit caec71aa71
No known key found for this signature in database
13 changed files with 456 additions and 37 deletions

View File

@ -294,3 +294,38 @@ async def pause_cloud_markers(node: Cloud = Depends(dep_node)) -> None:
async def resume_cloud_markers(node: Cloud = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT
)
async def delete_cloud_marker_capture(
marker_name: str,
link_id: str = "",
node: Cloud = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put("/{node_id}/markers/{marker_name}/rebuild")
async def rebuild_cloud_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: Cloud = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -450,3 +450,42 @@ async def pause_docker_markers(node: DockerVM = Depends(dep_node)) -> None:
async def resume_docker_markers(node: DockerVM = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT,
dependencies=[Depends(compute_authentication)]
)
async def delete_docker_marker_capture(
marker_name: str,
link_id: str = "",
node: DockerVM = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put(
"/{node_id}/markers/{marker_name}/rebuild",
dependencies=[Depends(compute_authentication)]
)
async def rebuild_docker_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: DockerVM = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -409,3 +409,42 @@ async def pause_dynamips_markers(node: Router = Depends(dep_node)) -> None:
async def resume_dynamips_markers(node: Router = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT,
dependencies=[Depends(compute_authentication)]
)
async def delete_dynamips_marker_capture(
marker_name: str,
link_id: str = "",
node: Router = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put(
"/{node_id}/markers/{marker_name}/rebuild",
dependencies=[Depends(compute_authentication)]
)
async def rebuild_dynamips_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: Router = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -388,3 +388,42 @@ async def pause_iou_markers(node: IOUVM = Depends(dep_node)) -> None:
async def resume_iou_markers(node: IOUVM = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT,
dependencies=[Depends(compute_authentication)]
)
async def delete_iou_marker_capture(
marker_name: str,
link_id: str = "",
node: IOUVM = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put(
"/{node_id}/markers/{marker_name}/rebuild",
dependencies=[Depends(compute_authentication)]
)
async def rebuild_iou_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: IOUVM = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -480,3 +480,42 @@ async def pause_qemu_markers(node: QemuVM = Depends(dep_node)) -> None:
async def resume_qemu_markers(node: QemuVM = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT,
dependencies=[Depends(compute_authentication)]
)
async def delete_qemu_marker_capture(
marker_name: str,
link_id: str = "",
node: QemuVM = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put(
"/{node_id}/markers/{marker_name}/rebuild",
dependencies=[Depends(compute_authentication)]
)
async def rebuild_qemu_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: QemuVM = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -387,3 +387,42 @@ async def pause_vpcs_markers(node: VPCSVM = Depends(dep_node)) -> None:
async def resume_vpcs_markers(node: VPCSVM = Depends(dep_node)) -> None:
await node._ubridge_marker_resume()
@router.delete(
"/{node_id}/markers/{marker_name}",
status_code=status.HTTP_204_NO_CONTENT,
dependencies=[Depends(compute_authentication)]
)
async def delete_vpcs_marker_capture(
marker_name: str,
link_id: str = "",
node: VPCSVM = Depends(dep_node)
) -> None:
"""
Delete a marker's capture pcap (called by the controller when the marker is
removed) so the file is cleaned up even with the node stopped.
"""
await node.delete_marker_capture(marker_name, link_id)
@router.put(
"/{node_id}/markers/{marker_name}/rebuild",
dependencies=[Depends(compute_authentication)]
)
async def rebuild_vpcs_marker(
marker_name: str,
rebuild_data: schemas.MarkerRebuild,
node: VPCSVM = Depends(dep_node)
) -> dict:
"""
Re-install a single marker filter with new BPF/tag/direction (delete + add,
no bridge reset) so sibling markers' pcaps stay open.
"""
await node.rebuild_marker_filter(
marker_name, rebuild_data.link_id, rebuild_data.bpf,
rebuild_data.tag, rebuild_data.direction, rebuild_data.enabled,
)
return {"marker_name": marker_name}

View File

@ -1130,6 +1130,61 @@ class BaseNode:
# bad expression must surface instead of being silently dropped.
await self._ubridge_send(cmd)
async def delete_marker_capture(self, name, link_id):
"""
Remove a marker from uBridge (fine-grained ``delete_packet_filter`` NOT
reset_packet_filters, so sibling markers' pcaps aren't closed/reopened)
and delete its capture pcap. Called by the controller when a marker is
removed; safe with the node stopped (filter removal is skipped, the file
is still unlinked). IOU overrides ``_ubridge_delete_marker_filter`` for
its ``iol_bridge`` command shape.
"""
bridge_name = self._marker_filter_bridges.pop((name, link_id), None)
if bridge_name is not None:
await self._ubridge_delete_marker_filter(bridge_name, name)
try:
markers_dir = self.project.markers_working_directory()
pcap_path = os.path.join(markers_dir, f"{self._id}_{link_id}_{name}.pcap")
os.remove(pcap_path)
except FileNotFoundError:
pass
except OSError as e:
log.warning("Could not remove marker pcap for '%s' on link %s: %s", name, link_id, e)
async def _ubridge_delete_marker_filter(self, bridge_name, name):
"""
Remove a single marker filter from uBridge with ``delete_packet_filter``
(not a bridge-wide reset) so other markers keep their pcaps open. A no-op
when uBridge isn't running — the pcap cleanup in the caller still proceeds.
"""
if not (self._ubridge_hypervisor and self._ubridge_hypervisor.is_running()):
return
try:
await self._ubridge_send(f"bridge delete_packet_filter {bridge_name} {name}")
except UbridgeError as e:
log.warning("Could not remove marker filter '%s' from %s: %s", name, bridge_name, e)
async def rebuild_marker_filter(self, name, link_id, bpf, tag=None, direction=None, enabled=True):
"""
Re-install a single marker filter with new params (delete + add), without
a bridge-wide reset so sibling markers keep their pcaps open. uBridge
reopens the marker's own pcap on re-add (a new capture session for the
new BPF), which is expected. No-op if the marker isn't installed (node
stopped) the next NIO reapply picks up the updated ``_markers``.
IOU needs no override: this calls ``_ubridge_delete_marker_filter`` /
``_ubridge_add_marker_filter`` / ``_ubridge_set_marker_filter_state``,
all of which IOU already overrides for ``iol_bridge``.
"""
bridge_name = self._marker_filter_bridges.get((name, link_id))
if bridge_name is None:
return
await self._ubridge_delete_marker_filter(bridge_name, name)
pcap_path = os.path.join(self.project.markers_working_directory(), f"{self._id}_{link_id}_{name}.pcap")
await self._ubridge_add_marker_filter(bridge_name, name, bpf, pcap_path, tag, link_id, direction=direction)
if not enabled:
await self._ubridge_set_marker_filter_state(name, enabled=False)
async def _ubridge_apply_markers(self, bridge_name, nio):
"""
(Re-)apply every traffic-insight marker carried by *nio* to the uBridge

View File

@ -1340,6 +1340,17 @@ class IOUVM(BaseNode):
if n == name:
await self._ubridge_send(f"iol_bridge enable_packet_filter {location} {name} {state}")
async def _ubridge_delete_marker_filter(self, location, name):
"""IOU override: remove a single marker filter via ``iol_bridge``
(location = ``{bridge} {bay} {unit}``), not a bridge-wide reset."""
if not (self._ubridge_hypervisor and self._ubridge_hypervisor.is_running()):
return
try:
await self._ubridge_send(f"iol_bridge delete_packet_filter {location} {name}")
except UbridgeError as e:
log.warning("Could not remove marker filter '%s' from %s: %s", name, location, e)
async def adapter_remove_nio_binding(self, adapter_number, port_number):
"""
Removes an adapter NIO binding.

View File

@ -427,17 +427,29 @@ class UDPLink(Link):
"Delete or update it via the marker-definitions API instead."
)
capture_node_id = self._markers[name].get("capture_node_id")
del self._markers[name]
if self._created:
await self.update()
# Remove the marker filter + its pcap on the capture node directly — NOT a
# full NIO reapply (which would reset_packet_filters and close/reopen every
# sibling marker's pcap). delete_packet_filter removes just this filter;
# the marker is already gone from _markers, so any later reapply (filter
# change, node restart) won't re-add it either.
if capture_node_id is not None:
side = next((s for s in self._nodes if str(s["node"].id) == str(capture_node_id)), None)
if side is not None:
try:
await side["node"].delete(f"/markers/{name}", params={"link_id": self._id})
except Exception:
pass # best-effort: old compute without the route leaves the file
self._project.emit_notification("link.updated", self.asdict())
self._project.dump()
async def update_marker(self, name, bpf=None, tag=None, enabled=None, direction=_UNSET, color=None, highlight_duration=None, inherited=False):
"""
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).
Update an existing marker's fields and push to uBridge fine-grained — no
full NIO reapply, so sibling markers' pcaps stay open. bpf/tag/direction
rebuild just this filter (delete + add); enabled is an instant toggle;
color/highlight_duration are UI-only (stored, never pushed).
:param name: filter name to update
:param bpf: new BPF expression (None = keep existing)
@ -460,33 +472,7 @@ class UDPLink(Link):
"Update it via the marker-definitions API instead."
)
# Instant toggle: when only `enabled` changes, send a single
# enable_packet_filter on|off to the capture node instead of rebuilding
# the whole NIO (no pcap flush, emitted counter preserved). Falls through
# to the full reset+reapply below if the compute route is unavailable.
only_enabled = (
enabled is not None
and bpf is None
and tag is None
and direction is _UNSET
and color is None
and highlight_duration is None
)
if only_enabled and self._created:
capture_node_id = marker_info.get("capture_node_id")
side = next((s for s in self._nodes if str(s["node"].id) == str(capture_node_id)), None)
if side is not None:
try:
await side["node"].put(f"/markers/{name}", data={"enabled": enabled})
marker_info["enabled"] = enabled
self._project.emit_notification("link.updated", self.asdict())
self._project.dump()
return
except Exception:
# Old compute without the toggle route / node down: fall
# through to the full NIO reset+reapply below.
pass
# Merge every changed field into the marker state first.
if bpf is not None and bpf != marker_info["bpf"]:
result = validate_bpf_syntax(bpf)
if not result.get("valid"):
@ -503,7 +489,34 @@ class UDPLink(Link):
if direction is not _UNSET:
marker_info["direction"] = direction # None = clear back to both directions
# Push to uBridge fine-grained — NO full NIO reapply (which would
# reset_packet_filters and close/reopen every sibling marker's pcap):
# * bpf/tag/direction changed → rebuild just this filter (delete + add),
# reopening only this marker's pcap (expected, new BPF)
# * only enabled changed → instant toggle (enable_packet_filter)
# * only UI fields changed → nothing to push to uBridge
if self._created:
await self.update()
ubridge_rebuild = (bpf is not None) or (tag is not None) or (direction is not _UNSET)
capture_node_id = marker_info.get("capture_node_id")
side = next((s for s in self._nodes if str(s["node"].id) == str(capture_node_id)), None)
if side is not None:
try:
if ubridge_rebuild:
await side["node"].put(
f"/markers/{name}/rebuild",
data={
"bpf": marker_info["bpf"],
"tag": marker_info.get("tag"),
"direction": marker_info.get("direction"),
"enabled": marker_info.get("enabled", True),
"link_id": self._id,
},
)
elif enabled is not None:
await side["node"].put(f"/markers/{name}", data={"enabled": enabled})
except Exception:
# Old compute without the route / node down: state is already
# correct in _markers; the next NIO reapply converges uBridge.
pass
self._project.emit_notification("link.updated", self.asdict())
self._project.dump()

View File

@ -91,7 +91,7 @@ from .controller.templates.dynamips_templates import (
)
# Compute schemas
from .compute.nios import UDPNIO, TAPNIO, EthernetNIO, MarkerToggle
from .compute.nios import UDPNIO, TAPNIO, EthernetNIO, MarkerToggle, MarkerRebuild
from .compute.atm_switch_nodes import ATMSwitchCreate, ATMSwitchUpdate, ATMSwitch
from .compute.cloud_nodes import CloudCreate, CloudUpdate, Cloud
from .compute.docker_nodes import DockerCreate, DockerUpdate, Docker

View File

@ -75,3 +75,19 @@ class MarkerToggle(BaseModel):
"""
enabled: bool
class MarkerRebuild(BaseModel):
"""
Body for the per-marker rebuild endpoint: re-install a single uBridge marker
filter with new BPF/tag/direction via ``delete_packet_filter`` + add (NOT a
bridge-wide reset), so sibling markers keep their pcaps open. The marker's
own pcap is reopened by uBridge on re-add (new capture session for the new
BPF), which is expected.
"""
bpf: str
tag: Optional[int] = None
direction: Optional[str] = None
enabled: bool = True
link_id: str = ""

View File

@ -15,6 +15,7 @@
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
import os
from collections import OrderedDict
import pytest
@ -227,3 +228,64 @@ async def test_apply_markers_turns_disabled_filter_off(compute_project, manager)
await node._ubridge_apply_markers("VPCS-10", nio)
node._ubridge_send.assert_any_call("bridge enable_packet_filter VPCS-10 m off")
assert node._marker_filter_bridges["m", "L1"] == "VPCS-10"
@pytest.mark.asyncio
async def test_delete_marker_capture_removes_pcap_and_entry(compute_project, manager):
# Deleting a marker's capture removes its pcap and forgets the bridge entry.
node = VPCSVM("test", "00010203-0405-0607-0809-0a0b0c0d0e0f", compute_project, manager)
markers_dir = compute_project.markers_working_directory()
os.makedirs(markers_dir, exist_ok=True)
node._marker_filter_bridges["m", "L1"] = "VPCS-10"
pcap = os.path.join(markers_dir, f"{node.id}_L1_m.pcap")
open(pcap, "wb").write(b"data")
await node.delete_marker_capture("m", "L1")
assert not os.path.exists(pcap)
assert ("m", "L1") not in node._marker_filter_bridges
@pytest.mark.asyncio
async def test_delete_marker_capture_idempotent_when_missing(compute_project, manager):
# No file on disk → must not raise, and still clears the entry.
node = VPCSVM("test", "00010203-0405-0607-0809-0a0b0c0d0e0f", compute_project, manager)
node._marker_filter_bridges["m", "L1"] = "VPCS-10"
await node.delete_marker_capture("m", "L1")
assert ("m", "L1") not in node._marker_filter_bridges
@pytest.mark.asyncio
async def test_delete_marker_capture_sends_delete_filter(compute_project, manager):
# With uBridge running, removing a marker issues a fine-grained
# delete_packet_filter (not a bridge-wide reset) so sibling pcaps stay open.
node = VPCSVM("test", "00010203-0405-0607-0809-0a0b0c0d0e0f", compute_project, manager)
node._ubridge_send = AsyncioMagicMock()
node._ubridge_hypervisor = MagicMock()
node._ubridge_hypervisor.is_running.return_value = True
node._marker_filter_bridges["m", "L1"] = "VPCS-10"
await node.delete_marker_capture("m", "L1")
node._ubridge_send.assert_any_call("bridge delete_packet_filter VPCS-10 m")
assert ("m", "L1") not in node._marker_filter_bridges
@pytest.mark.asyncio
async def test_rebuild_marker_filter_delete_then_add(compute_project, manager):
# rebuild re-installs a single filter (delete_packet_filter + add) with the
# new params, no bridge reset; enabled=False turns it off after re-add.
node = VPCSVM("test", "00010203-0405-0607-0809-0a0b0c0d0e0f", compute_project, manager)
node._ubridge_send = AsyncioMagicMock()
node._ubridge_hypervisor = MagicMock()
node._ubridge_hypervisor.is_running.return_value = True
node._marker_filter_bridges["m", "L1"] = "VPCS-10"
await node.rebuild_marker_filter("m", "L1", "tcp", tag=7, direction="rx", enabled=False)
cmds = [c.args[0] for c in node._ubridge_send.call_args_list]
assert any("delete_packet_filter VPCS-10 m" in c for c in cmds)
assert any("add_packet_filter VPCS-10 m mark" in c and "tcp" in c for c in cmds)
assert any("enable_packet_filter VPCS-10 m off" in c for c in cmds)

View File

@ -578,8 +578,9 @@ async def test_update_marker_enabled_only_hits_toggle_route(project):
@pytest.mark.asyncio
async def test_update_marker_with_bpf_still_rebuilds_nio(project):
# A non-enabled-only change falls through to the NIO reset+reapply path.
async def test_update_marker_with_bpf_rebuilds_single_filter(project):
# A bpf change rebuilds just this marker's filter (delete + add), NOT a full
# NIO reapply, so sibling markers' pcaps stay open.
with _valid_bpf():
link = await _make_link(project)
node = link._nodes[0]["node"]
@ -588,7 +589,23 @@ async def test_update_marker_with_bpf_still_rebuilds_nio(project):
compute.put.reset_mock()
await link.update_marker("m", bpf="tcp")
paths = [c.args[0] for c in compute.put.call_args_list]
assert any(p.endswith("/nio") for p in paths)
assert any(p.endswith("/markers/m/rebuild") for p in paths)
assert not any(p.endswith("/nio") for p in paths) # no full NIO reapply
@pytest.mark.asyncio
async def test_update_marker_ui_only_does_not_push(project):
# color/highlight_duration are UI-only — stored, never pushed to uBridge.
with _valid_bpf():
link = await _make_link(project)
node = link._nodes[0]["node"]
await link.start_marker("m", "icmp")
compute = node.compute
compute.put.reset_mock()
await link.update_marker("m", color="#ffffff", highlight_duration=1500)
assert compute.put.call_args_list == [] # nothing pushed to uBridge
assert link.markers["m"]["color"] == "#ffffff"
assert link.markers["m"]["highlight_duration"] == 1500
# ---------------------------------------------------------------------------
@ -663,3 +680,18 @@ async def test_update_marker_definition_rejects_directional(project):
await project.update_marker_definition("arp", color="#ffffff")
await project.update_marker_definition("arp", direction=None)
assert project.marker_definitions["arp"]["direction"] is None
@pytest.mark.asyncio
async def test_stop_marker_deletes_capture_pcap(project):
# Removing a marker asks the capture node's compute to delete its pcap, so
# the file is cleaned up even with the node stopped (the NIO reapply path
# only runs while uBridge is up).
with _valid_bpf():
link = await _make_link(project)
capture = link._nodes[0]["node"]
await link.start_marker("icmp", "icmp", capture_node_id=capture.id)
capture.delete = AsyncioMagicMock()
await link.stop_marker("icmp")
capture.delete.assert_called_once_with("/markers/icmp", params={"link_id": link.id})