From 4ed15c1a38735e9443e5bbb96de727890756c748 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 19:07:47 +0100 Subject: [PATCH] feat(simulation): add LocalBridge IPC over Unix socket (#78) Length-prefixed MessagePack framing (4-byte BE length + payload), LocalBridge struct implementing SimBridge trait over Unix domain sockets, BridgeResource wrapper for ECS integration. Adds Io error variant to BridgeError. Two integration tests verify snapshot and input round-trips over real sockets. Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/framing.rs | 109 ++++++++++++++++++++++++++++++ server/src/bridge/local.rs | 118 ++++++++++++++++++++++++++++++++ server/src/bridge/mod.rs | 27 ++++++++ server/tests/bridge_ipc.rs | 126 +++++++++++++++++++++++++++++++++++ 4 files changed, 380 insertions(+) create mode 100644 server/src/bridge/framing.rs create mode 100644 server/src/bridge/local.rs create mode 100644 server/tests/bridge_ipc.rs diff --git a/server/src/bridge/framing.rs b/server/src/bridge/framing.rs new file mode 100644 index 000000000..1ec3f928e --- /dev/null +++ b/server/src/bridge/framing.rs @@ -0,0 +1,109 @@ +// MessagePack framing protocol +// 4-byte big-endian length prefix + payload +// Implements D-020 IPC transport layer + +use std::io::{self, Read, Write}; + +/// Maximum message size: 16 MB +const MAX_MESSAGE_SIZE: u32 = 16 * 1024 * 1024; + +/// Write a length-prefixed message to a writer. +/// Format: [4-byte BE length][payload] +pub fn write_framed(writer: &mut impl Write, payload: &[u8]) -> io::Result<()> { + let len = payload.len() as u32; + if len > MAX_MESSAGE_SIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!( + "message too large: {} bytes (max {})", + len, MAX_MESSAGE_SIZE + ), + )); + } + + writer.write_all(&len.to_be_bytes())?; + writer.write_all(payload)?; + writer.flush()?; + Ok(()) +} + +/// Read a length-prefixed message from a reader. +/// Returns Ok(None) on clean EOF (connection closed). +/// Returns error on incomplete/corrupted reads. +pub fn read_framed(reader: &mut impl Read) -> io::Result>> { + // Read 4-byte length prefix + let mut len_bytes = [0u8; 4]; + match reader.read_exact(&mut len_bytes) { + Ok(()) => {} + Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(None), + Err(e) => return Err(e), + } + + let len = u32::from_be_bytes(len_bytes); + + if len > MAX_MESSAGE_SIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "message too large: {} bytes (max {})", + len, MAX_MESSAGE_SIZE + ), + )); + } + + // Read payload + let mut payload = vec![0u8; len as usize]; + reader.read_exact(&mut payload)?; + Ok(Some(payload)) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Cursor; + + #[test] + fn write_then_read_roundtrip() { + let payload = b"hello, world!"; + let mut buffer = Vec::new(); + + write_framed(&mut buffer, payload).expect("write failed"); + + let mut cursor = Cursor::new(buffer); + let result = read_framed(&mut cursor).expect("read failed"); + + assert_eq!(result.unwrap(), payload); + } + + #[test] + fn empty_payload_roundtrip() { + let payload = b""; + let mut buffer = Vec::new(); + + write_framed(&mut buffer, payload).expect("write failed"); + + let mut cursor = Cursor::new(buffer); + let result = read_framed(&mut cursor).expect("read failed"); + + assert_eq!(result.unwrap(), payload); + } + + #[test] + fn rejects_oversized_message() { + let mut buffer = Vec::new(); + let oversized_payload = vec![0u8; (MAX_MESSAGE_SIZE + 1) as usize]; + + let result = write_framed(&mut buffer, &oversized_payload); + assert!(result.is_err()); + assert!(result.unwrap_err().to_string().contains("too large")); + } + + #[test] + fn eof_returns_none() { + let buffer = Vec::new(); + let mut cursor = Cursor::new(buffer); + + let result = read_framed(&mut cursor).expect("read failed"); + assert_eq!(result, None); + } +} diff --git a/server/src/bridge/local.rs b/server/src/bridge/local.rs new file mode 100644 index 000000000..f558a135c --- /dev/null +++ b/server/src/bridge/local.rs @@ -0,0 +1,118 @@ +// LocalBridge - Unix socket IPC implementation +// Implements D-020 subprocess/IPC architecture +// Deterministic client-server communication via Unix domain sockets + +use super::{BridgeError, ObserverSnapshot, PlayerInput, SimBridge}; +use crate::bridge::framing::{read_framed, write_framed}; +use std::fs; +use std::io::{BufReader, BufWriter}; +use std::os::unix::net::{UnixListener, UnixStream}; +use std::path::{Path, PathBuf}; +use std::sync::Mutex; + +/// LocalBridge: Unix socket transport for client-server IPC +pub struct LocalBridge { + reader: Mutex>, + writer: Mutex>, + socket_path: PathBuf, +} + +impl LocalBridge { + /// Server-side: create a Unix socket listener and accept one connection. + /// Removes any stale socket file before binding. + pub fn accept(path: &Path) -> Result { + // Remove stale socket if it exists + if path.exists() { + fs::remove_file(path).map_err(|e| { + BridgeError::Transport(format!("failed to remove stale socket: {}", e)) + })?; + } + + tracing::info!("LocalBridge listening on {:?}", path); + + let listener = UnixListener::bind(path) + .map_err(|e| BridgeError::Transport(format!("failed to bind Unix socket: {}", e)))?; + + // Accept one connection + let (stream, _addr) = listener + .accept() + .map_err(|e| BridgeError::Transport(format!("failed to accept connection: {}", e)))?; + + tracing::info!("LocalBridge accepted connection on {:?}", path); + + // Clone stream for reader and writer + let reader_stream = stream.try_clone().map_err(|e| { + BridgeError::Transport(format!("failed to clone stream for reader: {}", e)) + })?; + + Ok(Self { + reader: Mutex::new(BufReader::new(reader_stream)), + writer: Mutex::new(BufWriter::new(stream)), + socket_path: path.to_path_buf(), + }) + } + + /// Client-side: connect to an existing Unix socket. + pub fn connect(path: &Path) -> Result { + tracing::info!("LocalBridge connecting to {:?}", path); + + let stream = UnixStream::connect(path).map_err(|e| { + BridgeError::Transport(format!("failed to connect to Unix socket: {}", e)) + })?; + + tracing::trace!("LocalBridge connected to {:?}", path); + + // Clone stream for reader and writer + let reader_stream = stream.try_clone().map_err(|e| { + BridgeError::Transport(format!("failed to clone stream for reader: {}", e)) + })?; + + Ok(Self { + reader: Mutex::new(BufReader::new(reader_stream)), + writer: Mutex::new(BufWriter::new(stream)), + socket_path: path.to_path_buf(), + }) + } + + pub fn socket_path(&self) -> &Path { + &self.socket_path + } +} + +impl SimBridge for LocalBridge { + fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> { + let payload = rmp_serde::to_vec(snapshot)?; + + let mut writer = self.writer.lock().unwrap(); + write_framed(writer.get_mut(), &payload)?; + + tracing::trace!("sent snapshot: tick={}", snapshot.tick); + Ok(()) + } + + fn receive_inputs(&self) -> Result, BridgeError> { + let mut reader = self.reader.lock().unwrap(); + + match read_framed(reader.get_mut())? { + Some(payload) => { + let inputs: Vec = rmp_serde::from_slice(&payload)?; + tracing::trace!("received {} inputs", inputs.len()); + Ok(inputs) + } + None => { + tracing::trace!("received EOF, returning empty input vec"); + Ok(Vec::new()) + } + } + } +} + +impl Drop for LocalBridge { + fn drop(&mut self) { + // Best-effort socket cleanup + if self.socket_path.exists() { + let _ = fs::remove_file(&self.socket_path); + tracing::trace!("removed socket file {:?}", self.socket_path); + } + } +} diff --git a/server/src/bridge/mod.rs b/server/src/bridge/mod.rs index 50619a4b1..05a398abe 100644 --- a/server/src/bridge/mod.rs +++ b/server/src/bridge/mod.rs @@ -3,7 +3,10 @@ // MessagePack serialization for Rust<->Godot communication use bevy_app::prelude::*; +use bevy_ecs::prelude::*; +pub mod framing; +pub mod local; pub mod types; pub use types::*; @@ -14,6 +17,8 @@ pub enum BridgeError { Serialization(#[from] rmp_serde::encode::Error), #[error("deserialization error: {0}")] Deserialization(#[from] rmp_serde::decode::Error), + #[error("io error: {0}")] + Io(#[from] std::io::Error), #[error("transport error: {0}")] Transport(String), } @@ -28,6 +33,28 @@ pub trait SimBridge: Send + Sync { fn receive_inputs(&self) -> Result, BridgeError>; } +/// BridgeResource: Bevy Resource wrapper for SimBridge trait object +#[derive(Resource)] +pub struct BridgeResource { + inner: Box, +} + +impl BridgeResource { + pub fn new(bridge: impl SimBridge + 'static) -> Self { + Self { + inner: Box::new(bridge), + } + } + + pub fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> { + self.inner.send_snapshot(snapshot) + } + + pub fn receive_inputs(&self) -> Result, BridgeError> { + self.inner.receive_inputs() + } +} + /// Bridge plugin for client-server communication /// Abstracts transport layer (LocalBridge/NetworkBridge) pub struct BridgePlugin; diff --git a/server/tests/bridge_ipc.rs b/server/tests/bridge_ipc.rs new file mode 100644 index 000000000..448144f57 --- /dev/null +++ b/server/tests/bridge_ipc.rs @@ -0,0 +1,126 @@ +//! Integration tests for LocalBridge over Unix sockets (D-030 Layer 2: IPC roundtrip). + +use settled_reach_server::bridge::framing::{read_framed, write_framed}; +use settled_reach_server::bridge::local::LocalBridge; +use settled_reach_server::bridge::types::*; +use settled_reach_server::bridge::SimBridge; +use std::os::unix::net::UnixStream; +use std::path::PathBuf; +use std::thread; +use std::time::Duration; + +/// Generate unique socket path for test isolation +fn test_socket_path(test_name: &str) -> PathBuf { + let pid = std::process::id(); + let timestamp = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis(); + PathBuf::from(format!( + "/tmp/sr-test-{}-{}-{}.sock", + test_name, pid, timestamp + )) +} + +#[test] +fn snapshot_roundtrip_over_unix_socket() { + let socket_path = test_socket_path("snapshot"); + + // Server thread: accept connection and send snapshot + let server_path = socket_path.clone(); + let server_handle = thread::spawn(move || { + let bridge = LocalBridge::accept(&server_path).expect("failed to accept"); + + let snapshot = ObserverSnapshot { + tick: 42, + entities: vec![VisibleEntity { + entity_id: 100, + x: 10.5, + y: 20.3, + z: 0, + kind: EntityKind::Npc, + }], + }; + + bridge + .send_snapshot(&snapshot) + .expect("failed to send snapshot"); + }); + + // Give server time to bind + thread::sleep(Duration::from_millis(50)); + + // Client: connect and receive snapshot + let stream = UnixStream::connect(&socket_path).expect("failed to connect"); + let mut reader = std::io::BufReader::new(stream); + + let payload = read_framed(&mut reader) + .expect("failed to read frame") + .expect("unexpected EOF"); + + let snapshot: ObserverSnapshot = + rmp_serde::from_slice(&payload).expect("failed to deserialize"); + + assert_eq!(snapshot.tick, 42); + assert_eq!(snapshot.entities.len(), 1); + assert_eq!(snapshot.entities[0].entity_id, 100); + assert_eq!(snapshot.entities[0].x, 10.5); + assert_eq!(snapshot.entities[0].y, 20.3); + + server_handle.join().expect("server thread panicked"); +} + +#[test] +fn input_roundtrip_over_unix_socket() { + let socket_path = test_socket_path("input"); + + // Server thread: accept connection and receive inputs + let server_path = socket_path.clone(); + let server_handle = thread::spawn(move || { + let bridge = LocalBridge::accept(&server_path).expect("failed to accept"); + + let inputs = bridge.receive_inputs().expect("failed to receive inputs"); + + assert_eq!(inputs.len(), 2); + assert_eq!(inputs[0].tick, 10); + assert_eq!(inputs[1].tick, 11); + + inputs + }); + + // Give server time to bind + thread::sleep(Duration::from_millis(50)); + + // Client: connect and send inputs + let stream = UnixStream::connect(&socket_path).expect("failed to connect"); + let mut writer = std::io::BufWriter::new(stream); + + let inputs = vec![ + PlayerInput { + tick: 10, + action: PlayerAction::MoveNorth, + }, + PlayerInput { + tick: 11, + action: PlayerAction::Interact, + }, + ]; + + let payload = rmp_serde::to_vec(&inputs).expect("failed to serialize"); + write_framed(&mut writer, &payload).expect("failed to write frame"); + + // Drop writer to close connection and signal EOF to server + drop(writer); + + let received_inputs = server_handle.join().expect("server thread panicked"); + + // Verify actions survived the round-trip + match &received_inputs[0].action { + PlayerAction::MoveNorth => {} + _ => panic!("expected MoveNorth action"), + } + match &received_inputs[1].action { + PlayerAction::Interact => {} + _ => panic!("expected Interact action"), + } +}