//! 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 { 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 { (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, } /// 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) -> 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::loopback() }, 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 = records .iter() .filter(|r| r["kind"] == "session") .map(|r| st(r, "session_key")) .collect(); let queue = Mutex::new(keys.into_iter().collect::>()); 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) -> Option { 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 { 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 = 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> { 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> = recs .iter() .filter(|r| r.dir == capture::Direction::Rx) .map(|r| r.frame.clone()) .collect(); let tx_caps: Vec> = 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 = 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::>(), without.iter().map(|f| hex(f)).collect::>(), "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 = whole.iter().map(|r| hex(&r.frame)).collect(); let f_frag: Vec = 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> = 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); }