diff --git a/.env.example b/.env.example index e954e1a..970aea3 100644 --- a/.env.example +++ b/.env.example @@ -32,6 +32,15 @@ NETBOT_DATA_DIR= NETBOT_LOG_DIR= NETBOT_LOG_LEVEL=info +# Flow-aware processing lanes. Metrics contain counters and safe reasons only. +NETBOT_FLOW_WORKERS_ENABLED=true +NETBOT_FLOW_WORKER_COUNT=4 +NETBOT_FLOW_WORKER_QUEUE_MAX=2000 +NETBOT_FLOW_WORKER_OVERFLOW_POLICY=drop_oldest +NETBOT_FLOW_WORKER_SHUTDOWN_TIMEOUT_SEC=5 +NETBOT_FLOW_WORKER_ERROR_THRESHOLD=25 +NETBOT_FLOW_WORKER_SLOW_JOB_MS=100 + # Bounded batch persistence. Counters and health are exposed in Ops Snapshot. NETBOT_PERSISTENCE_BATCH_ENABLED=true NETBOT_PERSISTENCE_PACKET_BATCH_SIZE=500 diff --git a/README.md b/README.md index 77c5267..4acbe41 100644 --- a/README.md +++ b/README.md @@ -69,9 +69,10 @@ or PCAP artifacts. | Deep Packet Inspection | Active MVP | Inspect renders a searchable layer tree, safe bytes view, streams, and Expert Info. | No TLS decryption; visible metadata and ASCII previews are centrally redacted. | | Display Filters | Active MVP | Safe packet and flow filter parser covers text, equality, range, and boolean operators. | Filters run on redacted metadata and never use Python `eval`. | | Offline PCAP Deep Analysis | Active MVP | Offline results include packet details, Expert Info, and stream summaries. | Previous API fields remain compatible; raw secrets are not exposed. | -| Bounded Packet Intake Queue | Foundation step | Queue pressure metrics, accepted/drop counters, overflow policies, worker liveness, high-water mark, and Ops Snapshot packet queue visibility are tested. | First engine-level performance hardening step; the Worker Pool is not implemented yet. | +| Bounded Packet Intake Queue | Foundation step | Queue pressure metrics, accepted/drop counters, overflow policies, worker liveness, high-water mark, and Ops Snapshot packet queue visibility are tested. | First engine-level intake hardening step; it is separate from processing workers. | +| Flow-aware Worker Pool | Foundation step | Stable per-flow lanes preserve ordering while different flows process concurrently; bounded queue pressure, latency, failures, and worker health are visible in Ops Snapshot. | This is not a benchmark claim or the complete performance engine. | | Batch Persistence / Storage Backpressure | Foundation step | Redacted packet, alert, and flow records use bounded batches with queue health, write latency, retry, backlog, failure, and Ops Snapshot visibility. | Audit and report exports remain synchronous; this is not a distributed storage engine or the complete performance pipeline. | -| WebSocket Event Aggregator | Foundation step | Realtime packet/alert batching, slow-client protection, WebSocket pressure metrics, and Ops Snapshot visibility are tested. | Realtime delivery batching only; the Worker Pool is not implemented yet. | +| WebSocket Event Aggregator | Foundation step | Realtime packet/alert batching, slow-client protection, WebSocket pressure metrics, and Ops Snapshot visibility are tested. | Realtime delivery batching only. | | Demo and operational QA | Ready | Token-safe demo and Agent script behavior is tested. | Demo launchers and status commands do not print raw tokens. | | Windows release packaging | Validated path | Desktop smoke, version consistency, and release workflow checks run in CI. | Versioned artifacts include SHA256 checksums. | | Linux desktop packaging | Staged | Build workflow exists; native production validation remains pending. | Publish only after native smoke and release QA. | @@ -174,13 +175,18 @@ runtime pressure are surfaced without exposing packet payloads or secrets. The bounded queue is the first performance-pipeline foundation step, not the complete performance pipeline. It adds queue pressure metrics, drop counters, overflow policies, worker liveness, and Ops Snapshot visibility before future -WebSocket batching, batch persistence, and worker-pool work. +WebSocket batching, batch persistence, and flow-aware worker lanes. -The WebSocket Event Aggregator is the next foundation step. It batches realtime +The WebSocket Event Aggregator batches realtime packet and alert updates, coalesces summary updates, protects against slow clients with a bounded outgoing queue, and exposes WebSocket pressure metrics in Ops Snapshot. It is still not the complete performance pipeline. +The Flow-aware Worker Pool sits after bounded packet intake. Stable +bidirectional flow keys select FIFO worker lanes, preserving order within a flow +while allowing different flows to process concurrently. It exposes bounded +backlog, drops/rejections, failures, worker liveness, and processing latency. + Packet intake queue tuning is controlled by: - `NETBOT_PACKET_QUEUE_MAX_SIZE`: default `2000`; use `1000` for small/local @@ -188,7 +194,17 @@ Packet intake queue tuning is controlled by: - `NETBOT_PACKET_QUEUE_OVERFLOW_POLICY`: default `drop_oldest`; allowed values are `drop_oldest` and `drop_newest`. - `NETBOT_PACKET_QUEUE_DRAIN_TIMEOUT_SEC`: default `5.0`; increase to `10.0` - when heavier capture should get more shutdown drain time. + for a longer graceful intake drain. +- `NETBOT_FLOW_WORKERS_ENABLED`: default `true`; set `false` for the existing + direct packet-processing fallback. +- `NETBOT_FLOW_WORKER_COUNT`: default `4` fixed worker lanes. +- `NETBOT_FLOW_WORKER_QUEUE_MAX`: default `2000` total pending processing jobs. +- `NETBOT_FLOW_WORKER_OVERFLOW_POLICY`: default `drop_oldest`; also supports + `drop_newest`, `reject_new`, and bounded `block_short`. +- `NETBOT_FLOW_WORKER_SHUTDOWN_TIMEOUT_SEC`: default `5` seconds. +- `NETBOT_FLOW_WORKER_ERROR_THRESHOLD`: default `25` failures or drops before + critical health. +- `NETBOT_FLOW_WORKER_SLOW_JOB_MS`: default `100` milliseconds. - `NETBOT_PERSISTENCE_BATCH_ENABLED`: batching toggle; default `true`. `false` uses compatible synchronous writes. - `NETBOT_PERSISTENCE_PACKET_BATCH_SIZE` / `PACKET_FLUSH_MS`: `500` / `1000`. diff --git a/backend/app/main.py b/backend/app/main.py index 20a2c70..665ade4 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -117,6 +117,7 @@ def _observability_snapshot() -> dict[str, Any]: "websocket": event_bus.websocket_stats(), "history": history_service.metrics(), "packet_queue": sniffer_service.packet_queue_stats(), + "flow_worker_pool": sniffer_service.flow_worker_pool_stats(), "persistence": sniffer_service.persistence_stats(), "auto_block": sniffer_service.auto_block_stats(), } diff --git a/backend/app/services/flow_worker_pool.py b/backend/app/services/flow_worker_pool.py new file mode 100644 index 0000000..cda7a39 --- /dev/null +++ b/backend/app/services/flow_worker_pool.py @@ -0,0 +1,512 @@ +from __future__ import annotations + +import hashlib +import logging +import os +import queue +import threading +import time +import uuid +from collections import deque +from dataclasses import dataclass, field +from datetime import datetime, timezone +from typing import Any, Callable + +logger = logging.getLogger(__name__) + +DEFAULT_WORKER_COUNT = 4 +DEFAULT_QUEUE_MAX = 2000 +DEFAULT_SHUTDOWN_TIMEOUT_SEC = 5.0 +DEFAULT_ERROR_THRESHOLD = 25 +DEFAULT_SLOW_JOB_MS = 100.0 +VALID_OVERFLOW_POLICIES = { + "drop_oldest", + "drop_newest", + "reject_new", + "block_short", +} + + +def _env_bool(name: str, default: bool) -> bool: + value = str(os.environ.get(name, str(default))).strip().lower() + if value in {"1", "true", "yes", "on"}: + return True + if value in {"0", "false", "no", "off"}: + return False + return default + + +def _env_int(name: str, default: int, minimum: int, maximum: int) -> int: + try: + value = int(os.environ.get(name, str(default))) + except (TypeError, ValueError): + return default + return min(maximum, max(minimum, value)) + + +def _env_float(name: str, default: float, minimum: float, maximum: float) -> float: + try: + value = float(os.environ.get(name, str(default))) + except (TypeError, ValueError): + return default + return min(maximum, max(minimum, value)) + + +def _percentile(values: list[float], percentile: float) -> float: + if not values: + return 0.0 + ordered = sorted(values) + index = max(0, min(len(ordered) - 1, int((len(ordered) - 1) * percentile))) + return round(ordered[index], 2) + + +@dataclass(frozen=True) +class FlowDispatchKey: + value: str + complete: bool + + @classmethod + def from_packet(cls, packet: dict[str, Any]) -> FlowDispatchKey: + src = str(packet.get("src") or packet.get("src_ip") or "").strip() + dst = str(packet.get("dst") or packet.get("dst_ip") or "").strip() + protocol = ( + str( + packet.get("proto") + or packet.get("transport") + or packet.get("protocol") + or "" + ) + .strip() + .upper() + ) + src_port = _safe_port(packet.get("sport", packet.get("src_port"))) + dst_port = _safe_port(packet.get("dport", packet.get("dst_port"))) + + if src and dst and protocol: + if protocol in {"TCP", "UDP"}: + endpoints = sorted(((src, src_port), (dst, dst_port))) + left, right = endpoints + return cls( + f"{protocol}|{left[0]}:{left[1]}|{right[0]}:{right[1]}", + True, + ) + endpoints = sorted((src, dst)) + return cls(f"{protocol}|{endpoints[0]}|{endpoints[1]}", True) + + parts = [ + protocol or "OTHER", + src or "-", + str(src_port), + dst or "-", + str(dst_port), + ] + if any(value not in {"", "-", "0", "OTHER"} for value in parts): + return cls("partial|" + "|".join(parts), False) + return cls("unknown-flow", False) + + +def _safe_port(value: Any) -> int: + try: + port = int(value or 0) + except (TypeError, ValueError): + return 0 + return port if 0 <= port <= 65535 else 0 + + +@dataclass(frozen=True) +class FlowWorkerJob: + job_id: str + created_at: str + flow_key: str + packet: dict[str, Any] + metadata: dict[str, str] + priority: str = "normal" + + +@dataclass +class _WorkerState: + worker_id: int + queue: queue.Queue[FlowWorkerJob] + thread: threading.Thread | None = None + processed_total: int = 0 + failed_total: int = 0 + dropped_total: int = 0 + latencies_ms: deque[float] = field(default_factory=lambda: deque(maxlen=2048)) + + +class FlowWorkerPool: + """Bounded ordered lanes for parallel packet processing by flow.""" + + def __init__( + self, + processor: Callable[[dict[str, Any]], None], + *, + enabled: bool = True, + worker_count: int = DEFAULT_WORKER_COUNT, + queue_max: int = DEFAULT_QUEUE_MAX, + overflow_policy: str = "drop_oldest", + shutdown_timeout_sec: float = DEFAULT_SHUTDOWN_TIMEOUT_SEC, + error_threshold: int = DEFAULT_ERROR_THRESHOLD, + slow_job_ms: float = DEFAULT_SLOW_JOB_MS, + block_timeout_sec: float = 0.05, + ) -> None: + self._processor = processor + self.enabled = bool(enabled) + self.worker_count = max(1, min(64, int(worker_count))) + self.queue_max_total = max(self.worker_count, int(queue_max)) + self.overflow_policy = ( + overflow_policy + if overflow_policy in VALID_OVERFLOW_POLICIES + else "drop_oldest" + ) + self.shutdown_timeout_sec = max(0.1, float(shutdown_timeout_sec)) + self.error_threshold = max(1, int(error_threshold)) + self.slow_job_ms = max(1.0, float(slow_job_ms)) + self.block_timeout_sec = max(0.001, min(1.0, float(block_timeout_sec))) + self._lock = threading.RLock() + self._stop = threading.Event() + self._accepting = True + self._jobs_received = 0 + self._jobs_processed = 0 + self._jobs_failed = 0 + self._jobs_dropped = 0 + self._jobs_rejected = 0 + self._unknown_flow_keys = 0 + self._slow_jobs = 0 + self._latencies_ms: deque[float] = deque(maxlen=4096) + self._last_error = "" + self._last_drop_reason = "" + self._last_slow_job_at = "" + self._workers = self._create_workers() + if self.enabled: + self._start_workers() + + @classmethod + def from_env(cls, processor: Callable[[dict[str, Any]], None]) -> FlowWorkerPool: + policy = str( + os.environ.get("NETBOT_FLOW_WORKER_OVERFLOW_POLICY", "drop_oldest") + ).strip() + if policy not in VALID_OVERFLOW_POLICIES: + policy = "drop_oldest" + return cls( + processor, + enabled=_env_bool("NETBOT_FLOW_WORKERS_ENABLED", True), + worker_count=_env_int( + "NETBOT_FLOW_WORKER_COUNT", DEFAULT_WORKER_COUNT, 1, 64 + ), + queue_max=_env_int( + "NETBOT_FLOW_WORKER_QUEUE_MAX", DEFAULT_QUEUE_MAX, 1, 1_000_000 + ), + overflow_policy=policy, + shutdown_timeout_sec=_env_float( + "NETBOT_FLOW_WORKER_SHUTDOWN_TIMEOUT_SEC", + DEFAULT_SHUTDOWN_TIMEOUT_SEC, + 0.1, + 120.0, + ), + error_threshold=_env_int( + "NETBOT_FLOW_WORKER_ERROR_THRESHOLD", + DEFAULT_ERROR_THRESHOLD, + 1, + 1_000_000, + ), + slow_job_ms=_env_float( + "NETBOT_FLOW_WORKER_SLOW_JOB_MS", DEFAULT_SLOW_JOB_MS, 1.0, 60_000.0 + ), + ) + + def _create_workers(self) -> list[_WorkerState]: + base, remainder = divmod(self.queue_max_total, self.worker_count) + return [ + _WorkerState( + worker_id=worker_id, + queue=queue.Queue(maxsize=base + (1 if worker_id < remainder else 0)), + ) + for worker_id in range(self.worker_count) + ] + + def _start_workers(self) -> None: + for state in self._workers: + state.thread = threading.Thread( + target=self._worker_loop, + args=(state,), + name=f"netbotpro-flow-worker-{state.worker_id}", + daemon=True, + ) + state.thread.start() + + def worker_index_for(self, packet: dict[str, Any]) -> int: + key = FlowDispatchKey.from_packet(packet) + digest = hashlib.sha256(key.value.encode("utf-8")).digest() + return int.from_bytes(digest[:8], "big") % self.worker_count + + def submit( + self, + packet: dict[str, Any], + *, + metadata: dict[str, str] | None = None, + priority: str = "normal", + ) -> bool: + with self._lock: + self._jobs_received += 1 + accepting = self._accepting + if not accepting: + self._record_rejection("worker_pool_closed") + return False + + key = FlowDispatchKey.from_packet(packet) + if not key.complete: + with self._lock: + self._unknown_flow_keys += 1 + + if not self.enabled: + self._process_inline(packet) + return True + + job = FlowWorkerJob( + job_id=uuid.uuid4().hex, + created_at=datetime.now(timezone.utc).isoformat(), + flow_key=key.value, + packet=dict(packet), + metadata=dict( + metadata or {"source": "live_capture", "capture_mode": "live"} + ), + priority=str(priority or "normal"), + ) + state = self._workers[self.worker_index_for(packet)] + return self._enqueue(state, job) + + def _enqueue(self, state: _WorkerState, job: FlowWorkerJob) -> bool: + try: + if self.overflow_policy == "block_short": + state.queue.put(job, timeout=self.block_timeout_sec) + else: + state.queue.put_nowait(job) + return True + except queue.Full: + pass + + if self.overflow_policy == "drop_oldest": + try: + state.queue.get_nowait() + state.queue.task_done() + self._record_drop(state, "flow_worker_queue_full_drop_oldest") + except queue.Empty: + pass + try: + state.queue.put_nowait(job) + return True + except queue.Full: + self._record_drop(state, "flow_worker_queue_full_after_drop_oldest") + return False + + if self.overflow_policy == "drop_newest": + self._record_drop(state, "flow_worker_queue_full_drop_newest") + return False + + reason = ( + "flow_worker_queue_block_timeout" + if self.overflow_policy == "block_short" + else "flow_worker_queue_full_reject_new" + ) + self._record_rejection(reason) + return False + + def _record_drop(self, state: _WorkerState, reason: str) -> None: + with self._lock: + state.dropped_total += 1 + self._jobs_dropped += 1 + self._last_drop_reason = reason + logger.warning( + "flow worker queue pressure; packet job dropped worker_id=%s reason=%s", + state.worker_id, + reason, + ) + + def _record_rejection(self, reason: str) -> None: + with self._lock: + self._jobs_rejected += 1 + self._last_drop_reason = reason + logger.warning( + "flow worker queue pressure; packet job rejected reason=%s", reason + ) + + def _worker_loop(self, state: _WorkerState) -> None: + while not self._stop.is_set() or not state.queue.empty(): + try: + job = state.queue.get(timeout=0.1) + except queue.Empty: + continue + started = time.perf_counter() + failed = False + try: + self._processor(job.packet) + except Exception as exc: # pragma: no cover - behavior asserted via metrics + failed = True + with self._lock: + state.failed_total += 1 + self._jobs_failed += 1 + self._last_error = type(exc).__name__ + logger.error( + "flow worker job failed worker_id=%s error_type=%s", + state.worker_id, + type(exc).__name__, + ) + finally: + latency_ms = (time.perf_counter() - started) * 1000.0 + self._record_completion(state, latency_ms, failed) + state.queue.task_done() + + def _process_inline(self, packet: dict[str, Any]) -> None: + started = time.perf_counter() + try: + self._processor(dict(packet)) + except Exception as exc: + with self._lock: + self._jobs_failed += 1 + self._last_error = type(exc).__name__ + logger.error( + "flow worker fallback failed error_type=%s", type(exc).__name__ + ) + return + latency_ms = (time.perf_counter() - started) * 1000.0 + with self._lock: + self._jobs_processed += 1 + self._latencies_ms.append(latency_ms) + self._record_slow_job(latency_ms) + + def _record_completion( + self, state: _WorkerState, latency_ms: float, failed: bool + ) -> None: + with self._lock: + self._latencies_ms.append(latency_ms) + state.latencies_ms.append(latency_ms) + if not failed: + state.processed_total += 1 + self._jobs_processed += 1 + self._record_slow_job(latency_ms) + + def _record_slow_job(self, latency_ms: float) -> None: + if latency_ms < self.slow_job_ms: + return + self._slow_jobs += 1 + self._last_slow_job_at = datetime.now(timezone.utc).isoformat() + + def wait_until_drained(self, timeout_sec: float | None = None) -> bool: + timeout = ( + self.shutdown_timeout_sec if timeout_sec is None else max(0.0, timeout_sec) + ) + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if all(state.queue.unfinished_tasks == 0 for state in self._workers): + return True + time.sleep(0.01) + return all(state.queue.unfinished_tasks == 0 for state in self._workers) + + def close(self, timeout_sec: float | None = None) -> bool: + with self._lock: + self._accepting = False + drained = self.wait_until_drained(timeout_sec) + self._stop.set() + deadline = time.monotonic() + ( + self.shutdown_timeout_sec if timeout_sec is None else max(0.1, timeout_sec) + ) + for state in self._workers: + thread = state.thread + if thread and thread.is_alive(): + thread.join(timeout=max(0.0, deadline - time.monotonic())) + return drained + + def stats(self) -> dict[str, Any]: + with self._lock: + queue_depth = sum(state.queue.qsize() for state in self._workers) + active_workers = sum( + 1 for state in self._workers if state.thread and state.thread.is_alive() + ) + utilization = round(queue_depth / self.queue_max_total * 100.0, 2) + latencies = list(self._latencies_ms) + pressure_reasons: list[str] = [] + if self.enabled and active_workers < self.worker_count: + pressure_reasons.append("flow_worker_not_alive") + if queue_depth and utilization >= 60.0: + pressure_reasons.append("flow_worker_queue_backlog") + if utilization >= 80.0: + pressure_reasons.append("flow_worker_high_utilization") + if self._slow_jobs: + pressure_reasons.append("flow_worker_slow_jobs") + if self._jobs_failed: + pressure_reasons.append("flow_worker_job_failures") + if self._jobs_dropped or self._jobs_rejected: + pressure_reasons.append("flow_worker_dropped_jobs") + + p95 = _percentile(latencies, 0.95) + health = "healthy" + if self.enabled and ( + active_workers < self.worker_count + or utilization >= 95.0 + or self._jobs_failed >= self.error_threshold + or self._jobs_dropped >= self.error_threshold + or p95 >= self.slow_job_ms * 5 + ): + health = "critical" + elif pressure_reasons: + health = "degraded" + + per_worker = [] + for state in self._workers: + worker_latencies = list(state.latencies_ms) + per_worker.append( + { + "worker_id": state.worker_id, + "worker_alive": bool(state.thread and state.thread.is_alive()), + "queue_depth": state.queue.qsize(), + "queue_max": state.queue.maxsize, + "processed_total": state.processed_total, + "failed_total": state.failed_total, + "dropped_total": state.dropped_total, + "avg_latency_ms": ( + round(sum(worker_latencies) / len(worker_latencies), 2) + if worker_latencies + else 0.0 + ), + "p95_latency_ms": _percentile(worker_latencies, 0.95), + } + ) + + return { + "enabled": self.enabled, + "health": health, + "worker_count": self.worker_count, + "active_workers": active_workers, + "queue_depth_total": queue_depth, + "queue_max_total": self.queue_max_total, + "utilization_percent": utilization, + "overflow_policy": self.overflow_policy, + "jobs_received_total": self._jobs_received, + "jobs_processed_total": self._jobs_processed, + "jobs_failed_total": self._jobs_failed, + "jobs_dropped_total": self._jobs_dropped, + "jobs_rejected_total": self._jobs_rejected, + "unknown_flow_key_total": self._unknown_flow_keys, + "slow_jobs_total": self._slow_jobs, + "avg_processing_latency_ms": ( + round(sum(latencies) / len(latencies), 2) if latencies else 0.0 + ), + "p95_processing_latency_ms": p95, + "max_processing_latency_ms": ( + round(max(latencies), 2) if latencies else 0.0 + ), + "per_worker": per_worker, + "last_error": self._last_error, + "last_drop_reason": self._last_drop_reason, + "last_slow_job_at": self._last_slow_job_at, + "pressure_reasons": pressure_reasons, + } + + +__all__ = [ + "FlowDispatchKey", + "FlowWorkerJob", + "FlowWorkerPool", + "VALID_OVERFLOW_POLICIES", +] diff --git a/backend/app/services/monitoring_service.py b/backend/app/services/monitoring_service.py index f459e6d..199e91e 100644 --- a/backend/app/services/monitoring_service.py +++ b/backend/app/services/monitoring_service.py @@ -116,6 +116,42 @@ def _safe_persistence_reasons(value: Any) -> list[str]: ] +def _safe_flow_worker_error(value: Any) -> str: + error_type = str(value or "") + if error_type and len(error_type) <= 80 and error_type.replace("_", "").isalnum(): + return error_type + return "" + + +def _safe_flow_worker_drop_reason(value: Any) -> str: + reason = str(value or "") + allowed = { + "flow_worker_queue_full_drop_oldest", + "flow_worker_queue_full_after_drop_oldest", + "flow_worker_queue_full_drop_newest", + "flow_worker_queue_full_reject_new", + "flow_worker_queue_block_timeout", + "worker_pool_closed", + } + return reason if reason in allowed else "" + + +def _safe_flow_worker_reasons(value: Any) -> list[str]: + allowed = { + "flow_worker_queue_backlog", + "flow_worker_high_utilization", + "flow_worker_slow_jobs", + "flow_worker_job_failures", + "flow_worker_dropped_jobs", + "flow_worker_not_alive", + } + return ( + [str(reason) for reason in value if str(reason) in allowed] + if isinstance(value, list) + else [] + ) + + def build_monitoring_metrics( *, sniffer_state: dict[str, Any], @@ -126,6 +162,7 @@ def build_monitoring_metrics( event_bus = dict(observability.get("event_bus") or {}) packet_queue = dict(observability.get("packet_queue") or {}) + flow_worker_pool = dict(observability.get("flow_worker_pool") or {}) event_aggregator = dict(observability.get("event_aggregator") or {}) websocket = dict(observability.get("websocket") or {}) persistence = dict(observability.get("persistence") or {}) @@ -215,6 +252,12 @@ def build_monitoring_metrics( or [] ) pressure_reasons: list[str] = list(packet_queue_pressure_reasons) + flow_worker_reasons = _safe_flow_worker_reasons( + flow_worker_pool.get("pressure_reasons") or [] + ) + pressure_reasons.extend( + reason for reason in flow_worker_reasons if reason not in pressure_reasons + ) pressure_reasons.extend( reason for reason in persistence_reasons if reason not in pressure_reasons ) @@ -268,6 +311,7 @@ def build_monitoring_metrics( or persistence.get("health") == "critical" or websocket_send_errors >= 3 or websocket_drops >= _int(event_aggregator.get("client_queue_max") or 1000) + or flow_worker_pool.get("health") == "critical" ): health = "critical" @@ -415,6 +459,57 @@ def build_monitoring_metrics( ), "pressure_reasons": packet_queue_pressure_reasons, }, + "flow_worker_pool": { + "enabled": bool(flow_worker_pool.get("enabled")), + "health": str(flow_worker_pool.get("health") or "healthy"), + "worker_count": _int(flow_worker_pool.get("worker_count")), + "active_workers": _int(flow_worker_pool.get("active_workers")), + "queue_depth_total": _int(flow_worker_pool.get("queue_depth_total")), + "queue_max_total": _int(flow_worker_pool.get("queue_max_total")), + "utilization_percent": _number(flow_worker_pool.get("utilization_percent")), + "overflow_policy": str( + flow_worker_pool.get("overflow_policy") or "drop_oldest" + ), + "jobs_received_total": _int(flow_worker_pool.get("jobs_received_total")), + "jobs_processed_total": _int(flow_worker_pool.get("jobs_processed_total")), + "jobs_failed_total": _int(flow_worker_pool.get("jobs_failed_total")), + "jobs_dropped_total": _int(flow_worker_pool.get("jobs_dropped_total")), + "jobs_rejected_total": _int(flow_worker_pool.get("jobs_rejected_total")), + "unknown_flow_key_total": _int( + flow_worker_pool.get("unknown_flow_key_total") + ), + "slow_jobs_total": _int(flow_worker_pool.get("slow_jobs_total")), + "avg_processing_latency_ms": _number( + flow_worker_pool.get("avg_processing_latency_ms") + ), + "p95_processing_latency_ms": _number( + flow_worker_pool.get("p95_processing_latency_ms") + ), + "max_processing_latency_ms": _number( + flow_worker_pool.get("max_processing_latency_ms") + ), + "per_worker": [ + { + "worker_id": _int(worker.get("worker_id")), + "worker_alive": bool(worker.get("worker_alive")), + "queue_depth": _int(worker.get("queue_depth")), + "queue_max": _int(worker.get("queue_max")), + "processed_total": _int(worker.get("processed_total")), + "failed_total": _int(worker.get("failed_total")), + "dropped_total": _int(worker.get("dropped_total")), + "avg_latency_ms": _number(worker.get("avg_latency_ms")), + "p95_latency_ms": _number(worker.get("p95_latency_ms")), + } + for worker in flow_worker_pool.get("per_worker", []) + if isinstance(worker, dict) + ], + "last_error": _safe_flow_worker_error(flow_worker_pool.get("last_error")), + "last_drop_reason": _safe_flow_worker_drop_reason( + flow_worker_pool.get("last_drop_reason") + ), + "last_slow_job_at": str(flow_worker_pool.get("last_slow_job_at") or ""), + "pressure_reasons": flow_worker_reasons, + }, "persistence": { "enabled": bool( persistence.get("persistence_enabled", persistence.get("enabled")) diff --git a/backend/app/services/sniffer_service.py b/backend/app/services/sniffer_service.py index a636636..180ac4d 100644 --- a/backend/app/services/sniffer_service.py +++ b/backend/app/services/sniffer_service.py @@ -11,12 +11,12 @@ from backend.app.services.capture_policy import current_capture_policy from backend.app.services.event_bus import EventBus from backend.app.services.flow_service import FlowService +from backend.app.services.flow_worker_pool import FlowWorkerPool from backend.app.services.packet_queue import BoundedPacketQueue from backend.app.services.service_attribution import enrich_service_attribution from backend.app.services.settings_service import get_settings_snapshot from backend.app.services.sniffer_dashboard_state import SnifferDashboardState -from backend.app.services.sniffer_detection_pipeline import \ - SnifferDetectionPipeline +from backend.app.services.sniffer_detection_pipeline import SnifferDetectionPipeline from backend.app.services.sniffer_event_publisher import SnifferEventPublisher from backend.app.services.sniffer_persistence import SnifferPersistence from core.capture import CaptureProvider, CaptureSession, SystemCaptureProvider @@ -74,6 +74,7 @@ def __init__( max_size=PACKET_QUEUE_MAX_SIZE, overflow_policy=PACKET_QUEUE_OVERFLOW_POLICY, ) + self._flow_worker_pool = FlowWorkerPool.from_env(self._process_packet) self._packet_worker_stop = threading.Event() self._packet_worker = threading.Thread( target=self._packet_worker_loop, @@ -205,7 +206,10 @@ def _packet_worker_loop(self) -> None: except queue.Empty: continue try: - self._process_packet(item.packet) + if self._flow_worker_pool.enabled: + self._flow_worker_pool.submit(item.packet) + else: + self._process_packet(item.packet) except Exception: logger.exception("Packet processing pipeline crashed") finally: @@ -214,7 +218,11 @@ def _packet_worker_loop(self) -> None: def drain_packet_queue( self, timeout_sec: float = PACKET_QUEUE_DRAIN_TIMEOUT_SEC ) -> bool: - return self._packet_queue.wait_until_drained(timeout_sec) + started = datetime.now().timestamp() + intake_drained = self._packet_queue.wait_until_drained(timeout_sec) + remaining = max(0.0, timeout_sec - (datetime.now().timestamp() - started)) + workers_drained = self._flow_worker_pool.wait_until_drained(remaining) + return intake_drained and workers_drained @staticmethod def _apply_payload_policy(row: dict[str, Any], settings: dict[str, Any]) -> None: @@ -257,6 +265,7 @@ def close(self) -> None: self._packet_worker_stop.set() if self._packet_worker.is_alive(): self._packet_worker.join(timeout=1.0) + self._flow_worker_pool.close() self._persistence.close() close_flow_service = getattr(self._flow_service, "close", None) if callable(close_flow_service): @@ -306,6 +315,9 @@ def persistence_stats(self) -> dict[str, Any]: def packet_queue_stats(self) -> dict[str, Any]: return self._packet_queue.stats(worker_alive=self._packet_worker.is_alive()) + def flow_worker_pool_stats(self) -> dict[str, Any]: + return self._flow_worker_pool.stats() + def auto_block_stats(self) -> dict[str, int | float]: return self._detection_pipeline.stats() @@ -313,6 +325,7 @@ def observability(self) -> dict[str, Any]: return { "event_bus": self._event_bus.stats(), "packet_queue": self.packet_queue_stats(), + "flow_worker_pool": self.flow_worker_pool_stats(), "persistence": self.persistence_stats(), "auto_block": self.auto_block_stats(), } diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 2d66f73..7babe85 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -51,7 +51,8 @@ flowchart TB subgraph Pipeline["Core Capture And Analysis Pipeline"] Provider["Capture Provider
Scapy / Npcap / libpcap"] IntakeQueue["Bounded Packet Intake Queue"] - QueueWorker["Packet Queue Worker"] + QueueWorker["Packet Queue Dispatcher"] + FlowWorkers["Flow-aware Worker Pool
Ordered Flow Lanes"] EventAggregator["Event Aggregator"] Parser["Packet Parser"] Layer7["Layer 7 / TLS Metadata"] @@ -108,8 +109,8 @@ flowchart TB Routes --> AgentAPI Events --> LiveClient CapturePolicy --> Provider - Provider --> IntakeQueue --> QueueWorker --> Parser --> Layer7 --> ProtocolIntel --> FlowEngine - QueueWorker --> EventAggregator --> Events + Provider --> IntakeQueue --> QueueWorker --> FlowWorkers --> Parser --> Layer7 --> ProtocolIntel --> FlowEngine + FlowWorkers --> EventAggregator --> Events ProtocolIntel --> TCPIntel ProtocolIntel --> DNSIntel ProtocolIntel --> HTTPIntel @@ -215,7 +216,19 @@ transactional SQLite batches; flow snapshots are coalesced by `flow_id` and written with batched upserts. Queue depth, utilization, dropped/failed writes, write latency, retry state, and worker liveness are visible in Ops metrics. User-triggered report files remain synchronous and outside the packet-rate -queue. Live ring buffer and flow-aware worker pool work remain future steps. +queue. + +The Flow-aware Worker Pool is the processing-pressure boundary after packet +intake. TCP and UDP use canonical bidirectional endpoint keys; other protocols +use normalized endpoint/protocol keys. Stable hashing selects a fixed FIFO lane, +so packets in one flow retain order while different flows can execute in +parallel. Incomplete metadata uses a safe fallback lane. Worker queues remain +bounded and expose backlog, utilization, failures, drops/rejections, processing +latency, and liveness to `/api/monitoring/metrics` and Ops Snapshot. + +The Event Aggregator remains the realtime delivery-pressure boundary after +processing, while Batch Persistence remains the storage-pressure boundary. +Live ring buffer and benchmark/soak validation remain future steps. ## WebSocket Event Aggregator diff --git a/docs/PERFORMANCE_PIPELINE.md b/docs/PERFORMANCE_PIPELINE.md index 28f35b0..6764df9 100644 --- a/docs/PERFORMANCE_PIPELINE.md +++ b/docs/PERFORMANCE_PIPELINE.md @@ -16,7 +16,7 @@ It was added to: - make packet drops visible instead of silent; - protect UI, flow analysis, detection, and storage from direct capture pressure; -- create the foundation for future batching and worker-pool work; +- create the foundation for bounded batching and flow-aware worker lanes; - expose queue pressure through `/api/monitoring/metrics` and Ops Snapshot. ## Current Pipeline @@ -24,7 +24,7 @@ It was added to: ```text Capture -> Bounded Packet Intake Queue --> Packet Queue Worker +-> Flow-aware Worker Pool -> Existing Packet Processing -> Bounded Batch Persistence (packets, alerts, flow snapshots) -> Event Aggregator @@ -33,12 +33,11 @@ Capture ``` The capture callback copies packet metadata into `BoundedPacketQueue` and -returns quickly. A single packet queue worker drains the queue and sends each -packet through the existing processing path: payload policy, protocol metadata, -flow ingestion, detection, dashboard state, persistence enqueue, and websocket -publishing. The Event Aggregator then batches high-frequency realtime updates -before websocket fan-out so the browser does not receive one message for every -processed packet. +returns quickly. Its dispatcher sends accepted metadata to bounded flow-aware +worker lanes. Each lane runs the existing processing path: payload policy, +protocol metadata, flow ingestion, detection, dashboard state, persistence +enqueue, and websocket publishing. The Event Aggregator then batches +high-frequency realtime updates before websocket fan-out. ## Environment Variables @@ -277,9 +276,53 @@ Reports stay synchronous so callers receive a definite result. Audit stays outside batching, writes immediately under its ordering lock, and is tested for ordered redacted output. -This step does not add a worker pool, database sharding, ClickHouse, benchmark -claims, or a new retention engine. Agent history remains on its existing safe -summary path until integration can preserve heartbeat reliability. +This step does not add database sharding, ClickHouse, benchmark claims, or a +new retention engine. Agent history remains on its existing safe summary path. + +## Flow-aware Worker Pool + +Step 5 places a bounded processing stage after packet intake and before packet, +flow, and DPI analysis. The intake queue protects capture from immediate +downstream pressure. The worker pool independently protects CPU-bound processing +and makes worker backlog, failures, drops, and latency visible. + +TCP and UDP packets use a canonical bidirectional key built from the transport +protocol and normalized endpoint/port pairs. Reverse-direction packets therefore +select the same worker. Other protocols use normalized IP endpoints plus the +protocol. Incomplete metadata uses a deterministic fallback lane and increments +`unknown_flow_key_total` without logging packet contents. + +Each stable key is hashed to one FIFO worker lane. Packets from one flow preserve +submission order. Different flows can map to different workers and process in +parallel. Hash collisions are safe: they reduce parallelism but do not corrupt +ordering. Existing FlowEngine and dashboard locks continue to protect shared +state. + +| Variable | Default | Allowed values | Purpose | +| --- | ---: | --- | --- | +| `NETBOT_FLOW_WORKERS_ENABLED` | `true` | boolean | Enable flow-aware dispatch; `false` keeps the existing direct processing path. | +| `NETBOT_FLOW_WORKER_COUNT` | `4` | `1` to `64` | Number of fixed worker lanes. | +| `NETBOT_FLOW_WORKER_QUEUE_MAX` | `2000` | positive integer | Total bounded capacity distributed across worker lanes. | +| `NETBOT_FLOW_WORKER_OVERFLOW_POLICY` | `drop_oldest` | `drop_oldest`, `drop_newest`, `reject_new`, `block_short` | Behavior when the selected lane is full. | +| `NETBOT_FLOW_WORKER_SHUTDOWN_TIMEOUT_SEC` | `5` | positive seconds | Maximum graceful drain wait. | +| `NETBOT_FLOW_WORKER_ERROR_THRESHOLD` | `25` | positive integer | Repeated failure/drop threshold for critical health. | +| `NETBOT_FLOW_WORKER_SLOW_JOB_MS` | `100` | positive milliseconds | Slow processing threshold. | + +Invalid configuration falls back to conservative defaults. `drop_oldest` keeps +the latest live work, `drop_newest` preserves queued work, `reject_new` refuses +the incoming job explicitly, and `block_short` waits for only a small bounded +interval. No policy busy-waits or retries indefinitely. + +The `flow_worker_pool` monitoring section exposes worker counts, total queue +depth/capacity, utilization, received/processed/failed/dropped/rejected counts, +unknown keys, slow jobs, average/p95/max latency, per-worker counters, safe last +error/drop fields, and pressure reasons. Metrics contain no packet payload, +header, credential, cookie, token, session, or secret data. + +Health is `healthy` under normal backlog and latency, `degraded` for growing +backlog, initial drops/failures, or slow jobs, and `critical` for missing workers, +near-full queues, repeated failures/drops, or extreme latency. These signals +contribute to overall Ops health and focused recommended actions. ## WebSocket Batching / Event Aggregator @@ -415,7 +458,6 @@ This is not the complete performance engine yet. Not implemented yet: -- Flow-aware Worker Pool; - Live Ring Buffer; - Benchmark / Soak Tests; - Optional ClickHouse or external metrics backend. @@ -427,13 +469,12 @@ behavior, or AI autonomous actions. ## Next Planned Steps -1. Flow-aware Worker Pool -2. Live Ring Buffer -3. Benchmark and Soak Tests -4. Performance Validation Report -5. Service Attribution / Destination Intelligence -6. Incident / Correlation Engine -7. Read-only AI Analyst +1. Live Ring Buffer +2. Benchmark and Soak Tests +3. Performance Validation Report +4. Service Attribution / Destination Intelligence +5. Incident / Correlation Engine +6. Read-only AI Analyst ### Recorded Product Direction diff --git a/frontend/src/components/OpsPanel.jsx b/frontend/src/components/OpsPanel.jsx index 1424970..282f46b 100644 --- a/frontend/src/components/OpsPanel.jsx +++ b/frontend/src/components/OpsPanel.jsx @@ -35,6 +35,8 @@ export function OpsPanel({ observability, operationalMetrics = null, isRefreshin const levelFor = (label, fallback = "healthy") => snapshot.summaryCards.find((card) => card.label === label)?.level || fallback; const packetQueue = snapshot.packetQueue; const packetQueueLastDropReason = snapshot.safeLastDropReason || "No drops recorded"; + const flowWorkerPool = snapshot.flowWorkerPool; + const flowWorkerLastDropReason = snapshot.safeFlowWorkerDropReason || "No drops recorded"; const persistence = snapshot.persistence; const eventBus = snapshot.eventBus; const eventAggregator = snapshot.eventAggregator; @@ -117,6 +119,44 @@ export function OpsPanel({ observability, operationalMetrics = null, isRefreshin + +
+ + + = 80 ? "degraded" : "healthy"} hint={`${Number(flowWorkerPool.utilization_percent || 0).toFixed(1)}% used`} /> + + + 0 ? "degraded" : "healthy"} hint={`Last error ${flowWorkerPool.last_error || "None"}`} /> + 0 ? "degraded" : "healthy"} hint={`Rejected ${flowWorkerPool.jobs_rejected_total || 0}`} /> + + 0 ? "warning" : "healthy"} hint={flowWorkerPool.last_slow_job_at || "None recorded"} /> + + = 100 ? "warning" : "healthy"} /> + + +
+ {(flowWorkerPool.per_worker || []).length ? ( +
+ {flowWorkerPool.per_worker.map((worker) => ( +
+
+ Worker {worker.worker_id} + {worker.worker_alive ? "Running" : "Stopped"} +
+ {worker.queue_depth || 0}/{worker.queue_max || 0} queued + {worker.processed_total || 0} processed | p95 {formatMs(worker.p95_latency_ms || 0)} +
+ ))} +
+ ) : null} +
+ { worker_alive: true, health: "healthy", }, + flow_worker_pool: { + enabled: true, + health: "healthy", + worker_count: 4, + active_workers: 4, + queue_depth_total: 0, + queue_max_total: 2000, + jobs_received_total: 42, + jobs_processed_total: 42, + per_worker: [ + { worker_id: 0, worker_alive: true, queue_depth: 0, queue_max: 500, processed_total: 12 }, + ], + }, event_aggregator: { packet_batch_ms: 500, packet_batch_max: 250, @@ -74,6 +87,8 @@ describe("OpsPanel", () => { expect(screen.getByText("60s")).toBeTruthy(); expect(screen.getByText("Capture and Flow Pressure")).toBeTruthy(); expect(screen.getByText("Packet Intake Queue")).toBeTruthy(); + expect(screen.getByText("Flow Worker Pool")).toBeTruthy(); + expect(screen.getByText("Worker 0")).toBeTruthy(); expect(screen.getByText("WebSocket Event Aggregator")).toBeTruthy(); expect(screen.getByText("Batches Sent")).toBeTruthy(); expect(screen.getAllByText("20/100").length).toBeGreaterThan(0); @@ -299,7 +314,7 @@ describe("OpsPanel", () => { expect(screen.getByText("Failed Writes")).toBeTruthy(); expect(screen.getByText("Write Latency p95")).toBeTruthy(); expect(screen.getByText("18.5 ms")).toBeTruthy(); - expect(screen.getByText("Pressure Reasons")).toBeTruthy(); + expect(screen.getAllByText("Pressure Reasons").length).toBeGreaterThan(0); expect(document.body.textContent).not.toMatch(/authorization|cookie|raw-secret/i); }); @@ -500,4 +515,88 @@ describe("OpsPanel", () => { expect(screen.getAllByText("No drops recorded").length).toBeGreaterThan(0); expect(document.body.textContent).not.toMatch(/raw-token|authorization/i); }); + + it.each(["healthy", "degraded", "critical"])( + "renders flow worker pool %s health", + (health) => { + render( + + ); + + expect(screen.getByText("Flow Worker Pool")).toBeTruthy(); + expect(screen.getAllByText(`Health ${health}`).length).toBeGreaterThan(0); + } + ); + + it("renders disabled flow worker pool state", () => { + render( + + ); + + expect(screen.getAllByText("Disabled").length).toBeGreaterThan(0); + }); + + it("renders flow worker pressure actions and sanitizes diagnostic text", () => { + render( + + ); + + expect(screen.getByText("Flow worker backlog is growing. Increase worker count, reduce capture pressure, or inspect slow packet processing.")).toBeTruthy(); + expect(screen.getByText("Packet processing is slow. Review DPI cost, protocol analysis, and worker count.")).toBeTruthy(); + expect(screen.getByText("Flow worker jobs are failing. Inspect backend logs and recent packet processing errors.")).toBeTruthy(); + expect(screen.getByText("Flow worker jobs were dropped due to processing pressure. Review worker queue size and overflow policy.")).toBeTruthy(); + expect(screen.getByText("A flow worker appears unhealthy. Restart capture or inspect backend runtime logs.")).toBeTruthy(); + expect(document.body.textContent).not.toMatch(/raw-token|authorization/i); + }); }); diff --git a/frontend/src/lib/opsHealth.js b/frontend/src/lib/opsHealth.js index b630530..9ff14c7 100644 --- a/frontend/src/lib/opsHealth.js +++ b/frontend/src/lib/opsHealth.js @@ -55,11 +55,47 @@ function safeWebSocketDropReason(value) { return allowed.has(value) ? value : ""; } +function safeFlowWorkerDropReason(value) { + const allowed = new Set([ + "flow_worker_queue_full_drop_oldest", + "flow_worker_queue_full_after_drop_oldest", + "flow_worker_queue_full_drop_newest", + "flow_worker_queue_full_reject_new", + "flow_worker_queue_block_timeout", + "worker_pool_closed", + ]); + return allowed.has(value) ? value : ""; +} + +function safeFlowWorkerError(value) { + const errorType = String(value || ""); + return errorType.length <= 80 && /^[A-Za-z0-9_]+$/.test(errorType) ? errorType : ""; +} + +function safeFlowWorkerReasons(value) { + const allowed = new Set([ + "flow_worker_queue_backlog", + "flow_worker_high_utilization", + "flow_worker_slow_jobs", + "flow_worker_job_failures", + "flow_worker_dropped_jobs", + "flow_worker_not_alive", + ]); + return Array.isArray(value) ? value.filter((reason) => allowed.has(reason)) : []; +} + export function buildOpsSnapshot(observability, operationalMetrics = null) { const eventBus = observability?.event_bus || {}; const eventAggregator = operationalMetrics?.event_aggregator || observability?.event_aggregator || {}; const websocket = operationalMetrics?.websocket || observability?.websocket || {}; const packetQueue = operationalMetrics?.packet_queue || observability?.packet_queue || {}; + const flowWorkerPool = operationalMetrics?.flow_worker_pool || observability?.flow_worker_pool || {}; + const safeFlowWorkerPool = { + ...flowWorkerPool, + last_error: safeFlowWorkerError(flowWorkerPool.last_error), + last_drop_reason: safeFlowWorkerDropReason(flowWorkerPool.last_drop_reason), + pressure_reasons: safeFlowWorkerReasons(flowWorkerPool.pressure_reasons), + }; const persistence = operationalMetrics?.persistence || observability?.persistence || {}; const history = observability?.history || {}; const autoBlock = observability?.auto_block || {}; @@ -81,6 +117,17 @@ export function buildOpsSnapshot(observability, operationalMetrics = null) { const packetQueueDropped = toNumber(packetQueue.dropped_total ?? packetQueue.dropped_packets); const packetQueueHighWater = toNumber(packetQueue.high_water_mark ?? packetQueue.queue_high_water_mark); const packetQueueWorkerAlive = packetQueue.worker_alive !== false; + const flowWorkersEnabled = flowWorkerPool.enabled === true; + const flowWorkerCount = toNumber(flowWorkerPool.worker_count); + const flowWorkerActive = toNumber(flowWorkerPool.active_workers); + const flowWorkerDepth = toNumber(flowWorkerPool.queue_depth_total); + const flowWorkerMax = toNumber(flowWorkerPool.queue_max_total); + const flowWorkerUtilization = toNumber(flowWorkerPool.utilization_percent); + const flowWorkerFailed = toNumber(flowWorkerPool.jobs_failed_total); + const flowWorkerDropped = toNumber(flowWorkerPool.jobs_dropped_total); + const flowWorkerRejected = toNumber(flowWorkerPool.jobs_rejected_total); + const flowWorkerSlowJobs = toNumber(flowWorkerPool.slow_jobs_total); + const flowWorkerP95 = toNumber(flowWorkerPool.p95_processing_latency_ms); const droppedWrites = toNumber(persistence.events_dropped_total ?? persistence.dropped_writes); const failedWrites = toNumber(persistence.events_failed_total ?? persistence.failed_writes); const persistenceUtilization = toNumber(persistence.utilization_percent ?? persistence.queue_utilization_percent); @@ -142,10 +189,19 @@ export function buildOpsSnapshot(observability, operationalMetrics = null) { : packetQueueDropped > 0 || packetQueueUtilization >= 80 || (packetQueueMaxSize > 0 && packetQueueDepth >= packetQueueMaxSize * 0.8) ? "warning" : "healthy"; + const flowWorkerLevel = flowWorkersEnabled + ? normalizeLevel(flowWorkerPool.health) !== "healthy" + ? normalizeLevel(flowWorkerPool.health) + : flowWorkerActive < flowWorkerCount || flowWorkerFailed > 0 || flowWorkerDropped > 0 + ? "degraded" + : flowWorkerUtilization >= 80 || flowWorkerRejected > 0 || flowWorkerSlowJobs > 0 || flowWorkerP95 >= 100 + ? "warning" + : "healthy" + : "healthy"; const freshnessLevel = ageSeconds == null || ageSeconds <= 120 ? "healthy" : ageSeconds <= 300 ? "warning" : "degraded"; const criticalFlows = toNumber(flows.risk_distribution?.critical); const highFlows = toNumber(flows.risk_distribution?.high); - const overall = worstLevel(backendLevel, freshnessLevel, packetQueueLevel, persistenceLevel, streamLevel, queryLevel, autoBlockLevel); + const overall = worstLevel(backendLevel, freshnessLevel, packetQueueLevel, flowWorkerLevel, persistenceLevel, streamLevel, queryLevel, autoBlockLevel); const recommendedActions = []; if (backendLevel !== "healthy") { @@ -166,6 +222,21 @@ export function buildOpsSnapshot(observability, operationalMetrics = null) { if (packetQueueMaxSize > 0 && packetQueueHighWater >= packetQueueMaxSize * 0.9) { recommendedActions.push("Queue pressure is approaching capacity. Consider increasing NETBOT_PACKET_QUEUE_MAX_SIZE."); } + if (flowWorkersEnabled && flowWorkerActive < flowWorkerCount) { + recommendedActions.push("A flow worker appears unhealthy. Restart capture or inspect backend runtime logs."); + } + if (flowWorkerUtilization >= 80) { + recommendedActions.push("Flow worker backlog is growing. Increase worker count, reduce capture pressure, or inspect slow packet processing."); + } + if (flowWorkerP95 >= 100 || flowWorkerSlowJobs > 0) { + recommendedActions.push("Packet processing is slow. Review DPI cost, protocol analysis, and worker count."); + } + if (flowWorkerFailed > 0) { + recommendedActions.push("Flow worker jobs are failing. Inspect backend logs and recent packet processing errors."); + } + if (flowWorkerDropped > 0 || flowWorkerRejected > 0) { + recommendedActions.push("Flow worker jobs were dropped due to processing pressure. Review worker queue size and overflow policy."); + } if (criticalFlows > 0) { recommendedActions.push("Review critical flows and related alerts first."); } else if (highFlows > 0) { @@ -248,6 +319,12 @@ export function buildOpsSnapshot(observability, operationalMetrics = null) { hint: `${packetQueueUtilization.toFixed(1)}% used | Drops ${packetQueueDropped}`, level: packetQueueLevel, }, + { + label: "Flow Workers", + value: flowWorkersEnabled ? `${flowWorkerActive}/${flowWorkerCount}` : "Disabled", + hint: `${flowWorkerDepth}/${flowWorkerMax} queued | P95 ${formatMs(flowWorkerP95)}`, + level: flowWorkerLevel, + }, { label: "Write Queue", value: String(queueSize), @@ -308,6 +385,9 @@ export function buildOpsSnapshot(observability, operationalMetrics = null) { packetQueue, safeLastDropReason: safeQueueDropReason(packetQueue.last_drop_reason), packetQueueLevel, + flowWorkerPool: safeFlowWorkerPool, + flowWorkerLevel, + safeFlowWorkerDropReason: safeFlowWorkerPool.last_drop_reason, persistence, history, autoBlock, diff --git a/frontend/src/styles.css b/frontend/src/styles.css index b38f6a1..a9d9d22 100644 --- a/frontend/src/styles.css +++ b/frontend/src/styles.css @@ -1968,6 +1968,36 @@ a { padding-left: 18px; } +.ops-worker-grid { + display: grid; + grid-template-columns: repeat(auto-fit, minmax(190px, 1fr)); + gap: 8px; + margin-top: 12px; +} + +.ops-worker-item { + min-width: 0; + border: 1px solid rgba(255, 255, 255, 0.08); + border-radius: 8px; + padding: 10px 12px; + display: grid; + gap: 6px; + background: rgba(255, 255, 255, 0.02); +} + +.ops-worker-item div { + display: flex; + align-items: center; + justify-content: space-between; + gap: 8px; +} + +.ops-worker-item span, +.ops-worker-item small { + color: var(--muted); + overflow-wrap: anywhere; +} + .graph-summary-card { border: 1px solid rgba(255, 255, 255, 0.08); border-radius: 18px; diff --git a/tests/test_flow_worker_pool.py b/tests/test_flow_worker_pool.py new file mode 100644 index 0000000..5141a29 --- /dev/null +++ b/tests/test_flow_worker_pool.py @@ -0,0 +1,284 @@ +import os +import threading +import time +import unittest +from unittest.mock import patch + +from backend.app.services.flow_worker_pool import FlowDispatchKey, FlowWorkerPool + + +def packet(index: int, *, src: str = "10.0.0.1", dst: str = "8.8.8.8"): + return { + "id": f"packet-{index}", + "sequence": index, + "src": src, + "dst": dst, + "sport": 50000, + "dport": 443, + "proto": "TCP", + } + + +class FlowWorkerPoolTests(unittest.TestCase): + def test_creates_configured_workers_and_accepts_jobs(self): + processed = [] + pool = FlowWorkerPool(processed.append, worker_count=3, queue_max=30) + try: + self.assertTrue(pool.submit(packet(1))) + self.assertTrue(pool.wait_until_drained(1.0)) + stats = pool.stats() + finally: + pool.close() + + self.assertEqual(stats["worker_count"], 3) + self.assertEqual(stats["active_workers"], 3) + self.assertEqual(stats["jobs_received_total"], 1) + self.assertEqual(stats["jobs_processed_total"], 1) + self.assertEqual(processed[0]["sequence"], 1) + + def test_bidirectional_flow_key_and_worker_assignment_are_stable(self): + forward = packet(1) + reverse = { + **forward, + "src": forward["dst"], + "dst": forward["src"], + "sport": forward["dport"], + "dport": forward["sport"], + } + pool = FlowWorkerPool(lambda _packet: None, worker_count=4, queue_max=40) + try: + self.assertEqual( + FlowDispatchKey.from_packet(forward), + FlowDispatchKey.from_packet(reverse), + ) + self.assertEqual( + pool.worker_index_for(forward), pool.worker_index_for(reverse) + ) + finally: + pool.close() + + def test_same_flow_preserves_processing_order(self): + processed = [] + pool = FlowWorkerPool( + lambda item: processed.append(item["sequence"]), + worker_count=4, + queue_max=200, + ) + try: + for index in range(30): + self.assertTrue(pool.submit(packet(index))) + self.assertTrue(pool.wait_until_drained(2.0)) + finally: + pool.close() + + self.assertEqual(processed, list(range(30))) + + def test_different_flows_can_process_concurrently(self): + started = set() + both_started = threading.Event() + release = threading.Event() + lock = threading.Lock() + + def processor(item): + with lock: + started.add(item["src"]) + if len(started) == 2: + both_started.set() + release.wait(1.0) + + pool = FlowWorkerPool(processor, worker_count=2, queue_max=20) + first = packet(1, src="10.0.0.1") + second = packet(2, src="10.0.0.2") + while pool.worker_index_for(second) == pool.worker_index_for(first): + suffix = int(second["src"].rsplit(".", 1)[1]) + 1 + second = packet(2, src=f"10.0.0.{suffix}") + try: + pool.submit(first) + pool.submit(second) + self.assertTrue(both_started.wait(1.0)) + finally: + release.set() + pool.close() + + def test_unknown_packet_uses_safe_fallback_lane(self): + pool = FlowWorkerPool(lambda _packet: None, worker_count=2, queue_max=10) + try: + self.assertTrue(pool.submit({"token": "never-expose-this"})) + self.assertTrue(pool.wait_until_drained(1.0)) + stats = pool.stats() + finally: + pool.close() + + self.assertEqual(stats["unknown_flow_key_total"], 1) + self.assertNotIn("never-expose-this", str(stats)) + + def test_drop_oldest_keeps_newest_queued_job(self): + processed, pool, release = self._blocked_pool("drop_oldest") + try: + self.assertTrue(pool.submit(packet(2))) + self.assertTrue(pool.submit(packet(3))) + stats = pool.stats() + self.assertEqual(stats["queue_depth_total"], 1) + self.assertEqual(stats["jobs_dropped_total"], 1) + self.assertEqual( + stats["last_drop_reason"], "flow_worker_queue_full_drop_oldest" + ) + finally: + release.set() + pool.close() + self.assertEqual(processed, [1, 3]) + + def test_drop_newest_rejects_latest_queued_job(self): + processed, pool, release = self._blocked_pool("drop_newest") + try: + self.assertTrue(pool.submit(packet(2))) + self.assertFalse(pool.submit(packet(3))) + stats = pool.stats() + self.assertEqual(stats["jobs_dropped_total"], 1) + finally: + release.set() + pool.close() + self.assertEqual(processed, [1, 2]) + + def test_reject_new_records_rejection_without_growing_queue(self): + processed, pool, release = self._blocked_pool("reject_new") + try: + self.assertTrue(pool.submit(packet(2))) + self.assertFalse(pool.submit(packet(3))) + stats = pool.stats() + self.assertEqual(stats["queue_depth_total"], 1) + self.assertEqual(stats["jobs_rejected_total"], 1) + self.assertEqual(stats["jobs_dropped_total"], 0) + finally: + release.set() + pool.close() + self.assertEqual(processed, [1, 2]) + + def test_block_short_returns_without_blocking_indefinitely(self): + _processed, pool, release = self._blocked_pool("block_short") + try: + self.assertTrue(pool.submit(packet(2))) + started = time.perf_counter() + self.assertFalse(pool.submit(packet(3))) + elapsed = time.perf_counter() - started + self.assertLess(elapsed, 0.25) + self.assertEqual(pool.stats()["jobs_rejected_total"], 1) + finally: + release.set() + pool.close() + + def test_failures_and_slow_jobs_are_counted_without_logging_secrets(self): + def processor(item): + if item["sequence"] == 1: + raise RuntimeError("Authorization: Bearer raw-secret") + time.sleep(0.02) + + pool = FlowWorkerPool( + processor, + worker_count=1, + queue_max=10, + error_threshold=1, + slow_job_ms=5, + ) + try: + with self.assertLogs( + "backend.app.services.flow_worker_pool", level="ERROR" + ) as logs: + pool.submit(packet(1)) + pool.submit(packet(2)) + self.assertTrue(pool.wait_until_drained(1.0)) + stats = pool.stats() + finally: + pool.close() + + self.assertEqual(stats["jobs_failed_total"], 1) + self.assertEqual(stats["slow_jobs_total"], 1) + self.assertEqual(stats["health"], "critical") + self.assertEqual(stats["last_error"], "RuntimeError") + self.assertNotIn("raw-secret", str(stats)) + self.assertNotIn("raw-secret", "\n".join(logs.output)) + + def test_backlog_degrades_health_and_queue_never_exceeds_max(self): + started = threading.Event() + release = threading.Event() + + def processor(_item): + started.set() + release.wait(1.0) + + pool = FlowWorkerPool(processor, worker_count=1, queue_max=10) + try: + pool.submit(packet(0)) + self.assertTrue(started.wait(1.0)) + for index in range(1, 7): + pool.submit(packet(index)) + stats = pool.stats() + self.assertLessEqual(stats["queue_depth_total"], 10) + self.assertEqual(stats["health"], "degraded") + self.assertIn("flow_worker_queue_backlog", stats["pressure_reasons"]) + finally: + release.set() + pool.close() + + def test_disabled_pool_processes_inline_as_legacy_fallback(self): + processed = [] + pool = FlowWorkerPool(processed.append, enabled=False, worker_count=4) + try: + self.assertTrue(pool.submit(packet(1))) + stats = pool.stats() + finally: + pool.close() + + self.assertEqual(processed[0]["sequence"], 1) + self.assertFalse(stats["enabled"]) + self.assertEqual(stats["active_workers"], 0) + self.assertEqual(stats["jobs_processed_total"], 1) + + def test_invalid_environment_values_use_safe_defaults(self): + env = { + "NETBOT_FLOW_WORKERS_ENABLED": "invalid", + "NETBOT_FLOW_WORKER_COUNT": "many", + "NETBOT_FLOW_WORKER_QUEUE_MAX": "none", + "NETBOT_FLOW_WORKER_OVERFLOW_POLICY": "unsafe", + } + with patch.dict(os.environ, env, clear=False): + pool = FlowWorkerPool.from_env(lambda _packet: None) + try: + stats = pool.stats() + finally: + pool.close() + + self.assertTrue(stats["enabled"]) + self.assertEqual(stats["worker_count"], 4) + self.assertEqual(stats["queue_max_total"], 2000) + self.assertEqual(stats["overflow_policy"], "drop_oldest") + + @staticmethod + def _blocked_pool(policy): + started = threading.Event() + release = threading.Event() + processed = [] + + def processor(item): + processed.append(item["sequence"]) + if item["sequence"] == 1: + started.set() + release.wait(1.0) + + pool = FlowWorkerPool( + processor, + worker_count=1, + queue_max=1, + overflow_policy=policy, + block_timeout_sec=0.02, + ) + pool.submit(packet(1)) + if not started.wait(1.0): + release.set() + pool.close() + raise AssertionError("worker did not start") + return processed, pool, release + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_monitoring_metrics.py b/tests/test_monitoring_metrics.py index fb19dbf..8921a11 100644 --- a/tests/test_monitoring_metrics.py +++ b/tests/test_monitoring_metrics.py @@ -5,7 +5,7 @@ from fastapi.testclient import TestClient -from backend.app.main import app +from backend.app.main import _observability_snapshot, app from backend.app.security import require_local_token, require_trusted_client from backend.app.services.monitoring_service import build_monitoring_metrics @@ -55,6 +55,24 @@ def test_build_monitoring_metrics_reports_healthy_snapshot(self): "worker_alive": True, "health": "healthy", }, + "flow_worker_pool": { + "enabled": True, + "health": "healthy", + "worker_count": 4, + "active_workers": 4, + "queue_depth_total": 0, + "queue_max_total": 2000, + "utilization_percent": 0, + "jobs_received_total": 12, + "jobs_processed_total": 12, + "per_worker": [ + { + "worker_id": 0, + "worker_alive": True, + "processed_total": 3, + } + ], + }, "persistence": { "enabled": True, "max_size": 5000, @@ -101,6 +119,9 @@ def test_build_monitoring_metrics_reports_healthy_snapshot(self): self.assertEqual(payload["packet_queue"]["current_depth"], 0) self.assertEqual(payload["packet_queue"]["health"], "healthy") self.assertEqual(payload["packet_queue"]["pressure_reasons"], []) + self.assertTrue(payload["flow_worker_pool"]["enabled"]) + self.assertEqual(payload["flow_worker_pool"]["worker_count"], 4) + self.assertEqual(payload["flow_worker_pool"]["jobs_processed_total"], 12) self.assertEqual(payload["event_aggregator"]["packet_batch_ms"], 500) self.assertEqual(payload["websocket"]["clients"], 1) self.assertEqual(payload["websocket"]["send_latency_ms_avg"], 2.0) @@ -222,6 +243,67 @@ def test_build_monitoring_metrics_reports_stopped_packet_worker(self): ) self.assertIn("packet_queue_worker_stopped", payload["pressure_reasons"]) + def test_flow_worker_pressure_contributes_to_overall_health(self): + payload = build_monitoring_metrics( + sniffer_state={"running": True}, + observability={ + "packet_queue": {"worker_alive": True}, + "flow_worker_pool": { + "enabled": True, + "health": "critical", + "worker_count": 4, + "active_workers": 3, + "queue_depth_total": 1900, + "queue_max_total": 2000, + "utilization_percent": 95, + "jobs_failed_total": 25, + "jobs_dropped_total": 1, + "last_error": "RuntimeError", + "last_drop_reason": "flow_worker_queue_full_drop_oldest", + "pressure_reasons": [ + "flow_worker_not_alive", + "flow_worker_high_utilization", + "flow_worker_job_failures", + "flow_worker_dropped_jobs", + ], + }, + }, + flow_summary={}, + ) + + self.assertEqual(payload["health"], "critical") + self.assertEqual(payload["flow_worker_pool"]["health"], "critical") + self.assertIn("flow_worker_not_alive", payload["pressure_reasons"]) + self.assertIn("flow_worker_high_utilization", payload["pressure_reasons"]) + + def test_flow_worker_metrics_do_not_expose_sensitive_values(self): + payload = build_monitoring_metrics( + sniffer_state={"running": True}, + observability={ + "packet_queue": {"worker_alive": True}, + "flow_worker_pool": { + "enabled": True, + "last_error": "Authorization: Bearer raw-token", + "last_drop_reason": "Cookie: session=raw-token", + "pressure_reasons": [ + "flow_worker_slow_jobs", + "token=raw-token", + ], + "per_worker": [{"worker_id": 0, "token": "raw-token"}], + }, + }, + flow_summary={}, + ) + + rendered = str(payload["flow_worker_pool"]) + self.assertNotIn("raw-token", rendered) + self.assertNotIn("Authorization", rendered) + self.assertNotIn("Cookie", rendered) + self.assertEqual( + payload["flow_worker_pool"]["pressure_reasons"], + ["flow_worker_slow_jobs"], + ) + def test_persistence_health_and_pressure_contribute_to_overall_ops(self): payload = build_monitoring_metrics( sniffer_state={"running": True}, @@ -342,8 +424,24 @@ def test_monitoring_metrics_endpoint_returns_compact_snapshot(self): self.assertIn("event_aggregator", payload) self.assertIn("websocket", payload) self.assertIn("persistence", payload) + self.assertIn("flow_worker_pool", payload) self.assertEqual(payload["health"], "healthy") + def test_runtime_observability_snapshot_includes_flow_worker_pool(self): + expected = { + "enabled": True, + "worker_count": 4, + "active_workers": 4, + "health": "healthy", + } + with patch( + "backend.app.main.sniffer_service.flow_worker_pool_stats", + return_value=expected, + ): + snapshot = _observability_snapshot() + + self.assertEqual(snapshot["flow_worker_pool"], expected) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_release_readiness.py b/tests/test_release_readiness.py index 9beef0e..3e98651 100644 --- a/tests/test_release_readiness.py +++ b/tests/test_release_readiness.py @@ -108,6 +108,10 @@ def test_performance_pipeline_docs_are_linked_and_scoped(self): self.assertIn("docs/PERFORMANCE_PIPELINE.md", readme) self.assertIn("Bounded Packet Intake Queue", readme) + self.assertIn("Flow-aware Worker Pool", readme) + self.assertIn("Flow-aware Worker Pool", performance) + self.assertIn("NETBOT_FLOW_WORKER_COUNT", performance) + self.assertIn("Flow-aware Worker Pool", architecture) self.assertIn("Queue pressure metrics", readme) self.assertIn("Ops Snapshot packet queue visibility", readme) self.assertIn("NETBOT_PACKET_QUEUE_MAX_SIZE", readme) diff --git a/tests/test_sniffer_service.py b/tests/test_sniffer_service.py index 99d496d..232dca8 100644 --- a/tests/test_sniffer_service.py +++ b/tests/test_sniffer_service.py @@ -1,8 +1,12 @@ +import os import unittest from unittest.mock import patch from backend.app.services.event_bus import EventBus -from backend.app.services.sniffer_service import CaptureStartUnavailableError, SnifferService +from backend.app.services.sniffer_service import ( + CaptureStartUnavailableError, + SnifferService, +) class _FakeCaptureSession: @@ -31,7 +35,11 @@ def create_session(self, packet_callback): return self.session def list_interfaces(self): - return {"recommended": "eth0", "recommended_label": "Ethernet", "items": [{"value": "eth0", "name": "Ethernet"}]} + return { + "recommended": "eth0", + "recommended_label": "Ethernet", + "items": [{"value": "eth0", "name": "Ethernet"}], + } def describe_interface(self, candidate): return "Ethernet" @@ -47,7 +55,14 @@ def to_dict(): "provider": "fake", "ready": True, "recommended_interface": "eth0", - "checks": [{"code": "interfaces_available", "ok": True, "severity": "error", "detail": "ok"}], + "checks": [ + { + "code": "interfaces_available", + "ok": True, + "severity": "error", + "detail": "ok", + } + ], } return _Report() @@ -55,7 +70,9 @@ def to_dict(): class _UnavailableCaptureProvider(_FakeCaptureProvider): def create_session(self, packet_callback): - raise AssertionError("create_session should not be called when preflight is not ready") + raise AssertionError( + "create_session should not be called when preflight is not ready" + ) def preflight(self): class _Report: @@ -64,7 +81,14 @@ def to_dict(): return { "provider": "fake", "ready": False, - "checks": [{"code": "interfaces_available", "ok": False, "severity": "error", "detail": "Detected 0 capture interface(s)."}], + "checks": [ + { + "code": "interfaces_available", + "ok": False, + "severity": "error", + "detail": "Detected 0 capture interface(s).", + } + ], } return _Report() @@ -119,7 +143,9 @@ def test_capture_interfaces_include_preflight(self): self.assertEqual(payload["preflight"]["provider"], "fake") def test_start_fails_fast_when_preflight_is_not_ready(self): - service = SnifferService(EventBus(), capture_provider=_UnavailableCaptureProvider()) + service = SnifferService( + EventBus(), capture_provider=_UnavailableCaptureProvider() + ) with self.assertRaises(CaptureStartUnavailableError) as ctx: service.start("iface=default") @@ -135,8 +161,13 @@ def test_start_rejects_unknown_remote_interface_name(self): self.assertIn("local interfaces", str(ctx.exception)) service.close() - @patch("backend.app.services.sniffer_service.get_settings_snapshot", return_value={"payload_capture_enabled": False, "alert_only_mode": False}) - def test_payload_preview_is_removed_by_default_before_state_and_persistence(self, _mock_settings): + @patch( + "backend.app.services.sniffer_service.get_settings_snapshot", + return_value={"payload_capture_enabled": False, "alert_only_mode": False}, + ) + def test_payload_preview_is_removed_by_default_before_state_and_persistence( + self, _mock_settings + ): service = SnifferService( EventBus(), capture_provider=_FakeCaptureProvider(), @@ -186,7 +217,9 @@ def test_alerts_are_linked_to_packet_flow_id(self): def test_start_times_out_when_capture_session_creation_hangs(self): service = SnifferService(EventBus(), capture_provider=_SlowCaptureProvider()) - with patch("backend.app.services.sniffer_service.CAPTURE_START_TIMEOUT_SEC", 0.01): + with patch( + "backend.app.services.sniffer_service.CAPTURE_START_TIMEOUT_SEC", 0.01 + ): with self.assertRaises(CaptureStartUnavailableError) as ctx: service.start("eth0") @@ -221,6 +254,21 @@ def test_close_stops_packet_queue_worker_cleanly(self): self.assertFalse(service._packet_worker.is_alive()) + def test_disabled_flow_workers_use_existing_packet_processing_path(self): + with patch.dict(os.environ, {"NETBOT_FLOW_WORKERS_ENABLED": "false"}): + service = SnifferService( + EventBus(), + capture_provider=_FakeCaptureProvider(), + flow_service=_FakeFlowService(), + ) + try: + service._on_packet({"src": "10.0.0.1", "dst": "8.8.8.8", "proto": "UDP"}) + self.assertTrue(service.drain_packet_queue(timeout_sec=1.0)) + self.assertEqual(len(service.recent_packets()), 1) + self.assertFalse(service.flow_worker_pool_stats()["enabled"]) + finally: + service.close() + if __name__ == "__main__": unittest.main()