Files
settled-reach/server/tests/ipc_bench.rs
T
jpmschweitzerandClaude Opus 4.6 c640563fc0 refactor(simulation): drop PROTOCOL_VERSION from bridge handshake (#874, D-192)
Removes the PROTOCOL_VERSION lockstep handshake across the bridge
transports. HandshakeMessage is now an empty marker indicating the
server is ready; the client sends StartupMessage next. Field-presence
and roundtrip behavioral tests are retained — tautological
version-literal assertions were deleted in Sprint 36.

After this lands, genuine schema drift surfaces as MessagePack
missing-field errors downstream, which is the intended signal per
D-192. Client ticket #875 is blocked by this commit and will ride
the same wave; server schema changes in #878 (character_archetype
removal) also depend on PROTOCOL_VERSION already being gone.

types.rs changes are bundled with the broader #877/#878 dead-code
sweep in the following commit, since they share the same file.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-22 08:55:19 +02:00

194 lines
6.6 KiB
Rust

//! IPC round-trip latency benchmark (#342, D-020)
//!
//! Measures end-to-end latency from client send (write_framed) to client receive
//! (read_framed) over the real subprocess IPC channel. Reports p50/p95/p99.
//!
//! Latency budget: p99 must be <= 5ms (D-020: "~1-5ms serialization latency per tick").
//!
//! Run with:
//! cargo test --release --test ipc_bench -- --ignored --nocapture
//!
//! Output: IPC_BENCH_RESULT:{json} on a single line for tooling to parse.
//!
//!
//! Spec references: D-020 (subprocess IPC, 5ms budget), D-030 (Layer 3)
use settled_reach_server::bridge::framing::{read_framed, write_framed};
use settled_reach_server::bridge::types::*;
use std::io::{BufRead, BufReader, BufWriter};
use std::net::TcpStream;
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};
/// Number of warmup round-trips before timing begins.
const WARMUP_ROUNDS: usize = 10;
/// Number of timed round-trips (N in the spec).
const MEASURE_ROUNDS: usize = 100;
/// Latency threshold (p99 must be below this). D-020: "~1-5ms per tick".
const THRESHOLD_MS: f64 = 5.0;
/// Timeout for the server to emit LISTENING:{port} on stdout.
const LISTEN_TIMEOUT: Duration = Duration::from_secs(15);
/// Timeout per round-trip read.
const ROUND_TRIP_TIMEOUT: Duration = Duration::from_secs(5);
fn percentile(sorted: &[f64], p: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let idx = ((sorted.len() - 1) as f64 * p).floor() as usize;
sorted[idx.min(sorted.len() - 1)]
}
#[test]
#[ignore]
fn ipc_round_trip_latency() {
// 1. Spawn server binary with --test-mode --port 0
let server_bin = env!("CARGO_BIN_EXE_settled-reach-server");
let mut child = Command::new(server_bin)
.args(["--test-mode", "--port", "0"])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("failed to spawn server binary");
let stdout = child.stdout.take().expect("stdout not captured");
let mut stdout_reader = BufReader::new(stdout);
// 2. Parse LISTENING:{port} from stdout
let port = {
let deadline = Instant::now() + LISTEN_TIMEOUT;
let mut line = String::new();
loop {
line.clear();
match stdout_reader.read_line(&mut line) {
Ok(0) => panic!("server stdout closed before LISTENING signal"),
Ok(_) => {
let trimmed = line.trim();
if let Some(port_str) = trimmed.strip_prefix("LISTENING:") {
break port_str
.parse::<u16>()
.unwrap_or_else(|e| panic!("invalid port '{}': {}", port_str, e));
}
}
Err(e) => panic!("failed to read server stdout: {}", e),
}
assert!(
Instant::now() < deadline,
"timed out waiting for LISTENING signal"
);
}
};
// 3. Connect via TCP
let addr = format!("127.0.0.1:{}", port);
let stream = TcpStream::connect(&addr)
.unwrap_or_else(|e| panic!("failed to connect to {}: {}", addr, e));
stream
.set_read_timeout(Some(ROUND_TRIP_TIMEOUT))
.expect("set read timeout");
let mut reader = BufReader::new(stream.try_clone().expect("clone stream"));
let mut writer = BufWriter::new(stream);
// 4. Handshake: read and validate HandshakeMessage before timing (#555/#556).
// Server sends HandshakeMessage { protocol_version } as the very first framed message.
let handshake_bytes = read_framed(&mut reader)
.expect("read handshake")
.expect("server closed before sending HandshakeMessage");
// D-192: HandshakeMessage carries no version field. Verify it deserialises cleanly.
let _handshake: HandshakeMessage =
rmp_serde::from_slice(&handshake_bytes).expect("deserialize HandshakeMessage");
let make_input = |tick: u64| PlayerInput {
tick,
action: PlayerAction::MoveNorth,
};
let mut round_trip_ms: Vec<f64> = Vec::with_capacity(WARMUP_ROUNDS + MEASURE_ROUNDS);
// 5. Warmup rounds (not timed)
for tick in 0..WARMUP_ROUNDS as u64 {
let payload =
rmp_serde::to_vec_named(&vec![make_input(tick)]).expect("serialize PlayerInput");
write_framed(&mut writer, &payload).expect("send warmup input");
let _ = read_framed(&mut reader)
.expect("read warmup snapshot")
.expect("server closed during warmup");
}
// 6. Timed measurement rounds
for tick in WARMUP_ROUNDS as u64..(WARMUP_ROUNDS + MEASURE_ROUNDS) as u64 {
let payload =
rmp_serde::to_vec_named(&vec![make_input(tick)]).expect("serialize PlayerInput");
let t_send = Instant::now();
write_framed(&mut writer, &payload).expect("send timed input");
let response = read_framed(&mut reader)
.expect("read timed snapshot")
.expect("server closed during measurement");
let elapsed_ms = t_send.elapsed().as_secs_f64() * 1000.0;
// Verify we received a valid snapshot (not just noise)
let _snapshot: ObserverSnapshot =
rmp_serde::from_slice(&response).expect("deserialize ObserverSnapshot");
round_trip_ms.push(elapsed_ms);
}
// 7. Clean up
drop(reader);
drop(writer);
let exit_deadline = Instant::now() + Duration::from_secs(5);
loop {
match child.try_wait() {
Ok(Some(_)) => break,
Ok(None) => {
if Instant::now() > exit_deadline {
child.kill().ok();
child.wait().ok();
break;
}
std::thread::sleep(Duration::from_millis(50));
}
Err(_) => {
child.kill().ok();
break;
}
}
}
// 8. Compute percentiles
let mut sorted = round_trip_ms.clone();
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap());
let p50 = percentile(&sorted, 0.50);
let p95 = percentile(&sorted, 0.95);
let p99 = percentile(&sorted, 0.99);
let passed = p99 <= THRESHOLD_MS;
let result = serde_json::json!({
"p50_ms": (p50 * 100.0).round() / 100.0,
"p95_ms": (p95 * 100.0).round() / 100.0,
"p99_ms": (p99 * 100.0).round() / 100.0,
"threshold_ms": THRESHOLD_MS,
"passed": passed,
"rounds": MEASURE_ROUNDS,
});
println!(
"IPC_BENCH_RESULT:{}",
serde_json::to_string(&result).unwrap()
);
// Fail the test if we exceed the latency budget
assert!(
passed,
"IPC latency budget exceeded: p99={:.2}ms > threshold={}ms",
p99, THRESHOLD_MS
);
}