Creare un'app in stile WeTransfer con Python

Creare un'app in stile WeTransfer con Python

WeTransfer ha reso popolare un'idea molto semplice: caricare uno o più file, ottenere un link, condividerlo, e lasciare che il link scada dopo qualche giorno. Dietro questa semplicità apparente si nascondono però diversi problemi tecnici interessanti: la gestione di upload di grandi dimensioni senza saturare la memoria del server, la generazione di identificatori non indovinabili, la scadenza automatica dei contenuti, la creazione di archivi ZIP al volo, il rate limiting e la protezione da abusi.

In questo articolo costruiremo un servizio completo di file transfer con Python, usando FastAPI come framework web, SQLAlchemy per la persistenza dei metadati, uno storage layer astratto (filesystem locale o S3-compatibile), Redis con ARQ per i task in background e Docker Compose per l'orchestrazione. Il risultato sarà un'applicazione pronta per essere messa dietro un reverse proxy Nginx.

Architettura del sistema

Prima di scrivere codice conviene chiarire il flusso funzionale. Un trasferimento (che chiameremo transfer) è un aggregato che contiene uno o più file. Il ciclo di vita è il seguente:

  1. Il client crea un transfer inviando i metadati (nomi dei file, dimensioni, destinatari, messaggio).
  2. Il server risponde con un identificatore del transfer e, per ogni file, un identificatore di upload.
  3. Il client carica il contenuto binario di ciascun file.
  4. Il client finalizza il transfer: il server verifica che tutti i file siano completi, genera lo slug pubblico e la data di scadenza, e invia le email.
  5. Chiunque possieda il link può scaricare i singoli file o l'intero pacchetto come ZIP, finché il transfer non scade.
  6. Un worker periodico cancella i transfer scaduti sia dal database sia dallo storage.

Questa separazione tra creazione, upload e finalizzazione è ciò che permette di gestire upload di svariati gigabyte in modo resiliente: se la connessione cade a metà del terzo file, il client può riprendere senza ricominciare da zero.

I componenti sono quindi:

  • API HTTP (FastAPI + Uvicorn): espone gli endpoint di creazione, upload, finalizzazione e download.
  • Database relazionale (PostgreSQL): conserva solo i metadati, mai i byte dei file.
  • Object storage: filesystem locale in sviluppo, MinIO o S3 in produzione.
  • Coda di task (Redis + ARQ): invio email e pulizia dei transfer scaduti.
  • Reverse proxy (Nginx): terminazione TLS, limiti sul body, eventuale accelerazione del download.

Requisiti e struttura del progetto

Il progetto richiede Python 3.12 o superiore. Creiamo l'ambiente virtuale e installiamo le dipendenze:

python3.12 -m venv .venv
source .venv/bin/activate
pip install "fastapi[standard]" "uvicorn[standard]" "sqlalchemy[asyncio]" asyncpg \
    alembic pydantic-settings python-multipart aiofiles arq redis \
    boto3 aiosmtplib jinja2 slowapi zipstream-ng pytest pytest-asyncio httpx

La struttura delle directory che adotteremo è la seguente:

filedrop/
├── app/
│   ├── __init__.py
│   ├── main.py
│   ├── config.py
│   ├── database.py
│   ├── models.py
│   ├── schemas.py
│   ├── security.py
│   ├── storage/
│   │   ├── __init__.py
│   │   ├── base.py
│   │   ├── local.py
│   │   └── s3.py
│   ├── services/
│   │   ├── __init__.py
│   │   ├── transfers.py
│   │   ├── mailer.py
│   │   └── packaging.py
│   ├── routers/
│   │   ├── __init__.py
│   │   ├── transfers.py
│   │   └── downloads.py
│   └── workers.py
├── migrations/
├── templates/
│   └── email_transfer.html
├── tests/
├── docker-compose.yml
├── Dockerfile
└── pyproject.toml

Configurazione

Usiamo pydantic-settings per centralizzare la configurazione e leggerla da variabili d'ambiente. Questo evita di disseminare os.environ nel codice e fornisce validazione gratuita all'avvio.

from functools import lru_cache
from pathlib import Path
from typing import Literal

from pydantic import Field
from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    model_config = SettingsConfigDict(env_file=".env", env_prefix="FILEDROP_")

    app_name: str = "FileDrop"
    base_url: str = "http://localhost:8000"
    debug: bool = False

    database_url: str = "postgresql+asyncpg://filedrop:filedrop@localhost/filedrop"
    redis_url: str = "redis://localhost:6379/0"

    # Backend di storage: locale in sviluppo, S3 in produzione
    storage_backend: Literal["local", "s3"] = "local"
    storage_path: Path = Path("./data/uploads")

    s3_endpoint_url: str | None = None
    s3_bucket: str = "filedrop"
    s3_access_key: str | None = None
    s3_secret_key: str | None = None
    s3_region: str = "eu-central-1"

    # Limiti applicativi
    max_file_size: int = 2 * 1024 * 1024 * 1024      # 2 GB per singolo file
    max_transfer_size: int = 5 * 1024 * 1024 * 1024  # 5 GB per transfer
    max_files_per_transfer: int = 50
    chunk_size: int = 1024 * 1024                    # 1 MB letti alla volta
    default_expiry_days: int = 7
    max_expiry_days: int = 30

    smtp_host: str = "localhost"
    smtp_port: int = 587
    smtp_user: str | None = None
    smtp_password: str | None = None
    smtp_from: str = "noreply@example.com"
    smtp_starttls: bool = True

    secret_key: str = Field(min_length=32)


@lru_cache
def get_settings() -> Settings:
    # La cache garantisce che il file .env venga letto una sola volta
    return Settings()

Il modello dei dati

Due tabelle bastano: transfers e transfer_files. La prima rappresenta il pacchetto condiviso, la seconda i singoli file al suo interno.

import enum
import uuid
from datetime import datetime, timezone

from sqlalchemy import (
    BigInteger, Boolean, DateTime, Enum, ForeignKey, Index, String, Text, func
)
from sqlalchemy.dialects.postgresql import UUID as PgUUID
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship


class Base(DeclarativeBase):
    pass


class TransferStatus(str, enum.Enum):
    PENDING = "pending"      # creato, upload non ancora completato
    READY = "ready"          # finalizzato e scaricabile
    EXPIRED = "expired"      # scaduto, in attesa di cancellazione
    DELETED = "deleted"      # contenuto rimosso dallo storage


