use std::time::Duration; use axum::Json; use axum::extract::{Multipart, Path, State}; use axum::http::StatusCode; use chrono::{DateTime, Utc}; use serde::Deserialize; use uuid::Uuid; use crate::auth::middleware::AuthUser; use crate::error::AppError; use crate::models::hashtag::{self, Hashtag}; use crate::models::upload::{Upload, UploadDto}; use crate::models::user::User; use crate::services::config; use crate::state::AppState; const MAX_CAPTION_LENGTH: usize = 2000; /// Allowlist of accepted media types, keyed by the MIME that `infer` derives from /// the file's magic bytes. The detected MIME (not the client-declared one) is what /// we trust, store, and hand to the compression pipeline — so a text-based payload /// (SVG/HTML/JS) can never be stored or served on-origin. Each entry maps to the /// server-controlled file extension we persist the original under. /// /// HEIC/HEIF are deliberately excluded: the preview pipeline (`image` crate, and /// the bundled ffmpeg 6.1) cannot decode them, so accepting them would store files /// that never get a thumbnail. iOS Safari already transcodes HEIC→JPEG when a photo /// is selected via a file input, so this rejects only the rare HEIC-preserving /// upload path — with a clear error rather than a silently broken post. const ALLOWED_MEDIA: &[(&str, &str)] = &[ ("image/jpeg", "jpg"), ("image/png", "png"), ("image/webp", "webp"), ("image/gif", "gif"), ("video/mp4", "mp4"), ("video/quicktime", "mov"), ("video/webm", "webm"), ]; pub async fn upload( State(state): State, auth: AuthUser, mut multipart: Multipart, ) -> Result<(StatusCode, Json), AppError> { // Rate limit: N uploads per hour per user. Gated by master + per-endpoint toggles. let rate_limits_on = config::get_bool(&state.config_cache, "rate_limits_enabled", true).await; let upload_rate_on = config::get_bool(&state.config_cache, "upload_rate_enabled", true).await; if rate_limits_on && upload_rate_on { let upload_rate = config::get_i64(&state.config_cache, "upload_rate_per_hour", 100).await as usize; if let Err(retry_after_secs) = state.rate_limiter.check_with_retry( format!("upload:{}", auth.user_id), upload_rate, Duration::from_secs(3600), ) { drain_multipart(multipart).await; return Err(AppError::TooManyRequests( "Du hast dein Upload-Limit für diese Stunde erreicht.".into(), Some(retry_after_secs), )); } } // Check if user is banned let user = User::find_by_id(&state.pool, auth.user_id) .await? .ok_or_else(|| AppError::NotFound("Benutzer nicht gefunden.".into()))?; if user.is_banned { drain_multipart(multipart).await; return Err(AppError::Forbidden("Du bist gesperrt.".into())); } // Check if uploads are locked let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug) .await? .ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?; if event.uploads_locked_at.is_some() { drain_multipart(multipart).await; // Reversible: a host can reopen the event, so the client keeps the queued blob and // retries on `event-opened` rather than purging it (UploadsLocked, not Forbidden). return Err(AppError::UploadsLocked("Uploads sind gesperrt.".into())); } // Belt-and-suspenders on top of the lock (release ⇒ lock): once the gallery is // released the export has been snapshotted, so a late upload could never make it into // the keepsake. Reject it explicitly rather than silently diverging the live feed. // Also reversible (reopen clears `export_released_at`), so likewise UploadsLocked. if event.export_released_at.is_some() { drain_multipart(multipart).await; return Err(AppError::UploadsLocked( "Galerie wurde bereits freigegeben.".into(), )); } // Read config limits from DB let max_image_mb: i64 = config::get_i64(&state.config_cache, "max_image_size_mb", 20).await; let max_video_mb: i64 = config::get_i64(&state.config_cache, "max_video_size_mb", 500).await; // The uploaded file is streamed straight to a temp file on disk (never buffered // whole in memory — a 500 MB video used to cost 500 MB of RAM per concurrent // upload). We only keep the first ≤512 bytes in memory for magic-byte sniffing. // On success the temp file is renamed into place under its detected extension. let upload_id = Uuid::new_v4(); let event_slug = &state.config.event_slug; let originals_dir = state .config .media_path .join(format!("originals/{event_slug}")); let temp_abs = originals_dir.join(format!("{upload_id}.tmp")); let mut streamed: Option<(i64, Vec)> = None; // (size, head bytes for sniffing) let mut caption: Option = None; let mut hashtags_csv: Option = None; // The client's idempotency key. Optional: an older client, or any other caller, simply // doesn't send one and gets the previous behaviour. let mut client_upload_id: Option = None; // Wrap the multipart read so any error after the temp file is created still cleans // it up (a mid-stream parse failure must not leave a stray `.tmp` on disk). let parse_result: Result<(), AppError> = async { while let Some(field) = multipart .next_field() .await .map_err(|e| AppError::BadRequest(e.to_string()))? { let name = field.name().unwrap_or_default().to_string(); match name.as_str() { "file" => { // The client-declared Content-Type does NOT determine the stored // MIME/extension — those come from the file's magic bytes below. The // declared type only picks the streaming cap so an oversized body is // aborted early; a mislabelled type only makes the cap *stricter* // (safe), and the authoritative per-class check still runs on the // detected type. let declared = field.content_type().unwrap_or("").to_string(); let cap_bytes = if declared.starts_with("video/") { (max_video_mb * 1024 * 1024) as usize } else if declared.starts_with("image/") { (max_image_mb * 1024 * 1024) as usize } else { (max_image_mb.max(max_video_mb) * 1024 * 1024) as usize }; tokio::fs::create_dir_all(&originals_dir) .await .map_err(|e| AppError::Internal(e.into()))?; streamed = Some(stream_field_to_file(field, &temp_abs, cap_bytes).await?); } "caption" => { caption = Some( field .text() .await .map_err(|e| AppError::BadRequest(e.to_string()))?, ); } "hashtags" => { hashtags_csv = Some( field .text() .await .map_err(|e| AppError::BadRequest(e.to_string()))?, ); } "client_upload_id" => { let raw = field .text() .await .map_err(|e| AppError::BadRequest(e.to_string()))?; // A malformed key is not worth rejecting an upload over — the photo is the // thing the guest cares about. Drop the key and lose only the retry // protection, which is exactly where we were before it existed. client_upload_id = Uuid::parse_str(raw.trim()).ok(); } _ => {} } } Ok(()) } .await; if let Err(e) = parse_result { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(e); } // Idempotency, fast path: this key already has a live upload, so the previous attempt DID // succeed and only its response was lost. Replay that response instead of storing the photo // a second time and charging the guest's quota twice. // // The body has necessarily already been streamed to disk — the key arrives as a multipart // field, so it cannot be known before the body is read. Re-sending the bytes is the client's // cost and it has already been paid by the time we get here; what has to be prevented is a // second ROW, and that is what this does. The concurrent case (two retries in flight at once) // is caught by the unique index inside the transaction below. if let Some(cid) = client_upload_id && let Some(existing) = Upload::find_by_client_upload_id(&state.pool, auth.user_id, cid) .await .map_err(AppError::from)? { let _ = tokio::fs::remove_file(&temp_abs).await; tracing::info!( client_upload_id = %cid, upload_id = %existing.id, "duplicate upload suppressed; replaying the original response" ); let dto = replay_upload_dto(&state, &existing, &user.display_name).await; return Ok((StatusCode::OK, Json(dto))); } // From here on the temp file may exist; every validation failure removes it before // returning so a rejected upload never leaves bytes behind. let (size, head) = match streamed { Some(s) => s, None => return Err(AppError::BadRequest("Keine Datei hochgeladen.".into())), }; // Validate caption length. Counted in chars (code points) to match the // "Zeichen" wording in the error message — `.len()` would be bytes and // reject perfectly valid German/emoji captions early. if let Some(ref cap) = caption && cap.chars().count() > MAX_CAPTION_LENGTH { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(AppError::BadRequest(format!( "Beschreibung ist zu lang. Maximum: {} Zeichen.", MAX_CAPTION_LENGTH ))); } // Determine the file type from its magic bytes and require it to be on the // allowlist. `infer` returns None for text-based payloads (SVG/HTML/JS), so // those are rejected outright — closing the stored-XSS vector. Both the MIME // we persist and the on-disk extension come from the detected type, never from // client-supplied values. let kind = match infer::get(&head) { Some(k) => k, None => { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(AppError::BadRequest( "Dateityp nicht erkannt oder nicht unterstützt.".into(), )); } }; let (mime, ext) = match ALLOWED_MEDIA .iter() .find(|(allowed, _)| *allowed == kind.mime_type()) .map(|(m, e)| ((*m).to_string(), *e)) { Some(v) => v, None => { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(AppError::BadRequest(format!( "Dateityp wird nicht unterstützt: {}.", kind.mime_type() ))); } }; // Validate file size against the authoritative per-detected-class limit. let max_bytes = if mime.starts_with("video/") { max_video_mb * 1024 * 1024 } else { max_image_mb * 1024 * 1024 }; if size > max_bytes { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(AppError::BadRequest(format!( "Datei ist zu groß. Maximum: {} MB.", max_bytes / (1024 * 1024) ))); } // Images only: refuse anything the compression worker could never decode, reading just // the header. Without this the upload is accepted with a 201 and then silently // soft-deleted minutes later when the worker gives up — the guest sees the photo // vanish with, at best, a vague "could not be processed". Rejecting here gives them a // reason at the door that they can act on, and it uses the SAME budget the worker // enforces, so admission and processing cannot disagree. // // Both probes open the file and run the codec's header parse — synchronous filesystem and // CPU work. They ran inline on the async task, which on this 2-vCPU box means tokio has // exactly two worker threads and every upload stalled half the runtime's request-serving // capacity. Everything else in the app that blocks (image encode, bcrypt) is already on the // blocking pool; this was the one that wasn't. if mime.starts_with("image/") { let probe_path = temp_abs.clone(); let probe = tokio::task::spawn_blocking(move || { let over = crate::services::imaging::exceeds_decode_budget(&probe_path); // Only pay for the second header read when it will actually be shown to the guest. let mp = over .then(|| crate::services::imaging::megapixels(&probe_path)) .flatten(); (over, mp) }) .await; // A join error is the blocking pool panicking or shutting down. That says nothing about // the image, so admit it and let the compression worker be the judge rather than // rejecting a photo for an infrastructure reason. let (over_budget, mp) = probe.unwrap_or_else(|e| { tracing::warn!(error = ?e, "decode-budget probe failed to run; admitting the upload"); (false, None) }); if over_budget { tracing::info!( %mime, megapixels = ?mp, "rejecting an image that exceeds the decode budget at admission" ); let _ = tokio::fs::remove_file(&temp_abs).await; let detail = mp.map_or(String::new(), |mp| format!(" (ca. {mp:.0} Megapixel)")); return Err(AppError::BadRequest(format!( "Bild hat zu viele Bildpunkte{detail} und kann nicht verarbeitet werden. \ Bitte verkleinere es und lade es erneut hoch." ))); } } // Per-user storage quota — dynamic formula based on available disk space and the // number of active uploaders. Gated by master + per-area toggles so the admin can // disable it on trusted instances. let quota_on = config::get_bool(&state.config_cache, "quota_enabled", true).await; let storage_quota_on = config::get_bool(&state.config_cache, "storage_quota_enabled", true).await; // When quota is enforced, this holds the byte ceiling so the increment UPDATE below can // enforce it atomically (`WHERE total + size <= limit`). Without that guard, two // concurrent uploads from the same user (e.g. phone + laptop) both pass this stale // pre-check and both increment, blowing past the quota. The pre-check stays as a // fast path that avoids the disk write when the user is already clearly over. let mut quota_limit: Option = None; if quota_on && storage_quota_on { let estimate = compute_storage_quota(&state).await; if let Some(limit) = estimate.limit_bytes { quota_limit = Some(limit); let prospective_total = user.total_upload_bytes.saturating_add(size); if prospective_total > limit { let _ = tokio::fs::remove_file(&temp_abs).await; return Err(AppError::QuotaExceeded( // Name the remedy, because the guest cannot see the number. Every quota // display is staff-gated by design, so a guest hitting this had no idea what // the limit was, how close they were, or what to do — and the one sentence // that tells them ("delete older posts") lived inside the staff-only block. "Du hast dein Upload-Limit für dieses Event erreicht. Lösche ältere eigene \ Beiträge, um wieder Platz zu schaffen." .into(), )); } } } // All checks passed — atomically move the temp file to its final, extension-correct // path (same directory, so the rename is cheap and atomic). let relative_path = format!("originals/{event_slug}/{upload_id}.{ext}"); let absolute_path = state.config.media_path.join(&relative_path); tokio::fs::rename(&temp_abs, &absolute_path) .await .map_err(|e| AppError::Internal(e.into()))?; // Process hashtags from caption and explicit CSV let mut tags: Vec = Vec::new(); if let Some(ref cap) = caption { tags.extend(hashtag::extract_hashtags(cap)); } if let Some(ref csv) = hashtags_csv { for tag in csv.split(',') { let t = tag.trim().trim_start_matches('#').to_lowercase(); if !t.is_empty() { tags.push(t); } } } tags.sort(); tags.dedup(); // Quota accounting, the upload row, and its hashtag links must be atomic: a // crash between the bytes increment and the insert would permanently charge // bytes with no row to reclaim them (silent quota erosion / spurious lockout). let tx_result: Result = async { let mut tx = state.pool.begin().await?; // RE-CHECK THE LOCK, UNDER A ROW LOCK, INSIDE THE COMMIT TX. // // The pre-flight check at the top of this handler ran BEFORE we streamed the body — which // for a 500 MB video is minutes. Trusting it here is a TOCTOU that silently loses photos // from the keepsake, and it is the real cause of the "stale keepsake" bug that survived // three rounds of fixes inside the export state machine: // // 1. guest starts a big upload; the lock check passes (event open) // 2. host releases the gallery → uploads lock, export workers snapshot the uploads table // 3. this upload commits AFTER that snapshot → it shows up in the live feed but is // MISSING from the downloaded keepsake, permanently (nothing ever regenerates it) // // `FOR SHARE` conflicts with the `UPDATE event` in `release_gallery`, which serializes us // against it. Either we take the lock first — and release (hence the export snapshot) is // strictly ordered after our commit, so the snapshot CONTAINS this upload — or release // commits first and we observe the lock here and reject. Either way the keepsake is // complete. `UploadsLocked` (not Forbidden) is reversible: the client keeps the blob and // resumes it when the host reopens. let (locked_at, released_at): (Option>, Option>) = sqlx::query_as( "SELECT uploads_locked_at, export_released_at FROM event WHERE id = $1 FOR SHARE", ) .bind(auth.event_id) .fetch_one(&mut *tx) .await?; if locked_at.is_some() || released_at.is_some() { return Err(AppError::UploadsLocked("Uploads sind gesperrt.".into())); } // Increment the user's byte total. When a quota is in force, guard it atomically // (`total + size <= limit`) so two concurrent uploads can't both slip past the // stale pre-check — the loser's UPDATE matches 0 rows and we abort with the same // terminal quota error (the tx rolls back on drop; the on-disk file is cleaned by // the error path below). let inc = if let Some(limit) = quota_limit { sqlx::query( "UPDATE \"user\" SET total_upload_bytes = total_upload_bytes + $2 WHERE id = $1 AND total_upload_bytes + $2 <= $3", ) .bind(auth.user_id) .bind(size) .bind(limit) .execute(&mut *tx) .await? } else { sqlx::query( "UPDATE \"user\" SET total_upload_bytes = total_upload_bytes + $2 WHERE id = $1", ) .bind(auth.user_id) .bind(size) .execute(&mut *tx) .await? }; if inc.rows_affected() == 0 { return Err(AppError::QuotaExceeded( // Name the remedy, because the guest cannot see the number. Every quota // display is staff-gated by design, so a guest hitting this had no idea what // the limit was, how close they were, or what to do — and the one sentence // that tells them ("delete older posts") lived inside the staff-only block. "Du hast dein Upload-Limit für dieses Event erreicht. Lösche ältere eigene \ Beiträge, um wieder Platz zu schaffen." .into(), )); } // `None` means a concurrent request already stored this key. The transaction — quota // increment included — is abandoned by returning here, and the caller replays the winning // row. This is the narrow race the fast path above cannot see: two retries of the same // photo in flight at the same moment. let Some(upload) = Upload::create( &mut *tx, auth.event_id, auth.user_id, &relative_path, &mime, size, caption.as_deref(), client_upload_id, ) .await? else { return Err(AppError::Conflict(DUPLICATE_UPLOAD_MARKER.into())); }; for tag in &tags { let h = Hashtag::upsert(&mut *tx, auth.event_id, tag).await?; Hashtag::link_to_upload(&mut *tx, upload.id, h.id).await?; } tx.commit().await?; Ok(upload) } .await; // The file is already on disk at `absolute_path`. If the transaction failed, no DB // row will ever reference it, so remove it now rather than orphan bytes on disk. let upload = match tx_result { Ok(u) => u, // The concurrent duplicate resolved inside the transaction. The winner's row is committed; // answer with it so both retries of the same photo get the same successful reply. Err(AppError::Conflict(ref marker)) if marker == DUPLICATE_UPLOAD_MARKER => { let _ = tokio::fs::remove_file(&absolute_path).await; let existing = match client_upload_id { Some(cid) => Upload::find_by_client_upload_id(&state.pool, auth.user_id, cid) .await .map_err(AppError::from)?, None => None, }; // If the winning row has vanished between the conflict and this lookup (deleted in // the intervening milliseconds), there is nothing to replay — report the conflict. let existing = existing.ok_or_else(|| { AppError::Conflict("Dieser Upload wurde bereits verarbeitet.".into()) })?; tracing::info!( upload_id = %existing.id, "concurrent duplicate upload resolved; replaying the stored row" ); let dto = replay_upload_dto(&state, &existing, &user.display_name).await; return Ok((StatusCode::OK, Json(dto))); } Err(e) => { let _ = tokio::fs::remove_file(&absolute_path).await; return Err(e); } }; // Spawn compression task state .compression .process(upload.id, relative_path, mime.clone()); // Broadcast SSE event let dto = UploadDto { id: upload.id, user_id: auth.user_id, uploader_name: user.display_name, preview_url: None, thumbnail_url: None, mime_type: mime, caption, hashtags: tags, like_count: 0, comment_count: 0, liked_by_me: false, created_at: upload.created_at, }; let _ = state.sse_tx.send(crate::state::SseEvent::new( "new-upload", serde_json::to_string(&dto).unwrap_or_default(), )); Ok((StatusCode::CREATED, Json(dto))) } #[derive(Deserialize)] pub struct EditUploadRequest { pub caption: Option, pub hashtags: Option>, } pub async fn edit_upload( State(state): State, auth: AuthUser, Path(upload_id): Path, Json(body): Json, ) -> Result { // Banned users keep read access but cannot mutate (USER_JOURNEYS §10). if auth.is_banned { return Err(AppError::Forbidden("Du bist gesperrt.".into())); } let upload = Upload::find_by_id_and_event(&state.pool, upload_id, auth.event_id) .await? .ok_or_else(|| AppError::NotFound("Upload nicht gefunden.".into()))?; if upload.user_id != auth.user_id { return Err(AppError::Forbidden("Nur eigene Uploads bearbeiten.".into())); } // Caption update + hashtag wipe-then-relink in one transaction, so a crash // mid-relink can't leave the upload with its hashtags stripped. // // Editing is intentionally allowed while uploads are locked or the gallery is released — like // comments and likes, the lock freezes *new uploads* only (USER_JOURNEYS §9.3). But a caption // is embedded in the HTML viewer keepsake (the ZIP holds media only — see export.rs), so an // edit AFTER release must regenerate the viewer, or the downloadable keepsake keeps showing the // old caption forever while the live feed shows the new one. Same atomicity as delete_upload: // the edit and its invalidation share one tx so a dropped handler can't leave them disagreeing. // `Affects::ViewerOnly` carries the finished ZIP forward (the media didn't change); when the // gallery isn't released, `invalidate_and_arm` returns None and this is a no-op. let mut tx = state.pool.begin().await?; if let Some(ref caption) = body.caption { Upload::update_caption(&mut *tx, upload_id, Some(caption)).await?; } if let Some(ref hashtags) = body.hashtags { Hashtag::unlink_all_from_upload(&mut *tx, upload_id).await?; // Sort + dedup before upserting, exactly as the upload path does. `Hashtag::upsert` // takes row locks, so two transactions touching the same two tags in OPPOSITE order // deadlock; Postgres aborts one after ~1s and the guest gets a 500. Here the order is // whatever the client sent, so it is genuinely attacker-free but genuinely unordered. // Sort on the NORMALISED form — that is the key `upsert` actually locks on. let mut tags: Vec<&String> = hashtags.iter().collect(); tags.sort_by_key(|t| t.trim().trim_start_matches('#').to_lowercase()); tags.dedup_by_key(|t| t.trim().trim_start_matches('#').to_lowercase()); for tag in tags { let h = Hashtag::upsert(&mut *tx, auth.event_id, tag).await?; Hashtag::link_to_upload(&mut *tx, upload_id, h.id).await?; } } let regen = crate::services::export::invalidate_and_arm( &mut tx, &state.config.event_slug, crate::services::export::Affects::ViewerOnly, ) .await?; tx.commit().await?; if let Some(r) = regen { crate::handlers::host::start_regen(&state, r); } Ok(StatusCode::OK) } pub async fn delete_upload( State(state): State, auth: AuthUser, Path(upload_id): Path, ) -> Result { // Banned users keep read access but cannot mutate (USER_JOURNEYS §10). if auth.is_banned { return Err(AppError::Forbidden("Du bist gesperrt.".into())); } let upload = Upload::find_by_id_and_event(&state.pool, upload_id, auth.event_id) .await? .ok_or_else(|| AppError::NotFound("Upload nicht gefunden.".into()))?; if upload.user_id != auth.user_id { return Err(AppError::Forbidden("Nur eigene Uploads löschen.".into())); } // Atomic with the keepsake invalidation: a guest removing their own photo must have it removed // from the downloadable archive too, and a half-applied delete would leave it there forever. let mut tx = state.pool.begin().await?; Upload::soft_delete_in_event(&mut tx, upload_id, auth.event_id).await?; let regen = crate::services::export::invalidate_and_arm( &mut tx, &state.config.event_slug, crate::services::export::Affects::Both, ) .await?; tx.commit().await?; if let Some(r) = regen { crate::handlers::host::start_regen(&state, r); } // Evict the card live on every other feed + the projector diashow — otherwise // a self-deleted post lingers until each viewer manually reloads. Same event // the host-delete path already emits and the frontend already handles. let _ = state.sse_tx.send(crate::state::SseEvent::new( "upload-deleted", serde_json::json!({ "upload_id": upload_id }).to_string(), )); Ok(StatusCode::NO_CONTENT) } /// Number of leading bytes retained in memory for magic-byte (`infer`) sniffing. Every /// allowed type's signature sits well within this; 512 is comfortably generous. const HEAD_SNIFF_BYTES: usize = 512; /// Stream a multipart field straight to `dest`, aborting with a 400 the moment it /// exceeds `max_bytes`. Only the first [`HEAD_SNIFF_BYTES`] bytes are kept in memory /// (for type detection); the rest goes chunk-by-chunk to disk, so peak memory is a /// single chunk rather than the whole file. Returns `(total_size, head_bytes)`. On any /// error the partial temp file is removed so no stray `.tmp` is left behind. async fn stream_field_to_file( mut field: axum::extract::multipart::Field<'_>, dest: &std::path::Path, max_bytes: usize, ) -> Result<(i64, Vec), AppError> { use tokio::io::AsyncWriteExt; let mut file = tokio::fs::File::create(dest) .await .map_err(|e| AppError::Internal(e.into()))?; let mut total: usize = 0; let mut head: Vec = Vec::with_capacity(HEAD_SNIFF_BYTES); loop { let chunk = match field.chunk().await { Ok(Some(c)) => c, Ok(None) => break, Err(e) => { let _ = file.shutdown().await; let _ = tokio::fs::remove_file(dest).await; return Err(AppError::BadRequest(format!( "Datei konnte nicht gelesen werden: {e}" ))); } }; total = total.saturating_add(chunk.len()); if total > max_bytes { let _ = file.shutdown().await; let _ = tokio::fs::remove_file(dest).await; return Err(AppError::BadRequest(format!( "Datei ist zu groß. Maximum: {} MB.", max_bytes / (1024 * 1024) ))); } if head.len() < HEAD_SNIFF_BYTES { let need = HEAD_SNIFF_BYTES - head.len(); head.extend_from_slice(&chunk[..need.min(chunk.len())]); } if let Err(e) = file.write_all(&chunk).await { let _ = tokio::fs::remove_file(dest).await; return Err(AppError::Internal(e.into())); } } if let Err(e) = file.flush().await { let _ = tokio::fs::remove_file(dest).await; return Err(AppError::Internal(e.into())); } Ok((total as i64, head)) } /// Sentinel for the duplicate detected INSIDE the commit transaction. It never reaches a client: /// the caller intercepts this exact `Conflict` and answers with the stored row. A marker rather /// than a new `AppError` variant because the condition is local to this one handler and returning /// early is the only way to abandon the transaction from inside the async block. const DUPLICATE_UPLOAD_MARKER: &str = "__duplicate_client_upload_id__"; /// Rebuild the response for an upload that already exists, so a retry is answered exactly as the /// original was. /// /// Reads the live state rather than assuming a fresh row: by the time a retry arrives — a /// reconnect can be minutes later — the derivatives may have been generated and the photo may /// already have been liked, and a response claiming otherwise would be wrong in a way the client /// has no way to detect. /// /// Every read here fails soft. This is the success path of an upload that is already safely /// stored; degrading to a sparser response is fine, failing the request is not. async fn replay_upload_dto(state: &AppState, upload: &Upload, uploader_name: &str) -> UploadDto { let hashtags: Vec = sqlx::query_scalar( "SELECT h.tag FROM upload_hashtag uh JOIN hashtag h ON h.id = uh.hashtag_id WHERE uh.upload_id = $1 ORDER BY h.tag", ) .bind(upload.id) .fetch_all(&state.pool) .await .unwrap_or_default(); let counts: Option<(i64, i64, bool)> = sqlx::query_as( "SELECT v.like_count, v.comment_count, EXISTS (SELECT 1 FROM \"like\" l WHERE l.upload_id = v.id AND l.user_id = $2) FROM v_feed v WHERE v.id = $1", ) .bind(upload.id) .bind(upload.user_id) .fetch_optional(&state.pool) .await .ok() .flatten(); let (like_count, comment_count, liked_by_me) = counts.unwrap_or((0, 0, false)); UploadDto { id: upload.id, user_id: upload.user_id, uploader_name: uploader_name.to_string(), preview_url: upload .preview_path .as_ref() .map(|_| format!("/api/v1/upload/{}/preview", upload.id)), thumbnail_url: upload .thumbnail_path .as_ref() .map(|_| format!("/api/v1/upload/{}/thumbnail", upload.id)), mime_type: upload.mime_type.clone(), caption: upload.caption.clone(), hashtags, like_count, comment_count, liked_by_me, created_at: upload.created_at, } } /// Drain a multipart body so the HTTP connection stays clean when returning an early error. /// Without draining, the client may still be sending the body after we've sent our response, /// which can corrupt the keep-alive connection for subsequent requests. async fn drain_multipart(mut mp: Multipart) { while let Ok(Some(mut field)) = mp.next_field().await { while field.chunk().await.ok().flatten().is_some() {} } } /// Snapshot of the dynamic per-user quota used both by the upload pre-check and the /// `GET /me/quota` endpoint. `limit_bytes = None` means quota enforcement is currently /// off (the frontend hides the widget in that case). pub struct QuotaEstimate { pub limit_bytes: Option, pub active_uploaders: i64, pub free_disk_bytes: i64, /// The tolerance factor the limit above was computed with. Carried on the snapshot so the /// number is self-describing; no caller reads it back today. #[allow(dead_code)] pub tolerance: f64, } /// Pure per-user quota formula: `floor((free_disk * tolerance) / max(active, 1))`. /// Extracted from `compute_storage_quota` so it's unit-testable without a DB or disk. fn quota_limit_bytes(free_disk: i64, tolerance: f64, active_uploaders: i64) -> i64 { let active = active_uploaders.max(1); ((free_disk as f64 * tolerance) / active as f64).floor() as i64 } /// Computes the per-user storage quota using /// `floor((free_disk * tolerance) / max(active_uploaders, 1))`. Returns `limit_bytes = /// None` whenever the storage quota is currently disabled — callers should skip the /// check (upload handler) or hide the UI (quota endpoint). pub async fn compute_storage_quota(state: &AppState) -> QuotaEstimate { let quota_on = config::get_bool(&state.config_cache, "quota_enabled", true).await; let storage_quota_on = config::get_bool(&state.config_cache, "storage_quota_enabled", true).await; let tolerance = config::get_f64(&state.config_cache, "quota_tolerance", 0.75).await; let (active_count,): (i64,) = sqlx::query_as("SELECT COUNT(DISTINCT user_id) FROM upload WHERE deleted_at IS NULL") .fetch_one(&state.pool) .await .unwrap_or((0,)); let active = active_count.max(1); // Cached disk reading. `None` means we couldn't resolve the media filesystem. let disk = state.disk_cache.snapshot(&state.config.media_path); let free_disk = disk.map(|d| d.free as i64).unwrap_or(0); let limit_bytes = if quota_on && storage_quota_on { match disk { Some(d) => Some(quota_limit_bytes(d.free as i64, tolerance, active)), // Fail OPEN, not closed: if the disk can't be read we don't know the real // free space, and enforcing a 0-byte limit would reject every upload with a // spurious "quota reached". Skip enforcement this round and warn instead. None => { tracing::warn!( "disk snapshot unavailable; skipping storage-quota enforcement this round" ); None } } } else { None }; QuotaEstimate { limit_bytes, active_uploaders: active, free_disk_bytes: free_disk, tolerance, } } /// Outcome of parsing a `Range` request header against a known file length. #[derive(Debug, PartialEq, Eq)] enum RangeSpec { /// No `Range` header, or one we deliberately don't honour (multi-range, non-`bytes` /// unit, malformed). RFC 9110 lets a server ignore a Range it can't process and reply /// 200 with the full body, which is what every one of these cases does. Full, /// A single satisfiable range, resolved to inclusive absolute offsets. Partial { start: u64, end: u64 }, /// Syntactically valid but starts beyond EOF — must be answered 416, not 200, or a /// player can loop re-requesting it. Unsatisfiable, } /// Parse a single-range `bytes=` header against `len`. /// /// Deliberately supports only the three forms a media element actually sends — /// `bytes=N-`, `bytes=N-M`, `bytes=-S` (suffix) — and treats everything else as `Full`. /// Multi-range responses need `multipart/byteranges`, which no `