diff --git a/sidecar/README.md b/sidecar/README.md index cb5a1af..bd774fe 100644 --- a/sidecar/README.md +++ b/sidecar/README.md @@ -41,10 +41,10 @@ So you can never accidentally run without auth. Rotate by editing the token and ## Protocol version -The wire protocol has a version (`PROTOCOL_VERSION`, currently **1**), so the website and sidecar detect a mismatch immediately instead of failing in strange ways when a message shape changes. +The wire protocol has a version (`PROTOCOL_VERSION`, currently **3**), so the website and sidecar detect a mismatch immediately instead of failing in strange ways when a message shape changes. -- Every response carries an `X-UOLink-Version: 1` header. -- `/health` and the WebSocket `ws.hello` include `"protocol": 1`. +- Every response carries an `X-UOLink-Version: 3` header. +- `/health` and the WebSocket `ws.hello` include `"protocol": 3`. - If a request sends `X-UOLink-Version` and it disagrees with the sidecar, the request is rejected **409 Conflict** with `{sidecar_protocol, client_protocol}` so the mismatch is obvious. Bump `PROTOCOL_VERSION` in `main.rs` whenever an event or endpoint's shape changes. diff --git a/sidecar/src/main.rs b/sidecar/src/main.rs index c6ce727..1a5a659 100644 --- a/sidecar/src/main.rs +++ b/sidecar/src/main.rs @@ -24,7 +24,12 @@ use tracing_subscriber::EnvFilter; /// v2 (Protocol 2.0): adds the account-provisioning verbs/endpoints (`POST /accounts/create`, /// `DELETE /link/:account`) and their events. Outbound event kinds are additive, so a v1 website /// keeps working against the live feed; the new *endpoints* require a v2 sidecar. -pub const PROTOCOL_VERSION: u32 = 2; +/// +/// v3 (Protocol 3.0): adds `world.ruleset`, `points.board` and `vendor.listing` / +/// `vendor.listing.remove`, with the `GET /ruleset`, `/points` and `/market` reads that serve them +/// from the store. Same shape as the v2 bump — the kinds are additive, the endpoints are not — and +/// there is deliberately no feature-negotiation array: v3 implies all three kinds. +pub const PROTOCOL_VERSION: u32 = 3; #[tokio::main] async fn main() -> anyhow::Result<()> { @@ -194,6 +199,79 @@ async fn main() -> anyhow::Result<()> { } } } + // Points/loyalty boards (Protocol 3.0): one row per point system, keyed by + // the shard's own PointsType name. The plugin only emits a system whose top N + // actually moved, so this is a sparse stream of overwrites — and there is no + // `points.remove` to handle, because the shard's set of systems is fixed at + // startup and cannot shrink. + "points.board" => { + if let Some(system) = ev.value.get("system").and_then(|s| s.as_str()) { + if let Err(e) = event_store + .upsert_points_board( + system, + ev.value.get("nameString").and_then(|n| n.as_str()), + &text, + t, + ) + .await + { + tracing::warn!(error = %e, "failed to upsert points board"); + } + } + } + // Player-vendor market index (Protocol 3.0). Each frame is authoritative for + // one vendor — the shard's round-robin sweep only emits a shop whose contents, + // prices or location actually moved — so this is a whole-row overwrite. + // + // Unlike the boards above there IS a remove: a vendor is dismissed, expires, or + // its owner switches off the in-game Vendor Search flag, and any of those must + // take the shop off the site. The last of the three is a privacy control, so + // dropping the row promptly is the point rather than housekeeping. + "vendor.listing" => { + if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) { + let loc = ev.value.get("location"); + let field = |k: &str| loc.and_then(|l| l.get(k)); + if let Err(e) = event_store + .upsert_vendor( + serial, + ev.value.get("shopName").and_then(|v| v.as_str()), + ev.value.get("ownerName").and_then(|v| v.as_str()), + field("map").and_then(|v| v.as_str()), + field("x").and_then(|v| v.as_i64()), + field("y").and_then(|v| v.as_i64()), + field("region").and_then(|v| v.as_str()), + ev.value.get("count").and_then(|v| v.as_i64()), + &text, + t, + ) + .await + { + tracing::warn!(error = %e, "failed to upsert vendor listing"); + } + } + } + "vendor.listing.remove" => { + if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) { + if let Err(e) = event_store.delete_vendor(serial).await { + tracing::warn!(error = %e, "failed to remove vendor listing"); + } + } + } + // Shard ruleset (Protocol 3.0): a singleton projection. The shard re-emits + // world.ruleset on every connect, so this row is simply overwritten; `rev` + // lets a reader tell a re-send from an actual config change. + "world.ruleset" => { + if let Err(e) = event_store + .upsert_ruleset( + ev.value.get("rev").and_then(|r| r.as_str()), + &text, + t, + ) + .await + { + tracing::warn!(error = %e, "failed to upsert ruleset"); + } + } _ => {} } } diff --git a/sidecar/src/store.rs b/sidecar/src/store.rs index 58f42cb..d1e55cc 100644 --- a/sidecar/src/store.rs +++ b/sidecar/src/store.rs @@ -293,6 +293,169 @@ impl Store { Ok(parse_json_column(rows)) } + // ---- shard ruleset (Protocol 3.0) ---- + + /// Stores the shard's published ruleset. A singleton (`id = 1`): the shard emits one + /// `world.ruleset` frame per connect describing how it is configured, and only the latest one + /// matters. `rev` is the shard's FNV-1a of the body, kept so a reader can tell "same ruleset, + /// re-sent on reconnect" from "the operator changed something" without diffing the JSON. + pub async fn upsert_ruleset(&self, rev: Option<&str>, json: &str, t: i64) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO ruleset (id, rev, json, updated_t) VALUES (1, ?, ?, ?) + ON CONFLICT(id) DO UPDATE SET rev = excluded.rev, json = excluded.json, updated_t = excluded.updated_t", + ) + .bind(rev) + .bind(json) + .bind(t) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// The stored ruleset, or `None` if the shard has never published one. Returning `None` rather + /// than an empty object is deliberate: "not published yet" and "published, everything off" are + /// different answers and the website renders them differently. + pub async fn ruleset(&self) -> anyhow::Result> { + let row = sqlx::query("SELECT json FROM ruleset WHERE id = 1") + .fetch_optional(&self.pool) + .await?; + Ok(row.and_then(|r| serde_json::from_str(&r.get::("json")).ok())) + } + + // ---- points / loyalty boards (Protocol 3.0) ---- + + /// Upserts one system's leaderboard, keyed by its `PointsType` name (`QueensLoyalty`, + /// `CleanUpBritannia`, …). Fed from `points.board`; one row per system, always the most recent + /// top-N snapshot. + /// + /// There is no matching delete, and that is deliberate rather than an omission: the shard's set + /// of point systems is fixed at startup by `PointsSystem.Configure`, so a system cannot vanish + /// at runtime and the plugin emits no `points.remove`. Same argument the governor board makes. + pub async fn upsert_points_board( + &self, + system: &str, + name: Option<&str>, + json: &str, + t: i64, + ) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO points_boards (system, name, json, updated_t) VALUES (?, ?, ?, ?) + ON CONFLICT(system) DO UPDATE SET name = excluded.name, json = excluded.json, updated_t = excluded.updated_t", + ) + .bind(system) + .bind(name) + .bind(json) + .bind(t) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Every system's latest board, ordered by display name then system key. Systems the shard has + /// never published are simply absent — the website renders the set it is given. + pub async fn points_boards_all(&self) -> anyhow::Result> { + let rows = sqlx::query("SELECT json FROM points_boards ORDER BY name, system") + .fetch_all(&self.pool) + .await?; + Ok(parse_json_column(rows)) + } + + /// One system's board, or `None` when that system has never published one. `None` is a real + /// answer (an unknown system name, or one the operator excluded via `Bridge.cfg PointsSystems`), + /// which the website turns into a 404 rather than an empty board. + pub async fn points_board(&self, system: &str) -> anyhow::Result> { + let row = sqlx::query("SELECT json FROM points_boards WHERE system = ?") + .bind(system) + .fetch_optional(&self.pool) + .await?; + Ok(row.and_then(|r| serde_json::from_str(&r.get::("json")).ok())) + } + + // ---- player-vendor market index (Protocol 3.0) ---- + + /// Upserts one vendor's whole listing, keyed by serial. Fed from `vendor.listing`, which the + /// shard emits as an authoritative per-vendor frame — so this replaces the row outright rather + /// than merging anything. + /// + /// The items ride inside `json` and are deliberately NOT normalized into a `vendor_items` + /// table. The sidecar's job for the market is outage resilience (`PROTOCOL_2.md` §12.2) — hand + /// the website back what the shard last said — not search. Search lives in MariaDB on the + /// website side, where the query surface, the indexes and the cliloc-resolved display names + /// already are; a second search implementation here would be one more thing to keep in step + /// with it for no reader. + #[allow(clippy::too_many_arguments)] + pub async fn upsert_vendor( + &self, + serial: &str, + shop_name: Option<&str>, + owner_name: Option<&str>, + map: Option<&str>, + x: Option, + y: Option, + region: Option<&str>, + count: Option, + json: &str, + t: i64, + ) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO vendors (serial, shop_name, owner_name, map, x, y, region, count, json, updated_t) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(serial) DO UPDATE SET shop_name = excluded.shop_name, + owner_name = excluded.owner_name, map = excluded.map, x = excluded.x, y = excluded.y, + region = excluded.region, count = excluded.count, json = excluded.json, + updated_t = excluded.updated_t", + ) + .bind(serial) + .bind(shop_name) + .bind(owner_name) + .bind(map) + .bind(x) + .bind(y) + .bind(region) + .bind(count) + .bind(json) + .bind(t) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Drops one vendor from the index. Fed from `vendor.listing.remove` — a vendor dismissed, + /// expired, or whose owner switched off its in-game Vendor Search flag. + pub async fn delete_vendor(&self, serial: &str) -> anyhow::Result<()> { + sqlx::query("DELETE FROM vendors WHERE serial = ?") + .bind(serial) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// One page of the index, ordered by serial. + /// + /// Paged where the other boards are not, and the ordering is why it can be: a whole-world + /// market is the one board that does not fit in a response. Ordering by SERIAL rather than by + /// shop name is deliberate — the page is a snapshot cursor for the website's reconnect + /// backfill, and a serial is stable while a shop name is renameable, so a rename mid-backfill + /// cannot make a vendor skip or repeat a page. + pub async fn vendors_page(&self, limit: i64, offset: i64) -> anyhow::Result> { + let limit = limit.clamp(1, 1000); + let offset = offset.max(0); + let rows = sqlx::query("SELECT json FROM vendors ORDER BY serial LIMIT ? OFFSET ?") + .bind(limit) + .bind(offset) + .fetch_all(&self.pool) + .await?; + Ok(parse_json_column(rows)) + } + + /// How many vendors the index holds, so a paging caller knows when to stop. + pub async fn vendors_count(&self) -> anyhow::Result { + let row = sqlx::query("SELECT COUNT(*) AS n FROM vendors") + .fetch_one(&self.pool) + .await?; + Ok(row.get::("n")) + } + // ---- Town Cryer news (Protocol 2.1) ---- /// Stores/replaces one external news article (the `news.add` command json), keyed by id. The @@ -392,4 +555,39 @@ CREATE TABLE IF NOT EXISTS news ( json TEXT NOT NULL, updated_t INTEGER NOT NULL ); + +-- Points/loyalty leaderboards (Protocol 3.0). One row per point system, keyed by the shard's +-- own PointsType name; `name` is the resolved display name, hoisted only for the ORDER BY. +CREATE TABLE IF NOT EXISTS points_boards ( + system TEXT PRIMARY KEY, + name TEXT, + json TEXT NOT NULL, + updated_t INTEGER NOT NULL +); + +-- Player-vendor market index (Protocol 3.0). One row per vendor, holding the whole authoritative +-- `vendor.listing` frame including its items. The hoisted columns exist for the ORDER BY and for +-- an operator eyeballing the table; nothing here is searched, because search is the website's job +-- (see upsert_vendor). Rows are dropped on `vendor.listing.remove`. +CREATE TABLE IF NOT EXISTS vendors ( + serial TEXT PRIMARY KEY, + shop_name TEXT, + owner_name TEXT, + map TEXT, + x INTEGER, + y INTEGER, + region TEXT, + count INTEGER, + json TEXT NOT NULL, + updated_t INTEGER NOT NULL +); + +-- The shard's published ruleset (Protocol 3.0). Singleton: the CHECK is what makes it one, +-- so an upsert can target id = 1 unconditionally and no second row can ever appear. +CREATE TABLE IF NOT EXISTS ruleset ( + id INTEGER PRIMARY KEY CHECK (id = 1), + rev TEXT, + json TEXT NOT NULL, + updated_t INTEGER NOT NULL +); "#; diff --git a/sidecar/src/web.rs b/sidecar/src/web.rs index b37b262..978716e 100644 --- a/sidecar/src/web.rs +++ b/sidecar/src/web.rs @@ -82,6 +82,20 @@ pub async fn serve(addr: &str, state: AppState) -> anyhow::Result<()> { .route("/governors", get(governors)) .route("/online", get(online)) .route("/houses", get(houses)) + // The shard ruleset (Protocol 3.0), likewise store-backed: the shard publishes it once per + // connect, so serving it from the store is what lets the site's rules page render while the + // shard is down. + .route("/ruleset", get(ruleset)) + // Points/loyalty leaderboards (Protocol 3.0), store-backed like the other boards: the + // whole set, or one system by its PointsType name. + .route("/points", get(points)) + .route("/points/:system", get(points_system)) + // The player-vendor market index (Protocol 3.0). `/market`, NOT `/vendors`: axum would + // route the latter fine, but `/vendors/:account` next door is the per-account RPC, and two + // routes a prefix apart that mean "this player's shops" and "every shop on the shard" is a + // readability trap nobody wins. The only PAGED read the sidecar serves — a whole-world + // market does not fit in one response. + .route("/market", get(market)) .route_layer(middleware::from_fn_with_state(state.clone(), gate)); let app = Router::new() @@ -790,6 +804,60 @@ async fn houses(State(st): State) -> impl IntoResponse { } } +/// The shard's published ruleset: expansion, which optional systems are on, skill/stat caps, +/// account and house limits, champion scroll rules, the save/restart schedule. Served from the +/// store, so it answers during a shard outage with the last-known ruleset — which is the whole +/// point, since a rules page that goes blank when the shard restarts is worse than a stale one. +/// +/// `{"ruleset": null}` means the shard has never published one (an old plugin, or +/// `Bridge.RulesetEnabled=false`), which the website renders differently from a published ruleset. +async fn ruleset(State(st): State) -> impl IntoResponse { + match st.store.ruleset().await { + Ok(r) => (StatusCode::OK, Json(json!({ "ruleset": r }))), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({"error": e.to_string()})), + ), + } +} + +/// Every points/loyalty leaderboard the shard publishes: one entry per point system, each with its +/// display name (literal and/or cliloc), max points, participant count and top N. Store-backed like +/// the other boards, so the site's leaderboards page renders during a shard outage — which matters +/// more here than elsewhere, since these are month-scale standings that a restart must not blank. +async fn points(State(st): State) -> impl IntoResponse { + match st.store.points_boards_all().await { + Ok(boards) => (StatusCode::OK, Json(json!({"boards": boards}))), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({"error": e.to_string()})), + ), + } +} + +/// One system's board by its `PointsType` name (`QueensLoyalty`, `CleanUpBritannia`, …). +/// +/// 404 rather than an empty board when the system is unknown: the shard publishes only the systems +/// it shows on the loyalty gump (or the explicit `Bridge.cfg PointsSystems` list), so "no such +/// board" and "a board with nobody on it" are different answers and the website renders them +/// differently. +async fn points_system( + State(st): State, + Path(system): Path, +) -> impl IntoResponse { + match st.store.points_board(&system).await { + Ok(Some(board)) => (StatusCode::OK, Json(board)), + Ok(None) => ( + StatusCode::NOT_FOUND, + Json(json!({"error": "unknown points system", "system": system})), + ), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({"error": e.to_string()})), + ), + } +} + /// The current online population: total plus per-facet and per-region counts. This is the most /// recent `presence.online` snapshot from the event store (so it survives a sidecar restart); the /// live `presence.online` stream keeps it current, and `GET /history?kind=presence.online` gives the @@ -810,6 +878,54 @@ async fn online(State(st): State) -> impl IntoResponse { } } +#[derive(Deserialize)] +struct PageQuery { + limit: Option, + offset: Option, +} + +/// The player-vendor market index: every vendor's shop name, owner, location and priced inventory, +/// as the shard last published it. Store-backed like the other boards, which is what lets the +/// website's market page render (labelled stale) while the shard is down. +/// +/// Paged — `?limit=&offset=`, limit clamped to 1..1000, default 200 — because this is the one board +/// that can be a whole world's inventory. `total` is returned alongside so the caller knows when to +/// stop rather than paging until it sees a short page, which would race a concurrent sweep. +/// +/// The frames are served VERBATIM, including owner names and coordinates. That is not an oversight: +/// the sidecar defines no audiences (docs/link/v3.md §3.2). Deciding who may see a vendor's owner +/// or whereabouts is the website's job and is admin-configurable there. +async fn market(State(st): State, Query(q): Query) -> impl IntoResponse { + let limit = q.limit.unwrap_or(200); + let offset = q.offset.unwrap_or(0); + + let total = match st.store.vendors_count().await { + Ok(n) => n, + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({"error": e.to_string()})), + ) + } + }; + + match st.store.vendors_page(limit, offset).await { + Ok(vendors) => ( + StatusCode::OK, + Json(json!({ + "vendors": vendors, + "total": total, + "limit": limit.clamp(1, 1000), + "offset": offset.max(0), + })), + ), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({"error": e.to_string()})), + ), + } +} + // ---- websocket ---- async fn ws_upgrade(ws: WebSocketUpgrade, State(state): State) -> impl IntoResponse {