Files
jpmschweitzerandClaude Opus 4.7 dfdc578429 feat(simulation): demux inbound bridge stream — receive() + atlas routing (#969, D-225)
Replaces the fixed-type SimBridge::receive_inputs() with a tagged
receive() -> Option<Inbound>, where Inbound is Inputs(Vec<PlayerInput>) or
AtlasRequest(AtlasLayerRequest). A shared decode_inbound() demuxes a frame by
shape (msgpack array = inputs, map = atlas request) — additive, no wire change
to existing input/snapshot frames. Adds send_atlas_response() to the trait
(both TcpBridge + LocalBridge impls). receive_bridge_inputs routes inputs to
the InputQueue as before; atlas requests to a new AtlasRequestBuffer (drained
by the serve system next). Integration tests (bridge_tcp/bridge_ipc) updated to
the tagged receive(); a demux unit test covers all three branches.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-24 18:00:30 +02:00

170 lines
5.5 KiB
Rust

//! 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::{Inbound, SimBridge};
use settled_reach_server::simulation::time::{DayPhase, TickRate};
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,
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");
});
// 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 = match bridge.receive().expect("failed to receive") {
Some(Inbound::Inputs(inputs)) => inputs,
other => panic!("expected inputs, got {other:?}"),
};
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 {
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");
// 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"),
}
}