From d1aba3d5542461f8d4ef799ab6a1d1625cff4c95 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Mon, 16 Feb 2026 23:25:45 +0100 Subject: [PATCH] fix(simulation): use non-blocking TCP to prevent game loop stall MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit read_framed() called read_exact() which blocked the entire bevy Update schedule waiting for client input — no systems ran until a keystroke arrived. Switch TcpStream to non-blocking mode so receive_inputs() returns Ok(vec![]) on WouldBlock instead of blocking. Toggle to blocking for snapshot writes (reliable delivery). Add 50ms frame throttle (~20 ticks/sec) since the non-blocking loop would otherwise spin. Downgrade input receive logging to trace, add error logging for missing bridge resource and failed observer queries. Fixes bug #1 (server never sends snapshots) and #2 (camera doesn't center until first keystroke). Co-Authored-By: Claude Opus 4.6 --- server/src/bridge/mod.rs | 7 +++-- server/src/bridge/tcp.rs | 28 +++++++++++++++++--- server/src/main.rs | 14 ++++++++-- server/src/perception/observer/mod.rs | 2 ++ server/tests/bridge_tcp.rs | 37 +++++++++++++++++++-------- 5 files changed, 69 insertions(+), 19 deletions(-) diff --git a/server/src/bridge/mod.rs b/server/src/bridge/mod.rs index f1803b5bc..35340576a 100644 --- a/server/src/bridge/mod.rs +++ b/server/src/bridge/mod.rs @@ -73,7 +73,7 @@ pub fn receive_bridge_inputs( match bridge.receive_inputs() { Ok(inputs) => { for input in &inputs { - tracing::debug!( + tracing::trace!( "Received input: tick={} action={:?}", input.tick, input.action @@ -115,7 +115,10 @@ pub fn send_bridge_snapshot( mut buffer: ResMut, mut running: ResMut, ) { - let Some(bridge) = bridge else { return }; + let Some(bridge) = bridge else { + tracing::error!("send_bridge_snapshot: no BridgeResource"); + return; + }; if let Some(snapshot) = buffer.snapshot.take() { if let Err(e) = bridge.send_snapshot(&snapshot) { match &e { diff --git a/server/src/bridge/tcp.rs b/server/src/bridge/tcp.rs index fe640f7c7..2a1bb56dd 100644 --- a/server/src/bridge/tcp.rs +++ b/server/src/bridge/tcp.rs @@ -44,6 +44,12 @@ impl TcpBridge { local_addr ); + // Set non-blocking so receive_inputs doesn't stall the game loop. + // read_framed handles WouldBlock by returning Ok(None). + stream.set_nonblocking(true).map_err(|e| { + BridgeError::Transport(format!("failed to set non-blocking: {}", e)) + })?; + // 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)) @@ -75,6 +81,10 @@ impl TcpBridge { local_addr ); + stream.set_nonblocking(true).map_err(|e| { + BridgeError::Transport(format!("failed to set non-blocking: {}", e)) + })?; + let reader_stream = stream.try_clone().map_err(|e| { BridgeError::Transport(format!("failed to clone stream for reader: {}", e)) })?; @@ -126,7 +136,14 @@ impl SimBridge for TcpBridge { .writer .lock() .map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?; - write_framed(writer.get_mut(), &payload)?; + + // Toggle to blocking for reliable write delivery. + // Single-threaded bevy guarantees no concurrent reads during this window. + let stream = writer.get_mut(); + stream.set_nonblocking(false).map_err(BridgeError::Io)?; + let result = write_framed(stream, &payload); + stream.set_nonblocking(true).map_err(BridgeError::Io)?; + result?; tracing::trace!("sent snapshot: tick={}", snapshot.tick); Ok(()) @@ -138,8 +155,8 @@ impl SimBridge for TcpBridge { .lock() .map_err(|e| BridgeError::MutexPoisoned(format!("reader: {}", e)))?; - match read_framed(reader.get_mut())? { - Some(payload) => match rmp_serde::from_slice::>(&payload) { + match read_framed(reader.get_mut()) { + Ok(Some(payload)) => match rmp_serde::from_slice::>(&payload) { Ok(inputs) => { tracing::trace!("received {} inputs", inputs.len()); Ok(inputs) @@ -160,7 +177,10 @@ impl SimBridge for TcpBridge { ))) } }, - None => Err(BridgeError::Disconnected), + Ok(None) => Err(BridgeError::Disconnected), + // Non-blocking socket: no data available this tick — not an error. + Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => Ok(vec![]), + Err(e) => Err(BridgeError::Io(e)), } } } diff --git a/server/src/main.rs b/server/src/main.rs index 750312267..acaf1266f 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -226,13 +226,23 @@ fn main() { tracing::info!("Simulation initialized, entering game loop"); - // Game loop: run until client disconnects + // Game loop: run until client disconnects. + // Targets ~20 ticks/sec (2 game-minutes/sec). The TCP bridge uses + // non-blocking reads, so without throttling this loop would spin. + // Remaining frame budget is available for NPC AI and pathfinding. + let target_frame_time = std::time::Duration::from_millis(50); loop { + let frame_start = std::time::Instant::now(); + app.update(); - // Check ServerRunning resource if !app.world().resource::().0 { break; } + + let elapsed = frame_start.elapsed(); + if elapsed < target_frame_time { + std::thread::sleep(target_frame_time - elapsed); + } } tracing::info!("Simulation server shutting down"); diff --git a/server/src/perception/observer/mod.rs b/server/src/perception/observer/mod.rs index e39429eb2..0adf37aa6 100644 --- a/server/src/perception/observer/mod.rs +++ b/server/src/perception/observer/mod.rs @@ -33,6 +33,7 @@ pub fn compute_visibility_geometry( mut geometry: ResMut, ) { let Ok((observer_pos, facing_opt)) = observer_query.single() else { + tracing::error!("compute_visibility_geometry: PlayerCharacter query failed"); return; }; @@ -90,6 +91,7 @@ pub fn compute_observer_snapshot( cognitive_delay_opt, )) = observer_query.single_mut() else { + tracing::error!("compute_observer_snapshot: PlayerCharacter query failed"); return; }; diff --git a/server/tests/bridge_tcp.rs b/server/tests/bridge_tcp.rs index bc30231e7..1c0168606 100644 --- a/server/tests/bridge_tcp.rs +++ b/server/tests/bridge_tcp.rs @@ -79,7 +79,18 @@ fn input_roundtrip_over_tcp() { let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); - let inputs = bridge.receive_inputs().expect("failed to receive inputs"); + // Non-blocking socket: retry until data arrives or timeout. + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + let inputs = loop { + match bridge.receive_inputs() { + Ok(inputs) if !inputs.is_empty() => break inputs, + Ok(_) => { + assert!(std::time::Instant::now() < deadline, "timed out waiting for inputs"); + thread::sleep(std::time::Duration::from_millis(1)); + } + Err(e) => panic!("failed to receive inputs: {}", e), + } + }; assert_eq!(inputs.len(), 2); assert_eq!(inputs[0].tick, 10); @@ -132,16 +143,20 @@ fn tcp_bridge_eof_returns_error() { let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); - let result = bridge.receive_inputs(); - - assert!(result.is_err(), "expected error on EOF"); - assert!( - matches!( - result, - Err(settled_reach_server::bridge::BridgeError::Disconnected) - ), - "expected Disconnected error" - ); + // Non-blocking socket: retry until we get Disconnected or timeout. + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + loop { + match bridge.receive_inputs() { + Ok(inputs) if inputs.is_empty() => { + // WouldBlock — client hasn't disconnected yet, retry + assert!(std::time::Instant::now() < deadline, "timed out waiting for EOF"); + thread::sleep(std::time::Duration::from_millis(1)); + } + Ok(inputs) => panic!("expected Disconnected error, got {} inputs", inputs.len()), + Err(settled_reach_server::bridge::BridgeError::Disconnected) => break, + Err(e) => panic!("expected Disconnected error, got: {}", e), + } + } }); // Client: connect and immediately disconnect without sending data