319 lines
13 KiB
Python
319 lines
13 KiB
Python
"""
|
|
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))
|