benchmark mode added test
This commit is contained in:
@@ -35,6 +35,13 @@ class SnifferStartRequest(BaseModel):
|
||||
example="tc_ebpf",
|
||||
description="Bridge capture mode: 'tc_ebpf' or 'af_packet'. Ignored for interface capture.",
|
||||
)
|
||||
benchmark_mode: bool = Field(
|
||||
False,
|
||||
description=(
|
||||
"Run capture hooks for measurement while skipping packet parsing, DB persistence, "
|
||||
"DPI enrichment, and live packet publication."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
class SnifferStartResponse(BaseModel):
|
||||
@@ -45,6 +52,7 @@ class SnifferStartResponse(BaseModel):
|
||||
target: str = Field(..., description="Started target name.")
|
||||
target_type: str = Field(..., description="Either 'bridge' or 'interface'.")
|
||||
capture_mode: str = Field(..., description="The effective capture mode used by the session.")
|
||||
benchmark_mode: bool = Field(False, description="Whether userspace packet processing is skipped.")
|
||||
|
||||
|
||||
class SnifferStopRequest(BaseModel):
|
||||
@@ -71,6 +79,7 @@ class InterfaceSnifferStatus(BaseModel):
|
||||
session_id: Optional[str] = Field(None, description="Owning capture session ID.")
|
||||
session_label: Optional[str] = Field(None, description="Human-readable session label.")
|
||||
capture_mode: Optional[str] = Field(None, description="Capture mode used by the owning session.")
|
||||
benchmark_mode: Optional[bool] = Field(None, description="Whether the owning session skips userspace processing.")
|
||||
|
||||
|
||||
class SnifferStatusResponse(BaseModel):
|
||||
@@ -90,13 +99,18 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse:
|
||||
|
||||
try:
|
||||
if req.interface:
|
||||
session_id = start_capture_session(req.interface, target_is_interface=True)
|
||||
session_id = start_capture_session(
|
||||
req.interface,
|
||||
target_is_interface=True,
|
||||
benchmark_mode=req.benchmark_mode,
|
||||
)
|
||||
return SnifferStartResponse(
|
||||
started=True,
|
||||
session_id=session_id,
|
||||
target=req.interface,
|
||||
target_type="interface",
|
||||
capture_mode="af_packet",
|
||||
benchmark_mode=req.benchmark_mode,
|
||||
)
|
||||
|
||||
effective_capture_mode = req.bridge_capture_mode or BRIDGE_CAPTURE_MODE_TC_EBPF
|
||||
@@ -106,6 +120,7 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse:
|
||||
req.bridge,
|
||||
target_is_interface=False,
|
||||
bridge_capture_mode=effective_capture_mode,
|
||||
benchmark_mode=req.benchmark_mode,
|
||||
)
|
||||
return SnifferStartResponse(
|
||||
started=True,
|
||||
@@ -113,6 +128,7 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse:
|
||||
target=req.bridge,
|
||||
target_type="bridge",
|
||||
capture_mode=effective_capture_mode,
|
||||
benchmark_mode=req.benchmark_mode,
|
||||
)
|
||||
except Exception as exc:
|
||||
raise HTTPException(status_code=500, detail=f"Failed to start sniffer: {exc}") from exc
|
||||
|
||||
@@ -630,8 +630,16 @@ def _sync_bridge_telemetry() -> None:
|
||||
for session_id, session in sessions.items()
|
||||
if session.get("is_bridge") and session.get("capture_mode") == BRIDGE_CAPTURE_MODE_TC_EBPF
|
||||
}
|
||||
benchmark_session_ids = {
|
||||
session_id
|
||||
for session_id, session in sessions.items()
|
||||
if session.get("benchmark_mode")
|
||||
}
|
||||
try:
|
||||
bridge_telemetry_manager.update_sessions(bridge_session_interfaces)
|
||||
bridge_telemetry_manager.update_sessions(
|
||||
bridge_session_interfaces,
|
||||
benchmark_session_ids=benchmark_session_ids,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to update bridge telemetry collector")
|
||||
|
||||
@@ -639,6 +647,7 @@ def _sync_bridge_telemetry() -> None:
|
||||
{
|
||||
iface
|
||||
for session in sessions.values()
|
||||
if not session.get("benchmark_mode")
|
||||
for iface in (
|
||||
list(session.get("capture_ifaces", []))
|
||||
+ (
|
||||
@@ -667,8 +676,9 @@ def _session_reader_loop(session_id: str) -> None:
|
||||
stop_event: threading.Event = session["stop_event"]
|
||||
sockets: Dict[str, socket.socket] = session["sockets"]
|
||||
label: str = session["label"]
|
||||
benchmark_mode = bool(session.get("benchmark_mode"))
|
||||
|
||||
logger.info("Session %s reader starting (label=%s)", session_id, label)
|
||||
logger.info("Session %s reader starting (label=%s benchmark_mode=%s)", session_id, label, benchmark_mode)
|
||||
sel = selectors.DefaultSelector()
|
||||
|
||||
# register existing sockets
|
||||
@@ -729,6 +739,11 @@ def _session_reader_loop(session_id: str) -> None:
|
||||
logger.exception("Recv error on %s in session %s", iface, session_id)
|
||||
continue
|
||||
|
||||
if benchmark_mode:
|
||||
session["benchmark_packets"] = int(session.get("benchmark_packets") or 0) + 1
|
||||
session["benchmark_bytes"] = int(session.get("benchmark_bytes") or 0) + len(raw)
|
||||
continue
|
||||
|
||||
# parse with scapy
|
||||
try:
|
||||
recv_ts = datetime.now(timezone.utc)
|
||||
@@ -775,6 +790,7 @@ def start_capture_session(
|
||||
target: str,
|
||||
target_is_interface: bool = False,
|
||||
bridge_capture_mode: str = BRIDGE_CAPTURE_MODE_TC_EBPF,
|
||||
benchmark_mode: bool = False,
|
||||
) -> str:
|
||||
"""
|
||||
Start a packet capture session. Returns session_id string.
|
||||
@@ -796,6 +812,9 @@ def start_capture_session(
|
||||
"label": target,
|
||||
"is_bridge": not target_is_interface,
|
||||
"capture_mode": effective_capture_mode,
|
||||
"benchmark_mode": bool(benchmark_mode),
|
||||
"benchmark_packets": 0,
|
||||
"benchmark_bytes": 0,
|
||||
"ports": [],
|
||||
"capture_ifaces": [],
|
||||
}
|
||||
@@ -834,10 +853,11 @@ def start_capture_session(
|
||||
session["thread"] = None
|
||||
_sync_bridge_telemetry()
|
||||
logger.info(
|
||||
"Started capture session %s label=%s capture_mode=%s ports=%s capture_ifaces=%s",
|
||||
"Started capture session %s label=%s capture_mode=%s benchmark_mode=%s ports=%s capture_ifaces=%s",
|
||||
session_id,
|
||||
target,
|
||||
effective_capture_mode,
|
||||
bool(benchmark_mode),
|
||||
ports,
|
||||
capture_ifaces,
|
||||
)
|
||||
@@ -946,6 +966,7 @@ def get_capture_session_status() -> Dict[str, Dict[str, object]]:
|
||||
"session_id": sid,
|
||||
"session_label": s.get("label"),
|
||||
"capture_mode": s.get("capture_mode"),
|
||||
"benchmark_mode": bool(s.get("benchmark_mode")),
|
||||
}
|
||||
if not s.get("sockets") and s.get("is_bridge"):
|
||||
for iface in s.get("ports", []):
|
||||
@@ -956,6 +977,7 @@ def get_capture_session_status() -> Dict[str, Dict[str, object]]:
|
||||
"session_id": sid,
|
||||
"session_label": s.get("label"),
|
||||
"capture_mode": s.get("capture_mode"),
|
||||
"benchmark_mode": bool(s.get("benchmark_mode")),
|
||||
}
|
||||
return out
|
||||
|
||||
@@ -970,6 +992,9 @@ def get_internal_debug_state() -> dict:
|
||||
"label": s.get("label"),
|
||||
"is_bridge": s.get("is_bridge"),
|
||||
"capture_mode": s.get("capture_mode"),
|
||||
"benchmark_mode": bool(s.get("benchmark_mode")),
|
||||
"benchmark_packets": int(s.get("benchmark_packets") or 0),
|
||||
"benchmark_bytes": int(s.get("benchmark_bytes") or 0),
|
||||
"ports": list(s.get("ports", [])),
|
||||
"capture_ifaces": list(s.get("capture_ifaces", [])),
|
||||
"sockets": list(s.get("sockets", {}).keys()),
|
||||
@@ -993,6 +1018,7 @@ def get_internal_debug_state() -> dict:
|
||||
}
|
||||
),
|
||||
"tshark": tshark_manager.get_debug_snapshot(),
|
||||
"bridge_telemetry": bridge_telemetry_manager.get_debug_snapshot(),
|
||||
"packet_tracker": packet_tracker.get_debug_snapshot(),
|
||||
}
|
||||
|
||||
@@ -1001,12 +1027,14 @@ def start_afpacket_sniffer(
|
||||
target: str,
|
||||
target_is_interface: bool = False,
|
||||
bridge_capture_mode: str = BRIDGE_CAPTURE_MODE_TC_EBPF,
|
||||
benchmark_mode: bool = False,
|
||||
) -> str:
|
||||
"""Backward-compatible wrapper for start_capture_session()."""
|
||||
return start_capture_session(
|
||||
target,
|
||||
target_is_interface=target_is_interface,
|
||||
bridge_capture_mode=bridge_capture_mode,
|
||||
benchmark_mode=benchmark_mode,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -27,6 +27,7 @@ class BridgeTelemetryManager:
|
||||
def __init__(self) -> None:
|
||||
self._interfaces: set[str] = set()
|
||||
self._session_ids_by_interface: dict[str, tuple[str, ...]] = {}
|
||||
self._benchmark_by_interface: dict[str, bool] = {}
|
||||
self._process: Optional[subprocess.Popen[str]] = None
|
||||
self._reader_thread: Optional[threading.Thread] = None
|
||||
self._event_queue: queue.Queue = queue.Queue(maxsize=settings.bridge_telemetry_event_queue_maxsize)
|
||||
@@ -38,13 +39,20 @@ class BridgeTelemetryManager:
|
||||
self._worker_thread.start()
|
||||
self._dropped_events = 0
|
||||
self._dropped_raw_payloads = 0
|
||||
self._benchmark_events = 0
|
||||
self._benchmark_raw_payloads = 0
|
||||
self._last_drop_log_at = 0.0
|
||||
self._suppressed_collector_messages = 0
|
||||
self._last_collector_warning_at = 0.0
|
||||
self._lock = threading.Lock()
|
||||
|
||||
def update_sessions(self, session_interfaces: Mapping[str, Iterable[str]]) -> None:
|
||||
def update_sessions(
|
||||
self,
|
||||
session_interfaces: Mapping[str, Iterable[str]],
|
||||
benchmark_session_ids: Optional[set[str]] = None,
|
||||
) -> None:
|
||||
"""Restart the collector when the active bridge interface set changes."""
|
||||
benchmark_session_ids = benchmark_session_ids or set()
|
||||
normalized: dict[str, set[str]] = {}
|
||||
for session_id, interfaces in session_interfaces.items():
|
||||
if not session_id:
|
||||
@@ -58,11 +66,17 @@ class BridgeTelemetryManager:
|
||||
iface: tuple(sorted(session_id for session_id, ifaces in normalized.items() if iface in ifaces))
|
||||
for iface in normalized_interfaces
|
||||
}
|
||||
benchmark_by_interface = {
|
||||
iface: bool(session_ids_by_interface.get(iface))
|
||||
and all(session_id in benchmark_session_ids for session_id in session_ids_by_interface.get(iface, ()))
|
||||
for iface in normalized_interfaces
|
||||
}
|
||||
|
||||
with self._lock:
|
||||
interfaces_changed = set(normalized_interfaces) != self._interfaces
|
||||
self._interfaces = set(normalized_interfaces)
|
||||
self._session_ids_by_interface = session_ids_by_interface
|
||||
self._benchmark_by_interface = benchmark_by_interface
|
||||
if not interfaces_changed:
|
||||
return
|
||||
self._restart_locked()
|
||||
@@ -192,6 +206,17 @@ class BridgeTelemetryManager:
|
||||
if not isinstance(event, dict):
|
||||
continue
|
||||
|
||||
iface = str(event.get("iface") or "")
|
||||
with self._lock:
|
||||
benchmark_mode = bool(self._benchmark_by_interface.get(iface))
|
||||
if benchmark_mode:
|
||||
raw_b64 = event.pop("raw_b64", None)
|
||||
with self._lock:
|
||||
self._benchmark_events += 1
|
||||
if raw_b64:
|
||||
self._benchmark_raw_payloads += 1
|
||||
continue
|
||||
|
||||
if event.get("event_type") == "ingress":
|
||||
self._handle_ingress_packet(event)
|
||||
|
||||
@@ -311,5 +336,17 @@ class BridgeTelemetryManager:
|
||||
if rc not in (0, None):
|
||||
logger.warning("Bridge telemetry collector exited with code %s", rc)
|
||||
|
||||
def get_debug_snapshot(self) -> dict[str, object]:
|
||||
with self._lock:
|
||||
return {
|
||||
"interfaces": sorted(self._interfaces),
|
||||
"benchmark_by_interface": dict(self._benchmark_by_interface),
|
||||
"queue_size": self._event_queue.qsize(),
|
||||
"dropped_events": self._dropped_events,
|
||||
"dropped_raw_payloads": self._dropped_raw_payloads,
|
||||
"benchmark_events": self._benchmark_events,
|
||||
"benchmark_raw_payloads": self._benchmark_raw_payloads,
|
||||
}
|
||||
|
||||
|
||||
bridge_telemetry_manager = BridgeTelemetryManager()
|
||||
|
||||
@@ -448,6 +448,7 @@ class Event(ct.Structure):
|
||||
|
||||
TARGET_INTERFACES: set[str] = set()
|
||||
IPR: IPRoute | None = None
|
||||
OMIT_RAW_PAYLOAD = False
|
||||
|
||||
|
||||
def _run_checked(cmd: list[str]) -> None:
|
||||
@@ -546,7 +547,7 @@ def _emit_ingress_event(cpu: int, data: int, size: int) -> None:
|
||||
return
|
||||
|
||||
raw_size = size - ct.sizeof(Event)
|
||||
if raw_size > 0:
|
||||
if raw_size > 0 and not OMIT_RAW_PAYLOAD:
|
||||
raw = ct.string_at(data + ct.sizeof(Event), min(raw_size, int(event.length)))
|
||||
payload["raw_b64"] = base64.b64encode(raw).decode("ascii")
|
||||
|
||||
@@ -597,6 +598,11 @@ def _parse_args() -> argparse.Namespace:
|
||||
default=128,
|
||||
help="Perf-buffer page count for metadata events.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--omit-raw-payload",
|
||||
action="store_true",
|
||||
help="Drain raw ingress events but do not base64-encode or print raw packet bytes.",
|
||||
)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
@@ -675,6 +681,8 @@ def _cleanup_tc(ifaces: Iterable[str]) -> None:
|
||||
|
||||
def main() -> int:
|
||||
args = _parse_args()
|
||||
global OMIT_RAW_PAYLOAD
|
||||
OMIT_RAW_PAYLOAD = bool(args.omit_raw_payload)
|
||||
raw_sample_every = max(0, args.raw_sample_every)
|
||||
meta_sample_every = max(0, args.meta_sample_every)
|
||||
ingress_pages = max(1, args.ingress_pages)
|
||||
@@ -706,6 +714,7 @@ def main() -> int:
|
||||
"meta_sample_every": meta_sample_every,
|
||||
"ingress_pages": ingress_pages,
|
||||
"meta_pages": meta_pages,
|
||||
"omit_raw_payload": OMIT_RAW_PAYLOAD,
|
||||
},
|
||||
separators=(",", ":"),
|
||||
),
|
||||
|
||||
@@ -17,6 +17,7 @@ import {
|
||||
Select,
|
||||
Space,
|
||||
Spin,
|
||||
Switch,
|
||||
Tag,
|
||||
Tooltip,
|
||||
Typography,
|
||||
@@ -62,8 +63,13 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement
|
||||
setIsModalOpen(false);
|
||||
};
|
||||
|
||||
const handleStartSubmit = async (values: { target?: string; bridgeCaptureMode?: 'tc_ebpf' | 'af_packet' }) => {
|
||||
const handleStartSubmit = async (values: {
|
||||
target?: string;
|
||||
bridgeCaptureMode?: 'tc_ebpf' | 'af_packet';
|
||||
benchmarkMode?: boolean;
|
||||
}) => {
|
||||
const target = values.target;
|
||||
const benchmarkMode = Boolean(values.benchmarkMode);
|
||||
if (!target) {
|
||||
notification.warning({ message: 'Warning', description: 'Please select a target to start capture on.' });
|
||||
return;
|
||||
@@ -72,12 +78,14 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement
|
||||
try {
|
||||
const payload =
|
||||
startMode === 'interface'
|
||||
? { interface: target }
|
||||
: { bridge: target, bridge_capture_mode: values.bridgeCaptureMode ?? bridgeCaptureMode };
|
||||
? { interface: target, benchmark_mode: benchmarkMode }
|
||||
: { bridge: target, bridge_capture_mode: values.bridgeCaptureMode ?? bridgeCaptureMode, benchmark_mode: benchmarkMode };
|
||||
const result = await startSniffer(payload);
|
||||
notification.success({
|
||||
message: 'Capture started',
|
||||
description: `Capture started on ${target} via ${result.capture_mode} (session ${result.session_id})`,
|
||||
description: `Capture started on ${target} via ${result.capture_mode}${
|
||||
result.benchmark_mode ? ' in benchmark mode' : ''
|
||||
} (session ${result.session_id})`,
|
||||
});
|
||||
await props.refreshAll();
|
||||
setIsModalOpen(false);
|
||||
@@ -241,6 +249,7 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement
|
||||
{status.capture_mode === 'af_packet' ? 'AF_PACKET' : 'tc/eBPF'}
|
||||
</Tag>
|
||||
)}
|
||||
{status.benchmark_mode && <Tag color="gold">benchmark</Tag>}
|
||||
</Space>
|
||||
}
|
||||
description={<Text type="secondary">interface: {name}</Text>}
|
||||
@@ -260,7 +269,12 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement
|
||||
confirmLoading={props.loading}
|
||||
okText="Start"
|
||||
>
|
||||
<Form form={form} layout="vertical" onFinish={handleStartSubmit} initialValues={{ target: undefined }}>
|
||||
<Form
|
||||
form={form}
|
||||
layout="vertical"
|
||||
onFinish={handleStartSubmit}
|
||||
initialValues={{ target: undefined, benchmarkMode: false }}
|
||||
>
|
||||
<Form.Item label="Mode" name="mode">
|
||||
<Radio.Group
|
||||
value={startMode}
|
||||
@@ -323,6 +337,19 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement
|
||||
</Radio.Group>
|
||||
</Form.Item>
|
||||
)}
|
||||
|
||||
<Form.Item
|
||||
label="Benchmark Mode"
|
||||
name="benchmarkMode"
|
||||
valuePropName="checked"
|
||||
extra={
|
||||
startMode === 'bridge' && bridgeCaptureMode === 'tc_ebpf'
|
||||
? 'Keeps tc/eBPF raw export and JSON emission, but skips backend parsing, telemetry merge, DB persistence, DPI enrichment, and live packet publishing.'
|
||||
: 'Receives packets for measurement, but skips packet parsing, packet tracking, DB persistence, DPI enrichment, and live packet publishing.'
|
||||
}
|
||||
>
|
||||
<Switch checkedChildren="Benchmark" unCheckedChildren="Normal" />
|
||||
</Form.Item>
|
||||
</Form>
|
||||
</Modal>
|
||||
</div>
|
||||
|
||||
@@ -6,6 +6,7 @@ export interface SnifferStartRequest {
|
||||
bridge?: string;
|
||||
interface?: string;
|
||||
bridge_capture_mode?: 'tc_ebpf' | 'af_packet';
|
||||
benchmark_mode?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -17,6 +18,7 @@ export interface SnifferStartResponse {
|
||||
target: string;
|
||||
target_type: 'bridge' | 'interface';
|
||||
capture_mode: 'tc_ebpf' | 'af_packet';
|
||||
benchmark_mode?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -47,6 +49,7 @@ export interface InterfaceSnifferStatus {
|
||||
session_id?: string | null;
|
||||
session_label?: string | null;
|
||||
capture_mode?: 'tc_ebpf' | 'af_packet' | null;
|
||||
benchmark_mode?: boolean | null;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user