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 <noreply@anthropic.com>
This commit is contained in:
@@ -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<Option<Vec<u8>>> {
|
||||
// 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);
|
||||
}
|
||||
}
|
||||
@@ -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<BufReader<UnixStream>>,
|
||||
writer: Mutex<BufWriter<UnixStream>>,
|
||||
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<Self, BridgeError> {
|
||||
// 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<Self, BridgeError> {
|
||||
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<Vec<PlayerInput>, BridgeError> {
|
||||
let mut reader = self.reader.lock().unwrap();
|
||||
|
||||
match read_framed(reader.get_mut())? {
|
||||
Some(payload) => {
|
||||
let inputs: Vec<PlayerInput> = 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<Vec<PlayerInput>, BridgeError>;
|
||||
}
|
||||
|
||||
/// BridgeResource: Bevy Resource wrapper for SimBridge trait object
|
||||
#[derive(Resource)]
|
||||
pub struct BridgeResource {
|
||||
inner: Box<dyn SimBridge>,
|
||||
}
|
||||
|
||||
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<Vec<PlayerInput>, BridgeError> {
|
||||
self.inner.receive_inputs()
|
||||
}
|
||||
}
|
||||
|
||||
/// Bridge plugin for client-server communication
|
||||
/// Abstracts transport layer (LocalBridge/NetworkBridge)
|
||||
pub struct BridgePlugin;
|
||||
|
||||
@@ -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"),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user