diff --git a/backend/src/handlers/host.rs b/backend/src/handlers/host.rs index ffc15cb..29ce39f 100644 --- a/backend/src/handlers/host.rs +++ b/backend/src/handlers/host.rs @@ -403,35 +403,49 @@ pub async fn release_gallery( .await? .ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?; - // Atomically claim the release (M8). Two concurrent presses can no longer - // both spawn export workers racing on the same files: only the request whose - // UPDATE matches a row proceeds. Re-release is allowed when the gallery was - // never released, OR every existing export job terminally failed (e.g. a - // crash mid-export) — otherwise the headline deliverable would be - // permanently unrecoverable. + // Atomically claim the release AND enqueue the jobs in one transaction (M8). + // The claim UPDATE row-locks the event row and the lock is held until commit + // — *after* the export_job rows exist — so a second concurrent press re-reads + // the now-present jobs and matches 0 rows (clean 400). Doing the UPDATE and + // INSERTs as separate autocommit statements left a cross-table TOCTOU window + // where both presses could win and race the same output files. + // + // Re-release guard (M8-3): allow when never released, OR nothing is currently + // in progress and at least one job terminally failed. The old "every job + // failed" guard left a one-sided failure (zip done, html failed) permanently + // unrecoverable. + let mut tx = state.pool.begin().await?; + let claim = sqlx::query( "UPDATE event SET export_released_at = NOW() WHERE slug = $1 AND ( export_released_at IS NULL - OR NOT EXISTS ( - SELECT 1 FROM export_job j - WHERE j.event_id = event.id AND j.status <> 'failed' + OR ( + NOT EXISTS ( + SELECT 1 FROM export_job j + WHERE j.event_id = event.id AND j.status IN ('running', 'pending') + ) + AND EXISTS ( + SELECT 1 FROM export_job j + WHERE j.event_id = event.id AND j.status = 'failed' + ) ) )", ) .bind(&state.config.event_slug) - .execute(&state.pool) + .execute(&mut *tx) .await?; if claim.rows_affected() == 0 { + tx.rollback().await.ok(); return Err(AppError::BadRequest( "Galerie wurde bereits freigegeben.".into(), )); } - // Enqueue export jobs, resetting any prior terminally-failed rows so the - // workers re-run cleanly. + // Enqueue export jobs, resetting any prior terminal rows so the workers + // re-run cleanly (including an already-`done` sibling on re-release). for export_type in ["zip", "html"] { sqlx::query( "INSERT INTO export_job (event_id, type) VALUES ($1, $2::export_type) @@ -441,11 +455,13 @@ pub async fn release_gallery( ) .bind(event.id) .bind(export_type) - .execute(&state.pool) + .execute(&mut *tx) .await?; } - // Spawn export workers + tx.commit().await?; + + // Spawn export workers only after the claim+enqueue is durably committed. crate::services::export::spawn_export_jobs( event.id, event.name, diff --git a/backend/src/handlers/upload.rs b/backend/src/handlers/upload.rs index cb392c9..7072eb9 100644 --- a/backend/src/handlers/upload.rs +++ b/backend/src/handlers/upload.rs @@ -87,6 +87,16 @@ pub async fn upload( let max_image_mb: i64 = config::get_i64(&state.pool, "max_image_size_mb", 20).await; let max_video_mb: i64 = config::get_i64(&state.pool, "max_video_size_mb", 500).await; + // Bound concurrent upload-body buffering to cap aggregate RAM (H3). Acquired + // *before* the multipart body is read, so requests waiting on a permit hold + // only a connection, not a buffered ~550 MB file. The permit is held for the + // rest of the handler (body read + write) and released on return. + let _upload_permit = state + .upload_limiter + .acquire() + .await + .map_err(|e| AppError::Internal(e.into()))?; + let mut file_data: Option = None; let mut caption: Option = None; let mut hashtags_csv: Option = None; diff --git a/backend/src/services/export.rs b/backend/src/services/export.rs index 0665f16..55a6594 100644 --- a/backend/src/services/export.rs +++ b/backend/src/services/export.rs @@ -94,8 +94,14 @@ pub fn spawn_export_jobs( let sse_tx2 = sse_tx.clone(); let event_name2 = event_name.clone(); + // Per-run id so two runs (e.g. a re-release after a crash mid-export, where + // startup_recovery marked the old run 'failed' but its task may still be + // winding down) write to distinct temp paths and can't truncate each other's + // output (M8). The final served filenames stay fixed. + let run_id = Uuid::new_v4(); + tokio::spawn(async move { - if let Err(e) = run_zip_export(event_id, &pool, &media_path, &sse_tx).await { + if let Err(e) = run_zip_export(event_id, run_id, &pool, &media_path, &sse_tx).await { tracing::error!("ZIP export failed for event {event_id}: {e:#}"); mark_failed(&pool, event_id, "zip", &e.to_string()).await; } @@ -104,7 +110,7 @@ pub fn spawn_export_jobs( tokio::spawn(async move { if let Err(e) = - run_html_export(event_id, &event_name2, &pool2, &media_path2, &sse_tx2).await + run_html_export(event_id, run_id, &event_name2, &pool2, &media_path2, &sse_tx2).await { tracing::error!("HTML export failed for event {event_id}: {e:#}"); mark_failed(&pool2, event_id, "html", &e.to_string()).await; @@ -117,6 +123,7 @@ pub fn spawn_export_jobs( async fn run_zip_export( event_id: Uuid, + run_id: Uuid, pool: &PgPool, media_path: &Path, sse_tx: &broadcast::Sender, @@ -129,7 +136,8 @@ async fn run_zip_export( let exports_dir = media_path.join("exports"); tokio::fs::create_dir_all(&exports_dir).await?; - let tmp_path = exports_dir.join("Gallery.zip.tmp"); + // Per-run temp name; final served name stays fixed. + let tmp_path = exports_dir.join(format!("Gallery.{run_id}.zip.tmp")); let out_path = exports_dir.join("Gallery.zip"); { @@ -190,6 +198,7 @@ async fn run_zip_export( async fn run_html_export( event_id: Uuid, + run_id: Uuid, event_name: &str, pool: &PgPool, media_path: &Path, @@ -208,8 +217,8 @@ async fn run_html_export( let exports_dir = media_path.join("exports"); tokio::fs::create_dir_all(&exports_dir).await?; - // 2. Create temp directory for media processing - let tmp_dir = exports_dir.join(format!("viewer_tmp_{event_id}")); + // 2. Create temp directory for media processing (per-run, see run_id). + let tmp_dir = exports_dir.join(format!("viewer_tmp_{event_id}_{run_id}")); let media_tmp = tmp_dir.join("media"); tokio::fs::create_dir_all(&media_tmp).await?; @@ -370,8 +379,8 @@ async fn run_html_export( update_progress(pool, event_id, "html", 72).await; - // 5. Create ZIP - let tmp_path = exports_dir.join("Memories.zip.tmp"); + // 5. Create ZIP (per-run temp name; final served name stays fixed). + let tmp_path = exports_dir.join(format!("Memories.{run_id}.zip.tmp")); let out_path = exports_dir.join("Memories.zip"); { diff --git a/backend/src/state.rs b/backend/src/state.rs index 31ffb5a..8cfc5af 100644 --- a/backend/src/state.rs +++ b/backend/src/state.rs @@ -1,11 +1,20 @@ +use std::sync::Arc; + use sqlx::PgPool; -use tokio::sync::broadcast; +use tokio::sync::{broadcast, Semaphore}; use crate::config::AppConfig; use crate::services::compression::CompressionWorker; use crate::services::rate_limiter::RateLimiter; use crate::services::sse_tickets::SseTicketStore; +/// Max concurrent in-flight `/upload` requests. Each holds its whole file in +/// memory while reading the multipart body, so this bounds aggregate upload RAM +/// to ~N × the 550 MB body cap (≈2.2 GB here) — keeping headroom for Postgres + +/// the app on an 8 GB box. The per-user rate limiter is a *count* limiter +/// (10/hr), not a concurrency cap, so this is the actual OOM guard. Tunable. +pub const UPLOAD_MAX_CONCURRENCY: usize = 4; + #[derive(Clone, Debug)] pub struct SseEvent { pub event_type: String, @@ -31,6 +40,8 @@ pub struct AppState { pub compression: CompressionWorker, pub rate_limiter: RateLimiter, pub sse_tickets: SseTicketStore, + /// Caps concurrent upload-body buffering (see UPLOAD_MAX_CONCURRENCY). + pub upload_limiter: Arc, } impl AppState { @@ -51,6 +62,7 @@ impl AppState { compression, rate_limiter: RateLimiter::new(), sse_tickets: SseTicketStore::new(), + upload_limiter: Arc::new(Semaphore::new(UPLOAD_MAX_CONCURRENCY)), } } } diff --git a/frontend/src/lib/components/LightboxModal.svelte b/frontend/src/lib/components/LightboxModal.svelte index 32a7a6f..6198d22 100644 --- a/frontend/src/lib/components/LightboxModal.svelte +++ b/frontend/src/lib/components/LightboxModal.svelte @@ -218,8 +218,8 @@
= 450 && newComment.length < COMMENT_MAX} class:dark:text-amber-400={newComment.length >= 450 && newComment.length < COMMENT_MAX} class:text-red-600={newComment.length >= COMMENT_MAX} diff --git a/frontend/src/lib/sse.ts b/frontend/src/lib/sse.ts index 3c897ad..c04e36c 100644 --- a/frontend/src/lib/sse.ts +++ b/frontend/src/lib/sse.ts @@ -38,6 +38,7 @@ const KNOWN_EVENTS = [ 'upload-deleted', 'like-update', 'new-comment', + 'user-banned', 'event-closed', 'event-opened', 'event-updated', diff --git a/frontend/src/routes/account/+page.svelte b/frontend/src/routes/account/+page.svelte index 531d702..3c96f77 100644 --- a/frontend/src/routes/account/+page.svelte +++ b/frontend/src/routes/account/+page.svelte @@ -219,7 +219,7 @@
{:else} -
+
PIN nicht gespeichert. Nutze die Wiederherstellungs-Seite, um dich mit deinem PIN anzumelden.
{/if} diff --git a/frontend/src/routes/diashow/+page.svelte b/frontend/src/routes/diashow/+page.svelte index e8d5bd6..5f44afa 100644 --- a/frontend/src/routes/diashow/+page.svelte +++ b/frontend/src/routes/diashow/+page.svelte @@ -35,7 +35,15 @@ clearTimer(); if (paused) return; // Videos: advance on `ended` or after `max(dwell, 12s)` — whichever first. - const ms = isVideo ? Math.max(dwellMs, 12000) : dwellMs; + let ms = isVideo ? Math.max(dwellMs, 12000) : dwellMs; + // Reduced-motion: slow auto-advance to a ≥30s floor so content doesn't + // auto-change rapidly (WCAG 2.2.2); the CSS already snaps the transitions. + if ( + typeof window !== 'undefined' && + window.matchMedia?.('(prefers-reduced-motion: reduce)').matches + ) { + ms = Math.max(ms, 30000); + } advanceTimer = setTimeout(advance, ms); }