"""Dateitransfer zu SSH-Zielen per SFTP (Upload/Download ueber den Jumphost). Groessenlimit, Sha256-Hashing und optionaler AV-Scan sind Pflicht (Konzept 6.6). Jeder Transfer wird in file_transfers + audit_log protokolliert. """ from __future__ import annotations import asyncio import hashlib import logging import asyncssh from fastapi import APIRouter, Depends, HTTPException, Query, Request, UploadFile, status from fastapi.responses import StreamingResponse from app.auth.deps import CurrentUser, get_current_user from app.config import settings from app.db import get_db from app.rbac import user_has_role_for_host from app.security.audit import write_audit_event from app.security.av_scan import scan_bytes 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.sftp") router = APIRouter(prefix="/ssh", tags=["file-transfer"]) MAX_UPLOAD_BYTES = 200 * 1024 * 1024 # 200 MiB, ueber Ansible-Variable konfigurierbar (siehe Konzept) # E4 (Umsetzungsauftrag Teil E): Ein-/Auslesen in Bloecken statt in einem # einzigen file.read(MAX_UPLOAD_BYTES + 1)/remote_file.read()-Aufruf -- # sowohl fuers Streaming zum/vom SFTP-Ziel als auch fuer inkrementelles # Hashing (ein einzelner hashlib.sha256(<200 MiB>)-Aufruf haelt den # Event-Loop zwar nur kurz, aber synchron und ohne jede Zwischen-await- # Gelegenheit an -- inkrementelles Update() je Chunk verteilt das). TRANSFER_CHUNK_BYTES = 1 * 1024 * 1024 # 1 MiB # E4: Obergrenze gleichzeitig laufender Transfers -- Prozess-weiter # In-Memory-Zustand wie login_rate_limiter/active_sessions (siehe E.0), # bewusst NICHT pro Worker (es gibt ohnehin nur einen, --workers waere laut # E.0 ohnehin keine Option). _transfer_semaphore = asyncio.Semaphore(settings.max_concurrent_transfers) def _client_ip(request: Request) -> str: return request.client.host if request.client else "unknown" async def _require_file_transfer(host_id: int, request: Request, user: CurrentUser = Depends(get_current_user)): 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="file_transfer" ): raise HTTPException(status.HTTP_403_FORBIDDEN, "Keine Filetransfer-Berechtigung fuer diesen Host") host = await load_host(conn, host_id) if not host["file_transfer_enabled"]: raise HTTPException(status.HTTP_403_FORBIDDEN, "Dateitransfer ist fuer diesen Host deaktiviert") return host async def _log_transfer(conn, *, user: CurrentUser, host_id: int, client_ip: str, direction: str, filename: str, size: int, sha256: str, av_result: str) -> None: cursor = await conn.execute( "INSERT INTO sessions (user_id, host_id, protocol, client_ip, ended_at, end_reason) " "VALUES (?, ?, 'ssh', ?, strftime('%Y-%m-%dT%H:%M:%fZ','now'), 'file_transfer')", (user.id, host_id, client_ip), ) session_id = cursor.lastrowid await conn.execute( "INSERT INTO file_transfers (session_id, direction, filename, size_bytes, sha256, av_scan_result) " "VALUES (?, ?, ?, ?, ?, ?)", (session_id, direction, filename, size, sha256, av_result), ) await write_audit_event( conn, event_type="file_transfer", user_id=user.id, client_ip=client_ip, details={ "host_id": host_id, "direction": direction, "filename": filename, "size_bytes": size, "sha256": sha256, "av_scan_result": av_result, }, ) await conn.commit() @router.post("/{host_id}/files/upload") async def upload_file( host_id: int, request: Request, remote_path: str = Query(..., max_length=1024), file: UploadFile = ..., host=Depends(_require_file_transfer), user: CurrentUser = Depends(get_current_user), ): # E4 (Umsetzungsauftrag Teil E): in Bloecken statt in einem einzigen # file.read(MAX_UPLOAD_BYTES + 1)-Aufruf lesen. Der AV-Scan (scan_bytes) # braucht weiterhin den vollstaendigen Inhalt in einem Stueck -- das # Limit wird deshalb WAEHREND des Einlesens durchgesetzt (frueher # Abbruch, sobald es ueberschritten ist, statt still bis MAX+1 zu lesen), # und das Hashing laeuft inkrementell mit, statt als ein einzelner # blockierender hashlib.sha256()-Aufruf am Ende. hasher = hashlib.sha256() chunks: list[bytes] = [] total = 0 while True: chunk = await file.read(TRANSFER_CHUNK_BYTES) if not chunk: break total += len(chunk) if total > MAX_UPLOAD_BYTES: raise HTTPException(status.HTTP_413_CONTENT_TOO_LARGE, "Datei zu gross") hasher.update(chunk) chunks.append(chunk) data = b"".join(chunks) del chunks sha256 = hasher.hexdigest() async with _transfer_semaphore: # E3 (Umsetzungsauftrag Teil E): scan_bytes() ruft synchron # subprocess.run(..., timeout=30) auf -- bis zu 30 Sekunden, in denen # der EINZIGE Event-Loop der Anwendung fuer ALLE Benutzer stillstand # (siehe E.0). asyncio.to_thread() lagert den blockierenden Aufruf in # einen Worker-Thread aus. av_result = await asyncio.to_thread(scan_bytes, data) if av_result.startswith("infected"): conn = get_db() await write_audit_event( conn, event_type="file_transfer_blocked_malware", user_id=user.id, client_ip=_client_ip(request), details={"host_id": host_id, "filename": file.filename, "av_scan_result": av_result}, ) await conn.commit() raise HTTPException(status.HTTP_400_BAD_REQUEST, f"Datei durch AV-Scan blockiert: {av_result}") conn = get_db() try: ssh_conn = await connect_to_host(conn, host_id, user_id=user.id) try: async with ssh_conn.start_sftp_client() as sftp: async with sftp.open(remote_path, "wb") as remote_file: # In Bloecken schreiben statt eines einzelnen # write(): der AV-Scan (scan_bytes) # braucht den vollstaendigen Inhalt vor dem Schreiben # in einem Stueck, `data` haelt die Datei deshalb # weiterhin komplett im Speicher (unveraendert # gegenueber vorher) -- echtes speicherbegrenztes # Streaming ist beim Upload durch das "erst scannen, # dann schreiben"-Erfordernis inhaerent nicht moeglich, # ohne den AV-Scan selbst auf Datei-/Stream-Basis # umzustellen. Der chunk-weise write() vermeidet # zumindest einen einzelnen sehr grossen asyncssh- # Aufruf und haelt jeden einzelnen await kurz. for offset in range(0, len(data), TRANSFER_CHUNK_BYTES): await remote_file.write(data[offset:offset + TRANSFER_CHUNK_BYTES]) finally: ssh_conn.close() except SSH_SETUP_ERRORS as exc: # Alles, was den Verbindungsaufbau verhindert und in der Konfiguration # begruendet ist (fehlender Benutzername/Schluessel, unbrauchbarer # Schluessel, nicht gepinnter oder abweichender Host-Key, Ziel nicht # 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_result=av_result, ) return {"status": "ok", "sha256": sha256, "size": len(data), "av_scan_result": av_result} @router.get("/{host_id}/files/download") async def download_file( host_id: int, request: Request, remote_path: str = Query(..., max_length=1024), host=Depends(_require_file_transfer), user: CurrentUser = Depends(get_current_user), ): # E4 (Umsetzungsauftrag Teil E): echtes Streaming statt # remote_file.read() der kompletten Datei in einen einzigen `data`-Puffer # gefolgt von StreamingResponse(yield data) (das war effektiv KEIN # Streaming -- der gesamte Speicher-Peak war identisch zu vorher, nur in # eine andere Form verpackt). Anders als beim Upload gibt es beim # Download KEINEN AV-Scan-Zwang, der den vollstaendigen Inhalt vor dem # Weiterreichen braucht -- echtes chunkweises Streaming direkt in die # HTTP-Antwort ist hier also tatsaechlich moeglich und sinnvoll. # # Verbindungsaufbau + stat() laufen bewusst NOCH VOR der # StreamingResponse (mit vollstaendiger Fehlerbehandlung wie bisher) -- # sobald die Antwort einmal zu streamen begonnen hat, kann der # HTTP-Statuscode nicht mehr geaendert werden. sftp/remote_file bleiben # ueber die gesamte Dauer des Streams offen und werden erst im # finally-Block des Generators geschlossen. conn = get_db() await _transfer_semaphore.acquire() try: ssh_conn = await connect_to_host(conn, host_id, user_id=user.id) try: sftp = await ssh_conn.start_sftp_client() try: stat = await sftp.stat(remote_path) if stat.size and stat.size > MAX_UPLOAD_BYTES: raise HTTPException(status.HTTP_413_CONTENT_TOO_LARGE, "Datei zu gross") remote_file = await sftp.open(remote_path, "rb") except Exception: sftp.exit() raise except Exception: ssh_conn.close() raise except SSH_SETUP_ERRORS as exc: # Alles, was den Verbindungsaufbau verhindert und in der Konfiguration # begruendet ist (fehlender Benutzername/Schluessel, unbrauchbarer # Schluessel, nicht gepinnter oder abweichender Host-Key, Ziel nicht # erreichbar): 400 mit Klartext statt eines unbehandelten 500. _transfer_semaphore.release() 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. _transfer_semaphore.release() logger.warning("SFTP-Fehler bei Host %s: %s", host_id, exc) raise HTTPException(status.HTTP_400_BAD_REQUEST, f"SFTP-Fehler: {exc}") except HTTPException: _transfer_semaphore.release() raise except Exception: _transfer_semaphore.release() logger.exception("Unerwarteter Fehler beim Datei-Download fuer Host %s", host_id) raise HTTPException(status.HTTP_500_INTERNAL_SERVER_ERROR, "Unerwarteter Fehler beim Dateitransfer") # Ab hier (kein except griff) ist der Stream erfolgreich eroeffnet -- # das Semaphor wird bewusst NICHT hier freigegeben, sondern erst im # finally-Block von _stream() unten, sobald der komplette Download # abgeschlossen (oder abgebrochen) ist. filename = remote_path.rsplit("/", 1)[-1] async def _stream(): hasher = hashlib.sha256() total = 0 try: while True: chunk = await remote_file.read(TRANSFER_CHUNK_BYTES) if not chunk: break total += len(chunk) hasher.update(chunk) yield chunk finally: try: # SFTPClientFile wird andernorts in diesem Modul immer als # `async with sftp.open(...) as remote_file:` verwendet -- # das bedeutet close() ist dort eine Koroutine (__aexit__ # ruft sie awaited auf). Hier manuell dasselbe nachgebildet, # da der Dateihandle ueber die gesamte Stream-Dauer offen # bleiben muss und daher nicht in einem `async with` um nur # den Lesevorgang herum verwaltet werden kann. await remote_file.close() except Exception: pass sftp.exit() ssh_conn.close() _transfer_semaphore.release() try: await _log_transfer( conn, user=user, host_id=host_id, client_ip=_client_ip(request), direction="download", filename=filename, size=total, sha256=hasher.hexdigest(), av_result="not_applicable_download", ) except Exception: # Der Download selbst ist zu diesem Zeitpunkt beim Client # bereits (teilweise) angekommen -- ein fehlgeschlagener # Logeintrag darf den bereits laufenden Stream nicht mehr # rueckwirkend als Fehler erscheinen lassen, muss aber sichtbar # sein. logger.exception( "Download-Protokollierung fuer Host %s / %s fehlgeschlagen", host_id, remote_path ) return StreamingResponse( _stream(), media_type="application/octet-stream", headers={"Content-Disposition": f'attachment; filename="{filename}"'}, )