From c152a955e5d51c9159aa8e5921cb696473bd521d Mon Sep 17 00:00:00 2001 From: randogoth Date: Tue, 29 Sep 2026 10:27:44 +0300 Subject: [PATCH] feat: add hostile harness, multi-session driver, and frame header deadline --- src/channel.rs | 14 +- src/harness.rs | 617 +++++++++++++++++++++++++++++++++++++++++++++++++ src/main.rs | 6 + src/proto.rs | 5 + src/server.rs | 11 +- 5 files changed, 649 insertions(+), 4 deletions(-) create mode 100644 src/harness.rs diff --git a/src/channel.rs b/src/channel.rs index fd67c93..bee9074 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -9,7 +9,10 @@ use std::net::TcpStream; use snow::{Builder, TransportState}; -use crate::proto::{ProtocolError, MAX_FRAME, NOISE_PARAMS, NOISE_PAYLOAD, PROLOGUE}; +use crate::proto::{ + ProtocolError, HEADER_TIMEOUT_SECS, IDLE_TIMEOUT_SECS, MAX_FRAME, NOISE_PARAMS, NOISE_PAYLOAD, + PROLOGUE, +}; /// Runs the Noise_NX responder handshake. The initiator stays anonymous; /// only we hold a static key. Returns the transport state and the @@ -72,6 +75,15 @@ impl Channel { } fn read_noise(&mut self) -> anyhow::Result> { + // A frame already started must finish quickly; only a client idle + // between frames earns the full session idle timeout. + let deadline = if self.buf.is_empty() { + IDLE_TIMEOUT_SECS + } else { + HEADER_TIMEOUT_SECS + }; + self.stream + .set_read_timeout(Some(std::time::Duration::from_secs(deadline)))?; let len = read_u16_len(&mut self.stream)?; let mut ciphertext = vec![0u8; len]; read_exact_into(&mut self.stream, &mut ciphertext)?; diff --git a/src/harness.rs b/src/harness.rs new file mode 100644 index 0000000..a8959a3 --- /dev/null +++ b/src/harness.rs @@ -0,0 +1,617 @@ +//! Hostile transport harness (RNS.md sec 9 roster): a client that completes +//! the real Noise_NX handshake and then deliberately misbehaves — garbage and +//! truncated handshakes, malformed frames, oversized payloads, absurd +//! cursors, mid-frame drops — against the same accept/handshake/session path +//! the TCP carrier serves. Every scenario ends by asserting the server's +//! response class and that a fresh connection still completes a handshake: +//! the carrier must die never, the session must die cleanly. +//! +//! Transport-level kills (mid-link, mid-spin, mid-resource-transfer) belong +//! to the RNS carrier and were live-fired against the deployed build during +//! the stress campaign; they are recorded in MAIL.md, not replayed here. + +#![cfg(test)] + +use std::io::{Read, Write}; +use std::net::TcpStream; + +use ed25519_dalek::{Signer, SigningKey}; +use snow::{Builder, TransportState}; + +use crate::crypto::generate_static_key; +use crate::proto::{ + AUTH_REQUIRED, ENVELOPE_MAGIC, ENVELOPE_VERSION, ID_LEN, LABEL_AUTH, LABEL_REGISTER, + MALFORMED, MAX_FRAME, NOISE_PARAMS, NOISE_PAYLOAD, OP_AUTH, OP_FETCH, OP_REGISTER, OP_RESOLVE, + OP_SEND, PROLOGUE, TOO_LARGE, +}; +use crate::ratelimit::RateLimiter; +use crate::server::handle_connection; +use crate::session::ServerConfig; +use crate::store::Store; + +/// One harness target: a real accept loop serving `handle_connection` on an +/// ephemeral port, one thread per connection, exactly like `server::run`. +struct Target { + port: u16, + static_public: [u8; 32], +} + +impl Target { + fn start(max_envelope: usize) -> Self { + let (static_private, static_public) = generate_static_key(); + let config: &'static ServerConfig = Box::leak(Box::new(ServerConfig { + max_envelope, + fetch_budget: 512 << 10, + main_quota: 8 << 20, + requests_quota: 2 << 20, + max_tokens: 16, + invite_token: None, + // Unlimited: the harness is hostile by content, not by volume. + conn_limiter: RateLimiter::new(0), + send_limiter: RateLimiter::new(0), + token_limiter: RateLimiter::new(0), + })); + static TARGETS: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0); + let seq = TARGETS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + let db_path = std::env::temp_dir() + .join(format!("bunshin-harness-{}-{seq}.db", std::process::id())) + .to_string_lossy() + .into_owned(); + let _ = std::fs::remove_file(&db_path); + Store::open(&db_path).unwrap(); + + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + std::thread::spawn(move || { + for stream in listener.incoming().flatten() { + let config = config; + let db_path = db_path.clone(); + std::thread::spawn(move || { + // The harness deliberately produces hostile connections; + // a panic in one must not take the accept loop with it. + let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + handle_connection( + stream, + config, + &static_private, + &db_path, + "127.0.0.1", + ) + })); + }); + } + }); + Target { + port, + static_public, + } + } + + /// Proof of life: a fresh connection completes the handshake and gets a + /// protocol answer (any status — UNKNOWN_USER counts; dead is no answer). + fn assert_alive(&self) { + let mut client = Client::connect(self.port); + client.handshake().expect("server survived"); + let (status, _) = client.request(OP_RESOLVE, &[7, b'h', b'a', b'r', b'n', b'e', b's', b's']); + assert_ne!(status, 0xFF, "server answered with a status byte"); + } +} + +/// The hostile client: a real Noise_NX initiator with raw-stream access so +/// scenarios can write garbage, truncated, or oversized wire data at will. +struct Client { + stream: TcpStream, + transport: Option, + handshake_hash: Vec, +} + +impl Client { + fn connect(port: u16) -> Self { + Client { + stream: TcpStream::connect(("127.0.0.1", port)).unwrap(), + transport: None, + handshake_hash: Vec::new(), + } + } + + /// The genuine client side of the server's `channel::handshake`. + fn handshake(&mut self) -> anyhow::Result<()> { + let params: snow::params::NoiseParams = NOISE_PARAMS.parse()?; + let mut noise = Builder::new(params).prologue(PROLOGUE).build_initiator()?; + let mut buf = [0u8; 65535]; + + let n = noise.write_message(&[], &mut buf)?; + self.stream.write_all(&(n as u16).to_be_bytes())?; + self.stream.write_all(&buf[..n])?; + + let mut len_buf = [0u8; 2]; + self.stream.read_exact(&mut len_buf)?; + let len = u16::from_be_bytes(len_buf) as usize; + let mut msg = vec![0u8; len]; + self.stream.read_exact(&mut msg)?; + let mut plaintext = [0u8; 65535]; + noise.read_message(&msg, &mut plaintext)?; + + anyhow::ensure!(noise.is_handshake_finished(), "handshake did not complete"); + self.handshake_hash = noise.get_handshake_hash().to_vec(); + self.transport = Some(noise.into_transport_mode()?); + Ok(()) + } + + fn write_noise(&mut self, payload: &[u8]) { + let transport = self.transport.as_mut().expect("handshake first"); + let mut packet = vec![0u8; payload.len() + 16]; + let n = transport.write_message(payload, &mut packet).unwrap(); + packet.truncate(n); + self.stream.write_all(&(packet.len() as u16).to_be_bytes()).unwrap(); + self.stream.write_all(&packet).unwrap(); + } + + fn read_noise(&mut self) -> Vec { + let transport = self.transport.as_mut().expect("handshake first"); + let mut len_buf = [0u8; 2]; + self.stream.read_exact(&mut len_buf).unwrap(); + let len = u16::from_be_bytes(len_buf) as usize; + let mut ciphertext = vec![0u8; len]; + self.stream.read_exact(&mut ciphertext).unwrap(); + let mut plaintext = vec![0u8; len]; + let n = transport.read_message(&ciphertext, &mut plaintext).unwrap(); + plaintext.truncate(n); + plaintext + } + + /// A well-formed frame — the baseline the hostile scenarios deviate from. + fn request(&mut self, op: u8, body: &[u8]) -> (u8, Vec) { + let mut frame = Vec::with_capacity(5 + body.len()); + frame.extend_from_slice(&((1 + body.len()) as u32).to_be_bytes()); + frame.push(op); + frame.extend_from_slice(body); + for chunk in frame.chunks(NOISE_PAYLOAD) { + self.write_noise(chunk); + } + self.read_frame() + } + + fn read_frame(&mut self) -> (u8, Vec) { + let mut buf = Vec::new(); + while buf.len() < 5 { + buf.extend_from_slice(&self.read_noise()); + } + let length = u32::from_be_bytes(buf[..4].try_into().unwrap()) as usize; + while buf.len() < 4 + length { + buf.extend_from_slice(&self.read_noise()); + } + let frame: Vec = buf[4..4 + length].to_vec(); + (frame[0], frame[1..].to_vec()) + } + + /// Raw bytes before the handshake — garbage on the wire. + fn write_raw(&mut self, bytes: &[u8]) { + self.stream.write_all(bytes).unwrap(); + } + + /// Claims a Noise message of `claimed` bytes but sends only `sent`, then + /// keeps the connection open briefly so the server sees a short read. + fn truncated_noise_len(&mut self, claimed: usize, sent: &[u8]) { + self.stream + .write_all(&(claimed as u16).to_be_bytes()) + .unwrap(); + self.stream.write_all(sent).unwrap(); + } + + /// One raw post-handshake Noise message carrying arbitrary plaintext. + fn raw_frame_bytes(&mut self, bytes: &[u8]) { + self.write_noise(bytes); + } +} + +/// SERVER-QUIET (the inverse): the server must never answer a connection +/// that sends nothing, and dropping it must cost nothing. +#[test] +fn silent_connection_costs_nothing_and_server_survives() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + std::thread::sleep(std::time::Duration::from_millis(50)); + drop(client); + target.assert_alive(); +} + +#[test] +fn short_garbage_handshake_is_refused_and_server_survives() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + // Below the 32-byte ephemeral an NX first message needs, so snow must + // reject it outright and the server closes without writing anything. + client.write_raw(&(10u16).to_be_bytes()); + client.write_raw(&[0xA5; 10]); + let mut sink = [0u8; 16]; + assert!(client.stream.read(&mut sink).unwrap_or(0) == 0, "server closed"); + target.assert_alive(); +} + +/// Noise_NX authenticates the server in message two; the initiator's first +/// message is unauthenticated by design, so 32+ bytes of random data can +/// parse as an ephemeral key and draw a handshake reply. That reply must be +/// the most a garbage sender ever gets: no session forms from it. +#[test] +fn long_garbage_first_message_never_forms_a_session() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.write_raw(&(100u16).to_be_bytes()); + client.write_raw(&[0xA5; 100]); + // Reply or close, either is protocol-legal; drop without completing. + drop(client); + target.assert_alive(); +} + +#[test] +fn truncated_handshake_is_survived() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.truncated_noise_len(1000, &[1, 2, 3, 4, 5, 6, 7, 8, 9, 10]); + drop(client); + target.assert_alive(); +} + +#[test] +fn unknown_op_gets_malformed_and_server_survives() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + let (op, body) = client.request(0x99, b""); + assert_eq!(op, 0x99); + assert_eq!(body[0], MALFORMED); + target.assert_alive(); +} + +#[test] +fn fetch_without_auth_gets_auth_required() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + let (op, body) = client.request(OP_FETCH, &[0u8; 40]); + assert_eq!(op, OP_FETCH); + assert_eq!(body[0], AUTH_REQUIRED); + target.assert_alive(); +} + +#[test] +fn garbage_resolve_body_gets_malformed() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + let (op, body) = client.request(OP_RESOLVE, &[0xEE; 50]); + assert_eq!(op, OP_RESOLVE); + assert_eq!(body[0], MALFORMED); + target.assert_alive(); +} + +#[test] +fn zero_length_frame_gets_malformed_then_close() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + // Four length bytes alone never reach the length check — read_frame + // buffers five bytes first — so the hostile frame is length + op byte. + client.raw_frame_bytes(&[0, 0, 0, 0, 0x00]); + let (op, body) = client.read_frame(); + assert_eq!(op, 0); + assert_eq!(body[0], MALFORMED); + let mut sink = [0u8; 16]; + assert!(client.stream.read(&mut sink).unwrap_or(0) == 0, "closed after bad frame"); + target.assert_alive(); +} + +#[test] +fn over_max_frame_length_gets_malformed_then_close() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + let claimed = (MAX_FRAME + 1) as u32; + let mut frame = claimed.to_be_bytes().to_vec(); + frame.push(0x00); + client.raw_frame_bytes(&frame); + let (op, body) = client.read_frame(); + assert_eq!(op, 0); + assert_eq!(body[0], MALFORMED); + target.assert_alive(); +} + +#[test] +fn oversized_envelope_gets_too_large() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + // Well-formed frame, well-formed envelope header, size past the cap. + let mut envelope = Vec::with_capacity(800 << 10); + envelope.extend_from_slice(ENVELOPE_MAGIC); + envelope.push(ENVELOPE_VERSION); + envelope.extend_from_slice(&[0u8; ID_LEN]); // recipient never looked up + envelope.resize(800 << 10, 0x41); + let mut body = vec![0u8]; // no accept-token MAC + body.extend_from_slice(&envelope); + let (op, resp) = client.request(OP_SEND, &body); + let body = &resp[..]; + assert_eq!(op, OP_SEND); + assert_eq!(body[0], TOO_LARGE); + target.assert_alive(); +} + +#[test] +fn bad_envelope_magic_gets_malformed() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + let mut envelope = vec![0x00; 128]; + envelope[4] = ENVELOPE_VERSION; + let (op, body) = client.request(OP_SEND, &envelope); + assert_eq!(op, OP_SEND); + assert_eq!(body[0], MALFORMED); + target.assert_alive(); +} + +#[test] +fn mid_frame_drop_is_survived() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + // Claim a frame, deliver a third of its header bytes, vanish. + client.raw_frame_bytes(&[0x00, 0x10, 0x00]); + drop(client); + target.assert_alive(); +} + +/// The well-behaved baseline plus the out-of-order cursor: a FETCH whose +/// cursor points past the end of time must answer an empty page, not an +/// error or a wrap-around into someone else's mail. +#[test] +fn registered_flow_survives_absurd_future_cursor() { + let target = Target::start(768 << 10); + let key = SigningKey::from_bytes(&[7u8; 32]); + let verifying = key.verifying_key(); + let identity: &[u8] = verifying.as_bytes(); + + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + + // REGISTER with a real proof of possession over the server's static key. + let mut pop = LABEL_REGISTER.to_vec(); + pop.extend_from_slice(&target.static_public); + pop.extend_from_slice(b"harness"); + pop.extend_from_slice(identity); + let mut body = vec![7]; + body.extend_from_slice(b"harness"); + body.extend_from_slice(identity); + body.extend_from_slice(&key.sign(&pop).to_bytes()); + body.push(0); // no invite token + body.push(0); // plain registration + let (op, resp) = client.request(OP_REGISTER, &body); + assert_eq!((op, resp.as_slice()), (OP_REGISTER, &[0u8][..])); + + // AUTH over the real handshake hash this session actually negotiated. + let mut auth_msg = LABEL_AUTH.to_vec(); + auth_msg.extend_from_slice(&client.handshake_hash); + let mut body = vec![7]; + body.extend_from_slice(b"harness"); + body.extend_from_slice(identity); + body.extend_from_slice(&key.sign(&auth_msg).to_bytes()); + body.push(0); // sync = 0 + body.extend_from_slice(&0u16.to_be_bytes()); + let (op, resp) = client.request(OP_AUTH, &body); + assert_eq!(op, OP_AUTH); + assert_eq!(resp[0], 0, "status OK"); + assert_eq!( + u16::from_be_bytes(resp[1..3].try_into().unwrap()), + 0, + "zero accepted tokens" + ); + + // FETCH with a cursor at the end of time: OK, zero messages. + let mut body = Vec::new(); + body.extend_from_slice(&i64::MAX.to_be_bytes()); + body.extend_from_slice(&[0xFF; ID_LEN]); + let (op, resp) = client.request(OP_FETCH, &body); + assert_eq!(op, OP_FETCH); + assert_eq!(resp[0], 0, "status OK"); + assert_eq!(u16::from_be_bytes(resp[1..3].try_into().unwrap()), 0, "empty page"); + + // And the registered name resolves — the store came through intact. + let (op, resp) = client.request(OP_RESOLVE, &[7, b'h', b'a', b'r', b'n', b'e', b's', b's']); + assert_eq!(op, OP_RESOLVE); + assert_eq!(resp[0], 0); + target.assert_alive(); +} + + +// --------------------------------------------------------------------------- +// Long-lived multi-session driver: the last agreed stress-campaign phase. +// A dozen authenticated sessions live on one server simultaneously, run +// concurrent FETCH cursors against one mailbox, and — in the slow scenario — +// are held past the session idle timeout to verify the server reaps idle +// sessions cleanly instead of leaking their threads. + +impl Client { + /// AUTH as an already-registered identity (the driver registers once). + fn auth(&mut self, username: &str, key: &SigningKey) { + let verifying = key.verifying_key(); + let identity: &[u8] = verifying.as_bytes(); + let mut auth_msg = LABEL_AUTH.to_vec(); + auth_msg.extend_from_slice(&self.handshake_hash); + let mut body = vec![username.len() as u8]; + body.extend_from_slice(username.as_bytes()); + body.extend_from_slice(identity); + body.extend_from_slice(&key.sign(&auth_msg).to_bytes()); + body.push(0); + body.extend_from_slice(&0u16.to_be_bytes()); + let (op, resp) = self.request(OP_AUTH, &body); + assert_eq!(op, OP_AUTH); + assert_eq!(resp[0], 0, "AUTH ok"); + } + + /// Pages FETCH from `cursor` to the end, returning every message id. + fn fetch_all(&mut self, cursor: (i64, [u8; 32])) -> Vec<[u8; 32]> { + let (mut after_at, mut after_id) = cursor; + let mut ids = Vec::new(); + loop { + let mut body = Vec::new(); + body.extend_from_slice(&after_at.to_be_bytes()); + body.extend_from_slice(&after_id); + let (op, resp) = self.request(OP_FETCH, &body); + assert_eq!(op, OP_FETCH); + assert_eq!(resp[0], 0, "FETCH ok"); + let count = u16::from_be_bytes(resp[1..3].try_into().unwrap()); + if count == 0 { + return ids; + } + let mut off = 3; + for _ in 0..count { + let id: [u8; 32] = resp[off..off + 32].try_into().unwrap(); + off += 32; + let received_at = i64::from_be_bytes(resp[off..off + 8].try_into().unwrap()); + off += 8; + off += 1; // unsolicited flag + let len = u32::from_be_bytes(resp[off..off + 4].try_into().unwrap()) as usize; + off += 4 + len; + ids.push(id); + // The cursor is (received_at, id) strictly after the last + // delivered record; a future-dated probe terminates cleanly. + after_at = received_at; + after_id = id; + } + } + } +} + +/// The driver's mailbox: an envelope addressed to a registered identity. +fn envelope_to(recipient: &[u8], payload_len: usize) -> Vec { + let mut envelope = Vec::with_capacity(69 + payload_len + 16); + envelope.extend_from_slice(ENVELOPE_MAGIC); + envelope.push(ENVELOPE_VERSION); + envelope.extend_from_slice(recipient); + envelope.resize(69 + payload_len, 0x42); // epk + ciphertext + tag + envelope +} + +#[test] +fn dozen_sessions_run_concurrent_fetch_cursors() { + let target = Target::start(768 << 10); + let key = SigningKey::from_bytes(&[8u8; 32]); + let verifying = key.verifying_key(); + let identity: &[u8] = verifying.as_bytes(); + + // Register the driver identity from a first session. + let mut reg = Client::connect(target.port); + reg.handshake().unwrap(); + let mut pop = LABEL_REGISTER.to_vec(); + pop.extend_from_slice(&target.static_public); + pop.extend_from_slice(b"driver"); + pop.extend_from_slice(identity); + let mut body = vec![6]; + body.extend_from_slice(b"driver"); + body.extend_from_slice(identity); + body.extend_from_slice(&key.sign(&pop).to_bytes()); + body.push(0); + body.push(0); + let (op, resp) = reg.request(OP_REGISTER, &body); + assert_eq!((op, resp.as_slice()), (OP_REGISTER, &[0u8][..])); + + // Seed 20 envelopes for the driver mailbox; ids are deterministic. + let mut expected = std::collections::HashSet::new(); + for i in 0..20u8 { + let envelope = envelope_to(identity, 40 + i as usize); + let mut body = vec![0u8]; // no accept-token MAC + body.extend_from_slice(&envelope); + let (op, resp) = reg.request(OP_SEND, &body); + assert_eq!((op, resp[0]), (OP_SEND, 0), "seed send {i}"); + let id: [u8; 32] = resp[1..33].try_into().unwrap(); + expected.insert(id); + } + + // A dozen sessions, each authenticated and fetching the whole mailbox + // concurrently; all must see the same 20 ids (nothing is deleted, so + // every session's cursor walk is independent and complete). + let mut handles = Vec::new(); + for _ in 0..12 { + let port = target.port; + let session_key = key.clone(); + handles.push(std::thread::spawn(move || { + let mut client = Client::connect(port); + client.handshake().unwrap(); + client.auth("driver", &session_key); + client.fetch_all((0, [0u8; 32])) + })); + } + for handle in handles { + let ids = handle.join().unwrap(); + let got: std::collections::HashSet<[u8; 32]> = ids.into_iter().collect(); + assert_eq!(got, expected, "every session saw the whole mailbox"); + } + target.assert_alive(); +} + +/// Held-open sessions must be closed by the server's idle timeout, not +/// leaked. Slow by nature: it waits out the full IDLE_TIMEOUT_SECS. +#[test] +#[ignore = "holds a dozen sessions past the 120 s idle timeout"] +fn idle_sessions_are_reaped_by_the_server() { + let target = Target::start(768 << 10); + let key = SigningKey::from_bytes(&[9u8; 32]); + let verifying = key.verifying_key(); + let identity: &[u8] = verifying.as_bytes(); + + let mut reg = Client::connect(target.port); + reg.handshake().unwrap(); + let mut pop = LABEL_REGISTER.to_vec(); + pop.extend_from_slice(&target.static_public); + pop.extend_from_slice(b"idler"); + pop.extend_from_slice(identity); + let mut body = vec![5]; + body.extend_from_slice(b"idler"); + body.extend_from_slice(identity); + body.extend_from_slice(&key.sign(&pop).to_bytes()); + body.push(0); + body.push(0); + let (op, resp) = reg.request(OP_REGISTER, &body); + assert_eq!((op, resp.as_slice()), (OP_REGISTER, &[0u8][..])); + + // A dozen sessions authenticate, then say nothing at all. + let mut clients = Vec::new(); + for _ in 0..12 { + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + client.auth("idler", &key); + clients.push(client); + } + + // Past the idle timeout the server closes them; the client sees EOF. + std::thread::sleep(std::time::Duration::from_secs(125)); + for client in &mut clients { + let mut sink = [0u8; 16]; + assert!( + client.stream.read(&mut sink).unwrap_or(0) == 0, + "server closed the idle session" + ); + } + target.assert_alive(); +} + +/// A frame started but not finished must not hold a thread for the full +/// session idle timeout: the header deadline (30 s) closes it first. +#[test] +#[ignore = "waits out the 30 s header deadline"] +fn partial_frame_is_dropped_after_the_header_deadline() { + let target = Target::start(768 << 10); + let mut client = Client::connect(target.port); + client.handshake().unwrap(); + // Three bytes of a frame header, then silence. Before this hardening + // the connection idled the full 120 s; the header deadline must close + // it well before that. + client.raw_frame_bytes(&[0x00, 0x10, 0x00]); + let mut sink = [0u8; 16]; + assert!( + client.stream.read(&mut sink).unwrap_or(0) == 0, + "server dropped the stalled frame" + ); + target.assert_alive(); +} diff --git a/src/main.rs b/src/main.rs index 34ab05c..5ec7d66 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,6 +3,8 @@ mod bind; mod channel; mod crypto; +#[cfg(test)] +mod harness; mod proto; mod ratelimit; #[cfg(feature = "rns")] @@ -101,6 +103,9 @@ struct RnsServeArgs { rns_rate_link_requests: u32, #[arg(long = "rns-rate-link-bytes", default_value_t = 1 << 20)] rns_rate_link_bytes: u64, + /// Reap RNS links idle for this many seconds (0 disables) + #[arg(long = "rns-link-idle", default_value_t = 300)] + rns_link_idle: u64, /// UDP interface to listen on, host[:port] #[arg(long = "rns-udp", default_value = "127.0.0.1:4242")] rns_udp: String, @@ -200,6 +205,7 @@ fn main() -> anyhow::Result<()> { max_links: rns.rns_max_links, rate_link_requests: rns.rns_rate_link_requests, rate_link_bytes: rns.rns_rate_link_bytes, + link_idle_secs: rns.rns_link_idle, udp_listen_host: listen_host, udp_listen_port: listen_port, udp_forward_host: forward_host, diff --git a/src/proto.rs b/src/proto.rs index f080386..d5c841b 100644 --- a/src/proto.rs +++ b/src/proto.rs @@ -51,6 +51,11 @@ pub const FETCH_BUDGET: usize = 512 * 1024; // must stay under MAX_FRAME #[cfg(feature = "rns")] pub const AUTH_FULL_TOKENS: usize = 164; pub const IDLE_TIMEOUT_SECS: u64 = 120; +/// Budget for completing a handshake or a frame already started. Without it +/// a connection that sends a partial frame header and goes silent holds a +/// thread for the full idle timeout — bounded by the connection limiter, +/// but a free hold-and-wait for the hostile kind of client. +pub const HEADER_TIMEOUT_SECS: u64 = 30; pub const PURGE_INTERVAL_SECS: u64 = 60; /// A peer sent something unparseable. Always answered with MALFORMED. diff --git a/src/server.rs b/src/server.rs index 1fdb5ef..50389ca 100644 --- a/src/server.rs +++ b/src/server.rs @@ -7,7 +7,7 @@ use std::time::Duration; use crate::bind::TransportBindValues; use crate::channel::{handshake, Channel}; use crate::crypto::{b32, derive_public}; -use crate::proto::{FETCH_BUDGET, IDLE_TIMEOUT_SECS, KEY_LEN, MALFORMED, PURGE_INTERVAL_SECS}; +use crate::proto::{FETCH_BUDGET, HEADER_TIMEOUT_SECS, KEY_LEN, MALFORMED, PURGE_INTERVAL_SECS}; use crate::ratelimit::RateLimiter; use crate::session::{ServerConfig, Session}; use crate::store::Store; @@ -94,14 +94,19 @@ pub fn run(args: ServeArgs) -> anyhow::Result<()> { Ok(()) } -fn handle_connection( +/// Also the hostile harness's entry point: it drives real connections +/// through the same accept/handshake/session path the TCP carrier serves. +pub(crate) fn handle_connection( mut stream: TcpStream, config: &ServerConfig, static_key: &[u8; KEY_LEN], db_path: &str, peer_ip: &str, ) -> anyhow::Result<()> { - stream.set_read_timeout(Some(Duration::from_secs(IDLE_TIMEOUT_SECS)))?; + // The handshake runs under the header budget; once it completes, + // Channel::read_noise picks between the header and idle deadlines per + // frame state. + stream.set_read_timeout(Some(Duration::from_secs(HEADER_TIMEOUT_SECS)))?; let handshake_started = std::time::Instant::now(); let (transport, handshake_hash) = match handshake(&mut stream, static_key) {