use std::path::{Path, PathBuf}; use anyhow::{Context, Result}; use async_zip::tokio::write::ZipFileWriter; use async_zip::{Compression, ZipEntryBuilder}; use chrono::{DateTime, Utc}; use futures::io::{copy as fcopy, AllowStdIo}; use include_dir::{include_dir, Dir}; use serde::Serialize; use sqlx::PgPool; use tokio::sync::broadcast; use tokio_util::compat::TokioAsyncReadCompatExt; use uuid::Uuid; use crate::state::SseEvent; // ── Embedded viewer assets (pre-built SvelteKit static output) ────────────── static VIEWER_DIR: Dir<'_> = include_dir!("$CARGO_MANIFEST_DIR/static/export-viewer"); // ── DB query rows ──────────────────────────────────────────────────────────── #[derive(sqlx::FromRow)] struct ExportUploadRow { id: Uuid, original_path: String, mime_type: String, caption: Option, uploader_name: String, like_count: i64, created_at: DateTime, } #[derive(sqlx::FromRow)] struct ExportCommentRow { upload_id: Uuid, uploader_name: String, body: String, created_at: DateTime, } // ── Viewer JSON structs (serialised to data.json) ─────────────────────────── #[derive(Serialize)] struct ViewerData { event: ViewerEvent, posts: Vec, } #[derive(Serialize)] struct ViewerEvent { name: String, exported_at: String, } #[derive(Serialize)] struct ViewerPost { id: String, uploader: String, caption: String, tags: Vec, timestamp: String, likes: i64, comments: Vec, media: ViewerMedia, } #[derive(Serialize)] struct ViewerComment { author: String, text: String, timestamp: String, } #[derive(Serialize)] struct ViewerMedia { #[serde(rename = "type")] media_type: String, thumb: String, full: String, } // ── Entry point ────────────────────────────────────────────────────────────── /// (Re)enqueue the export jobs for an event and spawn the workers. Safe to call on a /// fresh release *and* on a re-release: the `export_job` rows are reset to a clean /// `pending` state (clearing any prior `failed`/`done` from an earlier run) and the /// event's `export_*_ready` flags are cleared so downloads reflect the regenerating /// export rather than the stale one. This is the single path used by `release_gallery` /// and by startup export recovery, so the two can't drift. pub async fn enqueue_and_spawn_exports( event_id: Uuid, event_name: String, pool: PgPool, media_path: PathBuf, export_path: PathBuf, sse_tx: broadcast::Sender, ) -> Result<()> { // Only (re)generate a type that isn't ALREADY complete. A (re)release always arrives // with both ready flags false (a fresh release starts false; `open_event` clears them // before a re-release), so both regenerate as intended. But startup recovery of a // half-finished export (e.g. ZIP done+ready, HTML crashed) must leave the good half // alone — resetting it would 404 an already-downloadable keepsake until it needlessly // regenerates. Skip the ready ones; their still-`done` rows keep their `file_path`, and // the worker we spawn for them below simply bails its claim (no `pending` row to grab). let (zip_ready, html_ready): (bool, bool) = sqlx::query_as( "SELECT export_zip_ready, export_html_ready FROM event WHERE id = $1", ) .bind(event_id) .fetch_one(&pool) .await?; for (export_type, ready) in [("zip", zip_ready), ("html", html_ready)] { if ready { continue; } // Reset the row to a clean 'pending' and BUMP `release_seq` — unconditionally, even // if a worker from a prior release is still 'running'. That worker captured the old // seq at claim time; bumping it here supersedes its generation, so when it finishes // its guarded finalize (`WHERE release_seq = `) matches nothing and it // discards its now-stale output instead of overwriting a fresh keepsake. The worker // we spawn below then claims this new generation and regenerates from a current // snapshot. Per-seq temp/final paths (see `run_zip_export`/`run_html_export`) keep // the old and new workers from racing the same file on disk. sqlx::query( "INSERT INTO export_job (event_id, type, status, progress_pct, release_seq) VALUES ($1, $2::export_type, 'pending', 0, 1) ON CONFLICT (event_id, type) DO UPDATE SET status = 'pending', progress_pct = 0, file_path = NULL, error_message = NULL, completed_at = NULL, release_seq = export_job.release_seq + 1", ) .bind(event_id) .bind(export_type) .execute(&pool) .await?; } spawn_export_jobs(event_id, event_name, pool, media_path, export_path, sse_tx); Ok(()) } /// Startup export recovery: re-spawn exports for any released event whose keepsake is /// not fully ready (a crash mid-export left `export_released_at` set but the ZIP/HTML /// jobs `failed`). Without this, `release_gallery` would reject a retry with "bereits /// freigegeben" and downloads would 404 forever. Runs once at boot, after `AppState` /// exists (it needs the media/export paths + SSE sender). pub async fn recover_exports( pool: PgPool, media_path: PathBuf, export_path: PathBuf, sse_tx: broadcast::Sender, ) { let rows = match sqlx::query_as::<_, (Uuid, String)>( "SELECT id, name FROM event WHERE export_released_at IS NOT NULL AND NOT (export_zip_ready AND export_html_ready)", ) .fetch_all(&pool) .await { Ok(r) => r, Err(e) => { tracing::error!("export recovery: failed to query released events: {e:#}"); return; } }; for (event_id, event_name) in rows { tracing::warn!("export recovery: re-spawning export jobs for event {event_id}"); if let Err(e) = enqueue_and_spawn_exports( event_id, event_name, pool.clone(), media_path.clone(), export_path.clone(), sse_tx.clone(), ) .await { tracing::error!("export recovery: failed to re-spawn for event {event_id}: {e:#}"); } } } pub fn spawn_export_jobs( event_id: Uuid, event_name: String, pool: PgPool, media_path: PathBuf, export_path: PathBuf, sse_tx: broadcast::Sender, ) { let pool2 = pool.clone(); let media_path2 = media_path.clone(); let export_path2 = export_path.clone(); let sse_tx2 = sse_tx.clone(); let event_name2 = event_name.clone(); tokio::spawn(async move { // Failure marking happens INSIDE the worker (guarded by the claimed `release_seq`), // so a superseded worker's error can't clobber the fresh generation's row. if let Err(e) = run_zip_export(event_id, &pool, &media_path, &export_path, &sse_tx).await { tracing::error!("ZIP export failed for event {event_id}: {e:#}"); } maybe_broadcast_complete(&pool, event_id, &sse_tx).await; }); tokio::spawn(async move { if let Err(e) = run_html_export(event_id, &event_name2, &pool2, &media_path2, &export_path2, &sse_tx2).await { tracing::error!("HTML export failed for event {event_id}: {e:#}"); } maybe_broadcast_complete(&pool2, event_id, &sse_tx2).await; }); } // ── ZIP export ─────────────────────────────────────────────────────────────── async fn run_zip_export( event_id: Uuid, pool: &PgPool, media_path: &Path, export_path: &Path, sse_tx: &broadcast::Sender, ) -> Result<()> { let seq = match claim_job(pool, event_id, "zip").await { Some(seq) => seq, // Another worker already owns this ZIP export — bail rather than race the temp file. None => return Ok(()), }; // Run the body under the captured `seq`; on error mark THIS generation failed (a no-op // if a re-release already superseded us). let res = run_zip_export_inner(seq, event_id, pool, media_path, export_path, sse_tx).await; if let Err(e) = &res { mark_failed(pool, event_id, "zip", seq, &e.to_string()).await; } res } async fn run_zip_export_inner( seq: i64, event_id: Uuid, pool: &PgPool, media_path: &Path, export_path: &Path, sse_tx: &broadcast::Sender, ) -> Result<()> { let uploads = query_uploads(pool, event_id).await?; let total = uploads.len().max(1) as f32; // Written OUTSIDE media_path: the public /media ServeDir must never reach these. let exports_dir = export_path.to_path_buf(); tokio::fs::create_dir_all(&exports_dir).await?; // Per-generation paths: a superseded worker (older `seq`) and the fresh worker never // share a file on disk, so neither can truncate or interleave the other's bytes. The // final name embeds the seq too; the download handler serves whatever `file_path` the // current 'done' row points at, so an orphaned older generation is never served. let tmp_path = exports_dir.join(format!("Gallery.zip.{seq}.tmp")); let out_name = format!("Gallery.{seq}.zip"); let out_path = exports_dir.join(&out_name); { let file = tokio::fs::File::create(&tmp_path).await?; let mut zip = ZipFileWriter::with_tokio(file); for (i, row) in uploads.iter().enumerate() { let src = media_path.join(&row.original_path); if !src.exists() { continue; } let ext = ext_from_path(&row.original_path); let date = row.created_at.format("%Y-%m-%d_%H-%M").to_string(); let name_safe = sanitize_name(&row.uploader_name); let folder = if row.mime_type.starts_with("video/") { "Videos" } else { "Photos" }; let entry_name = format!("{folder}/{date}_{name_safe}_{}.{ext}", row.id); let builder = ZipEntryBuilder::new(entry_name.into(), Compression::Stored); let mut entry = zip.write_entry_stream(builder).await?; let mut f = tokio::fs::File::open(&src).await?.compat(); fcopy(&mut f, &mut entry).await?; entry.close().await?; let pct = ((i + 1) as f32 / total * 100.0) as i16; update_progress(pool, event_id, "zip", seq, pct.min(99)).await; } zip.close().await?; } tokio::fs::rename(&tmp_path, &out_path).await?; // Guarded finalize: commit ONLY if our generation is still current. If a reopen→ // re-release bumped `release_seq` while we were streaming, we lost — discard our (now // stale) archive and let the fresh worker own the keepsake. This is the fix for the // stale-keepsake data loss: a superseded worker can no longer flip `export_zip_ready`. if !finalize_job(pool, event_id, "zip", seq, &format!("exports/{out_name}")).await { let _ = tokio::fs::remove_file(&out_path).await; tracing::info!("ZIP export for event {event_id} superseded by a newer release; discarded"); return Ok(()); } // Flip the ready flag only while still current AND still released. The `release_seq` // EXISTS check alone is NOT enough: a reopen (`open_event`) landing in the window between // our `finalize_job` above and this UPDATE clears `export_released_at` + both ready flags // but leaves our `export_job` row `done` at this same seq — so EXISTS would still match and // we'd resurrect `export_zip_ready = TRUE` on a keepsake that predates the reopen. The next // re-release would then read that stale TRUE, skip regeneration (`if ready { continue }`), // and serve a snapshot missing every upload added during the reopen window — the exact // stale-keepsake data loss migration 012 exists to prevent. Anchoring on // `export_released_at IS NOT NULL` makes the flip a no-op once a reopen has landed. sqlx::query( "UPDATE event SET export_zip_ready = TRUE WHERE id = $1 AND export_released_at IS NOT NULL AND EXISTS (SELECT 1 FROM export_job WHERE event_id = $1 AND type = 'zip'::export_type AND release_seq = $2 AND status = 'done')", ) .bind(event_id) .bind(seq) .execute(pool) .await?; prune_stale_export_files(&exports_dir, "Gallery", event_id, seq).await; let _ = sse_tx.send(SseEvent { event_type: "export-progress".to_string(), data: serde_json::json!({ "type": "zip", "progress_pct": 100 }).to_string(), }); tracing::info!("ZIP export complete for event {event_id}"); Ok(()) } // ── HTML viewer export ────────────────────────────────────────────────────── /// Where a media entry's bytes come from at ZIP-writing time. Derived variants /// (thumbnails, compressed images) are staged under the temp dir; original-fidelity /// variants (videos, small images) are streamed straight from the source on disk so /// the export never transiently doubles disk usage by copying large files to temp. enum MediaSource { Temp(PathBuf), Original(PathBuf), } impl MediaSource { fn path(&self) -> &Path { match self { MediaSource::Temp(p) | MediaSource::Original(p) => p, } } } async fn run_html_export( event_id: Uuid, event_name: &str, pool: &PgPool, media_path: &Path, export_path: &Path, sse_tx: &broadcast::Sender, ) -> Result<()> { let seq = match claim_job(pool, event_id, "html").await { Some(seq) => seq, // Another worker already owns this HTML export — bail rather than race the temp dir. None => return Ok(()), }; let res = run_html_export_inner(seq, event_id, event_name, pool, media_path, export_path, sse_tx).await; if let Err(e) = &res { mark_failed(pool, event_id, "html", seq, &e.to_string()).await; } res } async fn run_html_export_inner( seq: i64, event_id: Uuid, event_name: &str, pool: &PgPool, media_path: &Path, export_path: &Path, sse_tx: &broadcast::Sender, ) -> Result<()> { // 1. Query data let uploads = query_uploads(pool, event_id).await?; let comments = query_comments(pool, event_id).await?; let hashtags_per_upload = query_hashtags(pool, event_id).await?; let total = uploads.len().max(1) as f32; update_progress(pool, event_id, "html", seq, 5).await; // Written OUTSIDE media_path: the public /media ServeDir must never reach these. let exports_dir = export_path.to_path_buf(); tokio::fs::create_dir_all(&exports_dir).await?; // 2. Create temp directory for media processing (per-generation — see run_zip_export). let tmp_dir = exports_dir.join(format!("viewer_tmp_{event_id}_{seq}")); let media_tmp = tmp_dir.join("media"); tokio::fs::create_dir_all(&media_tmp).await?; // 3. Process media and build post data let mut viewer_posts: Vec = Vec::new(); // (zip entry name under media/, where its bytes come from). Built here, streamed // into the ZIP in step 5 — so we also know the exact file count without a rescan. let mut media_manifest: Vec<(String, MediaSource)> = Vec::new(); for (i, row) in uploads.iter().enumerate() { let src = media_path.join(&row.original_path); if !src.exists() { continue; } let is_video = row.mime_type.starts_with("video/"); let id_str = row.id.to_string(); // Generate thumbnail and full variant. `full_source` says where the full-res // bytes come from at ZIP time — for videos and small images that's the original // on disk (streamed directly, never copied to temp). let (thumb_name, full_name, full_source) = if is_video { let thumb = format!("{id_str}_thumb.jpg"); let full_ext = ext_from_path(&row.original_path); let full = format!("{id_str}.{full_ext}"); // Video thumbnail via ffmpeg let thumb_path = media_tmp.join(&thumb); let ffmpeg_result = tokio::process::Command::new("ffmpeg") .args([ "-i", src.to_str().unwrap_or_default(), "-vframes", "1", "-ss", "00:00:01", "-vf", "scale=400:-1", "-y", thumb_path.to_str().unwrap_or_default(), ]) .output() .await; match ffmpeg_result { Ok(output) if output.status.success() => {} _ => { tracing::warn!("ffmpeg thumbnail failed for upload {}, skipping thumb", row.id); // Missing thumb entry — viewer handles missing thumbs gracefully. } } // Stream the video full-res straight from the original at ZIP time — no // copy to temp (that used to transiently double disk usage per video). (thumb, full, MediaSource::Original(src.clone())) } else { let thumb = format!("{id_str}_thumb.jpg"); let ext = ext_from_path(&row.original_path); let full = format!("{id_str}_full.{ext}"); // Image thumbnail: resize to 400px wide let src_clone = src.clone(); let thumb_path = media_tmp.join(&thumb); let thumb_path_clone = thumb_path.clone(); let thumb_result = tokio::task::spawn_blocking(move || -> Result<()> { let img = image::open(&src_clone).context("failed to open image for thumbnail")?; let resized = img.resize(400, 400, image::imageops::FilterType::Lanczos3); resized .save_with_format(&thumb_path_clone, image::ImageFormat::Jpeg) .context("failed to save thumbnail")?; Ok(()) }) .await?; if let Err(e) = thumb_result { tracing::warn!("thumbnail generation failed for upload {}: {e:#}", row.id); } // Full variant: compress to temp if >5MB, otherwise stream the original // as-is (no temp copy). let src_meta = tokio::fs::metadata(&src).await?; let full_source = if src_meta.len() > 5_000_000 { let src_clone = src.clone(); let full_path = media_tmp.join(&full); let full_path_clone = full_path.clone(); let compress_result = tokio::task::spawn_blocking(move || -> Result<()> { let img = image::open(&src_clone).context("failed to open image for compression")?; let resized = img.resize(2000, 2000, image::imageops::FilterType::Lanczos3); resized .save_with_format(&full_path_clone, image::ImageFormat::Jpeg) .context("failed to save compressed full image")?; Ok(()) }) .await?; match compress_result { Ok(()) => MediaSource::Temp(full_path), Err(e) => { tracing::warn!( "compression failed for upload {}, using original: {e:#}", row.id ); MediaSource::Original(src.clone()) } } } else { MediaSource::Original(src.clone()) }; (thumb, full, full_source) }; // Register this post's two media entries. Thumbnails always come from temp // (they're freshly generated); the full variant's source was decided above. media_manifest.push(( thumb_name.clone(), MediaSource::Temp(media_tmp.join(&thumb_name)), )); media_manifest.push((full_name.clone(), full_source)); // Build comments for this upload let post_comments: Vec = comments .iter() .filter(|c| c.upload_id == row.id) .map(|c| ViewerComment { author: c.uploader_name.clone(), text: c.body.clone(), timestamp: c.created_at.to_rfc3339(), }) .collect(); // Build tags for this upload let tags: Vec = hashtags_per_upload .iter() .filter(|(uid, _)| *uid == row.id) .map(|(_, tag)| tag.clone()) .collect(); viewer_posts.push(ViewerPost { id: id_str, uploader: row.uploader_name.clone(), caption: row.caption.clone().unwrap_or_default(), tags, timestamp: row.created_at.to_rfc3339(), likes: row.like_count, comments: post_comments, media: ViewerMedia { media_type: if is_video { "video".to_string() } else { "image".to_string() }, thumb: format!("media/{thumb_name}"), full: format!("media/{full_name}"), }, }); let pct = 10 + ((i + 1) as f32 / total * 60.0) as i16; update_progress(pool, event_id, "html", seq, pct.min(69)).await; } // 4. Build data.json let viewer_data = ViewerData { event: ViewerEvent { name: event_name.to_string(), exported_at: Utc::now().to_rfc3339(), }, posts: viewer_posts, }; let data_json = serde_json::to_string_pretty(&viewer_data).context("failed to serialize data.json")?; update_progress(pool, event_id, "html", seq, 72).await; // 5. Create ZIP (per-generation paths — see run_zip_export) let tmp_path = exports_dir.join(format!("Memories.zip.{seq}.tmp")); let out_name = format!("Memories.{seq}.zip"); let out_path = exports_dir.join(&out_name); { let file = tokio::fs::File::create(&tmp_path).await?; let mut zip = ZipFileWriter::with_tokio(file); // Write embedded viewer assets (index.html, _app/*, etc.) write_dir_to_zip(&VIEWER_DIR, &mut zip).await?; update_progress(pool, event_id, "html", seq, 75).await; // Write data.json { let builder = ZipEntryBuilder::new("data.json".into(), Compression::Deflate); let mut entry = zip.write_entry_stream(builder).await?; let mut cursor = AllowStdIo::new(std::io::Cursor::new(data_json.as_bytes())); fcopy(&mut cursor, &mut entry).await?; entry.close().await?; } // Write README.txt { let builder = ZipEntryBuilder::new("README.txt".into(), Compression::Deflate); let mut entry = zip.write_entry_stream(builder).await?; let mut cursor = AllowStdIo::new(std::io::Cursor::new(README_TEXT.as_bytes())); fcopy(&mut cursor, &mut entry).await?; entry.close().await?; } update_progress(pool, event_id, "html", seq, 78).await; // Write media files from the manifest built in step 3. Thumbnails and derived // image variants stream from temp; video and small-image full variants stream // straight from the original on disk. Sources that don't exist (e.g. a thumb // whose ffmpeg step failed) are skipped — the viewer tolerates gaps. let file_total = media_manifest.len().max(1) as f32; let mut files_written = 0u32; for (name, source) in &media_manifest { let path = source.path(); if !path.exists() { continue; } let entry_name = format!("media/{name}"); let builder = ZipEntryBuilder::new(entry_name.into(), Compression::Stored); let mut zip_entry = zip.write_entry_stream(builder).await?; let mut f = tokio::fs::File::open(path).await?.compat(); fcopy(&mut f, &mut zip_entry).await?; zip_entry.close().await?; files_written += 1; let pct = 78 + (files_written as f32 / file_total * 20.0) as i16; update_progress(pool, event_id, "html", seq, pct.min(98)).await; } zip.close().await?; } // 6. Finalise tokio::fs::rename(&tmp_path, &out_path).await?; // Clean up temp directory let _ = tokio::fs::remove_dir_all(&tmp_dir).await; // Guarded finalize — discard if a re-release superseded us mid-export (see run_zip_export). if !finalize_job(pool, event_id, "html", seq, &format!("exports/{out_name}")).await { let _ = tokio::fs::remove_file(&out_path).await; tracing::info!("HTML export for event {event_id} superseded by a newer release; discarded"); return Ok(()); } // Same released-anchored guard as the ZIP flip (see run_zip_export): a reopen between our // finalize and here must not resurrect `export_html_ready` on a pre-reopen snapshot. sqlx::query( "UPDATE event SET export_html_ready = TRUE WHERE id = $1 AND export_released_at IS NOT NULL AND EXISTS (SELECT 1 FROM export_job WHERE event_id = $1 AND type = 'html'::export_type AND release_seq = $2 AND status = 'done')", ) .bind(event_id) .bind(seq) .execute(pool) .await?; prune_stale_export_files(&exports_dir, "Memories", event_id, seq).await; let _ = sse_tx.send(SseEvent { event_type: "export-progress".to_string(), data: serde_json::json!({ "type": "html", "progress_pct": 100 }).to_string(), }); tracing::info!("HTML viewer export complete for event {event_id}"); Ok(()) } // ── DB helpers ─────────────────────────────────────────────────────────────── async fn query_uploads(pool: &PgPool, event_id: Uuid) -> Result> { Ok(sqlx::query_as::<_, ExportUploadRow>( "SELECT u.id, u.original_path, u.mime_type, u.caption, usr.display_name AS uploader_name, COUNT(DISTINCT l.user_id) AS like_count, u.created_at FROM upload u JOIN \"user\" usr ON usr.id = u.user_id LEFT JOIN \"like\" l ON l.upload_id = u.id WHERE u.event_id = $1 AND u.deleted_at IS NULL AND usr.uploads_hidden = FALSE AND usr.is_banned = FALSE GROUP BY u.id, usr.display_name ORDER BY u.created_at ASC", ) .bind(event_id) .fetch_all(pool) .await?) } async fn query_comments(pool: &PgPool, event_id: Uuid) -> Result> { Ok(sqlx::query_as::<_, ExportCommentRow>( "SELECT c.upload_id, usr.display_name AS uploader_name, c.body, c.created_at FROM comment c JOIN \"user\" usr ON usr.id = c.user_id JOIN upload u ON u.id = c.upload_id WHERE u.event_id = $1 AND c.deleted_at IS NULL AND u.deleted_at IS NULL ORDER BY c.created_at ASC", ) .bind(event_id) .fetch_all(pool) .await?) } async fn query_hashtags(pool: &PgPool, event_id: Uuid) -> Result> { let rows: Vec<(Uuid, String)> = sqlx::query_as( "SELECT uh.upload_id, h.tag FROM upload_hashtag uh JOIN hashtag h ON h.id = uh.hashtag_id JOIN upload u ON u.id = uh.upload_id WHERE h.event_id = $1 AND u.deleted_at IS NULL", ) .bind(event_id) .fetch_all(pool) .await?; Ok(rows) } /// Atomically claim a pending export job for this worker and return the generation /// (`release_seq`) it claimed. Returns `Some(seq)` only if we won (status flipped /// `pending`→`running` in one statement); `None` when another worker already owns it — e.g. /// a reopen→re-release spawned a second worker while a run from a prior release is still in /// flight. The row-level lock serializes the two UPDATEs, so exactly one sees /// `status = 'pending'`. The caller threads the returned seq into per-generation temp/final /// paths and into the guarded finalize, so a superseded worker never corrupts or resurrects /// a stale keepsake. async fn claim_job(pool: &PgPool, event_id: Uuid, export_type: &str) -> Option { sqlx::query_as::<_, (i64,)>( "UPDATE export_job SET status = 'running' WHERE event_id = $1 AND type = $2::export_type AND status = 'pending' RETURNING release_seq", ) .bind(event_id) .bind(export_type) .fetch_optional(pool) .await .ok() .flatten() .map(|(seq,)| seq) } /// Guarded finalize: mark this export `done` and record its file path ONLY if our /// generation is still current (`release_seq = seq` and still `running`). Returns `true` if /// we won the finalize, `false` if a re-release superseded us mid-export (in which case the /// caller must discard its output — the fresh worker owns the keepsake now). Atomic, so it /// can't race a concurrent re-release. async fn finalize_job( pool: &PgPool, event_id: Uuid, export_type: &str, seq: i64, file_path: &str, ) -> bool { sqlx::query( "UPDATE export_job SET status = 'done', progress_pct = 100, file_path = $3, completed_at = NOW() WHERE event_id = $1 AND type = $2::export_type AND release_seq = $4 AND status = 'running'", ) .bind(event_id) .bind(export_type) .bind(file_path) .bind(seq) .execute(pool) .await .map(|r| r.rows_affected() > 0) .unwrap_or(false) } /// Parse the trailing generation number out of `` (e.g. /// `Gallery..zip`, `Gallery.zip..tmp`, `viewer_tmp__`). Returns None if the /// name doesn't fit the shape, so unrelated files are left untouched. fn parse_gen_seq(name: &str, prefix: &str, suffix: &str) -> Option { name.strip_prefix(prefix)?.strip_suffix(suffix)?.parse::().ok() } /// Best-effort removal of stale per-generation export artifacts for one export type. Deletes /// ONLY strictly-older generations (`n < keep_seq`) — never `keep_seq`'s own current file, /// and never a NEWER generation that a concurrent re-release may already be producing (that /// "delete all but mine" would let a lagging older winner nuke a newer live file). Covers the /// finished archive (`..zip`), its temp (`.zip..tmp`), and — for html — /// the `viewer_tmp__` staging dir, so crash-orphaned per-generation temps don't /// accumulate. Runs after a worker wins its finalize; older generations can never be served /// again (their `file_path` was nulled by the re-release that superseded them). Tolerates /// races and IO errors — purely disk hygiene. async fn prune_stale_export_files(exports_dir: &Path, prefix: &str, event_id: Uuid, keep_seq: i64) { let final_prefix = format!("{prefix}."); // Gallery. / Memories. let temp_prefix = format!("{prefix}.zip."); // Gallery.zip. / Memories.zip. let viewer_prefix = format!("viewer_tmp_{event_id}_"); let mut rd = match tokio::fs::read_dir(exports_dir).await { Ok(rd) => rd, Err(_) => return, }; while let Ok(Some(entry)) = rd.next_entry().await { let name = entry.file_name(); let name = name.to_string_lossy(); // Try each artifact shape; a strictly-older seq in any of them marks it for deletion. // Check the temp shape before the final shape: `Gallery.zip..tmp` also starts with // `Gallery.` but isn't a `.zip`, so its final-shape parse returns None anyway. let seq = parse_gen_seq(&name, &temp_prefix, ".tmp") .or_else(|| parse_gen_seq(&name, &final_prefix, ".zip")) .or_else(|| { if prefix == "Memories" { parse_gen_seq(&name, &viewer_prefix, "") } else { None } }); if seq.is_some_and(|n| n < keep_seq) { let path = entry.path(); let is_dir = entry.file_type().await.map(|t| t.is_dir()).unwrap_or(false); let _ = if is_dir { tokio::fs::remove_dir_all(&path).await } else { tokio::fs::remove_file(&path).await }; } } } /// Mark this worker's export `failed` — but ONLY for the generation it was running /// (`release_seq = seq` and still `running`). Without the seq guard, a superseded worker /// erroring out after a re-release would clobber the FRESH worker's `running`/`pending` row /// to `failed`, stranding the new keepsake. A superseded worker's failure is a no-op here. async fn mark_failed(pool: &PgPool, event_id: Uuid, export_type: &str, seq: i64, msg: &str) { let _ = sqlx::query( "UPDATE export_job SET status = 'failed', error_message = $3 WHERE event_id = $1 AND type = $2::export_type AND release_seq = $4 AND status = 'running'", ) .bind(event_id) .bind(export_type) .bind(msg) .bind(seq) .execute(pool) .await; } /// Update the progress bar for THIS generation only (`release_seq = seq`). Seq-guarded so a /// superseded worker still streaming can't overwrite the fresh generation's progress and make /// the bar jump backwards. Cosmetic, but keeps the displayed percent monotonic per release. async fn update_progress(pool: &PgPool, event_id: Uuid, export_type: &str, seq: i64, pct: i16) { let _ = sqlx::query( "UPDATE export_job SET progress_pct = $3 WHERE event_id = $1 AND type = $2::export_type AND release_seq = $4", ) .bind(event_id) .bind(export_type) .bind(pct) .bind(seq) .execute(pool) .await; } async fn maybe_broadcast_complete( pool: &PgPool, event_id: Uuid, sse_tx: &broadcast::Sender, ) { let row: Option<(bool, bool)> = sqlx::query_as( "SELECT export_zip_ready, export_html_ready FROM event WHERE id = $1", ) .bind(event_id) .fetch_optional(pool) .await .unwrap_or(None); if let Some((zip_ready, html_ready)) = row { if zip_ready && html_ready { let _ = sse_tx.send(SseEvent { event_type: "export-available".to_string(), data: serde_json::json!({ "types": ["zip", "html"] }).to_string(), }); } } } /// Recursively write all files from an embedded `include_dir::Dir` into a ZIP. async fn write_dir_to_zip( dir: &include_dir::Dir<'_>, zip: &mut ZipFileWriter, ) -> Result<()> { for file in dir.files() { let path = file.path().to_string_lossy().to_string(); let contents = file.contents(); let builder = ZipEntryBuilder::new(path.into(), Compression::Deflate); let mut entry = zip.write_entry_stream(builder).await?; let mut cursor = AllowStdIo::new(std::io::Cursor::new(contents)); fcopy(&mut cursor, &mut entry).await?; entry.close().await?; } for sub_dir in dir.dirs() { Box::pin(write_dir_to_zip(sub_dir, zip)).await?; } Ok(()) } fn ext_from_path(path: &str) -> &str { path.rsplit('.').next().unwrap_or("bin") } fn sanitize_name(name: &str) -> String { name.chars() .map(|c| if c.is_alphanumeric() || c == '-' { c } else { '_' }) .collect() } // ── Static content ─────────────────────────────────────────────────────────── const README_TEXT: &str = "EventSnap Offline-Galerie\n\ \n\ So geht's:\n\ 1. Entpacke diese ZIP-Datei\n\ (Windows: Rechtsklick > \"Alle extrahieren\"; Mac: Doppelklick;\n\ Handy: Dateimanager-App verwenden).\n\ 2. Öffne \"index.html\" im Browser\n\ (z. B. Chrome, Safari oder Firefox).\n\ 3. Stöbere durch alle Fotos und Videos.\n\ Du kannst zwischen Listen- und Rasteransicht wechseln,\n\ nach Hashtags filtern und nach Nutzern suchen.\n\ 4. Eine Internetverbindung ist nicht nötig.\n\ Alles ist lokal auf deinem Gerät gespeichert.\n\ \n\ Viel Freude mit den Erinnerungen!\n";