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:
- Il client crea un transfer inviando i metadati (nomi dei file, dimensioni, destinatari, messaggio).
- Il server risponde con un identificatore del transfer e, per ogni file, un identificatore di upload.
- Il client carica il contenuto binario di ciascun file.
- 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.
- Chiunque possieda il link può scaricare i singoli file o l'intero pacchetto come ZIP, finché il transfer non scade.
- 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,nosniffeattachment, 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.