Merge branch 'tomerge'

This commit is contained in:
2026-08-29 20:11:13 +02:00
5 changed files with 319 additions and 14 deletions
+2 -1
View File
@@ -9,11 +9,12 @@ axum = "0.8"
rusqlite = { version = "0.40", features = ["bundled"] } rusqlite = { version = "0.40", features = ["bundled"] }
serde = { version = "1", features = ["derive"] } serde = { version = "1", features = ["derive"] }
serde_json = "1" serde_json = "1"
tokio = { version = "1", features = ["macros", "net", "rt-multi-thread"] } tokio = { version = "1", features = ["macros", "net", "rt-multi-thread", "sync"] }
[dev-dependencies] [dev-dependencies]
http-body-util = "0.1" http-body-util = "0.1"
tempfile = "3" tempfile = "3"
tokio = { version = "1", features = ["time"] }
tower = { version = "0.5", features = ["util"] } tower = { version = "0.5", features = ["util"] }
[lints.clippy] [lints.clippy]
+6 -3
View File
@@ -23,9 +23,12 @@ TDKPIN_HIGHSCORE_DB=/var/lib/tdkpin/highscores.sqlite3 \
cargo run --manifest-path highscore-server/Cargo.toml cargo run --manifest-path highscore-server/Cargo.toml
``` ```
Place [nginx.conf.example](nginx.conf.example) inside the public site's Copy the rate and connection zone declarations from
existing `server` block. The browser client expects the API at [nginx.conf.example](nginx.conf.example) into the existing `http` block, then
`/api/highscores` on the same origin as the game. 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 The crate inherits the parent [`rustfmt.toml`](../rustfmt.toml); run
`just fmt-highscore-server` when formatting it directly. `just fmt-highscore-server` when formatting it directly.
+22 -1
View File
@@ -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 { 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_pass http://127.0.0.1:3000;
proxy_http_version 1.1; proxy_http_version 1.1;
proxy_set_header Host $host; proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; 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. # Optional health check for local monitoring.
+243 -9
View File
@@ -6,16 +6,19 @@ use std::{
use axum::{ use axum::{
Json, Json,
Router, Router,
extract::State, extract::{DefaultBodyLimit, State},
http::StatusCode, http::StatusCode,
response::{IntoResponse, Response}, response::{IntoResponse, Response},
routing::get, routing::get,
}; };
use rusqlite::{Connection, params, types::Type}; use rusqlite::{Connection, params, types::Type};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use tokio::{sync::Semaphore, task};
const MAX_HIGH_SCORES: usize = 10; const MAX_HIGH_SCORES: usize = 10;
const MAX_NAME_CHARS: usize = 21; 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)] #[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)]
pub struct HighScore { pub struct HighScore {
@@ -127,6 +130,50 @@ impl HighScoreStore {
} }
} }
#[derive(Clone)]
struct AppState {
store: HighScoreStore,
database_slots: Arc<Semaphore>,
}
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<T, F>(
state: &AppState,
operation: F,
) -> Result<T, DatabaseRequestError>
where
T: Send + 'static,
F: FnOnce(&HighScoreStore) -> Result<T, StoreError> + 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)] #[derive(Debug, Deserialize)]
struct SubmitRequest { struct SubmitRequest {
name: String, name: String,
@@ -146,6 +193,19 @@ fn invalid_request(message: &'static str) -> Response {
.into_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<HighScore, &'static str> { fn validate_request(request: &SubmitRequest) -> Result<HighScore, &'static str> {
let name = request.name.trim(); let name = request.name.trim();
if name.is_empty() { if name.is_empty() {
@@ -167,39 +227,47 @@ async fn health() -> &'static str {
"ok" "ok"
} }
async fn list_high_scores(State(store): State<HighScoreStore>) -> Response { async fn list_high_scores(State(state): State<AppState>) -> Response {
match store.list() { match run_database_operation(&state, HighScoreStore::list).await {
Ok(scores) => Json(scores).into_response(), Ok(scores) => Json(scores).into_response(),
Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), Err(error) => database_error_response(&error),
} }
} }
async fn submit_high_score( async fn submit_high_score(
State(store): State<HighScoreStore>, State(state): State<AppState>,
Json(request): Json<SubmitRequest>, Json(request): Json<SubmitRequest>,
) -> Response { ) -> Response {
let entry = match validate_request(&request) { let entry = match validate_request(&request) {
Ok(entry) => entry, Ok(entry) => entry,
Err(message) => return invalid_request(message), 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(), 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 { pub fn router(store: HighScoreStore) -> Router {
router_with_state(AppState::new(store))
}
fn router_with_state(state: AppState) -> Router {
Router::new() Router::new()
.route("/healthz", get(health)) .route("/healthz", get(health))
.route( .route(
"/api/highscores", "/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)] #[cfg(test)]
mod tests { mod tests {
use std::{sync::mpsc, time::Duration};
use axum::{ use axum::{
body::Body, body::Body,
http::{Request, StatusCode}, http::{Request, StatusCode},
@@ -229,6 +297,44 @@ mod tests {
.status() .status()
} }
async fn request(app: &Router, request: Request<Body>) -> 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<HighScore> { async fn list(app: &Router) -> Vec<HighScore> {
let response = app let response = app
.clone() .clone()
@@ -284,6 +390,134 @@ mod tests {
assert_eq!(submit(&app, "ok\nno", 10).await, StatusCode::BAD_REQUEST); 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] #[tokio::test]
async fn health_endpoint_is_available_for_nginx() { async fn health_endpoint_is_available_for_nginx() {
let app = router(HighScoreStore::open_in_memory().expect("database should open")); let app = router(HighScoreStore::open_in_memory().expect("database should open"));
+46
View File
@@ -10,6 +10,8 @@
let highScoreRequestInFlight = false; let highScoreRequestInFlight = false;
let highScoreGeneration = 0; let highScoreGeneration = 0;
let highScoreSubmissionWarningShown = false; let highScoreSubmissionWarningShown = false;
let highScoreRetryDelayMs = highScoreRetryInitialDelayMs;
let highScoreRetryAt = 0;
function browserStorage() { function browserStorage() {
try { 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() { async function flushHighScoreSubmission() {
if (highScoreRequestInFlight) { if (highScoreRequestInFlight) {
return; return;
} }
if (Date.now() < highScoreRetryAt) {
return;
}
const revision = wasm_exports.tdkpin_browser_high_score_revision(); const revision = wasm_exports.tdkpin_browser_high_score_revision();
if (revision === lastHighScoreRevision) { if (revision === lastHighScoreRevision) {
return; return;
@@ -102,6 +144,7 @@
} }
highScoreRequestInFlight = true; highScoreRequestInFlight = true;
let retryAfterMs = null;
try { try {
const response = await fetch(highScoreApi, { const response = await fetch(highScoreApi, {
method: "POST", method: "POST",
@@ -110,6 +153,7 @@
keepalive: true, keepalive: true,
}); });
if (!response.ok) { if (!response.ok) {
retryAfterMs = retryAfterDelay(response);
throw new Error(`high-score submission failed (${response.status})`); throw new Error(`high-score submission failed (${response.status})`);
} }
const json = await response.text(); const json = await response.text();
@@ -118,7 +162,9 @@
wasm_exports.tdkpin_browser_high_score_ack(revision); wasm_exports.tdkpin_browser_high_score_ack(revision);
lastHighScoreRevision = revision; lastHighScoreRevision = revision;
highScoreSubmissionWarningShown = false; highScoreSubmissionWarningShown = false;
resetHighScoreRetry();
} catch (error) { } catch (error) {
scheduleHighScoreRetry(retryAfterMs);
if (!highScoreSubmissionWarningShown) { if (!highScoreSubmissionWarningShown) {
console.warn("shared high-score submission failed; will retry", error); console.warn("shared high-score submission failed; will retry", error);
highScoreSubmissionWarningShown = true; highScoreSubmissionWarningShown = true;