class Transfer(Base):
    __tablename__ = "transfers"

    id: Mapped[uuid.UUID] = mapped_column(
        PgUUID(as_uuid=True), primary_key=True, default=uuid.uuid4
    )
    # Slug pubblico presente nell'URL condiviso
    slug: Mapped[str] = mapped_column(String(32), unique=True, index=True)
    # Token segreto che autorizza upload e cancellazione da parte del mittente
    owner_token: Mapped[str] = mapped_column(String(64))
    status: Mapped[TransferStatus] = mapped_column(
        Enum(TransferStatus), default=TransferStatus.PENDING, index=True
    )

    sender_email: Mapped[str | None] = mapped_column(String(320))
    recipients: Mapped[str | None] = mapped_column(Text)  # lista separata da virgole
    message: Mapped[str | None] = mapped_column(Text)

    # Hash della password opzionale che protegge il download
    password_hash: Mapped[str | None] = mapped_column(String(255))

    total_size: Mapped[int] = mapped_column(BigInteger, default=0)
    download_count: Mapped[int] = mapped_column(BigInteger, default=0)
    max_downloads: Mapped[int | None] = mapped_column(BigInteger)

    created_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=True), server_default=func.now()
    )
    expires_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True)

    files: Mapped[list["TransferFile"]] = relationship(
        back_populates="transfer", cascade="all, delete-orphan", lazy="selectin"
    )

    @property
    def is_expired(self) -> bool:
        return datetime.now(timezone.utc) >= self.expires_at

    @property
    def is_exhausted(self) -> bool:
        # Un transfer può avere un numero massimo di download consentiti
        return (
            self.max_downloads is not None
            and self.download_count >= self.max_downloads
        )


class TransferFile(Base):
    __tablename__ = "transfer_files"

    id: Mapped[uuid.UUID] = mapped_column(
        PgUUID(as_uuid=True), primary_key=True, default=uuid.uuid4
    )
    transfer_id: Mapped[uuid.UUID] = mapped_column(
        ForeignKey("transfers.id", ondelete="CASCADE"), index=True
    )
    filename: Mapped[str] = mapped_column(String(255))
    content_type: Mapped[str] = mapped_column(
        String(127), default="application/octet-stream"
    )
    declared_size: Mapped[int] = mapped_column(BigInteger)
    stored_size: Mapped[int] = mapped_column(BigInteger, default=0)
    # Chiave dell'oggetto nello storage, indipendente dal nome originale
    object_key: Mapped[str] = mapped_column(String(255), unique=True)
    checksum: Mapped[str | None] = mapped_column(String(64))
    is_complete: Mapped[bool] = mapped_column(Boolean, default=False)

    transfer: Mapped[Transfer] = relationship(back_populates="files")


# Indice composto usato dal worker di pulizia
Index("ix_transfers_status_expires", Transfer.status, Transfer.expires_at)

Vale la pena notare due scelte progettuali. La prima: il nome del file mostrato all'utente e la chiave nello storage sono campi distinti. Questo elimina alla radice ogni rischio di path traversal, perché il nome fornito dal client non tocca mai il filesystem. La seconda: declared_size è quanto il client dichiara, stored_size è quanto il server ha effettivamente scritto; confrontarli in fase di finalizzazione permette di individuare upload troncati.

Il livello di storage

Astraiamo lo storage dietro un protocollo, così da poter passare dal filesystem locale a S3 senza toccare la logica applicativa.

from typing import AsyncIterator, Protocol


class StorageBackend(Protocol):
    async def save(self, key: str, stream: AsyncIterator[bytes]) -> int:
        """Scrive lo stream e restituisce il numero di byte salvati."""
        ...

    def open(
        self, key: str, start: int = 0, end: int | None = None
    ) -> AsyncIterator[bytes]:
        """Restituisce un iteratore sui byte dell'oggetto, con supporto ai range."""
        ...

    async def size(self, key: str) -> int: ...

    async def delete(self, key: str) -> None: ...

    async def exists(self, key: str) -> bool: ...

L'implementazione locale usa aiofiles per non bloccare il loop di eventi durante l'I/O su disco:

from pathlib import Path
from typing import AsyncIterator

import aiofiles
import aiofiles.os


class LocalStorage:
    def __init__(self, root: Path, chunk_size: int = 1024 * 1024) -> None:
        self.root = root.resolve()
        self.chunk_size = chunk_size
        self.root.mkdir(parents=True, exist_ok=True)

    def _path(self, key: str) -> Path:
        # Le chiavi sono generate dal server e contengono solo esadecimali e slash,
        # ma normalizziamo comunque per difesa in profondità
        target = (self.root / key).resolve()
        if not target.is_relative_to(self.root):
            raise ValueError("Chiave di storage non valida")
        return target

    async def save(self, key: str, stream: AsyncIterator[bytes]) -> int:
        path = self._path(key)
        await aiofiles.os.makedirs(path.parent, exist_ok=True)
        written = 0
        async with aiofiles.open(path, "wb") as handle:
            async for chunk in stream:
                await handle.write(chunk)
                written += len(chunk)
        return written

    async def open(
        self, key: str, start: int = 0, end: int | None = None
    ) -> AsyncIterator[bytes]:
        path = self._path(key)
        remaining = None if end is None else (end - start + 1)
        async with aiofiles.open(path, "rb") as handle:
            await handle.seek(start)
            while True:
                size = self.chunk_size
                if remaining is not None:
                    if remaining <= 0:
                        break
                    size = min(size, remaining)
                chunk = await handle.read(size)
                if not chunk:
                    break
                if remaining is not None:
                    remaining -= len(chunk)
                yield chunk

    async def size(self, key: str) -> int:
        stat = await aiofiles.os.stat(self._path(key))
        return stat.st_size

    async def delete(self, key: str) -> None:
        try:
            await aiofiles.os.remove(self._path(key))
        except FileNotFoundError:
            # La cancellazione deve essere idempotente
            pass

    async def exists(self, key: str) -> bool:
        return await aiofiles.os.path.exists(self._path(key))

Il backend S3 sfrutta invece il multipart upload di boto3. Poiché boto3 è sincrono, deleghiamo le chiamate a un thread pool con asyncio.to_thread per non bloccare il loop:

import asyncio
from typing import AsyncIterator

import boto3
from botocore.config import Config
from botocore.exceptions import ClientError

MIN_PART_SIZE = 5 * 1024 * 1024  # S3 impone parti di almeno 5 MB, tranne l'ultima


