feat: add hostile harness, multi-session driver, and frame header deadline
This commit is contained in:
parent
8ea6865e7e
commit
c152a955e5
5 changed files with 649 additions and 4 deletions
|
|
@ -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<Vec<u8>> {
|
||||
// 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)?;
|
||||
|
|
|
|||
617
src/harness.rs
Normal file
617
src/harness.rs
Normal file
|
|
@ -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<TransportState>,
|
||||
handshake_hash: Vec<u8>,
|
||||
}
|
||||
|
||||
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<u8> {
|
||||
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<u8>) {
|
||||
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<u8>) {
|
||||
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<u8> = 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<u8> {
|
||||
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();
|
||||
}
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue