""" Session-Aufzeichnung mit Hash-Verkettung (siehe Konzept 6.5). Jede Session schreibt eine oder mehrere JSONL-Dateien unter settings.recordings_dir. Jede Zeile verkettet sich mit der vorherigen (gleiches Prinzip wie das Audit-Log, app/security/audit.py), damit nachtraegliche Manipulation der Aufzeichnung erkennbar ist. Umsetzungsauftrag Teil A (D2) / Teil E (E2): `record()` wurde vorher synchron im Event-Loop aufgerufen und hat bei JEDER Instruktion `write()` + `flush()` gemacht -- beim RDP-Bildstrom also praktisch bei jedem Frame. Das haelt den einzigen Event-Loop der Anwendung fuer ALLE Benutzer an (siehe E.0) und kann bei laufenden Sitzungen unbegrenzt Plattenplatz verbrauchen. Design der Behebung -- zwei klar getrennte Zustaendigkeiten: * Event-Loop-Thread (synchron, `record()`): berechnet die Hash-Kette UND entscheidet ueber Rotation/Groessenbegrenzung -- beides reine CPU-/Zaehlerarbeit ohne I/O, daher unbedenklich synchron. Jede fertige Zeile landet mit ihrem Ziel-Teildateiindex in einem In-Memory-Puffer. * Executor-Thread (`asyncio.to_thread`, `_flush_loop`): fuehrt NUR noch das eigentliche blockierende Datei-I/O aus (open/write/flush/ close) -- trifft keine Entscheidungen und teilt sich daher keinen veraenderlichen Zustand mit dem Event-Loop-Thread ausser den fertig vorbereiteten (Teilindex, Zeile)-Tupeln. Damit bleibt die Hash-Kette deterministisch in Aufrufreihenfolge korrekt, unabhaengig davon, wann tatsaechlich geschrieben wird. Jede Teildatei (Rotation) hat ihre eigene, in sich geschlossene Kette ab GENESIS_HASH -- das genuegt fuer Manipulationserkennung je Datei und vermeidet krossen Zustand zwischen den Threads. Bekannte, bewusst akzeptierte Einschraenkung: `close()`/`aclose()` versucht einen Abschluss-Flush, aber bei hartem Prozessabsturz (kein regulaeres Sitzungsende) koennen die letzten <= FLUSH_INTERVAL_S Sekunden bzw. <= MAX_BUFFERED_ENTRIES gepufferte Eintraege verloren gehen. Die Kette selbst bleibt in jedem Fall bis zum letzten geschriebenen Eintrag gueltig (kein "gebrochener" Zustand wie beim alten Audit-Log-Bug E1) -- es fehlt hoechstens ein Rest am Ende. """ from __future__ import annotations import asyncio import hashlib import json import logging import time from pathlib import Path from app.config import settings logger = logging.getLogger("jumphost.recordings") GENESIS_HASH = "0" * 64 # Grobe Rahmenkosten (JSON-Huelle {"entry":...,"prev_hash":...,"hash":...} # plus Zeilenumbruch) fuer die Vorab-Groessenschaetzung in record() -- # muss nicht exakt sein, nur konservativ genug, um Rotation/Truncation # rechtzeitig auszuloesen, bevor eine Teildatei stark ueberschritten wird. _FRAME_OVERHEAD_BYTES = 96 class SessionRecorder: #: Wie oft der Hintergrund-Flush hoechstens schlaeft, bevor er den #: Puffer erneut prueft (er wacht frueher auf, wenn der Puffer voll wird). FLUSH_INTERVAL_S: float = settings.recording_flush_interval_s #: Ab wie vielen gepufferten Eintraegen sofort (statt erst nach #: FLUSH_INTERVAL_S) geflusht wird -- begrenzt den Speicherbedarf des #: Puffers bei sehr schnellen Sitzungen (RDP-Bildstrom). MAX_BUFFERED_ENTRIES: int = 200 def __init__(self, session_id: int) -> None: self.session_id = session_id self._base_path = settings.recordings_dir / f"session_{session_id}.jsonl" # Nach aussen (DB recording_path, Admin-UI) bleibt dies der stabile # "Ankerpfad" der Sitzung -- Teil 1 der ggf. rotierten Sequenz. self.path = self._base_path self._max_part_bytes = settings.recording_max_part_bytes self._max_total_bytes = settings.recording_max_total_bytes self._start_ts = time.time() self._prev_hash = GENESIS_HASH self._part_index = 1 self._part_queued_bytes = 0 self._total_queued_bytes = 0 self._truncated = False self._closed = False # Nur vom Event-Loop-Thread beruehrt (record()/_enqueue*). self._buffer: list[tuple[int, str]] = [] # Nur vom Executor-Thread beruehrt (_write_batch_sync und Freunde). self._writer_fh = self._open_part(self._part_index) self._writer_part_index = self._part_index self._flush_event = asyncio.Event() self._flush_task: asyncio.Task | None = asyncio.get_event_loop().create_task( self._flush_loop(), name=f"recorder-flush-{session_id}" ) # -- Pfade ----------------------------------------------------------- def _part_path(self, part_index: int) -> Path: if part_index == 1: return self._base_path return self._base_path.with_suffix(self._base_path.suffix + f".{part_index}") def _open_part(self, part_index: int): path = self._part_path(part_index) fh = open(path, "a", encoding="utf-8") try: path.chmod(0o600) except OSError: pass return fh # -- Event-Loop-Seite (synchron, kein I/O) ---------------------------- def record(self, direction: str, data: str) -> None: """direction: 'input' (Tastatureingabe) oder 'output' (Terminal-/RDP-Ausgabe). Bleibt bewusst eine SYNCHRONE Methode (kein `async def`/`await`): alle bestehenden Aufrufstellen (app/rdp_proxy/ws_tunnel.py, app/ssh_proxy/terminal_ws.py) rufen sie ohne `await` aus Hot-Path-Code auf. Sie darf daher niemals blockierendes I/O ausfuehren -- das eigentliche Schreiben passiert ausschliesslich im Hintergrund-Task (_flush_loop) im Executor. """ if self._closed or self._truncated: return offset = round(time.time() - self._start_ts, 4) entry = {"t": offset, "dir": direction, "data": data} entry_json = json.dumps(entry, ensure_ascii=False, sort_keys=True) estimated_bytes = len(entry_json.encode("utf-8")) + _FRAME_OVERHEAD_BYTES if self._total_queued_bytes + estimated_bytes > self._max_total_bytes: self._enqueue_truncation_marker() return if self._part_queued_bytes > 0 and self._part_queued_bytes + estimated_bytes > self._max_part_bytes: # Rotation: naechste Teildatei beginnt mit einer eigenen, # frischen Genesis-Hash-Kette -- unabhaengig von der vorherigen. self._part_index += 1 self._part_queued_bytes = 0 self._prev_hash = GENESIS_HASH logger.info( "Sitzungsaufzeichnung %s: rotiere auf Teil %d (Groessenbegrenzung %d Bytes erreicht)", self.session_id, self._part_index, self._max_part_bytes, ) self._append_line(entry, entry_json) def _append_line(self, entry: dict, entry_json: str) -> None: entry_hash = hashlib.sha256((self._prev_hash + "|" + entry_json).encode()).hexdigest() line = json.dumps({"entry": entry, "prev_hash": self._prev_hash, "hash": entry_hash}) line_bytes = len(line.encode("utf-8")) + 1 self._prev_hash = entry_hash self._part_queued_bytes += line_bytes self._total_queued_bytes += line_bytes self._buffer.append((self._part_index, line)) if len(self._buffer) >= self.MAX_BUFFERED_ENTRIES: self._flush_event.set() def _enqueue_truncation_marker(self) -> None: if self._truncated: return self._truncated = True logger.warning( "Sitzungsaufzeichnung %s: Gesamtgroessenbegrenzung (%d Bytes) erreicht, " "Aufzeichnung wird ab hier abgeschnitten (Sitzung laeuft normal weiter).", self.session_id, self._max_total_bytes, ) offset = round(time.time() - self._start_ts, 4) entry = {"t": offset, "dir": "system", "data": "[recording truncated: max size reached]"} entry_json = json.dumps(entry, ensure_ascii=False, sort_keys=True) self._append_line(entry, entry_json) self._flush_event.set() # -- Executor-Seite (blockierendes Datei-I/O, nie im Event-Loop) ----- def _write_batch_sync(self, items: list[tuple[int, str]]) -> None: """Laeuft ausschliesslich via asyncio.to_thread; die einzige Stelle, die self._writer_fh/self._writer_part_index anfasst -- sequentiell, da _flush_loop jeden to_thread-Aufruf abwartet, bevor der naechste beginnt (kein zweiter Executor-Thread gleichzeitig fuer denselben Recorder).""" try: for part_index, line in items: if part_index != self._writer_part_index: self._writer_fh.close() self._writer_part_index = part_index self._writer_fh = self._open_part(part_index) self._writer_fh.write(line + "\n") self._writer_fh.flush() except OSError: logger.exception( "Sitzungsaufzeichnung %s: Schreibfehler im Aufzeichnungs-Flush", self.session_id ) # -- Hintergrund-Task -------------------------------------------------- async def _flush_loop(self) -> None: try: while True: try: await asyncio.wait_for(self._flush_event.wait(), timeout=self.FLUSH_INTERVAL_S) except asyncio.TimeoutError: pass self._flush_event.clear() await self._drain() if self._closed and not self._buffer: return except asyncio.CancelledError: # Letzter bestmoeglicher Abschluss-Flush vor dem Beenden des Tasks. await self._drain() raise async def _drain(self) -> None: if not self._buffer: return items, self._buffer = self._buffer, [] await asyncio.to_thread(self._write_batch_sync, items) # -- Abschluss ----------------------------------------------------------- async def aclose(self) -> None: """Beendet die Aufzeichnung: signalisiert dem Hintergrund-Task das Ende, wartet auf dessen letzten Flush (bis zu FLUSH_INTERVAL_S plus Schreibzeit) und schliesst die Datei. MUSS aus einer Coroutine aufgerufen werden (siehe Aufrufstellen in ws_tunnel.py/terminal_ws.py, beide bereits in einem `finally:`-Block einer async-Funktion).""" if self._closed: return self._closed = True self._flush_event.set() if self._flush_task is not None: try: await asyncio.wait_for(self._flush_task, timeout=self.FLUSH_INTERVAL_S + 5.0) except asyncio.TimeoutError: logger.warning( "Sitzungsaufzeichnung %s: Abschluss-Flush-Task reagierte nicht rechtzeitig, " "breche ab (letzte gepufferte Eintraege koennen fehlen).", self.session_id, ) self._flush_task.cancel() except asyncio.CancelledError: pass try: self._writer_fh.flush() self._writer_fh.close() except OSError: pass def _iter_part_paths(base_path: Path): """Liefert alle Teildateien einer (ggf. rotierten) Aufzeichnung in Reihenfolge: zuerst base_path selbst, danach .2, .3, ... solange sie existieren.""" if base_path.exists(): yield base_path part_index = 2 while True: candidate = base_path.with_suffix(base_path.suffix + f".{part_index}") if not candidate.exists(): return yield candidate part_index += 1 def verify_recording(path: Path) -> bool: """Prueft die Hash-Kette EINER Teildatei (Rueckwaertskompatibel: eine unrotierte Aufzeichnung besteht nur aus dieser einen Datei). Fuer eine vollstaendig rotierte Aufzeichnung siehe verify_recording_set().""" prev_hash = GENESIS_HASH with open(path, encoding="utf-8") as fh: for line in fh: if not line.strip(): continue row = json.loads(line) if row["prev_hash"] != prev_hash: return False entry_json = json.dumps(row["entry"], ensure_ascii=False, sort_keys=True) expected = hashlib.sha256((prev_hash + "|" + entry_json).encode()).hexdigest() if expected != row["hash"]: return False prev_hash = row["hash"] return True def verify_recording_set(base_path: Path) -> bool: """Prueft ALLE Teildateien einer (ggf. rotierten) Aufzeichnung. Jede Teildatei hat ihre eigene, unabhaengige Hash-Kette (siehe Moduldoc) -- das Gesamtergebnis ist gueltig, wenn jede einzelne Teildatei fuer sich gueltig ist.""" parts = list(_iter_part_paths(base_path)) if not parts: return False return all(verify_recording(p) for p in parts) def iter_recording_entries(base_path: Path): """Liefert die geparsten 'entry'-Objekte ALLER Teildateien in chronologischer Reihenfolge (Teildateien sind bereits zeitlich fortlaufend, da Rotation nur beim Erreichen der Groessenbegrenzung ausgeloest wird).""" for part_path in _iter_part_paths(base_path): with open(part_path, encoding="utf-8") as fh: for line in fh: line = line.strip() if not line: continue yield json.loads(line)["entry"] def count_recording_entries(base_path: Path) -> int: return sum(1 for _ in iter_recording_entries(base_path))