class S3Storage:
    def __init__(
        self,
        bucket: str,
        endpoint_url: str | None = None,
        access_key: str | None = None,
        secret_key: str | None = None,
        region: str = "eu-central-1",
    ) -> None:
        self.bucket = bucket
        self.client = boto3.client(
            "s3",
            endpoint_url=endpoint_url,
            aws_access_key_id=access_key,
            aws_secret_access_key=secret_key,
            region_name=region,
            config=Config(signature_version="s3v4", retries={"max_attempts": 3}),
        )

    async def save(self, key: str, stream: AsyncIterator[bytes]) -> int:
        upload = await asyncio.to_thread(
            self.client.create_multipart_upload, Bucket=self.bucket, Key=key
        )
        upload_id = upload["UploadId"]
        parts: list[dict] = []
        buffer = bytearray()
        part_number = 1
        written = 0

        try:
            async for chunk in stream:
                buffer.extend(chunk)
                written += len(chunk)
                # Accumuliamo finché non raggiungiamo la dimensione minima di parte
                if len(buffer) >= MIN_PART_SIZE:
                    parts.append(
                        await self._upload_part(key, upload_id, part_number, bytes(buffer))
                    )
                    part_number += 1
                    buffer.clear()

            if buffer or not parts:
                parts.append(
                    await self._upload_part(key, upload_id, part_number, bytes(buffer))
                )

            await asyncio.to_thread(
                self.client.complete_multipart_upload,
                Bucket=self.bucket,
                Key=key,
                UploadId=upload_id,
                MultipartUpload={"Parts": parts},
            )
        except Exception:
            # Un upload multipart interrotto continua a occupare spazio: va abortito
            await asyncio.to_thread(
                self.client.abort_multipart_upload,
                Bucket=self.bucket,
                Key=key,
                UploadId=upload_id,
            )
            raise

        return written

    async def _upload_part(
        self, key: str, upload_id: str, number: int, body: bytes
    ) -> dict:
        response = await asyncio.to_thread(
            self.client.upload_part,
            Bucket=self.bucket,
            Key=key,
            UploadId=upload_id,
            PartNumber=number,
            Body=body,
        )
        return {"ETag": response["ETag"], "PartNumber": number}

    async def open(
        self, key: str, start: int = 0, end: int | None = None
    ) -> AsyncIterator[bytes]:
        range_header = f"bytes={start}-" if end is None else f"bytes={start}-{end}"
        response = await asyncio.to_thread(
            self.client.get_object, Bucket=self.bucket, Key=key, Range=range_header
        )
        body = response["Body"]
        while True:
            chunk = await asyncio.to_thread(body.read, 1024 * 1024)
            if not chunk:
                break
            yield chunk

    async def size(self, key: str) -> int:
        head = await asyncio.to_thread(
            self.client.head_object, Bucket=self.bucket, Key=key
        )
        return head["ContentLength"]

    async def delete(self, key: str) -> None:
        await asyncio.to_thread(self.client.delete_object, Bucket=self.bucket, Key=key)

    async def exists(self, key: str) -> bool:
        try:
            await asyncio.to_thread(self.client.head_object, Bucket=self.bucket, Key=key)
            return True
        except ClientError:
            return False

La factory che sceglie il backend in base alla configurazione è banale:

from functools import lru_cache

from app.config import get_settings
from app.storage.local import LocalStorage
from app.storage.s3 import S3Storage


@lru_cache
def get_storage():
    settings = get_settings()
    if settings.storage_backend == "s3":
        return S3Storage(
            bucket=settings.s3_bucket,
            endpoint_url=settings.s3_endpoint_url,
            access_key=settings.s3_access_key,
            secret_key=settings.s3_secret_key,
            region=settings.s3_region,
        )
    return LocalStorage(settings.storage_path, settings.chunk_size)

Slug, token e password

Lo slug pubblico è l'unica cosa che separa un estraneo dal contenuto del transfer: deve essere generato con un generatore crittograficamente sicuro e avere entropia sufficiente. Con 16 caratteri su un alfabeto di 33 simboli otteniamo circa 80 bit, più che sufficienti per rendere impraticabile qualsiasi enumerazione.

import hashlib
import hmac
import re
import secrets
import unicodedata

ALPHABET = "abcdefghijkmnopqrstuvwxyz23456789"  # senza caratteri ambigui: l, 1, 0, o


def generate_slug(length: int = 16) -> str:
    return "".join(secrets.choice(ALPHABET) for _ in range(length))


def generate_token() -> str:
    return secrets.token_urlsafe(32)


def hash_password(password: str, salt: bytes | None = None) -> str:
    # PBKDF2 è sufficiente per una password monouso a vita breve;
    # per credenziali persistenti sarebbe preferibile Argon2
    salt = salt or secrets.token_bytes(16)
    derived = hashlib.pbkdf2_hmac("sha256", password.encode(), salt, 240_000)
    return f"pbkdf2_sha256$240000${salt.hex()}${derived.hex()}"


def verify_password(password: str, encoded: str) -> bool:
    try:
        _, iterations, salt_hex, hash_hex = encoded.split("$")
    except ValueError:
        return False
    derived = hashlib.pbkdf2_hmac(
        "sha256", password.encode(), bytes.fromhex(salt_hex), int(iterations)
    )
    # Il confronto a tempo costante evita attacchi di timing
    return hmac.compare_digest(derived.hex(), hash_hex)


def sanitize_filename(name: str) -> str:
    # Normalizziamo l'unicode e rimuoviamo separatori di percorso e caratteri di controllo
    name = unicodedata.normalize("NFKC", name)
    name = name.replace("\\", "/").split("/")[-1]
    name = re.sub(r"[\x00-\x1f\x7f]", "", name).strip()
    name = re.sub(r'[<>:"|?*]', "_", name)
    if name in {"", ".", ".."}:
        name = "file"
    return name[:255]

Gli schemi Pydantic

Gli schemi definiscono il contratto dell'API e applicano i limiti configurati prima ancora che la richiesta raggiunga il database.

import uuid
from datetime import datetime

from pydantic import BaseModel, EmailStr, Field, field_validator, model_validator

from app.config import get_settings
from app.security import sanitize_filename

settings = get_settings()


class FileDeclaration(BaseModel):
    filename: str = Field(min_length=1, max_length=255)
    size: int = Field(gt=0, le=settings.max_file_size)
    content_type: str = Field(default="application/octet-stream", max_length=127)

    @field_validator("filename")
    @classmethod
    def clean_filename(cls, value: str) -> str:
        return sanitize_filename(value)


class TransferCreate(BaseModel):
    files: list[FileDeclaration] = Field(
        min_length=1, max_length=settings.max_files_per_transfer
    )
    sender_email: EmailStr | None = None
    recipients: list[EmailStr] = Field(default_factory=list, max_length=20)
    message: str | None = Field(default=None, max_length=2000)
    password: str | None = Field(default=None, min_length=6, max_length=128)
    expiry_days: int = Field(
        default=settings.default_expiry_days, ge=1, le=settings.max_expiry_days
    )
    max_downloads: int | None = Field(default=None, ge=1, le=10_000)

    @model_validator(mode="after")
    def check_total_size(self) -> "TransferCreate":
        total = sum(item.size for item in self.files)
        if total > settings.max_transfer_size:
            raise ValueError("La dimensione complessiva del transfer supera il limite")
        if self.recipients and not self.sender_email:
            raise ValueError("Per inviare notifiche serve l'indirizzo del mittente")
        return self


class FileSlot(BaseModel):
    id: uuid.UUID
    filename: str
    upload_url: str


class TransferCreated(BaseModel):
    id: uuid.UUID
    owner_token: str
    files: list[FileSlot]


class FileInfo(BaseModel):
    id: uuid.UUID
    filename: str
    size: int
    content_type: str


class TransferPublic(BaseModel):
    slug: str
    message: str | None
    total_size: int
    expires_at: datetime
    download_count: int
    requires_password: bool
    files: list[FileInfo]

Il servizio dei transfer

Concentriamo la logica di dominio in un modulo dedicato, mantenendo i router sottili. Il servizio conosce il database e lo storage, ma non conosce HTTP.

