Performance: - Cache the runtime `config` table in-memory (ConfigCache) with synchronous invalidation on every write (admin PATCH + test reseed). Was re-reading each key from Postgres on every request (~8 round-trips per upload). - Stream uploads chunk-by-chunk to a temp file instead of buffering the whole body in RAM (peak was up to the per-class cap, e.g. 500 MB/video); only 512 sniff-bytes are kept for magic-byte detection, then atomic rename into place. - Cache the media-filesystem disk snapshot (DiskCache, 15s TTL) shared by the quota check and admin stats; drop the discarded System::refresh_all(). - HTML export streams video (and small-image) originals straight into the ZIP via a manifest instead of copying them to a temp dir first (removed the transient 2x disk usage) and drops the double directory scan. - Auth extractor resolves session -> live user in one JOIN (was two queries), touching last_seen_at by token hash. Stability: - SSE: on broadcast lag, emit a `resync` event so the client runs a delta fetch instead of silently losing events; frontend reconciles adds, deletions, and (via an in-place refresh) like/comment counts on visible cards. - Storage quota fails OPEN when the disk can't be read (was a 0-byte limit that locked out all uploads). - Graceful shutdown drains in-flight requests on SIGTERM/SIGINT, bounded by a 10s backstop so open SSE streams can't stall a deploy. - Upload removes the persisted file if the DB transaction fails (no orphaned bytes with no row to reclaim them). Tests: - New pure select_disk() with 5 unit tests (longest-prefix, fallbacks, fail-open). - New e2e export-video spec covering the HTML export's video-streaming branch. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
422 lines
15 KiB
Rust
422 lines
15 KiB
Rust
use std::collections::HashMap;
|
||
use std::time::Duration;
|
||
|
||
use axum::extract::{Query, State};
|
||
use axum::http::{HeaderMap, StatusCode};
|
||
use axum::Json;
|
||
use serde::{Deserialize, Serialize};
|
||
|
||
use crate::auth::middleware::RequireAdmin;
|
||
use crate::error::AppError;
|
||
use crate::services::config;
|
||
use crate::services::rate_limiter::client_ip;
|
||
use crate::state::AppState;
|
||
|
||
// ── DTOs ─────────────────────────────────────────────────────────────────────
|
||
|
||
#[derive(Serialize)]
|
||
pub struct StatsDto {
|
||
pub user_count: i64,
|
||
pub upload_count: i64,
|
||
pub comment_count: i64,
|
||
pub disk_total_bytes: u64,
|
||
pub disk_used_bytes: u64,
|
||
pub disk_free_bytes: u64,
|
||
}
|
||
|
||
#[derive(Serialize, sqlx::FromRow)]
|
||
pub struct ExportJobDto {
|
||
pub id: uuid::Uuid,
|
||
pub r#type: String,
|
||
pub status: String,
|
||
pub progress_pct: i16,
|
||
pub error_message: Option<String>,
|
||
pub created_at: chrono::DateTime<chrono::Utc>,
|
||
pub completed_at: Option<chrono::DateTime<chrono::Utc>>,
|
||
}
|
||
|
||
// ── Handlers ─────────────────────────────────────────────────────────────────
|
||
|
||
pub async fn get_stats(
|
||
State(state): State<AppState>,
|
||
RequireAdmin(_auth): RequireAdmin,
|
||
) -> Result<Json<StatsDto>, AppError> {
|
||
let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug)
|
||
.await?
|
||
.ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?;
|
||
|
||
let (user_count,): (i64,) =
|
||
sqlx::query_as("SELECT COUNT(*) FROM \"user\" WHERE event_id = $1")
|
||
.bind(event.id)
|
||
.fetch_one(&state.pool)
|
||
.await?;
|
||
|
||
let (upload_count,): (i64,) = sqlx::query_as(
|
||
"SELECT COUNT(*) FROM upload WHERE event_id = $1 AND deleted_at IS NULL",
|
||
)
|
||
.bind(event.id)
|
||
.fetch_one(&state.pool)
|
||
.await?;
|
||
|
||
let (comment_count,): (i64,) = sqlx::query_as(
|
||
"SELECT COUNT(*) FROM comment c
|
||
JOIN upload u ON u.id = c.upload_id
|
||
WHERE u.event_id = $1 AND c.deleted_at IS NULL",
|
||
)
|
||
.bind(event.id)
|
||
.fetch_one(&state.pool)
|
||
.await?;
|
||
|
||
// Disk usage from the shared cache (unknown mount → zeros, same as before).
|
||
let (disk_total, disk_free) = state
|
||
.disk_cache
|
||
.snapshot(&state.config.media_path)
|
||
.map(|d| (d.total, d.free))
|
||
.unwrap_or((0, 0));
|
||
|
||
let disk_used = disk_total.saturating_sub(disk_free);
|
||
|
||
Ok(Json(StatsDto {
|
||
user_count,
|
||
upload_count,
|
||
comment_count,
|
||
disk_total_bytes: disk_total,
|
||
disk_used_bytes: disk_used,
|
||
disk_free_bytes: disk_free,
|
||
}))
|
||
}
|
||
|
||
pub async fn get_config(
|
||
State(state): State<AppState>,
|
||
RequireAdmin(_auth): RequireAdmin,
|
||
) -> Result<Json<HashMap<String, String>>, AppError> {
|
||
let rows: Vec<(String, String)> =
|
||
sqlx::query_as("SELECT key, value FROM config ORDER BY key")
|
||
.fetch_all(&state.pool)
|
||
.await?;
|
||
|
||
Ok(Json(rows.into_iter().collect()))
|
||
}
|
||
|
||
#[derive(Deserialize)]
|
||
pub struct PatchConfigRequest(pub HashMap<String, String>);
|
||
|
||
pub async fn patch_config(
|
||
State(state): State<AppState>,
|
||
RequireAdmin(_auth): RequireAdmin,
|
||
Json(body): Json<HashMap<String, String>>,
|
||
) -> Result<StatusCode, AppError> {
|
||
// Numeric keys validated as f64; boolean keys validated as truthy strings; the
|
||
// privacy note is free text. Splitting these explicitly is verbose but makes the
|
||
// failure mode for typos obvious (`Unbekannter Schlüssel: ...`).
|
||
// (key, integer_only, min, max). Ranges reject values that `parse::<f64>` would
|
||
// accept but that silently revert to the hardcoded default at read time
|
||
// (get_usize/get_i64 can't parse negatives/NaN/fractionals). `compression_concurrency`
|
||
// is intentionally absent — it's read once at boot, so a live edit was a no-op.
|
||
const NUMERIC_SPECS: &[(&str, bool, f64, f64)] = &[
|
||
("max_image_size_mb", true, 1.0, 1024.0),
|
||
("max_video_size_mb", true, 1.0, 10240.0),
|
||
("upload_rate_per_hour", true, 1.0, 100_000.0),
|
||
("feed_rate_per_min", true, 1.0, 100_000.0),
|
||
("export_rate_per_day", true, 1.0, 100_000.0),
|
||
("quota_tolerance", false, 0.0, 1.0),
|
||
("estimated_guest_count", true, 1.0, 1_000_000.0),
|
||
];
|
||
const BOOL_KEYS: &[&str] = &[
|
||
"rate_limits_enabled",
|
||
"upload_rate_enabled",
|
||
"feed_rate_enabled",
|
||
"export_rate_enabled",
|
||
"join_rate_enabled",
|
||
"quota_enabled",
|
||
"storage_quota_enabled",
|
||
"upload_count_quota_enabled",
|
||
];
|
||
const TEXT_KEYS: &[&str] = &["privacy_note"];
|
||
const PRIVACY_NOTE_MAX_LEN: usize = 16 * 1024; // 16 KiB free text is plenty
|
||
|
||
let mut privacy_note_changed = false;
|
||
|
||
// Validate every key first so a bad value in the batch can't leave a partial
|
||
// update behind — validation must fully precede any write.
|
||
for (key, value) in &body {
|
||
let key_str = key.as_str();
|
||
if let Some(&(_, integer_only, min, max)) =
|
||
NUMERIC_SPECS.iter().find(|(k, ..)| *k == key_str)
|
||
{
|
||
let n = value.trim().parse::<f64>().ok().filter(|n| n.is_finite());
|
||
let n = match n {
|
||
Some(n) => n,
|
||
None => {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Ungültiger Wert für {key}: muss eine Zahl sein."
|
||
)))
|
||
}
|
||
};
|
||
if integer_only && n.fract() != 0.0 {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Ungültiger Wert für {key}: muss eine ganze Zahl sein."
|
||
)));
|
||
}
|
||
if n < min || n > max {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Wert für {key} liegt außerhalb des zulässigen Bereichs ({min}–{max})."
|
||
)));
|
||
}
|
||
} else if BOOL_KEYS.contains(&key_str) {
|
||
match value.trim().to_ascii_lowercase().as_str() {
|
||
"true" | "false" | "1" | "0" | "yes" | "no" | "on" | "off" => {}
|
||
_ => {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Ungültiger Wert für {key}: muss true oder false sein."
|
||
)));
|
||
}
|
||
}
|
||
} else if TEXT_KEYS.contains(&key_str) {
|
||
// Count characters, not bytes — the message says "Zeichen" and a
|
||
// multi-byte grapheme shouldn't count against the limit multiple times.
|
||
if value.chars().count() > PRIVACY_NOTE_MAX_LEN {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Wert für {key} ist zu lang (max. {PRIVACY_NOTE_MAX_LEN} Zeichen)."
|
||
)));
|
||
}
|
||
if key_str == "privacy_note" {
|
||
privacy_note_changed = true;
|
||
}
|
||
} else {
|
||
return Err(AppError::BadRequest(format!(
|
||
"Unbekannter Konfigurationsschlüssel: {key}"
|
||
)));
|
||
}
|
||
}
|
||
|
||
// Apply all writes in one transaction — the batch is all-or-nothing.
|
||
let mut tx = state.pool.begin().await?;
|
||
for (key, value) in &body {
|
||
sqlx::query(
|
||
"INSERT INTO config (key, value, updated_at) VALUES ($1, $2, NOW())
|
||
ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, updated_at = NOW()",
|
||
)
|
||
.bind(key)
|
||
.bind(value)
|
||
.execute(&mut *tx)
|
||
.await?;
|
||
}
|
||
tx.commit().await?;
|
||
|
||
// The config cache must reflect this write on the very next read (tests PATCH then
|
||
// immediately assert the new value takes effect). Invalidate synchronously here —
|
||
// the TTL is only a backstop and must not be relied on for correctness.
|
||
state.config_cache.invalidate();
|
||
|
||
// Notify all clients that a publicly-readable config value changed so their stores
|
||
// (e.g. the privacy note in My Account) refresh without a manual reload.
|
||
if privacy_note_changed {
|
||
let _ = state.sse_tx.send(crate::state::SseEvent::new(
|
||
"event-updated",
|
||
serde_json::json!({ "keys": ["privacy_note"] }).to_string(),
|
||
));
|
||
}
|
||
|
||
Ok(StatusCode::NO_CONTENT)
|
||
}
|
||
|
||
pub async fn get_export_jobs(
|
||
State(state): State<AppState>,
|
||
RequireAdmin(_auth): RequireAdmin,
|
||
) -> Result<Json<Vec<ExportJobDto>>, AppError> {
|
||
let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug)
|
||
.await?
|
||
.ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?;
|
||
|
||
let jobs = sqlx::query_as::<_, ExportJobDto>(
|
||
"SELECT id, type::text, status::text, progress_pct, error_message, created_at, completed_at
|
||
FROM export_job
|
||
WHERE event_id = $1
|
||
ORDER BY created_at DESC",
|
||
)
|
||
.bind(event.id)
|
||
.fetch_all(&state.pool)
|
||
.await?;
|
||
|
||
Ok(Json(jobs))
|
||
}
|
||
|
||
// ── Export download endpoints (authenticated guests) ─────────────────────────
|
||
|
||
#[derive(Deserialize)]
|
||
pub struct DownloadQuery {
|
||
pub ticket: String,
|
||
}
|
||
|
||
/// Mint a short-lived ticket for a browser-driven export download. The download
|
||
/// is a top-level navigation so the multi-GB ZIP streams straight to disk instead
|
||
/// of being buffered in memory by `fetch()` + `blob()` — but a navigation can't
|
||
/// carry an `Authorization` header, so the client exchanges its Bearer token for
|
||
/// an opaque ticket here, then hits `/export/zip?ticket=...`. Reuses the same
|
||
/// single-use, 30s-TTL store as the SSE stream.
|
||
pub async fn export_ticket(
|
||
State(state): State<AppState>,
|
||
auth: crate::auth::middleware::AuthUser,
|
||
) -> Json<serde_json::Value> {
|
||
// NOTE: intentionally NOT gated on `is_banned`. A banned user keeps *read* access
|
||
// by design (USER_JOURNEYS §10.3, FEATURES: "Can still download the export once
|
||
// released — Spec design choice"). The export is read-only, so it stays available
|
||
// to them, consistent with the read-only-ban model.
|
||
let ticket = state.sse_tickets.issue(auth.token_hash);
|
||
Json(serde_json::json!({ "ticket": ticket }))
|
||
}
|
||
|
||
/// Validate a download ticket (single-use) and confirm its session still exists.
|
||
async fn authenticate_download_ticket(state: &AppState, ticket: &str) -> Result<(), AppError> {
|
||
let token_hash = state
|
||
.sse_tickets
|
||
.consume(ticket)
|
||
.ok_or_else(|| AppError::Unauthorized("Ticket ungültig oder abgelaufen.".into()))?;
|
||
crate::models::session::Session::find_by_token_hash(&state.pool, &token_hash)
|
||
.await
|
||
.map_err(|e| AppError::Internal(e.into()))?
|
||
.ok_or_else(|| AppError::Unauthorized("Sitzung nicht gefunden.".into()))?;
|
||
Ok(())
|
||
}
|
||
|
||
pub async fn download_zip(
|
||
State(state): State<AppState>,
|
||
Query(q): Query<DownloadQuery>,
|
||
headers: HeaderMap,
|
||
) -> Result<axum::response::Response, AppError> {
|
||
authenticate_download_ticket(&state, &q.ticket).await?;
|
||
enforce_export_rate(&state, &headers).await?;
|
||
|
||
let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug)
|
||
.await?
|
||
.ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?;
|
||
|
||
if !event.export_zip_ready {
|
||
return Err(AppError::NotFound(
|
||
"Der ZIP-Export ist noch nicht verfügbar.".into(),
|
||
));
|
||
}
|
||
|
||
let path = state.config.export_path.join("Gallery.zip");
|
||
if !path.exists() {
|
||
return Err(AppError::NotFound("Exportdatei nicht gefunden.".into()));
|
||
}
|
||
|
||
serve_file(path, "Gallery.zip", "application/zip").await
|
||
}
|
||
|
||
pub async fn download_html(
|
||
State(state): State<AppState>,
|
||
Query(q): Query<DownloadQuery>,
|
||
headers: HeaderMap,
|
||
) -> Result<axum::response::Response, AppError> {
|
||
authenticate_download_ticket(&state, &q.ticket).await?;
|
||
enforce_export_rate(&state, &headers).await?;
|
||
|
||
let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug)
|
||
.await?
|
||
.ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?;
|
||
|
||
if !event.export_html_ready {
|
||
return Err(AppError::NotFound(
|
||
"Der HTML-Export ist noch nicht verfügbar.".into(),
|
||
));
|
||
}
|
||
|
||
let path = state.config.export_path.join("Memories.zip");
|
||
if !path.exists() {
|
||
return Err(AppError::NotFound("Exportdatei nicht gefunden.".into()));
|
||
}
|
||
|
||
serve_file(path, "Memories.zip", "application/zip").await
|
||
}
|
||
|
||
async fn serve_file(
|
||
path: std::path::PathBuf,
|
||
filename: &str,
|
||
content_type: &str,
|
||
) -> Result<axum::response::Response, AppError> {
|
||
use axum::body::Body;
|
||
use axum::http::{header, Response, StatusCode};
|
||
use tokio_util::io::ReaderStream;
|
||
|
||
let file = tokio::fs::File::open(&path)
|
||
.await
|
||
.map_err(|e| AppError::Internal(e.into()))?;
|
||
let metadata = file
|
||
.metadata()
|
||
.await
|
||
.map_err(|e| AppError::Internal(e.into()))?;
|
||
let stream = ReaderStream::new(file);
|
||
|
||
let disposition = format!("attachment; filename=\"{filename}\"");
|
||
|
||
let response = Response::builder()
|
||
.status(StatusCode::OK)
|
||
.header(header::CONTENT_TYPE, content_type)
|
||
.header(header::CONTENT_DISPOSITION, disposition)
|
||
.header(header::CONTENT_LENGTH, metadata.len())
|
||
.body(Body::from_stream(stream))
|
||
.map_err(|e| AppError::Internal(e.into()))?;
|
||
|
||
Ok(response)
|
||
}
|
||
|
||
/// Also expose export status to all authenticated users (guests need it for the export page)
|
||
pub async fn export_status(
|
||
State(state): State<AppState>,
|
||
_auth: crate::auth::middleware::AuthUser,
|
||
) -> Result<Json<serde_json::Value>, AppError> {
|
||
let event = crate::models::event::Event::find_by_slug(&state.pool, &state.config.event_slug)
|
||
.await?
|
||
.ok_or_else(|| AppError::NotFound("Event nicht gefunden.".into()))?;
|
||
|
||
let released = event.export_released_at.is_some();
|
||
|
||
let jobs: Vec<(String, String, i16)> = sqlx::query_as(
|
||
"SELECT type::text, status::text, progress_pct FROM export_job WHERE event_id = $1",
|
||
)
|
||
.bind(event.id)
|
||
.fetch_all(&state.pool)
|
||
.await?;
|
||
|
||
let job_status = |type_name: &str| {
|
||
jobs.iter()
|
||
.find(|(t, _, _)| t == type_name)
|
||
.map(|(_, status, pct)| {
|
||
serde_json::json!({ "status": status, "progress_pct": pct })
|
||
})
|
||
.unwrap_or_else(|| serde_json::json!({ "status": "locked", "progress_pct": 0 }))
|
||
};
|
||
|
||
Ok(Json(serde_json::json!({
|
||
"released": released,
|
||
"zip": job_status("zip"),
|
||
"html": job_status("html"),
|
||
})))
|
||
}
|
||
|
||
/// Centralised guard for the export rate limit. Same pattern as upload/feed: master
|
||
/// switch + per-endpoint switch + numeric value, all stored in `config` and read on
|
||
/// each request.
|
||
async fn enforce_export_rate(state: &AppState, headers: &HeaderMap) -> Result<(), AppError> {
|
||
let rate_limits_on = config::get_bool(&state.config_cache, "rate_limits_enabled", true).await;
|
||
let export_rate_on = config::get_bool(&state.config_cache, "export_rate_enabled", true).await;
|
||
if !(rate_limits_on && export_rate_on) {
|
||
return Ok(());
|
||
}
|
||
let ip = client_ip(headers, "unknown");
|
||
let limit = config::get_usize(&state.config_cache, "export_rate_per_day", 3).await;
|
||
if !state
|
||
.rate_limiter
|
||
.check(format!("export:{ip}"), limit, Duration::from_secs(86400))
|
||
{
|
||
return Err(AppError::TooManyRequests(
|
||
"Zu viele Anfragen. Bitte warte kurz und versuche es erneut.".into(),
|
||
None,
|
||
));
|
||
}
|
||
Ok(())
|
||
}
|