""" Browser-Terminal <-> SSH-Ziel per WebSocket (xterm.js-kompatibel). Framing: JSON-Textframes. Client -> Server: {"type":"input","data":""} | {"type":"resize","cols":n,"rows":n} Server -> Client: {"type":"output","data":""} | {"type":"error","message":"..."} | {"type":"closed"} Jede Session wird aufgezeichnet (app.recordings.recorder) und im Audit-Log mit Start/Ende vermerkt (Konzept 4.2, 4.7, 6.5). """ from __future__ import annotations import asyncio import base64 import logging import time import asyncssh from fastapi import APIRouter, WebSocket, WebSocketDisconnect from app.auth.deps import get_current_user_ws from app.config import settings from app.db import get_db from app.rbac import user_has_role_for_host from app.recordings.recorder import SessionRecorder from app.security import active_sessions from app.security.audit import write_audit_event from app.ssh_proxy.proxy import ( SSH_SETUP_ERRORS, HostNotConfiguredError, PrivateKeyUnusableError, connect_to_host, describe_connection_error, load_host, ) logger = logging.getLogger("jumphost.ssh_proxy.ws") router = APIRouter() # E10/E11 (Umsetzungsauftrag Teil E): frueher feste Modul-Konstanten, jetzt # konfigurierbar (siehe app/config.py). MAX_SESSION_SECONDS wurde vorher # nirgends ausgewertet -- es gab de facto GAR KEINE absolute Obergrenze. IDLE_TIMEOUT_SECONDS = settings.ssh_idle_timeout_s MAX_SESSION_SECONDS = settings.ssh_max_session_seconds MAX_SESSION_WARNING_S = settings.ssh_max_session_warning_s # Wie oft die Haupt-Schleife hoechstens "blind" auf eine Benutzereingabe # wartet, bevor sie den gemeinsamen Aktivitaets-/Laufzeitstand neu prueft. # Niedrig genug, um Idle-Timeout und Sitzungsobergrenze zeitnah durchsetzen # zu koennen, aber hoch genug, um nicht sinnlos oft zu pollen. _POLL_INTERVAL_S = 20 async def _reject(websocket: WebSocket, code: int, reason: str, *, accepted: bool) -> None: """Beendet eine SSH-Sitzung vor ihrem eigentlichen Beginn -- mit einem fuer den Benutzer lesbaren Grund als WebSocket-Close-Reason, analog zu app/rdp_proxy/ws_tunnel.py::_reject (Troubleshooting-Verbesserung: vorher endeten diese Pfade in einem nackten `websocket.close(code=...)`, static/js/terminal.js zeigte dann nur ein generisches 'Verbindung beendet' ohne jeden Grund an). Voraussetzung fuer eine sichtbare Reason ist ein zustande gekommener Handshake -- vor accept() sieht der Browser nur einen HTTP-/WS-Fehler ohne Text (siehe die bewusste Ausnahme fuer den Nicht-angemeldet-Fall unten).""" logger.warning("SSH-Verbindung abgelehnt (code=%s): %s", code, reason) reason_bytes = reason.encode("utf-8")[:123] if not accepted: await websocket.accept() await websocket.close(code=code, reason=reason_bytes.decode("utf-8", errors="ignore")) async def _pump_ssh_to_ws( process: asyncssh.SSHClientProcess, websocket: WebSocket, recorder: SessionRecorder, session_id: int, activity: dict, ): try: while True: data = await process.stdout.read(65536) if not data: break if isinstance(data, str): data = data.encode("utf-8", errors="replace") # E10 (Umsetzungsauftrag Teil E): Ausgabe vom Ziel zaehlt genauso # als Aktivitaet wie eine Benutzereingabe -- ein laufendes # `tail -f`/langer Build haelt die Sitzung damit am Leben, auch # wenn niemand tippt. `activity` wird mit der Haupt-Schleife # unten geteilt (dieselbe Coroutine-Ausfuehrung, kein Lock noetig). activity["t"] = time.monotonic() recorder.record("output", base64.b64encode(data).decode()) # Live-Mitschau (GET /ws/sessions/{id}/watch, siehe unten): jeder # Chunk geht zusaetzlich an alle aktuell zuschauenden Superadmins. # Best-effort/read-only -- ein Beobachter kann diese Sitzung nicht # beeinflussen, und sein Fehlen/Trennen hat keinen Einfluss hierauf. active_sessions.broadcast(session_id, data) await websocket.send_json({"type": "output", "data": base64.b64encode(data).decode()}) except (asyncssh.Error, ConnectionResetError): pass @router.websocket("/ws/ssh/{host_id}") async def ssh_terminal(websocket: WebSocket, host_id: int): user = await get_current_user_ws(websocket) if user is None: # Auch die Ablehnungen VOR dem Sitzungsbeginn gehoeren protokolliert: # bisher schloss dieser Pfad das Socket kommentarlos, im # Verbindungslog war der Fehlversuch dadurch unsichtbar. logger.warning("SSH-Verbindung abgelehnt: keine gueltige Sitzung (host_id=%s)", host_id) await websocket.close(code=4401) return conn = get_db() if not user.is_admin and not await user_has_role_for_host( conn, user_id=user.id, host_id=host_id, role_name="ssh_connect" ): await _reject( websocket, 4403, f"Keine Berechtigung 'ssh_connect' fuer Host {host_id}", accepted=False, ) return # E.4 (Umsetzungsauftrag Teil E): Obergrenzen je Benutzer und global, # VOR jedem Ressourcenverbrauch geprueft (analog app/rdp_proxy/ws_tunnel.py). if not user.is_admin and active_sessions.count_for_user(user.id) >= settings.max_sessions_per_user: await _reject( websocket, 4429, f"Sie haben bereits {settings.max_sessions_per_user} Sitzungen offen " "(Obergrenze je Benutzer erreicht).", accepted=False, ) return if active_sessions.count_total() >= settings.max_sessions_global: await _reject( websocket, 4429, "Der Server hat die maximale Anzahl gleichzeitiger Sitzungen erreicht, " "bitte spaeter erneut versuchen.", accepted=False, ) return await websocket.accept() client_ip = websocket.client.host if websocket.client else "unknown" try: host = await load_host(conn, host_id) except HostNotConfiguredError as exc: # Grund SOWOHL als JSON-Frame (falls der Client schon zuhoert) ALS # AUCH als Close-Reason senden (falls nicht) -- terminal.js zeigt # beides an, je nachdem, was zuerst ankommt. await websocket.send_json({"type": "error", "message": str(exc)}) await _reject(websocket, 4404, str(exc), accepted=True) return cursor = await conn.execute( "INSERT INTO sessions (user_id, host_id, protocol, client_ip) VALUES (?, ?, 'ssh', ?)", (user.id, host_id, client_ip), ) session_id = cursor.lastrowid recorder = SessionRecorder(session_id) await conn.execute( "UPDATE sessions SET recording_path = ? WHERE id = ?", (str(recorder.path), session_id) ) await write_audit_event( conn, event_type="ssh_session_start", user_id=user.id, client_ip=client_ip, details={"host_id": host_id, "hostname": host["hostname"], "session_id": session_id}, ) await conn.commit() logger.debug( "SSH-Sitzung %s gestartet: user=%s host=%s (%s:%s) client_ip=%s", session_id, user.username, host["hostname"], host["address"], host["port"], client_ip, ) active_sessions.register(session_id, asyncio.current_task(), user.id) end_reason = "logout" ssh_conn = None process = None pump_task = None session_start = time.monotonic() # E10: von _pump_ssh_to_ws() UND der Haupt-Schleife hier gemeinsam # aktualisiert -- ein dict-Eintrag statt einer einfachen Variable, damit # beide Coroutinen denselben veraenderlichen Zustand sehen (Closures # koennen keine Nicht-lokalen einfachen Namen neu binden). Kein Lock # noetig: reine Zuweisungen im selben Event-Loop-Thread. activity = {"t": session_start} warned_max_duration = False try: ssh_conn = await connect_to_host(conn, host_id, user_id=user.id) logger.debug("SSH-Sitzung %s: Verbindung zu %s hergestellt", session_id, host["hostname"]) process = await ssh_conn.create_process(term_type="xterm-256color") pump_task = asyncio.create_task(_pump_ssh_to_ws(process, websocket, recorder, session_id, activity)) while True: try: msg = await asyncio.wait_for(websocket.receive_json(), timeout=_POLL_INTERVAL_S) except asyncio.TimeoutError: msg = None now = time.monotonic() if msg is not None: activity["t"] = now # E10: Inaktivitaet wird ueber `activity` in BEIDEN Richtungen # gemessen (Benutzereingabe hier, Zielausgabe in # _pump_ssh_to_ws()) -- vorher zaehlte nur receive_json(), ein # rein ausgabelastiges `tail -f` o.ae. wurde nach 15 Minuten ohne # Tastendruck getrennt, obwohl die Sitzung erkennbar aktiv war. if now - activity["t"] > IDLE_TIMEOUT_SECONDS: end_reason = "idle_timeout" try: await websocket.send_json({ "type": "error", "message": ( f"Sitzung wegen Inaktivitaet beendet " f"(> {IDLE_TIMEOUT_SECONDS // 60} Minuten ohne Ein-/Ausgabe)." ), }) except Exception: logger.debug("Idle-Timeout-Meldung konnte nicht mehr gesendet werden", exc_info=True) break # E11: MAX_SESSION_SECONDS war vorher eine definierte, aber # nirgends ausgewertete Konstante -- es gab de facto GAR KEINE # absolute Obergrenze fuer eine SSH-Sitzung. Jetzt aktiv # durchgesetzt, mit Vorwarnung statt eines ueberraschenden # sofortigen Abbruchs. running_for = now - session_start if running_for > MAX_SESSION_SECONDS: end_reason = "max_duration_exceeded" try: await websocket.send_json({ "type": "error", "message": ( f"Sitzung nach Erreichen der maximalen Sitzungsdauer " f"({MAX_SESSION_SECONDS // 3600} Stunden) beendet." ), }) except Exception: logger.debug("Sitzungsende-Meldung konnte nicht mehr gesendet werden", exc_info=True) break if not warned_max_duration and running_for > MAX_SESSION_SECONDS - MAX_SESSION_WARNING_S: warned_max_duration = True try: await websocket.send_json({ "type": "warning", "message": ( f"Diese Sitzung wird in ca. {MAX_SESSION_WARNING_S} Sekunden wegen der " "maximalen Sitzungsdauer automatisch beendet." ), }) except Exception: logger.debug("Vorwarnung konnte nicht mehr gesendet werden", exc_info=True) if msg is None: continue if msg.get("type") == "input": raw = base64.b64decode(msg.get("data", "")) recorder.record("input", base64.b64encode(raw).decode()) process.stdin.write(raw.decode("utf-8", errors="replace")) elif msg.get("type") == "resize": cols, rows = int(msg.get("cols", 80)), int(msg.get("rows", 24)) process.change_terminal_size(cols, rows) except WebSocketDisconnect: end_reason = "logout" except asyncssh.Error as exc: logger.warning("SSH-Sessionfehler (session_id=%s): %s", session_id, exc) end_reason = "error" try: await websocket.send_json({"type": "error", "message": "Verbindung zum Zielsystem fehlgeschlagen"}) except Exception: # Best-Effort-Fehlermeldung an einen ggf. bereits getrennten Client; # der eigentliche Fehler ist bereits oben geloggt (logger.warning). logger.debug("Fehlermeldung konnte nicht mehr an Client gesendet werden", exc_info=True) except PrivateKeyUnusableError as exc: # Der haeufigste Grund, warum eine SSH-Sitzung nie zustande kam: der # hinterlegte Private Key ist passphrasegeschuetzt und die Passphrase # fehlt oder passt nicht. asyncssh wirft dafuer einen KeyImportError # (ein ValueError, KEIN asyncssh.Error), der frueher an allen # Handlern vorbei bis aus der Route hinauslief -- der Browser sah nur # ein wortloses Verbindungsende. Der Text ist bewusst konkret und # nennt die Stelle im Adminbereich, an der es zu beheben ist. logger.warning("SSH-Sitzung %s: Schluessel unbrauchbar: %s", session_id, exc) end_reason = "error" try: await websocket.send_json({"type": "error", "message": str(exc)}) except Exception: logger.debug("Fehlermeldung konnte nicht mehr an Client gesendet werden", exc_info=True) except SSH_SETUP_ERRORS as exc: # Alles, was den Sitzungsaufbau verhindert und der Benutzer selbst # einordnen kann: fehlender Benutzername/Schluessel, nicht gepinnter # oder abweichender Host-Key, unerreichbares Ziel. Bewusst mit # Klartextmeldung statt "interner Fehler" -- der Grund steht in der # Konfiguration, nicht im Code. logger.warning("SSH-Sitzung %s abgebrochen: %s", session_id, exc) end_reason = "error" try: await websocket.send_json( {"type": "error", "message": describe_connection_error(exc)} ) except Exception: logger.debug("Fehlermeldung konnte nicht mehr an Client gesendet werden", exc_info=True) except asyncio.CancelledError: # Zwangs-Beendigung durch einen Superadmin ueber die Sessionview # (POST /admin/sessions/{id}/terminate, siehe app/security/active_sessions.py). end_reason = "terminated_by_admin" raise except Exception as exc: # Auffangnetz fuer alles, was NICHT asyncssh.Error ist (z.B. # asyncssh.KeyImportError bei einem defekten/nicht mehr passenden # Schluessel, oder cryptography.exceptions.InvalidTag bei # decrypt_secret() falls der KEK nicht mehr zum verschluesselten # Schluessel passt) -- diese Typen sind KEINE asyncssh.Error- # Unterklassen und liefen bisher unbehandelt durch, wodurch die # Sitzung ohne jede Fehlermeldung abrupt getrennt wurde und im Log # irrefuehrend als reason=logout erschien (Default-Wert, der nie # aktualisiert wurde). Mit logger.exception landet der volle # Traceback jetzt im Verbindungslog. logger.exception("Unerwarteter Fehler in SSH-Sitzung %s: %s", session_id, exc) end_reason = "error" try: await websocket.send_json({ "type": "error", "message": f"Verbindung zum Zielsystem fehlgeschlagen (interner Fehler: {exc.__class__.__name__})", }) except Exception: logger.debug("Fehlermeldung konnte nicht mehr an Client gesendet werden", exc_info=True) finally: active_sessions.unregister(session_id) if pump_task: pump_task.cancel() if process: process.close() if ssh_conn: ssh_conn.close() await recorder.aclose() logger.debug("SSH-Sitzung %s beendet: reason=%s", session_id, end_reason) await conn.execute( "UPDATE sessions SET ended_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'), end_reason = ? " "WHERE id = ?", (end_reason, session_id), ) await write_audit_event( conn, event_type="ssh_session_end", user_id=user.id, client_ip=client_ip, details={"host_id": host_id, "session_id": session_id, "reason": end_reason}, ) await conn.commit() try: await websocket.close() except Exception: logger.debug("WebSocket war beim Schliessen bereits getrennt", exc_info=True) @router.websocket("/ws/sessions/{session_id}/watch") async def watch_ssh_session(websocket: WebSocket, session_id: int): """Live-Mitschau einer laufenden SSH-Sitzung fuer Superadmins ('Aktive Session mitschauen', Sessionview im Adminbereich). Rein lesend: der Beobachter bekommt exakt dieselben "output"-Frames wie der Sitzungsinhaber selbst (siehe _pump_ssh_to_ws()/active_sessions. broadcast()), kann aber keine Eingaben in die fremde Sitzung schicken -- dafuer gibt es hier bewusst keinen "input"-Pfad. Eigener Top-Level-Router unter /ws/... (siehe app/admin/log_ws.py fuer die ausfuehrliche Begruendung: der nginx-Reverse-Proxy setzt WebSocket-Upgrade-Header nur fuer /ws/-Pfade). Nur auf demselben Server-Prozess moeglich wie 'Beenden' -- active_sessions.add_watcher() gibt None zurueck, wenn die Sitzung hier nicht (mehr) laeuft.""" user = await get_current_user_ws(websocket) if user is None or not user.is_admin: await websocket.close(code=4403) return conn = get_db() row = await (await conn.execute( "SELECT protocol, ended_at FROM sessions WHERE id = ?", (session_id,) )).fetchone() if row is None or row[0] != "ssh" or row[1] is not None: await websocket.close(code=4404) return queue = active_sessions.add_watcher(session_id) if queue is None: # Sitzung existiert, laeuft aber nicht (mehr) auf diesem Worker- # Prozess -- analog zum 'killable'-Feld bei 'Beenden'. await websocket.close(code=4409) return await websocket.accept() client_ip = websocket.client.host if websocket.client else "unknown" await write_audit_event( conn, event_type="session_watch_started", user_id=user.id, client_ip=client_ip, details={"session_id": session_id}, ) await conn.commit() logger.debug("Superadmin %s beobachtet SSH-Sitzung %s live", user.username, session_id) try: while True: data = await queue.get() await websocket.send_json({"type": "output", "data": base64.b64encode(data).decode()}) except WebSocketDisconnect: pass finally: active_sessions.remove_watcher(session_id, queue) try: await websocket.close() except Exception: logger.debug("Watch-WebSocket war beim Schliessen bereits getrennt", exc_info=True)