fix(export): close all 13 findings from the multi-agent review

A 6-lens multi-agent review (103 agents; every finding adversarially verified by three
refute-by-default skeptics) confirmed the epoch redesign is sound — the core invariant
holds — and found 13 real defects around it. All are fixed here.

The redesign's invariant was re-verified and stands: an accepted upload is always in the
keepsake, and a downloadable keepsake always reflects the most recent release.

But a WEAKER adjacent invariant did NOT hold — "a downloadable keepsake reflects the
current content" — and that is the theme of the biggest fixes.

CONTENT REMOVAL NOW ALWAYS REACHES THE KEEPSAKE
  Regeneration was wired only to host deletes. Ban, unban, guest self-delete and guest
  comment-delete did not trigger it — yet the export query filters `is_banned = FALSE`,
  so the pipeline already agreed that content must not be there. Ban someone for abusive
  content after release and every guest kept downloading an archive containing it.
  All five removal paths now invalidate and rebuild. `query_comments` also gained the
  banned/hidden filter it was missing entirely (the ZIP had it; the viewer did not).

  Each removal is now ATOMIC with its invalidation (one tx, models made executor-generic).
  Previously the delete committed in its own tx and the regeneration in a second: a dropped
  handler future between them left the taken-down photo permanently downloadable, and
  nothing could notice — the keepsake still looked complete and the host could no longer
  even find the upload to retry.

NO MORE TAKEDOWN STORM
  Every removal spawned two uncancellable full exports. A superseded worker was INERT, not
  STOPPED: it still ran every ffmpeg spawn and image resize and wrote a whole archive before
  discovering it had lost. Five deletes left ten workers alive, each holding a
  full-gallery-sized temp, in a 1 GB container.
  Now: a debounce before `claim_job` (a superseded worker fails its claim and does ZERO
  work, collapsing a burst into one export), plus `update_progress` — which already ran
  exactly the right liveness predicate and threw the answer away — returning it so both file
  loops bail the instant they are retired. Moderating a comment no longer rebuilds the
  multi-GB ZIP either: the ZIP is media-only, so it is carried forward, with prune taught to
  protect any file a live row still references.

AN ESCAPE HATCH
  A failed export was terminal at runtime: release_gallery refuses an already-released event,
  recovery only runs at boot, and there was no retry route — the only outs were restarting the
  container or reopening uploads to every guest. Added POST /host/export/rebuild. Also widened
  mark_failed's guard to `status IN ('running','pending')`: when claim_job itself ERRORED, the
  row was still `pending`, so the failure was never recorded and the job sat at 0% forever.

SILENT INVALIDATION
  Retiring the epoch 404s the download instantly, but nothing told the clients — guests kept
  seeing an enabled download button through the whole rebuild. Now broadcasts an invalidation.
  And the download was a top-level <a href>: a 404 (or a 429 from the per-IP export limit that
  everyone on the venue WiFi shares) NAVIGATED THE USER OUT OF THE PWA onto raw JSON. It now
  targets a hidden iframe, so an error response is harmless.

ALSO
  - The HTML export still had the exists()-then-open TOCTOU the ZIP path was fixed to remove,
    at three sites, two with a hard `?` that failed the ENTIRE viewer. It is reachable: the
    compression worker hard-deletes an original when a transcode fails and can still be running
    when the gallery is released.
  - Crash-orphaned .tmp files were immortal (prune deliberately skips .tmp). Swept at boot,
    where it is unambiguously safe.
  - export_status read the epoch in one statement and used it in the next; now one query.
  - DiskCache::snapshot IGNORED its path argument on a cache hit, returning whatever filesystem
    was measured last. Latent only because both callers pass media_path — it would have silently
    misreported the moment anyone measured the exports volume. Now keyed by path.
  - Corrected two comments that asserted safety properties the code does not have (claim_job
    does NOT prevent a retired-epoch claim — retirement is enforced at READ time; and a
    multi-replica reap is NOT contained, it can corrupt the keepsake). Added a rollback runbook
    to migration 014, the repo's first destructive migration.

TESTS
  The regression guard for the headline FOR SHARE fix was probabilistic — it could pass on
  broken code if the interleaving fell the wrong way. It is now STRUCTURAL: the upload is held
  open mid-body with its pre-flight check already passed, so the release provably lands inside
  the exact window. Verified 5/5 FAILING (201 instead of 403) with the fix reverted, 5/5
  passing with it.

Verified: backend 40 tests, svelte-check 0 errors, e2e typecheck clean,
e2e 160 passed / 1 skipped on chromium-desktop.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
fabi
2026-07-14 20:35:47 +02:00
parent 5affef47cf
commit 9666d74a46
14 changed files with 632 additions and 147 deletions

View File

@@ -1,5 +1,12 @@
-- Export generation, unified.
--
-- ⚠ ROLLBACK RUNBOOK. This is the repo's first DESTRUCTIVE migration (it DROPs two columns). The
-- previous image will NOT boot against this schema: sqlx runs migrations before serving and errors
-- with VersionMissing on an unknown version, so you get a crash loop, not a degraded service.
-- To roll back to the previous release you must run 014_export_epoch.down.sql BY HAND FIRST, then
-- deploy the old image. The down migration re-derives the old ready flags from the epoch state, so
-- no keepsake is lost.
--
-- Before this migration, "which generation of the keepsake is current?" had no single answer.
-- It was assembled at runtime from three separately-written pieces of state across two tables:
-- * event.export_released_at — a release marker with NO identity (you cannot ask *which* release)

View File

