Files
OpenFUT/openfut-blaze-host/tests/live_transport.rs
T
funman300 23374312bc blaze-host: opt-in raw frame capture + auditable sanitizer
Evidence infrastructure, not protocol functionality. Built before gates 7-10
because those sessions cannot be reproduced -- a later run is a different
session, and the migration-validation runs happen once. Gates 5-6 already went
past without their bytes being recorded.

TWO LAYERS

  live FIFA traffic
    ├── raw capture      exact RX/TX bytes, mode 0600, gitignored
    └── blaze-sanitize   → repository-safe, replayable fixtures

CAPTURE. Off unless OPENFUT_BLAZE_CAPTURE names a file. Deterministic
big-endian container: 20-byte file header, then per-frame records carrying
connection id, a global monotonic sequence, timestamp, direction and the EXACT
frame bytes. RX is recorded as received; TX only AFTER a successful write, so a
record means the bytes were sent rather than intended.

Component/command/msgNum/msgType/payload length are deliberately NOT stored
beside the frame: they are already in its 16-byte header, and a redundant copy
can disagree with the bytes, leaving a reader unable to tell which is true.
Record::header() derives them, so every field the requirements name is
available without duplicating it.

SANITIZER. Redacts only the named tags in SENSITIVE_TAGS (KEY, AUTH, SESS,
MAIL, PML) and reports every substitution with path, kind and length.
Replacement is LENGTH-PRESERVING, so the TDF varint, payload length and Fire2
header are unchanged and the sanitized frame is exactly the size of the
captured one -- asserted per frame, failing rather than emitting a subtly
different conversation. Frames with nothing sensitive keep their exact wire
bytes. Payloads that will not decode are passed through and REPORTED, so a
reader knows they were never inspected rather than assuming they were checked.

TESTS. 39 in this crate. All nine required cases: capture disabled produces no
artefact; RX and TX captured exactly; ordering preserved; fragmented input
(one byte at a time) reconstructs the same frames as a single write; coalesced
input is captured as separate frames, not per-read; capture does not alter wire
output; sanitization removes a real session key from a real captured login;
malformed/truncated/wrong-version captures fail clearly; every listed sensitive
tag is provably reachable.

MUTATION TESTED. Dropping TX capture, truncating captured frames to their
header, and removing KEY from the sensitive list were each verified to turn the
suite red. One mutation was NOT caught: moving the TX capture above the write.
It is indistinguishable while writes succeed and only diverges when one fails.
That invariant is held by code placement and a comment saying so, not by a
test, and the code says as much rather than implying coverage it does not have.

Python oracle unchanged.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-11 02:14:25 +00:00

751 lines
25 KiB
Rust