import uuid
from datetime import datetime, timedelta, timezone
from typing import AsyncIterator

from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession

from app.config import get_settings
from app.models import Transfer, TransferFile, TransferStatus
from app.schemas import TransferCreate
from app.security import generate_slug, generate_token, hash_password
from app.storage import get_storage

settings = get_settings()


class TransferError(Exception):
    pass


class QuotaExceeded(TransferError):
    pass


class UploadIncomplete(TransferError):
    pass


async def create_transfer(session: AsyncSession, payload: TransferCreate) -> Transfer:
    transfer = Transfer(
        id=uuid.uuid4(),
        slug=await _unique_slug(session),
        owner_token=generate_token(),
        status=TransferStatus.PENDING,
        sender_email=payload.sender_email,
        recipients=",".join(payload.recipients) or None,
        message=payload.message,
        password_hash=hash_password(payload.password) if payload.password else None,
        total_size=sum(item.size for item in payload.files),
        max_downloads=payload.max_downloads,
        expires_at=datetime.now(timezone.utc) + timedelta(days=payload.expiry_days),
    )

    for declaration in payload.files:
        file_id = uuid.uuid4()
        transfer.files.append(
            TransferFile(
                id=file_id,
                filename=declaration.filename,
                content_type=declaration.content_type,
                declared_size=declaration.size,
                # La chiave è partizionata per evitare directory con milioni di file
                object_key=f"{transfer.id.hex[:2]}/{transfer.id.hex}/{file_id.hex}",
            )
        )

    session.add(transfer)
    await session.commit()
    await session.refresh(transfer)
    return transfer


async def _unique_slug(session: AsyncSession, attempts: int = 5) -> str:
    for _ in range(attempts):
        candidate = generate_slug()
        existing = await session.scalar(
            select(Transfer.id).where(Transfer.slug == candidate)
        )
        if existing is None:
            return candidate
    raise TransferError("Impossibile generare uno slug univoco")


async def store_file_content(
    session: AsyncSession,
    transfer: Transfer,
    file_record: TransferFile,
    stream: AsyncIterator[bytes],
) -> TransferFile:
    if file_record.is_complete:
        raise TransferError("Il file è già stato caricato")

    storage = get_storage()
    # Il controllo sulla dimensione avviene durante lo streaming: non possiamo
    # fidarci dell'header Content-Length inviato dal client
    guarded = _size_guard(stream, file_record.declared_size)
    written = await storage.save(file_record.object_key, guarded)

    file_record.stored_size = written
    file_record.is_complete = written == file_record.declared_size
    await session.commit()
    return file_record


async def _size_guard(stream: AsyncIterator[bytes], limit: int) -> AsyncIterator[bytes]:
    total = 0
    async for chunk in stream:
        total += len(chunk)
        if total > limit:
            raise QuotaExceeded("Il file supera la dimensione dichiarata")
        yield chunk


async def finalize_transfer(session: AsyncSession, transfer: Transfer) -> Transfer:
    missing = [item.filename for item in transfer.files if not item.is_complete]
    if missing:
        raise UploadIncomplete(f"File non completi: {', '.join(missing)}")

    transfer.status = TransferStatus.READY
    transfer.total_size = sum(item.stored_size for item in transfer.files)
    await session.commit()
    return transfer


async def get_downloadable(session: AsyncSession, slug: str) -> Transfer | None:
    transfer = await session.scalar(select(Transfer).where(Transfer.slug == slug))
    if transfer is None or transfer.status != TransferStatus.READY:
        return None
    if transfer.is_expired or transfer.is_exhausted:
        return None
    return transfer


async def register_download(session: AsyncSession, transfer_id: uuid.UUID) -> None:
    # L'incremento atomico evita race condition tra download concorrenti
    await session.execute(
        update(Transfer)
        .where(Transfer.id == transfer_id)
        .values(download_count=Transfer.download_count + 1)
    )
    await session.commit()


async def purge_transfer(session: AsyncSession, transfer: Transfer) -> None:
    storage = get_storage()
    for item in transfer.files:
        await storage.delete(item.object_key)
    transfer.status = TransferStatus.DELETED
    await session.commit()

La funzione _size_guard merita attenzione: è un generatore asincrono che avvolge lo stream in ingresso e conta i byte man mano che passano. È l'unico modo affidabile per far rispettare un limite di dimensione, perché Content-Length è un valore fornito dal client e un attaccante può semplicemente mentire.

L'endpoint di upload in streaming

Questo è il punto più delicato dell'applicazione. La tentazione è usare UploadFile di FastAPI e chiamare await file.read(), ma così facendo l'intero file finirebbe in memoria. FastAPI usa in realtà uno SpooledTemporaryFile che oltre una certa soglia rovescia su disco, il che è già meglio, ma comporta comunque una scrittura temporanea inutile. Leggiamo invece direttamente dallo stream della richiesta.

import secrets
import uuid
from typing import AsyncIterator

from fastapi import APIRouter, Depends, HTTPException, Path, Request, status
from sqlalchemy.ext.asyncio import AsyncSession

from app.config import get_settings
from app.database import get_session
from app.models import Transfer, TransferFile
from app.schemas import FileSlot, TransferCreate, TransferCreated
from app.services import transfers as service
from app.services.transfers import QuotaExceeded, UploadIncomplete

router = APIRouter(prefix="/api/transfers", tags=["transfers"])
settings = get_settings()


@router.post("", response_model=TransferCreated, status_code=status.HTTP_201_CREATED)
async def create(payload: TransferCreate, session: AsyncSession = Depends(get_session)):
    transfer = await service.create_transfer(session, payload)
    return TransferCreated(
        id=transfer.id,
        owner_token=transfer.owner_token,
        files=[
            FileSlot(
                id=item.id,
                filename=item.filename,
                upload_url=f"{settings.base_url}/api/transfers/{transfer.id}/files/{item.id}",
            )
            for item in transfer.files
        ],
    )


@router.put("/{transfer_id}/files/{file_id}", status_code=status.HTTP_204_NO_CONTENT)
async def upload(
    request: Request,
    transfer_id: uuid.UUID = Path(...),
    file_id: uuid.UUID = Path(...),
    session: AsyncSession = Depends(get_session),
):
    transfer, file_record = await _resolve(session, transfer_id, file_id, request)

    async def body_stream() -> AsyncIterator[bytes]:
        # request.stream() consegna i chunk così come arrivano dal server ASGI,
        # senza mai materializzare l'intero corpo in memoria
        async for chunk in request.stream():
            yield chunk

    try:
        await service.store_file_content(session, transfer, file_record, body_stream())
    except QuotaExceeded as error:
        raise HTTPException(status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, str(error))