@@ -388,16 +388,21 @@ pub async fn export_status(
let released = event.export_released_at.is_some();
// Only report jobs at the CURRENT epoch. A row left behind by a retired generation (a worker
// that was superseded mid-run) is meaningless — surfacing its frozen `running`/77% would show
// the host a progress bar that never moves for a keepsake nobody is building. A retired row
// reads as "locked" (no current job), which is exactly what it is.
// ONE statement: the epoch comparison happens inside the query, against a single snapshot.
// Binding the epoch read by a previous statement would let a release/regeneration commit in
// between and yield `{released: true, zip: locked, html: locked}` — a state that never existed.
//
// Only jobs at the CURRENT epoch are reported. A row left behind by a retired generation (a
// worker superseded mid-run) is meaningless — surfacing its frozen `running`/77% would show a
// progress bar that never moves for a keepsake nobody is building. It reads as "locked" (no
// current job), which is exactly what it is.
let jobs: Vec<(String, String, i16)> = sqlx::query_as(
"SELECT type::text, status::text, progress_pct FROM export_job
WHERE event_id = $1 AND epoch = $2",
"SELECT j.type::text, j.status::text, j.progress_pct
FROM export_job j
JOIN event e ON e.id = j.event_id
WHERE e.id = $1 AND j.epoch = e.export_epoch",
)
.bind(event.id)
.bind(event.export_epoch)
.fetch_all(&state.pool)
.await?;

View File

@@ -12,6 +12,7 @@ use crate::models::event::Event;
use crate::models::session::Session;
use crate::models::upload::Upload;
use crate::models::user::UserRole;
use crate::services::export::Affects;
use crate::state::{AppState, SseEvent};
// ── DTOs ─────────────────────────────────────────────────────────────────────
@@ -148,6 +149,9 @@ pub async fn ban_user(
// (enforced live on the write handlers + Require{Host,Admin}). Revoking sessions would
// contradict that model, break the documented "banned guest can still download the
// keepsake" flow, and be ineffective anyway (the user could just /recover a new session).
//
// The ban and the keepsake invalidation are ONE transaction — see `host_delete_upload`.
let mut tx = state.pool.begin().await?;
sqlx::query(
"UPDATE \"user\"
SET is_banned = TRUE, uploads_hidden = TRUE, uploads_hidden_at = NOW()
@@ -155,9 +159,21 @@ pub async fn ban_user(
)
.bind(user_id)
.bind(auth.event_id)
.execute(&state.pool)
.execute(&mut *tx)
.await?;
// A ban hides the user's uploads EVERYWHERE — and the keepsake is the place that matters most,
// because it is the copy people keep. The export already filters `is_banned = FALSE`, so a
// FUTURE export excludes them; without this, an ALREADY-RELEASED archive would keep serving a
// banned user's photos forever. Same class as a takedown, so same treatment.
let regen =
crate::services::export::invalidate_and_arm(&mut tx, &state.config.event_slug, Affects::Both)
.await?;
tx.commit().await?;
if let Some(r) = regen {
start_regen(&state, r);
}
// Evict their content live from every feed + the diashow so it disappears without
// each viewer having to reload. (Their own SSE stream is separately dropped by the
// is_banned revalidation in `sse::stream`, so they stop receiving live pushes while
@@ -205,6 +221,7 @@ pub async fn unban_user(
// Unban restores visibility too: ban set `uploads_hidden = TRUE`, so clearing only
// `is_banned` would leave their content invisible. Clear all three (the timestamp too,
// so a future ban stamps a fresh `uploads_hidden_at` the reconnect delta will replay).
let mut tx = state.pool.begin().await?;
let result = sqlx::query(
"UPDATE \"user\"
SET is_banned = FALSE, uploads_hidden = FALSE, uploads_hidden_at = NULL
@@ -212,11 +229,22 @@ pub async fn unban_user(
)
.bind(user_id)
.bind(auth.event_id)
.execute(&state.pool)
.execute(&mut *tx)
.await?;
if result.rows_affected() == 0 {
return Err(AppError::NotFound("Benutzer nicht gefunden.".into()));
}
// The mirror of the ban case: an unban RESTORES their uploads to the export query, so an
// already-released keepsake is now missing content it should contain. Rebuild it.
let regen =
crate::services::export::invalidate_and_arm(&mut tx, &state.config.event_slug, Affects::Both)
.await?;
tx.commit().await?;
if let Some(r) = regen {
start_regen(&state, r);
}
tracing::info!(
actor_user_id = %auth.user_id,
target_user_id = %user_id,
@@ -226,6 +254,51 @@ pub async fn unban_user(
Ok(StatusCode::NO_CONTENT)
}
/// Force a keepsake rebuild. The ESCAPE HATCH.
///
/// Without this, a failed or stranded export is terminal at runtime: `release_gallery` refuses an
/// already-released event ("bereits freigegeben"), `recover_exports` only runs at boot, and there is
/// no other retry path — so the host's only options were restarting the container or reopening the
/// event (which unlocks uploads to every guest and discards the release). This is also the recovery
/// path for a keepsake that went stale for any reason we haven't thought of.
pub async fn rebuild_export(
State(state): State<AppState>,
RequireHost(auth): RequireHost,
) -> Result<StatusCode, AppError> {
let mut tx = state.pool.begin().await?;
let regen =
crate::services::export::invalidate_and_arm(&mut tx, &state.config.event_slug, Affects::Both)
.await?;
tx.commit().await?;
let Some(r) = regen else {
return Err(AppError::BadRequest(
"Die Galerie ist nicht freigegeben — es gibt nichts neu zu erzeugen.".into(),
));
};
// No debounce: this is an explicit, deliberate host action, not a burst.
for export_type in ["zip", "html"] {
let _ = state.sse_tx.send(SseEvent::new(
"export-progress",
serde_json::json!({ "type": export_type, "progress_pct": 0 }).to_string(),
));
}
crate::services::export::spawn_export_jobs(
r.event_id,
r.event_name,
r.epoch,
std::time::Duration::ZERO,
state.pool.clone(),
state.config.media_path.clone(),
state.config.export_path.clone(),
state.sse_tx.clone(),
);
tracing::info!(actor_user_id = %auth.user_id, "host: rebuild_export");
Ok(StatusCode::NO_CONTENT)
}
pub async fn set_role(
State(state): State<AppState>,
RequireHost(auth): RequireHost,
@@ -445,36 +518,33 @@ pub async fn dismiss_pin_reset_request(
/// stops being downloadable the instant the delete commits, and a fresh worker rebuilds it without
/// the removed content. The download 404s in the meantime, which is the correct answer — serving
/// the old archive would serve the deleted photo.
async fn regenerate_keepsake_if_released(state: &AppState) -> Result<(), AppError> {
let mut tx = state.pool.begin().await?;
let bumped: Option<(Uuid, String, i64)> = sqlx::query_as(
"UPDATE event SET export_epoch = export_epoch + 1
WHERE slug = $1 AND export_released_at IS NOT NULL
RETURNING id, name, export_epoch",
)
.bind(&state.config.event_slug)
.fetch_optional(&mut *tx)
.await?;
// Not released → nothing to regenerate; the export will be built fresh at release time.
let Some((event_id, event_name, epoch)) = bumped else {
return Ok(());
};
crate::services::export::enqueue_jobs_at_epoch(&mut tx, event_id, epoch).await?;
tx.commit().await?;
/// Start the workers for a regeneration that was armed inside a (now-committed) transaction, and
/// tell every client the current keepsake just became undownloadable.
///
/// The SSE matters: bumping the epoch retires the archive INSTANTLY, so `/export/zip` starts 404ing
/// the moment the change commits. Without a nudge, a guest sitting on `/export` keeps rendering an
/// enabled "download" button for the whole rebuild. Both the nav badge and the export page already
/// refetch `/export/status` on `export-progress`, so a 0% tick is the cheapest correct signal.
pub fn start_regen(state: &AppState, regen: crate::services::export::PendingRegen) {
for export_type in ["zip", "html"] {
let _ = state.sse_tx.send(SseEvent::new(
"export-progress",
serde_json::json!({ "type": export_type, "progress_pct": 0 }).to_string(),
));
}
crate::services::export::spawn_export_jobs(
event_id,
event_name,
epoch,
regen.event_id,
regen.event_name,
regen.epoch,
// Debounced: a takedown pass is a burst, and each request retires the last generation. The
// delay lets superseded workers fail their claim and do zero work instead of each building
// a full archive. See export::REGEN_DEBOUNCE.
crate::services::export::REGEN_DEBOUNCE,
state.pool.clone(),
state.config.media_path.clone(),
state.config.export_path.clone(),
state.sse_tx.clone(),
);
Ok(())
}
pub async fn host_delete_upload(
@@ -486,18 +556,26 @@ pub async fn host_delete_upload(
.await?
.ok_or_else(|| AppError::NotFound("Upload nicht gefunden.".into()))?;
let deleted = Upload::soft_delete_in_event(&state.pool, upload_id, auth.event_id).await?;
// The delete and the keepsake invalidation are ONE transaction: if the delete committed and the
// invalidation didn't, the taken-down photo would stay downloadable forever and nothing would
// notice (the keepsake still looks complete, and the host can no longer find the upload to retry).
let mut tx = state.pool.begin().await?;
let deleted = Upload::soft_delete_in_event(&mut tx, upload_id, auth.event_id).await?;
if !deleted {
return Err(AppError::NotFound("Upload nicht gefunden.".into()));
}
let regen =
crate::services::export::invalidate_and_arm(&mut tx, &state.config.event_slug, Affects::Both)
.await?;
tx.commit().await?;
let _ = state.sse_tx.send(SseEvent::new(
"upload-deleted",
serde_json::json!({ "upload_id": upload.id }).to_string(),
));
// A released keepsake must not keep serving a photo the host just took down.
regenerate_keepsake_if_released(&state).await?;
if let Some(r) = regen {
start_regen(&state, r);
}
tracing::info!(
actor_user_id = %auth.user_id,
@@ -514,17 +592,29 @@ pub async fn host_delete_comment(
RequireHost(auth): RequireHost,
Path(comment_id): Path<Uuid>,
) -> Result<StatusCode, AppError> {
let deleted =
Comment::soft_delete_in_event(&state.pool, comment_id, auth.event_id).await?;
let mut tx = state.pool.begin().await?;
let deleted = Comment::soft_delete_in_event(&mut tx, comment_id, auth.event_id).await?;
if !deleted {
return Err(AppError::NotFound("Kommentar nicht gefunden.".into()));
}
// Only the HTML viewer embeds comments — the ZIP is media-only, so it is carried forward rather
// than rebuilt. Otherwise moderating one comment would 404 the photo download for minutes to
// change nothing in it.
let regen = crate::services::export::invalidate_and_arm(
&mut tx,
&state.config.event_slug,
Affects::ViewerOnly,
)
.await?;
tx.commit().await?;
let _ = state.sse_tx.send(SseEvent::new(
"comment-deleted",
serde_json::json!({ "comment_id": comment_id }).to_string(),
));
// The HTML viewer embeds comments — rebuild it so a deleted comment doesn't live on there.
regenerate_keepsake_if_released(&state).await?;
if let Some(r) = regen {
start_regen(&state, r);
}
tracing::info!(
actor_user_id = %auth.user_id,
event_id = %auth.event_id,
@@ -647,6 +737,7 @@ pub async fn release_gallery(
event_id,
event_name,
epoch,
std::time::Duration::ZERO,
state.pool.clone(),
state.config.media_path.clone(),
state.config.export_path.clone(),

View File

@@ -220,10 +220,22 @@ pub async fn delete_comment(
// Event-scope: soft_delete_in_event only matches comments whose upload is in
// the caller's event, so a cross-event comment_id resolves to a 404 here.
let deleted = Comment::soft_delete_in_event(&state.pool, comment_id, auth.event_id).await?;
let mut tx = state.pool.begin().await?;
let deleted = Comment::soft_delete_in_event(&mut tx, comment_id, auth.event_id).await?;
if !deleted {
return Err(AppError::NotFound("Kommentar nicht gefunden.".into()));
}
// Comments live only in the HTML viewer, so the ZIP is carried forward, not rebuilt.
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);
}
let _ = state.sse_tx.send(crate::state::SseEvent::new(
"comment-deleted",
serde_json::json!({ "comment_id": comment_id, "upload_id": comment.upload_id }).to_string(),

View File

@@ -449,7 +449,20 @@ pub async fn delete_upload(
return Err(AppError::Forbidden("Nur eigene Uploads löschen.".into()));
}
Upload::soft_delete_in_event(&state.pool, upload_id, auth.event_id).await?;
// 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

View File

@@ -126,6 +126,9 @@ async fn main() -> Result<()> {
.route("/api/v1/host/event/close", post(handlers::host::close_event))
.route("/api/v1/host/event/open", post(handlers::host::open_event))
.route("/api/v1/host/gallery/release", post(handlers::host::release_gallery))
// Escape hatch: force a keepsake rebuild. Without it a failed export is terminal at runtime
// (release_gallery refuses an already-released event; recovery only runs at boot).
.route("/api/v1/host/export/rebuild", post(handlers::host::rebuild_export))
.route("/api/v1/host/users", get(handlers::host::list_users))
.route("/api/v1/host/users/{id}/ban", post(handlers::host::ban_user))
.route("/api/v1/host/users/{id}/unban", post(handlers::host::unban_user))

View File

@@ -97,8 +97,10 @@ impl Comment {
/// Event-scoped variant of [`Self::soft_delete`]. Returns `false` if the
/// comment doesn't exist or belongs to a different event.
/// Executor-generic so the delete and the keepsake regeneration can share one transaction
/// (see `Upload::soft_delete_in_event` for why that must be atomic).
pub async fn soft_delete_in_event(
pool: &PgPool,
conn: &mut sqlx::PgConnection,
id: Uuid,
event_id: Uuid,
) -> Result<bool, sqlx::Error> {
@@ -111,7 +113,7 @@ impl Comment {
)
.bind(id)
.bind(event_id)
.execute(pool)
.execute(conn)
.await?;
Ok(result.rows_affected() > 0)
}

View File

@@ -188,12 +188,17 @@ impl Upload {
/// Event-scoped variant of [`Self::soft_delete`]. Returns `false` if no row
/// matched (already deleted, wrong event, or unknown id) so host handlers
/// can return a clean 404 instead of silently no-op'ing.
/// Executor-generic so a caller can run the delete and the keepsake regeneration in ONE
/// transaction. They must be atomic: if the delete commits and the regeneration doesn't (a
/// dropped handler future, a failed second tx), the taken-down photo stays in the downloadable
/// archive forever, and recovery can't tell — the keepsake still looks complete at the current
/// epoch, and the host can no longer even find the upload to retry.
pub async fn soft_delete_in_event(
pool: &PgPool,
conn: &mut sqlx::PgConnection,
id: Uuid,
event_id: Uuid,
) -> Result<bool, sqlx::Error> {
let mut tx = pool.begin().await?;
let tx = conn;
let row: Option<(Uuid, i64)> = sqlx::query_as(
"UPDATE upload
SET deleted_at = NOW()
@@ -218,7 +223,6 @@ impl Upload {
} else {
false
};
tx.commit().await?;
Ok(deleted)
}

View File

@@ -7,7 +7,7 @@
//! barely moves. [`DiskCache`] refreshes it at most once per [`TTL`] and serves the
//! rest from memory.
use std::path::Path;
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
use std::time::{Duration, Instant};
@@ -24,7 +24,7 @@ pub struct DiskInfo {
/// `AppState`.
#[derive(Clone)]
pub struct DiskCache {
inner: Arc<RwLock<Option<(DiskInfo, Instant)>>>,
inner: Arc<RwLock<Option<(PathBuf, DiskInfo, Instant)>>>,
}
impl DiskCache {
@@ -39,14 +39,21 @@ impl DiskCache {
/// Returns `None` when the mount can't be resolved — callers MUST treat that as
/// "unknown", never "zero free". (The quota path in particular fails *open* on
/// `None`: enforcing a 0-byte limit would lock every user out of uploading.)
pub fn snapshot(&self, media_path: &Path) -> Option<DiskInfo> {
if let Some((info, at)) = *self.inner.read().unwrap() {
if at.elapsed() < TTL {
return Some(info);
/// Cached free-space reading for `path`.
///
/// The cache is keyed BY PATH. It used to hold a single slot and ignore its argument on a hit,
/// so it would happily return the media volume's numbers for any other path within the TTL.
/// That was invisible only because every caller happened to pass `media_path` — the first
/// caller to ask about a different volume (e.g. the exports volume, which is a separate mount)
/// would have silently got the wrong filesystem's free space.
pub fn snapshot(&self, path: &Path) -> Option<DiskInfo> {
if let Some((cached_path, info, at)) = self.inner.read().unwrap().as_ref() {
if cached_path == path && at.elapsed() < TTL {
return Some(*info);
}
}
let info = read_disk_for_path(media_path)?;
*self.inner.write().unwrap() = Some((info, Instant::now()));
let info = read_disk_for_path(path)?;
*self.inner.write().unwrap() = Some((path.to_path_buf(), info, Instant::now()));
Some(info)
}
}

View File

@@ -1,4 +1,5 @@
use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::{Context, Result};
use async_zip::tokio::write::ZipFileWriter;
@@ -100,7 +101,17 @@ pub async fn enqueue_jobs_at_epoch(
event_id: Uuid,
epoch: i64,
) -> Result<()> {
for export_type in ["zip", "html"] {
enqueue_types_at_epoch(conn, event_id, epoch, &["zip", "html"]).await
}
/// Arm a specific subset of the export types at `epoch` (see [`enqueue_jobs_at_epoch`]).
pub async fn enqueue_types_at_epoch(
conn: &mut sqlx::PgConnection,
event_id: Uuid,
epoch: i64,
types: &[&str],
) -> Result<()> {
for export_type in types {
sqlx::query(
"INSERT INTO export_job (event_id, type, status, progress_pct, epoch)
VALUES ($1, $2::export_type, 'pending', 0, $3)
@@ -119,6 +130,85 @@ pub async fn enqueue_jobs_at_epoch(
Ok(())
}
/// The export artifacts a content change can invalidate.
///
/// The ZIP holds only the media originals; the HTML viewer additionally embeds captions, likes and
/// comments. So a comment moderation needs only the viewer rebuilt — rebuilding the multi-GB ZIP for
/// it would 404 the photo download for minutes to change nothing in it.
#[derive(Clone, Copy, PartialEq, Eq)]
pub enum Affects {
/// Media changed (an upload removed/hidden) — both artifacts are stale.
Both,
/// Only viewer content changed (a comment) — carry the ZIP forward untouched.
ViewerOnly,
}
/// Handle for a regeneration that a caller has ARMED inside its transaction but not yet started.
/// The workers must only be spawned AFTER that transaction commits, or they could snapshot the
/// database before the change that triggered them is visible.
pub struct PendingRegen {
pub event_id: Uuid,
pub event_name: String,
pub epoch: i64,
}
/// Invalidate the current keepsake and arm a rebuild — IN THE CALLER'S TRANSACTION.
///
/// Any content removal that happens after a release (a takedown, a ban, a guest deleting their own
/// photo) must reach the keepsake: the archive is the artifact people keep forever, and "please take
/// my photo out" is the one request that most needs to be honoured there. Bumping the epoch retires
/// the current archive instantly (readiness is derived from the epoch), so the stale download stops
/// being served the moment the change commits, and a fresh worker rebuilds without the content.
///
/// Runs in the caller's tx ON PURPOSE. If the removal committed but the regeneration didn't, the
/// taken-down photo would stay downloadable forever and nothing would notice: the keepsake still
/// looks complete at the current epoch, so recovery skips it, and the host can no longer even find
/// the upload to retry. Returns `None` when the event isn't released (nothing to invalidate — the
/// export is built fresh at release time).
pub async fn invalidate_and_arm(
conn: &mut sqlx::PgConnection,
event_slug: &str,
affects: Affects,
) -> Result<Option<PendingRegen>> {
let bumped: Option<(Uuid, String, i64)> = sqlx::query_as(
"UPDATE event SET export_epoch = export_epoch + 1
WHERE slug = $1 AND export_released_at IS NOT NULL
RETURNING id, name, export_epoch",
)
.bind(event_slug)
.fetch_optional(&mut *conn)
.await?;
let Some((event_id, event_name, epoch)) = bumped else {
return Ok(None);
};
if affects == Affects::ViewerOnly {
// Carry the finished ZIP into the new epoch instead of rebuilding it: same file, same
// `file_path`, now stamped current so `export_current` keeps serving it. Only the viewer is
// re-armed below. (`prune_stale_export_files` is told to protect any file a current-epoch
// row still points at, so the carried archive isn't swept for having an older epoch in its
// name.)
sqlx::query(
"UPDATE export_job SET epoch = $2
WHERE event_id = $1 AND type = 'zip'::export_type
AND status = 'done' AND epoch = $2 - 1",
)
.bind(event_id)
.bind(epoch)
.execute(&mut *conn)
.await?;
}
let types: &[&str] = match affects {
Affects::Both => &["zip", "html"],
Affects::ViewerOnly => &["html"],
};
enqueue_types_at_epoch(&mut *conn, event_id, epoch, types).await?;
Ok(Some(PendingRegen { event_id, event_name, epoch }))
}
/// Startup export recovery: re-arm exports for any released event whose keepsake isn't fully
/// present at the current epoch (a crash mid-export, or a `done` row whose file has since gone
/// missing). Without this, `release_gallery` would reject a retry with "bereits freigegeben" and
@@ -147,6 +237,12 @@ pub async fn recover_exports(
}
};
// Crash-orphaned temps are immortal otherwise: `prune_stale_export_files` deliberately never
// touches `.tmp` (a superseded worker may still be streaming into one), and nothing else sweeps
// them. A hard kill mid-export — especially one whose epoch has since moved on — strands a
// full-gallery-sized file forever. Boot is the one moment this is unambiguously safe.
sweep_orphan_temps(&export_path).await;
for (event_id, event_name, epoch) in rows {
// Retire any `done` row at the current epoch whose file is missing or empty, so the
// upsert below re-arms it instead of leaving an undownloadable "ready" keepsake.
@@ -183,6 +279,7 @@ pub async fn recover_exports(
event_id,
event_name,
epoch,
Duration::ZERO,
pool.clone(),
media_path.clone(),
export_path.clone(),
@@ -191,6 +288,37 @@ pub async fn recover_exports(
}
}
/// Remove every export temp artifact at boot. Called from `recover_exports`, after the startup sweep
/// has marked all `running` jobs `failed` and before any worker is spawned — so no live worker can
/// own one of these. A worker that is re-armed will recreate its temp from scratch.
async fn sweep_orphan_temps(export_path: &Path) {
let mut rd = match tokio::fs::read_dir(export_path).await {
Ok(rd) => rd,
Err(_) => return,
};
let mut removed = 0u32;
while let Ok(Some(entry)) = rd.next_entry().await {
let name = entry.file_name();
let name = name.to_string_lossy();
let is_temp = name.ends_with(".tmp") || name.starts_with("viewer_tmp_");
if !is_temp {
continue;
}
let is_dir = entry.file_type().await.map(|t| t.is_dir()).unwrap_or(false);
let r = if is_dir {
tokio::fs::remove_dir_all(entry.path()).await
} else {
tokio::fs::remove_file(entry.path()).await
};
if r.is_ok() {
removed += 1;
}
}
if removed > 0 {
tracing::warn!("export recovery: swept {removed} orphaned export temp artifact(s)");
}
}
/// Reset any `done` job at the current epoch whose archive is absent or zero-length. Guards
/// against DB/disk divergence (lost volume, truncated write) that recovery would otherwise
/// mistake for a finished keepsake.
@@ -237,10 +365,23 @@ async fn invalidate_missing_files(
Ok(())
}
/// How long a regeneration waits before claiming. A takedown pass ("remove these five photos") is a
/// burst of independent requests, each of which retires the previous generation. Without a delay,
/// each one immediately spawns two full exports — and a superseded worker is INERT, not STOPPED, so
/// it still runs every ffmpeg spawn and image resize and writes an entire archive before discovering
/// it lost. Five deletes would leave ten workers alive, each holding a full-gallery-sized temp file,
/// inside a 1 GB container.
///
/// Sleeping before `claim_job` collapses the burst for free: a worker whose epoch is retired during
/// the delay fails its claim and does ZERO work. Release and boot-recovery pass ZERO — those must
/// start immediately.
pub const REGEN_DEBOUNCE: Duration = Duration::from_secs(20);
pub fn spawn_export_jobs(
event_id: Uuid,
event_name: String,
epoch: i64,
delay: Duration,
pool: PgPool,
media_path: PathBuf,
export_path: PathBuf,
@@ -256,6 +397,9 @@ pub fn spawn_export_jobs(
// The worker is BORN with its epoch — it never learns it from the DB. A worker whose epoch
// has been retired is inert by construction: every write it makes is `epoch`-guarded on the
// row it is updating, so it simply matches nothing. No cross-table guard, no race.
if !delay.is_zero() {
tokio::time::sleep(delay).await;
}
if let Err(e) = run_zip_export(event_id, epoch, &pool, &media_path, &export_path, &sse_tx).await {
tracing::error!("ZIP export failed for event {event_id} @ epoch {epoch}: {e:#}");
mark_failed(&pool, event_id, "zip", epoch, &e.to_string()).await;
@@ -264,6 +408,9 @@ pub fn spawn_export_jobs(
});
tokio::spawn(async move {
if !delay.is_zero() {
tokio::time::sleep(delay).await;
}
if let Err(e) =
run_html_export(event_id, epoch, &event_name2, &pool2, &media_path2, &export_path2, &sse_tx2)
.await
@@ -299,7 +446,22 @@ async fn run_zip_export(
let _ = tokio::fs::remove_file(export_path.join(gen_name(event_id, "Gallery", epoch, ".tmp")))
.await;
}
res
abandon_if_superseded("ZIP", event_id, epoch, res)
}
/// A worker that discovers mid-export that its generation was retired has not FAILED — it simply
/// lost. Swallow the sentinel so we don't log an error or (pointlessly) try to mark a row we no
/// longer own as failed. Its temp artifacts were already removed by the caller.
fn abandon_if_superseded(kind: &str, event_id: Uuid, epoch: i64, res: Result<()>) -> Result<()> {
match res {
Err(e) if e.downcast_ref::<Superseded>().is_some() => {
tracing::info!(
"{kind} export for event {event_id} superseded mid-export (epoch {epoch} retired); abandoned"
);
Ok(())
}
other => other,
}
}
/// Per-generation, per-EVENT artifact name. The event id matters: all events share one exports
@@ -369,7 +531,11 @@ async fn run_zip_export_inner(
entry.close().await?;
let pct = ((i + 1) as f32 / total * 100.0) as i16;
update_progress(pool, event_id, "zip", epoch, pct.min(99)).await;
// Also our liveness check: if we've been retired, stop NOW rather than grinding through
// the rest of the gallery to build an archive we would immediately delete.
if !update_progress(pool, event_id, "zip", epoch, pct.min(99)).await {
return Err(Superseded.into());
}
}
// FLUSH AND FSYNC BEFORE THE RENAME.
@@ -401,7 +567,7 @@ async fn run_zip_export_inner(
return Ok(());
}
prune_stale_export_files(&exports_dir, "Gallery", event_id, epoch).await;
prune_stale_export_files(pool, &exports_dir, "Gallery", event_id, epoch).await;
let _ = sse_tx.send(SseEvent {
event_type: "export-progress".to_string(),
@@ -448,14 +614,14 @@ async fn run_html_export(
let res =
run_html_export_inner(epoch, event_id, event_name, pool, media_path, export_path, sse_tx).await;
if res.is_err() {
// Clean up this generation's temp artifacts so a failing export can't leak them (the leak
// is what fills the disk, which is what corrupts the next archive).
// Clean up this generation's temp artifacts so a failing (or abandoned) export can't leak
// them — the leak is what fills the disk, which is what corrupts the next archive.
let _ =
tokio::fs::remove_file(export_path.join(gen_name(event_id, "Memories", epoch, ".tmp"))).await;
let _ = tokio::fs::remove_dir_all(export_path.join(format!("viewer_tmp_{event_id}_{epoch}")))
.await;
}
res
abandon_if_superseded("HTML", event_id, epoch, res)
}
async fn run_html_export_inner(
@@ -473,7 +639,7 @@ async fn run_html_export_inner(
let hashtags_per_upload = query_hashtags(pool, event_id).await?;
let total = uploads.len().max(1) as f32;
update_progress(pool, event_id, "html", epoch, 5).await;
let _ = update_progress(pool, event_id, "html", epoch, 5).await;
// Written OUTSIDE media_path: the public /media ServeDir must never reach these.
let exports_dir = export_path.to_path_buf();
@@ -492,9 +658,23 @@ async fn run_html_export_inner(
for (i, row) in uploads.iter().enumerate() {
let src = media_path.join(&row.original_path);
if !src.exists() {
// Stat ONCE, up front, and skip this upload if the source is gone. The old code probed with
// `exists()` here and then did `metadata(&src).await?` further down — a TOCTOU whose `?`
// aborted the ENTIRE keepsake if the file vanished in between. It genuinely can: the
// compression worker hard-deletes an original when its transcode fails, and it can still be
// running when the gallery is released. A missing source must degrade one entry, never the
// whole archive (which, once released, the host cannot rebuild without reopening uploads).
let src_meta = match tokio::fs::metadata(&src).await {
Ok(m) => m,
Err(e) => {
tracing::warn!(
"HTML export: skipping upload {} — cannot stat {}: {e}",
row.id,
src.display()
);
continue;
}
};
let is_video = row.mime_type.starts_with("video/");
let id_str = row.id.to_string();
@@ -561,8 +741,7 @@ async fn run_html_export_inner(
}
// Full variant: compress to temp if >5MB, otherwise stream the original
// as-is (no temp copy).
let src_meta = tokio::fs::metadata(&src).await?;
// as-is (no temp copy). `src_meta` was stat'd once at the top of the loop.
let full_source = if src_meta.len() > 5_000_000 {
let src_clone = src.clone();
let full_path = media_tmp.join(&full);
@@ -642,7 +821,9 @@ async fn run_html_export_inner(
});
let pct = 10 + ((i + 1) as f32 / total * 60.0) as i16;
update_progress(pool, event_id, "html", epoch, pct.min(69)).await;
if !update_progress(pool, event_id, "html", epoch, pct.min(69)).await {
return Err(Superseded.into());
}
}
// 4. Build data.json
@@ -656,7 +837,7 @@ async fn run_html_export_inner(
let data_json =
serde_json::to_string_pretty(&viewer_data).context("failed to serialize data.json")?;
update_progress(pool, event_id, "html", epoch, 72).await;
let _ = update_progress(pool, event_id, "html", epoch, 72).await;
// 5. Create ZIP (per-generation paths — see run_zip_export)
let tmp_path = exports_dir.join(gen_name(event_id, "Memories", epoch, ".tmp"));
@@ -670,7 +851,7 @@ async fn run_html_export_inner(
// Write embedded viewer assets (index.html, _app/*, etc.)
write_dir_to_zip(&VIEWER_DIR, &mut zip).await?;
update_progress(pool, event_id, "html", epoch, 75).await;
let _ = update_progress(pool, event_id, "html", epoch, 75).await;
// Write data.json
{
@@ -690,7 +871,7 @@ async fn run_html_export_inner(
entry.close().await?;
}
update_progress(pool, event_id, "html", epoch, 78).await;
let _ = update_progress(pool, event_id, "html", epoch, 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
@@ -701,20 +882,29 @@ async fn run_html_export_inner(
for (name, source) in &media_manifest {
let path = source.path();
if !path.exists() {
// Open-first: a source that disappeared between the manifest being built and now (the
// compression worker deletes originals on transcode failure) must skip this entry, not
// fail the whole viewer. Opening collapses the check and the use into one operation.
let src_file = match tokio::fs::File::open(path).await {
Ok(f) => f,
Err(e) => {
tracing::warn!("HTML export: skipping media {name} — cannot read {}: {e}", path.display());
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();
let mut f = src_file.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", epoch, pct.min(98)).await;
if !update_progress(pool, event_id, "html", epoch, pct.min(98)).await {
return Err(Superseded.into());
}
}
// Flush + fsync before the rename — see run_zip_export for why dropping the file here
@@ -739,7 +929,7 @@ async fn run_html_export_inner(
return Ok(());
}
prune_stale_export_files(&exports_dir, "Memories", event_id, epoch).await;
prune_stale_export_files(pool, &exports_dir, "Memories", event_id, epoch).await;
let _ = sse_tx.send(SseEvent {
event_type: "export-progress".to_string(),
@@ -778,6 +968,7 @@ async fn query_comments(pool: &PgPool, event_id: Uuid) -> Result<Vec<ExportComme
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
AND usr.uploads_hidden = FALSE AND usr.is_banned = FALSE
ORDER BY c.created_at ASC",
)
.bind(event_id)
@@ -808,8 +999,16 @@ async fn query_hashtags(pool: &PgPool, event_id: Uuid) -> Result<Vec<(Uuid, Stri
/// (`EXISTS (SELECT 1 FROM event WHERE ... export_released_at IS NOT NULL)`) was therefore unsound:
/// a claim blocking behind a concurrent reopen could see the pre-reopen event snapshot, pass the
/// check, and return the post-bump seq — handing the worker a LIVE generation on an event that was
/// already reopened. A same-row `epoch = $3` predicate is re-evaluated correctly by EPQ, so a
/// retired epoch can never be claimed.
/// already reopened. A same-row `epoch = $3` predicate is re-evaluated correctly by EPQ.
///
/// NOTE what this does and does NOT guarantee. It compares our birth epoch to the JOB ROW's epoch —
/// not to `event.export_epoch`. `open_event` bumps the event's epoch and writes nothing to
/// `export_job` (that is the point of the design: one write retires everything). So after a reopen
/// the row is still `pending` at our epoch and this claim SUCCEEDS: the worker will build an archive
/// nobody can ever see, because retirement is enforced at READ time (`export_current` requires
/// `j.epoch = e.export_epoch`), not at write time. That is wasted work, not incorrectness — and the
/// `update_progress` liveness check bails such a worker out early. Do not "optimise" this into a
/// cross-table check: that is exactly the unsound guard we removed.
///
/// Errors are distinguished from a lost claim: silently treating a pool timeout as "someone else
/// owns it" left the row `pending` at 0% with no live worker and no error — a spinner forever.
@@ -874,7 +1073,28 @@ fn parse_gen_seq(name: &str, prefix: &str, suffix: &str) -> Option<i64> {
/// 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) {
async fn prune_stale_export_files(
pool: &PgPool,
exports_dir: &Path,
prefix: &str,
event_id: Uuid,
keep_seq: i64,
) {
// Files still referenced by a live (current-epoch) job row are OFF LIMITS regardless of the
// epoch in their name. A ViewerOnly regeneration carries the finished ZIP forward by re-stamping
// its row to the new epoch WITHOUT renaming the file — so `Gallery.<event>.<older>.zip` is still
// the served archive, and deleting it by filename-epoch would 404 the download.
let protected: Vec<String> = sqlx::query_scalar::<_, String>(
"SELECT split_part(j.file_path, '/', -1) FROM export_job j
JOIN event e ON e.id = j.event_id
WHERE j.event_id = $1 AND j.epoch = e.export_epoch
AND j.status = 'done' AND j.file_path IS NOT NULL",
)
.bind(event_id)
.fetch_all(pool)
.await
.unwrap_or_default();
// EVERY shape is event-scoped. All events share one exports volume, so a name keyed only by
// generation would let event A's prune delete event B's live keepsake (and let two events
// collide on the same archive path). `viewer_tmp_` was already scoped; the archives were not.
@@ -899,6 +1119,9 @@ async fn prune_stale_export_files(exports_dir: &Path, prefix: &str, event_id: Uu
None
}
});
if protected.iter().any(|p| p == name.as_ref()) {
continue;
}
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);
@@ -911,14 +1134,19 @@ async fn prune_stale_export_files(exports_dir: &Path, prefix: &str, event_id: Uu
}
}
/// Mark this worker's export `failed` — ONLY for the generation it was running (`epoch` matches and
/// it still holds it `running`). A superseded worker's failure is a no-op, so it can't clobber the
/// fresh generation's row and strand the new keepsake.
/// Mark this worker's export `failed` — ONLY for the generation it owns (`epoch` matches). A
/// superseded worker's failure is a no-op, so it can't clobber the fresh generation's row.
///
/// The status guard admits `pending` as well as `running`: if `claim_job` itself ERRORS (pool
/// timeout, connection reset) the row is still `pending`, and a guard of `status = 'running'` would
/// match nothing — leaving the job sitting at 0% with no worker and no error message, a spinner
/// forever. The epoch guard is what keeps this safe; the status guard is only there to avoid
/// stomping a `done` row.
async fn mark_failed(pool: &PgPool, event_id: Uuid, export_type: &str, epoch: 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 epoch = $4 AND status = 'running'",
AND epoch = $4 AND status IN ('running', 'pending')",
)
.bind(event_id)
.bind(export_type)
@@ -928,11 +1156,23 @@ async fn mark_failed(pool: &PgPool, event_id: Uuid, export_type: &str, epoch: i6
.await;
}
/// Update the progress bar for THIS generation only. Epoch-guarded (and status-guarded) so a
/// superseded worker still streaming can't overwrite the fresh generation's progress or scribble on
/// a row the startup sweep already marked failed.
async fn update_progress(pool: &PgPool, event_id: Uuid, export_type: &str, epoch: i64, pct: i16) {
let _ = sqlx::query(
/// Update the progress bar for THIS generation only, and report whether we're STILL THE LIVE
/// GENERATION.
///
/// The WHERE clause (`epoch = ours AND status = 'running'`) is exactly the liveness predicate, so
/// this write already tells us for free whether we've been superseded — it just used to throw the
/// answer away. Returning it lets the file loops bail the instant a reopen/re-release/takedown
/// retires us, instead of grinding through every remaining image and writing a whole archive we
/// will then delete. `false` = we lost (or the row is gone); a DB error is reported as still-live,
/// because a transient blip must not abandon a legitimate export.
async fn update_progress(
pool: &PgPool,
event_id: Uuid,
export_type: &str,
epoch: i64,
pct: i16,
) -> bool {
match sqlx::query(
"UPDATE export_job SET progress_pct = $3
WHERE event_id = $1 AND type = $2::export_type
AND epoch = $4 AND status = 'running'",
@@ -942,9 +1182,31 @@ async fn update_progress(pool: &PgPool, event_id: Uuid, export_type: &str, epoch
.bind(pct)
.bind(epoch)
.execute(pool)
.await;
.await
{
Ok(r) => r.rows_affected() > 0,
// Don't abandon a good export over a transient pool hiccup — the epoch-guarded finalize is
// still the authority, so at worst we do some extra work.
Err(e) => {
tracing::warn!("progress update failed for {export_type} @ epoch {epoch}: {e:#}");
true
}
}
}
/// Sentinel error for "our generation was retired mid-export". The caller discards its output and
/// returns Ok — this is a normal, expected outcome, not a failure worth marking on the job row.
#[derive(Debug)]
struct Superseded;
impl std::fmt::Display for Superseded {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "export generation superseded")
}
}
impl std::error::Error for Superseded {}
/// Broadcast `export-available` once BOTH halves are downloadable. Reads the derived predicate
/// (a `done` row at the event's current epoch, on a released event) rather than a stored flag, so
/// it can never advertise a keepsake that a reopen has already invalidated.

View File

@@ -54,11 +54,13 @@ pub async fn startup_recovery(pool: &PgPool) {
// this process's. That is correct for the supported deployment (one app container; see
// docker-compose.yml), where no other process can own a job at boot. If EventSnap is ever run
// with more than one replica, a booting replica would mark a peer's live workers `failed`.
// The epoch model CONTAINS that (the victim's epoch-guarded finalize simply no-ops, and
// recovery re-arms the job at the SAME epoch, so no keepsake is corrupted — the work is just
// redone), but it would waste an export. Fixing it properly means an owner/lease/heartbeat on
// export_job; that machinery is not worth it for a single-container app, so this is a
// documented constraint rather than a hidden one.
// This is NOT merely wasteful — do not assume the epoch model contains it. Recovery re-arms the
// job at the SAME epoch, so the reaped-but-still-alive worker and the fresh one would BOTH hold
// that epoch; the guards carry no owner/lease token, so they share the same artifact paths and
// either can win the other's `finalize_job`. That can corrupt or delete the keepsake. Running
// more than one replica therefore requires an owner/lease/heartbeat on export_job first. That
// machinery is not worth building for a single-container app — so this is a documented
// CONSTRAINT, not a hidden assumption: do not scale this service horizontally as-is.
match sqlx::query(
"UPDATE export_job
SET status = 'failed',

View File

@@ -46,3 +46,67 @@ export const JPEG_MAGIC = new Uint8Array([0xff, 0xd8, 0xff, 0xe0, 0x00, 0x10, 0x
export const PNG_MAGIC = new Uint8Array([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
/** ELF header — the magic bytes of a Linux executable. `infer` reports `application/x-executable`. */
export const ELF_MAGIC = new Uint8Array([0x7f, 0x45, 0x4c, 0x46]);
/**
* Upload whose request body PAUSES mid-stream until you release it.
*
* This exists to make the release-vs-upload race deterministic instead of hoped-for. The bug it
* guards is a TOCTOU: `upload` checks `uploads_locked_at` BEFORE reading the body (minutes, for a
* big video) and commits the row afterwards — so a release landing in between used to let the photo
* commit AFTER the export snapshot, leaving it in the live feed but permanently absent from the
* keepsake.
*
* Racing N normal uploads against a release can't prove anything: if they all happen to commit
* first (or all get rejected), the assertion passes on broken code too. Here the server is *stuck*
* mid-body with its pre-flight check already passed, so the release provably lands inside the
* window. Await `checkPassed`, fire the release, then `finish()`.
*/
export function uploadPausedMidStream(
token: string,
head: Buffer,
tail: Buffer,
opts: UploadOptions = {}
): { checkPassed: Promise<void>; finish: () => void; response: Promise<Response> } {
const boundary = '----eventsnapPausedBoundary' + Math.random().toString(16).slice(2);
const filename = opts.filename ?? 'paused.jpg';
const contentType = opts.contentType ?? 'image/jpeg';
const preamble = Buffer.from(
`--${boundary}\r\n` +
`Content-Disposition: form-data; name="file"; filename="${filename}"\r\n` +
`Content-Type: ${contentType}\r\n\r\n`
);
const epilogue = Buffer.from(`\r\n--${boundary}--\r\n`);
let releaseBody: () => void;
const gate = new Promise<void>((r) => (releaseBody = r));
let markCheckPassed: () => void;
const checkPassed = new Promise<void>((r) => (markCheckPassed = r));
const body = new ReadableStream<Uint8Array>({
async pull(controller) {
controller.enqueue(new Uint8Array(preamble));
controller.enqueue(new Uint8Array(head));
// The server has now read the part headers and the first bytes — which means it is PAST its
// pre-flight lock check and is sitting in the body loop. This is the window.
markCheckPassed();
await gate;
controller.enqueue(new Uint8Array(tail));
controller.enqueue(new Uint8Array(epilogue));
controller.close();
},
});
const response = fetch(`${BASE}/api/v1/upload`, {
method: 'POST',
headers: {
Authorization: `Bearer ${token}`,
'Content-Type': `multipart/form-data; boundary=${boundary}`,
},
body,
// Node's fetch requires this for a streaming request body.
duplex: 'half',
} as RequestInit & { duplex: 'half' });
return { checkPassed, finish: () => releaseBody(), response };
}

View File

@@ -30,8 +30,9 @@
* generation and regenerates, rather than leaving the stale worker to finalize.
* - supersession (reopen): a reopen ALONE bumps the event epoch, so a worker holding the
* pre-reopen epoch is retired and can publish nothing.
* - acceptance: every upload the server ACCEPTED (201) is in the keepsake — the upload-path
* TOCTOU that was the real cause of the "stale keepsake" bug.
* - acceptance: an upload held open mid-body across a release is REJECTED, not silently dropped
* from the keepsake — the upload-path TOCTOU that was the real cause of the "stale keepsake" bug.
* Deterministic: the body is paused with the pre-flight check already passed.
* - takedown: deleting a photo AFTER release regenerates the keepsake without it.
*/
import { test, expect } from '../../fixtures/test';
@@ -41,7 +42,7 @@ import { writeFileSync, mkdtempSync, readFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { seedUpload } from '../../helpers/seed';
import { uploadRaw } from '../../helpers/upload-client';
import { uploadRaw, uploadPausedMidStream } from '../../helpers/upload-client';
import { SseListener } from '../../helpers/sse-listener';
const BASE = process.env.E2E_FRONTEND_URL ?? 'http://localhost:3101';
@@ -342,51 +343,51 @@ test.describe('Flow re-review — reopen→re-release export integrity + stale-k
expect((await downloadZipEntries(host.jwt)).length).toBe(1);
});
test('EVERY upload the server accepted is in the keepsake (upload-path release TOCTOU)', async ({
test('an upload STUCK MID-BODY when the release lands is rejected, not silently lost', async ({
host,
}) => {
// THE contract, and the bug that survived three rounds of fixes to the export state machine
// because it was never in the state machine:
// THE bug that survived three rounds of fixes to the export state machine — because it was never
// in the state machine. `upload` checked `uploads_locked_at` BEFORE streaming the body (minutes,
// for a 500 MB video) and committed the row without re-checking. A release landing mid-upload let
// the row commit AFTER the export snapshot: the photo showed up in the live feed and was
// PERMANENTLY missing from the keepsake, with nothing to regenerate it.
//
// `upload` checked `uploads_locked_at` BEFORE streaming the body (minutes, for a big video)
// and then committed the row WITHOUT re-checking. So a release landing mid-upload let the row
// commit AFTER the export snapshot: the photo appeared in the live feed and was PERMANENTLY
// missing from the keepsake. Nothing ever regenerated it.
//
// The invariant is absolute and does not depend on winning any race: if the server said 201,
// that photo is in the keepsake. If it said 403, it isn't. Race uploads against a release and
// hold the server to exactly that.
const bytes = readFileSync(join(process.cwd(), 'fixtures', 'media', 'sample.jpg'));
// One upload that is definitely committed before the race, so the keepsake is never trivially
// empty (the release can otherwise win every race and the assertion would prove nothing).
// This test is DETERMINISTIC, not a hoped-for race. The upload is held open mid-body with its
// pre-flight check already passed, so the release provably lands inside the exact window. The
// server must then reject it at commit time (`FOR SHARE` re-check) — anything else means a photo
// the server accepted is not in the archive.
await seedUpload(host.jwt, { caption: 'pre-race' });
const inflight = Array.from({ length: 10 }, (_, i) =>
uploadRaw(host.jwt, bytes, { filename: `race${i}.jpg`, contentType: 'image/jpeg' })
);
// Fire the release while those are in flight.
const releasing = post('/api/v1/host/gallery/release', host.jwt);
const img = readFileSync(join(process.cwd(), 'fixtures', 'media', 'sample.jpg'));
const split = Math.floor(img.length / 2);
const paused = uploadPausedMidStream(host.jwt, img.subarray(0, split), img.subarray(split), {
filename: 'stuck.jpg',
contentType: 'image/jpeg',
});
const results = await Promise.all(inflight);
await releasing;
// Wait until the server is genuinely mid-body (past its pre-flight lock check).
await paused.checkPassed;
const accepted = results.filter((r) => r.status >= 200 && r.status < 300).length;
const rejected = results.filter((r) => r.status === 403).length;
expect(accepted + rejected, 'every upload either succeeded or was locked out').toBe(
results.length
);
// Land the release squarely inside that window.
expect((await post('/api/v1/host/gallery/release', host.jwt)).status).toBe(204);
// Now let the body finish. The pre-flight check passed, so ONLY the commit-time re-check can
// save this photo from becoming a phantom.
paused.finish();
const res = await paused.response;
expect(
res.status,
'an upload whose body completed after the release MUST be rejected (403 uploads_locked). ' +
'A 201 here means the server accepted a photo it will never put in the keepsake — silent, ' +
'permanent data loss.'
).toBe(403);
expect((await res.json()).error).toBe('uploads_locked');
// And the keepsake holds exactly the one upload that was genuinely committed before the release.
await waitExportDone(host.jwt);
const listing = await downloadZipEntries(host.jwt);
// +1 for the pre-race upload.
expect(
listing.length,
`server accepted ${accepted} racing upload(s) (+1 pre-race) and locked out ${rejected}, but ` +
`the keepsake holds ${listing.length}. An accepted photo missing from the keepsake is ` +
`silent, permanent data loss — exactly the bug this guards.`
).toBe(accepted + 1);
expect(listing.length, `zip entries:\n${listing.join('\n')}`).toBe(1);
});
test('deleting a photo AFTER release regenerates the keepsake without it (takedown)', async ({

View File

@@ -81,21 +81,28 @@
}
}
// Stream the (potentially multi-GB) archive straight to disk via a top-level
// download navigation instead of fetch()+blob() — buffering the whole ZIP in
// memory crashes mobile Safari/Chrome. The navigation can't send the Bearer
// header, so we first exchange the token for a single-use ticket and pass it
// as a query param; the server responds with Content-Disposition: attachment,
// so the browser downloads it and we stay on the page.
// Stream the (potentially multi-GB) archive straight to disk instead of fetch()+blob() —
// buffering the whole ZIP in memory crashes mobile Safari/Chrome. The request can't carry the
// Bearer header, so we exchange the token for a single-use ticket and pass it as a query param.
//
// The navigation targets a HIDDEN IFRAME, not the top-level page. On success the server sends
// `Content-Disposition: attachment` and the browser downloads it either way — but on an ERROR
// (404 because a takedown/ban just retired this generation and it's rebuilding, or 429 from the
// per-IP export rate limit, which everyone on the venue WiFi shares) there is no attachment
// header, so a top-level navigation would UNLOAD THE APP and dump the raw `{"error":…}` JSON on
// the guest. Inside an iframe the error is swallowed harmlessly and the PWA stays put.
//
// We deliberately don't probe with HEAD first: the ticket is single-use and axum serves HEAD from
// the GET route, so a probe would consume the ticket and burn a rate-limit slot.
let downloadFrame: HTMLIFrameElement;
async function downloadFile(endpoint: string) {
try {
const { ticket } = await api.post<{ ticket: string }>('/export/ticket');
const a = document.createElement('a');
a.href = `${endpoint}?ticket=${encodeURIComponent(ticket)}`;
a.rel = 'noopener';
document.body.appendChild(a);
a.click();
a.remove();
downloadFrame.src = `${endpoint}?ticket=${encodeURIComponent(ticket)}`;
// If the archive was retired between the page render and this click, the download
// silently does nothing — so re-sync the status to flip the UI to "generating".
setTimeout(() => void loadStatus(), 1500);
} catch (e) {
toastError(e);
}
@@ -120,6 +127,11 @@
}
</script>
<!-- Download sink. Archive downloads are navigations (they can be multi-GB, so fetch()+blob() would
OOM mobile browsers), and a navigation to an error response would unload the whole PWA. Pointing
them at this hidden iframe keeps a 404/429 harmless. -->
<iframe bind:this={downloadFrame} title="Download" class="hidden" aria-hidden="true"></iframe>
<!-- HTML guide modal -->
{#if showHtmlGuide}
<div class="fixed inset-0 z-50 flex items-center justify-center bg-black/50 p-4">