From f6ce7d4046ab7d5faf88454b90b6e6629fa79548 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:01:15 +0100 Subject: [PATCH 1/8] fix(bridge): use named MessagePack format for wire compatibility Client-side Protocol.gd expects rmp_serde::to_vec_named() (maps with string keys), but LocalBridge was using to_vec() (compact positional arrays). Fix send_snapshot and update all test serialization calls to match actual wire format. Also change EOF from Ok(vec![]) to BridgeError::Transport so bridge systems can detect disconnects. Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/local.rs | 7 ++----- server/tests/bridge_ipc.rs | 2 +- server/tests/serialization.rs | 21 +++++++++++++-------- 3 files changed, 16 insertions(+), 14 deletions(-) diff --git a/server/src/bridge/local.rs b/server/src/bridge/local.rs index 6669ea773..62c36f873 100644 --- a/server/src/bridge/local.rs +++ b/server/src/bridge/local.rs @@ -81,7 +81,7 @@ impl LocalBridge { impl SimBridge for LocalBridge { fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> { - let payload = rmp_serde::to_vec(snapshot)?; + let payload = rmp_serde::to_vec_named(snapshot)?; let mut writer = self.writer.lock().expect("writer mutex poisoned"); write_framed(writer.get_mut(), &payload)?; @@ -99,10 +99,7 @@ impl SimBridge for LocalBridge { tracing::trace!("received {} inputs", inputs.len()); Ok(inputs) } - None => { - tracing::trace!("received EOF, returning empty input vec"); - Ok(Vec::new()) - } + None => Err(BridgeError::Transport("client disconnected (EOF)".into())), } } } diff --git a/server/tests/bridge_ipc.rs b/server/tests/bridge_ipc.rs index 448144f57..3d110b8bd 100644 --- a/server/tests/bridge_ipc.rs +++ b/server/tests/bridge_ipc.rs @@ -106,7 +106,7 @@ fn input_roundtrip_over_unix_socket() { }, ]; - let payload = rmp_serde::to_vec(&inputs).expect("failed to serialize"); + let payload = rmp_serde::to_vec_named(&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 diff --git a/server/tests/serialization.rs b/server/tests/serialization.rs index abe61f56a..af60a0d75 100644 --- a/server/tests/serialization.rs +++ b/server/tests/serialization.rs @@ -15,7 +15,7 @@ fn observer_snapshot_roundtrip() { }], }; - let bytes = rmp_serde::to_vec(&snapshot).expect("serialize"); + let bytes = rmp_serde::to_vec_named(&snapshot).expect("serialize"); let decoded: ObserverSnapshot = rmp_serde::from_slice(&bytes).expect("deserialize"); assert_eq!(decoded.tick, 42); @@ -30,7 +30,7 @@ fn player_input_roundtrip() { action: PlayerAction::MoveNorth, }; - let bytes = rmp_serde::to_vec(&input).expect("serialize"); + let bytes = rmp_serde::to_vec_named(&input).expect("serialize"); let decoded: PlayerInput = rmp_serde::from_slice(&bytes).expect("deserialize"); assert_eq!(decoded.tick, 100); @@ -43,7 +43,7 @@ fn empty_snapshot_roundtrip() { entities: vec![], }; - let bytes = rmp_serde::to_vec(&snapshot).expect("serialize"); + let bytes = rmp_serde::to_vec_named(&snapshot).expect("serialize"); let decoded: ObserverSnapshot = rmp_serde::from_slice(&bytes).expect("deserialize"); assert_eq!(decoded.tick, 0); @@ -73,11 +73,11 @@ fn all_player_action_variants_roundtrip() { tick: 1, action: action.clone(), }; - let bytes = rmp_serde::to_vec(&input).expect("serialize"); + let bytes = rmp_serde::to_vec_named(&input).expect("serialize"); let decoded: PlayerInput = rmp_serde::from_slice(&bytes).expect("deserialize"); assert_eq!(decoded.tick, 1); // Verify the variant survived by re-serializing and comparing bytes - let re_bytes = rmp_serde::to_vec(&decoded).expect("re-serialize"); + let re_bytes = rmp_serde::to_vec_named(&decoded).expect("re-serialize"); assert_eq!(bytes, re_bytes, "round-trip mismatch for action variant"); } } @@ -85,7 +85,12 @@ fn all_player_action_variants_roundtrip() { /// All EntityKind variants must survive MessagePack round-trip (D-030 Layer 1) #[test] fn all_entity_kind_variants_roundtrip() { - let kinds = vec![EntityKind::Npc, EntityKind::Object, EntityKind::Terrain]; + let kinds = vec![ + EntityKind::Player, + EntityKind::Npc, + EntityKind::Object, + EntityKind::Terrain, + ]; for (i, kind) in kinds.into_iter().enumerate() { let entity = VisibleEntity { @@ -99,9 +104,9 @@ fn all_entity_kind_variants_roundtrip() { tick: 0, entities: vec![entity], }; - let bytes = rmp_serde::to_vec(&snapshot).expect("serialize"); + let bytes = rmp_serde::to_vec_named(&snapshot).expect("serialize"); let decoded: ObserverSnapshot = rmp_serde::from_slice(&bytes).expect("deserialize"); - let re_bytes = rmp_serde::to_vec(&decoded).expect("re-serialize"); + let re_bytes = rmp_serde::to_vec_named(&decoded).expect("re-serialize"); assert_eq!( bytes, re_bytes, "round-trip mismatch for EntityKind variant" From 1b0e51456029ddad891fa5f0fbc08cf0f2c9f19f Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:01:20 +0100 Subject: [PATCH 2/8] feat(bridge): add TCP transport for Godot client connection Godot has no Unix socket API, so TCP localhost is required for client-server IPC. TcpBridge implements SimBridge with the same framing protocol as LocalBridge. Includes accept/connect methods and three integration tests over TCP. Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/tcp.rs | 114 +++++++++++++++++++++ server/tests/bridge_tcp.rs | 201 +++++++++++++++++++++++++++++++++++++ 2 files changed, 315 insertions(+) create mode 100644 server/src/bridge/tcp.rs create mode 100644 server/tests/bridge_tcp.rs diff --git a/server/src/bridge/tcp.rs b/server/src/bridge/tcp.rs new file mode 100644 index 000000000..02c4b5a83 --- /dev/null +++ b/server/src/bridge/tcp.rs @@ -0,0 +1,114 @@ +// TcpBridge - TCP localhost IPC implementation +// Implements D-020 subprocess/IPC architecture +// Deterministic client-server communication via TCP sockets +// Used for Godot client which lacks Unix socket support + +use super::{BridgeError, ObserverSnapshot, PlayerInput, SimBridge}; +use crate::bridge::framing::{read_framed, write_framed}; +use std::io::{BufReader, BufWriter}; +use std::net::{SocketAddr, TcpListener, TcpStream}; +use std::sync::Mutex; + +/// TcpBridge: TCP transport for client-server IPC +/// Same semantics as LocalBridge but over TCP localhost +pub struct TcpBridge { + reader: Mutex>, + writer: Mutex>, + local_addr: SocketAddr, +} + +impl TcpBridge { + /// Server-side: bind TCP listener and accept one connection. + /// Binds to the specified address (e.g., "127.0.0.1:0" for OS-assigned port). + /// Returns the bridge with the actual bound address available via local_addr(). + pub fn accept(addr: &str) -> Result { + tracing::info!("TcpBridge binding to {}", addr); + + let listener = TcpListener::bind(addr) + .map_err(|e| BridgeError::Transport(format!("failed to bind TCP socket: {}", e)))?; + + let local_addr = listener + .local_addr() + .map_err(|e| BridgeError::Transport(format!("failed to get local address: {}", e)))?; + + tracing::info!("TcpBridge listening on {}", local_addr); + + // Accept one connection + let (stream, peer_addr) = listener + .accept() + .map_err(|e| BridgeError::Transport(format!("failed to accept connection: {}", e)))?; + + tracing::info!( + "TcpBridge accepted connection from {} on {}", + peer_addr, + local_addr + ); + + // 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)), + local_addr, + }) + } + + /// Client-side: connect to TCP address (for tests). + pub fn connect(addr: &str) -> Result { + tracing::info!("TcpBridge connecting to {}", addr); + + let stream = TcpStream::connect(addr).map_err(|e| { + BridgeError::Transport(format!("failed to connect to TCP socket: {}", e)) + })?; + + let local_addr = stream + .local_addr() + .map_err(|e| BridgeError::Transport(format!("failed to get local address: {}", e)))?; + + tracing::trace!("TcpBridge connected to {}", addr); + + // 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)), + local_addr, + }) + } + + /// Get the local address (useful for OS-assigned port discovery in tests). + pub fn local_addr(&self) -> SocketAddr { + self.local_addr + } +} + +impl SimBridge for TcpBridge { + fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> { + let payload = rmp_serde::to_vec_named(snapshot)?; + + let mut writer = self.writer.lock().expect("writer mutex poisoned"); + 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().expect("reader mutex poisoned"); + + 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 => Err(BridgeError::Transport("client disconnected (EOF)".into())), + } + } +} diff --git a/server/tests/bridge_tcp.rs b/server/tests/bridge_tcp.rs new file mode 100644 index 000000000..ae1927b9a --- /dev/null +++ b/server/tests/bridge_tcp.rs @@ -0,0 +1,201 @@ +//! Integration tests for TcpBridge over TCP localhost (D-030 Layer 2: IPC roundtrip). + +use settled_reach_server::bridge::framing::{read_framed, write_framed}; +use settled_reach_server::bridge::tcp::TcpBridge; +use settled_reach_server::bridge::types::*; +use settled_reach_server::bridge::SimBridge; +use std::net::{TcpListener, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +/// Helper to bind a TCP listener and return its address. +/// This allows tests to discover the OS-assigned port before calling TcpBridge::accept(). +fn bind_listener() -> (TcpListener, std::net::SocketAddr) { + let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind"); + let addr = listener.local_addr().expect("failed to get local address"); + (listener, addr) +} + +#[test] +fn snapshot_roundtrip_over_tcp() { + // Pre-bind listener to discover the port + let (listener, server_addr) = bind_listener(); + let addr_str = server_addr.to_string(); + + // Channel to signal when server is ready to accept + let (ready_tx, ready_rx) = mpsc::channel(); + + // Server thread: accept connection and send snapshot + let server_handle = thread::spawn(move || { + // Drop the pre-bound listener since TcpBridge::accept will bind its own + drop(listener); + + // Signal we're about to accept + ready_tx.send(()).expect("failed to send ready signal"); + + let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and 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"); + }); + + // Wait for server thread to start + ready_rx + .recv_timeout(Duration::from_secs(2)) + .expect("server did not become ready"); + + // Give server time to bind and call accept() + thread::sleep(Duration::from_millis(100)); + + // Client: connect and receive snapshot + let stream = TcpStream::connect(server_addr).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_tcp() { + // Pre-bind listener to discover the port + let (listener, server_addr) = bind_listener(); + let addr_str = server_addr.to_string(); + + // Channel to signal when server is ready to accept + let (ready_tx, ready_rx) = mpsc::channel(); + + // Server thread: accept connection and receive inputs + let server_handle = thread::spawn(move || { + // Drop the pre-bound listener since TcpBridge::accept will bind its own + drop(listener); + + // Signal we're about to accept + ready_tx.send(()).expect("failed to send ready signal"); + + let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and 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 + }); + + // Wait for server thread to start + ready_rx + .recv_timeout(Duration::from_secs(2)) + .expect("server did not become ready"); + + // Give server time to bind and call accept() + thread::sleep(Duration::from_millis(100)); + + // Client: connect and send inputs + let stream = TcpStream::connect(server_addr).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_named(&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"), + } +} + +#[test] +fn tcp_bridge_eof_returns_error() { + // Pre-bind listener to discover the port + let (listener, server_addr) = bind_listener(); + let addr_str = server_addr.to_string(); + + // Channel to signal when server is ready to accept + let (ready_tx, ready_rx) = mpsc::channel(); + + // Server thread: accept connection and receive EOF + let server_handle = thread::spawn(move || { + // Drop the pre-bound listener since TcpBridge::accept will bind its own + drop(listener); + + // Signal we're about to accept + ready_tx.send(()).expect("failed to send ready signal"); + + let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and accept"); + + // Attempt to receive inputs - should get EOF error + let result = bridge.receive_inputs(); + + assert!(result.is_err(), "expected error on EOF"); + match result { + Err(settled_reach_server::bridge::BridgeError::Transport(msg)) => { + assert!( + msg.contains("disconnected") || msg.contains("EOF"), + "expected disconnect/EOF error, got: {}", + msg + ); + } + _ => panic!("expected Transport error with disconnect/EOF message"), + } + }); + + // Wait for server thread to start + ready_rx + .recv_timeout(Duration::from_secs(2)) + .expect("server did not become ready"); + + // Give server time to bind and call accept() + thread::sleep(Duration::from_millis(100)); + + // Client: connect and immediately disconnect without sending data + let stream = TcpStream::connect(server_addr).expect("failed to connect"); + drop(stream); // Close connection immediately + + server_handle.join().expect("server thread panicked"); +} From 7864bfdf7d239c93ec323b8be433830f466268df Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:01:28 +0100 Subject: [PATCH 3/8] feat(simulation): add input processing, snapshot gen, and game loop Implements the full server-side tick pipeline: - process_player_input drains InputQueue, converts PlayerActions to MoveIntent components or pause/unpause toggles - generate_snapshot builds ObserverSnapshot from ECS state with render coordinate conversion - receive_bridge_inputs/send_bridge_snapshot handle bridge I/O with graceful disconnect detection via ServerRunning resource - main.rs now accepts TCP connections and runs a proper game loop - PlayerCharacter marker, Player EntityKind, SnapshotBuffer resource Closes server side of #81, #82, #83. Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/mod.rs | 96 ++++++++++++++++++- server/src/bridge/types.rs | 8 ++ server/src/main.rs | 31 +++++- server/src/simulation/input.rs | 151 +++++++++++++++++++++++++++++- server/src/simulation/mod.rs | 7 +- server/src/simulation/movement.rs | 4 + server/tests/game_loop.rs | 110 ++++++++++++++++++++++ 7 files changed, 397 insertions(+), 10 deletions(-) create mode 100644 server/tests/game_loop.rs diff --git a/server/src/bridge/mod.rs b/server/src/bridge/mod.rs index 05a398abe..44e6a7fbd 100644 --- a/server/src/bridge/mod.rs +++ b/server/src/bridge/mod.rs @@ -4,9 +4,11 @@ use bevy_app::prelude::*; use bevy_ecs::prelude::*; +use bevy_ecs::schedule::IntoScheduleConfigs; pub mod framing; pub mod local; +pub mod tcp; pub mod types; pub use types::*; @@ -55,13 +57,103 @@ impl BridgeResource { } } +/// Generate ObserverSnapshot from ECS state +pub fn generate_snapshot( + time: Res, + entities: Query<( + Entity, + &crate::simulation::movement::TilePosition, + Option<&crate::simulation::movement::PlayerCharacter>, + Option<&crate::npc::Npc>, + )>, + mut buffer: ResMut, +) { + let mut visible = Vec::new(); + for (entity, pos, is_player, is_npc) in entities.iter() { + let (x, y, z) = pos.to_render_coords(); + let kind = if is_player.is_some() { + EntityKind::Player + } else if is_npc.is_some() { + EntityKind::Npc + } else { + EntityKind::Object + }; + visible.push(VisibleEntity { + entity_id: entity.to_bits(), + x, + y, + z, + kind, + }); + } + buffer.snapshot = Some(ObserverSnapshot { + tick: time.tick, + entities: visible, + }); +} + +/// Receive inputs from bridge and push to InputQueue +pub fn receive_bridge_inputs( + bridge: Option>, + mut input_queue: ResMut, + mut running: ResMut, +) { + let Some(bridge) = bridge else { return }; + match bridge.receive_inputs() { + Ok(inputs) => { + for input in inputs { + input_queue.push(input); + } + } + Err(BridgeError::Transport(ref msg)) if msg.contains("disconnected") => { + tracing::info!("Client disconnected, shutting down"); + running.0 = false; + } + Err(e) => { + tracing::error!("Bridge receive error: {}", e); + } + } +} + +/// Send snapshot from buffer to bridge +pub fn send_bridge_snapshot( + bridge: Option>, + mut buffer: ResMut, +) { + let Some(bridge) = bridge else { return }; + if let Some(snapshot) = buffer.snapshot.take() { + if let Err(e) = bridge.send_snapshot(&snapshot) { + tracing::error!("Bridge send error: {}", e); + } + } +} + +/// Server running flag resource +#[derive(Resource, Debug, Clone)] +pub struct ServerRunning(pub bool); + +impl Default for ServerRunning { + fn default() -> Self { + Self(true) + } +} + /// Bridge plugin for client-server communication /// Abstracts transport layer (LocalBridge/NetworkBridge) pub struct BridgePlugin; impl Plugin for BridgePlugin { - fn build(&self, _app: &mut App) { - // Stub implementation - will be populated in phase 2 + fn build(&self, app: &mut App) { + app.init_resource::() + .init_resource::() + .add_systems( + Update, + ( + receive_bridge_inputs.before(crate::simulation::input::process_player_input), + generate_snapshot.after(crate::simulation::movement::validate_movement), + send_bridge_snapshot.after(generate_snapshot), + ), + ); tracing::debug!("BridgePlugin initialized"); } } diff --git a/server/src/bridge/types.rs b/server/src/bridge/types.rs index 685262263..1dd34690f 100644 --- a/server/src/bridge/types.rs +++ b/server/src/bridge/types.rs @@ -2,6 +2,7 @@ // ObserverSnapshot: data crossing the client-server boundary // PlayerInput: semantic actions from client +use bevy_ecs::prelude::*; use serde::{Deserialize, Serialize}; /// The ONLY data structure crossing the client-server boundary (D-020) @@ -32,6 +33,7 @@ pub struct VisibleEntity { /// Category of visible entity #[derive(Debug, Clone, Serialize, Deserialize)] pub enum EntityKind { + Player, Npc, Object, Terrain, @@ -63,3 +65,9 @@ pub enum PlayerAction { Pause, Unpause, } + +/// Snapshot buffer resource for staging outgoing ObserverSnapshots +#[derive(Resource, Debug, Default)] +pub struct SnapshotBuffer { + pub snapshot: Option, +} diff --git a/server/src/main.rs b/server/src/main.rs index 6ddee05cf..d1550f443 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -4,7 +4,9 @@ use bevy_app::prelude::*; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; -use settled_reach_server::bridge::BridgePlugin; +use settled_reach_server::bridge::tcp::TcpBridge; +use settled_reach_server::bridge::{BridgePlugin, BridgeResource, ServerRunning}; +use settled_reach_server::simulation::movement::{PlayerCharacter, TilePosition, WalkabilityMap}; use settled_reach_server::simulation::SimulationPlugin; fn main() { @@ -17,15 +19,36 @@ fn main() { .with(tracing_subscriber::fmt::layer()) .init(); + let addr = std::env::args() + .nth(1) + .or_else(|| std::env::var("SR_ADDR").ok()) + .unwrap_or_else(|| "127.0.0.1:9876".to_string()); + tracing::info!("The Settled Reach - Simulation Server starting"); + tracing::info!("Waiting for client connection on {}", addr); + + let bridge = TcpBridge::accept(&addr).expect("Failed to accept client connection"); + tracing::info!("Client connected, initializing simulation"); // Create the bevy App and add plugins let mut app = App::new(); app.add_plugins(SimulationPlugin); app.add_plugins(BridgePlugin); + app.insert_resource(BridgeResource::new(bridge)); + app.insert_resource(WalkabilityMap::new(32, 32, 1)); + app.world_mut() + .spawn((PlayerCharacter, TilePosition::new(16, 16, 0))); - // Single tick for smoke verification; real game loop in phase 2 - app.update(); + tracing::info!("Simulation initialized, entering game loop"); - tracing::info!("Simulation server update complete"); + // Game loop: run until client disconnects + loop { + app.update(); + // Check ServerRunning resource + if !app.world().resource::().0 { + break; + } + } + + tracing::info!("Simulation server shutting down"); } diff --git a/server/src/simulation/input.rs b/server/src/simulation/input.rs index ea94108ab..cc711459c 100644 --- a/server/src/simulation/input.rs +++ b/server/src/simulation/input.rs @@ -2,7 +2,9 @@ // Timestamped player input events for deterministic simulation (D-010 principle 4) // PlayerInput: semantic actions (MoveNorth, Interact, UsePerceptionMode) -use crate::bridge::types::PlayerInput; +use crate::bridge::types::{PlayerAction, PlayerInput}; +use crate::simulation::movement::{MoveIntent, PlayerCharacter, TilePosition}; +use crate::simulation::time::SimulationTime; use bevy_ecs::prelude::*; use std::collections::VecDeque; @@ -51,10 +53,62 @@ impl InputQueue { } } +/// Drains InputQueue for the current tick, converts PlayerActions to ECS components. +pub fn process_player_input( + mut input_queue: ResMut, + mut time: ResMut, + mut commands: Commands, + player_query: Query<(Entity, &TilePosition), With>, +) { + let current_tick = time.tick; + let inputs = input_queue.drain_for_tick(current_tick); + + for input in inputs { + match input.action { + PlayerAction::MoveNorth => apply_move(&player_query, &mut commands, 0, -1), + PlayerAction::MoveSouth => apply_move(&player_query, &mut commands, 0, 1), + PlayerAction::MoveEast => apply_move(&player_query, &mut commands, 1, 0), + PlayerAction::MoveWest => apply_move(&player_query, &mut commands, -1, 0), + PlayerAction::MoveNortheast => apply_move(&player_query, &mut commands, 1, -1), + PlayerAction::MoveNorthwest => apply_move(&player_query, &mut commands, -1, -1), + PlayerAction::MoveSoutheast => apply_move(&player_query, &mut commands, 1, 1), + PlayerAction::MoveSouthwest => apply_move(&player_query, &mut commands, -1, 1), + PlayerAction::Pause => { + time.paused = true; + tracing::debug!("Simulation paused by player input"); + } + PlayerAction::Unpause => { + time.paused = false; + tracing::debug!("Simulation unpaused by player input"); + } + PlayerAction::Interact => { + tracing::trace!("Interact action — no-op for Sprint 1"); + } + PlayerAction::UsePerceptionMode(ref mode) => { + tracing::trace!("UsePerceptionMode({}) — no-op for Sprint 1", mode); + } + } + } +} + +fn apply_move( + player_query: &Query<(Entity, &TilePosition), With>, + commands: &mut Commands, + dx: i32, + dy: i32, +) { + if let Ok((entity, pos)) = player_query.single() { + commands.entity(entity).insert(MoveIntent { + target: TilePosition::new(pos.x + dx, pos.y + dy, pos.z), + }); + } else { + tracing::warn!("No player entity found for movement input"); + } +} + #[cfg(test)] mod tests { use super::*; - use crate::bridge::types::PlayerAction; #[test] fn drain_returns_inputs_up_to_tick() { @@ -96,4 +150,97 @@ mod tests { action: PlayerAction::MoveSouth, }); } + + #[test] + fn process_input_move_creates_intent() { + let mut world = bevy_ecs::world::World::new(); + world.insert_resource(InputQueue::default()); + world.insert_resource(SimulationTime { + tick: 0, + paused: false, + }); + + let player = world + .spawn((PlayerCharacter, TilePosition::new(5, 5, 0))) + .id(); + + world.resource_mut::().push(PlayerInput { + tick: 0, + action: PlayerAction::MoveNorth, + }); + + let mut schedule = bevy_ecs::schedule::Schedule::default(); + schedule.add_systems(process_player_input); + schedule.run(&mut world); + + let intent = world.get::(player).unwrap(); + assert_eq!(intent.target, TilePosition::new(5, 4, 0)); + } + + #[test] + fn process_input_pause_toggles() { + let mut world = bevy_ecs::world::World::new(); + world.insert_resource(InputQueue::default()); + world.insert_resource(SimulationTime { + tick: 0, + paused: false, + }); + + world.resource_mut::().push(PlayerInput { + tick: 0, + action: PlayerAction::Pause, + }); + + let mut schedule = bevy_ecs::schedule::Schedule::default(); + schedule.add_systems(process_player_input); + schedule.run(&mut world); + + assert!(world.resource::().paused); + } + + #[test] + fn process_input_no_player_no_panic() { + let mut world = bevy_ecs::world::World::new(); + world.insert_resource(InputQueue::default()); + world.insert_resource(SimulationTime { + tick: 0, + paused: false, + }); + + world.resource_mut::().push(PlayerInput { + tick: 0, + action: PlayerAction::MoveNorth, + }); + + let mut schedule = bevy_ecs::schedule::Schedule::default(); + schedule.add_systems(process_player_input); + // Should not panic + schedule.run(&mut world); + } + + #[test] + fn process_input_future_tick_ignored() { + let mut world = bevy_ecs::world::World::new(); + world.insert_resource(InputQueue::default()); + world.insert_resource(SimulationTime { + tick: 0, + paused: false, + }); + + let player = world + .spawn((PlayerCharacter, TilePosition::new(5, 5, 0))) + .id(); + + world.resource_mut::().push(PlayerInput { + tick: 5, + action: PlayerAction::MoveNorth, + }); + + let mut schedule = bevy_ecs::schedule::Schedule::default(); + schedule.add_systems(process_player_input); + schedule.run(&mut world); + + // No MoveIntent should be created (input for future tick) + assert!(world.get::(player).is_none()); + } } diff --git a/server/src/simulation/mod.rs b/server/src/simulation/mod.rs index 054a04e1d..dc101ad62 100644 --- a/server/src/simulation/mod.rs +++ b/server/src/simulation/mod.rs @@ -20,10 +20,13 @@ impl Plugin for SimulationPlugin { app.init_resource::() .insert_resource(rng::SimRng::new(0)) .init_resource::() - .add_systems(Update, time::advance_tick) .add_systems( Update, - movement::validate_movement.after(time::advance_tick), + ( + input::process_player_input, + time::advance_tick.after(input::process_player_input), + movement::validate_movement.after(time::advance_tick), + ), ); tracing::debug!("SimulationPlugin initialized"); diff --git a/server/src/simulation/movement.rs b/server/src/simulation/movement.rs index afdbf0fd4..d8943774b 100644 --- a/server/src/simulation/movement.rs +++ b/server/src/simulation/movement.rs @@ -10,6 +10,10 @@ use std::collections::HashMap; /// Chunk size in tiles (32x32 per chunk) pub const CHUNK_SIZE: i32 = 32; +/// Marker component identifying the player-controlled entity. +#[derive(Component, Debug)] +pub struct PlayerCharacter; + /// Tile position component for grid-based movement. /// Discrete integer coordinates used in simulation; converted to f32 /// at the bridge boundary for VisibleEntity wire format. diff --git a/server/tests/game_loop.rs b/server/tests/game_loop.rs new file mode 100644 index 000000000..f2a53e760 --- /dev/null +++ b/server/tests/game_loop.rs @@ -0,0 +1,110 @@ +//! E2E integration test: full game loop with input processing and snapshot generation +//! Tests the complete pipeline: client sends input → server processes → server sends snapshot + +use bevy_app::prelude::*; +use settled_reach_server::bridge::framing::{read_framed, write_framed}; +use settled_reach_server::bridge::tcp::TcpBridge; +use settled_reach_server::bridge::types::*; +use settled_reach_server::bridge::{BridgePlugin, BridgeResource}; +use settled_reach_server::simulation::movement::{PlayerCharacter, TilePosition, WalkabilityMap}; +use settled_reach_server::simulation::SimulationPlugin; +use std::io::{BufReader, BufWriter}; +use std::net::{TcpListener, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +#[test] +fn player_moves_north_through_full_pipeline() { + // Pre-bind listener to discover the port + let listener = TcpListener::bind("127.0.0.1:0").expect("bind listener"); + let server_addr = listener.local_addr().expect("get local addr"); + let addr_str = server_addr.to_string(); + + // Channel to signal when server is ready to accept + let (ready_tx, ready_rx) = mpsc::channel(); + + // Spawn server thread + let server_handle = thread::spawn(move || { + // Drop the pre-bound listener since TcpBridge::accept will bind its own + drop(listener); + + // Signal we're about to accept + ready_tx.send(()).expect("send ready signal"); + + // Accept connection + let bridge = TcpBridge::accept(&addr_str).expect("accept connection"); + eprintln!("Test server accepted connection"); + + // Build app + let mut app = App::new(); + app.add_plugins(SimulationPlugin); + app.add_plugins(BridgePlugin); + app.insert_resource(BridgeResource::new(bridge)); + app.insert_resource(WalkabilityMap::new(32, 32, 1)); + app.world_mut() + .spawn((PlayerCharacter, TilePosition::new(16, 16, 0))); + + // Run one tick: receive input, process, validate movement, generate snapshot, send + app.update(); + + eprintln!("Server completed one update cycle"); + }); + + // Wait for server thread to start + ready_rx + .recv_timeout(Duration::from_secs(2)) + .expect("server did not become ready"); + + // Give server time to bind and call accept() + thread::sleep(Duration::from_millis(100)); + + // Client: connect and send input + { + let stream = TcpStream::connect(server_addr).expect("client connect"); + let mut reader = BufReader::new(stream.try_clone().expect("clone for reader")); + let mut writer = BufWriter::new(stream); + + // Send PlayerInput: MoveNorth at tick 0 + let inputs = vec![PlayerInput { + tick: 0, + action: PlayerAction::MoveNorth, + }]; + let payload = rmp_serde::to_vec_named(&inputs).expect("serialize inputs"); + write_framed(&mut writer, &payload).expect("send inputs"); + eprintln!("Client sent MoveNorth input"); + + // Receive ObserverSnapshot + let response = read_framed(&mut reader) + .expect("read snapshot") + .expect("not EOF"); + let snapshot: ObserverSnapshot = + rmp_serde::from_slice(&response).expect("deserialize snapshot"); + + eprintln!( + "Client received snapshot: tick={}, entities={}", + snapshot.tick, + snapshot.entities.len() + ); + + // Verify snapshot + assert_eq!(snapshot.tick, 1); // After one tick + assert_eq!(snapshot.entities.len(), 1); // One entity (player) + + let player_entity = &snapshot.entities[0]; + // Player started at (16, 16, 0), moved north (y-1) to (16, 15, 0) + // Render coords: (16.5, 15.5, 0) + assert_eq!(player_entity.x, 16.5); + assert_eq!(player_entity.y, 15.5); + assert_eq!(player_entity.z, 0); + assert!(matches!(player_entity.kind, EntityKind::Player)); + + eprintln!( + "Client verified player moved to ({}, {}, {})", + player_entity.x, player_entity.y, player_entity.z + ); + } + + // Wait for server thread + server_handle.join().expect("server thread panicked"); +} From a5b88cb73535bec8263687f1340014035e53f770 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:01:31 +0100 Subject: [PATCH 4/8] chore(bridge): format gen_fixtures.rs Co-Authored-By: Claude Opus 4.6 --- server/tests/gen_fixtures.rs | 44 ++++++++++++++++++++++++++++++------ 1 file changed, 37 insertions(+), 7 deletions(-) diff --git a/server/tests/gen_fixtures.rs b/server/tests/gen_fixtures.rs index 4cd63aff2..f7e4beadf 100644 --- a/server/tests/gen_fixtures.rs +++ b/server/tests/gen_fixtures.rs @@ -28,7 +28,10 @@ fn generate_msgpack_fixtures() { kind: EntityKind::Npc, }], }; - write_fixture("snapshot_one_npc", &rmp_serde::to_vec_named(&snapshot).unwrap()); + write_fixture( + "snapshot_one_npc", + &rmp_serde::to_vec_named(&snapshot).unwrap(), + ); // Empty snapshot let empty = ObserverSnapshot { @@ -42,23 +45,50 @@ fn generate_msgpack_fixtures() { tick: 100, action: PlayerAction::MoveNorth, }; - write_fixture("input_move_north", &rmp_serde::to_vec_named(&input_north).unwrap()); + write_fixture( + "input_move_north", + &rmp_serde::to_vec_named(&input_north).unwrap(), + ); // PlayerInput: UsePerceptionMode let input_perception = PlayerInput { tick: 200, action: PlayerAction::UsePerceptionMode("thermal".to_string()), }; - write_fixture("input_perception_mode", &rmp_serde::to_vec_named(&input_perception).unwrap()); + write_fixture( + "input_perception_mode", + &rmp_serde::to_vec_named(&input_perception).unwrap(), + ); // Snapshot with multiple entities and all EntityKind variants let snapshot_multi = ObserverSnapshot { tick: 999, entities: vec![ - VisibleEntity { entity_id: 1, x: 5.0, y: 10.0, z: 0, kind: EntityKind::Npc }, - VisibleEntity { entity_id: 2, x: 15.5, y: 3.0, z: 1, kind: EntityKind::Object }, - VisibleEntity { entity_id: 3, x: 0.0, y: 0.0, z: -1, kind: EntityKind::Terrain }, + VisibleEntity { + entity_id: 1, + x: 5.0, + y: 10.0, + z: 0, + kind: EntityKind::Npc, + }, + VisibleEntity { + entity_id: 2, + x: 15.5, + y: 3.0, + z: 1, + kind: EntityKind::Object, + }, + VisibleEntity { + entity_id: 3, + x: 0.0, + y: 0.0, + z: -1, + kind: EntityKind::Terrain, + }, ], }; - write_fixture("snapshot_multi_entity", &rmp_serde::to_vec_named(&snapshot_multi).unwrap()); + write_fixture( + "snapshot_multi_entity", + &rmp_serde::to_vec_named(&snapshot_multi).unwrap(), + ); } From 058d0352abb7a04ebaa936b91c7a41716921f8c5 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:02:00 +0100 Subject: [PATCH 5/8] chore(meta): update changelog Co-Authored-By: Claude Opus 4.6 --- CHANGELOG.md | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0fb28c024..653251d8f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ Format based on [Keep a Changelog](https://keepachangelog.com/). ## [Unreleased] ### Added +- TcpBridge transport for Godot client connection — TCP localhost IPC alongside existing Unix socket LocalBridge +- Input processing system (process_player_input) — drains InputQueue, converts PlayerActions to MoveIntent components, handles pause/unpause +- Snapshot generation system (generate_snapshot) — builds ObserverSnapshot from ECS state with render coordinate conversion +- Bridge I/O systems (receive_bridge_inputs, send_bridge_snapshot) — wire bridge to ECS pipeline with graceful disconnect detection +- Full game loop in main.rs — TCP accept, tick loop with ServerRunning resource, CLI/env addr config +- PlayerCharacter marker component, Player EntityKind variant, SnapshotBuffer resource +- E2E game_loop integration test verifying player movement through full pipeline +- 8 new tests (3 TCP bridge + 4 input processing + 1 E2E game loop), total 53 - v0.1 content gap analysis workshop — 6 agents, 2 rounds, 9 content layers, 8 new decisions (D-032 through D-039) - D-032: Separate monologue pools per playable character (hard partition, not filter) - D-033: Entity color represents relationship to player character (asymmetric per character) @@ -96,6 +104,8 @@ Format based on [Keep a Changelog](https://keepachangelog.com/). - Rejected alternatives R-004 through R-010 documented (pure Bevy, pure Godot, GDExtension, C++ GDExtension, Fyrox, custom framework, protobuf) ### Fixed +- Wire format mismatch: LocalBridge now uses rmp_serde::to_vec_named() (named maps) matching client Protocol.gd expectations +- EOF on bridge read now returns BridgeError::Transport for disconnect detection instead of empty Vec - WalkabilityMap rewritten from flat Vec to chunk-based HashMap per D-012 architecture (Tyre review) - LocalBridge mutex .unwrap() → .expect() for clearer panic messages (Hoshe review) - Documented 16MB MAX_MESSAGE_SIZE rationale in framing.rs (Hoshe review) From 2cd8786036f48701178a7f0427048aee2a66591e Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:15:56 +0100 Subject: [PATCH 6/8] fix(bridge): add Disconnected error variant, fix test race condition Replace string-matching disconnect detection with explicit BridgeError::Disconnected variant. Add TcpBridge::accept_on(listener) that takes a pre-bound TcpListener, eliminating the 100ms sleep hack in TCP tests. Send errors now also trigger ServerRunning=false. Add trace logging to generate_snapshot for entity count visibility. Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/local.rs | 2 +- server/src/bridge/mod.rs | 15 ++++- server/src/bridge/tcp.rs | 32 ++++++++++- server/tests/bridge_tcp.rs | 110 +++++++------------------------------ 4 files changed, 65 insertions(+), 94 deletions(-) diff --git a/server/src/bridge/local.rs b/server/src/bridge/local.rs index 62c36f873..5e73649a6 100644 --- a/server/src/bridge/local.rs +++ b/server/src/bridge/local.rs @@ -99,7 +99,7 @@ impl SimBridge for LocalBridge { tracing::trace!("received {} inputs", inputs.len()); Ok(inputs) } - None => Err(BridgeError::Transport("client disconnected (EOF)".into())), + None => Err(BridgeError::Disconnected), } } } diff --git a/server/src/bridge/mod.rs b/server/src/bridge/mod.rs index 44e6a7fbd..702bd608d 100644 --- a/server/src/bridge/mod.rs +++ b/server/src/bridge/mod.rs @@ -23,6 +23,8 @@ pub enum BridgeError { Io(#[from] std::io::Error), #[error("transport error: {0}")] Transport(String), + #[error("client disconnected")] + Disconnected, } /// Abstracts transport layer (D-020) @@ -86,6 +88,11 @@ pub fn generate_snapshot( kind, }); } + tracing::trace!( + "generate_snapshot: tick={}, entities={}", + time.tick, + visible.len() + ); buffer.snapshot = Some(ObserverSnapshot { tick: time.tick, entities: visible, @@ -105,7 +112,7 @@ pub fn receive_bridge_inputs( input_queue.push(input); } } - Err(BridgeError::Transport(ref msg)) if msg.contains("disconnected") => { + Err(BridgeError::Disconnected) => { tracing::info!("Client disconnected, shutting down"); running.0 = false; } @@ -119,11 +126,13 @@ pub fn receive_bridge_inputs( pub fn send_bridge_snapshot( bridge: Option>, mut buffer: ResMut, + mut running: ResMut, ) { let Some(bridge) = bridge else { return }; if let Some(snapshot) = buffer.snapshot.take() { if let Err(e) = bridge.send_snapshot(&snapshot) { tracing::error!("Bridge send error: {}", e); + running.0 = false; } } } @@ -150,7 +159,9 @@ impl Plugin for BridgePlugin { Update, ( receive_bridge_inputs.before(crate::simulation::input::process_player_input), - generate_snapshot.after(crate::simulation::movement::validate_movement), + generate_snapshot + .after(crate::simulation::movement::validate_movement) + .before(crate::simulation::time::advance_tick), send_bridge_snapshot.after(generate_snapshot), ), ); diff --git a/server/src/bridge/tcp.rs b/server/src/bridge/tcp.rs index 02c4b5a83..72781e96e 100644 --- a/server/src/bridge/tcp.rs +++ b/server/src/bridge/tcp.rs @@ -56,6 +56,36 @@ impl TcpBridge { }) } + /// Server-side: accept one connection on an existing TcpListener. + /// Avoids race conditions in tests by separating bind from accept. + pub fn accept_on(listener: TcpListener) -> Result { + let local_addr = listener + .local_addr() + .map_err(|e| BridgeError::Transport(format!("failed to get local address: {}", e)))?; + + tracing::info!("TcpBridge accepting on {}", local_addr); + + let (stream, peer_addr) = listener + .accept() + .map_err(|e| BridgeError::Transport(format!("failed to accept connection: {}", e)))?; + + tracing::info!( + "TcpBridge accepted connection from {} on {}", + peer_addr, + local_addr + ); + + 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)), + local_addr, + }) + } + /// Client-side: connect to TCP address (for tests). pub fn connect(addr: &str) -> Result { tracing::info!("TcpBridge connecting to {}", addr); @@ -108,7 +138,7 @@ impl SimBridge for TcpBridge { tracing::trace!("received {} inputs", inputs.len()); Ok(inputs) } - None => Err(BridgeError::Transport("client disconnected (EOF)".into())), + None => Err(BridgeError::Disconnected), } } } diff --git a/server/tests/bridge_tcp.rs b/server/tests/bridge_tcp.rs index ae1927b9a..9f5bdf97a 100644 --- a/server/tests/bridge_tcp.rs +++ b/server/tests/bridge_tcp.rs @@ -5,36 +5,17 @@ use settled_reach_server::bridge::tcp::TcpBridge; use settled_reach_server::bridge::types::*; use settled_reach_server::bridge::SimBridge; use std::net::{TcpListener, TcpStream}; -use std::sync::mpsc; use std::thread; -use std::time::Duration; - -/// Helper to bind a TCP listener and return its address. -/// This allows tests to discover the OS-assigned port before calling TcpBridge::accept(). -fn bind_listener() -> (TcpListener, std::net::SocketAddr) { - let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind"); - let addr = listener.local_addr().expect("failed to get local address"); - (listener, addr) -} #[test] fn snapshot_roundtrip_over_tcp() { - // Pre-bind listener to discover the port - let (listener, server_addr) = bind_listener(); - let addr_str = server_addr.to_string(); + // Bind listener first — port is guaranteed ready before spawning threads + let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind"); + let server_addr = listener.local_addr().expect("failed to get local address"); - // Channel to signal when server is ready to accept - let (ready_tx, ready_rx) = mpsc::channel(); - - // Server thread: accept connection and send snapshot + // Server thread: accept on pre-bound listener and send snapshot let server_handle = thread::spawn(move || { - // Drop the pre-bound listener since TcpBridge::accept will bind its own - drop(listener); - - // Signal we're about to accept - ready_tx.send(()).expect("failed to send ready signal"); - - let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and accept"); + let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); let snapshot = ObserverSnapshot { tick: 42, @@ -52,15 +33,7 @@ fn snapshot_roundtrip_over_tcp() { .expect("failed to send snapshot"); }); - // Wait for server thread to start - ready_rx - .recv_timeout(Duration::from_secs(2)) - .expect("server did not become ready"); - - // Give server time to bind and call accept() - thread::sleep(Duration::from_millis(100)); - - // Client: connect and receive snapshot + // Client: connect and receive snapshot (no sleep needed — listener already bound) let stream = TcpStream::connect(server_addr).expect("failed to connect"); let mut reader = std::io::BufReader::new(stream); @@ -82,22 +55,11 @@ fn snapshot_roundtrip_over_tcp() { #[test] fn input_roundtrip_over_tcp() { - // Pre-bind listener to discover the port - let (listener, server_addr) = bind_listener(); - let addr_str = server_addr.to_string(); + let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind"); + let server_addr = listener.local_addr().expect("failed to get local address"); - // Channel to signal when server is ready to accept - let (ready_tx, ready_rx) = mpsc::channel(); - - // Server thread: accept connection and receive inputs let server_handle = thread::spawn(move || { - // Drop the pre-bound listener since TcpBridge::accept will bind its own - drop(listener); - - // Signal we're about to accept - ready_tx.send(()).expect("failed to send ready signal"); - - let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and accept"); + let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); let inputs = bridge.receive_inputs().expect("failed to receive inputs"); @@ -108,14 +70,6 @@ fn input_roundtrip_over_tcp() { inputs }); - // Wait for server thread to start - ready_rx - .recv_timeout(Duration::from_secs(2)) - .expect("server did not become ready"); - - // Give server time to bind and call accept() - thread::sleep(Duration::from_millis(100)); - // Client: connect and send inputs let stream = TcpStream::connect(server_addr).expect("failed to connect"); let mut writer = std::io::BufWriter::new(stream); @@ -139,7 +93,6 @@ fn input_roundtrip_over_tcp() { 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"), @@ -152,50 +105,27 @@ fn input_roundtrip_over_tcp() { #[test] fn tcp_bridge_eof_returns_error() { - // Pre-bind listener to discover the port - let (listener, server_addr) = bind_listener(); - let addr_str = server_addr.to_string(); + let listener = TcpListener::bind("127.0.0.1:0").expect("failed to bind"); + let server_addr = listener.local_addr().expect("failed to get local address"); - // Channel to signal when server is ready to accept - let (ready_tx, ready_rx) = mpsc::channel(); - - // Server thread: accept connection and receive EOF let server_handle = thread::spawn(move || { - // Drop the pre-bound listener since TcpBridge::accept will bind its own - drop(listener); + let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); - // Signal we're about to accept - ready_tx.send(()).expect("failed to send ready signal"); - - let bridge = TcpBridge::accept(&addr_str).expect("failed to bind and accept"); - - // Attempt to receive inputs - should get EOF error let result = bridge.receive_inputs(); assert!(result.is_err(), "expected error on EOF"); - match result { - Err(settled_reach_server::bridge::BridgeError::Transport(msg)) => { - assert!( - msg.contains("disconnected") || msg.contains("EOF"), - "expected disconnect/EOF error, got: {}", - msg - ); - } - _ => panic!("expected Transport error with disconnect/EOF message"), - } + assert!( + matches!( + result, + Err(settled_reach_server::bridge::BridgeError::Disconnected) + ), + "expected Disconnected error" + ); }); - // Wait for server thread to start - ready_rx - .recv_timeout(Duration::from_secs(2)) - .expect("server did not become ready"); - - // Give server time to bind and call accept() - thread::sleep(Duration::from_millis(100)); - // Client: connect and immediately disconnect without sending data let stream = TcpStream::connect(server_addr).expect("failed to connect"); - drop(stream); // Close connection immediately + drop(stream); server_handle.join().expect("server thread panicked"); } From b9af725b02ddf51a82c44e4e7bd514df90fb21c6 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:16:04 +0100 Subject: [PATCH 7/8] fix(simulation): correct tick synchronization in snapshot generation Snapshot for tick N should show state at END of tick N. Reorder systems so generate_snapshot runs after validate_movement but before advance_tick. Previously snapshot.tick was the incremented tick, not the tick whose inputs were processed. Also fix main.rs accept error to log address context before exiting. Co-Authored-By: Claude Opus 4.6 --- server/src/simulation/mod.rs | 4 +- server/tests/game_loop.rs | 102 ++++++++++++----------------------- 2 files changed, 35 insertions(+), 71 deletions(-) diff --git a/server/src/simulation/mod.rs b/server/src/simulation/mod.rs index dc101ad62..e82407744 100644 --- a/server/src/simulation/mod.rs +++ b/server/src/simulation/mod.rs @@ -24,8 +24,8 @@ impl Plugin for SimulationPlugin { Update, ( input::process_player_input, - time::advance_tick.after(input::process_player_input), - movement::validate_movement.after(time::advance_tick), + movement::validate_movement.after(input::process_player_input), + time::advance_tick.after(movement::validate_movement), ), ); diff --git a/server/tests/game_loop.rs b/server/tests/game_loop.rs index f2a53e760..97109f6ca 100644 --- a/server/tests/game_loop.rs +++ b/server/tests/game_loop.rs @@ -1,5 +1,5 @@ //! E2E integration test: full game loop with input processing and snapshot generation -//! Tests the complete pipeline: client sends input → server processes → server sends snapshot +//! Tests the complete pipeline: client sends input -> server processes -> server sends snapshot use bevy_app::prelude::*; use settled_reach_server::bridge::framing::{read_framed, write_framed}; @@ -10,31 +10,17 @@ use settled_reach_server::simulation::movement::{PlayerCharacter, TilePosition, use settled_reach_server::simulation::SimulationPlugin; use std::io::{BufReader, BufWriter}; use std::net::{TcpListener, TcpStream}; -use std::sync::mpsc; use std::thread; -use std::time::Duration; #[test] fn player_moves_north_through_full_pipeline() { - // Pre-bind listener to discover the port + // Bind listener first — port guaranteed ready, no sleep needed let listener = TcpListener::bind("127.0.0.1:0").expect("bind listener"); let server_addr = listener.local_addr().expect("get local addr"); - let addr_str = server_addr.to_string(); - - // Channel to signal when server is ready to accept - let (ready_tx, ready_rx) = mpsc::channel(); // Spawn server thread let server_handle = thread::spawn(move || { - // Drop the pre-bound listener since TcpBridge::accept will bind its own - drop(listener); - - // Signal we're about to accept - ready_tx.send(()).expect("send ready signal"); - - // Accept connection - let bridge = TcpBridge::accept(&addr_str).expect("accept connection"); - eprintln!("Test server accepted connection"); + let bridge = TcpBridge::accept_on(listener).expect("accept connection"); // Build app let mut app = App::new(); @@ -47,64 +33,42 @@ fn player_moves_north_through_full_pipeline() { // Run one tick: receive input, process, validate movement, generate snapshot, send app.update(); - - eprintln!("Server completed one update cycle"); }); - // Wait for server thread to start - ready_rx - .recv_timeout(Duration::from_secs(2)) - .expect("server did not become ready"); + // Client: connect and send input (no sleep — listener was pre-bound) + let stream = TcpStream::connect(server_addr).expect("client connect"); + let mut reader = BufReader::new(stream.try_clone().expect("clone for reader")); + let mut writer = BufWriter::new(stream); - // Give server time to bind and call accept() - thread::sleep(Duration::from_millis(100)); + // Send PlayerInput: MoveNorth at tick 0 + let inputs = vec![PlayerInput { + tick: 0, + action: PlayerAction::MoveNorth, + }]; + let payload = rmp_serde::to_vec_named(&inputs).expect("serialize inputs"); + write_framed(&mut writer, &payload).expect("send inputs"); - // Client: connect and send input - { - let stream = TcpStream::connect(server_addr).expect("client connect"); - let mut reader = BufReader::new(stream.try_clone().expect("clone for reader")); - let mut writer = BufWriter::new(stream); + // Receive ObserverSnapshot + let response = read_framed(&mut reader) + .expect("read snapshot") + .expect("not EOF"); + let snapshot: ObserverSnapshot = + rmp_serde::from_slice(&response).expect("deserialize snapshot"); - // Send PlayerInput: MoveNorth at tick 0 - let inputs = vec![PlayerInput { - tick: 0, - action: PlayerAction::MoveNorth, - }]; - let payload = rmp_serde::to_vec_named(&inputs).expect("serialize inputs"); - write_framed(&mut writer, &payload).expect("send inputs"); - eprintln!("Client sent MoveNorth input"); + // Snapshot captures state at end of tick 0 (before advance_tick increments to 1) + assert_eq!(snapshot.tick, 0); + assert_eq!(snapshot.entities.len(), 1); - // Receive ObserverSnapshot - let response = read_framed(&mut reader) - .expect("read snapshot") - .expect("not EOF"); - let snapshot: ObserverSnapshot = - rmp_serde::from_slice(&response).expect("deserialize snapshot"); + let player_entity = &snapshot.entities[0]; + // Player started at (16, 16, 0), moved north (y-1) to (16, 15, 0) + // Render coords: (16.5, 15.5, 0) + assert_eq!(player_entity.x, 16.5); + assert_eq!(player_entity.y, 15.5); + assert_eq!(player_entity.z, 0); + assert!(matches!(player_entity.kind, EntityKind::Player)); - eprintln!( - "Client received snapshot: tick={}, entities={}", - snapshot.tick, - snapshot.entities.len() - ); - - // Verify snapshot - assert_eq!(snapshot.tick, 1); // After one tick - assert_eq!(snapshot.entities.len(), 1); // One entity (player) - - let player_entity = &snapshot.entities[0]; - // Player started at (16, 16, 0), moved north (y-1) to (16, 15, 0) - // Render coords: (16.5, 15.5, 0) - assert_eq!(player_entity.x, 16.5); - assert_eq!(player_entity.y, 15.5); - assert_eq!(player_entity.z, 0); - assert!(matches!(player_entity.kind, EntityKind::Player)); - - eprintln!( - "Client verified player moved to ({}, {}, {})", - player_entity.x, player_entity.y, player_entity.z - ); - } - - // Wait for server thread + // Clean up + drop(reader); + drop(writer); server_handle.join().expect("server thread panicked"); } From 3f4c663242b366c57baea84fa11a32d5d92156e3 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 11 Feb 2026 21:16:08 +0100 Subject: [PATCH 8/8] fix(bridge): improve accept error reporting with address context Replace expect() with unwrap_or_else that logs the bind address and error via tracing before exiting. Co-Authored-By: Claude Opus 4.6 --- server/src/main.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/server/src/main.rs b/server/src/main.rs index d1550f443..887c8a367 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -27,7 +27,10 @@ fn main() { tracing::info!("The Settled Reach - Simulation Server starting"); tracing::info!("Waiting for client connection on {}", addr); - let bridge = TcpBridge::accept(&addr).expect("Failed to accept client connection"); + let bridge = TcpBridge::accept(&addr).unwrap_or_else(|e| { + tracing::error!("Failed to accept client connection on {}: {}", addr, e); + std::process::exit(1); + }); tracing::info!("Client connected, initializing simulation"); // Create the bevy App and add plugins