//! SQLite persistence: the event history, and the boards holding what is true right now. //! //! This is what lets the website read the past without asking the game, and what survives a sidecar //! restart. The event loop writes every live event here as it broadcasts it; REST reads query here //! instead of round-tripping the plugin. //! //! **The sidecar defines no schema for a frame's contents.** Events are persisted whole, as the //! JSON text that arrived, with only the columns it must *index* lifted out. That is the //! dumb-forwarder property doing real work: a protocol version that adds fields to an event needs //! no change here, and only a version that adds a new indexed column ever needs a migration. //! //! Protocol 2 is the first version that needed one — `server_id` and `wipe_id` (PROTOCOL.md §8.9) //! — and it is applied the way this project applies every schema change: as an `ALTER` guarded by a //! column check, never as an edit to the `CREATE`, because `CREATE TABLE IF NOT EXISTS` does //! nothing at all against a database that already has the table and an edited column would reach //! fresh installs only. use std::path::Path; use serde_json::{json, Value}; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; use sqlx::{Row, SqlitePool}; use tracing::{info, warn}; const SCHEMA: &str = " CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, t INTEGER NOT NULL, kind TEXT NOT NULL, server_id TEXT, wipe_id TEXT, json TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_events_kind_id ON events (kind, id DESC); CREATE INDEX IF NOT EXISTS idx_events_t ON events (t); -- Boards: current state, one row per kind, replaced whole. A board in chapter 4's sense — state -- with exactly one producer, re-sent on every connect — rather than a history. Protocol 1 had one -- of these hard-coded as `server_state`; protocol 2 has two and will have more, so the kind is a -- key rather than a table name. CREATE TABLE IF NOT EXISTS boards ( kind TEXT PRIMARY KEY, t INTEGER NOT NULL, json TEXT NOT NULL ); -- Protocol 1's single board. Kept so that a sidecar upgraded in place can carry its contents over -- (see `migrate`); nothing writes to it any more. CREATE TABLE IF NOT EXISTS server_state ( id INTEGER PRIMARY KEY CHECK (id = 1), t INTEGER NOT NULL, json TEXT NOT NULL ); "; /// The board holding the last `server.hello`. Named once here rather than spelled at four call /// sites, because it is the one board with a REST route of its own. pub const SERVER_BOARD: &str = "server.hello"; /// One row of the ingest feed: a stored event, with the identity a cursor needs. #[derive(Debug, Clone)] pub struct FeedItem { pub id: i64, pub t: i64, pub kind: String, pub frame: Value, } impl FeedItem { pub fn to_json(&self) -> Value { json!({ "id": self.id, "t": self.t, "kind": self.kind, "frame": self.frame }) } } #[derive(Clone)] pub struct Store { pool: SqlitePool, } impl Store { /// Opens (creating if absent) the SQLite database and ensures the schema exists. /// /// `path` is a filesystem path, handed to sqlx as one. It is deliberately **not** formatted /// into a `sqlite://` URL first: that spelling is parsed as a URL, so it percent-decodes the /// path and splits it on `?`. Under an installed layout the path is absolute and chosen by the /// operator — `C:\ProgramData\RunicGateway\rust-link.db`, or something under a home directory /// with a `%` or `#` in it — and a URL round-trip silently opens a *different* file. pub async fn open(path: &str) -> anyhow::Result { // A service unit can name a data directory that does not exist yet; creating it here means // one less way for a fresh install to fail on first start. if let Some(dir) = Path::new(path).parent() { if !dir.as_os_str().is_empty() && !dir.exists() { std::fs::create_dir_all(dir)?; } } let opts = SqliteConnectOptions::new() .filename(path) .create_if_missing(true); // An in-memory database is **per connection**, not per process: every connection the pool // opens gets its own empty one, so a second pooled connection finds none of the schema the // first created. It presents as `no such table` from a random subset of queries, which is // as confusing a failure as this file has. A single connection is the only coherent // reading of `:memory:`, and it is what makes it usable at all. let max = if is_in_memory(path) { 1 } else { 4 }; let pool = SqlitePoolOptions::new() .max_connections(max) .connect_with(opts) .await?; sqlx::query(SCHEMA).execute(&pool).await?; let store = Self { pool }; store.migrate().await?; info!(%path, "store ready"); Ok(store) } /// Brings a database created by an older protocol up to this one. /// /// Two steps, both idempotent, both safe to run on a fresh database where they do nothing: /// add the columns protocol 2 indexes on, and carry protocol 1's single board into `boards`. /// /// The board carry-over matters more than it looks: without it, an upgraded sidecar answers /// `204` for `/server` until the game next connects, and the website reads that as *this /// server has never been heard from* — the site loses a server it has been rendering for /// weeks, at the exact moment somebody upgraded the bridge. async fn migrate(&self) -> anyhow::Result<()> { for (column, ddl) in [ ("server_id", "ALTER TABLE events ADD COLUMN server_id TEXT"), ("wipe_id", "ALTER TABLE events ADD COLUMN wipe_id TEXT"), ] { if !self.has_column("events", column).await? { sqlx::query(ddl).execute(&self.pool).await?; info!(column, "events: column added"); } } // Indexed after the columns exist, and in the same idempotent spirit. sqlx::query("CREATE INDEX IF NOT EXISTS idx_events_wipe_id ON events (wipe_id, id DESC)") .execute(&self.pool) .await?; let carried: Option<(i64, String)> = sqlx::query_as( "SELECT t, json FROM server_state WHERE id = 1 AND NOT EXISTS (SELECT 1 FROM boards WHERE kind = ?)", ) .bind(SERVER_BOARD) .fetch_optional(&self.pool) .await?; if let Some((t, json)) = carried { self.put_board(SERVER_BOARD, t, &json).await?; info!("carried the protocol 1 server board into boards"); } Ok(()) } async fn has_column(&self, table: &str, column: &str) -> anyhow::Result { // `PRAGMA table_info` does not take a bind parameter for the table name, which is why this // is formatted. Both call sites pass a literal; nothing here is reachable from a request. let rows = sqlx::query(&format!("PRAGMA table_info({table})")) .fetch_all(&self.pool) .await?; Ok(rows .iter() .any(|r| r.get::("name").eq_ignore_ascii_case(column))) } /// Cheap liveness check for the health endpoint. pub async fn ping(&self) -> anyhow::Result<()> { sqlx::query("SELECT 1").execute(&self.pool).await?; Ok(()) } /// Appends one live event. Failures are logged by the caller; persistence must never block the /// live feed. /// /// `server_id` and `wipe_id` are lifted out of the frame by the caller and stored as columns as /// well as remaining in the JSON. Duplicated deliberately: the column is what an index and a /// `WHERE` can reach, and the JSON is what stays correct when the columns change. pub async fn insert_event( &self, t: i64, kind: &str, server_id: Option<&str>, wipe_id: Option<&str>, json: &str, ) -> anyhow::Result<()> { sqlx::query( "INSERT INTO events (t, kind, server_id, wipe_id, json) VALUES (?, ?, ?, ?, ?)", ) .bind(t) .bind(kind) .bind(server_id) .bind(wipe_id) .bind(json) .execute(&self.pool) .await?; Ok(()) } /// Most-recent events, **newest first**, optionally filtered by kind and by wipe. /// /// For a human, an admin screen, or a point-in-time look. A consumer that must not miss a row /// wants [`Store::feed`] instead — see its documentation for why these are two functions and /// not one with a flag. pub async fn recent( &self, kind: Option<&str>, wipe: Option<&str>, limit: i64, ) -> anyhow::Result> { let limit = limit.clamp(1, 1000); // Built rather than branched four ways: two optional filters is four combinations, and the // fourth is always the one nobody tested. The bindings stay parameterised. let mut sql = String::from("SELECT json FROM events WHERE 1 = 1"); if kind.is_some() { sql.push_str(" AND kind = ?"); } if wipe.is_some() { sql.push_str(" AND wipe_id = ?"); } sql.push_str(" ORDER BY id DESC LIMIT ?"); let mut query = sqlx::query(&sql); if let Some(k) = kind { query = query.bind(k); } if let Some(w) = wipe { query = query.bind(w); } let rows = query.bind(limit).fetch_all(&self.pool).await?; Ok(parse_json_column(rows)) } /// The ingest cursor: events **after** `since`, **oldest first**. /// /// This is a separate function from [`Store::recent`], and the route on top of it is a separate /// route, for one reason: a single route whose ordering depends on a query parameter serves the /// other ordering to every caller that forgets it, and for the ingesting caller that means /// advancing its cursor past rows it never read. Silently, and once per deployment mistake. /// /// Returns the page and whether it filled — a consumer an hour behind drains at its own pace /// rather than guessing from a count. pub async fn feed(&self, since: i64, limit: i64) -> anyhow::Result<(Vec, bool)> { let limit = limit.clamp(1, 1000); let rows = sqlx::query( "SELECT id, t, kind, json FROM events WHERE id > ? ORDER BY id ASC LIMIT ?", ) .bind(since) .bind(limit) .fetch_all(&self.pool) .await?; let more = rows.len() as i64 == limit; let items = rows .into_iter() .filter_map(|r| { let json: String = r.get("json"); serde_json::from_str(&json).ok().map(|frame| FeedItem { id: r.get("id"), t: r.get("t"), kind: r.get("kind"), frame, }) }) .collect(); Ok((items, more)) } /// The highest event id in the store, or 0 when it is empty. /// /// A consumer starting from nothing uses this to begin at the *end* rather than replaying the /// whole history it has no use for — a fresh module against a sidecar that has been running for /// a month wants what happens next, not a fortnight of deaths. pub async fn last_event_id(&self) -> anyhow::Result { let row = sqlx::query("SELECT COALESCE(MAX(id), 0) AS id FROM events") .fetch_one(&self.pool) .await?; Ok(row.get("id")) } /// Replaces one board. Called for every snapshot frame, which the plugin re-sends on every /// connect and on a cadence — so this is an upsert by construction, not by accident. pub async fn put_board(&self, kind: &str, t: i64, json: &str) -> anyhow::Result<()> { sqlx::query( "INSERT INTO boards (kind, t, json) VALUES (?, ?, ?) ON CONFLICT(kind) DO UPDATE SET t = excluded.t, json = excluded.json", ) .bind(kind) .bind(t) .bind(json) .execute(&self.pool) .await?; Ok(()) } /// One board, or `None` if the game has never sent it. /// /// This is the read that makes the website render while the game is off, which is the whole /// reason the sidecar holds a database at all. pub async fn board(&self, kind: &str) -> anyhow::Result> { let row = sqlx::query("SELECT json FROM boards WHERE kind = ?") .bind(kind) .fetch_optional(&self.pool) .await?; Ok(row.and_then(|r| serde_json::from_str(&r.get::("json")).ok())) } /// Every board, keyed by kind. What a consumer reads once on connect to know the present /// before it starts following the story. pub async fn all_boards(&self) -> anyhow::Result> { let rows = sqlx::query("SELECT kind, json FROM boards ORDER BY kind") .fetch_all(&self.pool) .await?; let mut out = serde_json::Map::new(); for row in rows { let kind: String = row.get("kind"); if let Ok(v) = serde_json::from_str::(&row.get::("json")) { out.insert(kind, v); } } Ok(out) } /// Deletes events older than `retain_days`, returning how many went. /// /// **Boards are never pruned**, and that asymmetry is the design rather than an oversight: a /// board is one row per kind holding what is true now, and deleting it would make a server the /// site has rendered for weeks look like one that has never connected. History is bounded /// because it grows; the present is not, because it does not. /// /// Safe to be aggressive here because the *permanent* record lives on the website — per-wipe /// rollups in the module's own tables (R12) — and this database sits on a game host whose disk /// belongs to the operator. pub async fn prune(&self, retain_days: i64) -> anyhow::Result { if retain_days <= 0 { return Ok(0); // retention off; an operator who wants everything keeps everything } let cutoff = now_ms() - retain_days * 86_400_000; let done = sqlx::query("DELETE FROM events WHERE t < ?") .bind(cutoff) .execute(&self.pool) .await?; let n = done.rows_affected(); if n > 0 { info!(pruned = n, retain_days, "pruned old events"); } Ok(n) } } fn now_ms() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_millis() as i64) .unwrap_or(0) } /// Whether this path names an in-memory database rather than a file. Covers the bare `:memory:` /// spelling and the `file:` URI form that carries `mode=memory`. fn is_in_memory(path: &str) -> bool { path == ":memory:" || (path.starts_with("file:") && path.contains("mode=memory")) } fn parse_json_column(rows: Vec) -> Vec { rows.into_iter() .filter_map(|r| serde_json::from_str(&r.get::("json")).ok()) .collect() } /// Warns once about a store write that failed. Persistence failures must never stop the live feed, /// so every caller logs and carries on; this keeps them saying the same thing. pub fn warn_write(what: &str, e: &anyhow::Error) { warn!(error = %e, "{what}"); } #[cfg(test)] mod tests { use super::*; use serde_json::json; async fn store() -> Store { Store::open(":memory:").await.unwrap() } async fn insert(s: &Store, t: i64, kind: &str, wipe: Option<&str>) { s.insert_event( t, kind, Some("main"), wipe, &json!({"kind": kind, "t": t}).to_string(), ) .await .unwrap(); } /// The trap this file's pool sizing exists for: a multi-connection pool over `:memory:` hands /// out empty databases. Asserting the *pool* is what makes the reason visible; asserting only /// that a query works would pass again the moment someone "tidied" the sizing back. #[tokio::test] async fn an_in_memory_store_uses_exactly_one_connection() { assert!(is_in_memory(":memory:")); assert!(is_in_memory("file:x?mode=memory&cache=shared")); assert!(!is_in_memory("rust-link.db")); assert!(!is_in_memory("file:/var/lib/rg/rust-link.db")); let s = store().await; assert_eq!(s.pool.options().get_max_connections(), 1); } /// `sqlx::query` over a multi-statement string is the kind of thing that quietly runs only the /// first statement. Every table and both reads have to work on a real file, under the pool size /// production uses. #[tokio::test] async fn the_whole_schema_is_created_on_a_pooled_file_store() { let dir = std::env::temp_dir().join(format!("rust-link-test-{}", std::process::id())); let path = dir.join("schema.db"); let _ = std::fs::remove_dir_all(&dir); let s = Store::open(path.to_str().unwrap()).await.unwrap(); assert_eq!(s.pool.options().get_max_connections(), 4); insert(&s, 1, "k", None).await; s.put_board(SERVER_BOARD, 1, "{}").await.unwrap(); assert_eq!(s.recent(None, None, 10).await.unwrap().len(), 1); assert!(s.board(SERVER_BOARD).await.unwrap().is_some()); drop(s); let _ = std::fs::remove_dir_all(&dir); } #[tokio::test] async fn events_come_back_newest_first_and_filter_by_kind() { let s = store().await; insert(&s, 1, "server.hello", None).await; insert(&s, 2, "other", None).await; insert(&s, 3, "server.hello", None).await; let all = s.recent(None, None, 10).await.unwrap(); assert_eq!(all.len(), 3); assert_eq!(all[0]["t"], 3); let hellos = s.recent(Some("server.hello"), None, 10).await.unwrap(); assert_eq!(hellos.len(), 2); assert_eq!(hellos[0]["t"], 3); } /// R12 in one test: a wipe splits the history without erasing any of it. #[tokio::test] async fn events_filter_by_wipe_without_losing_the_other_wipe() { let s = store().await; insert(&s, 1, "player.death", Some("w-a")).await; insert(&s, 2, "player.death", Some("w-a")).await; insert(&s, 3, "player.death", Some("w-b")).await; assert_eq!(s.recent(None, Some("w-a"), 10).await.unwrap().len(), 2); assert_eq!(s.recent(None, Some("w-b"), 10).await.unwrap().len(), 1); assert_eq!(s.recent(None, None, 10).await.unwrap().len(), 3); // Both filters at once is the combination that is easy to build wrong. assert_eq!( s.recent(Some("player.death"), Some("w-b"), 10) .await .unwrap() .len(), 1 ); } /// The cursor's two properties, and they are the ones a consumer's correctness rests on: /// oldest first, and strictly after the id it was given. #[tokio::test] async fn the_feed_is_a_cursor_and_runs_oldest_first() { let s = store().await; for i in 1..=5 { insert(&s, i, "player.death", None).await; } let (page, more) = s.feed(0, 2).await.unwrap(); assert_eq!(page.len(), 2); assert!(more, "a full page must say there is more"); assert_eq!(page[0].t, 1); assert_eq!(page[1].t, 2); let (page, more) = s.feed(page[1].id, 10).await.unwrap(); assert_eq!(page.len(), 3); assert!(!more, "a short page is the end of the queue"); assert_eq!(page[0].t, 3); // The cursor is exclusive; re-reading from the last id delivers nothing twice. let (page, _) = s.feed(page[2].id, 10).await.unwrap(); assert!(page.is_empty()); } #[tokio::test] async fn a_fresh_consumer_can_start_at_the_end() { let s = store().await; assert_eq!(s.last_event_id().await.unwrap(), 0); for i in 1..=3 { insert(&s, i, "k", None).await; } let last = s.last_event_id().await.unwrap(); assert_eq!(last, 3); assert!(s.feed(last, 10).await.unwrap().0.is_empty()); } /// An absent board reads as `None`, not as an empty object. A caller must be able to tell /// "the game has never connected" from "the game connected and said nothing" — collapsing the /// two is how a site ends up rendering a server that does not exist. #[tokio::test] async fn a_board_is_absent_until_a_snapshot_arrives() { let s = store().await; assert!(s.board(SERVER_BOARD).await.unwrap().is_none()); s.put_board(SERVER_BOARD, 1, &json!({"serverId": "main"}).to_string()) .await .unwrap(); assert_eq!( s.board(SERVER_BOARD).await.unwrap().unwrap()["serverId"], "main" ); } /// A board holds exactly one row however many times the plugin reconnects, and the boards are /// independent of one another. #[tokio::test] async fn a_second_snapshot_replaces_the_first_of_its_own_kind_only() { let s = store().await; s.put_board(SERVER_BOARD, 1, &json!({"bootId": "a"}).to_string()) .await .unwrap(); s.put_board(SERVER_BOARD, 2, &json!({"bootId": "b"}).to_string()) .await .unwrap(); s.put_board("players.online", 2, &json!({"count": 4}).to_string()) .await .unwrap(); assert_eq!(s.board(SERVER_BOARD).await.unwrap().unwrap()["bootId"], "b"); assert_eq!( s.board("players.online").await.unwrap().unwrap()["count"], 4 ); let all = s.all_boards().await.unwrap(); assert_eq!(all.len(), 2); } #[tokio::test] async fn the_limit_is_clamped_rather_than_trusted() { let s = store().await; for i in 0..5 { insert(&s, i, "k", None).await; } // 0 and negatives would otherwise mean "no rows" and "SQLite's unlimited" respectively. assert_eq!(s.recent(None, None, 0).await.unwrap().len(), 1); assert_eq!(s.recent(None, None, -1).await.unwrap().len(), 1); assert_eq!(s.recent(None, None, 100_000).await.unwrap().len(), 5); assert_eq!(s.feed(0, 0).await.unwrap().0.len(), 1); } /// Retention deletes history and leaves the present alone. The second half is the half worth /// asserting: a pruned board is a server that has "never connected". #[tokio::test] async fn pruning_bounds_the_history_and_never_touches_a_board() { let s = store().await; let old = now_ms() - 30 * 86_400_000; insert(&s, old, "player.death", None).await; insert(&s, now_ms(), "player.death", None).await; s.put_board(SERVER_BOARD, old, &json!({"serverId": "main"}).to_string()) .await .unwrap(); assert_eq!(s.prune(14).await.unwrap(), 1); assert_eq!(s.recent(None, None, 10).await.unwrap().len(), 1); assert!(s.board(SERVER_BOARD).await.unwrap().is_some()); // Retention off keeps everything, which is a supported configuration rather than a bug. insert(&s, old, "player.death", None).await; assert_eq!(s.prune(0).await.unwrap(), 0); assert_eq!(s.recent(None, None, 10).await.unwrap().len(), 2); } /// The upgrade path, on a real file because that is the only place it can happen: a protocol 1 /// database has `server_state` and no `wipe_id`, and opening it with this build must produce a /// store that still knows which server it is holding. #[tokio::test] async fn a_protocol_1_database_is_migrated_in_place() { let dir = std::env::temp_dir().join(format!("rust-link-migrate-{}", std::process::id())); let path = dir.join("old.db"); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(&dir).unwrap(); // Exactly protocol 1's schema, written by hand so the test does not depend on this file // still being able to produce it. let opts = SqliteConnectOptions::new() .filename(&path) .create_if_missing(true); let pool = SqlitePoolOptions::new() .max_connections(1) .connect_with(opts) .await .unwrap(); sqlx::query( "CREATE TABLE events (id INTEGER PRIMARY KEY AUTOINCREMENT, t INTEGER NOT NULL, kind TEXT NOT NULL, json TEXT NOT NULL); CREATE TABLE server_state (id INTEGER PRIMARY KEY CHECK (id = 1), t INTEGER NOT NULL, json TEXT NOT NULL);", ) .execute(&pool) .await .unwrap(); sqlx::query("INSERT INTO events (t, kind, json) VALUES (1, 'server.hello', '{\"t\":1}')") .execute(&pool) .await .unwrap(); sqlx::query( "INSERT INTO server_state (id, t, json) VALUES (1, 1, '{\"serverId\":\"legacy\"}')", ) .execute(&pool) .await .unwrap(); pool.close().await; let s = Store::open(path.to_str().unwrap()).await.unwrap(); // The columns arrived, the old rows survived with them empty, and the board came across — // so an upgraded sidecar does not report a server it has been serving for weeks as one it // has never heard of. assert!(s.has_column("events", "wipe_id").await.unwrap()); assert_eq!(s.recent(None, None, 10).await.unwrap().len(), 1); assert_eq!( s.board(SERVER_BOARD).await.unwrap().unwrap()["serverId"], "legacy" ); // And it is idempotent: opening again must not fail on an ALTER that already ran. drop(s); let again = Store::open(path.to_str().unwrap()).await.unwrap(); assert!(again.board(SERVER_BOARD).await.unwrap().is_some()); drop(again); let _ = std::fs::remove_dir_all(&dir); } }