Files
Rust-Link/sidecar/src/game.rs
wtclaude 3f5ed059b8
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 4m1s
fix(sidecar): refuse a plugin that names another server (D155)
`[game].server_id` was a cross-check that WARNED and kept the plugin's id.
The phase 18 walk showed what that costs: a second server's plugin,
parked in this listener's backlog by a plugin bug (Rust-Plugins, D154),
was accepted the moment the first server's plugin reloaded, and the
website showed server "alpha" with beta's hostname and wipe.

Now, with `server_id` set, the connection is closed on the first frame
that names another server, BEFORE that frame reaches the store or the
feed, and both ids are logged at ERROR. The command channel is installed
only once a frame has named this server, so no website command (a grant,
a world write) can reach a plugin about to be refused, and /health reports
the plugin connected only from then. Blank `server_id`: nothing checked,
as before.

The egg's launcher now hands the sidecar the plugin config's ServerId once
that file exists. The plugin reads RUSTLINK_SERVER_ID only at its first
config write (D150); without this, a variable edited after the first boot
would be refused instead of changing nothing, as INSTALL.md promises.

Walked: a plugin aimed at another server's sidecar is refused with an
ERROR naming both ids, and the site keeps the right server's identity.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01E14m6SuuY6i1vASFeGDBeY
2026-09-26 01:19:50 -05:00

533 lines
22 KiB
Rust