@router.post("/{transfer_id}/finalize")
async def finalize(
    request: Request,
    transfer_id: uuid.UUID,
    session: AsyncSession = Depends(get_session),
):
    transfer = await _authorize(session, transfer_id, request)
    try:
        transfer = await service.finalize_transfer(session, transfer)
    except UploadIncomplete as error:
        raise HTTPException(status.HTTP_409_CONFLICT, str(error))

    # L'invio delle email è delegato al worker: non deve rallentare la risposta
    if transfer.recipients:
        await request.app.state.queue.enqueue_job(
            "send_transfer_emails", str(transfer.id)
        )

    return {
        "slug": transfer.slug,
        "url": f"{settings.base_url}/d/{transfer.slug}",
        "expires_at": transfer.expires_at,
    }


async def _authorize(
    session: AsyncSession, transfer_id: uuid.UUID, request: Request
) -> Transfer:
    token = request.headers.get("X-Owner-Token", "")
    transfer = await session.get(Transfer, transfer_id)
    if transfer is None:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "Transfer inesistente")
    if not secrets.compare_digest(token, transfer.owner_token):
        raise HTTPException(status.HTTP_403_FORBIDDEN, "Token non valido")
    return transfer


async def _resolve(
    session: AsyncSession, transfer_id: uuid.UUID, file_id: uuid.UUID, request: Request
) -> tuple[Transfer, TransferFile]:
    transfer = await _authorize(session, transfer_id, request)
    file_record = next((item for item in transfer.files if item.id == file_id), None)
    if file_record is None:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "File inesistente")
    return transfer, file_record

Notare la scelta di PUT anziché POST per l'upload: il corpo della richiesta è il contenuto grezzo del file, senza incapsulamento multipart. Questo semplifica enormemente lo streaming ed è idempotente per definizione, il che aiuta in caso di retry automatici.

Download con supporto ai range

Un servizio di file transfer serio deve supportare le richieste parziali (Range), altrimenti i download interrotti non possono essere ripresi e i player multimediali non riescono a fare seeking.

import re
import uuid
from urllib.parse import quote

from fastapi import APIRouter, Depends, Header, HTTPException, Request, status
from fastapi.responses import StreamingResponse
from sqlalchemy.ext.asyncio import AsyncSession

from app.database import get_session
from app.models import Transfer
from app.security import verify_password
from app.services import transfers as service
from app.storage import get_storage

router = APIRouter(prefix="/d", tags=["downloads"])

RANGE_PATTERN = re.compile(r"bytes=(\d*)-(\d*)")


def parse_range(header: str | None, size: int) -> tuple[int, int] | None:
    if not header:
        return None
    match = RANGE_PATTERN.fullmatch(header.strip())
    if not match:
        return None
    start_raw, end_raw = match.groups()
    if start_raw == "" and end_raw == "":
        return None
    if start_raw == "":
        # Forma "bytes=-500": gli ultimi 500 byte
        length = int(end_raw)
        start = max(size - length, 0)
        end = size - 1
    else:
        start = int(start_raw)
        end = int(end_raw) if end_raw else size - 1
    if start > end or start >= size:
        raise HTTPException(status.HTTP_416_REQUESTED_RANGE_NOT_SATISFIABLE)
    return start, min(end, size - 1)


def content_disposition(filename: str) -> str:
    # RFC 6266: forniamo sia il fallback ASCII sia la versione UTF-8 percent-encoded
    ascii_name = filename.encode("ascii", "ignore").decode() or "download"
    return f'attachment; filename="{ascii_name}"; filename*=UTF-8\'\'{quote(filename)}'


async def _authorized_transfer(
    slug: str, session: AsyncSession, password: str | None
) -> Transfer:
    transfer = await service.get_downloadable(session, slug)
    if transfer is None:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "Link non valido o scaduto")
    if transfer.password_hash:
        if not password or not verify_password(password, transfer.password_hash):
            raise HTTPException(status.HTTP_401_UNAUTHORIZED, "Password errata")
    return transfer


@router.get("/{slug}/files/{file_id}")
async def download_file(
    slug: str,
    file_id: uuid.UUID,
    range_header: str | None = Header(default=None, alias="Range"),
    x_password: str | None = Header(default=None, alias="X-Transfer-Password"),
    session: AsyncSession = Depends(get_session),
):
    transfer = await _authorized_transfer(slug, session, x_password)
    record = next((item for item in transfer.files if item.id == file_id), None)
    if record is None:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "File inesistente")

    storage = get_storage()
    size = record.stored_size
    window = parse_range(range_header, size)

    headers = {
        "Content-Disposition": content_disposition(record.filename),
        "Accept-Ranges": "bytes",
        "X-Content-Type-Options": "nosniff",
    }

    if window is None:
        headers["Content-Length"] = str(size)
        body = storage.open(record.object_key)
        code = status.HTTP_200_OK
    else:
        start, end = window
        headers["Content-Range"] = f"bytes {start}-{end}/{size}"
        headers["Content-Length"] = str(end - start + 1)
        body = storage.open(record.object_key, start, end)
        code = status.HTTP_206_PARTIAL_CONTENT

    await service.register_download(session, transfer.id)
    return StreamingResponse(
        body,
        status_code=code,
        headers=headers,
        # Forziamo un content type generico: servire il MIME dichiarato dall'utente
        # aprirebbe la porta ad attacchi XSS tramite file HTML caricati
        media_type="application/octet-stream",
    )

Il commento sull'ultimo parametro non è pedanteria. Se un servizio serve un file HTML caricato da un utente con Content-Type: text/html sullo stesso dominio dell'applicazione, quel file può eseguire JavaScript nel contesto di sicurezza del sito. La difesa canonica è duplice: forzare sempre application/octet-stream con Content-Disposition: attachment e X-Content-Type-Options: nosniff, e possibilmente servire i download da un dominio separato senza cookie.

Creazione dello ZIP al volo

Comprimere in anticipo l'intero pacchetto sarebbe uno spreco: la maggior parte dei transfer non viene mai scaricata come archivio, e i file più pesanti (video, immagini, PDF) sono già compressi. Generiamo quindi lo ZIP in streaming, senza mai scriverlo su disco, usando zipstream-ng.

from typing import AsyncIterator

from zipstream import ZIP_STORED, AsyncZipStream

from app.models import Transfer
from app.storage import get_storage


async def stream_zip(transfer: Transfer) -> AsyncIterator[bytes]:
    storage = get_storage()
    # ZIP_STORED evita di ricomprimere dati già compressi: risparmia CPU
    # e permette di conoscere in anticipo la dimensione totale dell'archivio
    archive = AsyncZipStream(sized=True, compress_type=ZIP_STORED)

    used_names: set[str] = set()
    for record in transfer.files:
        name = _deduplicate(record.filename, used_names)
        archive.add(
            storage.open(record.object_key),
            arcname=name,
            size=record.stored_size,
        )

    async for chunk in archive:
        yield chunk


def _deduplicate(name: str, used: set[str]) -> str:
    # Due file possono avere lo stesso nome: lo ZIP li accetterebbe, ma
    # l'estrazione ne sovrascriverebbe uno
    if name not in used:
        used.add(name)
        return name
    stem, dot, extension = name.rpartition(".")
    counter = 1
    while True:
        candidate = f"{stem} ({counter}){dot}{extension}" if dot else f"{name} ({counter})"
        if candidate not in used:
            used.add(candidate)
            return candidate
        counter += 1

