The release workflow's `cargo fmt --check` gate failed on the Protocol 2.0 additions: upsert_guild/upsert_house signatures and their call sites in main.rs exceeded rustfmt's default width. Reformatted with `cargo fmt` — no behavioral change. `cargo fmt --check`, `cargo check`, and `cargo test` all clean. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
396 lines
13 KiB
Rust
396 lines
13 KiB
Rust
//! SQLite persistence: event history, the economy series, cached profiles, and the account-link
|
|
//! map. This is what lets the website read the past without asking the shard, and what survives a
|
|
//! sidecar restart.
|
|
//!
|
|
//! The event loop writes every live event here as it broadcasts it; REST read endpoints query here
|
|
//! instead of round-tripping the shard. Links and profiles are written from the REST reply paths
|
|
//! (`link.ok`, `char.profile`), which are RPC replies and never hit the broadcast stream.
|
|
|
|
use std::str::FromStr;
|
|
|
|
use serde_json::Value;
|
|
use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
|
|
use sqlx::{Row, SqlitePool};
|
|
use tracing::info;
|
|
|
|
#[derive(Clone)]
|
|
pub struct Store {
|
|
pool: SqlitePool,
|
|
}
|
|
|
|
impl Store {
|
|
/// Opens (creating if absent) the SQLite database and ensures the schema exists.
|
|
pub async fn open(path: &str) -> anyhow::Result<Self> {
|
|
let opts =
|
|
SqliteConnectOptions::from_str(&format!("sqlite://{path}"))?.create_if_missing(true);
|
|
|
|
let pool = SqlitePoolOptions::new()
|
|
.max_connections(4)
|
|
.connect_with(opts)
|
|
.await?;
|
|
|
|
sqlx::query(SCHEMA).execute(&pool).await?;
|
|
info!(%path, "store ready");
|
|
Ok(Self { pool })
|
|
}
|
|
|
|
/// 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.
|
|
pub async fn insert_event(&self, t: i64, kind: &str, json: &str) -> anyhow::Result<()> {
|
|
sqlx::query("INSERT INTO events (t, kind, json) VALUES (?, ?, ?)")
|
|
.bind(t)
|
|
.bind(kind)
|
|
.bind(json)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Most-recent events, newest first, optionally filtered by kind.
|
|
pub async fn recent(&self, kind: Option<&str>, limit: i64) -> anyhow::Result<Vec<Value>> {
|
|
let limit = limit.clamp(1, 1000);
|
|
let rows = match kind {
|
|
Some(k) => {
|
|
sqlx::query("SELECT json FROM events WHERE kind = ? ORDER BY id DESC LIMIT ?")
|
|
.bind(k)
|
|
.bind(limit)
|
|
.fetch_all(&self.pool)
|
|
.await?
|
|
}
|
|
None => {
|
|
sqlx::query("SELECT json FROM events ORDER BY id DESC LIMIT ?")
|
|
.bind(limit)
|
|
.fetch_all(&self.pool)
|
|
.await?
|
|
}
|
|
};
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
|
|
/// The money-supply series: the `economy.supply` events, newest first.
|
|
pub async fn economy(&self, limit: i64) -> anyhow::Result<Vec<Value>> {
|
|
self.recent(Some("economy.supply"), limit).await
|
|
}
|
|
|
|
pub async fn record_link(
|
|
&self,
|
|
account: &str,
|
|
website_user_id: &str,
|
|
t: i64,
|
|
) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO links (account, website_user_id, linked_t) VALUES (?, ?, ?)
|
|
ON CONFLICT(account) DO UPDATE SET website_user_id = excluded.website_user_id, linked_t = excluded.linked_t",
|
|
)
|
|
.bind(account)
|
|
.bind(website_user_id)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn get_link(&self, account: &str) -> anyhow::Result<Option<String>> {
|
|
let row = sqlx::query("SELECT website_user_id FROM links WHERE account = ?")
|
|
.bind(account)
|
|
.fetch_optional(&self.pool)
|
|
.await?;
|
|
Ok(row.map(|r| r.get::<String, _>("website_user_id")))
|
|
}
|
|
|
|
/// Drops the mirrored link row so event attribution stops immediately, without waiting on the
|
|
/// shard. Returns the number of rows removed (0 if the account was not linked here).
|
|
pub async fn record_unlink(&self, account: &str) -> anyhow::Result<u64> {
|
|
let res = sqlx::query("DELETE FROM links WHERE account = ?")
|
|
.bind(account)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(res.rows_affected())
|
|
}
|
|
|
|
pub async fn cache_profile(
|
|
&self,
|
|
serial: &str,
|
|
account: Option<&str>,
|
|
name: Option<&str>,
|
|
json: &str,
|
|
t: i64,
|
|
) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO profiles (serial, account, name, json, updated_t) VALUES (?, ?, ?, ?, ?)
|
|
ON CONFLICT(serial) DO UPDATE SET account = excluded.account, name = excluded.name, json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(serial)
|
|
.bind(account)
|
|
.bind(name)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn get_cached_profile(&self, serial: &str) -> anyhow::Result<Option<Value>> {
|
|
let row = sqlx::query("SELECT json FROM profiles WHERE serial = ?")
|
|
.bind(serial)
|
|
.fetch_optional(&self.pool)
|
|
.await?;
|
|
Ok(row.and_then(|r| serde_json::from_str(&r.get::<String, _>("json")).ok()))
|
|
}
|
|
|
|
/// Upserts one champion-spawn's latest state, keyed by serial. Fed from the `champ.update`
|
|
/// stream; this table is the live board the website reads, so there is exactly one row per
|
|
/// spawn and it always holds the most recent snapshot.
|
|
pub async fn upsert_champ(
|
|
&self,
|
|
serial: &str,
|
|
status: Option<&str>,
|
|
name: Option<&str>,
|
|
json: &str,
|
|
t: i64,
|
|
) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO champs (serial, status, name, json, updated_t) VALUES (?, ?, ?, ?, ?)
|
|
ON CONFLICT(serial) DO UPDATE SET status = excluded.status, name = excluded.name, json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(serial)
|
|
.bind(status)
|
|
.bind(name)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Drops one spawn from the board. Fed from the `champ.remove` stream: a controller that was
|
|
/// deleted, or a transient sea boss that was slain, leaves the board this way.
|
|
pub async fn delete_champ(&self, serial: &str) -> anyhow::Result<()> {
|
|
sqlx::query("DELETE FROM champs WHERE serial = ?")
|
|
.bind(serial)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// The full champion-spawn board: every spawn's latest snapshot. Ordered by name so the site
|
|
/// gets a stable list.
|
|
pub async fn champs_all(&self) -> anyhow::Result<Vec<Value>> {
|
|
let rows = sqlx::query("SELECT json FROM champs ORDER BY name, serial")
|
|
.fetch_all(&self.pool)
|
|
.await?;
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
|
|
// ---- guild board (Protocol 2.0) ----
|
|
|
|
/// Upserts one guild's latest state, keyed by guild id. Fed from `guild.update`; one row per
|
|
/// guild, always the most recent snapshot. This is the board the website reads on load.
|
|
pub async fn upsert_guild(
|
|
&self,
|
|
id: i64,
|
|
name: Option<&str>,
|
|
json: &str,
|
|
t: i64,
|
|
) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO guilds (id, name, json, updated_t) VALUES (?, ?, ?, ?)
|
|
ON CONFLICT(id) DO UPDATE SET name = excluded.name, json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(id)
|
|
.bind(name)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Drops one guild from the board. Fed from `guild.remove` (a disband or a removed guild).
|
|
pub async fn delete_guild(&self, id: i64) -> anyhow::Result<()> {
|
|
sqlx::query("DELETE FROM guilds WHERE id = ?")
|
|
.bind(id)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// The full guild board: every guild's latest snapshot, ordered by name.
|
|
pub async fn guilds_all(&self) -> anyhow::Result<Vec<Value>> {
|
|
let rows = sqlx::query("SELECT json FROM guilds ORDER BY name, id")
|
|
.fetch_all(&self.pool)
|
|
.await?;
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
|
|
// ---- governor board (Protocol 2.0) ----
|
|
|
|
/// Upserts one city's latest governance state, keyed by city name. Fed from `city.update`.
|
|
pub async fn upsert_governor(&self, city: &str, json: &str, t: i64) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO governors (city, json, updated_t) VALUES (?, ?, ?)
|
|
ON CONFLICT(city) DO UPDATE SET json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(city)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// The full governor board: every city's latest governance snapshot, ordered by city.
|
|
pub async fn governors_all(&self) -> anyhow::Result<Vec<Value>> {
|
|
let rows = sqlx::query("SELECT json FROM governors ORDER BY city")
|
|
.fetch_all(&self.pool)
|
|
.await?;
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
|
|
// ---- house registry (Protocol 2.0) ----
|
|
|
|
/// Upserts one house's latest state, keyed by serial. Fed from `house.update`.
|
|
pub async fn upsert_house(
|
|
&self,
|
|
serial: &str,
|
|
name: Option<&str>,
|
|
json: &str,
|
|
t: i64,
|
|
) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO houses (serial, name, json, updated_t) VALUES (?, ?, ?, ?)
|
|
ON CONFLICT(serial) DO UPDATE SET name = excluded.name, json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(serial)
|
|
.bind(name)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Drops one house from the registry. Fed from `house.remove` (demolished / traded away).
|
|
pub async fn delete_house(&self, serial: &str) -> anyhow::Result<()> {
|
|
sqlx::query("DELETE FROM houses WHERE serial = ?")
|
|
.bind(serial)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// The full house registry: every house's latest snapshot, ordered by name then serial.
|
|
pub async fn houses_all(&self) -> anyhow::Result<Vec<Value>> {
|
|
let rows = sqlx::query("SELECT json FROM houses ORDER BY name, serial")
|
|
.fetch_all(&self.pool)
|
|
.await?;
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
|
|
// ---- Town Cryer news (Protocol 2.1) ----
|
|
|
|
/// Stores/replaces one external news article (the `news.add` command json), keyed by id. The
|
|
/// website is the source of truth; this lets the sidecar replay the set to the shard on reconnect
|
|
/// (the shard does not persist NewsEntries across a reboot).
|
|
pub async fn upsert_news(&self, id: &str, json: &str, t: i64) -> anyhow::Result<()> {
|
|
sqlx::query(
|
|
"INSERT INTO news (id, json, updated_t) VALUES (?, ?, ?)
|
|
ON CONFLICT(id) DO UPDATE SET json = excluded.json, updated_t = excluded.updated_t",
|
|
)
|
|
.bind(id)
|
|
.bind(json)
|
|
.bind(t)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Removes one external news article.
|
|
pub async fn delete_news(&self, id: &str) -> anyhow::Result<()> {
|
|
sqlx::query("DELETE FROM news WHERE id = ?")
|
|
.bind(id)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Every stored external news article (as its `news.add` command), oldest first so a replay
|
|
/// re-inserts them in the same order the website added them.
|
|
pub async fn news_all(&self) -> anyhow::Result<Vec<Value>> {
|
|
let rows = sqlx::query("SELECT json FROM news ORDER BY updated_t")
|
|
.fetch_all(&self.pool)
|
|
.await?;
|
|
Ok(parse_json_column(rows))
|
|
}
|
|
}
|
|
|
|
fn parse_json_column(rows: Vec<sqlx::sqlite::SqliteRow>) -> Vec<Value> {
|
|
rows.into_iter()
|
|
.filter_map(|r| serde_json::from_str(&r.get::<String, _>("json")).ok())
|
|
.collect()
|
|
}
|
|
|
|
const SCHEMA: &str = r#"
|
|
CREATE TABLE IF NOT EXISTS events (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
t INTEGER NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
json TEXT NOT NULL
|
|
);
|
|
CREATE INDEX IF NOT EXISTS idx_events_kind ON events(kind, id);
|
|
|
|
CREATE TABLE IF NOT EXISTS links (
|
|
account TEXT PRIMARY KEY,
|
|
website_user_id TEXT NOT NULL,
|
|
linked_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS profiles (
|
|
serial TEXT PRIMARY KEY,
|
|
account TEXT,
|
|
name TEXT,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS champs (
|
|
serial TEXT PRIMARY KEY,
|
|
status TEXT,
|
|
name TEXT,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS guilds (
|
|
id INTEGER PRIMARY KEY,
|
|
name TEXT,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS governors (
|
|
city TEXT PRIMARY KEY,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS houses (
|
|
serial TEXT PRIMARY KEY,
|
|
name TEXT,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
|
|
CREATE TABLE IF NOT EXISTS news (
|
|
id TEXT PRIMARY KEY,
|
|
json TEXT NOT NULL,
|
|
updated_t INTEGER NOT NULL
|
|
);
|
|
"#;
|