connectiopn fix round 2 3
This commit is contained in:
@ -1921,13 +1921,20 @@ async def list_sessions(
|
||||
result = []
|
||||
for r in rows:
|
||||
recording_path = r[11]
|
||||
is_active = r[8] is None
|
||||
killable = r[0] in active_ids
|
||||
result.append({
|
||||
"id": r[0], "user_id": r[1], "username": r[2], "host_id": r[3],
|
||||
"hostname": r[4], "host_group_name": r[5], "protocol": r[6],
|
||||
"started_at": r[7], "ended_at": r[8], "client_ip": r[9], "end_reason": r[10],
|
||||
"is_active": r[8] is None,
|
||||
"killable": r[0] in active_ids,
|
||||
"is_active": is_active,
|
||||
"killable": killable,
|
||||
"has_recording": bool(recording_path) and Path(recording_path).exists(),
|
||||
# Live-Mitschau (GET /ws/sessions/{id}/watch): wie 'killable' nur
|
||||
# moeglich, wenn die Sitzung auf DIESEM Worker-Prozess laeuft --
|
||||
# und bisher nur fuer SSH umgesetzt (RDP haette dafuer eine eigene
|
||||
# Multiplexing-Loesung fuer den guacd-Binaer-Tunnel noetig).
|
||||
"watchable": is_active and killable and r[6] == "ssh",
|
||||
})
|
||||
return result
|
||||
|
||||
|
||||
@ -14,17 +14,27 @@ 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."""
|
||||
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
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
|
||||
@dataclass
|
||||
class ActiveSession:
|
||||
session_id: int
|
||||
task: asyncio.Task
|
||||
watchers: set[asyncio.Queue] = field(default_factory=set)
|
||||
|
||||
|
||||
_active: dict[int, ActiveSession] = {}
|
||||
@ -44,3 +54,38 @@ def get(session_id: int) -> ActiveSession | None:
|
||||
|
||||
def all_ids() -> set[int]:
|
||||
return set(_active.keys())
|
||||
|
||||
|
||||
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
|
||||
|
||||
@ -8,6 +8,7 @@ from __future__ import annotations
|
||||
import hashlib
|
||||
import logging
|
||||
|
||||
import asyncssh
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, Request, UploadFile, status
|
||||
from fastapi.responses import StreamingResponse
|
||||
|
||||
@ -111,10 +112,33 @@ async def upload_file(
|
||||
# erreichbar): 400 mit Klartext statt eines unbehandelten 500.
|
||||
logger.warning("Dateitransfer fuer Host %s nicht moeglich: %s", host_id, exc)
|
||||
raise HTTPException(status.HTTP_400_BAD_REQUEST, describe_connection_error(exc))
|
||||
except asyncssh.Error as exc:
|
||||
# Bugfix: SSH_SETUP_ERRORS deckt nur Fehler VOR der Anmeldung ab
|
||||
# (Pinning/Konfiguration). asyncssh.Error ist asyncssh's gemeinsame
|
||||
# Basisklasse -- das schliesst sowohl einen echten Verbindungsfehler
|
||||
# WAEHREND asyncssh.connect() ein (z.B. Verbindung abgelehnt/abgebrochen,
|
||||
# von connect_to_host() bewusst unuebersetzt weitergereicht) als auch
|
||||
# SFTP-Fehler NACH erfolgreicher Anmeldung (z.B. SFTPNoSuchFile, wenn
|
||||
# das Zielverzeichnis nicht existiert, oder SFTPPermissionDenied). Beides
|
||||
# war hier bisher NICHT gefangen und lief unbehandelt bis Starlette
|
||||
# durch, das bei einer unbehandelten Exception (debug=False) eine
|
||||
# KLARTEXT-500-Antwort "Internal Server Error" liefert statt JSON --
|
||||
# daher der Frontend-Fehler "Unexpected token 'I', 'Internal S'... is
|
||||
# not valid JSON" beim Upload.
|
||||
logger.warning("SFTP-Fehler bei Host %s: %s", host_id, exc)
|
||||
raise HTTPException(status.HTTP_400_BAD_REQUEST, f"SFTP-Fehler: {exc}")
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception:
|
||||
# Letztes Auffangnetz (gleiches Muster wie terminal_ws.py): niemals
|
||||
# eine unbehandelte Ausnahme bis zu Starlettes Klartext-500 durchreichen
|
||||
# -- das Frontend erwartet hier immer eine JSON-Antwort.
|
||||
logger.exception("Unerwarteter Fehler beim Datei-Upload fuer Host %s", host_id)
|
||||
raise HTTPException(status.HTTP_500_INTERNAL_SERVER_ERROR, "Unerwarteter Fehler beim Dateitransfer")
|
||||
|
||||
await _log_transfer(
|
||||
conn, user=user, host_id=host_id, client_ip=_client_ip(request), direction="upload",
|
||||
filename=file.filename or remote_path, size=len(data), sha256=sha256, av_scan_result=av_result,
|
||||
filename=file.filename or remote_path, size=len(data), sha256=sha256, av_result=av_result,
|
||||
)
|
||||
return {"status": "ok", "sha256": sha256, "size": len(data), "av_scan_result": av_result}
|
||||
|
||||
@ -146,12 +170,24 @@ async def download_file(
|
||||
# erreichbar): 400 mit Klartext statt eines unbehandelten 500.
|
||||
logger.warning("Dateitransfer fuer Host %s nicht moeglich: %s", host_id, exc)
|
||||
raise HTTPException(status.HTTP_400_BAD_REQUEST, describe_connection_error(exc))
|
||||
except asyncssh.Error as exc:
|
||||
# Siehe ausfuehrlicher Kommentar in upload_file() weiter oben -- gleicher
|
||||
# Bugfix: SFTP-Fehler nach erfolgreicher Anmeldung (z.B. Datei nicht
|
||||
# gefunden, keine Leseberechtigung) liefen bisher unbehandelt bis zu
|
||||
# Starlettes Klartext-500 durch.
|
||||
logger.warning("SFTP-Fehler bei Host %s: %s", host_id, exc)
|
||||
raise HTTPException(status.HTTP_400_BAD_REQUEST, f"SFTP-Fehler: {exc}")
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception:
|
||||
logger.exception("Unerwarteter Fehler beim Datei-Download fuer Host %s", host_id)
|
||||
raise HTTPException(status.HTTP_500_INTERNAL_SERVER_ERROR, "Unerwarteter Fehler beim Dateitransfer")
|
||||
|
||||
sha256 = hashlib.sha256(data).hexdigest()
|
||||
filename = remote_path.rsplit("/", 1)[-1]
|
||||
await _log_transfer(
|
||||
conn, user=user, host_id=host_id, client_ip=_client_ip(request), direction="download",
|
||||
filename=filename, size=len(data), sha256=sha256, av_scan_result="not_applicable_download",
|
||||
filename=filename, size=len(data), sha256=sha256, av_result="not_applicable_download",
|
||||
)
|
||||
|
||||
def _iter():
|
||||
|
||||
@ -56,7 +56,9 @@ async def _reject(websocket: WebSocket, code: int, reason: str, *, accepted: boo
|
||||
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):
|
||||
async def _pump_ssh_to_ws(
|
||||
process: asyncssh.SSHClientProcess, websocket: WebSocket, recorder: SessionRecorder, session_id: int
|
||||
):
|
||||
try:
|
||||
while True:
|
||||
data = await process.stdout.read(65536)
|
||||
@ -65,6 +67,11 @@ async def _pump_ssh_to_ws(process: asyncssh.SSHClientProcess, websocket: WebSock
|
||||
if isinstance(data, str):
|
||||
data = data.encode("utf-8", errors="replace")
|
||||
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
|
||||
@ -133,7 +140,7 @@ async def ssh_terminal(websocket: WebSocket, host_id: int):
|
||||
ssh_conn = await connect_to_host(conn, host_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))
|
||||
pump_task = asyncio.create_task(_pump_ssh_to_ws(process, websocket, recorder, session_id))
|
||||
|
||||
while True:
|
||||
try:
|
||||
@ -237,3 +244,60 @@ async def ssh_terminal(websocket: WebSocket, host_id: int):
|
||||
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)
|
||||
|
||||
Reference in New Issue
Block a user