Files
settled-reach/server/tests/bridge_tcp.rs
T

402 lines
15 KiB
Rust

//! 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<u8> {
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<Vec<PlayerInput>> {
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, HandshakeState, ServerRunning,
};
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: three frames back-to-back in one tick window — two input
// batches plus one atlas request (the D-225 demux path). Returns the
// stream so it stays open until assertions complete (no EOF race).
let client_handle = thread::spawn(move || {
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");
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::<InputQueue>();
world.init_resource::<ServerRunning>();
world.insert_resource(HandshakeState::Complete);
world.init_resource::<SimErrorBuffer>();
world.init_resource::<AtlasRequestBuffer>();
world
.run_system_once(receive_bridge_inputs)
.expect("receive_bridge_inputs failed to run");
assert_eq!(
world.resource::<InputQueue>().len(),
2,
"both input batches must drain in a single tick"
);
assert_eq!(
world.resource::<AtlasRequestBuffer>().0.len(),
1,
"the atlas request must drain in the same tick"
);
assert!(
world.resource::<ServerRunning>().0,
"draining must not shut the server down"
);
}