Merge branch 'fix/deploy-unattended-blockers' into main

Two independent lines of production hardening diverged at 7d0334b and attacked
overlapping problems. Neither was a superset, so this is a merge of substance
rather than a fast-forward: every conflict was resolved on the merits, and the
losing side's intent was re-checked against the winner rather than assumed.

MIGRATIONS. The branch's 021/022/023 collided with main's already-DEPLOYED
021_hashtag_counts_respect_bans and 022_client_upload_idempotency. Renumbered to
023/024/025 in a prior commit — main's versions are applied in production, so
their version numbers are immutable and the branch's had to move. Verified by
running the full sqlx::test suite, which applies the whole chain from scratch.

RESOLVED IN MAIN'S FAVOUR (the branch would have regressed these):
  * upload-queue.ts wholesale — the branch's copy has ZERO client_upload_id
    references, so taking it would have silently destroyed end-to-end upload
    idempotency, the one thing standing between a lost response and a duplicate
    photo charged twice against the guest's quota.
  * maintenance.rs supervisor — the branch replaced it with a bare tokio::spawn,
    where one panic silently stops session pruning, media reclaim, the temp
    sweep and both HashMap prunes, permanently and with no log line.
  * The decode-budget probe on spawn_blocking, not inline on the async runtime.
  * feed/+page.svelte's 8s debounce + jitter + max-wait + hidden-tab deferral,
    against the branch's naive 800ms — at 100 guests the branch's version walks
    straight into the per-user feed rate limit.
  * db.rs pool tuning, /uploaders, and the docker-compose deployment story.
  * ONE /health, still DB-backed. The branch's split (dependency-free liveness +
    DB-backed readiness) is defensible, but a constant-"ok" /health is the exact
    defect faea555 fixed and verified live, its motive (Caddy's boot gate) is
    already covered by app depends_on db: service_healthy, and the two handlers
    were the same SELECT 1 under two names.

TAKEN FROM THE BRANCH:
  * The large-PNG OOM guard and its bounded-retry counter (023). Together these
    turn a single upload that can OOM-kill a 1G container into a bounded failure
    instead of an infinite restart loop under `restart: unless-stopped`.
  * 024_feed_scalar_counts — the feed no longer aggregates the whole event per
    page. Pure SQL; column names, order and types are unchanged by design.
  * The admin-lockout fix: look the admin up BY ROLE, never by name. 025 also
    frees any guest already squatting on a reserved name.
  * PIN lockout tier ordering, bounded caption/hashtag reads, SSE ticket caps,
    PoolTimedOut -> 503 + Retry-After, and the ffmpeg stderr drain.
  * backfill_video_posters, which main lacked entirely.
  * TempFileGuard, plus sweep_orphan_originals wired into main's SUPERVISED loop
    (not the branch's bare one) — it reclaims final-named originals whose commit
    never happened, a class main's .tmp-only sweep structurally cannot see.
  * shouldAbortForStall, hand-ported into main's upload-queue.ts since that file
    was resolved to main. Widens the watchdog at loadend instead of disarming it,
    bounding a half-open socket at 2 minutes rather than handing the window to
    xhr.timeout (5-60 min) with the whole queue's `processing` latch held.

ALSO: RUST_LOG and EXPORT_PATH pinned in compose. The code fallback was
`debug` (a line per request, all night) and EXPORT_PATH was the one path with a
mount-shaped default that nothing validated.

Verified: cargo check --all-targets, cargo clippy (clean), 144/144 backend tests
against a live Postgres including upload_idempotency and upload_concurrency,
51/51 vitest, svelte-check 0 errors, eslint clean, vite build.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
MechaCat02
2026-08-08 21:16:31 +02:00
30 changed files with 2192 additions and 439 deletions

192
backend/Cargo.lock generated
View File

@@ -65,56 +65,6 @@ dependencies = [
"libc",
]
[[package]]
name = "anstream"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d"
dependencies = [
"anstyle",
"anstyle-parse",
"anstyle-query",
"anstyle-wincon",
"colorchoice",
"is_terminal_polyfill",
"utf8parse",
]
[[package]]
name = "anstyle"
version = "1.0.14"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000"
[[package]]
name = "anstyle-parse"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e"
dependencies = [
"utf8parse",
]
[[package]]
name = "anstyle-query"
version = "1.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc"
dependencies = [
"windows-sys 0.61.2",
]
[[package]]
name = "anstyle-wincon"
version = "3.0.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d"
dependencies = [
"anstyle",
"once_cell_polyfill",
"windows-sys 0.61.2",
]
[[package]]
name = "anyhow"
version = "1.0.102"
@@ -554,46 +504,12 @@ dependencies = [
"inout",
]
[[package]]
name = "clap"
version = "4.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b193af5b67834b676abd72466a96c1024e6a6ad978a1f484bd90b85c94041351"
dependencies = [
"clap_builder",
]
[[package]]
name = "clap_builder"
version = "4.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "714a53001bf66416adb0e2ef5ac857140e7dc3a0c48fb28b2f10762fc4b5069f"
dependencies = [
"anstream",
"anstyle",
"clap_lex",
"strsim",
"terminal_size",
]
[[package]]
name = "clap_lex"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9"
[[package]]
name = "color_quant"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3d7b894f5411737b7867f4827955924d7c254fc9f4d91a6aad6b097804b1018b"
[[package]]
name = "colorchoice"
version = "1.0.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570"
[[package]]
name = "compression-codecs"
version = "0.4.37"
@@ -677,15 +593,6 @@ dependencies = [
"cfg-if",
]
[[package]]
name = "crossbeam-channel"
version = "0.5.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "82b8f8f868b36967f9606790d1903570de9ceaf870a7bf9fbbd3016d636a2cb2"
dependencies = [
"crossbeam-utils",
]
[[package]]
name = "crossbeam-deque"
version = "0.8.6"
@@ -816,27 +723,6 @@ dependencies = [
"cfg-if",
]
[[package]]
name = "env_filter"
version = "1.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32e90c2accc4b07a8456ea0debdc2e7587bdd890680d71173a15d4ae604f6eef"
dependencies = [
"log",
]
[[package]]
name = "env_logger"
version = "0.11.10"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0621c04f2196ac3f488dd583365b9c09be011a4ab8b9f37248ffcc8f6198b56a"
dependencies = [
"anstream",
"anstyle",
"env_filter",
"log",
]
[[package]]
name = "equator"
version = "0.4.2"
@@ -1229,12 +1115,6 @@ dependencies = [
"weezl",
]
[[package]]
name = "glob"
version = "0.3.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0cc23270f6e1808e30a928bdc84dea0b9b4136a8bc82338574f23baf47bbd280"
[[package]]
name = "governor"
version = "0.6.3"
@@ -1622,7 +1502,6 @@ checksum = "7714e70437a7dc3ac8eb7e6f8df75fd8eb422675fc7678aff7364301092b1017"
dependencies = [
"equivalent",
"hashbrown 0.16.1",
"rayon",
"serde",
"serde_core",
]
@@ -1656,12 +1535,6 @@ dependencies = [
"syn",
]
[[package]]
name = "is_terminal_polyfill"
version = "1.70.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695"
[[package]]
name = "itertools"
version = "0.14.0"
@@ -1795,12 +1668,6 @@ dependencies = [
"vcpkg",
]
[[package]]
name = "linux-raw-sys"
version = "0.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53"
[[package]]
name = "litemap"
version = "0.8.1"
@@ -2089,12 +1956,6 @@ version = "1.21.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50"
[[package]]
name = "once_cell_polyfill"
version = "1.70.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe"
[[package]]
name = "oxipng"
version = "9.1.5"
@@ -2102,18 +1963,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "26c613f0f566526a647c7473f6a8556dbce22c91b13485ee4b4ec7ab648e4973"
dependencies = [
"bitvec",
"clap",
"crossbeam-channel",
"env_logger",
"filetime",
"glob",
"indexmap",
"libdeflater",
"log",
"rayon",
"rgb",
"rustc-hash",
"zopfli",
]
[[package]]
@@ -2607,19 +2462,6 @@ version = "2.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "94300abf3f1ae2e2b8ffb7b58043de3d399c73fa6f4b73826402a5c457614dbe"
[[package]]
name = "rustix"
version = "1.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190"
dependencies = [
"bitflags",
"errno",
"libc",
"linux-raw-sys",
"windows-sys 0.61.2",
]
[[package]]
name = "rustversion"
version = "1.0.22"
@@ -3060,12 +2902,6 @@ dependencies = [
"unicode-properties",
]
[[package]]
name = "strsim"
version = "0.11.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f"
[[package]]
name = "subtle"
version = "2.6.1"
@@ -3120,16 +2956,6 @@ version = "1.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "55937e1799185b12863d447f42597ed69d9928686b8d88a1df17376a097d8369"
[[package]]
name = "terminal_size"
version = "0.4.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "230a1b821ccbd75b185820a1f1ff7b14d21da1e442e22c0863ea5f08771a8874"
dependencies = [
"rustix",
"windows-sys 0.61.2",
]
[[package]]
name = "thiserror"
version = "1.0.69"
@@ -3505,12 +3331,6 @@ version = "1.0.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be"
[[package]]
name = "utf8parse"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821"
[[package]]
name = "uuid"
version = "1.23.0"
@@ -4187,18 +4007,6 @@ version = "1.0.21"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b8848ee67ecc8aedbaf3e4122217aff892639231befc6a1b58d29fff4c2cabaa"
[[package]]
name = "zopfli"
version = "0.8.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f05cd8797d63865425ff89b5c4a48804f35ba0ce8d125800027ad6017d2b5249"
dependencies = [
"bumpalo",
"crc32fast",
"log",
"simd-adler32",
]
[[package]]
name = "zstd"
version = "0.13.3"

View File

@@ -27,7 +27,17 @@ tracing-subscriber = { version = "0.3", features = ["env-filter"] }
dotenvy = "0.15"
sysinfo = "0.32"
image = "0.25"
oxipng = "9"
# default-features = false drops "parallel", which is what actually bounds oxipng's memory:
# with rayon it evaluates row filters concurrently, each trial holding its own full-size
# buffer, and there is no Options knob to cap that. Without the feature, lib.rs swaps in a
# sequential shim (oxipng's own supported path) so peak scales with ONE trial, not N.
# PNG optimisation gets slower; it is a background, best-effort, lossless size saving.
#
# "filetime" must be KEPT: without it OutFile::Path { preserve_attrs: true } silently no-ops.
# Dropping "binary" also removes clap/glob/env_logger — a CLI's dependencies that were being
# compiled into a server image — and "zopfli", which preset 2 does not use (it selects
# Deflaters::Libdeflater, which is not feature-gated).
oxipng = { version = "9", default-features = false, features = ["filetime"] }
async_zip = { version = "0.0.17", features = ["tokio", "deflate"] }
include_dir = "0.7"
infer = "0.15"

View File

@@ -0,0 +1,3 @@
DROP INDEX IF EXISTS idx_upload_derivative_backfill;
ALTER TABLE upload DROP COLUMN IF EXISTS derivative_last_error;
ALTER TABLE upload DROP COLUMN IF EXISTS derivative_attempts;

View File

@@ -0,0 +1,23 @@
-- Bound how many times a permanently-failing upload can be re-processed.
--
-- Without this, one poisoned row is an outage. The upload row is committed BEFORE compression
-- starts, `derivatives_rev` defaults to 0, and `set_derivatives_rev` only runs on success — so
-- a row whose processing kills the container survives at rev 0, the unconditional startup
-- backfill re-selects it on the next boot, and `restart: unless-stopped` turns that into an
-- infinite kill loop. Every restart also drops every SSE stream and truncates every in-flight
-- upload. That was reachable via a single large PNG (see services/compression.rs), but the
-- shape is general: any input that can kill or hang the worker repeats forever.
--
-- The counter is incremented WRITE-AHEAD, before the work is attempted, because the failure
-- mode being defended against is a SIGKILL — no error is returned, no handler runs, no Drop
-- fires. A counter bumped in an error path increments zero times per crash and changes nothing.
ALTER TABLE upload ADD COLUMN IF NOT EXISTS derivative_attempts SMALLINT NOT NULL DEFAULT 0;
-- Last failure text, so a row that has given up can be diagnosed without reproducing it.
-- Nothing reads this in code; it exists for the operator.
ALTER TABLE upload ADD COLUMN IF NOT EXISTS derivative_last_error TEXT;
-- Serves the backfill selection, which now filters on both columns.
CREATE INDEX IF NOT EXISTS idx_upload_derivative_backfill
ON upload (derivatives_rev, derivative_attempts)
WHERE deleted_at IS NULL;

View File

@@ -0,0 +1,26 @@
-- Restore the 016 definition verbatim.
DROP VIEW IF EXISTS v_feed;
CREATE VIEW v_feed AS
SELECT
u.id,
u.event_id,
u.user_id,
usr.display_name AS uploader_name,
usr.is_banned,
usr.uploads_hidden,
u.preview_path,
u.thumbnail_path,
u.display_path,
u.mime_type,
u.caption,
u.created_at,
COUNT(DISTINCT l.user_id) AS like_count,
COUNT(DISTINCT c.id) AS comment_count
FROM upload u
JOIN "user" usr ON u.user_id = usr.id
LEFT JOIN "like" l ON l.upload_id = u.id
LEFT JOIN comment c ON c.upload_id = u.id AND c.deleted_at IS NULL
WHERE u.deleted_at IS NULL
AND usr.uploads_hidden = FALSE
AND usr.is_banned = FALSE
GROUP BY u.id, usr.display_name, usr.is_banned, usr.uploads_hidden;

View File

@@ -0,0 +1,55 @@
-- Make a feed page cost a page, not the whole event.
--
-- The previous definition (016) computed like_count/comment_count with LEFT JOINs and a
-- GROUP BY. Postgres CAN push the `event_id = $1` qual and the keyset predicate through the
-- view — verified with EXPLAIN, it uses idx_upload_event_created_id — but it CANNOT push
-- ORDER BY ... LIMIT across a GroupAggregate. So every feed request aggregated every upload in
-- the event (times its likes and comments) and only then sorted and took 21 rows. The cost
-- grew with the event, not with the page, and page 1 — the most expensive one — is exactly
-- what refreshFeedInPlace refetches on every completed upload, from every open feed in the
-- venue.
--
-- Correlated scalar subqueries move the counts ABOVE the Limit in the plan: they are evaluated
-- once per returned row, so 21 index lookups instead of a full aggregation.
--
-- The rewrite is EXACTLY equivalent, not merely close:
-- * "like" is keyed (upload_id, user_id), so COUNT(DISTINCT l.user_id) == count(*).
-- * comment.id is the primary key, so COUNT(DISTINCT c.id) == count(*).
-- * one row per upload either way — the GROUP BY was on u.id.
-- Column names, order and types are unchanged (count(*) and COUNT(DISTINCT ...) are both
-- bigint), so no Rust code changes.
--
-- No new index needed: idx_like_upload plus the (upload_id, user_id) PK serve the like
-- subquery, and idx_comment_upload ... WHERE deleted_at IS NULL matches the comment
-- subquery's predicate exactly.
--
-- One thing a future editor needs to know: the hashtag-filtered feed joins upload_hashtag
-- against this view. That was safe before only because the GROUP BY collapsed the join
-- fan-out; it is safe now because the view is one row per upload and upload_hashtag is keyed
-- (upload_id, hashtag_id) with a single tag filtered. Adding a second tag filter would need
-- fresh thought.
-- Not CASCADE: if something ever comes to depend on this view, the migration should fail
-- loudly rather than silently drop it.
DROP VIEW IF EXISTS v_feed;
CREATE VIEW v_feed AS
SELECT
u.id,
u.event_id,
u.user_id,
usr.display_name AS uploader_name,
usr.is_banned,
usr.uploads_hidden,
u.preview_path,
u.thumbnail_path,
u.display_path,
u.mime_type,
u.caption,
u.created_at,
(SELECT count(*) FROM "like" l WHERE l.upload_id = u.id) AS like_count,
(SELECT count(*) FROM comment c WHERE c.upload_id = u.id AND c.deleted_at IS NULL) AS comment_count
FROM upload u
JOIN "user" usr ON u.user_id = usr.id
WHERE u.deleted_at IS NULL
AND usr.uploads_hidden = FALSE
AND usr.is_banned = FALSE;

View File

@@ -0,0 +1,11 @@
-- NOTE: the reserved-name rename in the up migration is NOT reversible. The original names
-- are not recorded anywhere, and reversing it would in any case re-create the state that
-- bricked admin login. Rolling back the schema does not roll back that data change.
DELETE FROM config WHERE key IN (
'recover_name_rate_per_15min',
'pin_reset_ip_rate_per_min',
'upload_edit_rate_per_min',
'upload_edit_rate_enabled'
);
ALTER TABLE "user" DROP COLUMN IF EXISTS last_failed_pin_at;

View File

@@ -0,0 +1,48 @@
-- Two independent auth defects that share a migration because they share a table.
-- 1. RESERVED NAMES — free any guest squatting on a name the admin path used to depend on.
--
-- Migration 007 made display_name unique per event case-insensitively, and `join` had no
-- reserved-name guard. So any guest could join as "admin"/"Admin"/"ADMIN" before the operator's
-- first admin login; admin_login then looked its user up BY NAME, missed (wrong role), fell
-- through to creating "Admin", violated that unique index, and returned a 500 — permanently,
-- with no in-app recovery. Moderation, config and gallery release all gone, fixed only by SQL.
--
-- The real fix is in code (look the admin up by role, never by name — see auth/handlers.rs).
-- This clears the state an already-deployed database may be carrying.
--
-- RENAMED, NEVER DELETED: the guest keeps their uploads, their PIN and their session. Only
-- non-admin rows are touched — a real admin row named "Admin" is the expected state.
UPDATE "user" u
SET display_name = u.display_name || ' (' || left(u.id::text, 8) || ')'
WHERE u.role <> 'admin'
AND lower(u.display_name) IN ('admin', 'administrator', 'host', 'eventsnap');
-- 2. PIN LOCKOUT DECAY.
--
-- failed_pin_attempts only ever cleared on a successful recovery or after a lockout expired, so
-- honest typos accumulated across days: a guest who fat-fingered their PIN twice last night
-- arrives today already two-thirds of the way to being locked out. With the threshold now
-- raised (see below) a decay window is what keeps that raise safe rather than merely lenient.
ALTER TABLE "user" ADD COLUMN IF NOT EXISTS last_failed_pin_at TIMESTAMPTZ;
-- Rate-limit knobs introduced with this release.
--
-- recover_name_rate_per_15min (4, was a hardcoded 5): the per-(IP, name) ceiling. It MUST stay
-- below the account-lock threshold, which is the whole defect — at 5-per-IP against a 3-strike
-- lock, three requests from one IP locked any guest whose name is visible on the feed, every 15
-- minutes, forever. The lock threshold moves to 12 in code, so locking a victim now needs at
-- least three distinct sources while an honest guest never comes close.
--
-- pin_reset_ip_rate_per_min (30): /recover/request was the one unauthenticated endpoint with no
-- per-IP ceiling at all — /join got one in 017 and /recover in 019, and this third one was
-- simply missed. Its per-name key is attacker-chosen, so cycling names minted a fresh bucket
-- every time and the per-IP cost was unbounded.
--
-- upload_edit_rate_per_min (30): PATCH /upload/{id} had no rate limit of any kind.
INSERT INTO config (key, value) VALUES
('recover_name_rate_per_15min', '4'),
('pin_reset_ip_rate_per_min', '30'),
('upload_edit_rate_per_min', '30'),
('upload_edit_rate_enabled', 'true')
ON CONFLICT (key) DO NOTHING;

View File

@@ -19,6 +19,45 @@ use crate::services::config;
use crate::services::rate_limiter::client_ip;
use crate::state::AppState;
/// Names a guest may not take.
///
/// Defence in depth only. The real fix for the admin-lockout defect is that `admin_login` now
/// resolves its user by ROLE rather than by name (see `User::find_admin_for_event`), which is
/// why homoglyph and zero-width bypasses of this list are not a concern: the name is no longer
/// load-bearing for anything. What this buys is that a guest cannot impersonate the host in the
/// feed's byline, and that "Admin" stays available for the admin row.
const RESERVED_DISPLAY_NAMES: &[&str] = &["admin", "administrator", "host", "eventsnap"];
fn is_reserved_display_name(name: &str) -> bool {
let name = name.trim().to_lowercase();
RESERVED_DISPLAY_NAMES.contains(&name.as_str())
}
/// Trim and bounds-check a display name.
///
/// Shared by `join`, `recover` and `request_pin_reset` so the length check happens BEFORE the
/// name is used to build a rate-limiter key. It was inline in `join` only, so on the other two
/// endpoints `format!("...:{ip}:{name_key}")` allocated from an unbounded, attacker-chosen
/// string and stored it in a HashMap pruned once an hour with a 24 h ceiling — turning the
/// limiter itself into the memory-exhaustion primitive it exists to prevent.
fn validate_display_name(raw: &str) -> Result<&str, AppError> {
let name = raw.trim();
let chars = name.chars().count();
if chars == 0 || chars > 50 {
return Err(AppError::BadRequest(
"Name muss zwischen 1 und 50 Zeichen lang sein.".into(),
));
}
// Postgres rejects 0x00 in TEXT columns with a 500. Catch it here so callers see a clean
// 400 instead of an internal error.
if name.contains('\0') {
return Err(AppError::BadRequest(
"Name enthält ungültige Zeichen.".into(),
));
}
Ok(name)
}
#[derive(Deserialize)]
pub struct JoinRequest {
pub display_name: String,
@@ -62,19 +101,13 @@ pub async fn join(
}
}
let display_name = body.display_name.trim();
let name_chars = display_name.chars().count();
if name_chars == 0 || name_chars > 50 {
return Err(AppError::BadRequest(
"Name muss zwischen 1 und 50 Zeichen lang sein.".into(),
));
}
// Postgres rejects 0x00 in TEXT columns with a 500. Catch it here so callers
// see a clean 400 instead of an internal error.
if display_name.contains('\0') {
return Err(AppError::BadRequest(
"Name enthält ungültige Zeichen.".into(),
));
let display_name = validate_display_name(&body.display_name)?;
if is_reserved_display_name(display_name) {
// 409, matching the name-taken response below, so the frontend's existing handling
// works unchanged. See RESERVED_DISPLAY_NAMES for why this exists.
return Err(AppError::Conflict(format!(
"Der Name \"{display_name}\" ist reserviert. Bitte wähle einen anderen."
)));
}
// Per-guest bucket, keyed like the `recover:{ip}:{name}` limiter below. This carries
@@ -151,6 +184,28 @@ pub async fn join(
))
}
/// Default for `recover_name_rate_per_15min` — wrong PINs allowed per (IP, name) per 15 min.
/// Mirrors migration 023; kept here so the invariant below can be asserted in a test.
const RECOVER_NAME_CEILING_DEFAULT: usize = 4;
/// Wrong PINs, from ALL sources, before the ACCOUNT itself is locked for 15 minutes.
///
/// This was 3, which sat BELOW the per-(IP, name) ceiling of 5 — and that ordering, not the
/// number, was the defect. Display names are public on the feed, so three requests from a single
/// IP locked any guest out of their own account, repeatable every 15 minutes, indefinitely. The
/// tier meant to protect a guest was the easiest way to attack them.
///
/// Raised deliberately far above the per-IP tier so the two do different jobs. The per-(IP, name)
/// bucket is what stops a guesser, and it costs the ATTACKER. This tier is the last line against
/// a DISTRIBUTED guesser, and it is the only one an attacker can turn on a victim — so reaching
/// it must require at least three distinct sources inside the decay window.
///
/// Brute-force cost is unchanged: 12 attempts per 15 minutes is 48/hour against one account, so
/// 10 000 four-digit PINs still take ~208 hours no matter how many IPs are used. An honest guest
/// fat-fingering a 4-digit PIN never comes close, and `increment_failed_pin` now decays the
/// streak after 15 minutes so yesterday's typos don't count toward today's.
const PIN_LOCK_THRESHOLD: i16 = 12;
#[derive(Deserialize)]
pub struct RecoverRequest {
pub display_name: String,
@@ -203,13 +258,15 @@ pub async fn recover(
headers: HeaderMap,
Json(body): Json<RecoverRequest>,
) -> Result<Json<RecoverResponse>, AppError> {
let display_name = body.display_name.trim();
// Validated BEFORE it is used as a rate-limiter key — see `validate_display_name`. The
// per-IP ceiling below is keyed only on the IP, so it is safe to run either side of this;
// the per-NAME bucket is not.
let display_name = validate_display_name(&body.display_name)?;
// Per-IP+name throttle BEFORE the per-user 3-strike counter. Without this
// an attacker who knows a display name (they're visible on the feed) can
// burn through 3 wrong PINs and lock the victim for 15 minutes — repeated
// every 15 minutes, indefinitely. 5 attempts per 15 minutes per (IP, name)
// softens that into a real cost.
// Per-IP+name throttle BEFORE the per-user lockout counter. Without this an attacker who
// knows a display name (they're visible on the feed) can burn through the victim's wrong-PIN
// budget and lock them out, repeatedly. The ceiling here MUST stay below
// PIN_LOCK_THRESHOLD — see the constant for why that ordering is the whole control.
let ip = client_ip(&headers, &peer.ip().to_string());
let rate_limits_on = config::get_bool(&state.config_cache, "rate_limits_enabled", true).await;
let recover_rate_on = config::get_bool(&state.config_cache, "recover_rate_enabled", true).await;
@@ -234,10 +291,16 @@ pub async fn recover(
));
}
let name_ceiling = config::get_usize(
&state.config_cache,
"recover_name_rate_per_15min",
RECOVER_NAME_CEILING_DEFAULT,
)
.await;
let name_key = display_name.to_lowercase();
if let Err(retry_after_secs) = state.rate_limiter.check_with_retry(
format!("recover:{ip}:{name_key}"),
5,
name_ceiling,
Duration::from_secs(15 * 60),
) {
return Err(AppError::TooManyRequests(
@@ -317,7 +380,7 @@ pub async fn recover(
attempts,
"recover: wrong PIN"
);
if attempts >= 3 {
if attempts >= PIN_LOCK_THRESHOLD {
let lockout = Utc::now() + chrono::Duration::minutes(15);
User::lock_pin(&state.pool, user.id, lockout).await?;
tracing::warn!(
@@ -400,27 +463,16 @@ pub async fn admin_login(
)
.await?;
// Find or create the admin user for this event
let admin_name = "Admin";
let users = User::find_by_event_and_name(&state.pool, event.id, admin_name).await?;
let admin_user = if let Some(u) = users.into_iter().find(|u| u.role == UserRole::Admin) {
u
} else {
// Admin authenticates via password, but the schema still requires a PIN
// hash. Generate a random unguessable PIN so the recovery path remains
// unusable as an escalation route even if the role flag ever got cleared.
let dummy_pin: String = (0..32)
.map(|_| rand::rng().random_range(b'a'..=b'z') as char)
.collect();
let dummy_hash = hash_password(dummy_pin.clone(), 4).await?;
let user = User::create(&state.pool, event.id, admin_name, &dummy_hash).await?;
sqlx::query("UPDATE \"user\" SET role = 'admin' WHERE id = $1")
.bind(user.id)
.execute(&state.pool)
.await?;
User::find_by_id(&state.pool, user.id)
.await?
.ok_or_else(|| AppError::Internal(anyhow::anyhow!("admin user creation failed")))?
// Find or create the admin user for this event — BY ROLE, never by name.
//
// The name lookup this replaces is what made admin login brickable. Migration 007 makes
// display_name unique per event case-insensitively and `join` had no reserved-name guard,
// so a guest joining as "admin" before the operator's first login made the lookup miss on
// role, the fallback `create("Admin")` violate that index, and `?` return a permanent 500 —
// taking out moderation, config and gallery release with no in-app recovery.
let admin_user = match User::find_admin_for_event(&state.pool, event.id).await? {
Some(u) => u,
None => create_admin_user(&state, event.id).await?,
};
tracing::info!(user_id = %admin_user.id, event_id = %event.id, ip = %ip, "admin_login: success");
@@ -445,6 +497,51 @@ pub async fn admin_login(
}))
}
/// Create this event's admin row on first successful admin login.
///
/// Prefers the name "Admin". If a legacy database has a guest squatting on it — the state
/// migration 023 renames away, but a row could also predate that or be created between
/// migrations — falls back to a suffixed name rather than failing the login.
///
/// PROMOTING THE SQUATTING ROW WOULD BE A SERIOUS MISTAKE, and is the obvious-looking fix, so
/// it is spelled out: that row carries a `recovery_pin_hash` the guest knows. Setting
/// `role = 'admin'` on it would hand them the admin dashboard through `/recover`, permanently,
/// via a path that needs no password. A separate row under an uglier name is worse UX and much
/// better security — and since the lookup is now by role, the fallback name never has to be
/// guessed again on a later login.
async fn create_admin_user(state: &AppState, event_id: Uuid) -> Result<User, AppError> {
// Admin authenticates via password, but the schema still requires a PIN hash. Generate a
// random unguessable one so the recovery path stays unusable as an escalation route even if
// the role flag were ever cleared.
let dummy_pin: String = (0..32)
.map(|_| rand::rng().random_range(b'a'..=b'z') as char)
.collect();
let dummy_hash = hash_password(dummy_pin, 4).await?;
match User::create_with_role(&state.pool, event_id, "Admin", &dummy_hash, UserRole::Admin).await
{
Ok(u) => Ok(u),
Err(sqlx::Error::Database(db)) if db.is_unique_violation() => {
let fallback = format!("Admin-{}", &Uuid::new_v4().to_string()[..8]);
tracing::warn!(
%event_id, %fallback,
"the name \"Admin\" is held by a non-admin user; creating the admin under a \
fallback name. Rename that guest to free it — do NOT promote their row, they \
know its recovery PIN."
);
Ok(User::create_with_role(
&state.pool,
event_id,
&fallback,
&dummy_hash,
UserRole::Admin,
)
.await?)
}
Err(e) => Err(e.into()),
}
}
pub async fn logout(State(state): State<AppState>, auth: AuthUser) -> Result<StatusCode, AppError> {
Session::delete_by_token_hash(&state.pool, &auth.token_hash).await?;
Ok(StatusCode::NO_CONTENT)
@@ -476,9 +573,38 @@ pub async fn request_pin_reset(
headers: HeaderMap,
Json(body): Json<PinResetRequestBody>,
) -> Result<StatusCode, AppError> {
let display_name = body.display_name.trim();
let ip = client_ip(&headers, &peer.ip().to_string());
let rate_limits_on = config::get_bool(&state.config_cache, "rate_limits_enabled", true).await;
// Coarse per-IP ceiling FIRST, keyed only on the IP so its key is bounded by construction.
// /join got one of these in migration 017 and /recover in 019; this third unauthenticated
// endpoint was simply missed — migration 019's own comment describes exactly this attack.
// Without it, the per-name bucket below is no ceiling at all: the name is attacker-chosen,
// so cycling names mints a fresh bucket every request.
if rate_limits_on {
let ip_ceiling =
config::get_usize(&state.config_cache, "pin_reset_ip_rate_per_min", 30).await;
if let Err(retry_after_secs) = state.rate_limiter.check_with_retry(
format!("pin_reset_ip:{ip}"),
ip_ceiling,
Duration::from_secs(60),
) {
return Err(AppError::TooManyRequests(
"Zu viele Anfragen. Bitte warte kurz und versuche es erneut.".into(),
Some(retry_after_secs),
));
}
}
// Validated BEFORE the per-name key is built, so an unbounded name can never be retained in
// the limiter map. NOTE the 204: this endpoint's contract is that it answers identically
// whether or not the name exists, so it cannot enumerate guests. A 400 here would be a new
// signal — it would distinguish a malformed name from a well-formed unknown one. Silence is
// the correct response, and matches what an empty name already did.
let Ok(display_name) = validate_display_name(&body.display_name) else {
return Ok(StatusCode::NO_CONTENT);
};
if rate_limits_on {
let name_key = display_name.to_lowercase();
if let Err(retry_after_secs) = state.rate_limiter.check_with_retry(
@@ -492,9 +618,6 @@ pub async fn request_pin_reset(
));
}
}
if display_name.is_empty() {
return Ok(StatusCode::NO_CONTENT);
}
// Single statement so the existing-name and unknown-name paths do IDENTICAL work
// (same event+user index scans, an INSERT that matches 0 rows for an unknown name) —
@@ -524,3 +647,62 @@ pub async fn request_pin_reset(
Ok(StatusCode::NO_CONTENT)
}
#[cfg(test)]
mod tests {
use super::*;
/// THE defect, stated as arithmetic: the account-lock threshold sat BELOW the per-(IP, name)
/// attempt ceiling, so a single IP could exhaust it and lock any guest whose display name is
/// visible on the feed — every 15 minutes, indefinitely. The tier meant to protect a guest
/// was the cheapest way to attack them.
///
/// The fix is the ORDERING, not either number on its own, so that is what this pins.
#[test]
fn one_ip_cannot_reach_the_account_lock() {
assert!(
PIN_LOCK_THRESHOLD as usize >= RECOVER_NAME_CEILING_DEFAULT * 3,
"locking a victim must require at least three distinct sources; \
threshold {PIN_LOCK_THRESHOLD} vs per-IP ceiling {RECOVER_NAME_CEILING_DEFAULT}"
);
}
/// Raising the threshold must not quietly weaken brute-force resistance. 4-digit PINs, and
/// the lockout window is 15 minutes, so an attacker gets PIN_LOCK_THRESHOLD tries per window.
#[test]
fn the_raised_threshold_still_makes_guessing_a_four_digit_pin_impractical() {
let attempts_per_hour = PIN_LOCK_THRESHOLD as u64 * 4; // four 15-minute windows
let hours_for_full_keyspace = 10_000 / attempts_per_hour;
assert!(
hours_for_full_keyspace >= 168,
"exhausting 10k PINs would take {hours_for_full_keyspace}h — under a week is too fast"
);
}
#[test]
fn reserved_names_are_matched_case_insensitively_and_trimmed() {
for name in ["admin", "Admin", "ADMIN", " Host ", "EventSnap"] {
assert!(is_reserved_display_name(name), "{name} must be reserved");
}
}
/// A substring match here would reject perfectly ordinary names, which is a worse outcome
/// than the impersonation the list guards against.
#[test]
fn names_that_merely_contain_a_reserved_word_are_allowed() {
for name in ["Administrata", "Hostess", "Adminah", "Ghost", "hosting"] {
assert!(!is_reserved_display_name(name), "{name} must be allowed");
}
}
#[test]
fn display_names_are_bounded_before_they_can_become_a_rate_limit_key() {
assert!(validate_display_name(" Lena ").is_ok());
assert_eq!(validate_display_name(" Lena ").unwrap(), "Lena");
// The case that made the limiter itself the exhaustion primitive.
assert!(validate_display_name(&"a".repeat(51)).is_err());
assert!(validate_display_name(&"a".repeat(2_000_000)).is_err());
assert!(validate_display_name(" ").is_err());
assert!(validate_display_name("bad\0name").is_err());
}
}

View File

@@ -19,6 +19,12 @@ pub enum AppError {
/// the client can treat it as *terminal* (413, no retry) instead of backing off and
/// retrying a permanently-failing upload forever.
QuotaExceeded(String),
/// The server is temporarily unable to serve this request — currently only pool
/// saturation. Distinct from `Internal` because it is TRANSIENT and the client should be
/// told so: a 500 reads as "this request is broken", while a 503 + Retry-After reads as
/// "come back shortly", which is what the upload queue's retry classifier needs to make
/// the right call. Second field: optional retry-after seconds.
ServiceUnavailable(String, Option<u64>),
Internal(anyhow::Error),
}
@@ -33,6 +39,9 @@ impl AppError {
Self::Conflict(_) => (StatusCode::CONFLICT, "conflict"),
Self::TooManyRequests(..) => (StatusCode::TOO_MANY_REQUESTS, "too_many_requests"),
Self::QuotaExceeded(_) => (StatusCode::PAYLOAD_TOO_LARGE, "quota_exceeded"),
Self::ServiceUnavailable(..) => {
(StatusCode::SERVICE_UNAVAILABLE, "service_unavailable")
}
Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "internal_error"),
}
}
@@ -46,6 +55,7 @@ impl AppError {
| Self::NotFound(msg)
| Self::Conflict(msg) => msg.clone(),
Self::TooManyRequests(msg, _) => msg.clone(),
Self::ServiceUnavailable(msg, _) => msg.clone(),
Self::QuotaExceeded(msg) => msg.clone(),
Self::Internal(err) => {
tracing::error!("internal error: {err:#}");
@@ -58,10 +68,13 @@ impl AppError {
impl IntoResponse for AppError {
fn into_response(self) -> Response {
let (status, code) = self.status_and_code();
let retry_after_secs = if let Self::TooManyRequests(_, Some(secs)) = &self {
Some(*secs)
} else {
None
// BOTH retry-carrying variants must be matched here. `message()` would fail to
// compile on a missing arm; this one would not — it would silently drop the header and
// the `retry_after_secs` body field, which is exactly the sort of omission that only
// shows up under the load the 503 exists for.
let retry_after_secs = match &self {
Self::TooManyRequests(_, secs) | Self::ServiceUnavailable(_, secs) => *secs,
_ => None,
};
let message = self.message();
@@ -93,6 +106,84 @@ impl From<anyhow::Error> for AppError {
impl From<sqlx::Error> for AppError {
fn from(err: sqlx::Error) -> Self {
Self::Internal(err.into())
match err {
// Pool saturation is load, not a bug. Reporting it as a 500 was actively harmful:
// the frontend's upload-queue classifier treats 5xx as transient and retries, so
// the retries piled straight back into the saturated pool with no Retry-After to
// pace them. A 503 says the same thing honestly and carries the backoff.
//
// `PoolClosed` stays `Internal` — it only happens during shutdown, where a 503
// would invite a retry against a server that is going away.
sqlx::Error::PoolTimedOut => {
tracing::warn!("database pool exhausted; shedding a request with 503");
Self::ServiceUnavailable(
"Server ist gerade ausgelastet. Bitte versuche es in ein paar Sekunden erneut."
.into(),
Some(POOL_TIMEOUT_RETRY_AFTER_SECS),
)
}
other => Self::Internal(other.into()),
}
}
}
/// Retry-After for a shed request. Short: pool saturation clears in seconds once the queue
/// drains, and a long value would make a brief spike feel like an outage.
const POOL_TIMEOUT_RETRY_AFTER_SECS: u64 = 3;
#[cfg(test)]
mod tests {
use super::*;
/// `into_response` extracts `retry_after_secs` by MATCHING ON VARIANTS, so unlike
/// `message()` a missing arm is not a compile error — it silently drops the header. Pin the
/// behaviour for both retry-carrying variants.
#[test]
fn both_retry_carrying_variants_emit_retry_after() {
for err in [
AppError::TooManyRequests("slow down".into(), Some(42)),
AppError::ServiceUnavailable("busy".into(), Some(3)),
] {
let expected = match &err {
AppError::TooManyRequests(_, Some(s)) | AppError::ServiceUnavailable(_, Some(s)) => {
s.to_string()
}
_ => unreachable!(),
};
let resp = err.into_response();
assert_eq!(
resp.headers()
.get(axum::http::header::RETRY_AFTER)
.and_then(|v| v.to_str().ok()),
Some(expected.as_str()),
"a shed/throttled client must be told when to come back"
);
}
}
/// Pool saturation is load, not a bug. A 500 makes the frontend's retry classifier pile
/// straight back into the saturated pool with no backoff to pace it.
#[test]
fn pool_exhaustion_sheds_with_503_but_shutdown_does_not() {
let shed: AppError = sqlx::Error::PoolTimedOut.into();
assert_eq!(
shed.status_and_code(),
(StatusCode::SERVICE_UNAVAILABLE, "service_unavailable")
);
// PoolClosed only happens during shutdown; a 503 there would invite a retry against a
// server that is going away.
let closing: AppError = sqlx::Error::PoolClosed.into();
assert_eq!(
closing.status_and_code(),
(StatusCode::INTERNAL_SERVER_ERROR, "internal_error")
);
// Everything else must keep its existing mapping.
let missing: AppError = sqlx::Error::RowNotFound.into();
assert_eq!(
missing.status_and_code(),
(StatusCode::INTERNAL_SERVER_ERROR, "internal_error")
);
}
}

View File

@@ -38,7 +38,26 @@ pub async fn issue_ticket(
State(state): State<AppState>,
auth: AuthUser,
) -> Result<Json<StreamTicketResponse>, AppError> {
let ticket = state.sse_tickets.issue(auth.token_hash);
// The endpoint had no rate limit at all. Authentication is not a bound here: one valid
// session could loop it freely. 60/min is far above a real client (one ticket per SSE
// (re)connect, and reconnects are backed off) while capping a loop.
if let Err(retry_after_secs) = state.rate_limiter.check_with_retry(
format!("sse_ticket:{}", auth.user_id),
60,
Duration::from_secs(60),
) {
return Err(AppError::TooManyRequests(
"Zu viele Verbindungsversuche. Bitte warte kurz.".into(),
Some(retry_after_secs),
));
}
let ticket = state.sse_tickets.issue(auth.token_hash).ok_or_else(|| {
AppError::ServiceUnavailable(
"Server ist gerade ausgelastet. Live-Updates folgen in Kürze.".into(),
Some(30),
)
})?;
let server_time = sqlx::query_scalar("SELECT NOW()")
.fetch_one(&state.pool)
.await?;

View File

@@ -17,6 +17,135 @@ use crate::state::AppState;
const MAX_CAPTION_LENGTH: usize = 2000;
/// Byte ceiling for the caption field, enforced WHILE reading it.
///
/// `Field::text()` buffers the entire field before returning, and this is the one route whose
/// `DefaultBodyLimit` is raised to 576 MiB (main.rs) — so `caption=<576 MiB of text>` allocated
/// 576 MiB of heap per concurrent request inside a 1 GiB container, and the
/// `MAX_CAPTION_LENGTH` check only ran afterwards, on a string that had already been built.
/// 4 bytes per code point is the worst case for UTF-8, so this can never reject a caption the
/// character limit would have accepted.
const MAX_CAPTION_BYTES: usize = MAX_CAPTION_LENGTH * 4;
/// Byte ceiling for the raw hashtag CSV. Generous next to what the tag caps below allow.
const MAX_HASHTAGS_BYTES: usize = 4 * 1024;
/// Hashtags stored per upload. The CSV was never length-checked at all and was split into an
/// unbounded `Vec`, then upserted TAG BY TAG inside the commit transaction — which holds a
/// `FOR SHARE` lock on the event row, so one request could stall every other upload behind
/// tens of thousands of round trips.
const MAX_HASHTAGS_PER_UPLOAD: usize = 30;
/// Characters per stored tag. `extract_hashtags` already self-bounds at 40; this covers the CSV
/// path, which had no bound of its own.
const MAX_HASHTAG_LENGTH: usize = 50;
/// Read a multipart text field, refusing it the moment it exceeds `max_bytes`.
///
/// The point is to fail DURING the read rather than after it — `Field::text()` cannot, because
/// it has already allocated the whole thing by the time it returns.
async fn read_text_field_bounded(
mut field: axum::extract::multipart::Field<'_>,
max_bytes: usize,
) -> Result<String, AppError> {
let mut buf: Vec<u8> = Vec::new();
while let Some(chunk) = field
.chunk()
.await
.map_err(|e| AppError::BadRequest(e.to_string()))?
{
if buf.len() + chunk.len() > max_bytes {
return Err(AppError::BadRequest(
"Eingabe ist zu lang.".to_string(),
));
}
buf.extend_from_slice(&chunk);
}
String::from_utf8(buf).map_err(|_| AppError::BadRequest("Ungültige Zeichenkodierung.".into()))
}
/// Normalise, dedupe and CAP the tags for one upload.
///
/// Extracted as a pure function so the caps are testable without standing up multipart, and
/// shared by the upload and edit paths — which previously disagreed: upload lowercased and
/// stripped `#`, while edit upserted raw strings, so `#Party` via edit and `party` via upload
/// became two different hashtag rows.
///
/// Truncates rather than rejecting. `extract_hashtags` legitimately derives tags from a
/// 2000-character caption, and 400-ing a guest for writing an enthusiastic caption would be a
/// worse outcome than silently keeping the first 30.
fn normalize_tags(caption_tags: Vec<String>, csv: Option<&str>) -> Vec<String> {
let mut tags = caption_tags;
if let Some(csv) = csv {
for tag in csv.split(',') {
let t = tag.trim().trim_start_matches('#').to_lowercase();
if !t.is_empty() {
tags.push(t);
}
}
}
tags.sort();
tags.dedup();
tags.retain(|t| t.chars().count() <= MAX_HASHTAG_LENGTH);
tags.truncate(MAX_HASHTAGS_PER_UPLOAD);
tags
}
/// Owns the bytes an in-flight upload has written to disk, and deletes them unless the
/// request reaches the point where a database row takes ownership.
///
/// Reclaim used to be a dozen explicit `remove_file` calls on the handler's return paths.
/// That covers every way the handler can FINISH, and none of the ways it can simply STOP:
/// when a client disconnects mid-body — a phone leaving wifi, iOS killing a backgrounded
/// PWA, the user hitting back — axum drops the handler future at a `.await` inside
/// `field.chunk()`, and no return path runs at all. The partial file then survives forever:
/// it has no upload row, so `cleanup_deleted_media` (which is row-driven) can never see it,
/// and no sweeper existed for the originals directory. Those bytes are also invisible to the
/// quota while still consuming the free disk that `quota_limit_bytes` divides among guests.
///
/// A drop guard is the only construct that survives cancellation, because dropping the future
/// is exactly what runs it.
struct TempFileGuard {
/// `None` once disarmed — a row now owns these bytes.
path: Option<std::path::PathBuf>,
}
impl TempFileGuard {
fn new(path: std::path::PathBuf) -> Self {
Self { path: Some(path) }
}
/// Follow the bytes to their new location after a rename.
///
/// NOT `disarm`. Between the rename and the commit the file exists under its FINAL name
/// with still no row pointing at it, so that window needs guarding just as much as the
/// `.tmp` did — arguably more, since a leftover final-named original looks legitimate.
fn retarget(&mut self, path: std::path::PathBuf) {
self.path = Some(path);
}
/// Hand ownership to the committed row. Only correct after `tx.commit()` succeeds.
fn disarm(&mut self) {
self.path = None;
}
}
impl Drop for TempFileGuard {
fn drop(&mut self) {
let Some(path) = self.path.take() else {
return;
};
// std::fs, not tokio::fs: `Drop` cannot await, and a runtime-dependent unlink is not
// guaranteed a live runtime here (shutdown drops in-flight tasks).
match std::fs::remove_file(&path) {
Ok(()) => tracing::debug!(path = %path.display(), "reclaimed an abandoned upload"),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => {
tracing::warn!(error = ?e, path = %path.display(), "failed to reclaim an abandoned upload")
}
}
}
}
/// Allowlist of accepted media types, keyed by the MIME that `infer` derives from
/// the file's magic bytes. The detected MIME (not the client-declared one) is what
/// we trust, store, and hand to the compression pipeline — so a text-based payload
@@ -107,6 +236,11 @@ pub async fn upload(
.media_path
.join(format!("originals/{event_slug}"));
let temp_abs = originals_dir.join(format!("{upload_id}.tmp"));
// Armed before anything can create the file, so there is no window in which bytes exist
// unowned. From here on, EVERY exit — return, `?`, panic, or the future being dropped
// mid-body by a client disconnect — reclaims them, and the explicit `remove_file` calls
// that used to be sprinkled over the return paths are gone. One owner, one rule.
let mut file_guard = TempFileGuard::new(temp_abs.clone());
let mut streamed: Option<(i64, Vec<u8>)> = None; // (size, head bytes for sniffing)
let mut caption: Option<String> = None;
@@ -115,8 +249,8 @@ pub async fn upload(
// doesn't send one and gets the previous behaviour.
let mut client_upload_id: Option<Uuid> = None;
// Wrap the multipart read so any error after the temp file is created still cleans
// it up (a mid-stream parse failure must not leave a stray `.tmp` on disk).
// The multipart read is wrapped so the field loop can use `?` freely; reclaiming the temp
// file on failure is `file_guard`'s job, not this block's.
let parse_result: Result<(), AppError> = async {
while let Some(field) = multipart
.next_field()
@@ -146,20 +280,10 @@ pub async fn upload(
streamed = Some(stream_field_to_file(field, &temp_abs, cap_bytes).await?);
}
"caption" => {
caption = Some(
field
.text()
.await
.map_err(|e| AppError::BadRequest(e.to_string()))?,
);
caption = Some(read_text_field_bounded(field, MAX_CAPTION_BYTES).await?);
}
"hashtags" => {
hashtags_csv = Some(
field
.text()
.await
.map_err(|e| AppError::BadRequest(e.to_string()))?,
);
hashtags_csv = Some(read_text_field_bounded(field, MAX_HASHTAGS_BYTES).await?);
}
"client_upload_id" => {
let raw = field
@@ -178,10 +302,7 @@ pub async fn upload(
}
.await;
if let Err(e) = parse_result {
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(e);
}
parse_result?;
// Idempotency, fast path: this key already has a live upload, so the previous attempt DID
// succeed and only its response was lost. Replay that response instead of storing the photo
@@ -192,12 +313,14 @@ pub async fn upload(
// cost and it has already been paid by the time we get here; what has to be prevented is a
// second ROW, and that is what this does. The concurrent case (two retries in flight at once)
// is caught by the unique index inside the transaction below.
//
// The temp file needs no explicit cleanup here: `file_guard` is armed and reclaims it when
// this early return drops it.
if let Some(cid) = client_upload_id
&& let Some(existing) = Upload::find_by_client_upload_id(&state.pool, auth.user_id, cid)
.await
.map_err(AppError::from)?
{
let _ = tokio::fs::remove_file(&temp_abs).await;
tracing::info!(
client_upload_id = %cid, upload_id = %existing.id,
"duplicate upload suppressed; replaying the original response"
@@ -206,8 +329,8 @@ pub async fn upload(
return Ok((StatusCode::OK, Json(dto)));
}
// From here on the temp file may exist; every validation failure removes it before
// returning so a rejected upload never leaves bytes behind.
// From here on the temp file may exist. Every exit reclaims it via `file_guard` — see
// TempFileGuard for why the explicit per-branch cleanup this replaced was not enough.
let (size, head) = match streamed {
Some(s) => s,
None => return Err(AppError::BadRequest("Keine Datei hochgeladen.".into())),
@@ -219,7 +342,6 @@ pub async fn upload(
if let Some(ref cap) = caption
&& cap.chars().count() > MAX_CAPTION_LENGTH
{
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(AppError::BadRequest(format!(
"Beschreibung ist zu lang. Maximum: {} Zeichen.",
MAX_CAPTION_LENGTH
@@ -234,7 +356,6 @@ pub async fn upload(
let kind = match infer::get(&head) {
Some(k) => k,
None => {
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(AppError::BadRequest(
"Dateityp nicht erkannt oder nicht unterstützt.".into(),
));
@@ -247,7 +368,6 @@ pub async fn upload(
{
Some(v) => v,
None => {
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(AppError::BadRequest(format!(
"Dateityp wird nicht unterstützt: {}.",
kind.mime_type()
@@ -262,7 +382,6 @@ pub async fn upload(
max_image_mb * 1024 * 1024
};
if size > max_bytes {
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(AppError::BadRequest(format!(
"Datei ist zu groß. Maximum: {} MB.",
max_bytes / (1024 * 1024)
@@ -331,7 +450,6 @@ pub async fn upload(
quota_limit = Some(limit);
let prospective_total = user.total_upload_bytes.saturating_add(size);
if prospective_total > limit {
let _ = tokio::fs::remove_file(&temp_abs).await;
return Err(AppError::QuotaExceeded(
// Name the remedy, because the guest cannot see the number. Every quota
// display is staff-gated by design, so a guest hitting this had no idea what
@@ -352,22 +470,20 @@ pub async fn upload(
tokio::fs::rename(&temp_abs, &absolute_path)
.await
.map_err(|e| AppError::Internal(e.into()))?;
// THERE MUST BE NO `.await` BETWEEN THE RENAME AND THIS LINE. Both statements resolve on
// the same poll, so the future cannot be dropped between them and the guard is never
// pointing at a path that no longer holds the bytes. If the rename fails the guard still
// owns `temp_abs`, which is why retargeting comes after it rather than before.
file_guard.retarget(absolute_path.clone());
// Process hashtags from caption and explicit CSV
let mut tags: Vec<String> = Vec::new();
if let Some(ref cap) = caption {
tags.extend(hashtag::extract_hashtags(cap));
}
if let Some(ref csv) = hashtags_csv {
for tag in csv.split(',') {
let t = tag.trim().trim_start_matches('#').to_lowercase();
if !t.is_empty() {
tags.push(t);
}
}
}
tags.sort();
tags.dedup();
// Process hashtags from caption and explicit CSV, capped — see `normalize_tags`.
let tags = normalize_tags(
caption
.as_deref()
.map(hashtag::extract_hashtags)
.unwrap_or_default(),
hashtags_csv.as_deref(),
);
// Quota accounting, the upload row, and its hashtag links must be atomic: a
// crash between the bytes increment and the insert would permanently charge
@@ -466,14 +582,17 @@ pub async fn upload(
}
.await;
// The file is already on disk at `absolute_path`. If the transaction failed, no DB
// row will ever reference it, so remove it now rather than orphan bytes on disk.
// The file is already on disk at `absolute_path`, and `file_guard` was retargeted to it
// above — so every path out of here that is NOT a successful commit leaves the guard armed
// and reclaims the bytes on the way out. That covers the concurrent-duplicate loser below
// as well as the plain error case, and unlike the explicit `remove_file` calls it replaces,
// it also covers axum dropping this future instead of returning.
let upload = match tx_result {
Ok(u) => u,
// The concurrent duplicate resolved inside the transaction. The winner's row is committed;
// answer with it so both retries of the same photo get the same successful reply.
// answer with it so both retries of the same photo get the same successful reply. The
// loser's bytes are reclaimed by the guard when this return drops it.
Err(AppError::Conflict(ref marker)) if marker == DUPLICATE_UPLOAD_MARKER => {
let _ = tokio::fs::remove_file(&absolute_path).await;
let existing = match client_upload_id {
Some(cid) => Upload::find_by_client_upload_id(&state.pool, auth.user_id, cid)
.await
@@ -492,11 +611,10 @@ pub async fn upload(
let dto = replay_upload_dto(&state, &existing, &user.display_name).await;
return Ok((StatusCode::OK, Json(dto)));
}
Err(e) => {
let _ = tokio::fs::remove_file(&absolute_path).await;
return Err(e);
}
Err(e) => return Err(e),
};
// The committed row now references these bytes — hand ownership over.
file_guard.disarm();
// Spawn compression task
state
@@ -551,6 +669,57 @@ pub async fn edit_upload(
return Err(AppError::Forbidden("Nur eigene Uploads bearbeiten.".into()));
}
// This endpoint had no rate limit of any kind, while every other mutating route has one.
let rate_limits_on = config::get_bool(&state.config_cache, "rate_limits_enabled", true).await;
let edit_rate_on = config::get_bool(&state.config_cache, "upload_edit_rate_enabled", true).await;
if rate_limits_on && edit_rate_on {
let edit_rate =
config::get_i64(&state.config_cache, "upload_edit_rate_per_min", 30).await as usize;
if let Err(retry_after_secs) = state.rate_limiter.check_with_retry(
format!("upload_edit:{}", auth.user_id),
edit_rate,
Duration::from_secs(60),
) {
return Err(AppError::TooManyRequests(
"Zu viele Änderungen. Bitte warte kurz.".into(),
Some(retry_after_secs),
));
}
}
// Validate to the same limits as the upload path. This route had none at all, so a caption
// rejected at upload could be set here instead, and the tags went in raw — meaning `#Party`
// via edit and `party` via upload became two different hashtag rows.
if let Some(ref caption) = body.caption
&& caption.chars().count() > MAX_CAPTION_LENGTH
{
return Err(AppError::BadRequest(format!(
"Beschreibung ist zu lang. Maximum: {MAX_CAPTION_LENGTH} Zeichen."
)));
}
let normalized_tags = body
.hashtags
.as_ref()
.map(|tags| normalize_tags(tags.clone(), None));
// A PATCH that changes nothing must not retire the keepsake generation.
//
// `invalidate_and_arm` below ran unconditionally, outside both `if let Some(...)` guards, so
// `PATCH {}` — which any authenticated guest can send in a loop against their own upload —
// bumped export_epoch and armed a fresh pair of full-gallery export workers every time.
// REGEN_DEBOUNCE bounds the rate of that, not the total work, so the keepsake could be kept
// permanently un-downloadable.
//
// Residual, deliberately not fixed: re-sending an IDENTICAL hashtag list still counts as a
// change. Comparing would need another query, and unlike `PATCH {}` it is not a free loop.
let caption_changed = match (&body.caption, &upload.caption) {
(Some(new), existing) => Some(new.as_str()) != existing.as_deref(),
(None, _) => false,
};
if !caption_changed && normalized_tags.is_none() {
return Ok(StatusCode::OK);
}
// Caption update + hashtag wipe-then-relink in one transaction, so a crash
// mid-relink can't leave the upload with its hashtags stripped.
//
@@ -566,7 +735,7 @@ pub async fn edit_upload(
if let Some(ref caption) = body.caption {
Upload::update_caption(&mut *tx, upload_id, Some(caption)).await?;
}
if let Some(ref hashtags) = body.hashtags {
if let Some(ref hashtags) = normalized_tags {
Hashtag::unlink_all_from_upload(&mut *tx, upload_id).await?;
// Sort + dedup before upserting, exactly as the upload path does. `Hashtag::upsert`
// takes row locks, so two transactions touching the same two tags in OPPOSITE order
@@ -1251,4 +1420,103 @@ mod tests {
fn full_tolerance_is_identity_for_a_single_uploader() {
assert_eq!(quota_limit_bytes(500, 1.0, 1), 500);
}
/// The guard is the only reclaim mechanism that survives a client disconnect, so its three
/// states have to be exactly right — a wrong `disarm` leaks bytes forever, a wrong `Drop`
/// deletes a committed guest's photo.
mod temp_file_guard {
use super::super::TempFileGuard;
fn scratch(name: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!("es-guard-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let p = dir.join(name);
std::fs::write(&p, b"bytes").unwrap();
p
}
#[test]
fn dropping_an_armed_guard_reclaims_the_file() {
let p = scratch("armed.tmp");
drop(TempFileGuard::new(p.clone()));
assert!(!p.exists(), "an abandoned upload must not survive the request");
}
#[test]
fn a_disarmed_guard_leaves_the_file_alone() {
let p = scratch("committed.jpg");
let mut g = TempFileGuard::new(p.clone());
g.disarm();
drop(g);
assert!(p.exists(), "a committed upload's bytes must never be deleted");
}
#[test]
fn retarget_follows_the_rename_and_forgets_the_old_path() {
let old = scratch("old.tmp");
let new = scratch("new.jpg");
let mut g = TempFileGuard::new(old.clone());
// The rename already moved the bytes; only the new path is at risk now.
std::fs::remove_file(&old).unwrap();
g.retarget(new.clone());
drop(g);
assert!(!new.exists(), "the final-named original is orphaned too until the row commits");
}
#[test]
fn a_guard_whose_file_is_already_gone_is_harmless() {
let p = scratch("vanished.tmp");
std::fs::remove_file(&p).unwrap();
drop(TempFileGuard::new(p)); // must not panic
}
}
/// The CSV hashtag path had no length check at all and was upserted tag-by-tag INSIDE the
/// commit transaction — which holds a FOR SHARE lock on the event row, so one request could
/// stall every other upload behind tens of thousands of round trips.
mod hashtag_caps {
use super::super::{MAX_HASHTAGS_PER_UPLOAD, MAX_HASHTAG_LENGTH, normalize_tags};
#[test]
fn a_huge_csv_is_capped_not_upserted_in_full() {
let csv = (0..10_000)
.map(|i| format!("tag{i}"))
.collect::<Vec<_>>()
.join(",");
let tags = normalize_tags(vec![], Some(&csv));
assert_eq!(tags.len(), MAX_HASHTAGS_PER_UPLOAD);
}
#[test]
fn an_overlong_tag_is_dropped_rather_than_stored() {
let long = "a".repeat(MAX_HASHTAG_LENGTH + 1);
let tags = normalize_tags(vec![], Some(&format!("ok,{long}")));
assert_eq!(tags, vec!["ok"]);
}
#[test]
fn tags_are_normalised_and_deduped_across_both_sources() {
// The upload path lowercased and stripped `#` while the edit path did not, so
// `#Party` and `party` became two different hashtag rows. One helper, one rule.
let tags = normalize_tags(vec!["party".into()], Some("#Party, PARTY ,tanz"));
assert_eq!(tags, vec!["party", "tanz"]);
}
#[test]
fn truncation_keeps_a_stable_prefix_not_an_arbitrary_one() {
// Sorted before truncation, so the same input always yields the same tags —
// otherwise an edit could silently shuffle which 30 survived.
let csv = "zulu,alpha,mike,bravo";
assert_eq!(
normalize_tags(vec![], Some(csv)),
normalize_tags(vec![], Some("bravo,mike,alpha,zulu"))
);
}
#[test]
fn an_empty_or_absent_csv_yields_nothing() {
assert!(normalize_tags(vec![], None).is_empty());
assert!(normalize_tags(vec![], Some(",, ,")).is_empty());
}
}
}

View File

@@ -27,8 +27,15 @@ async fn main() -> Result<()> {
tracing_subscriber::registry()
.with(
// `info`, not `debug`. A stock deploy sets RUST_LOG nowhere (it is absent from
// .env.example and was absent from docker-compose.yml), so this fallback IS the
// production level — and at `debug` the TraceLayer below emits a line per request
// AND per response, into a log file that had no rotation. `tower_http=warn`
// rather than `info` states the intent: those spans are diagnostics, not an
// access log, and a future `DefaultOnResponse::new().level(Level::INFO)` should
// not silently re-enable them.
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "eventsnap_backend=debug,tower_http=debug".into()),
.unwrap_or_else(|_| "eventsnap_backend=info,tower_http=warn".into()),
)
.with(tracing_subscriber::fmt::layer())
.init();
@@ -51,6 +58,12 @@ async fn main() -> Result<()> {
// originals are never touched, so a failure just retries on the next start.
state.compression.backfill_stale_derivatives().await;
// Re-extract poster frames for videos a restart interrupted. `startup_recovery` above
// marks their compression `failed` but nothing re-enqueued them, so `thumbnail_path`
// stayed NULL for the rest of the event. Shares the attempt budget with the image
// backfill, so a clip that genuinely yields no frame stops being retried.
state.compression.backfill_video_posters().await;
// Re-spawn exports for events that were released but whose keepsake never finished
// (crash mid-export). Needs the media/export paths + SSE sender, so it runs here
// rather than inside `startup_recovery`. Fire-and-forget: the workers run in the
@@ -249,6 +262,23 @@ async fn main() -> Result<()> {
// four subtrees. Deleting the route removes the vector outright rather than racing the
// decoder; `/media/**` now 404s regardless of encoding.
let router = Router::new()
// ONE probe, and it touches the database. The merge of the unattended-blockers work
// brought a competing design — a dependency-free `/health` for the compose gate plus
// a DB-backed `/health/ready` for an external monitor. That split is defensible, and
// it was rejected deliberately:
//
// * `/health` returning a constant "ok" is the exact defect faea555 fixed and
// verified live (stop Postgres → 503 → start → 200, with no app restart). Every
// request path touches the database, so a constant probe reports healthy while
// the app is useless — the disk-full endgame stayed green all the way down.
// * The split's motive was that Caddy's `depends_on: app: service_healthy` would
// be blocked by a Postgres hiccup at boot. But `app` itself already gates on
// `db: service_healthy`, so the DB is up before this probe ever runs, and the
// healthcheck carries a 20s start_period plus 5 retries on top.
// * The two handlers were the same `SELECT 1` with the same 2s timeout under two
// names, so keeping both bought nothing.
//
// The external uptime monitor the runbook now calls for points at this route.
.route("/health", get(health))
.merge(api)
.layer(TraceLayer::new_for_http())

View File

@@ -185,10 +185,62 @@ impl Upload {
/// Stamp which revision of the derivative pipeline produced this row's preview/display,
/// so the startup backfill can find rows generated by an older one exactly once.
///
/// Also clears the attempt counter: success is the only thing that resets it, and folding
/// the reset in here means both the live path and the backfill get it with no extra call
/// site to forget.
pub async fn set_derivatives_rev(pool: &PgPool, id: Uuid, rev: i16) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE upload SET derivatives_rev = $2 WHERE id = $1")
sqlx::query(
"UPDATE upload
SET derivatives_rev = $2, derivative_attempts = 0, derivative_last_error = NULL
WHERE id = $1",
)
.bind(id)
.bind(rev)
.execute(pool)
.await?;
Ok(())
}
/// Record that derivative processing is ABOUT to be attempted, returning the new count.
///
/// WRITE-AHEAD ON PURPOSE. The failure this bounds is a cgroup SIGKILL: the process
/// vanishes mid-work, so no `Err` is returned, no error handler runs and no `Drop` fires.
/// A counter incremented after a failure would increment zero times per crash and the
/// boot loop would be unchanged. Counting the ATTEMPT is the only thing that survives the
/// process dying. The cost is that a genuinely transient failure also burns an attempt —
/// acceptable, because the retry budget is per-boot-loop, not per-request, and success
/// resets it to zero.
/// `None` when the row no longer exists (hard-deleted, or an e2e TRUNCATE landed while the
/// task waited on the semaphore) — the caller should abandon quietly rather than treat a
/// missing row as a processing failure.
pub async fn begin_derivative_attempt(
pool: &PgPool,
id: Uuid,
) -> Result<Option<i16>, sqlx::Error> {
sqlx::query_scalar(
"UPDATE upload
SET derivative_attempts = derivative_attempts + 1
WHERE id = $1
RETURNING derivative_attempts",
)
.bind(id)
.fetch_optional(pool)
.await
}
/// Store why the last derivative attempt failed. Diagnostics only — nothing branches on it.
pub async fn record_derivative_failure(
pool: &PgPool,
id: Uuid,
error: &str,
) -> Result<(), sqlx::Error> {
// Bounded: an anyhow chain can be long, and this is written on a failure path that may
// repeat across every row of a bad batch.
let truncated: String = error.chars().take(500).collect();
sqlx::query("UPDATE upload SET derivative_last_error = $2 WHERE id = $1")
.bind(id)
.bind(rev)
.bind(truncated)
.execute(pool)
.await?;
Ok(())

View File

@@ -60,6 +60,54 @@ impl User {
.await
}
/// Create a user with an explicit role, in ONE statement.
///
/// `create` + a separate `UPDATE ... SET role` is not equivalent: a crash or a pool error
/// between the two leaves a GUEST row holding a reserved name, which is exactly the
/// poisoned state that bricked admin login — now self-inflicted, and invisible to a
/// role-based lookup, so the next login would create yet another.
pub async fn create_with_role(
pool: &PgPool,
event_id: Uuid,
display_name: &str,
pin_hash: &str,
role: UserRole,
) -> Result<Self, sqlx::Error> {
sqlx::query_as::<_, Self>(
"INSERT INTO \"user\" (event_id, display_name, recovery_pin_hash, role)
VALUES ($1, $2, $3, $4)
RETURNING *",
)
.bind(event_id)
.bind(display_name)
.bind(pin_hash)
.bind(role)
.fetch_one(pool)
.await
}
/// The event's admin, looked up BY ROLE.
///
/// The name is not the identity and never was. Looking the admin up by `display_name`
/// meant any guest who joined as "Admin" first made the lookup miss, and the fallback
/// `create` then violated the case-insensitive unique index from migration 007 — a
/// permanent 500 on admin login, recoverable only by hand-editing the database.
///
/// `ORDER BY created_at` so a database that somehow acquired two admin rows resolves to a
/// stable one rather than alternating between them.
pub async fn find_admin_for_event(
pool: &PgPool,
event_id: Uuid,
) -> Result<Option<Self>, sqlx::Error> {
sqlx::query_as::<_, Self>(
"SELECT * FROM \"user\" WHERE event_id = $1 AND role = 'admin'
ORDER BY created_at ASC LIMIT 1",
)
.bind(event_id)
.fetch_optional(pool)
.await
}
pub async fn find_by_id(pool: &PgPool, id: Uuid) -> Result<Option<Self>, sqlx::Error> {
sqlx::query_as::<_, Self>("SELECT * FROM \"user\" WHERE id = $1")
.bind(id)
@@ -96,14 +144,31 @@ impl User {
Ok(row.0)
}
/// Window after which a failed-PIN streak is forgotten. Matches the lockout duration, so
/// "wait out the cooldown" and "start clean" are the same interval to a guest.
const PIN_ATTEMPT_DECAY_MINUTES: i64 = 15;
/// Record a wrong PIN and return the CURRENT streak length.
///
/// The counter decays: before this, it only ever cleared on a successful recovery or after
/// a lockout expired, so ordinary typos accumulated across days and a guest could arrive at
/// an event already most of the way to being locked out by mistakes made the night before.
/// Decay is what makes the raised lock threshold safe rather than merely lenient.
pub async fn increment_failed_pin(pool: &PgPool, id: Uuid) -> Result<i16, sqlx::Error> {
let row: (i16,) = sqlx::query_as(
"UPDATE \"user\"
SET failed_pin_attempts = failed_pin_attempts + 1
SET failed_pin_attempts = CASE
WHEN last_failed_pin_at IS NULL
OR last_failed_pin_at < NOW() - ($2 || ' minutes')::interval
THEN 1
ELSE failed_pin_attempts + 1
END,
last_failed_pin_at = NOW()
WHERE id = $1
RETURNING failed_pin_attempts",
)
.bind(id)
.bind(Self::PIN_ATTEMPT_DECAY_MINUTES.to_string())
.fetch_one(pool)
.await?;
Ok(row.0)
@@ -124,7 +189,9 @@ impl User {
pub async fn reset_pin_attempts(pool: &PgPool, id: Uuid) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE \"user\" SET failed_pin_attempts = 0, pin_locked_until = NULL WHERE id = $1",
"UPDATE \"user\"
SET failed_pin_attempts = 0, pin_locked_until = NULL, last_failed_pin_at = NULL
WHERE id = $1",
)
.bind(id)
.execute(pool)

View File

@@ -13,6 +13,9 @@ use crate::state::SseEvent;
#[derive(Clone)]
pub struct CompressionWorker {
semaphore: Arc<Semaphore>,
/// Serialises the memory-heavy image jobs — see `HEAVY_IMAGE_BYTES`. Separate from
/// `semaphore` so ordinary photos keep full concurrency.
heavy: Arc<Semaphore>,
pool: PgPool,
media_path: PathBuf,
sse_tx: broadcast::Sender<SseEvent>,
@@ -31,6 +34,7 @@ impl CompressionWorker {
) -> Self {
Self {
semaphore: Arc::new(Semaphore::new(concurrency)),
heavy: Arc::new(Semaphore::new(1)),
pool,
media_path,
sse_tx,
@@ -58,6 +62,21 @@ impl CompressionWorker {
/// next start. Rev 1 = EXIF orientation is applied.
const DERIVATIVES_REV: i16 = 1;
/// How many times derivative generation may be ATTEMPTED for one upload before it is left
/// alone. Counted write-ahead and reset on success — see `Upload::begin_derivative_attempt`.
///
/// This is what turns a fatal input from an outage into a blemish. The startup backfill
/// runs unconditionally on every boot, so before this bound a row whose processing killed
/// the process was re-selected and re-run forever, and `restart: unless-stopped` made that
/// an infinite loop that also dropped every SSE stream and truncated every in-flight
/// upload on each cycle. Three attempts absorbs genuinely transient infrastructure
/// failures (an ENOSPC spike, a pool blip) without ever becoming unbounded.
const MAX_DERIVATIVE_ATTEMPTS: i16 = 3;
/// Rows regenerated per boot. Bounds both the query and the amount of work a single start
/// can queue; whatever is left is picked up on the next boot.
const BACKFILL_BATCH: i64 = 200;
/// Spawn a background task to process an uploaded file.
pub fn process(&self, upload_id: Uuid, original_path: String, mime_type: String) {
let worker = self.clone();
@@ -193,6 +212,21 @@ impl CompressionWorker {
let original = self.media_path.join(original_path);
if mime_type.starts_with("image/") {
// Count the attempt BEFORE doing the work — see `begin_derivative_attempt`. If this
// input is the one that kills the container, this write is the only record that
// survives, and it is what stops the boot backfill replaying it forever.
match Upload::begin_derivative_attempt(&self.pool, upload_id).await? {
Some(attempts) if attempts > Self::MAX_DERIVATIVE_ATTEMPTS => {
anyhow::bail!(
"derivative generation gave up after {} attempt(s)",
attempts - 1
);
}
Some(_) => {}
// The row vanished while this task waited on the semaphore. Nothing to do, and
// reporting a failure would broadcast into a stream that no longer has a card.
None => return Ok(()),
}
let (preview_rel, display_rel) = self
.generate_image_derivatives(upload_id, &original, mime_type)
.await?;
@@ -256,6 +290,38 @@ impl CompressionWorker {
/// Longest edge of the phone-feed "preview" (data-saver default).
const PREVIEW_MAX_EDGE: u32 = 800;
/// Above this pixel count the PNG original is stored as uploaded, unoptimised.
///
/// oxipng's peak memory scales with PIXELS, not file size: it decodes the PNG itself and
/// then evaluates row filters, each trial holding a full-size buffer. That is why a 2.82
/// MiB file could measure 1250 MiB of peak RSS inside a 1 GiB container — smooth,
/// synthetic content compresses to almost nothing on disk while still being 8000x8000.
/// 8 MP covers every real phone photo; beyond it we decline the (lossless, cosmetic)
/// saving rather than risk the OOM kill.
const OXIPNG_MAX_PIXELS: u64 = 8_000_000;
/// Estimated peak heap above which an image job takes the exclusive `heavy` permit.
///
/// `compression_concurrency` (default 2) bounds how many jobs run at once, but says
/// nothing about how much memory each one costs, and the container gets 1 GiB total. A
/// single 8000x8000 original measures ~516 MiB peak even with the decode correctly scoped
/// — two of those overlapping is 1032 MiB and another OOM kill, from nothing more exotic
/// than two guests uploading big photos at the same moment.
///
/// 150 MiB sits far above a normal phone photo (a 12 MP JPEG costs ~50 MiB all-in) so the
/// common path never serialises, and far below the point where two jobs stop fitting.
/// Throughput is unaffected for everything except the rare giant, which is exactly the
/// case that must not run in parallel with another giant.
const HEAVY_IMAGE_BYTES: u64 = 150 * 1024 * 1024;
/// Wall-clock ceiling for one oxipng run.
///
/// Bounds TIME, NOT MEMORY — oxipng checks the deadline between trials, so a single trial
/// still allocates in full. The pixel gate above and the sequential build (see
/// `default-features = false` in Cargo.toml) are what bound memory. Do not treat this
/// constant as the OOM fix.
const OXIPNG_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(20);
/// Decode the image ONCE and emit both derivatives — the 800px `preview` (phone feed)
/// and the 2048px `display` (diashow). Returns `(preview_rel, display_rel)`.
async fn generate_image_derivatives(
@@ -274,53 +340,28 @@ impl CompressionWorker {
let display_path = displays_dir.join(&filename);
let original = original.to_path_buf();
let mime_owned = mime_type.to_string();
let preview_max = Self::PREVIEW_MAX_EDGE;
let display_max = Self::DISPLAY_MAX_EDGE;
// Estimate the peak from the HEADER (no pixels decoded — the same kind of cheap probe
// the upload handler already does via `exceeds_decode_budget`) and, if this job is a
// giant, take the exclusive permit so it cannot overlap another giant. Held for the
// whole blocking section, released on drop including on error.
let estimate =
crate::services::imaging::estimated_processing_peak_bytes(&original, Self::DISPLAY_MAX_EDGE);
let _heavy_permit = match estimate {
Some(bytes) if bytes > Self::HEAVY_IMAGE_BYTES => {
tracing::debug!(
%upload_id,
estimated_mib = bytes / (1024 * 1024),
"waiting for the heavy-image permit"
);
Some(self.heavy.acquire().await)
}
_ => None,
};
// Run blocking image operations in a spawn_blocking task
tokio::task::spawn_blocking(move || -> Result<()> {
// Decompression-bomb limits + EXIF orientation, both in one place — see
// services::imaging for why neither may be skipped.
let img = crate::services::imaging::decode_oriented(&original)?;
// Preview: max 800px, preserving aspect ratio (data-saver feed).
img.resize(
preview_max,
preview_max,
image::imageops::FilterType::Lanczos3,
)
.save_with_format(&preview_path, image::ImageFormat::Jpeg)
.context("failed to save preview")?;
// Display: max 2048px for the diashow. Only DOWNSCALE — never upscale a smaller
// original (that adds bytes with no quality gain); re-encode it as JPEG as-is.
let display = if img.width() > display_max || img.height() > display_max {
img.resize(
display_max,
display_max,
image::imageops::FilterType::Lanczos3,
)
} else {
img
};
display
.save_with_format(&display_path, image::ImageFormat::Jpeg)
.context("failed to save display")?;
// If the original is PNG, try lossless compression in-place
if mime_owned == "image/png" {
let opts = oxipng::Options::from_preset(2);
let _ = oxipng::optimize(
&oxipng::InFile::Path(original),
&oxipng::OutFile::Path {
path: None,
preserve_attrs: true,
},
&opts,
);
}
Ok(())
tokio::task::spawn_blocking(move || {
write_image_derivatives(upload_id, &original, &mime_owned, &preview_path, &display_path)
})
.await??;
@@ -341,17 +382,29 @@ impl CompressionWorker {
///
/// Unlike the failure path in `process`, a backfill error is logged and skipped — it must
/// NEVER destroy or soft-delete an upload that already has a working preview.
///
/// Bounded in three ways, all of them load-bearing on a box that restarts itself:
/// `derivative_attempts` stops a fatal row being replayed on every boot, `BACKFILL_BATCH`
/// stops one start queueing unbounded work, and the whole thing runs as ONE task walking
/// the rows sequentially rather than N tasks racing for the same semaphore.
pub async fn backfill_stale_derivatives(&self) {
// `original_path IS NOT NULL` was dead — the column is NOT NULL. What actually needs
// excluding is the blanked path `cleanup_deleted_media` leaves behind.
let rows = sqlx::query_as::<_, (Uuid, String, String)>(
"SELECT id, original_path, mime_type FROM upload
WHERE deleted_at IS NULL AND mime_type LIKE 'image/%'
AND original_path IS NOT NULL
AND original_path <> ''
AND derivative_attempts < $2
AND (
(display_path IS NULL AND preview_path IS NOT NULL)
OR derivatives_rev < $1
)",
)
ORDER BY created_at DESC
LIMIT $3",
)
.bind(Self::DERIVATIVES_REV)
.bind(Self::MAX_DERIVATIVE_ATTEMPTS)
.bind(Self::BACKFILL_BATCH)
.fetch_all(&self.pool)
.await;
let rows = match rows {
@@ -361,14 +414,32 @@ impl CompressionWorker {
return;
}
};
self.report_exhausted_derivatives().await;
if rows.is_empty() {
return;
}
tracing::info!("regenerating derivatives for {} upload(s)", rows.len());
for (id, original_path, mime_type) in rows {
let worker = self.clone();
tokio::spawn(async move {
// ONE task for the whole batch. The previous shape spawned a task per row, so a large
// backlog created thousands of live tasks that each held a pool handle and queued on
// the same two semaphore permits, competing with live uploads for the entire boot.
let worker = self.clone();
tokio::spawn(async move {
for (id, original_path, mime_type) in rows {
let _permit = worker.semaphore.acquire().await;
// Write-ahead, exactly as in the live path: if this row is the one that kills
// the process, this increment is the only thing that outlives the SIGKILL.
match Upload::begin_derivative_attempt(&worker.pool, id).await {
Ok(Some(n)) if n > Self::MAX_DERIVATIVE_ATTEMPTS => continue,
Ok(Some(_)) => {}
Ok(None) => continue,
Err(e) => {
tracing::warn!(error = ?e, %id, "could not record a backfill attempt; skipping");
continue;
}
}
let original = worker.media_path.join(&original_path);
match worker
.generate_image_derivatives(id, &original, &mime_type)
@@ -377,6 +448,8 @@ impl CompressionWorker {
Ok((preview_rel, display_rel)) => {
let _ = Upload::set_preview_path(&worker.pool, id, &preview_rel).await;
let _ = Upload::set_display_path(&worker.pool, id, &display_rel).await;
// Clears derivative_attempts too, so a row that failed transiently is
// not one boot closer to being abandoned.
let _ =
Upload::set_derivatives_rev(&worker.pool, id, Self::DERIVATIVES_REV)
.await;
@@ -384,12 +457,133 @@ impl CompressionWorker {
}
Err(e) => {
// Leave the existing derivatives and the original intact; this row is
// simply retried on the next start. The rev stays behind, which is the
// marker that it still needs doing.
// retried on the next start until its attempt budget runs out. The rev
// stays behind, which is the marker that it still needs doing.
tracing::warn!(error = ?e, %id, "derivative backfill failed; leaving as-is");
let _ =
Upload::record_derivative_failure(&worker.pool, id, &format!("{e:#}"))
.await;
}
}
});
}
});
}
/// Re-extract poster frames for videos that never got one.
///
/// A video interrupted by a restart is stranded: `startup_recovery` flips its
/// `compression_status` from `processing` to `failed` and nothing re-enqueues it, so
/// `thumbnail_path` stays NULL forever while the clip itself plays fine. The feed shows a
/// posterless tile for the rest of the event, and after
/// `FAILED_ORIGINAL_RETENTION_DAYS` the reclaim sweep is entitled to the original.
///
/// Shares `derivative_attempts` with the image backfill on purpose. Note the consequence,
/// which is intended rather than a bug to fix later: `extract_poster_frame` returning
/// `Ok(false)` is a NORMAL, permanent outcome for a sub-second clip (Live Photos,
/// mis-taps), and since the counter is write-ahead and only cleared by a real success,
/// those clips stop being re-ffmpeg'd on every boot once the budget is spent.
pub async fn backfill_video_posters(&self) {
let rows = sqlx::query_as::<_, (Uuid, String)>(
"SELECT id, original_path FROM upload
WHERE deleted_at IS NULL AND mime_type LIKE 'video/%'
AND thumbnail_path IS NULL
AND original_path <> ''
AND derivative_attempts < $1
ORDER BY created_at DESC
LIMIT $2",
)
.bind(Self::MAX_DERIVATIVE_ATTEMPTS)
.bind(Self::BACKFILL_BATCH)
.fetch_all(&self.pool)
.await;
let rows = match rows {
Ok(r) => r,
Err(e) => {
tracing::warn!(error = ?e, "video poster backfill query failed");
return;
}
};
if rows.is_empty() {
return;
}
tracing::info!("re-extracting posters for {} video(s)", rows.len());
let worker = self.clone();
tokio::spawn(async move {
for (id, original_path) in rows {
let _permit = worker.semaphore.acquire().await;
match Upload::begin_derivative_attempt(&worker.pool, id).await {
Ok(Some(n)) if n > Self::MAX_DERIVATIVE_ATTEMPTS => continue,
Ok(Some(_)) => {}
Ok(None) => continue,
Err(e) => {
tracing::warn!(error = ?e, %id, "could not record a poster attempt; skipping");
continue;
}
}
let original = worker.media_path.join(&original_path);
match worker.generate_video_thumbnail(id, &original).await {
Ok(Some(thumb_rel)) => {
if Upload::set_thumbnail_path(&worker.pool, id, &thumb_rel)
.await
.is_ok()
{
// Clears the attempt counter: a video that eventually succeeded
// must not carry a budget scar into a future pipeline revision.
let _ = Upload::set_derivatives_rev(
&worker.pool,
id,
Self::DERIVATIVES_REV,
)
.await;
tracing::info!("poster regenerated for upload {id}");
}
}
// No frame at all — normal for a very short clip. The tile stays
// posterless and the attempt is spent, which is what stops the retry.
Ok(None) => {
tracing::debug!(%id, "still no poster frame; leaving the tile as-is");
}
Err(e) => {
tracing::warn!(error = ?e, %id, "poster backfill failed; leaving as-is");
let _ =
Upload::record_derivative_failure(&worker.pool, id, &format!("{e:#}"))
.await;
}
}
}
});
}
/// Say out loud, once per boot, that some uploads have stopped being retried.
///
/// Without this the give-up is invisible: the loop stops (which is the point) but the
/// affected photos keep a stale or missing derivative forever with nothing to notice. The
/// originals are untouched, so this is recoverable once the cause is fixed — reset
/// `derivative_attempts` to 0 and restart.
async fn report_exhausted_derivatives(&self) {
let exhausted: Result<i64, _> = sqlx::query_scalar(
"SELECT count(*) FROM upload
WHERE deleted_at IS NULL AND mime_type LIKE 'image/%'
AND derivative_attempts >= $2
AND (
(display_path IS NULL AND preview_path IS NOT NULL)
OR derivatives_rev < $1
)",
)
.bind(Self::DERIVATIVES_REV)
.bind(Self::MAX_DERIVATIVE_ATTEMPTS)
.fetch_one(&self.pool)
.await;
if let Ok(count) = exhausted
&& count > 0
{
tracing::error!(
count,
"{count} upload(s) exhausted derivative regeneration and will no longer be \
retried; their originals are intact — see upload.derivative_last_error, fix \
the cause, then reset derivative_attempts to 0 and restart"
);
}
}
@@ -413,3 +607,250 @@ impl CompressionWorker {
Ok(produced.then(|| format!("thumbnails/{thumb_filename}")))
}
}
/// The blocking half of [`CompressionWorker::generate_image_derivatives`]: decode once, write
/// both derivatives, then optionally shrink a PNG original in place.
///
/// A free function rather than an inline closure so its memory behaviour is directly testable —
/// this is the code path that OOM-killed the container, and the fix is a scoping property that a
/// future edit could silently undo.
fn write_image_derivatives(
upload_id: Uuid,
original: &Path,
mime_type: &str,
preview_path: &Path,
display_path: &Path,
) -> Result<()> {
let preview_max = CompressionWorker::PREVIEW_MAX_EDGE;
let display_max = CompressionWorker::DISPLAY_MAX_EDGE;
// THE FULL-SIZE DECODE IS SCOPED TO THIS BLOCK ON PURPOSE, and the block yields the
// DISPLAY derivative rather than the original.
//
// `img` is up to 256 MiB (imaging::decode_limits max_alloc) and `resize` only BORROWS it,
// so it used to stay alive through both resizes AND the oxipng call below — which decodes
// the PNG a second time and holds a full-size buffer per filter trial. That measured
// ~1250 MiB of peak RSS for a 2.8 MiB input, inside a 1 GiB cgroup: the container was
// SIGKILLed, taking every SSE stream and every in-flight upload with it.
//
// A block rather than a bare `drop(img)` because a `drop` call is one careless edit away
// from being removed as redundant-looking — and note the `else` arm MOVES `img` out, which
// is what makes "the block's value is the only survivor" true in both arms.
let (display, width, height) = {
// Decompression-bomb limits + EXIF orientation, both in one place — see
// services::imaging for why neither may be skipped.
let img = crate::services::imaging::decode_oriented(original)?;
let (width, height) = (img.width(), img.height());
// Display: max 2048px for the diashow. Only DOWNSCALE — never upscale a smaller
// original (that adds bytes with no quality gain); re-encode it as JPEG as-is.
let display = if width > display_max || height > display_max {
img.resize(
display_max,
display_max,
image::imageops::FilterType::Lanczos3,
)
} else {
img
};
(display, width, height)
};
display
.save_with_format(display_path, image::ImageFormat::Jpeg)
.context("failed to save display")?;
// Preview: max 800px, derived from the DISPLAY, not from the original.
//
// Both derivatives used to resize the full-size decode independently, so a 8000x8000
// original paid for two full-size Lanczos passes and their intermediates — measured 520
// MiB peak even after the scoping fix above, which two concurrent workers cannot fit in a
// 1 GiB container. Chaining 8000 -> 2048 -> 800 makes the second pass operate on 2048px
// input, and the full-size buffer is already freed by the time it runs. Quality is not the
// trade-off here: a staged Lanczos3 downscale to 800px is visually indistinguishable from
// a single-step one (and is a standard technique for large ratios).
display
.resize(
preview_max,
preview_max,
image::imageops::FilterType::Lanczos3,
)
.save_with_format(preview_path, image::ImageFormat::Jpeg)
.context("failed to save preview")?;
drop(display);
let pixels = u64::from(width) * u64::from(height);
// If the original is PNG, try lossless compression in place — but only when its pixel count
// is inside the budget, and never for longer than OXIPNG_TIMEOUT. This is a best-effort size
// saving: declining it costs disk, while attempting it unbounded cost the whole container.
if mime_type == "image/png" {
if pixels <= CompressionWorker::OXIPNG_MAX_PIXELS {
let mut opts = oxipng::Options::from_preset(2);
opts.timeout = Some(CompressionWorker::OXIPNG_TIMEOUT);
let _ = oxipng::optimize(
&oxipng::InFile::Path(original.to_path_buf()),
&oxipng::OutFile::Path {
path: None,
preserve_attrs: true,
},
&opts,
);
} else {
tracing::info!(
%upload_id, pixels,
"skipping oxipng: above the pixel budget; the original is stored as uploaded"
);
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
/// Peak resident set of THIS process, in bytes, from `/proc/self/status`.
fn peak_rss_bytes() -> u64 {
let status = std::fs::read_to_string("/proc/self/status").expect("procfs");
let line = status
.lines()
.find(|l| l.starts_with("VmHWM:"))
.expect("VmHWM");
let kb: u64 = line
.split_whitespace()
.nth(1)
.and_then(|v| v.parse().ok())
.expect("VmHWM value");
kb * 1024
}
/// Reset the kernel's peak-RSS watermark so the measurement covers only what follows.
/// Linux 4.0+; writing "5" to `clear_refs` resets `VmHWM` to the current RSS.
fn reset_peak_rss() {
let _ = std::fs::write("/proc/self/clear_refs", "5");
}
/// The pixel gate has to sit below what the axis limits allow, or it can never fire.
#[test]
fn the_oxipng_gate_is_reachable_within_the_decode_limits() {
const _: () = {
// imaging::decode_limits permits 12_000 x 12_000 = 144 MP. A gate above that would
// never skip anything.
assert!(CompressionWorker::OXIPNG_MAX_PIXELS < 12_000 * 12_000);
// ...and it must stay above a 48 MP camera, so real photos still get optimised.
assert!(CompressionWorker::OXIPNG_MAX_PIXELS >= 8_000_000);
};
}
/// The heavy-image gate has to classify the two cases the way the sizing assumed:
/// an ordinary phone photo must NOT serialise, and the giant must.
#[test]
fn the_heavy_gate_separates_a_phone_photo_from_a_giant() {
let dir = std::env::temp_dir().join(format!("es-heavy-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
// 12 MP, the shape of a default phone capture.
let ordinary = dir.join("ordinary.jpg");
image::RgbImage::new(4032, 3024).save(&ordinary).unwrap();
let ordinary_peak = crate::services::imaging::estimated_processing_peak_bytes(
&ordinary,
CompressionWorker::DISPLAY_MAX_EDGE,
)
.expect("header readable");
assert!(
ordinary_peak <= CompressionWorker::HEAVY_IMAGE_BYTES,
"a 12 MP photo estimated at {} MiB would serialise the common path",
ordinary_peak / 1048576
);
// The 64 MP RGBA case that measured ~516 MiB peak.
let giant = dir.join("giant.png");
image::RgbaImage::new(8000, 8000).save(&giant).unwrap();
let giant_peak = crate::services::imaging::estimated_processing_peak_bytes(
&giant,
CompressionWorker::DISPLAY_MAX_EDGE,
)
.expect("header readable");
assert!(
giant_peak > CompressionWorker::HEAVY_IMAGE_BYTES,
"an 8000x8000 RGBA original estimated at only {} MiB would be allowed to run \
concurrently with another one — 2x its real ~516 MiB peak does not fit in 1 GiB",
giant_peak / 1048576
);
// The estimate must also be in the right ballpark, not merely on the right side of the
// threshold: 244 MiB decode + 262 MiB f32 resize intermediate.
assert!(
(400..700).contains(&(giant_peak / 1048576)),
"estimate {} MiB is far from the measured ~516 MiB peak",
giant_peak / 1048576
);
let _ = std::fs::remove_dir_all(&dir);
}
/// The OOM that took the container down, measured rather than argued.
///
/// An 8000x8000 RGBA PNG passes admission: 256,000,000 bytes is just under the 256 MiB
/// `max_alloc`, and smooth content is a few MB on disk, far under any size cap. The old
/// code kept that ~244 MiB decode alive across an unbounded, multi-threaded oxipng run and
/// peaked at ~1250 MiB — inside a 1 GiB cgroup. Being SIGKILLed there is not a blip: the
/// row was already committed, so the boot backfill replayed the identical workload on every
/// restart.
///
/// `#[ignore]` because it allocates ~250 MiB and takes a few seconds. Run explicitly:
/// cargo test --release oom -- --ignored --nocapture --test-threads=1
/// It must run ALONE — `VmHWM` is per process, so a concurrent test would pollute it.
#[test]
#[ignore = "heavy: allocates ~250 MiB; run with --ignored --test-threads=1"]
fn a_large_png_stays_far_below_the_container_limit() {
const EDGE: u32 = 8_000;
let dir = std::env::temp_dir().join(format!("es-oom-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let original = dir.join("big.png");
// Smooth gradient: ~244 MiB decoded, a couple of MB on disk. That gap is the whole
// point — file size tells you nothing about what a PNG costs to process.
{
let mut buf = image::RgbaImage::new(EDGE, EDGE);
for (x, y, px) in buf.enumerate_pixels_mut() {
*px = image::Rgba([(x >> 5) as u8, (y >> 5) as u8, ((x + y) >> 6) as u8, 255]);
}
buf.save(&original).unwrap();
}
// Everything above is fixture setup, not the code under test.
reset_peak_rss();
let before = peak_rss_bytes();
write_image_derivatives(
Uuid::new_v4(),
&original,
"image/png",
&dir.join("preview.jpg"),
&dir.join("display.jpg"),
)
.expect("derivatives");
let peak = peak_rss_bytes();
let on_disk = std::fs::metadata(&original).unwrap().len();
eprintln!(
"input {:.2} MiB on disk ({EDGE}x{EDGE}); peak RSS {:.0} MiB (was {:.0} MiB before)",
on_disk as f64 / 1048576.0,
peak as f64 / 1048576.0,
before as f64 / 1048576.0
);
assert!(dir.join("preview.jpg").exists() && dir.join("display.jpg").exists());
// The container gets 1 GiB and runs two of these concurrently. 600 MiB is a generous
// ceiling that the old code (~1250 MiB) could not have met.
assert!(
peak < 600 * 1024 * 1024,
"peak RSS {} MiB — the decode is being held across oxipng again, or the pixel \
gate stopped firing",
peak / 1048576
);
let _ = std::fs::remove_dir_all(&dir);
}
}

View File

@@ -118,6 +118,40 @@ fn decoder_within_budget(path: &Path) -> Result<impl image::ImageDecoder> {
Ok(decoder)
}
/// Rough peak heap an image will cost to turn into derivatives, read from the HEADER only —
/// no pixels are decoded. `None` when the header can't be read or the image is over budget
/// (the caller is about to fail on it anyway).
///
/// Two terms, and the second is the one that surprises:
///
/// - the decoded buffer, `width * height * channels`; and
/// - the resize intermediate. `image`'s Lanczos3 path accumulates in `f32`, so the buffer
/// between the horizontal and vertical passes is `new_width * old_height * 4 channels * 4
/// bytes` — 16 bytes per pixel-row-slot, not the 4 the output uses. For an 8000x8000
/// original that is 262 MiB on top of a 244 MiB decode, measured. It is bigger than the
/// decode for any tall image, which is why "the decode is bounded by max_alloc" was never
/// the whole story.
///
/// Used to decide whether an image is heavy enough to need exclusive use of the box's memory
/// headroom, NOT to reject anything.
pub fn estimated_processing_peak_bytes(path: &Path, display_edge: u32) -> Option<u64> {
let decoder = decoder_within_budget(path).ok()?;
let (width, height) = decoder.dimensions();
let decoded = decoder.total_bytes();
// Aspect-preserving fit into `display_edge`, matching DynamicImage::resize. No downscale
// means no intermediate at all.
let intermediate = if width > display_edge || height > display_edge {
let ratio = f64::from(display_edge) / f64::from(width.max(height));
let new_width = (f64::from(width) * ratio).round().max(1.0) as u64;
new_width * u64::from(height) * 16
} else {
0
};
Some(decoded.saturating_add(intermediate))
}
/// Megapixels an image would decode to, or `None` if its header can't be read. Used only
/// to put a concrete number in the message the guest sees.
pub fn megapixels(path: &Path) -> Option<f64> {

View File

@@ -120,6 +120,19 @@ pub async fn startup_recovery(pool: &PgPool) {
}
}
/// How long a file in `originals/` may exist without a database row before it is treated as
/// abandoned.
///
/// This window is the ONLY thing making the sweep safe, because the upload handler renames the
/// temp file into its final path BEFORE committing the row: for a short moment a perfectly
/// healthy upload legitimately looks exactly like an orphan. Six hours is far beyond any live
/// request (a 576 MiB body over a bad venue uplink is minutes, and the request itself is bounded
/// by the reverse proxy) while still reclaiming the leak inside a single event.
///
/// DO NOT SHORTEN THIS to make a test faster — a value below the longest possible in-flight
/// upload deletes photos out from under the request that is committing them.
const ORPHAN_UPLOAD_RETENTION_HOURS: u64 = 6;
/// Spawns a background task that periodically:
/// - deletes session rows whose `expires_at` is more than a day in the past
/// - prunes the in-memory rate-limiter HashMap of empty windows
@@ -176,6 +189,13 @@ async fn periodic_loop(
cleanup_sessions(&pool).await;
cleanup_deleted_media(&pool, &media_path).await;
sweep_orphan_upload_temps(&media_path).await;
// Runs AFTER the .tmp sweep, and covers the class that one structurally cannot see:
// an original that was renamed to its final name but whose transaction never
// committed. Those have no row, so `cleanup_deleted_media` (row-driven) can never
// find them, and `sweep_orphan_upload_temps` skips them because they no longer end
// in `.tmp` — they were permanently unowned, silently shrinking the free disk that
// `compute_storage_quota` divides among guests.
sweep_orphan_originals(&pool, &media_path).await;
rate_limiter.prune();
sse_tickets.prune();
}
@@ -379,6 +399,126 @@ async fn cleanup_sessions(pool: &PgPool) {
}
}
/// Reclaim files in `originals/` that no upload row references.
///
/// The backstop behind [`TempFileGuard`](crate::handlers::upload). The guard covers the
/// process that is running; this covers the process that was killed — a SIGKILL, an OOM, or a
/// power cut leaves whatever bytes had been written with no `Drop` to reclaim them, and those
/// files are then permanently invisible: they have no row, so `cleanup_deleted_media` (which is
/// row-driven) can never see them, and they are not counted against any quota while still
/// consuming the free disk that `compute_storage_quota` divides among guests. On a single box
/// where all three volumes share a filesystem, that ends with Postgres unable to write WAL.
///
/// Two classes:
/// - `*.tmp` — an upload that never got as far as being renamed. Always safe past the window.
/// - everything else — a final-named original whose commit never happened.
async fn sweep_orphan_originals(pool: &PgPool, media_path: &std::path::Path) {
let originals = media_path.join("originals");
let cutoff = Duration::from_secs(ORPHAN_UPLOAD_RETENTION_HOURS * 3600);
// originals/{event_slug}/{uuid}.{ext} — one level of per-event directories.
let mut event_dirs = match tokio::fs::read_dir(&originals).await {
Ok(rd) => rd,
// Nothing uploaded yet; the directory is created lazily by the upload handler.
Err(_) => return,
};
let mut candidates: Vec<(String, std::path::PathBuf)> = Vec::new();
let mut temps_removed = 0u32;
while let Ok(Some(event_dir)) = event_dirs.next_entry().await {
if !event_dir
.file_type()
.await
.map(|t| t.is_dir())
.unwrap_or(false)
{
continue;
}
let slug = event_dir.file_name().to_string_lossy().to_string();
let Ok(mut files) = tokio::fs::read_dir(event_dir.path()).await else {
continue;
};
while let Ok(Some(entry)) = files.next_entry().await {
let Ok(meta) = entry.metadata().await else {
continue;
};
if !meta.is_file() {
continue;
}
// Too young to judge: an upload committing RIGHT NOW is indistinguishable from an
// orphan, because the rename precedes the commit.
let recent = meta
.modified()
.ok()
.and_then(|m| m.elapsed().ok())
.is_none_or(|age| age < cutoff);
if recent {
continue;
}
let name = entry.file_name().to_string_lossy().to_string();
if name.ends_with(".tmp") {
// A `.tmp` never has a row by construction — no DB check needed.
if tokio::fs::remove_file(entry.path()).await.is_ok() {
temps_removed += 1;
}
continue;
}
candidates.push((format!("originals/{slug}/{name}"), entry.path()));
}
}
if temps_removed > 0 {
tracing::warn!(
"reclaimed {temps_removed} abandoned upload temp file(s) older than \
{ORPHAN_UPLOAD_RETENTION_HOURS}h"
);
}
if candidates.is_empty() {
return;
}
// One query per batch, not one per file: a backlog of thousands of orphans must not turn
// into thousands of round trips on an hourly timer.
let mut orphans_removed = 0u32;
for chunk in candidates.chunks(500) {
let paths: Vec<String> = chunk.iter().map(|(rel, _)| rel.clone()).collect();
// NO `deleted_at IS NULL` FILTER HERE. A soft-deleted row still points at its file
// during its retention window, and reclaiming that file is `cleanup_deleted_media`'s
// job — filtering here would race the two sweeps and destroy the exact files the
// recovery window exists to preserve.
let unreferenced: Result<Vec<(String,)>, _> = sqlx::query_as(
"SELECT p FROM unnest($1::text[]) AS p
WHERE NOT EXISTS (SELECT 1 FROM upload u WHERE u.original_path = p)",
)
.bind(&paths)
.fetch_all(pool)
.await;
let unreferenced = match unreferenced {
Ok(rows) => rows,
Err(e) => {
tracing::warn!(error = ?e, "orphan-original sweep query failed");
return;
}
};
for (rel,) in unreferenced {
if let Some((_, abs)) = chunk.iter().find(|(r, _)| *r == rel)
&& tokio::fs::remove_file(abs).await.is_ok()
{
orphans_removed += 1;
}
}
}
if orphans_removed > 0 {
tracing::warn!(
"reclaimed {orphans_removed} original(s) with no upload row, older than \
{ORPHAN_UPLOAD_RETENTION_HOURS}h"
);
}
}
#[cfg(test)]
mod tests {
use super::*;

View File

@@ -13,6 +13,16 @@ use rand::Rng;
/// stream open. Tickets are consumed on use and expire after `TTL`.
const TTL: Duration = Duration::from_secs(30);
/// Ceiling on outstanding tickets across the whole process.
///
/// Not really about the bytes (~120 each) — about `issue` having had no bound of any kind.
/// Sized well above a real event: ~1000 concurrent clients each holding one live 30 s ticket.
const MAX_TICKETS: usize = 4096;
/// Live tickets one session may hold. Above 1 because two tabs sharing a token open their
/// EventSources concurrently; 4 absorbs that without letting a reconnect loop accumulate.
const MAX_TICKETS_PER_SESSION: usize = 4;
#[derive(Clone)]
pub struct SseTicketStore {
inner: Arc<Mutex<HashMap<String, Entry>>>,
@@ -39,9 +49,47 @@ impl SseTicketStore {
}
/// Mint a new ticket bound to the caller's session (identified by token hash).
pub fn issue(&self, token_hash: String) -> String {
///
/// `None` when the store is at capacity — the caller should answer 503, not evict.
///
/// Three bounds, because `issue` had none: no size cap, no per-caller cap, and no rate
/// limit on the endpoint, while `prune` ran only hourly against a 30-second TTL. So any
/// authenticated session could loop the endpoint and grow the map for an hour.
pub fn issue(&self, token_hash: String) -> Option<String> {
let ticket = random_ticket();
let mut map = self.inner.lock().unwrap();
// Prune on issue rather than only hourly. This alone changes the bound from "tickets
// minted since the last maintenance tick" to "tickets live at once", which is what the
// 30 s TTL was always meant to express.
map.retain(|_, e| e.issued_at.elapsed() <= TTL);
// Cap the caller's own outstanding tickets, evicting their oldest. NOT one-per-session:
// two tabs sharing a token open their EventSources concurrently, and having tab B
// invalidate tab A's unconsumed ticket looks exactly like a flaky SSE connection.
let mut mine: Vec<(String, Instant)> = map
.iter()
.filter(|(_, e)| e.token_hash == token_hash)
.map(|(k, e)| (k.clone(), e.issued_at))
.collect();
if mine.len() >= MAX_TICKETS_PER_SESSION {
mine.sort_by_key(|(_, issued)| *issued);
for (key, _) in mine.iter().take(mine.len() - MAX_TICKETS_PER_SESSION + 1) {
map.remove(key);
}
}
// At capacity, REFUSE — never evict a stranger's ticket. Evicting would let one
// misbehaving client deny SSE to the whole venue, which is worse than failing the
// request that hit the ceiling.
if map.len() >= MAX_TICKETS {
tracing::warn!(
outstanding = map.len(),
"SSE ticket store at capacity; refusing to mint"
);
return None;
}
map.insert(
ticket.clone(),
Entry {
@@ -49,7 +97,7 @@ impl SseTicketStore {
issued_at: Instant::now(),
},
);
ticket
Some(ticket)
}
/// Consume a ticket. Returns `Some(token_hash)` if the ticket exists and is
@@ -84,10 +132,16 @@ fn random_ticket() -> String {
mod tests {
use super::*;
/// `issue` now returns `Option`; in every test below the store is far from capacity, so an
/// `expect` here documents that refusing is exceptional rather than routine.
fn issue(store: &SseTicketStore, hash: &str) -> String {
store.issue(hash.into()).expect("store has capacity")
}
#[test]
fn issue_then_consume_returns_the_hash_exactly_once() {
let store = SseTicketStore::new();
let ticket = store.issue("hash-1".into());
let ticket = issue(&store, "hash-1");
assert_eq!(store.consume(&ticket).as_deref(), Some("hash-1"));
// Single-use: a replay of the same ticket is rejected.
assert_eq!(
@@ -106,8 +160,8 @@ mod tests {
#[test]
fn issued_tickets_are_unique_and_hex() {
let store = SseTicketStore::new();
let a = store.issue("h".into());
let b = store.issue("h".into());
let a = issue(&store, "h");
let b = issue(&store, "h");
assert_ne!(a, b, "each ticket must be unique");
assert_eq!(a.len(), 48, "24 random bytes → 48 hex chars");
assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
@@ -116,29 +170,104 @@ mod tests {
#[test]
fn fresh_ticket_survives_prune() {
let store = SseTicketStore::new();
let ticket = store.issue("h".into());
let ticket = issue(&store, "h");
store.prune(); // not expired → kept
assert_eq!(store.consume(&ticket).as_deref(), Some("h"));
}
#[test]
fn expired_ticket_consumes_to_none() {
// Construct an entry that is already past the TTL and confirm consume() rejects it.
let store = SseTicketStore::new();
let stale = "stale-ticket".to_string();
/// Build an entry that is already past the TTL.
fn insert_stale(store: &SseTicketStore, key: &str, token_hash: &str) {
store.inner.lock().unwrap().insert(
stale.clone(),
key.to_string(),
Entry {
token_hash: "h".into(),
token_hash: token_hash.into(),
issued_at: Instant::now()
.checked_sub(TTL + Duration::from_secs(1))
.expect("host uptime should exceed the ticket TTL"),
},
);
}
#[test]
fn expired_ticket_consumes_to_none() {
let store = SseTicketStore::new();
insert_stale(&store, "stale-ticket", "h");
assert_eq!(
store.consume(&stale),
store.consume("stale-ticket"),
None,
"an expired ticket must not authenticate"
);
}
/// The TTL is 30 s but `prune` only ran hourly, so the map was really bounded by "tickets
/// minted in the last hour" — which is unbounded for a client in a loop.
#[test]
fn issuing_prunes_expired_entries() {
let store = SseTicketStore::new();
insert_stale(&store, "stale-a", "someone-else");
insert_stale(&store, "stale-b", "someone-else");
issue(&store, "h");
assert_eq!(
store.inner.lock().unwrap().len(),
1,
"issue must reclaim expired slots, not merely add to them"
);
}
/// Two tabs sharing a token is normal, so the per-session cap must be above 1 — but a
/// reconnect loop must not accumulate. The caller's OWN oldest is what gets evicted.
#[test]
fn a_session_is_capped_and_evicts_only_its_own_oldest() {
let store = SseTicketStore::new();
let stranger = issue(&store, "other-session");
let mut mine: Vec<String> = Vec::new();
for _ in 0..MAX_TICKETS_PER_SESSION + 2 {
mine.push(issue(&store, "mine"));
}
let live = mine
.iter()
.filter(|t| store.inner.lock().unwrap().contains_key(*t))
.count();
assert_eq!(live, MAX_TICKETS_PER_SESSION, "one session, bounded");
assert!(
store.inner.lock().unwrap().contains_key(&mine[mine.len() - 1]),
"the newest ticket is the one the caller is about to use"
);
assert_eq!(
store.consume(&stranger).as_deref(),
Some("other-session"),
"another session's ticket must survive — evicting it would let one client deny \
SSE to the venue"
);
}
/// At capacity the store REFUSES rather than evicting a stranger. Refusing fails the one
/// request that hit the ceiling; evicting would break an unrelated client's live stream.
#[test]
fn at_capacity_the_store_refuses_instead_of_evicting() {
let store = SseTicketStore::new();
{
let mut map = store.inner.lock().unwrap();
for i in 0..MAX_TICKETS {
map.insert(
format!("filler-{i}"),
Entry {
token_hash: format!("session-{i}"),
issued_at: Instant::now(),
},
);
}
}
assert_eq!(
store.issue("newcomer".into()),
None,
"a full store must refuse, so the caller can answer 503"
);
assert!(
store.inner.lock().unwrap().contains_key("filler-0"),
"no existing ticket may be sacrificed to make room"
);
}
}

View File

@@ -76,7 +76,7 @@ pub async fn extract_poster_frame(src: &Path, dest: &Path, width: u32) -> Result
/// Run one ffmpeg attempt. A non-zero exit is NOT an error here — the artifact check above is the
/// authority, and a corrupt input that fails at 1 s may still yield a frame at 0.
async fn run_ffmpeg(src: &Path, dest: &Path, width: u32, seek: &str) -> Result<()> {
let mut child = tokio::process::Command::new("ffmpeg")
let child = tokio::process::Command::new("ffmpeg")
.args([
// BEFORE -i: an input-side seek. See SEEK_POSITIONS.
"-ss",
@@ -90,24 +90,56 @@ async fn run_ffmpeg(src: &Path, dest: &Path, width: u32, seek: &str) -> Result<(
"-y",
dest.to_str().unwrap_or_default(),
])
.stdout(std::process::Stdio::piped())
// ffmpeg writes the poster to `dest` itself; nothing here ever reads stdout, so
// giving it a pipe only created something that could fill.
.stdout(std::process::Stdio::null())
// stderr IS piped — it is the only diagnostic when a clip yields no frame — but it
// must be DRAINED, which is the whole point of `wait_with_output` below.
.stderr(std::process::Stdio::piped())
.kill_on_drop(true)
.spawn()
.context("failed to spawn ffmpeg")?;
match tokio::time::timeout(FFMPEG_TIMEOUT, child.wait()).await {
Ok(res) => {
res.context("ffmpeg wait failed")?;
}
Err(_) => {
let _ = child.kill().await;
anyhow::bail!("ffmpeg timed out after {}s", FFMPEG_TIMEOUT.as_secs());
}
// `wait_with_output`, NOT `wait`. ffmpeg is verbose on stderr (banner, stream info,
// per-frame progress) and `wait()` reads neither pipe — so once the ~64 KiB pipe buffer
// filled, ffmpeg blocked writing, `wait()` never returned, and the call burned the full
// timeout. That is not merely slow: the timeout is an `Err`, so after 2 seek positions x
// 3 compression attempts the caller soft-deletes a perfectly playable video for a
// poster-frame failure. `wait_with_output` polls the pipe and the exit status together.
//
// It also CONSUMES the child, so the explicit `child.kill()` that used to sit on the
// timeout arm cannot exist here — and is not needed: `kill_on_drop(true)` is set above,
// and dropping the future on timeout drops the child with it.
let out = match tokio::time::timeout(FFMPEG_TIMEOUT, child.wait_with_output()).await {
Ok(res) => res.context("ffmpeg wait failed")?,
Err(_) => anyhow::bail!("ffmpeg timed out after {}s", FFMPEG_TIMEOUT.as_secs()),
};
// A non-zero exit is not an error (see the doc comment) — the artifact check in
// `extract_poster_frame` is the authority. Log the tail so a systematically failing
// format is diagnosable without turning it into data loss.
if !out.status.success() {
tracing::debug!(
seek,
status = ?out.status,
stderr = %tail_lines(&out.stderr, 10),
"ffmpeg exited non-zero; the artifact check decides"
);
}
Ok(())
}
/// Last `n` lines of a child's stderr, lossily decoded.
///
/// Bounded on purpose: ffmpeg's stderr is unbounded, and the reason we now drain it is that
/// unbounded output used to be a hazard. Emitting all of it into a log line — into container
/// logs that are themselves size-capped — would just move the problem.
fn tail_lines(bytes: &[u8], n: usize) -> String {
let text = String::from_utf8_lossy(bytes);
let lines: Vec<&str> = text.lines().filter(|l| !l.trim().is_empty()).collect();
lines[lines.len().saturating_sub(n)..].join(" | ")
}
#[cfg(test)]
mod tests {
use super::*;
@@ -167,4 +199,65 @@ mod tests {
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_stderr_tail_is_bounded_and_survives_invalid_utf8() {
let noisy: Vec<u8> = (0..500)
.map(|i| format!("line {i}\n"))
.collect::<String>()
.into_bytes();
let got = tail_lines(&noisy, 3);
assert_eq!(got, "line 497 | line 498 | line 499");
// ffmpeg emits filenames verbatim, so its stderr is not guaranteed to be UTF-8.
assert_eq!(tail_lines(&[b'o', b'k', 0xff], 5), "ok\u{fffd}");
assert_eq!(tail_lines(b"", 5), "");
}
/// A real extraction must finish in a small fraction of `FFMPEG_TIMEOUT`.
///
/// Wall-clock is the ONLY observable of the bug this guards: piping stderr and then
/// calling `wait()` (which drains nothing) blocks ffmpeg on a full pipe buffer until the
/// timeout fires, and the timeout is an `Err`, so the upload is soft-deleted. The
/// assertion is deliberately on elapsed time, not on the exit status.
///
/// Honest limitation: our fixture is quiet enough not to fill a 64 KiB pipe on its own,
/// so this catches a regression to `wait()` only in combination with a verbose input. It
/// is still worth pinning — a reverted drain plus any chatty clip is data loss.
#[tokio::test]
async fn a_real_clip_yields_a_poster_well_inside_the_timeout() {
if tokio::process::Command::new("ffmpeg")
.arg("-version")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status()
.await
.is_err()
{
eprintln!("skipping: ffmpeg not on PATH");
return;
}
let src = Path::new("../e2e/fixtures/media/sample.mp4");
if !src.exists() {
eprintln!("skipping: {} missing", src.display());
return;
}
let dir = std::env::temp_dir().join(format!("es-video-ok-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let dest = dir.join("poster.jpg");
let started = std::time::Instant::now();
let got = extract_poster_frame(src, &dest, 400).await;
let elapsed = started.elapsed();
assert!(matches!(got, Ok(true)), "expected a poster, got {got:?}");
assert!(dest.metadata().unwrap().len() > 0);
assert!(
elapsed < FFMPEG_TIMEOUT / 4,
"extraction took {elapsed:?}; a drained stderr finishes in well under \
{FFMPEG_TIMEOUT:?} — this is the pipe-deadlock regression guard"
);
let _ = std::fs::remove_dir_all(&dir);
}
}