315 lines
15 KiB
Python
315 lines
15 KiB
Python
"""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(<gesamte Datei>)-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(<gesamte Datei>): 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}"'},
|
|
)
|