Il flag sized=True combinato con ZIP_STORED è ciò che permette di calcolare in anticipo la lunghezza esatta dell'archivio e quindi di inviare un header Content-Length corretto. Senza di esso il browser mostrerebbe un download di dimensione ignota, con barra di progresso inutilizzabile.

from app.services.packaging import stream_zip


@router.get("/{slug}/archive")
async def download_archive(
    slug: str,
    x_password: str | None = Header(default=None, alias="X-Transfer-Password"),
    session: AsyncSession = Depends(get_session),
):
    transfer = await _authorized_transfer(slug, session, x_password)
    if not transfer.files:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "Nessun file nel transfer")

    filename = f"{transfer.slug}.zip"
    await service.register_download(session, transfer.id)
    return StreamingResponse(
        stream_zip(transfer),
        media_type="application/zip",
        headers={
            "Content-Disposition": content_disposition(filename),
            "X-Content-Type-Options": "nosniff",
        },
    )

Notifiche via email

Le email vengono inviate dal worker, mai dal ciclo richiesta-risposta: un server SMTP lento trasformerebbe una finalizzazione istantanea in un'attesa di dieci secondi.

from email.message import EmailMessage

import aiosmtplib
from jinja2 import Environment, FileSystemLoader, select_autoescape

from app.config import get_settings
from app.models import Transfer

settings = get_settings()

# L'autoescape è indispensabile: il messaggio è testo fornito dall'utente
environment = Environment(
    loader=FileSystemLoader("templates"),
    autoescape=select_autoescape(["html"]),
)


def render_email(transfer: Transfer, recipient: str) -> EmailMessage:
    template = environment.get_template("email_transfer.html")
    html = template.render(
        transfer=transfer,
        download_url=f"{settings.base_url}/d/{transfer.slug}",
        app_name=settings.app_name,
    )

    message = EmailMessage()
    # Il From resta il nostro dominio: cambiarlo farebbe fallire SPF e DMARC
    message["From"] = settings.smtp_from
    message["To"] = recipient
    message["Reply-To"] = transfer.sender_email or settings.smtp_from
    message["Subject"] = f"{transfer.sender_email} ti ha inviato dei file"
    message.set_content(
        f"Scarica i file da {settings.base_url}/d/{transfer.slug}\n"
        f"Il link scade il {transfer.expires_at:%d/%m/%Y}."
    )
    message.add_alternative(html, subtype="html")
    return message


async def send_message(message: EmailMessage) -> None:
    await aiosmtplib.send(
        message,
        hostname=settings.smtp_host,
        port=settings.smtp_port,
        username=settings.smtp_user,
        password=settings.smtp_password,
        start_tls=settings.smtp_starttls,
    )

Un dettaglio spesso trascurato: l'header From deve restare quello del nostro dominio, altrimenti SPF e DMARC faranno finire il messaggio nello spam. L'indirizzo del mittente reale va in Reply-To.

I worker in background

ARQ fornisce sia le code sia lo scheduler cron, con un'API asincrona che si integra naturalmente con il resto dello stack.

import uuid
from datetime import datetime, timedelta, timezone

from arq import cron
from arq.connections import RedisSettings
from sqlalchemy import select

from app.config import get_settings
from app.database import async_session_factory
from app.models import Transfer, TransferStatus
from app.services import transfers as service
from app.services.mailer import render_email, send_message

settings = get_settings()


async def send_transfer_emails(ctx: dict, transfer_id: str) -> int:
    async with async_session_factory() as session:
        transfer = await session.get(Transfer, uuid.UUID(transfer_id))
        if transfer is None or not transfer.recipients:
            return 0
        sent = 0
        for recipient in transfer.recipients.split(","):
            # Un destinatario che fallisce non deve bloccare gli altri
            try:
                await send_message(render_email(transfer, recipient.strip()))
                sent += 1
            except Exception as error:
                ctx["logger"].warning("Invio fallito a %s: %s", recipient, error)
        return sent


async def expire_transfers(ctx: dict) -> int:
    now = datetime.now(timezone.utc)
    async with async_session_factory() as session:
        result = await session.scalars(
            select(Transfer).where(
                Transfer.status == TransferStatus.READY,
                Transfer.expires_at <= now,
            )
        )
        count = 0
        for transfer in result:
            transfer.status = TransferStatus.EXPIRED
            count += 1
        await session.commit()
        return count


async def purge_expired(ctx: dict) -> int:
    async with async_session_factory() as session:
        result = await session.scalars(
            select(Transfer).where(Transfer.status == TransferStatus.EXPIRED)
        )
        count = 0
        for transfer in result:
            await service.purge_transfer(session, transfer)
            count += 1
        return count


async def purge_abandoned(ctx: dict) -> int:
    # I transfer creati ma mai finalizzati occupano spazio: li rimuoviamo dopo 24 ore
    threshold = datetime.now(timezone.utc) - timedelta(hours=24)
    async with async_session_factory() as session:
        result = await session.scalars(
            select(Transfer).where(
                Transfer.status == TransferStatus.PENDING,
                Transfer.created_at <= threshold,
            )
        )
        count = 0
        for transfer in result:
            await service.purge_transfer(session, transfer)
            count += 1
        return count


class WorkerSettings:
    functions = [send_transfer_emails]
    cron_jobs = [
        cron(expire_transfers, minute={0, 15, 30, 45}),
        cron(purge_expired, hour={3}, minute={0}),
        cron(purge_abandoned, hour={4}, minute={0}),
    ]
    redis_settings = RedisSettings.from_dsn(settings.redis_url)

La separazione tra expire_transfers e purge_expired non è ridondanza. Marcare come scaduto è un'operazione istantanea e sicura che rende immediatamente inaccessibile il contenuto; la cancellazione fisica è irreversibile e può essere rimandata alle ore notturne, lasciando una finestra durante la quale un errore è ancora recuperabile.

Assemblare l'applicazione

from contextlib import asynccontextmanager

from arq import create_pool
from arq.connections import RedisSettings
from fastapi import FastAPI, Request
from fastapi.responses import JSONResponse
from slowapi import Limiter
from slowapi.errors import RateLimitExceeded
from slowapi.util import get_remote_address

from app.config import get_settings
from app.routers import downloads, transfers

settings = get_settings()
limiter = Limiter(key_func=get_remote_address, storage_uri=settings.redis_url)


@asynccontextmanager
async def lifespan(app: FastAPI):
    # La connessione a Redis viene aperta una volta sola e riusata
    app.state.queue = await create_pool(RedisSettings.from_dsn(settings.redis_url))
    yield
    await app.state.queue.close()


app = FastAPI(title=settings.app_name, lifespan=lifespan)
app.state.limiter = limiter
app.include_router(transfers.router)
app.include_router(downloads.router)


