//! 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::{Inbound, SimBridge}; use settled_reach_server::simulation::time::{DayPhase, TickRate}; use std::net::{TcpListener, TcpStream}; use std::thread; #[test] fn snapshot_roundtrip_over_tcp() { // 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"); // Server thread: accept on pre-bound listener and send snapshot let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); let snapshot = ObserverSnapshot { tick: 42, game_time: GameTime { day: 0, time_of_day: 0, day_phase: DayPhase::Morning, tick_rate: TickRate::Full, }, player_facing: FacingDirection::North, player_stance: MovementStance::default(), player_inventory: vec![], entities: vec![VisibleEntity { entity_id: 100, x: 10.5, y: 20.3, z: 0, kind: EntityKind::Npc, visibility: VisibilitySector::Forward, relationship: RelationshipState::Unknown, observation: EntityVisibility::Visible, tell_state: None, }], visible_tiles: vec![], nearby_interactions: vec![], current_monologue: None, pending_recognitions: vec![], dialogue_response: None, blocked_entities: vec![], scan_events: vec![], sound_events: vec![], follow_state: None, character_pressure: None, rng_seed: None, poi_list: vec![], examine_result: None, player_knowledge: None, save_result: None, triangle_crisis_events: vec![], state_hash: None, debug_response: None, sim_errors: vec![], current_ticker: None, settings_response: None, economy_snapshot: None, bookmark_catalog: None, }; bridge .send_snapshot(&snapshot) .expect("failed to send 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); 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() { 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"); let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); // 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() { Ok(Some(Inbound::Inputs(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); assert_eq!(inputs[1].tick, 11); inputs }); // 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 { target_entity_id: None, verb: None, }, }, ]; 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"); 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() { 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"); let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); // 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() { Ok(None) => { // No data yet — client hasn't disconnected, retry assert!( std::time::Instant::now() < deadline, "timed out waiting for EOF" ); thread::sleep(std::time::Duration::from_millis(1)); } Ok(Some(_)) => panic!("expected Disconnected error, got a message"), Err(settled_reach_server::bridge::BridgeError::Disconnected) => break, Err(e) => panic!("expected Disconnected error, got: {}", e), } } }); // Client: connect and immediately disconnect without sending data let stream = TcpStream::connect(server_addr).expect("failed to connect"); drop(stream); server_handle.join().expect("server thread panicked"); } /// T-1045 regression: EOF *mid-frame* (peer dies after a partial write) must /// escalate to `Disconnected` like clean EOF — not surface an Io error every /// tick forever against a frame that can never complete. #[test] fn tcp_bridge_eof_mid_frame_escalates_to_disconnected() { 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"); let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); loop { match bridge.receive() { Ok(None) => { assert!( std::time::Instant::now() < deadline, "timed out waiting for mid-frame EOF" ); thread::sleep(std::time::Duration::from_millis(1)); } Ok(Some(_)) => panic!("expected Disconnected, got a message"), Err(settled_reach_server::bridge::BridgeError::Disconnected) => break, Err(e) => panic!("expected Disconnected, got: {}", e), } } }); // Client: write 2 of the 4 length-prefix bytes, then die. use std::io::Write; let mut stream = TcpStream::connect(server_addr).expect("failed to connect"); stream .write_all(&[0x00, 0x00]) .expect("partial write failed"); stream.flush().expect("flush failed"); thread::sleep(std::time::Duration::from_millis(50)); drop(stream); server_handle.join().expect("server thread panicked"); } /// Build a raw frame (4-byte BE length prefix + payload) for split-write tests. fn raw_frame(payload: &[u8]) -> Vec { let mut frame = (payload.len() as u32).to_be_bytes().to_vec(); frame.extend_from_slice(payload); frame } /// Poll the non-blocking bridge until `count` input batches arrive. fn collect_input_batches(bridge: &TcpBridge, count: usize) -> Vec> { let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); let mut batches = Vec::new(); while batches.len() < count { match bridge.receive() { Ok(Some(Inbound::Inputs(inputs))) if !inputs.is_empty() => batches.push(inputs), Ok(_) => { assert!( std::time::Instant::now() < deadline, "timed out waiting for {} input batches (got {})", count, batches.len() ); thread::sleep(std::time::Duration::from_millis(1)); } Err(e) => panic!("failed to receive inputs: {}", e), } } batches } /// T-1045 regression: a frame delivered in two TCP writes (split at /// `split_point(frame_len)`) must decode, and the NEXT frame must also decode /// — i.e. a partial read mid-frame must not desync the stream. The server /// polls receive() during the gap, so it observes WouldBlock mid-frame. fn assert_split_frame_does_not_desync(split_point: impl Fn(usize) -> usize) { use std::io::Write; 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"); let server_handle = thread::spawn(move || { let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); collect_input_batches(&bridge, 2) }); let mut stream = TcpStream::connect(server_addr).expect("failed to connect"); let inputs1 = vec![PlayerInput { tick: 10, action: PlayerAction::MoveNorth, }]; let payload1 = rmp_serde::to_vec_named(&inputs1).expect("failed to serialize"); let frame1 = raw_frame(&payload1); let split = split_point(frame1.len()); assert!(split > 0 && split < frame1.len(), "split must be mid-frame"); // First half, then a gap long enough for the server to poll mid-frame. stream .write_all(&frame1[..split]) .expect("write first half"); stream.flush().expect("flush first half"); thread::sleep(std::time::Duration::from_millis(50)); stream .write_all(&frame1[split..]) .expect("write second half"); stream.flush().expect("flush second half"); // A second frame in a single write — decodes only if the stream is in sync. let inputs2 = vec![PlayerInput { tick: 11, action: PlayerAction::MoveSouth, }]; let payload2 = rmp_serde::to_vec_named(&inputs2).expect("failed to serialize"); stream .write_all(&raw_frame(&payload2)) .expect("write second frame"); stream.flush().expect("flush second frame"); let batches = server_handle.join().expect("server thread panicked"); assert_eq!(batches[0].len(), 1); assert_eq!(batches[0][0].tick, 10); assert!(matches!(batches[0][0].action, PlayerAction::MoveNorth)); assert_eq!(batches[1].len(), 1); assert_eq!(batches[1][0].tick, 11); assert!(matches!(batches[1][0].action, PlayerAction::MoveSouth)); } #[test] fn frame_split_inside_prefix_does_not_desync() { // Split inside the 4-byte length prefix. assert_split_frame_does_not_desync(|_| 2); } #[test] fn frame_split_inside_payload_does_not_desync() { // Split midway through the payload (prefix is 4 bytes). assert_split_frame_does_not_desync(|frame_len| 4 + (frame_len - 4) / 2); } /// T-1045 drain: one receive_bridge_inputs run (= one 50 ms tick) must drain /// every complete frame buffered on the socket, not one frame per tick. #[test] fn single_tick_drains_all_ready_inbound_frames() { use bevy_ecs::system::RunSystemOnce; use settled_reach_server::atlas::cascade::CascadeLayer; use settled_reach_server::atlas::layer_proxy::AtlasLayerRequest; use settled_reach_server::bridge::{ receive_bridge_inputs, AtlasRequestBuffer, BridgeResource, CityNamesRequestBuffer, HandshakeState, ServerRunning, StarMapRequestBuffer, }; use settled_reach_server::simulation::input::InputQueue; 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"); // Client: five frames back-to-back in one tick window — two input // batches plus one of EACH request shape (atlas, star-map, city-names: // the full D-225/T-949 demux surface over the real framing/poll path — // PR #176 review H6). Returns the stream so it stays open until // assertions complete (no EOF race). let client_handle = thread::spawn(move || { use settled_reach_server::atlas::atlas_data_proxy::{CityNamesRequest, StarMapRequest}; let mut stream = TcpStream::connect(server_addr).expect("failed to connect"); for tick in [20u64, 21] { let inputs = vec![PlayerInput { tick, action: PlayerAction::MoveNorth, }]; let payload = rmp_serde::to_vec_named(&inputs).expect("failed to serialize"); write_framed(&mut stream, &payload).expect("write input frame"); } let req = AtlasLayerRequest { body_id: "GJ1c".into(), up_to: CascadeLayer::Topography, }; let payload = rmp_serde::to_vec_named(&req).expect("failed to serialize"); write_framed(&mut stream, &payload).expect("write atlas frame"); let sm = StarMapRequest { star_map: true }; let payload = rmp_serde::to_vec_named(&sm).expect("failed to serialize star map"); write_framed(&mut stream, &payload).expect("write star map frame"); let cn = CityNamesRequest { city_names: true, body_id: "GJ1c".into(), }; let payload = rmp_serde::to_vec_named(&cn).expect("failed to serialize city names"); write_framed(&mut stream, &payload).expect("write city names frame"); stream }); let bridge = TcpBridge::accept_on(listener).expect("failed to accept"); let _stream = client_handle.join().expect("client thread panicked"); // Writes are flushed and joined; small grace period for loopback delivery. thread::sleep(std::time::Duration::from_millis(100)); let mut world = bevy_ecs::world::World::new(); world.insert_resource(BridgeResource::new(bridge)); world.init_resource::(); world.init_resource::(); world.insert_resource(HandshakeState::Complete); world.init_resource::(); world.init_resource::(); world.init_resource::(); world.init_resource::(); world .run_system_once(receive_bridge_inputs) .expect("receive_bridge_inputs failed to run"); assert_eq!( world.resource::().len(), 2, "both input batches must drain in a single tick" ); assert_eq!( world.resource::().0.len(), 1, "the atlas request must drain in the same tick" ); assert_eq!( world.resource::().0.len(), 1, "the star-map request must drain in the same tick (H6: real wire path)" ); let city_names = &world.resource::().0; assert_eq!( city_names.len(), 1, "the city-names request must drain in the same tick (H6: real wire path)" ); assert_eq!(city_names[0].body_id, "GJ1c"); assert!( world.resource::().0, "draining must not shut the server down" ); }