umbau 1.0
This commit is contained in:
@ -1,52 +1,278 @@
|
||||
"""
|
||||
Session-Aufzeichnung mit Hash-Verkettung (siehe Konzept 6.5).
|
||||
|
||||
Jede Session schreibt eine eigene JSONL-Datei 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.
|
||||
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.path = settings.recordings_dir / f"session_{session_id}.jsonl"
|
||||
self._prev_hash = GENESIS_HASH
|
||||
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._fh = open(self.path, "a", encoding="utf-8")
|
||||
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:
|
||||
self.path.chmod(0o600)
|
||||
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)."""
|
||||
"""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})
|
||||
self._fh.write(line + "\n")
|
||||
self._fh.flush()
|
||||
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 close(self) -> None:
|
||||
if not self._fh.closed:
|
||||
self._fh.close()
|
||||
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:
|
||||
@ -61,3 +287,32 @@ def verify_recording(path: Path) -> bool:
|
||||
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))
|
||||
|
||||
63
app/recordings/recorder.py.bak_1788208152
Normal file
63
app/recordings/recorder.py.bak_1788208152
Normal file
@ -0,0 +1,63 @@
|
||||
"""
|
||||
Session-Aufzeichnung mit Hash-Verkettung (siehe Konzept 6.5).
|
||||
|
||||
Jede Session schreibt eine eigene JSONL-Datei 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.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
from app.config import settings
|
||||
|
||||
GENESIS_HASH = "0" * 64
|
||||
|
||||
|
||||
class SessionRecorder:
|
||||
def __init__(self, session_id: int) -> None:
|
||||
self.session_id = session_id
|
||||
self.path = settings.recordings_dir / f"session_{session_id}.jsonl"
|
||||
self._prev_hash = GENESIS_HASH
|
||||
self._start_ts = time.time()
|
||||
self._fh = open(self.path, "a", encoding="utf-8")
|
||||
try:
|
||||
self.path.chmod(0o600)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
def record(self, direction: str, data: str) -> None:
|
||||
"""direction: 'input' (Tastatureingabe) oder 'output' (Terminal-/RDP-Ausgabe)."""
|
||||
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)
|
||||
entry_hash = hashlib.sha256((self._prev_hash + "|" + entry_json).encode()).hexdigest()
|
||||
line = json.dumps({"entry": entry, "prev_hash": self._prev_hash, "hash": entry_hash})
|
||||
self._fh.write(line + "\n")
|
||||
self._fh.flush()
|
||||
self._prev_hash = entry_hash
|
||||
|
||||
def close(self) -> None:
|
||||
if not self._fh.closed:
|
||||
self._fh.close()
|
||||
|
||||
|
||||
def verify_recording(path: Path) -> bool:
|
||||
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
|
||||
Reference in New Issue
Block a user