@app.exception_handler(RateLimitExceeded)
async def rate_limit_handler(request: Request, exc: RateLimitExceeded):
    return JSONResponse(
        status_code=429,
        content={"detail": "Troppe richieste, riprova più tardi"},
    )


@app.middleware("http")
async def security_headers(request: Request, call_next):
    response = await call_next(request)
    response.headers["X-Content-Type-Options"] = "nosniff"
    response.headers["X-Frame-Options"] = "DENY"
    response.headers["Referrer-Policy"] = "no-referrer"
    return response


@app.get("/health")
async def health():
    return {"status": "ok"}

Il rate limiting va applicato in modo selettivo: la creazione di transfer è l'operazione più costosa e va limitata con severità, mentre il download di un file già pubblicato può essere più permissivo.

@router.post("", response_model=TransferCreated, status_code=status.HTTP_201_CREATED)
@limiter.limit("10/hour")
async def create(
    request: Request,
    payload: TransferCreate,
    session: AsyncSession = Depends(get_session),
):
    ...

Il client di upload

Sul frontend, il caricamento sfrutta la Fetch API. Per avere una barra di progresso affidabile con file di grandi dimensioni conviene però ricorrere a XMLHttpRequest, che espone l'evento progress sull'upload, cosa che fetch non fa in modo portabile.

async function createTransfer(files, options = {}) {
  const response = await fetch("/api/transfers", {
    method: "POST",
    headers: { "Content-Type": "application/json" },
    body: JSON.stringify({
      files: files.map((file) => ({
        filename: file.name,
        size: file.size,
        content_type: file.type || "application/octet-stream",
      })),
      ...options,
    }),
  });

  if (!response.ok) {
    const error = await response.json();
    throw new Error(error.detail ?? "Creazione del transfer fallita");
  }
  return response.json();
}

function uploadFile(file, url, ownerToken, onProgress) {
  return new Promise((resolve, reject) => {
    const request = new XMLHttpRequest();
    request.open("PUT", url);
    request.setRequestHeader("X-Owner-Token", ownerToken);
    request.setRequestHeader("Content-Type", "application/octet-stream");

    // L'evento upload.progress è l'unico modo affidabile per una barra di avanzamento
    request.upload.addEventListener("progress", (event) => {
      if (event.lengthComputable) {
        onProgress(event.loaded / event.total);
      }
    });

    request.addEventListener("load", () => {
      if (request.status >= 200 && request.status < 300) resolve();
      else reject(new Error(`Upload fallito con stato ${request.status}`));
    });
    request.addEventListener("error", () => reject(new Error("Errore di rete")));
    request.addEventListener("abort", () => reject(new Error("Upload annullato")));

    request.send(file);
  });
}

async function sendFiles(files, options, onProgress) {
  const transfer = await createTransfer(files, options);
  const totalBytes = files.reduce((sum, file) => sum + file.size, 0);
  const uploaded = new Array(files.length).fill(0);

  // Gli upload avvengono in sequenza: il parallelismo su file grandi
  // satura la banda in upstream senza alcun guadagno reale
  for (let index = 0; index < files.length; index += 1) {
    const slot = transfer.files[index];
    await uploadFile(files[index], slot.upload_url, transfer.owner_token, (ratio) => {
      uploaded[index] = ratio * files[index].size;
      const done = uploaded.reduce((sum, value) => sum + value, 0);
      onProgress(done / totalBytes);
    });
    uploaded[index] = files[index].size;
  }

  const finalize = await fetch(`/api/transfers/${transfer.id}/finalize`, {
    method: "POST",
    headers: { "X-Owner-Token": transfer.owner_token },
  });

  if (!finalize.ok) throw new Error("Finalizzazione fallita");
  return finalize.json();
}

Configurazione di Nginx

Il reverse proxy va configurato con attenzione, altrimenti vanificherà tutto il lavoro fatto sullo streaming. In particolare il buffering delle richieste è attivo di default e farebbe scrivere ogni upload su un file temporaneo di Nginx prima di inoltrarlo all'applicazione.

server {
    listen 443 ssl;
    http2 on;
    server_name filedrop.example.com;

    ssl_certificate     /etc/letsencrypt/live/filedrop.example.com/fullchain.pem;
    ssl_certificate_key /etc/letsencrypt/live/filedrop.example.com/privkey.pem;

    # 0 disabilita del tutto il limite: il controllo è applicativo
    client_max_body_size 0;

    location / {
        proxy_pass http://127.0.0.1:8000;
        proxy_http_version 1.1;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
        proxy_set_header X-Forwarded-Proto $scheme;

        # Senza questi due parametri Nginx bufferizza l'intero corpo su disco
        # prima di inoltrarlo, annullando i vantaggi dello streaming
        proxy_request_buffering off;
        proxy_buffering off;

        # Gli upload lunghi non devono cadere per timeout
        proxy_read_timeout 3600s;
        proxy_send_timeout 3600s;
    }
}

Con lo storage locale è possibile fare un ulteriore passo avanti usando X-Accel-Redirect: l'applicazione autorizza la richiesta e restituisce un header che dice a Nginx quale file servire, delegandogli l'intero trasferimento. Il processo Python si libera immediatamente invece di rimanere occupato per l'intera durata del download.

from urllib.parse import quote

from fastapi import Response


@router.get("/{slug}/files/{file_id}/accel")
async def download_accelerated(
    slug: str,
    file_id: uuid.UUID,
    x_password: str | None = Header(default=None, alias="X-Transfer-Password"),
    session: AsyncSession = Depends(get_session),
):
    transfer = await _authorized_transfer(slug, session, x_password)
    record = next((item for item in transfer.files if item.id == file_id), None)
    if record is None:
        raise HTTPException(status.HTTP_404_NOT_FOUND, "File inesistente")

    await service.register_download(session, transfer.id)
    # Il percorso punta a una location interna, non raggiungibile dall'esterno.
    # Nginx si occupa di range, sendfile e throttling al posto nostro
    return Response(
        status_code=200,
        headers={
            "X-Accel-Redirect": f"/protected/{quote(record.object_key)}",
            "Content-Disposition": content_disposition(record.filename),
            "Content-Type": "application/octet-stream",
        },
    )
location /protected/ {
    internal;
    alias /var/lib/filedrop/uploads/;
}

Docker Compose

