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
This commit is contained in:
2026-09-10 15:04:24 +00:00
3 changed files with 429 additions and 20 deletions

View File

@@ -91,7 +91,32 @@ use tracing_subscriber::EnvFilter;
/// the dumb-forwarder property again.
///
/// **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.
///
/// **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
// 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
//! accept one shard connection at a time and re-accept when it drops (the shard reconnects on its
//! 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 serde_json::Value;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpListener;
use tokio::sync::{mpsc, Mutex};
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.
#[derive(Debug, Clone)]
pub struct ShardEvent {
@@ -112,31 +227,41 @@ async fn handle_connection(
handle.set(Some(cmd_tx)).await;
let mut reader = BufReader::new(read_half);
let mut line = String::new();
let mut lines = LineReader::default();
loop {
tokio::select! {
// Inbound: a line from the shard.
result = reader.read_line(&mut line) => {
let n = result?;
if n == 0 {
return Ok(()); // clean EOF: shard closed
}
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 });
result = lines.next(&mut reader) => {
match result? {
Line::Eof => return Ok(()), // clean EOF: shard closed
Line::TooLong(bytes) => {
// Deliberately not a disconnect. See MAX_INBOUND_LINE_BYTES: a reply lost
// this way times out on the caller's side and is re-requested, and tearing
// the link down would take the live event feed with it.
warn!(
bytes,
cap = MAX_INBOUND_LINE_BYTES,
"inbound line over the cap; discarded"
);
}
Line::Complete(line) => {
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.
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,12 @@ pub async fn serve(addr: &str, state: AppState) -> anyhow::Result<()> {
// 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))
// 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))
.route_layer(middleware::from_fn_with_state(state.clone(), gate));
let app = Router::new()
@@ -1296,6 +1302,76 @@ 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" {
let reason = value
.get("reason")
.and_then(|r| r.as_str())
.unwrap_or("request rejected");
let code = if reason.contains("disabled") {
StatusCode::FORBIDDEN
} else {
StatusCode::BAD_REQUEST
};
(code, 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"})),
),
}
}
// ---- websocket ----
async fn ws_upgrade(ws: WebSocketUpgrade, State(state): State<AppState>) -> impl IntoResponse {
@@ -1366,6 +1442,84 @@ mod tests {
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);
}
/// 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