//! Transport parity: the recorded conversations over a real TCP socket.
//!
//! The adapter's own suite calls `dispatch()` directly. That proves what the
//! adapter decides, and nothing about what a socket does. This suite runs the
//! actual host process-internally, connects over real TCP, and replays the same
//! recorded conversations — which is what exercises the things fixtures cannot:
//!
//! * a stream that delivers bytes, not frames
//! * requests split across reads and coalesced into one
//! * many frames per connection, and session state surviving between them
//! * multiple frames written back in order for a single request
//! * connection lifecycle and close reasons
//!
//! Byte-exactness is only achievable because `Hooks` injects the session key
//! and clock; with the real ones, no live run could reproduce a recording, and
//! the guarantee would drop to structural.
use std::io::{Read, Write};
use std::net::TcpStream;
use std::sync::Mutex;
use std::time::Duration;
use openfut_adapter_fifa17::blaze::{AdapterConfig, Endpoints, Identity};
use openfut_blaze_host::{bind, capture, config::HostConfig, sanitize, Hooks};
use openfut_protocol_blaze::fire2::{Frame, Header, HEADER_LEN};
use serde_json::Value as J;
fn records() -> Vec<J> {
let path = format!(
"{}/../openfut-adapter-fifa17/fixtures/blaze_transactions.jsonl",
env!("CARGO_MANIFEST_DIR")
);
std::fs::read_to_string(&path)
.unwrap_or_else(|e| panic!("cannot read {path}: {e}"))
.lines()
.filter(|l| !l.trim().is_empty())
.map(|l| serde_json::from_str(l).expect("valid JSON"))
.collect()
}
fn unhex(s: &str) -> Vec<u8> {
(0..s.len())
.step_by(2)
.map(|i| u8::from_str_radix(&s[i..i + 2], 16).expect("valid hex"))
.collect()
}
fn hex(b: &[u8]) -> String {
b.iter().map(|x| format!("{x:02x}")).collect()
}
fn st(j: &J, k: &str) -> String {
j[k].as_str().expect("string").to_string()
}
fn num(j: &J, k: &str) -> i64 {
j[k].as_i64().expect("number")
}
struct Harness {
addr: String,
records: Vec<J>,
}
/// Start the real host on an ephemeral port with the recording's config.
///
/// Session keys are handed out in the order the recording declares them, so the
/// Nth connection this test opens gets the Nth recorded key. That is why the
/// connection order below must match the fixture's session order.
fn start() -> Harness {
start_with_capture(None)
}
fn start_with_capture(capture_path: Option<String>) -> Harness {
let records = records();
let cfg_rec = records
.iter()
.find(|r| r["kind"] == "config")
.expect("config record")
.clone();
let id = &cfg_rec["identity"];
let adapter_cfg = AdapterConfig {
identity: Identity {
persona_id: num(id, "persona_id"),
persona_name: st(id, "persona_name"),
user_id: num(id, "user_id"),
ext_id: num(id, "ext_id"),
email: st(id, "email"),
namespace: st(id, "namespace"),
client_platform: num(id, "client_platform"),
persona_status: num(id, "persona_status"),
user_session_type: num(id, "user_session_type"),
account_locale: num(id, "account_locale_int"),
locale: st(id, "locale"),
content_id: st(id, "content_id"),
entitlement_tag: st(id, "entitlement_tag"),
entitlement_group: st(id, "entitlement_group"),
title_id: st(id, "title_id"),
client_id: st(id, "client_id"),
platform: st(id, "platform"),
},
endpoints: Endpoints {
advertise: st(&cfg_rec, "advertise"),
bind: st(&cfg_rec, "bind"),
pow_content_host: st(&cfg_rec, "pow_content_host"),
pow_host: st(&cfg_rec, "pow_host"),
..Endpoints::default()
},
server_version: st(id, "server_version"),
};
let host_cfg = HostConfig {
// Port 0: the OS picks a free one, so this can never collide with the
// running Python backend or with a parallel test.
listen_addr: "127.0.0.1".into(),
listen_port: 0,
idle_timeout_secs: 10,
max_payload_bytes: 4 * 1024 * 1024,
trace_path: None,
capture_path,
adapter: adapter_cfg,
};
let keys: Vec<String> = records
.iter()
.filter(|r| r["kind"] == "session")
.map(|r| st(r, "session_key"))
.collect();
let queue = Mutex::new(keys.into_iter().collect::<std::collections::VecDeque<_>>());
let now = num(&cfg_rec, "now");
let hooks = Hooks {
session_key: Box::new(move || {
queue
.lock()
.expect("key queue")
.pop_front()
.expect("a recorded session key for each connection")
}),
now: Box::new(move || now),
};
let server = bind(host_cfg, hooks).expect("bind ephemeral port");
let addr = server.local_addr.to_string();
std::thread::spawn(move || {
let _ = server.run();
});
Harness { addr, records }
}
fn connect(addr: &str) -> TcpStream {
let s = TcpStream::connect(addr).expect("connect to host");
s.set_read_timeout(Some(Duration::from_secs(10))).unwrap();
s.set_nodelay(true).unwrap();
s
}
fn read_frame(stream: &mut TcpStream, buf: &mut Vec<u8>) -> Option<Frame> {
let mut chunk = [0u8; 65536];
while buf.len() < HEADER_LEN {
match stream.read(&mut chunk) {
Ok(0) | Err(_) => return None,
Ok(got) => buf.extend_from_slice(&chunk[..got]),
}
}
let total = Header::parse(&buf[..HEADER_LEN]).ok()?.frame_len();
while buf.len() < total {
match stream.read(&mut chunk) {
Ok(0) | Err(_) => return None,
Ok(got) => buf.extend_from_slice(&chunk[..got]),
}
}
let (frame, used) = Frame::parse(&buf[..total]).ok()?;
buf.drain(..used);
Some(frame)
}
/// Transactions for one recorded session, in order.
fn transactions<'a>(records: &'a [J], session: &str) -> Vec<&'a J> {
records
.iter()
.filter(|r| r["kind"] == "tx" && r["session"] == session)
.collect()
}
fn session_ids(records: &[J]) -> Vec<String> {
records
.iter()
.filter(|r| r["kind"] == "session")
.map(|r| st(r, "id"))
.collect()
}
/// The headline transport test: every recorded conversation, over TCP,
/// byte-for-byte, one connection per recorded session.
#[test]
fn recorded_conversations_replay_byte_for_byte_over_tcp() {
let h = start();
let mut total_frames = 0usize;
for sid in session_ids(&h.records) {
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
for tx in transactions(&h.records, &sid) {
let name = st(tx, "name");
let request = unhex(&st(tx, "request_hex"));
let expected: Vec<String> = tx["responses"]
.as_array()
.unwrap()
.iter()
.map(|v| v.as_str().unwrap().to_string())
.collect();
stream.write_all(&request).expect("send request");
stream.flush().unwrap();
for (i, want) in expected.iter().enumerate() {
let frame = read_frame(&mut stream, &mut buf)
.unwrap_or_else(|| panic!("{sid}/{name}: no frame {i}, expected one"));
assert_eq!(
hex(&frame.encode()),
*want,
"\n{sid}/{name}: frame {i} differs over the wire"
);
total_frames += 1;
}
}
// A clean close between frames must be recognised as such.
drop(stream);
}
assert!(
total_frames >= 50,
"expected the full recorded conversation, saw {total_frames} frames"
);
}
/// The same conversation with every request written ONE BYTE AT A TIME.
///
/// This is the fragmentation case: a real client's frames arrive split across
/// reads, and a host that assumed one read per frame would pass every fixture
/// test and then fail against FIFA.
#[test]
fn requests_split_across_reads_are_reassembled() {
let h = start();
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
for tx in transactions(&h.records, &sid).into_iter().take(6) {
let name = st(tx, "name");
let request = unhex(&st(tx, "request_hex"));
for byte in &request {
stream.write_all(&[*byte]).expect("dribble a byte");
stream.flush().unwrap();
}
for (i, want) in tx["responses"].as_array().unwrap().iter().enumerate() {
let frame = read_frame(&mut stream, &mut buf)
.unwrap_or_else(|| panic!("{name}: no frame {i} after fragmented send"));
assert_eq!(
hex(&frame.encode()),
want.as_str().unwrap(),
"{name} frame {i}"
);
}
}
}
/// Several requests written as ONE write, so the host sees them coalesced in a
/// single read and must split them itself.
#[test]
fn coalesced_requests_in_one_write_are_split() {
let h = start();
let sid = session_ids(&h.records).into_iter().next().unwrap();
let txs: Vec<&J> = transactions(&h.records, &sid).into_iter().take(4).collect();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
let mut blob = Vec::new();
for tx in &txs {
blob.extend_from_slice(&unhex(&st(tx, "request_hex")));
}
stream.write_all(&blob).expect("one big write");
stream.flush().unwrap();
for tx in &txs {
let name = st(tx, "name");
for (i, want) in tx["responses"].as_array().unwrap().iter().enumerate() {
let frame = read_frame(&mut stream, &mut buf)
.unwrap_or_else(|| panic!("{name}: no frame {i} after coalesced send"));
assert_eq!(
hex(&frame.encode()),
want.as_str().unwrap(),
"{name} frame {i}"
);
}
}
}
/// Login must deliver four frames in order on a real socket, not just from
/// `dispatch()`. This is the ordering guarantee the client depends on.
#[test]
fn login_burst_arrives_in_order_over_the_wire() {
let h = start();
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
for tx in transactions(&h.records, &sid) {
let request = unhex(&st(tx, "request_hex"));
stream.write_all(&request).unwrap();
stream.flush().unwrap();
let expected = tx["responses"].as_array().unwrap().len();
let mut got = Vec::new();
for _ in 0..expected {
got.push(read_frame(&mut stream, &mut buf).expect("frame"));
}
if st(tx, "name") == "login" {
assert_eq!(got.len(), 4, "reply + three pushes");
let routes: Vec<(u16, u16)> = got
.iter()
.map(|f| (f.header.component, f.header.command))
.collect();
assert_eq!(
routes,
vec![
(0x0001, 0x000A), // Authentication::login REPLY
(0x7802, 0x0008), // UserAuthenticated
(0x7802, 0x0001), // UserSessionExtendedDataUpdate
(0x7802, 0x0002), // UserAdded
]
);
return;
}
}
panic!("login transaction never reached");
}
/// Session state must persist across frames on one connection: the auth code
/// login records has to still be there when getAuthToken asks for it later.
#[test]
fn session_state_persists_across_frames_on_one_connection() {
let h = start();
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
let txs = transactions(&h.records, &sid);
let mut saw_before = false;
let mut saw_after = false;
for tx in txs {
let name = st(tx, "name");
stream.write_all(&unhex(&st(tx, "request_hex"))).unwrap();
stream.flush().unwrap();
for want in tx["responses"].as_array().unwrap() {
let frame = read_frame(&mut stream, &mut buf).expect("frame");
let got = hex(&frame.encode());
assert_eq!(got, want.as_str().unwrap(), "{name}");
if name == "get_auth_token_before_login" {
let body = openfut_protocol_blaze::heat2::decode(&frame.payload).unwrap();
let tok = body.get("AUTH").and_then(|v| v.as_str()).unwrap();
assert!(
tok.starts_with("OPENFUT-"),
"synthesised before login: {tok}"
);
saw_before = true;
}
if name == "get_auth_token_after_login" {
let body = openfut_protocol_blaze::heat2::decode(&frame.payload).unwrap();
let tok = body.get("AUTH").and_then(|v| v.as_str()).unwrap();
assert_eq!(tok, "OPENFUT-TEST-AUTHCODE", "echoes the login's code");
saw_after = true;
}
}
}
assert!(
saw_before && saw_after,
"both getAuthToken phases must be exercised"
);
}
/// A separate connection must get a separate session — no leakage.
#[test]
fn a_second_connection_gets_a_fresh_session() {
let h = start();
let ids = session_ids(&h.records);
// Connection 1: log in.
let mut c1 = connect(&h.addr);
let mut b1 = Vec::new();
for tx in transactions(&h.records, &ids[0]) {
c1.write_all(&unhex(&st(tx, "request_hex"))).unwrap();
c1.flush().unwrap();
for _ in 0..tx["responses"].as_array().unwrap().len() {
read_frame(&mut c1, &mut b1).expect("frame");
}
}
// Connection 2 asks for an auth token without logging in. If session state
// leaked it would echo connection 1's code instead of synthesising one.
let mut c2 = connect(&h.addr);
let mut b2 = Vec::new();
let get_token = transactions(&h.records, &ids[0])
.into_iter()
.find(|t| st(t, "name") == "get_auth_token_before_login")
.expect("fixture present");
c2.write_all(&unhex(&st(get_token, "request_hex"))).unwrap();
c2.flush().unwrap();
let frame = read_frame(&mut c2, &mut b2).expect("frame");
let body = openfut_protocol_blaze::heat2::decode(&frame.payload).unwrap();
let tok = body.get("AUTH").and_then(|v| v.as_str()).unwrap();
assert!(
tok.starts_with("OPENFUT-"),
"a fresh connection must not inherit a login: {tok}"
);
assert_ne!(tok, "OPENFUT-TEST-AUTHCODE");
}
/// A header claiming an absurd payload must drop the connection rather than
/// try to allocate for it.
#[test]
fn an_absurd_payload_length_drops_the_connection() {
let h = start();
let mut stream = connect(&h.addr);
let mut header = [0u8; HEADER_LEN];
header[0..4].copy_from_slice(&0x7FFF_FFFFu32.to_be_bytes()); // ~2 GiB
header[6..8].copy_from_slice(&0x0009u16.to_be_bytes());
header[8..10].copy_from_slice(&0x0002u16.to_be_bytes());
stream.write_all(&header).unwrap();
stream.flush().unwrap();
let mut buf = Vec::new();
assert!(
read_frame(&mut stream, &mut buf).is_none(),
"the host must close, not answer, an absurd frame"
);
}
/// A body that will not decode as TDF must not kill the connection: the oracle
/// logs it and dispatches with no fields.
#[test]
fn an_undecodable_body_still_gets_a_reply() {
let h = start();
let mut stream = connect(&h.addr);
// Util::ping with a garbage payload. ping ignores its body, so the reply
// must arrive exactly as if the body had been empty.
let mut frame = Frame::new(
0x0009,
0x0002,
1,
openfut_protocol_blaze::fire2::MsgType::Message,
vec![0xFF; 8],
);
frame.header.payload_len = 8;
stream.write_all(&frame.encode()).unwrap();
stream.flush().unwrap();
let mut buf = Vec::new();
let reply = read_frame(&mut stream, &mut buf).expect("a reply despite the bad body");
assert_eq!(reply.header.component, 0x0009);
assert_eq!(reply.header.command, 0x0002);
assert!(!reply.payload.is_empty(), "ping still answers with STIM");
}
// ───────────────────────────── raw frame capture ─────────────────────────────
//
// Capture is evidence infrastructure. These tests exist because the value of a
// capture is entirely in its fidelity: a capture that quietly drops, reorders
// or alters frames is worse than none, since it would be trusted.
fn cap_path(name: &str) -> String {
let dir = std::env::var("TMPDIR").unwrap_or_else(|_| "/tmp".into());
format!("{dir}/ofcap-it-{}-{name}.ofcap", std::process::id())
}
/// Drive a conversation over TCP and return (captured records, wire responses).
fn converse(h: &Harness, take: usize) -> Vec<Vec<u8>> {
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
let mut got = Vec::new();
for tx in transactions(&h.records, &sid).into_iter().take(take) {
stream.write_all(&unhex(&st(tx, "request_hex"))).unwrap();
stream.flush().unwrap();
for _ in 0..tx["responses"].as_array().unwrap().len() {
got.push(read_frame(&mut stream, &mut buf).expect("frame").encode());
}
}
got
}
#[test]
fn capture_disabled_produces_no_artefact() {
let path = cap_path("disabled");
let _ = std::fs::remove_file(&path);
let h = start(); // capture_path: None
converse(&h, 4);
assert!(
!std::path::Path::new(&path).exists(),
"no capture file may appear when capture is off"
);
}
#[test]
fn capture_records_rx_and_tx_frames_exactly_and_in_order() {
let path = cap_path("exact");
let _ = std::fs::remove_file(&path);
let h = start_with_capture(Some(path.clone()));
let sid = session_ids(&h.records).into_iter().next().unwrap();
let txs: Vec<&J> = transactions(&h.records, &sid).into_iter().take(6).collect();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
let mut sent = Vec::new();
let mut received = Vec::new();
for tx in &txs {
let req = unhex(&st(tx, "request_hex"));
stream.write_all(&req).unwrap();
stream.flush().unwrap();
sent.push(req);
for _ in 0..tx["responses"].as_array().unwrap().len() {
received.push(read_frame(&mut stream, &mut buf).expect("frame").encode());
}
}
drop(stream);
std::thread::sleep(std::time::Duration::from_millis(200));
let recs = capture::read_file(&path).expect("capture parses");
let rx: Vec<Vec<u8>> = recs
.iter()
.filter(|r| r.dir == capture::Direction::Rx)
.map(|r| r.frame.clone())
.collect();
let tx_caps: Vec<Vec<u8>> = recs
.iter()
.filter(|r| r.dir == capture::Direction::Tx)
.map(|r| r.frame.clone())
.collect();
assert_eq!(
rx, sent,
"captured RX must equal the exact bytes sent by the client"
);
assert_eq!(
tx_caps, received,
"captured TX must equal the exact bytes the client received"
);
// Total order is monotonic, and RX precedes its own TX.
let seqs: Vec<u64> = recs.iter().map(|r| r.seq).collect();
let mut sorted = seqs.clone();
sorted.sort();
assert_eq!(seqs, sorted, "records are written in order");
assert_eq!(recs[0].dir, capture::Direction::Rx);
let _ = std::fs::remove_file(&path);
}
#[test]
fn capture_does_not_alter_the_wire_output() {
// The same conversation with and without capture must produce identical
// response bytes. Capture must be observation, not participation.
let without = converse(&start(), 6);
let path = cap_path("noalter");
let _ = std::fs::remove_file(&path);
let with = converse(&start_with_capture(Some(path.clone())), 6);
assert_eq!(
with.iter().map(|f| hex(f)).collect::<Vec<_>>(),
without.iter().map(|f| hex(f)).collect::<Vec<_>>(),
"enabling capture changed the wire output"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn fragmented_input_reconstructs_the_same_captured_frames() {
// A frame dribbled one byte at a time must be captured as ONE whole frame,
// identical to the same frame sent in a single write.
let sid_take = 3;
let whole_path = cap_path("whole");
let _ = std::fs::remove_file(&whole_path);
converse(&start_with_capture(Some(whole_path.clone())), sid_take);
std::thread::sleep(std::time::Duration::from_millis(150));
let whole = capture::read_file(&whole_path).unwrap();
let frag_path = cap_path("frag");
let _ = std::fs::remove_file(&frag_path);
let h = start_with_capture(Some(frag_path.clone()));
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
for tx in transactions(&h.records, &sid).into_iter().take(sid_take) {
for byte in unhex(&st(tx, "request_hex")) {
stream.write_all(&[byte]).unwrap();
stream.flush().unwrap();
}
for _ in 0..tx["responses"].as_array().unwrap().len() {
read_frame(&mut stream, &mut buf).expect("frame");
}
}
drop(stream);
std::thread::sleep(std::time::Duration::from_millis(200));
let frag = capture::read_file(&frag_path).unwrap();
let f_whole: Vec<String> = whole.iter().map(|r| hex(&r.frame)).collect();
let f_frag: Vec<String> = frag.iter().map(|r| hex(&r.frame)).collect();
assert_eq!(
f_frag, f_whole,
"fragmentation must not change captured frames"
);
let _ = std::fs::remove_file(&whole_path);
let _ = std::fs::remove_file(&frag_path);
}
#[test]
fn coalesced_input_is_captured_as_separate_frames() {
let path = cap_path("coalesced");
let _ = std::fs::remove_file(&path);
let h = start_with_capture(Some(path.clone()));
let sid = session_ids(&h.records).into_iter().next().unwrap();
let txs: Vec<&J> = transactions(&h.records, &sid).into_iter().take(4).collect();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
let mut blob = Vec::new();
for tx in &txs {
blob.extend_from_slice(&unhex(&st(tx, "request_hex")));
}
stream.write_all(&blob).unwrap(); // ONE write
stream.flush().unwrap();
for tx in &txs {
for _ in 0..tx["responses"].as_array().unwrap().len() {
read_frame(&mut stream, &mut buf).expect("frame");
}
}
drop(stream);
std::thread::sleep(std::time::Duration::from_millis(200));
let recs = capture::read_file(&path).unwrap();
let rx: Vec<Vec<u8>> = recs
.iter()
.filter(|r| r.dir == capture::Direction::Rx)
.map(|r| r.frame.clone())
.collect();
assert_eq!(rx.len(), txs.len(), "one record per frame, not per read");
for (got, tx) in rx.iter().zip(txs.iter()) {
assert_eq!(
hex(got),
st(tx, "request_hex"),
"frame boundaries preserved"
);
}
let _ = std::fs::remove_file(&path);
}
#[test]
fn sanitization_removes_the_session_key_from_a_real_captured_login() {
let path = cap_path("sanitize");
let _ = std::fs::remove_file(&path);
let h = start_with_capture(Some(path.clone()));
let sid = session_ids(&h.records).into_iter().next().unwrap();
let mut stream = connect(&h.addr);
let mut buf = Vec::new();
for tx in transactions(&h.records, &sid) {
stream.write_all(&unhex(&st(tx, "request_hex"))).unwrap();
stream.flush().unwrap();
for _ in 0..tx["responses"].as_array().unwrap().len() {
read_frame(&mut stream, &mut buf).expect("frame");
}
if st(tx, "name") == "login" {
break;
}
}
drop(stream);
std::thread::sleep(std::time::Duration::from_millis(200));
let recs = capture::read_file(&path).unwrap();
let session_key = h
.records
.iter()
.find(|r| r["kind"] == "session" && r["id"] == sid.as_str())
.map(|r| st(r, "session_key"))
.expect("recorded session key");
// The raw capture DOES contain the key — that is what makes it forensic
// evidence, and why it is gitignored.
let raw_has_key = recs
.iter()
.any(|r| String::from_utf8_lossy(&r.frame).contains(&session_key));
assert!(
raw_has_key,
"raw capture should contain the real session key"
);
let (out, report) = sanitize::sanitize(&recs).expect("sanitizes");
assert!(
report.frames_redacted > 0,
"login frames carry a session key"
);
for (_, bytes) in &out {
assert!(
!String::from_utf8_lossy(bytes).contains(&session_key),
"sanitized output still contains the session key"
);
}
// Structure is untouched: same frame count, same sizes, same order.
assert_eq!(out.len(), recs.len());
for ((rec, bytes), orig) in out.iter().zip(recs.iter()) {
assert_eq!(
bytes.len(),
orig.frame.len(),
"frame size must be preserved"
);
assert_eq!(rec.seq, orig.seq, "order must be preserved");
}
// And every redaction is reported, not silent.
assert!(
report.redactions.iter().any(|r| r.path.ends_with("KEY")),
"{:?}",
report.redactions
);
let _ = std::fs::remove_file(&path);
}