services:
  api:
    build: .
    command: uvicorn app.main:app --host 0.0.0.0 --port 8000 --workers 4
    environment:
      FILEDROP_DATABASE_URL: postgresql+asyncpg://filedrop:secret@db/filedrop
      FILEDROP_REDIS_URL: redis://cache:6379/0
      FILEDROP_STORAGE_BACKEND: s3
      FILEDROP_S3_ENDPOINT_URL: http://minio:9000
      FILEDROP_S3_ACCESS_KEY: minioadmin
      FILEDROP_S3_SECRET_KEY: minioadmin
      FILEDROP_SECRET_KEY: ${SECRET_KEY}
    depends_on: [db, cache, minio]
    ports: ["8000:8000"]

  worker:
    build: .
    command: arq app.workers.WorkerSettings
    environment:
      FILEDROP_DATABASE_URL: postgresql+asyncpg://filedrop:secret@db/filedrop
      FILEDROP_REDIS_URL: redis://cache:6379/0
      FILEDROP_SECRET_KEY: ${SECRET_KEY}
    depends_on: [db, cache]

  db:
    image: postgres:17-alpine
    environment:
      POSTGRES_USER: filedrop
      POSTGRES_PASSWORD: secret
      POSTGRES_DB: filedrop
    volumes: ["pgdata:/var/lib/postgresql/data"]
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U filedrop"]
      interval: 10s

  cache:
    image: redis:7-alpine
    command: redis-server --appendonly yes
    volumes: ["redisdata:/data"]

  minio:
    image: minio/minio
    command: server /data --console-address ":9001"
    environment:
      MINIO_ROOT_USER: minioadmin
      MINIO_ROOT_PASSWORD: minioadmin
    volumes: ["miniodata:/data"]
    ports: ["9001:9001"]

volumes:
  pgdata:
  redisdata:
  miniodata:

Test

I test più preziosi sono quelli che verificano i confini: i limiti di dimensione, la scadenza, la protezione con password.

import pytest
from httpx import ASGITransport, AsyncClient

from app.main import app


@pytest.fixture
async def client():
    transport = ASGITransport(app=app)
    async with AsyncClient(transport=transport, base_url="http://test") as instance:
        yield instance


@pytest.mark.asyncio
async def test_upload_and_download_roundtrip(client):
    payload = {"files": [{"filename": "note.txt", "size": 10}]}
    created = (await client.post("/api/transfers", json=payload)).json()
    token = created["owner_token"]
    slot = created["files"][0]

    upload = await client.put(
        f"/api/transfers/{created['id']}/files/{slot['id']}",
        content=b"ciao mondo",
        headers={"X-Owner-Token": token},
    )
    assert upload.status_code == 204

    finalized = await client.post(
        f"/api/transfers/{created['id']}/finalize",
        headers={"X-Owner-Token": token},
    )
    slug = finalized.json()["slug"]

    download = await client.get(f"/d/{slug}/files/{slot['id']}")
    assert download.content == b"ciao mondo"


@pytest.mark.asyncio
async def test_upload_bigger_than_declared_is_rejected(client):
    payload = {"files": [{"filename": "note.txt", "size": 5}]}
    created = (await client.post("/api/transfers", json=payload)).json()
    slot = created["files"][0]

    # Il client dichiara 5 byte ma ne invia molti di più
    response = await client.put(
        f"/api/transfers/{created['id']}/files/{slot['id']}",
        content=b"x" * 1000,
        headers={"X-Owner-Token": created["owner_token"]},
    )
    assert response.status_code == 413


@pytest.mark.asyncio
async def test_finalize_without_upload_fails(client):
    payload = {"files": [{"filename": "note.txt", "size": 10}]}
    created = (await client.post("/api/transfers", json=payload)).json()

    response = await client.post(
        f"/api/transfers/{created['id']}/finalize",
        headers={"X-Owner-Token": created["owner_token"]},
    )
    assert response.status_code == 409


@pytest.mark.asyncio
async def test_unknown_slug_returns_404(client):
    response = await client.get("/d/inesistente/archive")
    assert response.status_code == 404

Considerazioni di sicurezza

Riassumiamo i punti su cui non è consigliabile fare compromessi.

  • Nomi dei file mai usati come percorsi. La chiave nello storage è un UUID generato dal server. Il nome originale vive solo nel database e nell'header Content-Disposition.
  • Content type forzato. Servire il MIME dichiarato dall'utente sullo stesso dominio dell'applicazione consente XSS. Forzare application/octet-stream, nosniff e attachment, oppure usare un dominio dedicato ai download.
  • Slug con entropia sufficiente. Ottanta bit rendono l'enumerazione impraticabile. Uno slug corto o derivato da un contatore è un invito a scoprire i file altrui.
  • Limiti applicati durante lo streaming. Content-Length è un'affermazione del client, non un fatto.
  • Confronti a tempo costante per token e password, tramite hmac.compare_digest.
  • Rate limiting sulla creazione, per impedire che il servizio venga usato come storage gratuito illimitato o come veicolo di distribuzione di malware.
  • Pulizia dei transfer abbandonati. Senza di essa, chi crea transfer e non li finalizza riempie il disco silenziosamente.
  • Nessun indice pubblico. L'endpoint che elenca i transfer non deve esistere; l'unico accesso è tramite slug.

Ottimizzazioni possibili

Il sistema descritto regge senza problemi carichi di piccola e media entità. Per scalare oltre, tre direzioni sono particolarmente efficaci.

La prima è l'upload diretto verso S3 tramite URL presigned: il server genera un URL firmato a breve scadenza e il browser carica direttamente sull'object storage, senza che un solo byte attraversi l'applicazione Python. È il modello che usano tutti i servizi di grandi dimensioni.

import asyncio

from app.config import get_settings
from app.storage import get_storage

settings = get_settings()


async def presign_upload(key: str, content_type: str, expires: int = 3600) -> dict:
    client = get_storage().client
    # Le condizioni limitano ciò che il client può caricare con l'URL firmato
    return await asyncio.to_thread(
        client.generate_presigned_post,
        Bucket=settings.s3_bucket,
        Key=key,
        Fields={"Content-Type": content_type},
        Conditions=[
            {"Content-Type": content_type},
            ["content-length-range", 1, settings.max_file_size],
        ],
        ExpiresIn=expires,
    )

La seconda è la deduplicazione per checksum: calcolando lo SHA-256 durante l'upload si può riconoscere che un contenuto identico esiste già e limitarsi a creare un riferimento. Su servizi dove gli stessi file circolano ripetutamente il risparmio di spazio è sostanziale, ma introduce un vincolo importante: la cancellazione deve diventare un conteggio di riferimenti, non una delete immediata.

La terza è la CDN davanti allo storage: se un transfer viene scaricato da cinquanta destinatari, servirlo cinquanta volte dall'origine è uno spreco di banda evitabile con URL firmati e TTL allineati alla scadenza del transfer.

Conclusioni

Costruire un clone di WeTransfer è un esercizio che mette alla prova quasi tutti gli aspetti dello sviluppo backend: I/O asincrono, gestione della memoria sotto carico, sicurezza applicativa, task in background, protocolli HTTP oltre il livello superficiale. Il codice presentato è deliberatamente completo nei punti che di solito vengono liquidati con un commento del tipo "in produzione andrebbe gestito meglio", perché è esattamente lì che si concentrano i problemi reali.

Il nucleo dell'architettura è la separazione tra metadati e contenuto: il database sa tutto e non contiene nulla, lo storage contiene tutto e non sa nulla. Tenere questi due mondi distinti è ciò che rende possibile sostituire il filesystem locale con S3 cambiando una variabile d'ambiente, mettere una CDN davanti ai download senza toccare la logica, e cancellare i contenuti scaduti senza corrompere lo stato applicativo.