diff --git a/.env.example b/.env.example index e5b347a..41dca19 100644 --- a/.env.example +++ b/.env.example @@ -44,6 +44,13 @@ POSTGRES_DB=eventsnap # OOM in Postgres doesn't degrade one feature, it takes the whole event down. DATABASE_MAX_CONNECTIONS=30 +# Log level. `info` is the right production default: at `debug` the tower-http trace +# layer writes a line per request AND per response, which on a busy event is a large +# multiple of the useful output. Container logs are capped at 10m x 3 per service +# (docker-compose.yml), so a chatty level buys you a shorter history, not more of it. +# To debug a live event: RUST_LOG=eventsnap_backend=debug docker compose up -d app +RUST_LOG=info + # ── Authentication ──────────────────────────────────────────────────────────── # Generate with: openssl rand -hex 64 JWT_SECRET=change_me_to_a_random_64_byte_hex_string diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 8396b9d..cb4cbea 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -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" diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 654a1ad..d43d288 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -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" diff --git a/backend/migrations/023_derivative_attempts.down.sql b/backend/migrations/023_derivative_attempts.down.sql new file mode 100644 index 0000000..6870c02 --- /dev/null +++ b/backend/migrations/023_derivative_attempts.down.sql @@ -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; diff --git a/backend/migrations/023_derivative_attempts.up.sql b/backend/migrations/023_derivative_attempts.up.sql new file mode 100644 index 0000000..b91a296 --- /dev/null +++ b/backend/migrations/023_derivative_attempts.up.sql @@ -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; diff --git a/backend/migrations/024_feed_scalar_counts.down.sql b/backend/migrations/024_feed_scalar_counts.down.sql new file mode 100644 index 0000000..207b02c --- /dev/null +++ b/backend/migrations/024_feed_scalar_counts.down.sql @@ -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; diff --git a/backend/migrations/024_feed_scalar_counts.up.sql b/backend/migrations/024_feed_scalar_counts.up.sql new file mode 100644 index 0000000..1ed7d61 --- /dev/null +++ b/backend/migrations/024_feed_scalar_counts.up.sql @@ -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; diff --git a/backend/migrations/025_reserved_names_and_pin_decay.down.sql b/backend/migrations/025_reserved_names_and_pin_decay.down.sql new file mode 100644 index 0000000..ea485dd --- /dev/null +++ b/backend/migrations/025_reserved_names_and_pin_decay.down.sql @@ -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; diff --git a/backend/migrations/025_reserved_names_and_pin_decay.up.sql b/backend/migrations/025_reserved_names_and_pin_decay.up.sql new file mode 100644 index 0000000..7a3e926 --- /dev/null +++ b/backend/migrations/025_reserved_names_and_pin_decay.up.sql @@ -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; diff --git a/backend/src/auth/handlers.rs b/backend/src/auth/handlers.rs index b16e15b..3c515d1 100644 --- a/backend/src/auth/handlers.rs +++ b/backend/src/auth/handlers.rs @@ -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, ) -> Result, 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 { + // 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, auth: AuthUser) -> Result { 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, ) -> Result { - 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()); + } +} diff --git a/backend/src/error.rs b/backend/src/error.rs index 30fa2ad..d61a356 100644 --- a/backend/src/error.rs +++ b/backend/src/error.rs @@ -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), 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 for AppError { impl From 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") + ); } } diff --git a/backend/src/handlers/sse.rs b/backend/src/handlers/sse.rs index 3fae809..4e9de30 100644 --- a/backend/src/handlers/sse.rs +++ b/backend/src/handlers/sse.rs @@ -38,7 +38,26 @@ pub async fn issue_ticket( State(state): State, auth: AuthUser, ) -> Result, 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?; diff --git a/backend/src/handlers/upload.rs b/backend/src/handlers/upload.rs index 06a820a..92e1e57 100644 --- a/backend/src/handlers/upload.rs +++ b/backend/src/handlers/upload.rs @@ -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 { + let mut buf: Vec = 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, csv: Option<&str>) -> Vec { + 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, +} + +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)> = None; // (size, head bytes for sniffing) let mut caption: Option = None; @@ -115,8 +249,8 @@ pub async fn upload( // doesn't send one and gets the previous behaviour. let mut client_upload_id: Option = None; - // Wrap the multipart read so any error after the temp file is created still cleans - // it up (a mid-stream parse failure must not leave a stray `.tmp` on disk). + // 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 = 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::>() + .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()); + } + } } diff --git a/backend/src/main.rs b/backend/src/main.rs index 1329e24..1ecb7a5 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -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()) diff --git a/backend/src/models/upload.rs b/backend/src/models/upload.rs index 2598930..d3b44a4 100644 --- a/backend/src/models/upload.rs +++ b/backend/src/models/upload.rs @@ -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, 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(()) diff --git a/backend/src/models/user.rs b/backend/src/models/user.rs index 8b1a7ae..cad4fb5 100644 --- a/backend/src/models/user.rs +++ b/backend/src/models/user.rs @@ -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 { + 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, 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, 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 { 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) diff --git a/backend/src/services/compression.rs b/backend/src/services/compression.rs index 24af3d1..8832a54 100644 --- a/backend/src/services/compression.rs +++ b/backend/src/services/compression.rs @@ -13,6 +13,9 @@ use crate::state::SseEvent; #[derive(Clone)] pub struct CompressionWorker { semaphore: Arc, + /// Serialises the memory-heavy image jobs — see `HEAVY_IMAGE_BYTES`. Separate from + /// `semaphore` so ordinary photos keep full concurrency. + heavy: Arc, pool: PgPool, media_path: PathBuf, sse_tx: broadcast::Sender, @@ -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 = 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); + } +} diff --git a/backend/src/services/imaging.rs b/backend/src/services/imaging.rs index 3b29c3a..513fa1a 100644 --- a/backend/src/services/imaging.rs +++ b/backend/src/services/imaging.rs @@ -118,6 +118,40 @@ fn decoder_within_budget(path: &Path) -> Result { 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 { + 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 { diff --git a/backend/src/services/maintenance.rs b/backend/src/services/maintenance.rs index b684499..30f3d35 100644 --- a/backend/src/services/maintenance.rs +++ b/backend/src/services/maintenance.rs @@ -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 = 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, _> = 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::*; diff --git a/backend/src/services/sse_tickets.rs b/backend/src/services/sse_tickets.rs index 535b27c..fd921c9 100644 --- a/backend/src/services/sse_tickets.rs +++ b/backend/src/services/sse_tickets.rs @@ -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>>, @@ -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 { 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 = 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" + ); + } } diff --git a/backend/src/services/video.rs b/backend/src/services/video.rs index 968121e..c31d31f 100644 --- a/backend/src/services/video.rs +++ b/backend/src/services/video.rs @@ -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 = (0..500) + .map(|i| format!("line {i}\n")) + .collect::() + .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); + } } diff --git a/docker-compose.yml b/docker-compose.yml index b9da24f..279799f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -56,9 +56,24 @@ services: logging: *default-logging env_file: .env environment: + # Default to info. Without this the code fallback in main.rs applies, which is + # `eventsnap_backend=debug,tower_http=debug` — a line per HTTP request AND per + # response, including every preview fetch, for a multi-day run. The x-logging cap + # above bounds the disk cost but not the CPU/IO one. + # + # Set here rather than only in `.env` because a stock deploy sets RUST_LOG nowhere, + # and this is the layer an operator will actually find when they need to raise it + # for a single event (`RUST_LOG=eventsnap_backend=debug docker compose up -d app`). + RUST_LOG: ${RUST_LOG:-info} # Activates the production secret guard in config.rs — refuses to boot with # placeholder JWT_SECRET / ADMIN_PASSWORD_HASH. APP_ENV: production + # Pinned beside MEDIA_PATH for the same reason, and because nothing validates it: + # config.rs defaults it to /exports but never checks that it is a mount or that it + # differs from media_path. A stray EXPORT_PATH in .env builds the keepsake into the + # container's writable layer, where it passes every health check and disk preflight + # and then evaporates on the next `up -d`. + EXPORT_PATH: /exports # The media volume is mounted at /media (below), so the app MUST write there. # Pin it here rather than trusting .env: if MEDIA_PATH in .env points elsewhere # (e.g. a host path used for running the backend natively) the container can't diff --git a/e2e/fixtures/db.ts b/e2e/fixtures/db.ts index 223e9a1..abde3a6 100644 --- a/e2e/fixtures/db.ts +++ b/e2e/fixtures/db.ts @@ -37,6 +37,52 @@ export const db = { ); }, + /** + * Is this user's account currently PIN-locked? + * + * Distinguishes the two ways /recover can answer 429 — the per-(IP, name) throttle, which + * costs the attacker, and the account lock, which costs the VICTIM. Only the second one is + * weaponizable, so a test asserting "a single IP cannot lock a guest out" has to look at the + * row, not at the status code. + */ + async isPinLocked(userId: string): Promise { + return withClient(async (c) => { + const r = await c.query<{ locked: boolean }>( + `SELECT (pin_locked_until IS NOT NULL AND pin_locked_until > NOW()) AS locked + FROM "user" WHERE id = $1`, + [userId] + ); + return r.rows[0]?.locked ?? false; + }); + }, + + /** + * Preload the wrong-PIN streak, standing in for failures that arrived from other IPs. + * + * The account lock is deliberately out of reach of any single source, so a test that wants to + * exercise it has to simulate the distributed case rather than hammer from one address. + * `last_failed_pin_at` is set to now so the 15-minute decay does not immediately reset it. + */ + async setFailedPinAttempts(userId: string, attempts: number) { + await withClient((c) => + c.query( + `UPDATE "user" SET failed_pin_attempts = $2, last_failed_pin_at = NOW() WHERE id = $1`, + [userId, attempts] + ) + ); + }, + + /** Current wrong-PIN streak. Decays after 15 minutes — see User::increment_failed_pin. */ + async failedPinAttempts(userId: string): Promise { + return withClient(async (c) => { + const r = await c.query<{ failed_pin_attempts: number }>( + `SELECT failed_pin_attempts FROM "user" WHERE id = $1`, + [userId] + ); + return r.rows[0]?.failed_pin_attempts ?? 0; + }); + }, + async expireSession(userId: string) { await withClient((c) => c.query(`UPDATE session SET expires_at = NOW() - interval '1 hour' WHERE user_id = $1`, [ diff --git a/e2e/fixtures/media/not-an-image.jpg b/e2e/fixtures/media/not-an-image.jpg new file mode 100644 index 0000000..4238434 --- /dev/null +++ b/e2e/fixtures/media/not-an-image.jpg @@ -0,0 +1,4 @@ +This is plain text, not an image at all. +This is plain text, not an image at all. +This is plain text, not an image at all. +This is plain text, not an image at all. diff --git a/e2e/specs/01-auth/join.spec.ts b/e2e/specs/01-auth/join.spec.ts index c5eadb9..fb02ac2 100644 --- a/e2e/specs/01-auth/join.spec.ts +++ b/e2e/specs/01-auth/join.spec.ts @@ -75,7 +75,11 @@ test.describe('Auth — join flow', () => { expect(storage.pin).toBe(original.pin); }); - test('wrong PIN three times locks the account for 15 minutes', async ({ page, guest, db }) => { + test('repeated wrong PINs are throttled without locking the guest out', async ({ + page, + guest, + db, + }) => { const dave = await guest('Dave'); await clearAllStorage(page); @@ -85,22 +89,33 @@ test.describe('Auth — join flow', () => { await join.submit(); await expect(join.recoveryPinInput).toBeVisible(); - // Wrong PIN (real one is dave.pin) + // Wrong PIN (real one is dave.pin), four times — one more than the OLD lock threshold of 3. + // Typed digit by digit so the 4th character auto-submits (see pin-auto-submit.spec.ts); + // clicking as well would double-submit and race the disabled state of the button. const wrong = dave.pin === '0000' ? '1111' : '0000'; - for (let i = 0; i < 3; i++) { - await join.recoveryPinInput.fill(wrong); - await join.recoverySubmit.click(); + for (let i = 0; i < 4; i++) { + await join.recoveryPinInput.fill(''); + await join.recoveryPinInput.pressSequentially(wrong, { delay: 30 }); await expect(join.recoveryError).toBeVisible(); + await expect(join.recoverySubmit).toBeEnabled(); } - // Fourth attempt should hit the 429 lockout (even with the correct PIN now) - await join.recoveryPinInput.fill(dave.pin); - await join.recoverySubmit.click(); - await expect(join.recoveryError).toContainText(/15 Minuten/); + // THE PROPERTY THIS TEST EXISTS FOR, stated the way a guest experiences it: Dave can still + // get into his own account. + // + // The lock threshold used to be 3, BELOW the per-(IP, name) ceiling — so these very + // keystrokes locked Dave out for 15 minutes, and anyone who can read his name off the feed + // could do it to him on repeat. Rate limits are disabled in this environment (see + // config `rate_limits_enabled`), so what is exercised here is purely the account-lock tier; + // the throttle tier is covered in 07-adversarial/auth-tampering.spec.ts. + expect( + await db.isPinLocked(dave.userId), + 'four wrong PINs from one device must not lock a guest out of their own account' + ).toBe(false); - // Sanity: DB row reflects the lock - // (The handler sets pin_locked_until directly — verify via API "recover" returning 429) - void db; // unused for now, documenting that db.lockUserPin exists if we want shortcut path + await join.recoveryPinInput.fill(''); + await join.recoveryPinInput.pressSequentially(dave.pin, { delay: 30 }); + await page.waitForURL('**/feed'); }); test('"Anderen Namen wählen" returns to the normal join form', async ({ page, guest }) => { diff --git a/e2e/specs/02-upload/burst-queue.spec.ts b/e2e/specs/02-upload/burst-queue.spec.ts index fb9762f..a64b773 100644 --- a/e2e/specs/02-upload/burst-queue.spec.ts +++ b/e2e/specs/02-upload/burst-queue.spec.ts @@ -199,9 +199,12 @@ test.describe('Upload — client queue under a burst', () => { // a closed tab / killed PWA. The remaining pending items live only in // IndexedDB now. await page.reload(); - // The queue only resumes where loadQueue() runs — the /upload route's - // onMount. Navigating there is the "reopen the composer" recovery path. - await page.goto('/upload'); + // Deliberately NOT navigating to /upload. Rehydration is now module-level and + // auth-gated (upload-queue.ts `hydrateQueue`), so the queue resumes wherever the + // reload lands. This assertion is the regression guard for the defect it replaced: + // `loadQueue()` used to have a single call site in the whole app — the /upload + // route's onMount — so a guest who reloaded anywhere else saw a 0 badge and their + // staged photos never left the phone, having already been shown a success. // (4) RESUME: every file ends up server-side without re-staging anything. // `>=` not `===`: the only imperfection possible is a DUPLICATE (an upload diff --git a/e2e/specs/07-adversarial/auth-tampering.spec.ts b/e2e/specs/07-adversarial/auth-tampering.spec.ts index 69e7f05..a076434 100644 --- a/e2e/specs/07-adversarial/auth-tampering.spec.ts +++ b/e2e/specs/07-adversarial/auth-tampering.spec.ts @@ -89,15 +89,43 @@ test.describe('Adversarial — JWT', () => { }); test.describe('Adversarial — PIN brute-force', () => { - test('sequential wrong-PIN attempts lock the account after 3 attempts', async ({ guest }) => { + /** + * The PIN defence has two tiers, and telling them apart is the whole point of these tests: + * + * - the per-(IP, name) throttle, which refuses the ATTACKER; and + * - the account lock, which refuses the VICTIM — the only tier that can be weaponised. + * + * Both answer 429, so the status code alone proves nothing. The regression these guard is that + * the lock threshold used to sit BELOW the throttle ceiling (3 vs 5), so three requests from a + * single IP locked any guest whose display name is readable off the feed, every 15 minutes, + * indefinitely. The tier meant to protect a guest was the cheapest way to attack them. + */ + // Restore the default. The first test turns the limiter on for the whole instance, and leaving + // it on would throttle unrelated specs sharing this stack. Runs even on failure. + test.afterEach(async ({ api, adminToken }) => { + await api.patchConfig(adminToken, { rate_limits_enabled: 'false' }); + }); + + test('a single IP is throttled without ever locking the victim out', async ({ + api, + adminToken, + guest, + db, + }) => { + // Rate limits are off by default in this environment, and the throttle IS the tier under + // test — without it the run would silently assert only half the property. + await api.patchConfig(adminToken, { + rate_limits_enabled: 'true', + recover_rate_enabled: 'true', + }); + const g = await guest('Brute'); const wrong = g.pin === '0000' ? '1111' : '0000'; - // Do them serially so the failed_pin_attempts counter increments - // monotonically. Parallel attempts race and may never accumulate to 3 in - // the current handler implementation — that's a separate finding. + // Serially, so the failed-PIN counter increments monotonically. Well past the per-(IP, name) + // ceiling of 4, and past the OLD lock threshold of 3. const statuses: number[] = []; - for (let i = 0; i < 4; i++) { + for (let i = 0; i < 8; i++) { const r = await fetch(`${BASE}/api/v1/recover`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, @@ -105,26 +133,52 @@ test.describe('Adversarial — PIN brute-force', () => { }); statuses.push(r.status); } - // First three are 401, fourth (or later) is 429. - expect(statuses.filter((s) => s === 200)).toHaveLength(0); - expect(statuses.some((s) => s === 429)).toBe(true); + expect(statuses.filter((s) => s === 200), 'a wrong PIN must never authenticate').toHaveLength( + 0 + ); + expect(statuses.some((s) => s === 429), 'the attacker must be throttled').toBe(true); - // Now even the correct PIN fails until lockout expires. + expect( + await db.isPinLocked(g.userId), + 'one IP must not be able to lock a guest out of their own account' + ).toBe(false); + }); + + test('the account still locks once the failure count is reached', async ({ guest, db }) => { + const g = await guest('BruteDistributed'); + const wrong = g.pin === '0000' ? '1111' : '0000'; + + // The per-IP throttle is what a single source hits first, so drive the counter the way a + // DISTRIBUTED attacker would — the tier this test covers is the last line against exactly + // that, and it must not have been removed while fixing the weaponisation above. + await db.setFailedPinAttempts(g.userId, 11); + + const r = await fetch(`${BASE}/api/v1/recover`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-Forwarded-For': '203.0.113.77' }, + body: JSON.stringify({ display_name: g.displayName, pin: wrong }), + }); + expect(r.status).toBe(401); + + expect( + await db.isPinLocked(g.userId), + 'a distributed guesser must still trip the account lock' + ).toBe(true); + + // And the lock holds even against the correct PIN, which is what makes it a real control. const correct = await fetch(`${BASE}/api/v1/recover`, { method: 'POST', - headers: { 'Content-Type': 'application/json' }, + headers: { 'Content-Type': 'application/json', 'X-Forwarded-For': '203.0.113.78' }, body: JSON.stringify({ display_name: g.displayName, pin: g.pin }), }); expect(correct.status).toBe(429); }); - test('parallel wrong-PIN attempts still lock the account (counter is not lost to the race)', async ({ - guest, - }) => { + test('the wrong-PIN streak is atomic under concurrency', async ({ guest, db }) => { const g = await guest('BruteParallel'); const wrong = g.pin === '0000' ? '1111' : '0000'; - const attempts = await Promise.all( + await Promise.all( Array.from({ length: 10 }, () => fetch(`${BASE}/api/v1/recover`, { method: 'POST', @@ -133,29 +187,14 @@ test.describe('Adversarial — PIN brute-force', () => { }) ) ); - const statuses = attempts.map((r) => r.status); - expect( - statuses.filter((s) => s === 200), - 'a wrong PIN must never authenticate' - ).toHaveLength(0); - // The in-flight requests all read `pin_locked_until` before any of them wrote it, so - // *which* of the 10 come back 429 is genuinely racy and can't be asserted. What is NOT - // racy — and is the property this test exists to guard — is the state left behind: - // `failed_pin_attempts` is incremented with an atomic `SET x = x + 1 ... RETURNING`, so - // 10 wrong PINs must push it past the 3-strike threshold and leave the account locked. - // - // We prove that with a follow-up request using the CORRECT pin: it must be refused with - // 429 (locked), not 200. Delete the lockout counter and this line goes 200 → red. - const correct = await fetch(`${BASE}/api/v1/recover`, { - method: 'POST', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ display_name: g.displayName, pin: g.pin }), - }); + // How many of the 10 get past the throttle is genuinely racy and cannot be asserted. What is + // NOT racy is that every one that DID reach the handler incremented the counter — it is a + // single `SET x = x + 1 ... RETURNING`, so none of them can be lost to the race. expect( - correct.status, - 'after 10 wrong PINs the account must be locked, even for the right PIN' - ).toBe(429); + await db.failedPinAttempts(g.userId), + 'concurrent wrong PINs must all be counted' + ).toBeGreaterThan(1); }); }); diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index 87d7e6a..dab35b3 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -46,6 +46,13 @@ async function request(method: string, path: string, body?: unknown): Promise // Abort hung requests so a dead connection surfaces as a friendly error // instead of a spinner that never resolves. + // + // The timer must stay armed until the BODY has been read, not just the headers. + // `fetch` resolves as soon as the response head arrives, so clearing it in a `finally` + // around the fetch left `res.text()` below completely uncovered — and no longer + // abortable, since the controller had already been disarmed. An upstream that sends + // headers and then stalls the body (the shape of a half-dead proxy, or of the pool + // saturation this same release adds shedding for) hung that call forever. const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), TIMEOUT_MS); @@ -57,6 +64,26 @@ async function request(method: string, path: string, body?: unknown): Promise body: body !== undefined ? JSON.stringify(body) : undefined, signal: controller.signal }); + } catch (e) { + clearTimeout(timer); + if (e instanceof DOMException && e.name === 'AbortError') { + throw new ApiError(0, 'timeout', 'Zeitüberschreitung – bitte erneut versuchen.'); + } + throw new ApiError(0, 'network', 'Netzwerkfehler – bitte Verbindung prüfen.'); + } + + if (res.status === 204) { + // Must clear on this path too, or every no-content request (logout, delete, like) + // leaks a live 20 s timer. + clearTimeout(timer); + return undefined as T; + } + + // A 5xx behind a proxy (or a crash page) can return HTML, not JSON — parsing + // it directly would throw an opaque SyntaxError. Read text, parse defensively. + let raw: string; + try { + raw = await res.text(); } catch (e) { if (e instanceof DOMException && e.name === 'AbortError') { throw new ApiError(0, 'timeout', 'Zeitüberschreitung – bitte erneut versuchen.'); @@ -66,13 +93,6 @@ async function request(method: string, path: string, body?: unknown): Promise clearTimeout(timer); } - if (res.status === 204) { - return undefined as T; - } - - // A 5xx behind a proxy (or a crash page) can return HTML, not JSON — parsing - // it directly would throw an opaque SyntaxError. Read text, parse defensively. - const raw = await res.text(); let data: { error?: string; message?: string } | unknown = null; if (raw) { try { diff --git a/frontend/src/lib/upload-queue.test.ts b/frontend/src/lib/upload-queue.test.ts index e861cde..b5b8072 100644 --- a/frontend/src/lib/upload-queue.test.ts +++ b/frontend/src/lib/upload-queue.test.ts @@ -1,5 +1,10 @@ import { describe, it, expect } from 'vitest'; -import { classifyUploadStatus, isReversibleLock, entryToQueueItem } from './upload-queue'; +import { + classifyUploadStatus, + isReversibleLock, + entryToQueueItem, + shouldAbortForStall +} from './upload-queue'; /** * Regression guard for the upload-queue retry policy (H2 + M1). The bug being locked out: @@ -107,3 +112,32 @@ describe('entryToQueueItem', () => { expect(item.hashtags).toBe(''); }); }); + +/** + * The upload XHR had no timeout of any kind while `processQueue` held the `isProcessing` + * latch across it. On a half-open socket neither `error` nor `abort` ever fires, so the + * latch was pinned forever and the whole queue wedged with no recovery but a reload. + * + * The policy that matters: bound SILENCE, not total duration. A 500 MB video over a venue + * uplink legitimately runs 30+ minutes while making steady progress, and a flat total cap + * would kill exactly the uploads worth keeping. + */ +describe('shouldAbortForStall', () => { + const now = 1_000_000; + + it('lets a long upload run as long as progress keeps arriving', () => { + // Two hours in, but progress landed a second ago. + expect(shouldAbortForStall(now - 1_000, now, false)).toBe(false); + }); + + it('aborts once the body stalls past the no-progress ceiling', () => { + expect(shouldAbortForStall(now - 89_000, now, false)).toBe(false); + expect(shouldAbortForStall(now - 91_000, now, false)).toBe(true); + }); + + it('applies the wider ceiling once the body is sent and progress goes quiet', () => { + // Silence that would abort mid-body is normal while waiting for the response. + expect(shouldAbortForStall(now - 91_000, now, true)).toBe(false); + expect(shouldAbortForStall(now - 121_000, now, true)).toBe(true); + }); +}); diff --git a/frontend/src/lib/upload-queue.ts b/frontend/src/lib/upload-queue.ts index d6d2c0f..2ad17bb 100644 --- a/frontend/src/lib/upload-queue.ts +++ b/frontend/src/lib/upload-queue.ts @@ -65,6 +65,35 @@ const MAX_RETRY_DELAY_MS = 5 * 60_000; const STALL_TIMEOUT_MS = 90_000; const STALL_CHECK_INTERVAL_MS = 5_000; +/** + * Ceiling for the window AFTER the last byte is sent. + * + * `upload.progress` is silent there BY DEFINITION — the server is sniffing the magic bytes, + * committing the row and writing the response — so the bytes-moved signal above has nothing + * to measure and the watchdog used to simply switch off at `loadend`. That left the wall-clock + * `xhr.timeout` as the only remaining bound: 5 minutes for a photo, up to 60 for a video. A + * half-open socket in that window (phone roams wifi→LTE while the server is still writing a + * 500 MB video) pinned `processing` for that entire time, so NOTHING else in the guest's queue + * drained either. + * + * Sized above the backend's worst-case commit path, not above compression — compression is + * spawned after the response and does not hold it open. + */ +const RESPONSE_TIMEOUT_MS = 120_000; + +/** + * Pure predicate behind the watchdog, extracted so the policy is unit-testable without + * standing up an XHR harness. + */ +export function shouldAbortForStall( + lastActivityAt: number, + now: number, + bodySent: boolean +): boolean { + const ceiling = bodySent ? RESPONSE_TIMEOUT_MS : STALL_TIMEOUT_MS; + return now - lastActivityAt > ceiling; +} + /** * Wall-clock cap for one attempt, scaled by file size assuming a floor of ~8 kB/s — a * deliberately pessimistic rate, because killing a slow-but-progressing upload would lose @@ -865,9 +894,10 @@ async function uploadItem(id: string): Promise { // connection that never errors and never completes. Only "no bytes moved" catches // that without also punishing a healthy slow link. let lastProgressAt = Date.now(); + let bodySent = false; let stalled = false; const stallTimer = setInterval(() => { - if (Date.now() - lastProgressAt < STALL_TIMEOUT_MS) return; + if (!shouldAbortForStall(lastProgressAt, Date.now(), bodySent)) return; stalled = true; xhr.abort(); }, STALL_CHECK_INTERVAL_MS); @@ -889,9 +919,19 @@ async function uploadItem(id: string): Promise { }); // Once the last byte is out the watchdog has nothing left to measure: the server may // legitimately sit on the request while it validates and stores the file, and no - // progress event fires in that window. Aborting there would re-send a whole video - // the server had already accepted, so hand over to `xhr.timeout` instead. - xhr.upload.addEventListener('loadend', () => clearInterval(stallTimer)); + // progress event fires in that window. Aborting on the 90s no-progress ceiling there + // would re-send a whole video the server had already accepted. + // + // But switching the watchdog OFF here (which is what this used to do) handed the + // window to `xhr.timeout` alone — 5 to 60 minutes, during which a half-open socket + // holds `processing` and the guest's whole queue stops draining. So instead of + // disarming, widen: `shouldAbortForStall` switches to RESPONSE_TIMEOUT_MS once + // `bodySent` flips, which bounds the wedge at 2 minutes without ever firing on a + // server that is legitimately still working. + xhr.upload.addEventListener('loadend', () => { + bodySent = true; + lastProgressAt = Date.now(); + }); xhr.addEventListener('load', () => { const body = (() => {