diff --git a/backend/Dockerfile b/backend/Dockerfile index 2ffc1ab..8c0305e 100644 --- a/backend/Dockerfile +++ b/backend/Dockerfile @@ -3,7 +3,7 @@ FROM python:3.11-slim WORKDIR /app RUN apt-get update \ - && apt-get install -y --no-install-recommends build-essential libpcap-dev pkg-config \ + && DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends build-essential libpcap-dev pkg-config tshark \ && rm -rf /var/lib/apt/lists/* COPY requirements.txt . diff --git a/backend/requirements.txt b/backend/requirements.txt index 0b13e45..79ba04c 100644 --- a/backend/requirements.txt +++ b/backend/requirements.txt @@ -12,7 +12,6 @@ pydantic==2.12.3 pydantic_core==2.41.4 pyroute2==0.9.5 python-multipart==0.0.22 -nfstream==6.5.4 scapy==2.6.1 sniffio==1.3.1 starlette==0.48.0 diff --git a/backend/src/config.py b/backend/src/config.py index 88ecea2..1acdb38 100644 --- a/backend/src/config.py +++ b/backend/src/config.py @@ -54,17 +54,12 @@ class BackendSettings: bridge_bpf_build_dir: str telemetry_process_stop_timeout_seconds: float telemetry_reader_join_timeout_seconds: float - nfstream_enabled: bool - nfstream_promiscuous_mode: bool - nfstream_idle_timeout_seconds: int - 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 - nfstream_process_stop_timeout_seconds: float + tshark_enabled: bool + tshark_display_filter: str + tshark_cache_ttl_seconds: float + tshark_match_window_ms: int + tshark_reader_join_timeout_seconds: float + tshark_process_stop_timeout_seconds: float def load_settings() -> BackendSettings: @@ -92,17 +87,12 @@ def load_settings() -> BackendSettings: bridge_bpf_build_dir=_env_str("BACKEND_BRIDGE_BPF_BUILD_DIR", "/tmp/mitm-bpf"), telemetry_process_stop_timeout_seconds=_env_float("BACKEND_TELEMETRY_PROCESS_STOP_TIMEOUT_SECONDS", 3.0), telemetry_reader_join_timeout_seconds=_env_float("BACKEND_TELEMETRY_READER_JOIN_TIMEOUT_SECONDS", 2.0), - nfstream_enabled=_env_bool("BACKEND_NFSTREAM_ENABLED", True), - nfstream_promiscuous_mode=_env_bool("BACKEND_NFSTREAM_PROMISCUOUS_MODE", True), - nfstream_idle_timeout_seconds=_env_int("BACKEND_NFSTREAM_IDLE_TIMEOUT_SECONDS", 120), - 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), - nfstream_process_stop_timeout_seconds=_env_float("BACKEND_NFSTREAM_PROCESS_STOP_TIMEOUT_SECONDS", 3.0), + tshark_enabled=_env_bool("BACKEND_TSHARK_ENABLED", True), + tshark_display_filter=_env_str("BACKEND_TSHARK_DISPLAY_FILTER", "http or tls or dns"), + tshark_cache_ttl_seconds=_env_float("BACKEND_TSHARK_CACHE_TTL_SECONDS", 5.0), + tshark_match_window_ms=_env_int("BACKEND_TSHARK_MATCH_WINDOW_MS", 1_500), + tshark_reader_join_timeout_seconds=_env_float("BACKEND_TSHARK_READER_JOIN_TIMEOUT_SECONDS", 2.0), + tshark_process_stop_timeout_seconds=_env_float("BACKEND_TSHARK_PROCESS_STOP_TIMEOUT_SECONDS", 3.0), ) diff --git a/backend/src/main.py b/backend/src/main.py index b93687e..95b420c 100644 --- a/backend/src/main.py +++ b/backend/src/main.py @@ -95,11 +95,11 @@ async def shutdown_event() -> None: logging.exception("Failed to stop bridge telemetry collector") try: - from src.utilities.nfstream_manager import nfstream_manager + from src.utilities.tshark_manager import tshark_manager - nfstream_manager.stop() + tshark_manager.stop() except Exception: - logging.exception("Failed to stop NFStream workers") + logging.exception("Failed to stop tshark workers") try: from src.utilities.packet_tracker import packet_tracker diff --git a/backend/src/network_sniffer.py b/backend/src/network_sniffer.py index bb899c7..6b9db47 100644 --- a/backend/src/network_sniffer.py +++ b/backend/src/network_sniffer.py @@ -33,9 +33,9 @@ from src.utilities.interface_bridge_helpers import ( ) from src.config import settings from src.utilities.bridge_telemetry import bridge_telemetry_manager -from src.utilities.nfstream_manager import nfstream_manager from src.utilities.packet_identity import build_packet_uid from src.utilities.packet_tracker import packet_tracker +from src.utilities.tshark_manager import tshark_manager from src.Models.etherType import EtherTypeEnum, ethertype_from_int from src.Models.ip_protocol import IPProtocolEnum, protocol_from_number @@ -371,13 +371,13 @@ def parse_packet(pkt, bridge_label: str, capture_metadata: Optional[Dict[str, An if Raw in pkt and not pkt_info.get("protocol_name"): pkt_info["protocol_name"] = "RAW" - # Flow-level enrichment using NFStream, if available. + # Packet-level enrichment using tshark, if available. try: - flow_info = nfstream_manager.lookup_packet(pkt, pkt_iface) - if flow_info: - _merge_enrichment(pkt_info, flow_info) + packet_info = tshark_manager.lookup_packet(pkt, pkt_iface) + if packet_info: + _merge_enrichment(pkt_info, packet_info) except Exception: - logger.exception("NFStream enrichment failed") + logger.exception("tshark enrichment failed") if ICMP in pkt: inner = pkt[ICMP].payload @@ -536,9 +536,9 @@ def _sync_bridge_telemetry() -> None: } ) try: - nfstream_manager.update_interfaces(active_enrichment_ifaces) + tshark_manager.update_interfaces(active_enrichment_ifaces) except Exception: - logger.exception("Failed to update NFStream enrichment workers") + logger.exception("Failed to update tshark enrichment workers") # ------------------------- @@ -839,6 +839,6 @@ def get_internal_debug_state() -> dict: "buffer_len": len(_PACKET_BUFFER), "bridge_capture_mode": "tc_ingress_raw" if any(s.get("is_bridge") for s in sessions.values()) else "af_packet", "telemetry_ports": sorted({iface for session in sessions.values() for iface in session.get("ports", [])}), - "nfstream": nfstream_manager.get_debug_snapshot(), + "tshark": tshark_manager.get_debug_snapshot(), "packet_tracker": packet_tracker.get_debug_snapshot(), } diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index 8e38816..bd71cec 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -267,7 +267,7 @@ class DatabasePool: except Exception: logger.exception("Failed to publish pkt_info to broadcaster") - async def backfill_flow_metadata( + async def backfill_packet_metadata( self, *, iface: str, @@ -276,17 +276,17 @@ class DatabasePool: src_port: int, dst_port: int, protocol: int, - first_seen_ms: int, - last_seen_ms: int, + length: int, + observed_at_ms: int, enrichment: Dict[str, Any], window_ms: int, ) -> int: - """Update recent packet rows for a flow after enrichment arrives asynchronously.""" + """Update recent packet rows after tshark metadata arrives asynchronously.""" if self._pool is None: await self.init_pool() - lower_bound = datetime.fromtimestamp(max(first_seen_ms - window_ms, 0) / 1000.0, tz=timezone.utc) - upper_bound = datetime.fromtimestamp(max(last_seen_ms + window_ms, 0) / 1000.0, tz=timezone.utc) + lower_bound = datetime.fromtimestamp(max(observed_at_ms - window_ms, 0) / 1000.0, tz=timezone.utc) + upper_bound = datetime.fromtimestamp(max(observed_at_ms + window_ms, 0) / 1000.0, tz=timezone.utc) dpi_metadata = enrichment.get("dpi_metadata") try: @@ -303,19 +303,19 @@ class DatabasePool: app_hostname = COALESCE(packets.app_hostname, $13), app_is_encrypted = COALESCE(packets.app_is_encrypted, $14), dpi_metadata = CASE - WHEN $15::jsonb IS NULL THEN packets.dpi_metadata - WHEN packets.dpi_metadata IS NULL THEN $15::jsonb - ELSE packets.dpi_metadata || $15::jsonb + WHEN $16::jsonb IS NULL THEN packets.dpi_metadata + WHEN packets.dpi_metadata IS NULL THEN $16::jsonb + ELSE packets.dpi_metadata || $16::jsonb END WHERE ip_proto_raw = $1 AND (capture_iface = $2 OR ingress_if = $2 OR egress_if = $2) - AND timestamp BETWEEN $7 AND $8 - AND ( - (src_ip = $3::inet AND dst_ip = $4::inet AND src_port = $5 AND dst_port = $6) - OR - (src_ip = $4::inet AND dst_ip = $3::inet AND src_port = $6 AND dst_port = $5) - ) + AND src_ip = $3::inet + AND dst_ip = $4::inet + AND src_port = $5 + AND dst_port = $6 + AND length = $7 + AND timestamp BETWEEN $8 AND $9 AND ( packets.app_protocol IS NULL OR packets.app_master_protocol IS NULL @@ -323,7 +323,7 @@ class DatabasePool: OR packets.app_confidence IS NULL OR packets.app_hostname IS NULL OR packets.app_is_encrypted IS NULL - OR ($15::jsonb IS NOT NULL) + OR ($16::jsonb IS NOT NULL) ) RETURNING * """, @@ -333,6 +333,7 @@ class DatabasePool: dst_ip, src_port, dst_port, + length, lower_bound, upper_bound, enrichment.get("app_protocol"), @@ -344,7 +345,7 @@ class DatabasePool: json.dumps(dpi_metadata) if dpi_metadata is not None else None, ) except Exception: - logger.exception("DB flow metadata backfill failed") + logger.exception("DB packet metadata backfill failed") return 0 if not rows: @@ -358,7 +359,7 @@ class DatabasePool: try: self.broadcaster.sync_publish(serialized) except Exception: - logger.exception("Failed to publish flow-enriched packet row") + logger.exception("Failed to publish tshark-enriched packet row") return updated_count diff --git a/backend/src/utilities/flow_identity.py b/backend/src/utilities/flow_identity.py deleted file mode 100644 index 66ef650..0000000 --- a/backend/src/utilities/flow_identity.py +++ /dev/null @@ -1,77 +0,0 @@ -"""Helpers for canonical bidirectional flow keys used by enrichers.""" - -from __future__ import annotations - -import time -from typing import Any, Optional, Tuple - -from scapy.all import IP, IPv6, TCP, UDP # type: ignore - - -def flow_key_from_endpoints( - ip_version: int, - protocol: int, - src_ip: Any, - src_port: Any, - dst_ip: Any, - dst_port: Any, -) -> Optional[str]: - """Build a stable bidirectional key for IPv4/IPv6 TCP/UDP flows.""" - if src_ip in (None, "") or dst_ip in (None, ""): - return None - - try: - left = (str(src_ip), int(src_port or 0)) - right = (str(dst_ip), int(dst_port or 0)) - ep1, ep2 = (left, right) if left <= right else (right, left) - return f"{int(ip_version)}|{int(protocol)}|{ep1[0]}|{ep1[1]}|{ep2[0]}|{ep2[1]}" - except Exception: - return None - - -def flow_key_from_packet(pkt: Any) -> Optional[str]: - """Extract a canonical flow key from a Scapy packet.""" - if IP in pkt: - ip_layer = pkt[IP] - src_ip = getattr(ip_layer, "src", None) - dst_ip = getattr(ip_layer, "dst", None) - protocol = int(getattr(ip_layer, "proto", 0) or 0) - ip_version = 4 - elif IPv6 in pkt: - ip_layer = pkt[IPv6] - src_ip = getattr(ip_layer, "src", None) - dst_ip = getattr(ip_layer, "dst", None) - protocol = int(getattr(ip_layer, "nh", 0) or 0) - ip_version = 6 - else: - return None - - if protocol not in (6, 17): - return None - - src_port = 0 - dst_port = 0 - if TCP in pkt: - src_port = int(getattr(pkt[TCP], "sport", 0) or 0) - dst_port = int(getattr(pkt[TCP], "dport", 0) or 0) - elif UDP in pkt: - src_port = int(getattr(pkt[UDP], "sport", 0) or 0) - dst_port = int(getattr(pkt[UDP], "dport", 0) or 0) - - return flow_key_from_endpoints(ip_version, protocol, src_ip, src_port, dst_ip, dst_port) - - -def packet_observed_at_ms(pkt: Any) -> int: - """Return packet timestamp in milliseconds with a real-time fallback.""" - try: - packet_time = float(getattr(pkt, "time", 0.0) or 0.0) - if packet_time > 0: - return int(packet_time * 1000) - except Exception: - pass - return int(time.time() * 1000) - - -def flow_cache_key(iface: str, flow_key: str) -> Tuple[str, str]: - """Key active flow enrichment by interface and canonical flow ID.""" - return (str(iface or ""), flow_key) diff --git a/backend/src/utilities/nfstream_flow_worker.py b/backend/src/utilities/nfstream_flow_worker.py deleted file mode 100644 index b8f2e9a..0000000 --- a/backend/src/utilities/nfstream_flow_worker.py +++ /dev/null @@ -1,168 +0,0 @@ -"""Run NFStream on one interface and emit flow metadata as JSON lines.""" - -from __future__ import annotations - -import argparse -import json -import logging -import signal -import sys -import time -from typing import Any, Optional - -logger = logging.getLogger("nfstream_flow_worker") -_STOP = False - - -def _handle_signal(signum: int, _frame: Any) -> None: - global _STOP - _STOP = True - logger.info("Received signal %s, stopping NFStream worker", signum) - - -def _safe_scalar(value: Any) -> Any: - if isinstance(value, (str, int, float, bool)) or value is None: - return value - return str(value) - - -def _safe_int(value: Any) -> Optional[int]: - try: - if value is None: - return None - return int(value) - except Exception: - return None - - -def _build_event(iface: str, flow: Any, event_name: str) -> Optional[dict[str, Any]]: - ip_version = _safe_int(getattr(flow, "ip_version", None)) - protocol = _safe_int(getattr(flow, "protocol", None)) - src_ip = getattr(flow, "src_ip", None) - dst_ip = getattr(flow, "dst_ip", None) - src_port = _safe_int(getattr(flow, "src_port", None)) - dst_port = _safe_int(getattr(flow, "dst_port", None)) - if ip_version is None or protocol is None or not src_ip or not dst_ip: - return None - - return { - "type": "flow_update", - "event": event_name, - "iface": iface, - "ip_version": ip_version, - "protocol": protocol, - "src_ip": str(src_ip), - "dst_ip": str(dst_ip), - "src_port": src_port or 0, - "dst_port": dst_port or 0, - "first_seen_ms": _safe_int(getattr(flow, "bidirectional_first_seen_ms", None)) - or _safe_int(getattr(flow, "src2dst_first_seen_ms", None)) - or int(time.time() * 1000), - "last_seen_ms": _safe_int(getattr(flow, "bidirectional_last_seen_ms", None)) - or _safe_int(getattr(flow, "src2dst_last_seen_ms", None)) - or int(time.time() * 1000), - "application_name": _safe_scalar(getattr(flow, "application_name", None)), - "application_category_name": _safe_scalar(getattr(flow, "application_category_name", None)), - "application_confidence": _safe_scalar(getattr(flow, "application_confidence", None)), - "requested_server_name": _safe_scalar(getattr(flow, "requested_server_name", None)), - "client_fingerprint": _safe_scalar(getattr(flow, "client_fingerprint", None)), - "server_fingerprint": _safe_scalar(getattr(flow, "server_fingerprint", None)), - "user_agent": _safe_scalar(getattr(flow, "user_agent", None)), - "content_type": _safe_scalar(getattr(flow, "content_type", None)), - "bidirectional_packets": _safe_int(getattr(flow, "bidirectional_packets", None)), - "bidirectional_bytes": _safe_int(getattr(flow, "bidirectional_bytes", None)), - } - - -def _emit(payload: dict[str, Any]) -> None: - sys.stdout.write(json.dumps(payload, separators=(",", ":")) + "\n") - sys.stdout.flush() - - -def main() -> int: - parser = argparse.ArgumentParser(description="NFStream flow worker") - parser.add_argument("--iface", required=True) - parser.add_argument("--idle-timeout", type=int, required=True) - 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() - - signal.signal(signal.SIGTERM, _handle_signal) - signal.signal(signal.SIGINT, _handle_signal) - - try: - from nfstream import NFPlugin, NFStreamer - except Exception as exc: - _emit( - { - "type": "worker_error", - "iface": args.iface, - "message": f"Failed to import nfstream: {exc!r}", - } - ) - return 1 - - class EmitFlowMetadata(NFPlugin): - def on_init(self, _packet: Any, flow: Any) -> None: - payload = _build_event(args.iface, flow, "init") - if payload: - _emit(payload) - - def on_update(self, _packet: Any, flow: Any) -> None: - payload = _build_event(args.iface, flow, "update") - if payload: - _emit(payload) - - def on_expire(self, flow: Any) -> None: - payload = _build_event(args.iface, flow, "expire") - if payload: - _emit(payload) - - try: - streamer = NFStreamer( - source=args.iface, - promiscuous_mode=bool(args.promiscuous_mode), - snapshot_length=args.snapshot_length, - 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(), - ) - except Exception as exc: - _emit( - { - "type": "worker_error", - "iface": args.iface, - "message": f"Failed to start NFStreamer: {exc!r}", - } - ) - return 1 - - _emit({"type": "worker_ready", "iface": args.iface}) - - try: - for _flow in streamer: - if _STOP: - break - except Exception as exc: - _emit( - { - "type": "worker_error", - "iface": args.iface, - "message": f"NFStreamer runtime error: {exc!r}", - } - ) - return 1 - - return 0 - - -if __name__ == "__main__": - logging.basicConfig(level=logging.INFO) - raise SystemExit(main()) diff --git a/backend/src/utilities/nfstream_manager.py b/backend/src/utilities/nfstream_manager.py deleted file mode 100644 index 559de9d..0000000 --- a/backend/src/utilities/nfstream_manager.py +++ /dev/null @@ -1,393 +0,0 @@ -"""Manage optional NFStream flow enrichment workers and packet lookups.""" - -from __future__ import annotations - -import asyncio -import json -import logging -import os -import signal -import subprocess -import sys -import threading -import time -from pathlib import Path -from typing import Any, Dict, Iterable, Optional - -import src.shared_objects as shared_objects -from src.config import settings -from src.utilities.flow_identity import flow_cache_key, flow_key_from_endpoints, flow_key_from_packet, packet_observed_at_ms - -logger = logging.getLogger("nfstream_manager") - - -def _merge_metadata(base: Optional[Dict[str, Any]], extra: Dict[str, Any]) -> Dict[str, Any]: - merged = dict(base or {}) - for key, value in extra.items(): - if value not in (None, "", [], {}): - merged[key] = value - return merged - - -def _safe_text(value: Any) -> Optional[str]: - if value is None: - return None - text = str(value).strip() - return text if text else None - - -def _normalized_label(value: Any) -> Optional[str]: - text = _safe_text(value) - if text is None: - return None - if text.isdigit(): - return None - if text.lower() in {"unknown", "other", "unclassified", "none", "null"}: - return None - return text - - -def _infer_is_encrypted(payload: Dict[str, Any]) -> Optional[bool]: - application_name = str(_normalized_label(payload.get("application_name")) or "").upper() - if any(token in application_name for token in ("TLS", "HTTPS", "QUIC", "SSL")): - return True - return None - - -def _build_enrichment(payload: Dict[str, Any], iface: str, flow_key: str) -> Dict[str, Any]: - app_protocol = _normalized_label(payload.get("application_name")) - if app_protocol is None and any(payload.get(key) for key in ("user_agent", "content_type")): - app_protocol = "HTTP" - - metadata = { - "iface": iface, - "flow_key": flow_key, - "first_seen_ms": int(payload.get("first_seen_ms") or 0) or None, - "last_seen_ms": int(payload.get("last_seen_ms") or 0) or None, - "event": payload.get("event"), - "application_name": _safe_text(payload.get("application_name")), - "application_category_name": _safe_text(payload.get("application_category_name")), - "application_confidence": _safe_text(payload.get("application_confidence")), - "bidirectional_packets": payload.get("bidirectional_packets"), - "bidirectional_bytes": payload.get("bidirectional_bytes"), - } - metadata = _merge_metadata( - metadata, - { - "requested_server_name": payload.get("requested_server_name"), - "client_fingerprint": payload.get("client_fingerprint"), - "server_fingerprint": payload.get("server_fingerprint"), - "user_agent": payload.get("user_agent"), - "content_type": payload.get("content_type"), - }, - ) - - return { - "app_protocol": app_protocol, - "app_master_protocol": app_protocol, - "app_category": _normalized_label(payload.get("application_category_name")), - "app_confidence": _safe_text(payload.get("application_confidence")), - "app_hostname": _safe_text(payload.get("requested_server_name")), - "app_is_encrypted": _infer_is_encrypted(payload), - "dpi_metadata": {"nfstream": metadata}, - } - - -def _has_useful_metadata(payload: Dict[str, Any]) -> bool: - return any( - payload.get(key) not in (None, "", [], {}) - for key in ( - "application_name", - "application_category_name", - "application_confidence", - "requested_server_name", - "client_fingerprint", - "server_fingerprint", - "user_agent", - "content_type", - ) - ) - - -class NFStreamManager: - """Own long-lived NFStream subprocesses and a short-lived flow metadata cache.""" - - def __init__(self) -> None: - self._enabled = settings.nfstream_enabled - self._workers: Dict[str, Dict[str, Any]] = {} - self._cache: Dict[tuple[str, str], Dict[str, Any]] = {} - self._last_error_by_iface: Dict[str, str] = {} - self._lock = threading.Lock() - - @property - def enabled(self) -> bool: - return self._enabled - - def update_interfaces(self, interfaces: Iterable[str]) -> None: - """Align worker set with currently active interfaces.""" - if not self._enabled: - return - - targets = {iface.strip() for iface in interfaces if iface and iface.strip()} - with self._lock: - current = set(self._workers.keys()) - for iface in sorted(current - targets): - self._stop_worker_locked(iface) - for iface in sorted(targets - current): - self._start_worker_locked(iface) - self._purge_cache_locked() - - def stop(self) -> None: - """Stop all workers and clear transient state.""" - with self._lock: - for iface in list(self._workers.keys()): - self._stop_worker_locked(iface) - self._cache.clear() - - def lookup_packet(self, pkt: Any, iface: Optional[str]) -> Dict[str, Any]: - """Return NFStream-derived enrichment for the packet if a recent flow is known.""" - if not self._enabled or not iface: - return {} - - flow_key = flow_key_from_packet(pkt) - if flow_key is None: - return {} - - now_ms = packet_observed_at_ms(pkt) - with self._lock: - self._purge_cache_locked(now_ms=now_ms) - payload = self._cache.get(flow_cache_key(iface, flow_key)) - if payload is None: - return {} - - first_seen_ms = int(payload.get("first_seen_ms") or 0) - last_seen_ms = int(payload.get("last_seen_ms") or 0) - window_ms = settings.nfstream_lookup_window_ms - if now_ms + window_ms < first_seen_ms or now_ms - window_ms > last_seen_ms: - return {} - return _build_enrichment(payload, iface, flow_key) - - def get_debug_snapshot(self) -> Dict[str, Any]: - """Expose current worker/cache state for API debugging.""" - with self._lock: - self._purge_cache_locked() - return { - "enabled": self._enabled, - "workers": { - iface: { - "running": bool(worker.get("process") and worker["process"].poll() is None), - "thread_alive": bool(worker.get("reader") and worker["reader"].is_alive()), - "last_error": self._last_error_by_iface.get(iface), - } - for iface, worker in self._workers.items() - }, - "cache_entries": len(self._cache), - "last_error_by_iface": dict(self._last_error_by_iface), - } - - def _start_worker_locked(self, iface: str) -> None: - helper = Path(__file__).with_name("nfstream_flow_worker.py") - python_bin = sys.executable or "python3" - - env = os.environ.copy() - env["PYTHONUNBUFFERED"] = "1" - backend_root = str(helper.parents[2]) - existing_pythonpath = env.get("PYTHONPATH", "") - env["PYTHONPATH"] = backend_root if not existing_pythonpath else f"{backend_root}:{existing_pythonpath}" - - cmd = [ - python_bin, - str(helper), - "--iface", - iface, - "--idle-timeout", - str(settings.nfstream_idle_timeout_seconds), - "--active-timeout", - str(settings.nfstream_active_timeout_seconds), - "--snapshot-length", - 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") - - logger.info("Starting NFStream worker for %s", iface) - try: - process = subprocess.Popen( - cmd, - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - bufsize=1, - env=env, - start_new_session=True, - ) - except Exception: - logger.exception("Failed to start NFStream worker for %s", iface) - self._last_error_by_iface[iface] = "failed_to_start_process" - return - - reader = threading.Thread( - target=self._read_loop, - args=(iface, process), - daemon=True, - name=f"nfstream-reader-{iface}", - ) - self._workers[iface] = { - "process": process, - "reader": reader, - } - reader.start() - - def _stop_worker_locked(self, iface: str) -> None: - worker = self._workers.pop(iface, None) - if worker is None: - return - - process = worker.get("process") - reader = worker.get("reader") - - if process is not None and process.poll() is None: - try: - os.killpg(os.getpgid(process.pid), signal.SIGTERM) - process.wait(timeout=settings.nfstream_process_stop_timeout_seconds) - except subprocess.TimeoutExpired: - 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) - - if reader is not None and reader.is_alive(): - reader.join(timeout=settings.nfstream_reader_join_timeout_seconds) - - for key in [cache_key for cache_key in self._cache.keys() if cache_key[0] == iface]: - self._cache.pop(key, None) - - def _read_loop(self, iface: str, process: subprocess.Popen[str]) -> None: - stdout = process.stdout - if stdout is None: - return - - for line in stdout: - text = line.strip() - if not text: - continue - try: - payload = json.loads(text) - except json.JSONDecodeError: - logger.info("nfstream[%s]: %s", iface, text) - continue - - payload_type = payload.get("type") - if payload_type == "worker_ready": - logger.info("NFStream worker ready on %s", iface) - continue - if payload_type == "worker_error": - message = str(payload.get("message") or "unknown error") - logger.warning("NFStream worker error on %s: %s", iface, message) - with self._lock: - self._last_error_by_iface[iface] = message - continue - if payload_type != "flow_update": - logger.debug("Ignoring NFStream payload on %s: %s", iface, payload) - continue - - flow_key = flow_key_from_endpoints( - payload.get("ip_version"), - payload.get("protocol"), - payload.get("src_ip"), - payload.get("src_port"), - payload.get("dst_ip"), - payload.get("dst_port"), - ) - if flow_key is None: - continue - - with self._lock: - previous = self._cache.get(flow_cache_key(iface, flow_key)) - self._cache[flow_cache_key(iface, flow_key)] = dict(payload) - self._purge_cache_locked() - - if self._should_backfill(previous, payload): - self._schedule_backfill(iface, flow_key, payload) - - rc = process.poll() - if rc not in (0, None): - logger.warning("NFStream worker for %s exited with code %s", iface, rc) - - def _purge_cache_locked(self, now_ms: Optional[int] = None) -> None: - if now_ms is None: - now_ms = int(time.time() * 1000) - expiry_ms = int(settings.nfstream_cache_ttl_seconds * 1000) - - stale_keys = [] - for key, payload in self._cache.items(): - last_seen_ms = int(payload.get("last_seen_ms") or 0) - if last_seen_ms and now_ms - last_seen_ms > expiry_ms: - stale_keys.append(key) - - for key in stale_keys: - self._cache.pop(key, None) - - def _should_backfill(self, previous: Optional[Dict[str, Any]], current: Dict[str, Any]) -> bool: - if not _has_useful_metadata(current): - return False - if previous is None: - return True - - interesting_keys = ( - "application_name", - "application_category_name", - "application_confidence", - "requested_server_name", - "client_fingerprint", - "server_fingerprint", - "user_agent", - "content_type", - ) - return any(previous.get(key) != current.get(key) for key in interesting_keys) - - def _schedule_backfill(self, iface: str, flow_key: str, payload: Dict[str, Any]) -> None: - web_db = getattr(shared_objects, "db", None) - web_loop = getattr(shared_objects, "web_loop", None) - if web_db is None or web_loop is None: - return - - enrichment = _build_enrichment(payload, iface, flow_key) - - async def _backfill() -> None: - updated = await web_db.backfill_flow_metadata( - iface=iface, - src_ip=str(payload.get("src_ip")), - dst_ip=str(payload.get("dst_ip")), - src_port=int(payload.get("src_port") or 0), - dst_port=int(payload.get("dst_port") or 0), - protocol=int(payload.get("protocol") or 0), - first_seen_ms=int(payload.get("first_seen_ms") or 0), - last_seen_ms=int(payload.get("last_seen_ms") or 0), - enrichment=enrichment, - window_ms=settings.nfstream_lookup_window_ms, - ) - if updated: - logger.debug("Backfilled NFStream metadata for %s packets on %s flow=%s", updated, iface, flow_key) - - try: - future = asyncio.run_coroutine_threadsafe(_backfill(), web_loop) - future.add_done_callback(self._log_backfill_result) - except Exception: - logger.exception("Failed to schedule NFStream metadata backfill for %s", flow_key) - - @staticmethod - def _log_backfill_result(future: Any) -> None: - try: - future.result() - except Exception: - logger.exception("NFStream metadata backfill task failed") - - -nfstream_manager = NFStreamManager() diff --git a/backend/src/utilities/tshark_manager.py b/backend/src/utilities/tshark_manager.py new file mode 100644 index 0000000..899b18c --- /dev/null +++ b/backend/src/utilities/tshark_manager.py @@ -0,0 +1,542 @@ +"""Manage optional tshark packet enrichment workers and matching.""" + +from __future__ import annotations + +import asyncio +import logging +import os +import signal +import subprocess +import threading +import time +from typing import Any, Dict, Iterable, List, Optional, Tuple + +import src.shared_objects as shared_objects +from scapy.all import IP, IPv6, TCP, UDP # type: ignore + +from src.config import settings + +logger = logging.getLogger("tshark_manager") + +_FIELDS: List[str] = [ + "frame.time_epoch", + "frame.interface_name", + "frame.len", + "ip.src", + "ipv6.src", + "ip.dst", + "ipv6.dst", + "ip.proto", + "ipv6.nxt", + "tcp.srcport", + "udp.srcport", + "tcp.dstport", + "udp.dstport", + "frame.protocols", + "http.request.method", + "http.request.uri", + "http.host", + "http.user_agent", + "http.response.code", + "http.response.phrase", + "http.server", + "http.content_type", + "tls.handshake.extensions_server_name", + "tls.handshake.version", + "dns.flags.response", + "dns.qry.name", + "dns.qry.type", + "dns.resp.name", + "dns.a", + "dns.aaaa", + "dns.cname", +] + + +def _safe_text(value: Any) -> Optional[str]: + if value is None: + return None + text = str(value).strip() + return text if text else None + + +def _safe_int(value: Any) -> Optional[int]: + try: + text = _safe_text(value) + if text is None: + return None + return int(text) + except Exception: + return None + + +def _safe_float(value: Any) -> Optional[float]: + try: + text = _safe_text(value) + if text is None: + return None + return float(text) + except Exception: + return None + + +def _safe_bool_flag(value: Any) -> Optional[bool]: + text = _safe_text(value) + if text is None: + return None + if text in {"1", "true", "True"}: + return True + if text in {"0", "false", "False"}: + return False + return None + + +def _jsonable(metadata: Dict[str, Any]) -> Dict[str, Any]: + out: Dict[str, Any] = {} + for key, value in metadata.items(): + if value in (None, "", [], {}): + continue + out[key] = value + return out + + +def _packet_signature( + iface: str, + protocol: int, + src_ip: str, + dst_ip: str, + src_port: int, + dst_port: int, + length: int, +) -> Tuple[str, int, str, str, int, int, int]: + return (str(iface or ""), int(protocol), str(src_ip), str(dst_ip), int(src_port), int(dst_port), int(length)) + + +def _signature_from_packet(pkt: Any, iface: str) -> Optional[Tuple[str, int, str, str, int, int, int]]: + if IP in pkt: + ip_layer = pkt[IP] + protocol = int(getattr(ip_layer, "proto", 0) or 0) + src_ip = str(getattr(ip_layer, "src", "") or "") + dst_ip = str(getattr(ip_layer, "dst", "") or "") + elif IPv6 in pkt: + ip_layer = pkt[IPv6] + protocol = int(getattr(ip_layer, "nh", 0) or 0) + src_ip = str(getattr(ip_layer, "src", "") or "") + dst_ip = str(getattr(ip_layer, "dst", "") or "") + else: + return None + + src_port = 0 + dst_port = 0 + if TCP in pkt: + src_port = int(getattr(pkt[TCP], "sport", 0) or 0) + dst_port = int(getattr(pkt[TCP], "dport", 0) or 0) + elif UDP in pkt: + src_port = int(getattr(pkt[UDP], "sport", 0) or 0) + dst_port = int(getattr(pkt[UDP], "dport", 0) or 0) + + if not src_ip or not dst_ip: + return None + return _packet_signature(iface, protocol, src_ip, dst_ip, src_port, dst_port, len(pkt)) + + +def _packet_observed_at_ms(pkt: Any) -> int: + try: + packet_time = float(getattr(pkt, "time", 0.0) or 0.0) + if packet_time > 0: + return int(packet_time * 1000) + except Exception: + pass + return int(time.time() * 1000) + + +def _parse_line(line: str, fallback_iface: str) -> Optional[Dict[str, Any]]: + parts = line.rstrip("\n").split("\t") + if len(parts) < len(_FIELDS): + parts.extend([""] * (len(_FIELDS) - len(parts))) + row = dict(zip(_FIELDS, parts)) + + timestamp = _safe_float(row["frame.time_epoch"]) + length = _safe_int(row["frame.len"]) + protocol = _safe_int(row["ip.proto"]) or _safe_int(row["ipv6.nxt"]) + src_ip = _safe_text(row["ip.src"]) or _safe_text(row["ipv6.src"]) + dst_ip = _safe_text(row["ip.dst"]) or _safe_text(row["ipv6.dst"]) + src_port = _safe_int(row["tcp.srcport"]) or _safe_int(row["udp.srcport"]) or 0 + dst_port = _safe_int(row["tcp.dstport"]) or _safe_int(row["udp.dstport"]) or 0 + iface = _safe_text(row["frame.interface_name"]) or fallback_iface + + if timestamp is None or length is None or protocol is None or src_ip is None or dst_ip is None: + return None + + observed_at_ms = int(timestamp * 1000) + return { + "iface": iface, + "observed_at_ms": observed_at_ms, + "length": length, + "protocol": protocol, + "src_ip": src_ip, + "dst_ip": dst_ip, + "src_port": src_port, + "dst_port": dst_port, + "frame_protocols": _safe_text(row["frame.protocols"]), + "http": _jsonable( + { + "method": _safe_text(row["http.request.method"]), + "uri": _safe_text(row["http.request.uri"]), + "host": _safe_text(row["http.host"]), + "user_agent": _safe_text(row["http.user_agent"]), + "response_code": _safe_int(row["http.response.code"]), + "response_phrase": _safe_text(row["http.response.phrase"]), + "server": _safe_text(row["http.server"]), + "content_type": _safe_text(row["http.content_type"]), + } + ), + "tls": _jsonable( + { + "server_name": _safe_text(row["tls.handshake.extensions_server_name"]), + "handshake_version": _safe_text(row["tls.handshake.version"]), + } + ), + "dns": _jsonable( + { + "is_response": _safe_bool_flag(row["dns.flags.response"]), + "query_name": _safe_text(row["dns.qry.name"]), + "query_type": _safe_text(row["dns.qry.type"]), + "response_name": _safe_text(row["dns.resp.name"]), + "a": _safe_text(row["dns.a"]), + "aaaa": _safe_text(row["dns.aaaa"]), + "cname": _safe_text(row["dns.cname"]), + } + ), + } + + +def _build_enrichment(event: Dict[str, Any]) -> Dict[str, Any]: + http_meta = dict(event.get("http") or {}) + tls_meta = dict(event.get("tls") or {}) + dns_meta = dict(event.get("dns") or {}) + protocols = str(event.get("frame_protocols") or "") + + app_protocol: Optional[str] = None + app_category: Optional[str] = None + app_hostname: Optional[str] = None + app_is_encrypted: Optional[bool] = None + + if http_meta: + app_protocol = "HTTP" + app_category = "Web" + app_hostname = http_meta.get("host") + app_is_encrypted = False + elif dns_meta: + app_protocol = "DNS" + app_category = "Infrastructure" + app_hostname = dns_meta.get("query_name") or dns_meta.get("response_name") + app_is_encrypted = False + elif tls_meta or "quic" in protocols.lower(): + app_protocol = "QUIC" if "quic" in protocols.lower() and not tls_meta else "TLS" + app_category = "Encrypted" + app_hostname = tls_meta.get("server_name") + app_is_encrypted = True + + tshark_meta = _jsonable( + { + "observed_at_ms": event.get("observed_at_ms"), + "frame_protocols": event.get("frame_protocols"), + "length": event.get("length"), + } + ) + + dpi_metadata = _jsonable( + { + "tshark": tshark_meta, + "http": http_meta, + "tls": tls_meta, + "dns": dns_meta, + } + ) + + return { + "app_protocol": app_protocol, + "app_master_protocol": app_protocol, + "app_category": app_category, + "app_confidence": "high" if app_protocol else None, + "app_hostname": app_hostname, + "app_is_encrypted": app_is_encrypted, + "dpi_metadata": dpi_metadata or None, + } + + +def _has_useful_enrichment(enrichment: Dict[str, Any]) -> bool: + return any( + enrichment.get(key) is not None + for key in ("app_protocol", "app_hostname", "app_is_encrypted", "dpi_metadata") + ) + + +class TsharkManager: + """Own long-lived tshark subprocesses and recent packet metadata cache.""" + + def __init__(self) -> None: + self._enabled = settings.tshark_enabled + self._workers: Dict[str, Dict[str, Any]] = {} + self._cache: Dict[Tuple[str, int, str, str, int, int, int], List[Dict[str, Any]]] = {} + self._last_error_by_iface: Dict[str, str] = {} + self._lock = threading.Lock() + + @property + def enabled(self) -> bool: + return self._enabled + + def update_interfaces(self, interfaces: Iterable[str]) -> None: + if not self._enabled: + return + + targets = {iface.strip() for iface in interfaces if iface and iface.strip()} + with self._lock: + current = set(self._workers.keys()) + for iface in sorted(current - targets): + self._stop_worker_locked(iface) + for iface in sorted(targets - current): + self._start_worker_locked(iface) + self._purge_cache_locked() + + def stop(self) -> None: + with self._lock: + for iface in list(self._workers.keys()): + self._stop_worker_locked(iface) + self._cache.clear() + + def lookup_packet(self, pkt: Any, iface: Optional[str]) -> Dict[str, Any]: + if not self._enabled or not iface: + return {} + + signature = _signature_from_packet(pkt, iface) + if signature is None: + return {} + + packet_time_ms = _packet_observed_at_ms(pkt) + with self._lock: + self._purge_cache_locked(now_ms=packet_time_ms) + entries = self._cache.get(signature) + if not entries: + return {} + + window_ms = settings.tshark_match_window_ms + best_index: Optional[int] = None + best_delta: Optional[int] = None + for index, event in enumerate(entries): + delta = abs(int(event.get("observed_at_ms") or 0) - packet_time_ms) + if delta > window_ms: + continue + if best_delta is None or delta < best_delta: + best_index = index + best_delta = delta + + if best_index is None: + return {} + + event = entries.pop(best_index) + if not entries: + self._cache.pop(signature, None) + + return _build_enrichment(event) + + def get_debug_snapshot(self) -> Dict[str, Any]: + with self._lock: + self._purge_cache_locked() + return { + "enabled": self._enabled, + "workers": { + iface: { + "running": bool(worker.get("process") and worker["process"].poll() is None), + "thread_alive": bool(worker.get("reader") and worker["reader"].is_alive()), + "last_error": self._last_error_by_iface.get(iface), + } + for iface, worker in self._workers.items() + }, + "cache_entries": sum(len(entries) for entries in self._cache.values()), + "last_error_by_iface": dict(self._last_error_by_iface), + } + + def _start_worker_locked(self, iface: str) -> None: + cmd = [ + "tshark", + "-l", + "-n", + "-Q", + "-i", + iface, + "-T", + "fields", + "-E", + "header=n", + "-E", + "separator=\t", + "-E", + "quote=n", + "-E", + "occurrence=f", + ] + if settings.tshark_display_filter: + cmd.extend(["-Y", settings.tshark_display_filter]) + for field in _FIELDS: + cmd.extend(["-e", field]) + + logger.info("Starting tshark worker for %s", iface) + try: + process = subprocess.Popen( + cmd, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + bufsize=1, + env=os.environ.copy(), + start_new_session=True, + ) + except Exception: + logger.exception("Failed to start tshark worker for %s", iface) + self._last_error_by_iface[iface] = "failed_to_start_process" + return + + reader = threading.Thread( + target=self._read_loop, + args=(iface, process), + daemon=True, + name=f"tshark-reader-{iface}", + ) + self._workers[iface] = { + "process": process, + "reader": reader, + } + reader.start() + + def _stop_worker_locked(self, iface: str) -> None: + worker = self._workers.pop(iface, None) + if worker is None: + return + + process = worker.get("process") + reader = worker.get("reader") + + if process is not None and process.poll() is None: + try: + os.killpg(os.getpgid(process.pid), signal.SIGTERM) + process.wait(timeout=settings.tshark_process_stop_timeout_seconds) + except subprocess.TimeoutExpired: + try: + os.killpg(os.getpgid(process.pid), signal.SIGKILL) + except ProcessLookupError: + pass + except Exception: + logger.exception("Failed to stop tshark worker for %s cleanly", iface) + + if reader is not None and reader.is_alive(): + reader.join(timeout=settings.tshark_reader_join_timeout_seconds) + + stale_keys = [key for key in self._cache.keys() if key[0] == iface] + for key in stale_keys: + self._cache.pop(key, None) + + def _read_loop(self, iface: str, process: subprocess.Popen[str]) -> None: + stdout = process.stdout + if stdout is None: + return + + for line in stdout: + text = line.rstrip("\n") + if not text: + continue + + event = _parse_line(text, iface) + if event is None: + logger.debug("Ignoring unparsable tshark line on %s: %s", iface, text) + continue + + signature = _packet_signature( + str(event["iface"]), + int(event["protocol"]), + str(event["src_ip"]), + str(event["dst_ip"]), + int(event["src_port"]), + int(event["dst_port"]), + int(event["length"]), + ) + + with self._lock: + entries = self._cache.setdefault(signature, []) + entries.append(event) + self._purge_cache_locked(now_ms=int(event["observed_at_ms"])) + + enrichment = _build_enrichment(event) + if _has_useful_enrichment(enrichment): + self._schedule_backfill(event, enrichment) + + rc = process.poll() + if rc not in (0, None): + logger.warning("tshark worker for %s exited with code %s", iface, rc) + + def _schedule_backfill(self, event: Dict[str, Any], enrichment: Dict[str, Any]) -> None: + web_db = getattr(shared_objects, "db", None) + web_loop = getattr(shared_objects, "web_loop", None) + if web_db is None or web_loop is None: + return + + async def _backfill() -> None: + updated = await web_db.backfill_packet_metadata( + iface=str(event["iface"]), + src_ip=str(event["src_ip"]), + dst_ip=str(event["dst_ip"]), + src_port=int(event["src_port"]), + dst_port=int(event["dst_port"]), + protocol=int(event["protocol"]), + length=int(event["length"]), + observed_at_ms=int(event["observed_at_ms"]), + enrichment=enrichment, + window_ms=settings.tshark_match_window_ms, + ) + if updated: + logger.debug( + "Backfilled tshark metadata for %s packets on %s %s:%s -> %s:%s", + updated, + event["iface"], + event["src_ip"], + event["src_port"], + event["dst_ip"], + event["dst_port"], + ) + + try: + future = asyncio.run_coroutine_threadsafe(_backfill(), web_loop) + future.add_done_callback(self._log_backfill_result) + except Exception: + logger.exception("Failed to schedule tshark metadata backfill") + + @staticmethod + def _log_backfill_result(future: Any) -> None: + try: + future.result() + except Exception: + logger.exception("tshark metadata backfill task failed") + + def _purge_cache_locked(self, now_ms: Optional[int] = None) -> None: + if now_ms is None: + now_ms = int(time.time() * 1000) + expiry_ms = int(settings.tshark_cache_ttl_seconds * 1000) + + stale_keys: List[Tuple[str, int, str, str, int, int, int]] = [] + for key, entries in self._cache.items(): + fresh_entries = [ + event + for event in entries + if now_ms - int(event.get("observed_at_ms") or 0) <= expiry_ms + ] + if fresh_entries: + self._cache[key] = fresh_entries + else: + stale_keys.append(key) + + for key in stale_keys: + self._cache.pop(key, None) + + +tshark_manager = TsharkManager() diff --git a/setup_build_server.sh b/setup_build_server.sh index 286c0fc..96dd131 100755 --- a/setup_build_server.sh +++ b/setup_build_server.sh @@ -11,6 +11,7 @@ BACKEND_DIR="$APP_DIR/backend" BACKEND_SERVICE="mitm-backend" NGINX_SITE="/etc/nginx/sites-available/mitm-webserver" USER_ROOT="root" +BACKEND_ENV_FILE="$BACKEND_DIR/.env" GITEA_RUNNER_URL="https://gitea.malmert.de//api/v1/repos/marcus/mitm-webserver/actions/runners/register" RUNNER_TOKEN="cThC2xmAZWaOAqRRMENuVVTJckxaHiJxVGx2NCQY" @@ -24,11 +25,11 @@ PYTHON_VERSION="3" # aktuelle Python 3 Version # ----------------------------- echo "==> Update & Upgrade" apt update && apt upgrade -y -apt install -y git curl build-essential nginx python3 python3-pip python3-venv unzip wget python3-dev \ +DEBIAN_FRONTEND=noninteractive apt install -y git curl build-essential nginx python3 python3-pip python3-venv unzip wget python3-dev \ libpcap-dev autoconf automake libtool pkg-config libjson-c-dev gettext flex bison libnuma-dev \ - libpcre2-dev libmaxminddb-dev librrd-dev python3-bpfcc bpfcc-tools linux-headers-generic \ + libpcre2-dev libmaxminddb-dev librrd-dev python3-bpfcc bpfcc-tools linux-headers-generic tshark \ clang llvm libelf-dev libbpf-dev iproute2 -apt install -y linux-headers-$(uname -r) || apt install -y linux-headers-generic +DEBIAN_FRONTEND=noninteractive apt install -y linux-headers-$(uname -r) || DEBIAN_FRONTEND=noninteractive apt install -y linux-headers-generic # ----------------------------- # Node.js installieren (LTS) @@ -70,6 +71,19 @@ cd "$BACKEND_DIR" pip install -r requirements.txt deactivate +# ----------------------------- +# Backend Environment Defaults +# ----------------------------- +echo "==> Write backend environment defaults" +cat >"$BACKEND_ENV_FILE" <