//! The loopback link to the Rust server's Oxide bridge plugin.
//!
//! The plugin is the TCP *client*: it dials out to us. So the sidecar owns the listener, and the
//! plugin's outbound socket is the only thing that ever connects. This is the whole reason the game
//! is never directly reachable from the website — it exposes no port for us.
//!
//! Framing is newline-delimited JSON, bidirectional: the plugin sends events (`kind`), we send
//! commands (`cmd`). We accept one plugin connection at a time and re-accept when it drops (the
//! plugin reconnects on its own, with a bounded backoff).
//!
//! **Loopback is the trust boundary.** There is no token on this link, exactly as on the ServUO
//! bridge: the plugin and the sidecar share a host, and the address to bind is `127.0.0.1`. Binding
//! `[game].bind` to anything routable puts an unauthenticated command channel on the network, and
//! the operator documentation says so in as many words.
//!
//! Inbound lines are **capped** (see [`MAX_INBOUND_LINE_BYTES`]) from the start rather than after
//! the first large frame arrives: an unbounded `read_line` facing a peer that will one day send a
//! map image is a memory-exhaustion shape we would be inventing ourselves.
use std::sync::Arc;
use serde_json::Value;
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpListener;
use tokio::sync::{mpsc, Mutex};
use tracing::{error, info, warn};
/// The longest line the sidecar will accept from the plugin, in bytes.
///
/// Over-long lines are **discarded, not buffered**, and the connection stays up: a single malformed
/// frame is not a reason to tear down a link that live events are flowing over. A dropped reply
/// simply times out on the caller's side and is re-requested.
pub const MAX_INBOUND_LINE_BYTES: usize = 1024 * 1024;
/// What one read off the plugin 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 plugin 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 plugin, parsed. `kind` is lifted out for routing.
#[derive(Debug, Clone)]
pub struct GameEvent {
pub kind: String,
pub value: Value,
}
/// A handle for sending command lines to the plugin. Cloneable and cheap.
///
/// Commands are dropped (with a warning) when no plugin is connected, rather than buffered: a
/// website query that arrives during an outage should fail fast and be retried, not silently queue
/// behind a reconnect. Live *events* are what must survive an outage, and those the plugin buffers
/// on its side.
#[derive(Clone)]
pub struct GameHandle {
tx: Arc<Mutex<Option<mpsc::UnboundedSender<String>>>>,
}
impl GameHandle {
fn new() -> Self {
Self {
tx: Arc::new(Mutex::new(None)),
}
}
async fn set(&self, sender: Option<mpsc::UnboundedSender<String>>) {
*self.tx.lock().await = sender;
}
/// Send one command line (a complete JSON object, no newline — we add the frame delimiter).
/// Returns false if no plugin is currently connected.
pub async fn send(&self, line: String) -> bool {
let guard = self.tx.lock().await;
match guard.as_ref() {
Some(sender) => sender.send(line).is_ok(),
None => {
warn!("dropping command; no plugin connected");
false
}
}
}
pub async fn is_connected(&self) -> bool {
self.tx.lock().await.is_some()
}
}
/// Binds the loopback listener and accepts plugin connections forever. Each accepted connection
/// runs until it drops, then we loop back to accept the next one. Events are forwarded to
/// `event_tx`; the returned handle sends commands to whichever plugin is currently connected.
///
/// The **bound** address is returned alongside the handle rather than assumed to be the one asked
/// for: `127.0.0.1:0` is a legitimate thing to configure (and what the tests use), and a log line
/// echoing the request rather than the result is the kind that is wrong exactly when it matters.
///
/// `server_id` is `[game].server_id`. When it is set, a plugin that names another server is
/// refused — see [`foreign_server_id`].
pub async fn serve(
addr: &str,
server_id: String,
event_tx: mpsc::UnboundedSender<GameEvent>,
) -> std::io::Result<(GameHandle, std::net::SocketAddr)> {
let listener = TcpListener::bind(addr).await?;
let bound = listener.local_addr()?;
info!(addr = %bound, "game link listening");
let handle = GameHandle::new();
let accept_handle = handle.clone();
tokio::spawn(async move {
loop {
match listener.accept().await {
Ok((stream, peer)) => {
info!(%peer, "plugin connected");
if let Err(e) =
handle_connection(stream, &server_id, &event_tx, &accept_handle).await
{
warn!(error = %e, "plugin connection ended");
} else {
info!("plugin disconnected");
}
accept_handle.set(None).await;
// A disconnect is a fact the website should see without polling, so it rides
// the same channel every other fact does. Nothing persists it: it is this
// process's own observation, not something the game said — which is exactly
// what `type: "control"` means (PROTOCOL.md §8.1). It is synthesised here
// rather than anywhere else because this is the only place that knows.
let _ = event_tx.send(GameEvent {
kind: "link.down".to_string(),
value: serde_json::json!({ "kind": "link.down", "type": "control" }),
});
}
Err(e) => {
warn!(error = %e, "accept failed");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
}
}
}
});
Ok((handle, bound))
}
/// The server a frame names, when it is not the one this sidecar is configured for.
///
/// `None` means the frame may pass: no `server_id` is configured, or the frame names none, or it
/// names this one. The plugin stamps `serverId` on every frame, so the first line of a connection
/// is enough to tell — and the check runs before that line reaches the store or the feed, because
/// by the time it is filed, one server's history is already under another's name (D155). That is
/// what the phase 18 walk saw when this was a warning: a second server's plugin, parked in this
/// listener's backlog, was accepted the moment the first one reloaded, and the website showed the
/// first server with the second's hostname and wipe.
fn foreign_server_id<'a>(configured: &str, frame: &'a Value) -> Option<&'a str> {
if configured.is_empty() {
return None;
}
match frame.get("serverId").and_then(|v| v.as_str()) {
Some(announced) if !announced.is_empty() && announced != configured => Some(announced),
_ => None,
}
}
async fn handle_connection(
stream: tokio::net::TcpStream,
server_id: &str,
event_tx: &mpsc::UnboundedSender<GameEvent>,
handle: &GameHandle,
) -> std::io::Result<()> {
stream.set_nodelay(true).ok();
let (read_half, mut write_half) = stream.into_split();
// The outbound command channel for this connection. With a `server_id` configured it is
// installed only once a frame has named this server: until then the peer could be a plugin
// about to be refused, and a command sent in that window — a permission grant, a world write —
// would land on the wrong server. `/health` reads the same handle, so it reports the plugin
// connected only from that moment too.
let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel::<String>();
let mut pending_tx = Some(cmd_tx);
if server_id.is_empty() {
handle.set(pending_tx.take()).await;
}
let mut reader = BufReader::new(read_half);
let mut lines = LineReader::default();
loop {
tokio::select! {
// Inbound: a line from the plugin.
result = lines.next(&mut reader) => {
match result? {
Line::Eof => return Ok(()), // clean EOF: plugin closed
Line::TooLong(bytes) => {
// Deliberately not a disconnect. See MAX_INBOUND_LINE_BYTES.
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) => {
if let Some(announced) = foreign_server_id(server_id, &value) {
error!(
configured = server_id,
announced,
"refusing a plugin that names another server: this \
sidecar is '{server_id}', the plugin says \
'{announced}'. Two servers are dialling one game \
port — check each plugin config's Port against its \
sidecar's [game].bind"
);
return Ok(());
}
if pending_tx.is_some()
&& value.get("serverId").and_then(|v| v.as_str())
== Some(server_id)
{
handle.set(pending_tx.take()).await;
}
let kind = value
.get("kind")
.and_then(|k| k.as_str())
.unwrap_or("")
.to_string();
let _ = event_tx.send(GameEvent { kind, value });
}
Err(e) => warn!(error = %e, line = %trimmed, "unparseable event"),
}
}
}
}
}
// Outbound: a command to write to the plugin.
cmd = cmd_rx.recv() => {
match cmd {
Some(mut c) => {
c.push('\n');
write_half.write_all(c.as_bytes()).await?;
write_half.flush().await?;
}
None => return Ok(()), // channel closed
}
}
}
}
}
#[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` would have done 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.
#[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 plugin 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}"]);
}
/// The end-to-end shape of the link, over a real socket: the plugin dials in, sends a hello,
/// and receives a command. This is the criterion phase 1 is judged on, minus the two peers.
#[tokio::test]
async fn a_dialling_plugin_is_accepted_and_can_be_commanded() {
use tokio::io::AsyncReadExt;
let (tx, mut rx) = mpsc::unbounded_channel();
let (handle, addr) = serve("127.0.0.1:0", String::new(), tx).await.unwrap();
let mut client = tokio::net::TcpStream::connect(addr).await.unwrap();
client
.write_all(b"{\"kind\":\"server.hello\",\"serverId\":\"main\"}\n")
.await
.unwrap();
let ev = rx.recv().await.unwrap();
assert_eq!(ev.kind, "server.hello");
assert_eq!(ev.value["serverId"], "main");
assert!(handle.is_connected().await);
assert!(handle.send("{\"cmd\":\"ping\"}".to_string()).await);
let mut buf = [0u8; 64];
let n = client.read(&mut buf).await.unwrap();
assert_eq!(&buf[..n], b"{\"cmd\":\"ping\"}\n");
}
#[test]
fn foreign_server_id_names_only_a_disagreement() {
let beta = serde_json::json!({ "kind": "server.hello", "serverId": "beta" });
let alpha = serde_json::json!({ "kind": "server.hello", "serverId": "alpha" });
let anonymous = serde_json::json!({ "kind": "pong" });
let blank = serde_json::json!({ "kind": "pong", "serverId": "" });
assert_eq!(foreign_server_id("alpha", &beta), Some("beta"));
assert_eq!(foreign_server_id("alpha", &alpha), None);
assert_eq!(foreign_server_id("alpha", &anonymous), None);
assert_eq!(foreign_server_id("alpha", &blank), None);
// Nothing configured: nothing to check against.
assert_eq!(foreign_server_id("", &beta), None);
}
/// D155, over a real socket: the plugin that names another server is closed on its first
/// frame, the frame is never forwarded, and no command could have reached it meanwhile.
#[tokio::test]
async fn a_plugin_naming_another_server_is_refused_before_anything_is_filed() {
use tokio::io::AsyncReadExt;
let (tx, mut rx) = mpsc::unbounded_channel();
let (handle, addr) = serve("127.0.0.1:0", "alpha".to_string(), tx).await.unwrap();
let mut client = tokio::net::TcpStream::connect(addr).await.unwrap();
// Accepted, but not yet vetted: no command may go to it.
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
assert!(!handle.is_connected().await);
client
.write_all(b"{\"kind\":\"server.hello\",\"type\":\"board\",\"serverId\":\"beta\"}\n")
.await
.unwrap();
let mut buf = [0u8; 16];
let n = tokio::time::timeout(std::time::Duration::from_secs(2), client.read(&mut buf))
.await
.expect("the sidecar closes a refused plugin")
.unwrap_or(0);
assert_eq!(n, 0);
// The only thing forwarded is the listener's own link.down, never beta's hello.
let ev = rx.recv().await.unwrap();
assert_eq!(ev.kind, "link.down");
assert!(!handle.is_connected().await);
}
#[tokio::test]
async fn the_command_channel_opens_on_the_first_frame_naming_this_server() {
let (tx, mut rx) = mpsc::unbounded_channel();
let (handle, addr) = serve("127.0.0.1:0", "alpha".to_string(), tx).await.unwrap();
let mut client = tokio::net::TcpStream::connect(addr).await.unwrap();
client
.write_all(b"{\"kind\":\"server.hello\",\"serverId\":\"alpha\"}\n")
.await
.unwrap();
let ev = rx.recv().await.unwrap();
assert_eq!(ev.value["serverId"], "alpha");
assert!(handle.is_connected().await);
}
}