""" In-Memory-Registry der laufenden SSH/RDP-Sitzungen dieses Server-Prozesses. Wird von app/ssh_proxy/terminal_ws.py und app/rdp_proxy/ws_tunnel.py beim Sitzungsstart befuellt (Referenz auf den eigenen asyncio.Task) und beim Sitzungsende wieder entfernt. Basis fuer die Superadmin-'Sessionview' (app/admin/routes.py: GET /admin/sessions, POST /admin/sessions/{id}/terminate) -- 'Beenden' funktioniert per asyncio.Task.cancel(), was den regulaeren finally-Cleanup-Pfad der jeweiligen WS-Route ausloest (DB-Update, Audit- Log-Eintrag, WebSocket schliessen), genau wie bei einem normalen Logout. Bewusst rein prozesslokal (kein Redis/DB-Backing): bei einem einzelnen uvicorn-Worker (Standard-Deployment dieses Projekts, siehe ansible/) sieht und beendet jeder Superadmin-Request alle laufenden Sitzungen. Bei mehreren Worker-Prozessen sind nur die Sitzungen DIESES Workers 'killable' -- die Sessionview zeigt das ueber das 'killable'-Feld pro Sitzung an, statt einen Beenden-Versuch fehlschlagen zu lassen. Seit der Live-Mitschau ("Aktive Session live ansehen", GET /ws/sessions/ {id}/watch, app/ssh_proxy/terminal_ws.py) traegt jeder Registry-Eintrag zusaetzlich eine Menge von Beobachter-Queues: app/ssh_proxy/terminal_ws.py:: _pump_ssh_to_ws() speist jeden vom Zielsystem gelesenen Byte-Chunk sowohl an den eigentlichen Sitzungsinhaber als auch (per broadcast()) an alle gerade zuschauenden Superadmins. Genau wie 'killable' gilt: nur auf demselben Worker-Prozess moeglich -- ein Beobachter kann sich nur an eine Sitzung haengen, die in DIESER Prozess-Registry steht.""" from __future__ import annotations import asyncio from dataclasses import dataclass, field @dataclass class ActiveSession: session_id: int task: asyncio.Task user_id: int watchers: set[asyncio.Queue] = field(default_factory=set) _active: dict[int, ActiveSession] = {} def register(session_id: int, task: asyncio.Task, user_id: int) -> None: _active[session_id] = ActiveSession(session_id=session_id, task=task, user_id=user_id) def unregister(session_id: int) -> None: _active.pop(session_id, None) def get(session_id: int) -> ActiveSession | None: return _active.get(session_id) def all_ids() -> set[int]: return set(_active.keys()) def count_total() -> int: """Anzahl aller laufenden Sitzungen auf diesem Prozess (E.4, Umsetzungsauftrag Teil E: globale Obergrenze).""" return len(_active) def count_for_user(user_id: int) -> int: """Anzahl laufender Sitzungen eines Benutzers auf diesem Prozess (E.4: Obergrenze je Benutzer). Bewusst prozesslokal wie die gesamte Registry (siehe Moduldoc) -- bei einem einzelnen uvicorn-Worker (Standard- Deployment) ist das gleichbedeutend mit "insgesamt", bei mehreren Workern zaehlt es nur die auf DIESEM Worker laufenden Sitzungen dieses Benutzers.""" return sum(1 for entry in _active.values() if entry.user_id == user_id) def add_watcher(session_id: int) -> asyncio.Queue | None: """Meldet einen Beobachter fuer eine laufende Sitzung an. Gibt None zurueck, wenn die Sitzung nicht (mehr) auf diesem Prozess laeuft -- der Aufrufer (WS-Route) beendet die Beobachter-Verbindung dann sofort.""" entry = _active.get(session_id) if entry is None: return None queue: asyncio.Queue = asyncio.Queue(maxsize=1000) entry.watchers.add(queue) return queue def remove_watcher(session_id: int, queue: asyncio.Queue) -> None: entry = _active.get(session_id) if entry is not None: entry.watchers.discard(queue) def broadcast(session_id: int, data: bytes) -> None: """Verteilt einen Ausgabe-Chunk an alle aktuell zuschauenden Beobachter dieser Sitzung. Bewusst best-effort wie log_stream.py: ein langsamer oder haengender Beobachter darf weder die eigentliche Sitzung verlangsamen noch blockieren -- volle Queues verwerfen den Chunk stattdessen (put_nowait/QueueFull), statt den Sitzungsinhaber warten zu lassen.""" entry = _active.get(session_id) if entry is None or not entry.watchers: return for queue in list(entry.watchers): try: queue.put_nowait(data) except asyncio.QueueFull: pass