diff --git a/tdkpin-rs/highscore-server/Cargo.toml b/tdkpin-rs/highscore-server/Cargo.toml index 0ea322f..6795d63 100644 --- a/tdkpin-rs/highscore-server/Cargo.toml +++ b/tdkpin-rs/highscore-server/Cargo.toml @@ -9,11 +9,12 @@ axum = "0.8" rusqlite = { version = "0.40", features = ["bundled"] } serde = { version = "1", features = ["derive"] } serde_json = "1" -tokio = { version = "1", features = ["macros", "net", "rt-multi-thread"] } +tokio = { version = "1", features = ["macros", "net", "rt-multi-thread", "sync"] } [dev-dependencies] http-body-util = "0.1" tempfile = "3" +tokio = { version = "1", features = ["time"] } tower = { version = "0.5", features = ["util"] } [lints.clippy] diff --git a/tdkpin-rs/highscore-server/README.md b/tdkpin-rs/highscore-server/README.md index 4aac5c6..b4a1ca6 100644 --- a/tdkpin-rs/highscore-server/README.md +++ b/tdkpin-rs/highscore-server/README.md @@ -23,9 +23,12 @@ TDKPIN_HIGHSCORE_DB=/var/lib/tdkpin/highscores.sqlite3 \ cargo run --manifest-path highscore-server/Cargo.toml ``` -Place [nginx.conf.example](nginx.conf.example) inside the public site's -existing `server` block. The browser client expects the API at -`/api/highscores` on the same origin as the game. +Copy the rate and connection zone declarations from +[nginx.conf.example](nginx.conf.example) into the existing `http` block, then +place its two `location` blocks inside the public site's `server` block. The +example bounds per-client and aggregate API traffic, request bodies, and proxy +waits. The browser client expects the API at `/api/highscores` on the same +origin as the game. The crate inherits the parent [`rustfmt.toml`](../rustfmt.toml); run `just fmt-highscore-server` when formatting it directly. diff --git a/tdkpin-rs/highscore-server/nginx.conf.example b/tdkpin-rs/highscore-server/nginx.conf.example index f243ddf..2c033eb 100644 --- a/tdkpin-rs/highscore-server/nginx.conf.example +++ b/tdkpin-rs/highscore-server/nginx.conf.example @@ -1,10 +1,31 @@ -# Add this block inside the nginx server block that serves the game. +# Add these directives inside the existing nginx http block. The server-wide +# zones bound aggregate traffic, while the address-keyed zones prevent one +# client from consuming the whole allowance. +limit_req_zone $binary_remote_addr zone=tdkpin_highscore_client_rate:10m rate=5r/s; +limit_req_zone $server_name zone=tdkpin_highscore_global_rate:1m rate=50r/s; +limit_conn_zone $binary_remote_addr zone=tdkpin_highscore_client_connections:10m; +limit_conn_zone $server_name zone=tdkpin_highscore_global_connections:1m; + +# Add these blocks inside the server block that serves the game. location /api/highscores { + limit_req zone=tdkpin_highscore_client_rate burst=10 nodelay; + limit_req zone=tdkpin_highscore_global_rate burst=25 nodelay; + limit_req_status 429; + limit_conn tdkpin_highscore_client_connections 10; + limit_conn tdkpin_highscore_global_connections 100; + limit_conn_status 429; + + client_max_body_size 1k; + client_body_timeout 5s; proxy_pass http://127.0.0.1:3000; proxy_http_version 1.1; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_connect_timeout 2s; + proxy_send_timeout 5s; + proxy_read_timeout 5s; + proxy_next_upstream off; } # Optional health check for local monitoring. diff --git a/tdkpin-rs/highscore-server/src/lib.rs b/tdkpin-rs/highscore-server/src/lib.rs index 6599c6b..9338861 100644 --- a/tdkpin-rs/highscore-server/src/lib.rs +++ b/tdkpin-rs/highscore-server/src/lib.rs @@ -6,16 +6,19 @@ use std::{ use axum::{ Json, Router, - extract::State, + extract::{DefaultBodyLimit, State}, http::StatusCode, response::{IntoResponse, Response}, routing::get, }; use rusqlite::{Connection, params, types::Type}; use serde::{Deserialize, Serialize}; +use tokio::{sync::Semaphore, task}; const MAX_HIGH_SCORES: usize = 10; const MAX_NAME_CHARS: usize = 21; +const MAX_SUBMISSION_BODY_BYTES: usize = 1024; +const MAX_CONCURRENT_DATABASE_OPERATIONS: usize = 1; #[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)] pub struct HighScore { @@ -127,6 +130,50 @@ impl HighScoreStore { } } +#[derive(Clone)] +struct AppState { + store: HighScoreStore, + database_slots: Arc, +} + +impl AppState { + fn new(store: HighScoreStore) -> Self { + Self { + store, + database_slots: Arc::new(Semaphore::new(MAX_CONCURRENT_DATABASE_OPERATIONS)), + } + } +} + +#[derive(Debug)] +enum DatabaseRequestError { + Busy, + Failed, +} + +async fn run_database_operation( + state: &AppState, + operation: F, +) -> Result +where + T: Send + 'static, + F: FnOnce(&HighScoreStore) -> Result + Send + 'static, +{ + let permit = state + .database_slots + .clone() + .try_acquire_owned() + .map_err(|_| DatabaseRequestError::Busy)?; + let store = state.store.clone(); + task::spawn_blocking(move || { + let _permit = permit; + operation(&store) + }) + .await + .map_err(|_| DatabaseRequestError::Failed)? + .map_err(|_| DatabaseRequestError::Failed) +} + #[derive(Debug, Deserialize)] struct SubmitRequest { name: String, @@ -146,6 +193,19 @@ fn invalid_request(message: &'static str) -> Response { .into_response() } +fn database_error_response(error: &DatabaseRequestError) -> Response { + match error { + DatabaseRequestError::Busy => ( + StatusCode::SERVICE_UNAVAILABLE, + Json(ErrorResponse { + error: "service is busy", + }), + ) + .into_response(), + DatabaseRequestError::Failed => StatusCode::INTERNAL_SERVER_ERROR.into_response(), + } +} + fn validate_request(request: &SubmitRequest) -> Result { let name = request.name.trim(); if name.is_empty() { @@ -167,39 +227,47 @@ async fn health() -> &'static str { "ok" } -async fn list_high_scores(State(store): State) -> Response { - match store.list() { +async fn list_high_scores(State(state): State) -> Response { + match run_database_operation(&state, HighScoreStore::list).await { Ok(scores) => Json(scores).into_response(), - Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), + Err(error) => database_error_response(&error), } } async fn submit_high_score( - State(store): State, + State(state): State, Json(request): Json, ) -> Response { let entry = match validate_request(&request) { Ok(entry) => entry, Err(message) => return invalid_request(message), }; - match store.submit(&entry) { + match run_database_operation(&state, move |store| store.submit(&entry)).await { Ok(scores) => (StatusCode::CREATED, Json(scores)).into_response(), - Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), + Err(error) => database_error_response(&error), } } pub fn router(store: HighScoreStore) -> Router { + router_with_state(AppState::new(store)) +} + +fn router_with_state(state: AppState) -> Router { Router::new() .route("/healthz", get(health)) .route( "/api/highscores", - get(list_high_scores).post(submit_high_score), + get(list_high_scores) + .post(submit_high_score) + .layer(DefaultBodyLimit::max(MAX_SUBMISSION_BODY_BYTES)), ) - .with_state(store) + .with_state(state) } #[cfg(test)] mod tests { + use std::{sync::mpsc, time::Duration}; + use axum::{ body::Body, http::{Request, StatusCode}, @@ -229,6 +297,44 @@ mod tests { .status() } + async fn request(app: &Router, request: Request) -> StatusCode { + app.clone() + .oneshot(request) + .await + .expect("router should respond") + .status() + } + + fn hold_database_lock( + store: HighScoreStore, + ) -> (mpsc::SyncSender<()>, std::thread::JoinHandle<()>) { + let (locked_sender, locked_receiver) = mpsc::sync_channel(0); + let (release_sender, release_receiver) = mpsc::sync_channel(0); + let lock_thread = std::thread::spawn(move || { + let _connection = store.0.lock().expect("database lock should succeed"); + locked_sender + .send(()) + .expect("test should observe the held database lock"); + release_receiver + .recv() + .expect("test should release the database lock"); + }); + locked_receiver + .recv() + .expect("database lock thread should start"); + (release_sender, lock_thread) + } + + async fn wait_until_database_is_busy(database_slots: &Semaphore) { + tokio::time::timeout(Duration::from_secs(1), async { + while database_slots.available_permits() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("first database request should acquire admission"); + } + async fn list(app: &Router) -> Vec { let response = app .clone() @@ -284,6 +390,134 @@ mod tests { assert_eq!(submit(&app, "ok\nno", 10).await, StatusCode::BAD_REQUEST); } + #[tokio::test] + async fn api_rejects_oversized_submission_bodies() { + let app = router(HighScoreStore::open_in_memory().expect("database should open")); + let body = serde_json::to_vec(&serde_json::json!({ + "name": "Player", + "score": 10, + "padding": "x".repeat(MAX_SUBMISSION_BODY_BYTES), + })) + .expect("request JSON should encode"); + + let status = request( + &app, + Request::post("/api/highscores") + .header("content-type", "application/json") + .body(Body::from(body)) + .expect("request should build"), + ) + .await; + + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + async fn busy_database_load_sheds_api_without_blocking_health() { + let store = HighScoreStore::open_in_memory().expect("database should open"); + let (release_sender, lock_thread) = hold_database_lock(store.clone()); + let state = AppState::new(store); + let database_slots = state.database_slots.clone(); + let app = router_with_state(state); + + let accepted_app = app.clone(); + let accepted = tokio::spawn(async move { submit(&accepted_app, "Player", 10).await }); + wait_until_database_is_busy(&database_slots).await; + + assert_eq!( + tokio::time::timeout( + Duration::from_millis(250), + request( + &app, + Request::get("/healthz") + .body(Body::empty()) + .expect("request should build"), + ), + ) + .await + .expect("health request should not wait for the database"), + StatusCode::OK + ); + assert_eq!( + submit(&app, "Other Player", 20).await, + StatusCode::SERVICE_UNAVAILABLE + ); + assert_eq!( + request( + &app, + Request::get("/api/highscores") + .body(Body::empty()) + .expect("request should build"), + ) + .await, + StatusCode::SERVICE_UNAVAILABLE + ); + assert_eq!( + request( + &app, + Request::builder() + .method("HEAD") + .uri("/api/highscores") + .body(Body::empty()) + .expect("request should build"), + ) + .await, + StatusCode::SERVICE_UNAVAILABLE + ); + + release_sender + .send(()) + .expect("database lock should be released"); + assert_eq!( + accepted.await.expect("accepted request should complete"), + StatusCode::CREATED + ); + lock_thread + .join() + .expect("database lock thread should stop"); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + async fn canceled_request_holds_admission_until_blocking_work_stops() { + let store = HighScoreStore::open_in_memory().expect("database should open"); + let (release_sender, lock_thread) = hold_database_lock(store.clone()); + let state = AppState::new(store); + let database_slots = state.database_slots.clone(); + let app = router_with_state(state); + + let canceled_app = app.clone(); + let canceled = tokio::spawn(async move { submit(&canceled_app, "Player", 10).await }); + wait_until_database_is_busy(&database_slots).await; + canceled.abort(); + assert!( + canceled + .await + .expect_err("request should be canceled") + .is_cancelled() + ); + + assert_eq!(database_slots.available_permits(), 0); + assert_eq!( + submit(&app, "Other Player", 20).await, + StatusCode::SERVICE_UNAVAILABLE + ); + + release_sender + .send(()) + .expect("database lock should be released"); + lock_thread + .join() + .expect("database lock thread should stop"); + tokio::time::timeout(Duration::from_secs(1), async { + while database_slots.available_permits() == 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("completed blocking work should release admission"); + assert_eq!(submit(&app, "Other Player", 20).await, StatusCode::CREATED); + } + #[tokio::test] async fn health_endpoint_is_available_for_nginx() { let app = router(HighScoreStore::open_in_memory().expect("database should open")); diff --git a/tdkpin-rs/web/storage.js b/tdkpin-rs/web/storage.js index 9970a64..3e6e9b2 100644 --- a/tdkpin-rs/web/storage.js +++ b/tdkpin-rs/web/storage.js @@ -10,6 +10,8 @@ let highScoreRequestInFlight = false; let highScoreGeneration = 0; let highScoreSubmissionWarningShown = false; + let highScoreRetryDelayMs = highScoreRetryInitialDelayMs; + let highScoreRetryAt = 0; function browserStorage() { try { @@ -86,10 +88,50 @@ } } + function retryAfterDelay(response) { + if (response.status !== 429 && response.status !== 503) { + return null; + } + const value = response.headers.get("Retry-After")?.trim(); + if (!value) { + return null; + } + + const seconds = Number(value); + if (Number.isFinite(seconds) && seconds >= 0) { + const delay = seconds * 1000; + return Number.isFinite(delay) ? delay : null; + } + + const timestamp = Date.parse(value); + return Number.isFinite(timestamp) + ? Math.max(0, timestamp - Date.now()) + : null; + } + + function scheduleHighScoreRetry(retryAfterMs) { + const jitteredDelay = highScoreRetryDelayMs * (0.5 + Math.random()); + const delay = + retryAfterMs === null ? jitteredDelay : retryAfterMs + jitteredDelay; + highScoreRetryAt = Date.now() + delay; + highScoreRetryDelayMs = Math.min( + highScoreRetryDelayMs * 2, + highScoreRetryMaxDelayMs, + ); + } + + function resetHighScoreRetry() { + highScoreRetryDelayMs = highScoreRetryInitialDelayMs; + highScoreRetryAt = 0; + } + async function flushHighScoreSubmission() { if (highScoreRequestInFlight) { return; } + if (Date.now() < highScoreRetryAt) { + return; + } const revision = wasm_exports.tdkpin_browser_high_score_revision(); if (revision === lastHighScoreRevision) { return; @@ -102,6 +144,7 @@ } highScoreRequestInFlight = true; + let retryAfterMs = null; try { const response = await fetch(highScoreApi, { method: "POST", @@ -110,6 +153,7 @@ keepalive: true, }); if (!response.ok) { + retryAfterMs = retryAfterDelay(response); throw new Error(`high-score submission failed (${response.status})`); } const json = await response.text(); @@ -118,7 +162,9 @@ wasm_exports.tdkpin_browser_high_score_ack(revision); lastHighScoreRevision = revision; highScoreSubmissionWarningShown = false; + resetHighScoreRetry(); } catch (error) { + scheduleHighScoreRetry(retryAfterMs); if (!highScoreSubmissionWarningShown) { console.warn("shared high-score submission failed; will retry", error); highScoreSubmissionWarningShown = true;