From 9fc2079e8628df8df781ec1bf7efb0f725924412 Mon Sep 17 00:00:00 2001 From: malmert Date: Sat, 7 Mar 2026 18:23:28 +0100 Subject: [PATCH] test better shutdown nfstream --- backend/src/config.py | 2 ++ backend/src/utilities/nfstream_flow_worker.py | 2 ++ backend/src/utilities/nfstream_manager.py | 10 ++++++++-- 3 files changed, 12 insertions(+), 2 deletions(-) diff --git a/backend/src/config.py b/backend/src/config.py index 95e60f1..88ecea2 100644 --- a/backend/src/config.py +++ b/backend/src/config.py @@ -60,6 +60,7 @@ class BackendSettings: nfstream_active_timeout_seconds: int nfstream_snapshot_length: int nfstream_n_dissections: int + nfstream_n_meters: int nfstream_cache_ttl_seconds: float nfstream_lookup_window_ms: int nfstream_reader_join_timeout_seconds: float @@ -97,6 +98,7 @@ def load_settings() -> BackendSettings: nfstream_active_timeout_seconds=_env_int("BACKEND_NFSTREAM_ACTIVE_TIMEOUT_SECONDS", 1800), nfstream_snapshot_length=_env_int("BACKEND_NFSTREAM_SNAPSHOT_LENGTH", 1536), nfstream_n_dissections=_env_int("BACKEND_NFSTREAM_N_DISSECTIONS", 20), + nfstream_n_meters=_env_int("BACKEND_NFSTREAM_N_METERS", 1), nfstream_cache_ttl_seconds=_env_float("BACKEND_NFSTREAM_CACHE_TTL_SECONDS", 10.0), nfstream_lookup_window_ms=_env_int("BACKEND_NFSTREAM_LOOKUP_WINDOW_MS", 5_000), nfstream_reader_join_timeout_seconds=_env_float("BACKEND_NFSTREAM_READER_JOIN_TIMEOUT_SECONDS", 2.0), diff --git a/backend/src/utilities/nfstream_flow_worker.py b/backend/src/utilities/nfstream_flow_worker.py index d995366..b8f2e9a 100644 --- a/backend/src/utilities/nfstream_flow_worker.py +++ b/backend/src/utilities/nfstream_flow_worker.py @@ -86,6 +86,7 @@ def main() -> int: parser.add_argument("--active-timeout", type=int, required=True) parser.add_argument("--snapshot-length", type=int, required=True) parser.add_argument("--n-dissections", type=int, required=True) + parser.add_argument("--n-meters", type=int, required=True) parser.add_argument("--promiscuous-mode", action="store_true") args = parser.parse_args() @@ -128,6 +129,7 @@ def main() -> int: idle_timeout=args.idle_timeout, active_timeout=args.active_timeout, n_dissections=args.n_dissections, + n_meters=args.n_meters, statistical_analysis=False, accounting_mode=0, udps=EmitFlowMetadata(), diff --git a/backend/src/utilities/nfstream_manager.py b/backend/src/utilities/nfstream_manager.py index 252b0a7..76b0e95 100644 --- a/backend/src/utilities/nfstream_manager.py +++ b/backend/src/utilities/nfstream_manager.py @@ -185,6 +185,8 @@ class NFStreamManager: str(settings.nfstream_snapshot_length), "--n-dissections", str(settings.nfstream_n_dissections), + "--n-meters", + str(settings.nfstream_n_meters), ] if settings.nfstream_promiscuous_mode: cmd.append("--promiscuous-mode") @@ -198,6 +200,7 @@ class NFStreamManager: text=True, bufsize=1, env=env, + start_new_session=True, ) except Exception: logger.exception("Failed to start NFStream worker for %s", iface) @@ -226,10 +229,13 @@ class NFStreamManager: if process is not None and process.poll() is None: try: - process.send_signal(signal.SIGTERM) + os.killpg(os.getpgid(process.pid), signal.SIGTERM) process.wait(timeout=settings.nfstream_process_stop_timeout_seconds) except subprocess.TimeoutExpired: - process.kill() + try: + os.killpg(os.getpgid(process.pid), signal.SIGKILL) + except ProcessLookupError: + pass except Exception: logger.exception("Failed to stop NFStream worker for %s cleanly", iface)