7 Commits

Author SHA1 Message Date
82872ffba7 Merge pull request 'feat(sidecar): protocol 8 — the asset plane (Asset Bridge cutover, 1 of 5)' (#44) from edge into main
Some checks failed
sync-project-tree / sync (push) Successful in 7s
SonarQube / analysis (push) Failing after -53s
Release sidecar / release (push) Successful in 12m2s
Reviewed-on: #44
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-09-14 23:09:12 +00:00
baa04e1a76 Merge pull request 'feat(sidecar): forward the asset manifest, the pixels and the body pass (Phase 3)' (#43) from feat/asset-bridge-p3 into edge
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 1m52s
Reviewed-on: #43
2026-09-10 23:57:55 +00:00
143f424867 feat(sidecar): forward the asset manifest, the pixels and the body pass (Phase 3)
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m24s
Three routes, forwarded verbatim like everything else on this link:

  GET  /assets/manifest?family=&cursor=
  POST /assets/fetch
  POST /assets/bodies

**The two POSTs are reads.** The method is the request body, not a side
effect -- a few hundred asset keys do not belong in a query string, and these
are the only reads on this link that take one. `assets_call` is `event_call`'s
shape with one difference that matters: it responds through `respond_assets`,
so `bridge.busy` is a 425 rather than an idempotency collision. On this plane
busy is the ORDINARY answer during an import, and a caller that read it as an
error would abandon a healthy transfer.

422 gains a second meaning here alongside "the shard cannot decode that file":
the mid-import guard. A manifest reply carries a `catalog` id the shard derives
from its own client files, and passing it back on a fetch makes the shard refuse
if those files moved in between -- without which an operator patching their
client halfway through an import gets one asset set stitched out of two, with no
error anywhere.

v8.md §16 listed phase 3 as servuo-plugins + module-uo. That was wrong: web.rs
routes every command explicitly and has no generic /assets/* forwarder, so this
repo is in the phase. The doc now says so.

64 tests pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016wDDVXWMDz82WqE1i969r4
2026-09-10 18:40:26 -05:00
60e6de55f3 Merge pull request 'feat(sidecar): forward the cliloc table, and stop reading refusals for meaning (Phase 2)' (#42) from feat/asset-bridge-p2 into edge
Reviewed-on: #42
2026-09-10 16:20:00 +00:00
b92393d224 feat(sidecar): forward the cliloc table, and stop reading refusals for meaning (Phase 2)
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 3m2s
`GET /cliloc` — the first protocol-8 family that carries content rather than a
manifest. The shard decompresses its own client's table and cuts it into pages;
this forwards them and keeps none of it, which matters more here than usual: the
payload is five megabytes of EA's strings out of the operator's own client, and
the one copy that should exist is the one the website imports.

Paging is the caller's, deliberately. `?cursor=` echoes back the previous reply's
cursor until one says `more: false`; a sidecar that helpfully assembled the pages
would be holding the whole table in memory to do it. `?lang=` selects the file
and defaults on the shard.

`asset_error_status` replaces phase 1's substring test. That test chose 403 or
400 by looking for the word "disabled" in an operator-facing sentence, so
rewording the message would silently turn a refusal into a bad request. The
overlay now sends a `code`: DISABLED 403, NOT_FOUND 404 (a client without the
file — an operator fact, not a bug), UNREADABLE 422 (a file it has and cannot
decode, where repeating the request cannot help), UNAVAILABLE 503, BAD_REQUEST
400. The substring check survives as a fallback, with a test, because an overlay
and a sidecar are deployed separately and a phase-1 shard must keep its 403.

No `PROTOCOL_VERSION` change: 8 already covers this family (v8.md §14).

Verified against a live shard: 12 pages, 67,496 rows, every page inside the
512 KiB budget (max 524,086 of 524,288) and well under the 1 MiB line cap, the
whole table in 1.4 s. Concurrent callers get 425 while one is served, which is
the flow control working rather than an error.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016wDDVXWMDz82WqE1i969r4
2026-09-10 11:13:03 -05:00
6d83df0a2c Merge pull request 'feat(sidecar): protocol 8 — the asset plane, and a bound on what the shard can send' (#41) from feat/asset-bridge-p1 into edge
Reviewed-on: #41
2026-09-10 15:04:24 +00:00
a8f1804de9 feat(sidecar): protocol 8 — the asset plane, and a bound on what the shard can send
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m39s
Asset Bridge phase 1, sidecar half (docs/link/v8.md §3.3, §14).
Shard half: RunicGateway/servuo-plugins#28. Docs half: RunicGateway/docs#236.

Three things, one of which is not additive.

## The inbound line cap (§3.3) — the one that matters

`read_line` had **no bound at all**. That was survivable only because the shard had
never had a reason to send a large line. Protocol 8 gives it one deliberately, and an
unbounded read facing a component that now sends megabytes is a memory-exhaustion
shape we would be inventing ourselves.

`MAX_INBOUND_LINE_BYTES` is **1 MiB** — symmetric with the cap `BridgeLink.cs` has
always applied to its own inbound lines, so both directions of this link now read the
same. The shard's batch budget is 512 KiB, and the factor of two is load-bearing: a
page always admits its first item even when that item alone exceeds the budget (the
alternative is an oversized item skipped for the budget on every page forever), so the
wire needs room for one overshoot.

An over-long line is **discarded and the connection kept** — `BridgeLink.cs`'s own
disposition in the other direction. Tearing the link down would take the live event
feed with it over one malformed frame, and the lost reply just times out and is
re-requested; everything on this plane is idempotent.

**`LineReader` holds its state in a struct rather than in locals, and that is the
subtle part.** This is polled inside a `tokio::select!`, so the future is dropped
whenever a command wins the race. A `discarding` flag in a local would be lost with
it — and losing it turns the tail of an over-long line into a line of its own, silently.
There is a test for exactly that, and another for an over-long line whose terminator
lands in the very chunk that crosses the cap.

## `GET /assets/sources`

Stage 1 of the import gate, forwarded verbatim like everything else. `respond_assets`
maps `bridge.busy` → **425** and a disabled plane → **403**.

425 deserves a note: on this plane it is not an idempotency collision, it is flow
control, and it is the **ordinary** answer mid-import rather than a rare one. The shard
serves one asset request at a time because its outbound queue is bounded in lines, not
bytes. A caller treating it as an error would abandon a healthy transfer.

403 for the same reason the event plane's gate is a 403: `Bridge.AssetsEnabled` off is
an operator declining to let the website read their client files, not a malformed
request, and 400 would send an administrator hunting a bug in a correct call.

## `PROTOCOL_VERSION` 7 → 8

Paired with `servuo-plugins/overlay.toml` in the linked PR — the installer refuses to
compose a bundle whose halves disagree, so a split bump fails silently at the next
release.

## Also

`docs/link/INTEGRATION.md` still advertised `X-UOLink-Version: 6`; it was already two
versions stale before this change. Fixed in the docs PR.

61 tests pass, `cargo fmt --check` and `cargo clippy -- -D warnings` clean. Verified
against the real shard: `/health` reports protocol 8, `/assets/sources` returns 200 with
`X-UOLink-Version: 8`, and live events kept flowing through the new reader with no
warnings logged.

- [x] AI-assisted — Claude Code (Opus 5)

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016wDDVXWMDz82WqE1i969r4
2026-09-10 08:32:39 -05:00
3 changed files with 688 additions and 20 deletions

View File

@@ -91,7 +91,40 @@ use tracing_subscriber::EnvFilter;
/// the dumb-forwarder property again. /// the dumb-forwarder property again.
/// ///
/// **No store migration.** Nothing gains a column; the new kind is persisted whole like every other. /// **No store migration.** Nothing gains a column; the new kind is persisted whole like every other.
pub const PROTOCOL_VERSION: u32 = 7; ///
/// # Protocol 8 — the Asset Bridge (docs/link/v8.md)
///
/// The shard starts sending the operator's own **client assets** over this link: the cliloc string
/// table, creature and item art, player models. The point is that an operator stops having to run
/// a GUI converter on a desktop to make their site render a bestiary, and the shard is the only
/// host that already has the client files — a ServUO server cannot boot without them.
///
/// Phase 1 is the transport, and the sidecar's share of it is three things:
///
/// * **A new command family, `assets.*`, forwarded verbatim** like every other. The first of them
/// is `assets.sources` — stage 1 of the import gate: what the client files currently are, and
/// what version of the shard's extractor would read them. No pixels cross on this call.
/// * **An inbound line cap** — [`shard::MAX_INBOUND_LINE_BYTES`]. This is the one change that is
/// not additive. `read_line` had no bound at all, which was survivable while the shard had no
/// reason to send a large line; protocol 8 gives it one deliberately, and an unbounded read
/// facing a component that now sends megabytes is a memory-exhaustion shape we would be
/// inventing ourselves.
/// * **Nothing else.** Assets ride the request/reply path, so `rpc::try_route` consumes them
/// before `app.rs` can persist them to the store and fan them out to every WebSocket
/// subscriber — which is what keeps a 512 KiB reply from being written to SQLite and broadcast
/// to every connected client. The dumb-forwarder property is doing real work here: the sidecar
/// does not know what an asset is, and must not learn.
///
/// Phase 2 adds the first family that actually carries content: **`cliloc.table`**, served at
/// `GET /cliloc`. It pages — the shard cuts at a byte budget and the caller echoes a cursor back —
/// and the sidecar forwards it without keeping any of it, which matters more here than usual: the
/// payload is five megabytes of EA's strings out of the operator's own client, and the one copy of
/// it that should exist is the one the website imports. That phase also gave `assets.error` a
/// `code`, so the status a refusal maps to stops depending on the wording of a human-facing
/// sentence (see [`web::asset_error_status`]).
///
/// **No store migration**, again: nothing on this plane is an event, so nothing is persisted.
pub const PROTOCOL_VERSION: u32 = 8;
// Not `#[tokio::main]`: on Windows the SCM dispatcher takes over this thread and starts the runtime // Not `#[tokio::main]`: on Windows the SCM dispatcher takes over this thread and starts the runtime
// itself, on its own thread, once the service actually begins. The runtime is built by whichever // itself, on its own thread, once the service actually begins. The runtime is built by whichever

View File

@@ -7,15 +7,130 @@
//! Framing is newline-delimited JSON, bidirectional: the shard sends events, we send commands. We //! Framing is newline-delimited JSON, bidirectional: the shard sends events, we send commands. We
//! accept one shard connection at a time and re-accept when it drops (the shard reconnects on its //! accept one shard connection at a time and re-accept when it drops (the shard reconnects on its
//! own, with a bounded backoff). //! own, with a bounded backoff).
//!
//! Inbound lines are **capped** (see [`MAX_INBOUND_LINE_BYTES`]). Until protocol 8 they were not:
//! `read_line` will buffer a line of any length, which was survivable only because the shard had
//! never had a reason to send a large one. The Asset Bridge gives it one, so the gap had to close
//! before it became a memory-exhaustion shape we invented ourselves.
use std::sync::Arc; use std::sync::Arc;
use serde_json::Value; use serde_json::Value;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpListener; use tokio::net::TcpListener;
use tokio::sync::{mpsc, Mutex}; use tokio::sync::{mpsc, Mutex};
use tracing::{info, warn}; use tracing::{info, warn};
/// The longest line the sidecar will accept from the shard, in bytes.
///
/// Set above the largest legal batch rather than at it: the shard cuts a batch when the next item
/// would take it past `Bridge.AssetBatchBytes` (512 KiB), and always admits the first item of a
/// page even when that item alone is bigger than the budget — so one page can legitimately
/// overshoot by one item. Doubling the budget to get this cap is what makes that overshoot safe
/// instead of a dropped reply.
///
/// Over-long lines are **discarded, not buffered**, and the connection stays up. That is the same
/// disposition `BridgeLink.cs` has always had for its own 1 MiB inbound cap in the other
/// direction, and it is the right one here: a single malformed frame is not a reason to tear down
/// a link that live events are flowing over. The dropped reply simply times out and is
/// re-requested, which is safe because everything on the asset plane is idempotent.
pub const MAX_INBOUND_LINE_BYTES: usize = 1024 * 1024;
/// What one read off the shard socket produced.
#[derive(Debug)]
enum Line {
/// A complete line, within the cap.
Complete(String),
/// A line that ran past the cap. Carries how many bytes were thrown away, for the log.
TooLong(usize),
/// The shard closed the connection.
Eof,
}
/// A cancel-safe, capped, newline-delimited reader.
///
/// Every piece of state that must survive a partial read lives here rather than in a local,
/// because this is polled inside a `tokio::select!`: the loop below drops the future whenever a
/// command wins the race, and a `discarding` flag or a half-filled buffer held in a local would be
/// lost with it. Losing the buffer corrupts the *next* line; losing `discarding` turns the tail of
/// an over-long line into a line of its own. Both are silent.
///
/// The only await point is `fill_buf`, and nothing is consumed until after it returns, so a
/// cancellation between the two can lose at most the wakeup.
#[derive(Default)]
struct LineReader {
buf: Vec<u8>,
discarding: bool,
discarded: usize,
}
impl LineReader {
async fn next<R: AsyncBufRead + Unpin>(&mut self, reader: &mut R) -> std::io::Result<Line> {
loop {
let consumed;
let outcome;
{
let available = reader.fill_buf().await?;
if available.is_empty() {
return Ok(Line::Eof);
}
match available.iter().position(|&b| b == b'\n') {
Some(at) => {
consumed = at + 1;
if self.discarding {
// The tail of a line we already gave up on. Swallow it, terminator
// included, and report the size once.
self.discarded += at;
let total = self.discarded;
self.discarding = false;
self.discarded = 0;
outcome = Some(Line::TooLong(total));
} else if self.buf.len() + at > MAX_INBOUND_LINE_BYTES {
// The cap is reached only now, on the chunk that also holds the
// terminator — so there is nothing left to discard.
let total = self.buf.len() + at;
self.buf.clear();
outcome = Some(Line::TooLong(total));
} else {
self.buf.extend_from_slice(&available[..at]);
let line = String::from_utf8_lossy(&self.buf).into_owned();
self.buf.clear();
outcome = Some(Line::Complete(line));
}
}
None => {
consumed = available.len();
if self.discarding {
self.discarded += consumed;
} else if self.buf.len() + consumed > MAX_INBOUND_LINE_BYTES {
// Refuse rather than buffer: this is the whole point of the cap.
// Everything up to the next newline is now dropped on the floor.
self.discarded = self.buf.len() + consumed;
self.buf.clear();
self.discarding = true;
} else {
self.buf.extend_from_slice(available);
}
outcome = None;
}
}
}
reader.consume(consumed);
if let Some(line) = outcome {
return Ok(line);
}
}
}
}
/// An event line received from the shard, parsed. `kind` is lifted out for routing. /// An event line received from the shard, parsed. `kind` is lifted out for routing.
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct ShardEvent { pub struct ShardEvent {
@@ -112,31 +227,41 @@ async fn handle_connection(
handle.set(Some(cmd_tx)).await; handle.set(Some(cmd_tx)).await;
let mut reader = BufReader::new(read_half); let mut reader = BufReader::new(read_half);
let mut line = String::new(); let mut lines = LineReader::default();
loop { loop {
tokio::select! { tokio::select! {
// Inbound: a line from the shard. // Inbound: a line from the shard.
result = reader.read_line(&mut line) => { result = lines.next(&mut reader) => {
let n = result?; match result? {
if n == 0 { Line::Eof => return Ok(()), // clean EOF: shard closed
return Ok(()); // clean EOF: shard closed Line::TooLong(bytes) => {
} // Deliberately not a disconnect. See MAX_INBOUND_LINE_BYTES: a reply lost
let trimmed = line.trim_end(); // this way times out on the caller's side and is re-requested, and tearing
if !trimmed.is_empty() { // the link down would take the live event feed with it.
match serde_json::from_str::<Value>(trimmed) { warn!(
Ok(value) => { bytes,
let kind = value cap = MAX_INBOUND_LINE_BYTES,
.get("kind") "inbound line over the cap; discarded"
.and_then(|k| k.as_str()) );
.unwrap_or("") }
.to_string(); Line::Complete(line) => {
let _ = event_tx.send(ShardEvent { kind, value }); let trimmed = line.trim_end();
if !trimmed.is_empty() {
match serde_json::from_str::<Value>(trimmed) {
Ok(value) => {
let kind = value
.get("kind")
.and_then(|k| k.as_str())
.unwrap_or("")
.to_string();
let _ = event_tx.send(ShardEvent { kind, value });
}
Err(e) => warn!(error = %e, line = %trimmed, "unparseable event"),
}
} }
Err(e) => warn!(error = %e, line = %trimmed, "unparseable event"),
} }
} }
line.clear();
} }
// Outbound: a command to write to the shard. // Outbound: a command to write to the shard.
cmd = cmd_rx.recv() => { cmd = cmd_rx.recv() => {
@@ -152,3 +277,108 @@ async fn handle_connection(
} }
} }
} }
#[cfg(test)]
mod tests {
use super::*;
/// Drives `LineReader` over a byte slice, returning every outcome up to EOF.
async fn read_all(input: &[u8]) -> Vec<Line> {
let mut reader = BufReader::with_capacity(64, input);
let mut lines = LineReader::default();
let mut out = Vec::new();
loop {
match lines.next(&mut reader).await.unwrap() {
Line::Eof => break,
other => out.push(other),
}
}
out
}
fn complete(lines: &[Line]) -> Vec<&str> {
lines
.iter()
.filter_map(|l| match l {
Line::Complete(s) => Some(s.as_str()),
_ => None,
})
.collect()
}
#[tokio::test]
async fn splits_on_newlines() {
let lines = read_all(b"{\"a\":1}\n{\"b\":2}\n").await;
assert_eq!(complete(&lines), vec!["{\"a\":1}", "{\"b\":2}"]);
}
/// The reader's buffer is 64 bytes here, so every one of these lines spans several
/// `fill_buf` chunks. Reassembly across chunks is the thing `read_line` did for us.
#[tokio::test]
async fn reassembles_across_chunks() {
let long = "x".repeat(500);
let input = format!("{}\n{}\n", long, long);
let lines = read_all(input.as_bytes()).await;
assert_eq!(complete(&lines), vec![long.as_str(), long.as_str()]);
}
/// The cap itself. The over-long line must be reported and thrown away, and — the part that
/// actually matters — the line *after* it must still arrive intact. A reader that lost its
/// `discarding` flag would emit the tail of the oversized line as a line of its own.
#[tokio::test]
async fn refuses_an_over_long_line_and_recovers() {
let mut input = Vec::new();
input.extend_from_slice(&b"a".repeat(MAX_INBOUND_LINE_BYTES + 10));
input.push(b'\n');
input.extend_from_slice(b"{\"kind\":\"pong\"}\n");
let lines = read_all(&input).await;
assert_eq!(lines.len(), 2);
assert!(
matches!(lines[0], Line::TooLong(n) if n >= MAX_INBOUND_LINE_BYTES),
"expected TooLong, got {:?}",
lines[0]
);
assert_eq!(complete(&lines), vec!["{\"kind\":\"pong\"}"]);
}
/// A line of exactly the cap is legal; one byte more is not. Checking both sides is what says
/// the comparison is `>` rather than `>=`, which would silently cost a byte of the budget.
#[tokio::test]
async fn the_cap_is_inclusive() {
let at_cap = "b".repeat(MAX_INBOUND_LINE_BYTES);
let lines = read_all(format!("{}\n", at_cap).as_bytes()).await;
assert_eq!(complete(&lines).len(), 1);
let over = "b".repeat(MAX_INBOUND_LINE_BYTES + 1);
let lines = read_all(format!("{}\n", over).as_bytes()).await;
assert!(complete(&lines).is_empty());
assert!(matches!(lines[0], Line::TooLong(_)));
}
/// An over-long line whose terminator lands in the very chunk that crosses the cap: the
/// reader must not leave itself in `discarding` and eat the next line as well.
#[tokio::test]
async fn over_long_line_terminating_in_the_crossing_chunk() {
let mut input = Vec::new();
input.extend_from_slice(&b"c".repeat(MAX_INBOUND_LINE_BYTES + 1));
input.extend_from_slice(b"\n{\"kind\":\"pong\"}\n");
let lines = read_all(&input).await;
assert!(matches!(lines[0], Line::TooLong(_)));
assert_eq!(complete(&lines), vec!["{\"kind\":\"pong\"}"]);
}
/// A partial line at EOF is dropped rather than delivered half-parsed. The shard reconnects
/// and re-sends; half a JSON object is not something to hand to the event fan-out.
#[tokio::test]
async fn trailing_partial_line_at_eof_is_dropped() {
let lines = read_all(b"{\"a\":1}\n{\"b\":").await;
assert_eq!(complete(&lines), vec!["{\"a\":1}"]);
}
}

View File

@@ -135,6 +135,28 @@ pub async fn serve(addr: &str, state: AppState) -> anyhow::Result<()> {
// readability trap nobody wins. The only PAGED read the sidecar serves — a whole-world // readability trap nobody wins. The only PAGED read the sidecar serves — a whole-world
// market does not fit in one response. // market does not fit in one response.
.route("/market", get(market)) .route("/market", get(market))
// The Asset Bridge (Protocol 8). Stage 1 of the two-stage import gate: what the shard's
// UO client files currently are. RPC, never store-backed — unlike the boards above there
// is nothing here worth serving stale, because the only question this answers is "have
// the files on that host changed since the last import", and a cached answer to that is
// worse than no answer.
.route("/assets/sources", get(assets_sources))
// Stage 2 (phase 3): what the shard could serve, hashed but without the pixels, so the
// website can ask only for what changed. Paged.
.route("/assets/manifest", get(assets_manifest))
// The pixels, for an explicit list of keys. POST rather than GET because the list is the
// request -- a few hundred keys do not belong in a query string, and this is the one place
// on this link where a read takes a body.
.route("/assets/fetch", post(assets_fetch))
// Slug -> body id. The one asset-plane call the shard answers ON ITS CORE THREAD, because
// it resolves a class name by constructing the creature and reading its body; the sidecar
// neither knows nor cares, which is the point of forwarding verbatim.
.route("/assets/bodies", post(assets_bodies))
// The cliloc table (Protocol 8, phase 2): UO's id -> display-string map, read out of the
// shard's own client and paged. RPC for the same reason as the source gate, and one more:
// it is five megabytes of somebody else's copyrighted strings, which this process has no
// business holding a copy of. It forwards them and forgets them.
.route("/cliloc", get(cliloc_table))
.route_layer(middleware::from_fn_with_state(state.clone(), gate)); .route_layer(middleware::from_fn_with_state(state.clone(), gate));
let app = Router::new() let app = Router::new()
@@ -1296,6 +1318,243 @@ async fn market(State(st): State<AppState>, Query(q): Query<PageQuery>) -> impl
} }
} }
// ---- the Asset Bridge (Protocol 8) ----
/// Stage 1 of the import gate: the shard's UO client files as they are right now — size, mtime and
/// content hash — plus the version of the extractor that would read them, and whether this host
/// can render an image at all.
///
/// The website diffs this against what it last imported and, in the overwhelmingly common case
/// that nothing changed, stops. That is the whole reason stage 1 exists separately from the asset
/// manifest: the normal case is a restart that changed nothing, and it has to cost nothing.
///
/// Forwarded verbatim, like everything else on this link. The sidecar does not know what a cliloc
/// or an anim file is, does not cache this, and has no opinion about what the website does with
/// the answer — the same dumb-forwarder property that keeps access control on the website where it
/// belongs.
async fn assets_sources(State(st): State<AppState>) -> impl IntoResponse {
let req_id = st.rpc.next_req_id();
let cmd = json!({"kind": "assets.sources", "reqId": req_id});
respond_assets(st.rpc.call(&st.shard, cmd, &req_id).await)
}
/// Maps an asset-plane reply to a status.
///
/// Two of these matter more than the rest and neither is the generic responder's answer:
///
/// **`bridge.busy` is a 425**, as everywhere else. On this plane it is not an idempotency
/// collision, it is flow control: the shard serves one asset request at a time on purpose, because
/// its outbound queue is bounded in *lines* and a queue of large replies is how the shard runs out
/// of memory. So it means "come back", it is entirely expected during an import, and a caller that
/// treated it as an error would abandon a perfectly healthy transfer.
///
/// **The plane being switched off is a 403.** `Bridge.AssetsEnabled` is an operator declining to
/// let the website read their client files off this host — a deliberate refusal, not a malformed
/// request — and answering 400 would send an administrator hunting a bug in a call that is written
/// correctly. Same argument the event plane's gate made in protocol 7.
fn respond_assets(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
match result {
Ok(value) => {
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
if kind == BUSY_KIND {
(BUSY_STATUS, Json(value))
} else if kind == "assets.error" {
(asset_error_status(&value), Json(value))
} else {
(StatusCode::OK, Json(value))
}
}
Err(RpcError::NoShard) => (
StatusCode::SERVICE_UNAVAILABLE,
Json(json!({"error": "shard not connected"})),
),
Err(RpcError::Timeout) => (
StatusCode::GATEWAY_TIMEOUT,
Json(json!({"error": "shard did not reply in time"})),
),
}
}
/// The status behind one `assets.error`.
///
/// Phase 2 gave the frame a `code`, and the reason is worth stating: phase 1 decided between 403
/// and 400 by looking for the word "disabled" **in the operator-facing sentence**. That works
/// until someone improves the wording, at which point a refusal quietly becomes a bad request and
/// an administrator goes hunting for a bug in a correctly-written call. The sentence is for a
/// human; the code is for this function.
///
/// The substring check survives as a fallback because an overlay is deployed independently of the
/// sidecar: a phase-1 shard paired with a phase-2 sidecar still sends the codeless frame, and it
/// must keep getting its 403.
fn asset_error_status(value: &Value) -> StatusCode {
match value.get("code").and_then(|c| c.as_str()) {
Some("DISABLED") => return StatusCode::FORBIDDEN,
// The shard has no such file. Not the caller's mistake and not a broken shard -- a client
// that does not carry what was asked for, which the website reports to its operator.
Some("NOT_FOUND") => return StatusCode::NOT_FOUND,
// It has the file and cannot decode it: truncated, hand-edited, or not what it claims to
// be. 422 rather than 400, because the request was fine and repeating it will not help.
Some("UNREADABLE") => return StatusCode::UNPROCESSABLE_ENTITY,
// The shard cannot do this right now (it could not start its asset worker, say). Same
// status as "no shard connected", because it means the same thing to a caller: come back.
Some("UNAVAILABLE") => return StatusCode::SERVICE_UNAVAILABLE,
Some("BAD_REQUEST") => return StatusCode::BAD_REQUEST,
_ => {}
}
let reason = value
.get("reason")
.and_then(|r| r.as_str())
.unwrap_or("request rejected");
if reason.contains("disabled") {
StatusCode::FORBIDDEN
} else {
StatusCode::BAD_REQUEST
}
}
/// The cliloc table, as the shard's own client holds it (docs/link/v8.md §9).
///
/// Until protocol 8 this table reached the website by hand: the operator installed UOFiddler,
/// built a converter against its `Ultima.dll`, ran it over their client's compressed `Cliloc.enu`
/// and copied the result to the web host. The shard could not help, because ServUO's bundled
/// `Ultima.StringList` cannot read a modern client's file either. Phase 2 put the decompressor in
/// the overlay, so the shard reads its own client and the operator installs nothing.
///
/// **Paged, and the caller drives the paging** -- `?cursor=` echoes back whatever the previous
/// reply's `cursor` was, until a reply says `more: false`. The pages are cut by byte budget on the
/// shard (512 KiB against the 1 MiB inbound line cap), so a stock English table arrives in about
/// eleven of them. The sidecar keeps none of it: it has no opinion about what a cliloc is, and a
/// cached copy of five megabytes of EA's strings is exactly what this process should not hold.
///
/// `?lang=` selects the file; it defaults to `enu` on the shard and the shard refuses anything its
/// `Ultima.Files` cannot resolve, which is a 404 rather than a 400.
async fn cliloc_table(
State(st): State<AppState>,
Query(q): Query<ClilocQuery>,
) -> impl IntoResponse {
let req_id = st.rpc.next_req_id();
let mut cmd = json!({"kind": "cliloc.table", "reqId": req_id});
if let Some(lang) = q.lang.as_deref().filter(|s| !s.is_empty()) {
cmd["lang"] = json!(lang);
}
if let Some(cursor) = q.cursor.as_deref().filter(|s| !s.is_empty()) {
cmd["cursor"] = json!(cursor);
}
respond_assets(st.rpc.call(&st.shard, cmd, &req_id).await)
}
#[derive(Deserialize)]
struct ClilocQuery {
lang: Option<String>,
cursor: Option<String>,
}
/// Stage 2 of the import gate (docs/link/v8.md §6, phase 3): every asset the shard could serve,
/// with a hash and a size and **no pixels**.
///
/// That separation is the whole difference between an Update and a re-download. The website holds
/// the hashes from last time, diffs this against them and asks `/assets/fetch` only for the keys
/// that moved — which, on the normal restart-that-changed-nothing, is none of them.
///
/// **Paged, and the caller drives the paging**, same envelope as `/cliloc`: echo the previous
/// reply's `cursor` until one says `more: false`, and read `cut` to learn *why* a page was the
/// last — only `end` means the manifest is complete. This family pages on the shard's wall clock
/// rather than on bytes, because its rows are tiny and building them means decoding hundreds of
/// sprites, so expect several pages of a few hundred rows each.
///
/// `?family=` selects which asset family; phase 3 serves `body` and the shard refuses anything
/// else by name rather than substituting a default.
async fn assets_manifest(
State(st): State<AppState>,
Query(q): Query<ManifestQuery>,
) -> impl IntoResponse {
let req_id = st.rpc.next_req_id();
let mut cmd = json!({"kind": "assets.manifest", "reqId": req_id});
if let Some(family) = q.family.as_deref().filter(|s| !s.is_empty()) {
cmd["family"] = json!(family);
}
if let Some(cursor) = q.cursor.as_deref().filter(|s| !s.is_empty()) {
cmd["cursor"] = json!(cursor);
}
respond_assets(st.rpc.call(&st.shard, cmd, &req_id).await)
}
#[derive(Deserialize)]
struct ManifestQuery {
family: Option<String>,
cursor: Option<String>,
}
/// The bytes, for keys the caller names: `{"keys":[…],"catalog":<opt>,"cursor":<opt>}`.
///
/// Each row carries the sprite as base64 PNG. The shard encodes it, and that is deliberate rather
/// than incidental: `System.Drawing` is already in its *decode* path (docs/link/v8.md §4.2), so
/// PNG costs it no new dependency, while asking the website to encode would put an image encoder
/// in Node and make the manifest's hash cover bytes nobody ever stores.
///
/// **`catalog` is the mid-import guard.** A manifest reply carries a `catalog` id derived from the
/// client files themselves; passing it back here makes the shard refuse (422) if those files moved
/// in between. Without it an operator who patched their client halfway through an import would get
/// one asset set stitched out of two, with no error anywhere.
///
/// A key the shard cannot serve comes back as a **row** with a `status`, not as a failed request —
/// a body this client has no art for is the expected answer for two thirds of the player bodies,
/// and failing the whole page over one would make an import impossible on a stock client.
async fn assets_fetch(State(st): State<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
assets_call(&st, "assets.fetch", body).await
}
/// Slug → body id: `{"types":["GiantSpider", …]}` (docs/link/v8.md §8, phase 3).
///
/// The spawn atlas knows a creature by the class name in `Spawns/*.xml`; the client knows it by a
/// body id; nothing in the ServUO tree declares the mapping as data. Only code running *inside*
/// ServUO can answer it — construct the type, read `Body.BodyID`, delete it — which is why this is
/// a request kind of its own rather than a step inside asset extraction: it runs on the shard's
/// Core thread, while every decode on this plane runs off it.
///
/// **The shard caps the batch and refuses rather than truncates** a longer list, because every
/// name in it costs a real constructor between two ticks of the world. Chunk the list; a 400 here
/// names the cap.
async fn assets_bodies(State(st): State<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
assets_call(&st, "assets.bodies", body).await
}
/// `event_call`'s shape for the asset plane: take the caller's object verbatim, stamp `kind` and
/// `reqId` on it, and map the reply through `respond_assets` rather than `respond_event`.
///
/// The two differ in exactly one way that matters, and it is the reason this is not `event_call`
/// with a different string: `bridge.busy` is a **425** here, and it is the ordinary answer during
/// an import rather than a rare collision. The shard serves one asset request at a time on purpose,
/// so a caller that read busy as an error would abandon a healthy transfer.
async fn assets_call(st: &AppState, kind: &str, body: Value) -> (StatusCode, Json<Value>) {
let mut obj = match body {
Value::Object(m) => m,
Value::Null => serde_json::Map::new(),
_ => {
return (
StatusCode::BAD_REQUEST,
Json(json!({"error": "body must be a JSON object"})),
)
}
};
let req_id = st.rpc.next_req_id();
obj.insert("kind".to_string(), json!(kind));
obj.insert("reqId".to_string(), json!(req_id));
respond_assets(st.rpc.call(&st.shard, Value::Object(obj), &req_id).await)
}
// ---- websocket ---- // ---- websocket ----
async fn ws_upgrade(ws: WebSocketUpgrade, State(state): State<AppState>) -> impl IntoResponse { async fn ws_upgrade(ws: WebSocketUpgrade, State(state): State<AppState>) -> impl IntoResponse {
@@ -1366,6 +1625,152 @@ mod tests {
StatusCode::TOO_EARLY StatusCode::TOO_EARLY
); );
assert_eq!(respond_event(reply("bridge.busy")).0, StatusCode::TOO_EARLY); assert_eq!(respond_event(reply("bridge.busy")).0, StatusCode::TOO_EARLY);
// Protocol 8. On the asset plane `bridge.busy` is not a keyed retry colliding with itself
// -- it is flow control, and it is the ORDINARY answer during an import rather than a rare
// one. The shard serves one asset request at a time because its outbound queue is bounded
// in lines, not bytes, so a queue of large replies is how it runs out of memory. A
// responder that answered 200 here would tell the website an import step succeeded and
// returned nothing.
assert_eq!(
respond_assets(reply("bridge.busy")).0,
StatusCode::TOO_EARLY
);
}
/// The asset plane's own gate, and it is a refusal rather than a mistake: an operator who has
/// not enabled `Bridge.AssetsEnabled` has declined to let the website read their UO client
/// files off the shard host. 403, for the same reason the event plane's switch is a 403.
#[test]
fn assets_disabled_is_a_403() {
let value = json!({
"kind": "assets.error",
"reqId": "r-9",
"reason": "asset extraction is disabled on this shard"
});
assert_eq!(respond_assets(Ok(value)).0, StatusCode::FORBIDDEN);
}
/// Anything else the shard refuses on this plane is the caller's mistake.
#[test]
fn other_asset_errors_are_400() {
let value = json!({
"kind": "assets.error",
"reason": "assets.sources requires a reqId"
});
assert_eq!(respond_assets(Ok(value)).0, StatusCode::BAD_REQUEST);
}
/// Phase 2's codes, each of which says something a 400 does not. `NOT_FOUND` is a client that
/// does not carry the file (an operator fact, not a bug); `UNREADABLE` is a file that is there
/// and cannot be decoded, where repeating the request cannot help; `UNAVAILABLE` is the shard
/// declining for now.
#[test]
fn an_asset_error_code_picks_the_status() {
let with = |code: &str| {
respond_assets(Ok(json!({
"kind": "assets.error",
"reqId": "r-1",
"code": code,
"reason": ""
})))
.0
};
assert_eq!(with("DISABLED"), StatusCode::FORBIDDEN);
assert_eq!(with("NOT_FOUND"), StatusCode::NOT_FOUND);
assert_eq!(with("UNREADABLE"), StatusCode::UNPROCESSABLE_ENTITY);
assert_eq!(with("UNAVAILABLE"), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(with("BAD_REQUEST"), StatusCode::BAD_REQUEST);
// An unknown code is not a reason to invent a status.
assert_eq!(with("SOMETHING_NEW"), StatusCode::BAD_REQUEST);
}
/// The overlay and the sidecar ship separately and an operator can run a phase-1 shard against
/// a phase-2 sidecar, so the codeless refusal must keep its 403. This is the whole reason the
/// substring check survives rather than being deleted with the code it preceded.
#[test]
fn a_codeless_disabled_refusal_is_still_a_403() {
let value = json!({
"kind": "assets.error",
"reqId": "r-9",
"reason": "asset extraction is disabled on this shard"
});
assert_eq!(respond_assets(Ok(value)).0, StatusCode::FORBIDDEN);
}
/// One page of the cliloc table. The assertion that matters is the envelope: `more` and
/// `cursor` reach the caller untouched, because the sidecar does not page this — the shard
/// cuts the pages and the website drives them, and a sidecar that "helpfully" assembled them
/// would be holding the whole table in memory to do it.
#[test]
fn a_cliloc_page_is_a_200_and_keeps_its_cursor() {
let value = json!({
"kind": "cliloc.table.ok",
"reqId": "r-2",
"lang": "enu",
"extractorVersion": 1,
"total": 67496,
"rows": [{"n": 1023721, "f": 0, "t": "quarter staff"}],
"more": true,
"cursor": "n:1023721",
"cut": "budget"
});
let (status, body) = respond_assets(Ok(value));
assert_eq!(status, StatusCode::OK);
assert_eq!(body.0.get("more").and_then(|v| v.as_bool()), Some(true));
assert_eq!(
body.0.get("cursor").and_then(|v| v.as_str()),
Some("n:1023721")
);
assert_eq!(body.0.get("total").and_then(|v| v.as_i64()), Some(67496));
}
/// A source manifest comes back whole. Worth asserting because `respond_assets` sniffs `kind`
/// and a family whose success kind ends in `.ok` sits one character away from the `.error`
/// suffix the generic responder matches on -- which is exactly why this plane has its own
/// responder and matches `assets.error` exactly rather than by suffix.
#[test]
fn a_source_manifest_is_a_200() {
let value = json!({
"kind": "assets.sources.ok",
"reqId": "r-9",
"extractorVersion": 1,
"imaging": {"ok": true},
"files": [],
"more": false,
"cut": "end"
});
let (status, body) = respond_assets(Ok(value));
assert_eq!(status, StatusCode::OK);
assert_eq!(
body.0.get("extractorVersion").and_then(|v| v.as_i64()),
Some(1)
);
}
/// A shard that is not connected is a 503 and a shard that did not answer in time is a 504,
/// and the asset plane needs the second one to stay distinct more than any other plane does:
/// hashing a 195 MB anim.mul is the one thing on this link that can genuinely outlast the
/// 10 s reply timeout, and the website's response to that is to poll again rather than to
/// declare the shard down.
#[test]
fn asset_transport_failures_keep_their_own_statuses() {
assert_eq!(
respond_assets(Err(RpcError::NoShard)).0,
StatusCode::SERVICE_UNAVAILABLE
);
assert_eq!(
respond_assets(Err(RpcError::Timeout)).0,
StatusCode::GATEWAY_TIMEOUT
);
} }
/// The event plane is the FIRST place `bridge.busy` is reachable on a live shard rather than /// The event plane is the FIRST place `bridge.busy` is reachable on a live shard rather than