diff --git a/protocols/turnloop-http/src/asynchronous/client.rs b/protocols/turnloop-http/src/asynchronous/client.rs index 31ec64e..8c7edc3 100644 --- a/protocols/turnloop-http/src/asynchronous/client.rs +++ b/protocols/turnloop-http/src/asynchronous/client.rs @@ -445,7 +445,9 @@ impl Client { http2::Event::Reset { stream, .. } if stream == id => { return Err(io::Error::other("HTTP/2 stream reset")); } - http2::Event::Goaway { last_stream, code } => { + http2::Event::Goaway { + last_stream, code, .. + } => { peer_draining = true; if code != 0 || last_stream < id { return Err(io::Error::other("HTTP/2 GOAWAY rejected request")); @@ -540,7 +542,9 @@ impl Client { http2::Event::Reset { stream, .. } if stream == id => { return Err(io::Error::other("HTTP/2 stream reset")); } - http2::Event::Goaway { last_stream, code } => { + http2::Event::Goaway { + last_stream, code, .. + } => { peer_draining = true; if code != 0 || last_stream < id { return Err(io::Error::other( diff --git a/protocols/turnloop-http/src/http2.rs b/protocols/turnloop-http/src/http2.rs index 8000085..f5b1d02 100644 --- a/protocols/turnloop-http/src/http2.rs +++ b/protocols/turnloop-http/src/http2.rs @@ -1,7 +1,7 @@ //! RFC 9113 connection/stream framing. Caller retains partial frames and acknowledges //! writes. DATA is borrowed, flow-control credit is returned explicitly by the host. use crate::{Error, Result, hpack, http1::Header}; -use std::time::Instant; +use std::{collections::VecDeque, time::Instant}; pub const PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Role { @@ -27,6 +27,127 @@ impl Default for Limits { } } } +/// SETTINGS parameters this endpoint sends: at most one value per identifier, +/// kept and serialised in ascending identifier order (the order Node writes +/// them in). An empty `Settings` is a valid, empty SETTINGS frame. +/// +/// Identifiers this crate does not know are sent as given, so a host can +/// advertise extension settings; the named constants are the ones RFC 9113 +/// section 6.5.2 and RFC 8441 define. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct Settings { + entries: Vec<(u16, u32)>, +} +impl Settings { + pub const HEADER_TABLE_SIZE: u16 = 1; + pub const ENABLE_PUSH: u16 = 2; + pub const MAX_CONCURRENT_STREAMS: u16 = 3; + pub const INITIAL_WINDOW_SIZE: u16 = 4; + pub const MAX_FRAME_SIZE: u16 = 5; + pub const MAX_HEADER_LIST_SIZE: u16 = 6; + pub const ENABLE_CONNECT_PROTOCOL: u16 = 8; + pub fn new() -> Self { + Self::default() + } + /// The three limits [`Connection::new`] advertises. + pub fn from_limits(limits: &Limits) -> Self { + let mut settings = Self::new(); + settings + .set(Self::MAX_CONCURRENT_STREAMS, limits.streams as u32) + .set(Self::MAX_FRAME_SIZE, limits.frame_size as u32) + .set(Self::MAX_HEADER_LIST_SIZE, limits.header_list as u32); + settings + } + /// Set `id` to `value`, replacing any earlier value for it. + pub fn set(&mut self, id: u16, value: u32) -> &mut Self { + match self.entries.binary_search_by_key(&id, |e| e.0) { + Ok(i) => self.entries[i].1 = value, + Err(i) => self.entries.insert(i, (id, value)), + } + self + } + pub fn remove(&mut self, id: u16) -> Option { + let i = self.entries.binary_search_by_key(&id, |e| e.0).ok()?; + Some(self.entries.remove(i).1) + } + pub fn get(&self, id: u16) -> Option { + let i = self.entries.binary_search_by_key(&id, |e| e.0).ok()?; + Some(self.entries[i].1) + } + /// `(identifier, value)` pairs in ascending identifier order. + pub fn iter(&self) -> impl Iterator + '_ { + self.entries.iter().copied() + } + pub fn len(&self) -> usize { + self.entries.len() + } + pub fn is_empty(&self) -> bool { + self.entries.is_empty() + } + fn merge(&mut self, from: impl Iterator) { + for (id, value) in from { + self.set(id, value); + } + } + fn encode(&self) -> Vec { + let mut payload = Vec::with_capacity(self.entries.len() * 6); + for (id, value) in self.iter() { + payload.extend_from_slice(&id.to_be_bytes()); + payload.extend_from_slice(&value.to_be_bytes()); + } + payload + } + /// Refuse what this endpoint cannot honour once the peer acts on it. + fn validate(&self) -> Result<()> { + for (id, value) in self.iter() { + let valid = match id { + // The HPACK decoder's dynamic table is fixed at 4096 octets; + // advertising more would let the peer's encoder outgrow it. + Self::HEADER_TABLE_SIZE => value <= 4096, + // Server push is not implemented: PUSH_PROMISE is always a + // connection error, so neither side may invite one. + Self::ENABLE_PUSH => value == 0, + Self::INITIAL_WINDOW_SIZE => value <= 0x7fffffff, + Self::MAX_FRAME_SIZE => (16384..=0xffffff).contains(&value), + Self::ENABLE_CONNECT_PROTOCOL => value <= 1, + _ => true, + }; + if !valid { + return Err(protocol("invalid SETTINGS value")); + } + } + Ok(()) + } +} +/// The parameters of one SETTINGS frame the peer sent, borrowed from the input +/// exactly as they were on the wire: in frame order, duplicates and unknown +/// identifiers included. RFC 9113 section 6.5.3 processes them in that order, +/// so for a repeated identifier the last value is the one in force. +#[derive(Debug, Clone, Copy)] +pub struct SettingsFrame<'a> { + payload: &'a [u8], +} +impl<'a> SettingsFrame<'a> { + pub fn iter(&self) -> impl Iterator + 'a { + self.payload + .as_chunks::<6>() + .0 + .iter() + .map(|s| (u16::from_be_bytes([s[0], s[1]]), u32be(&s[2..]))) + } + /// The value in force for `id` after this frame, if the frame carried it. + pub fn get(&self, id: u16) -> Option { + self.iter().filter(|e| e.0 == id).last().map(|e| e.1) + } + pub fn is_empty(&self) -> bool { + self.payload.is_empty() + } + pub fn to_settings(&self) -> Settings { + let mut settings = Settings::new(); + settings.merge(self.iter()); + settings + } +} #[derive(Debug, Clone, Copy)] pub struct Frame<'a> { pub kind: u8, @@ -96,7 +217,16 @@ pub enum HeadersKind { } #[derive(Debug)] pub enum Event<'a> { - Settings, + /// The peer's SETTINGS, already applied and acknowledged. + Settings(SettingsFrame<'a>), + /// The peer acknowledged the oldest SETTINGS frame this endpoint still had + /// outstanding; its parameters are returned here and are now in force. + /// Acknowledgements arrive in the order the frames were sent (RFC 9113 + /// section 6.5.3), and the first is always for the frame the constructor + /// queued, so a host that timestamps construction and each + /// [`Connection::settings`] call can pair them in FIFO order to measure a + /// round trip. No clock is read here. + SettingsAck(Settings), Headers { stream: u32, headers: Vec
, @@ -136,6 +266,11 @@ pub enum Event<'a> { Goaway { last_stream: u32, code: u32, + /// The opaque debug data (RFC 9113 section 6.8). `None` when the frame + /// carried none - Node reports `undefined` there, not an empty buffer. + /// The wire cannot express "present but empty", so this is never + /// `Some` of an empty slice. + debug: Option<&'a [u8]>, }, Ping { ack: bool, @@ -153,7 +288,7 @@ pub enum Event<'a> { /// | `consumed` | `event` | meaning | /// |---|---|---| /// | `0` | `None` | **stop.** A partial preface or a partial frame; read more input before calling again. | -/// | `> 0` | `None` | **keep going.** Progress with nothing for the host: the client preface, a SETTINGS acknowledgement, PRIORITY, an unknown frame type, or a frame the peer had in flight for a stream that is already gone. | +/// | `> 0` | `None` | **keep going.** Progress with nothing for the host: the client preface, PRIORITY, an unknown frame type, or a frame the peer had in flight for a stream that is already gone. | /// | `> 0` | `Some` | an event. In HTTP/2 an event always consumes; `consumed == 0` is only ever the stop case. | /// /// So the loop condition is `consumed > 0 || event.is_some()`: @@ -249,7 +384,21 @@ pub struct Connection { peer_streams: usize, preface: bool, settings_received: bool, - settings_awaiting_ack: bool, + /// SETTINGS frames sent and not yet acknowledged, oldest first. RFC 9113 + /// section 6.5.3 acknowledges them in order, so this is a queue: a second + /// frame sent while one is outstanding is ordinary, not an error. + settings_pending: VecDeque, + /// Every local parameter the peer has acknowledged, merged in order. + local_settings: Settings, + /// Every parameter the peer has sent, merged in order. + remote_settings: Settings, + /// The limits the acknowledged frames establish. `limits`, the ones + /// enforced, is this relaxed by every outstanding frame: see + /// `enforce_settings`. + acked_limits: Limits, + /// Our SETTINGS_INITIAL_WINDOW_SIZE as acknowledged, and as enforced. + acked_recv_initial: i64, + recv_initial: i64, decoder: hpack::Decoder, encoder: hpack::Encoder, block: Vec, @@ -266,13 +415,35 @@ pub struct Connection { settings_deadline: Option, } impl Connection { + /// A connection whose initial SETTINGS advertises `limits`: the + /// concurrent-stream, frame-size and header-list limits, in that + /// (ascending) identifier order, plus ENABLE_PUSH = 0 for a client. pub fn new(role: Role, limits: Limits) -> Result { + Self::with_settings(role, limits, &Settings::from_limits(&limits)) + } + /// A connection whose initial SETTINGS frame carries exactly `settings`, + /// serialised in ascending identifier order. `Settings::new()` sends an + /// empty frame, which is what Node's server sends. + /// + /// `limits` is what this endpoint enforces on its own. A parameter in + /// `settings` that corresponds to a limit replaces it once the peer + /// acknowledges the frame (and at once, if that relaxes it: the peer may + /// act on the new value before its ack reaches us). One parameter is + /// added: a client that does not name ENABLE_PUSH sends ENABLE_PUSH = 0, + /// because push defaults to enabled and this crate treats PUSH_PROMISE + /// as a connection error. + pub fn with_settings(role: Role, limits: Limits, settings: &Settings) -> Result { if !(16384..=0xffffff).contains(&limits.frame_size) || limits.streams == 0 || limits.streams > u32::MAX as usize { return Err(protocol("invalid limits")); } + let mut settings = settings.clone(); + if role == Role::Client && settings.get(Settings::ENABLE_PUSH).is_none() { + settings.set(Settings::ENABLE_PUSH, 0); + } + settings.validate()?; let mut result = Self { role, limits, @@ -287,7 +458,12 @@ impl Connection { peer_streams: usize::MAX, preface: role == Role::Client, settings_received: false, - settings_awaiting_ack: true, + settings_pending: VecDeque::new(), + local_settings: Settings::new(), + remote_settings: Settings::new(), + acked_limits: limits, + acked_recv_initial: 65535, + recv_initial: 65535, decoder: hpack::Decoder::new(4096, limits.header_list), encoder: hpack::Encoder::new(4096), block: Vec::new(), @@ -303,27 +479,100 @@ impl Connection { if role == Role::Client { result.output.extend_from_slice(PREFACE); } - let mut settings = Vec::new(); - for (id, value) in [ - (3, limits.streams as u32), - (5, limits.frame_size as u32), - (6, limits.header_list as u32), - ] { - settings.extend_from_slice(&(id as u16).to_be_bytes()); - settings.extend_from_slice(&value.to_be_bytes()); + result.send_settings(settings)?; + Ok(result) + } + /// Send a SETTINGS frame carrying `settings` to a live connection - + /// Node's `session.settings()`. It may be called while earlier frames + /// are still unacknowledged; each is acknowledged in turn with + /// [`Event::SettingsAck`]. + /// + /// A parameter that loosens what this endpoint enforces - a larger + /// MAX_FRAME_SIZE, INITIAL_WINDOW_SIZE, MAX_HEADER_LIST_SIZE or + /// MAX_CONCURRENT_STREAMS - applies at once, because the peer may act on + /// it before its acknowledgement arrives. One that tightens applies when + /// the peer acknowledges the frame (RFC 9113 section 6.5.3), so frames the + /// peer sent under the old value are never treated as errors. + /// + /// Values this endpoint could not honour are refused before anything is + /// sent: ENABLE_PUSH other than 0, HEADER_TABLE_SIZE above 4096, and the + /// out-of-range values RFC 9113 section 6.5.2 forbids. + pub fn settings(&mut self, settings: &Settings) -> Result<()> { + if self.failed { + return Err(protocol("failed connection")); } - if role == Role::Client { - settings.extend_from_slice(&[0, 2, 0, 0, 0, 0]); + settings.validate()?; + self.send_settings(settings.clone()) + } + fn send_settings(&mut self, settings: Settings) -> Result<()> { + let payload = settings.encode(); + if payload.len() > self.peer_frame { + return Err(frame_error()); } - result.frame(4, 0, 0, &settings)?; - Ok(result) + self.frame(4, 0, 0, &payload)?; + self.settings_pending.push_back(settings); + self.enforce_settings(); + Ok(()) + } + /// Recompute the enforced limits: the acknowledged ones, loosened by every + /// frame still outstanding. The peer may already be acting on any of them. + fn enforce_settings(&mut self) { + let mut limits = self.acked_limits; + let mut recv_initial = self.acked_recv_initial; + for settings in &self.settings_pending { + for (id, value) in settings.iter() { + let value = value as usize; + match id { + Settings::MAX_CONCURRENT_STREAMS => limits.streams = limits.streams.max(value), + Settings::INITIAL_WINDOW_SIZE => recv_initial = recv_initial.max(value as i64), + Settings::MAX_FRAME_SIZE => limits.frame_size = limits.frame_size.max(value), + Settings::MAX_HEADER_LIST_SIZE => { + limits.header_list = limits.header_list.max(value) + } + _ => {} + } + } + } + // RFC 9113 section 6.9.2: a change to the initial window adjusts every + // open stream's window by the difference. + let delta = recv_initial - self.recv_initial; + if delta != 0 { + for s in &mut self.streams { + s.recv_window += delta; + } + self.recv_initial = recv_initial; + } + self.decoder.max_list_size = limits.header_list; + self.limits = limits; + } + /// What this endpoint enforces now, including the effect of SETTINGS it + /// has sent (see [`settings`](Self::settings)). + pub fn limits(&self) -> Limits { + self.limits + } + /// Every parameter this endpoint has sent and the peer has acknowledged - + /// Node's `session.localSettings`, less the defaults nobody sent. + pub fn local_settings(&self) -> &Settings { + &self.local_settings + } + /// Every parameter the peer has sent - Node's `session.remoteSettings`, + /// less the defaults the peer did not send. + pub fn remote_settings(&self) -> &Settings { + &self.remote_settings + } + /// SETTINGS frames sent and not yet acknowledged. + pub fn pending_settings(&self) -> usize { + self.settings_pending.len() } - /// Host-supplied SETTINGS acknowledgement deadline; no clock is sampled. + /// Host-supplied deadline for the oldest outstanding SETTINGS frame's + /// acknowledgement; no clock is sampled. Each acknowledgement clears it, so + /// a host with further frames outstanding sets the next one after each + /// [`Event::SettingsAck`]. pub fn set_settings_deadline(&mut self, deadline: Option) { self.settings_deadline = deadline; } pub fn next_timeout(&self) -> Option { - if self.settings_awaiting_ack && !self.failed { + if !self.settings_pending.is_empty() && !self.failed { self.settings_deadline } else { None @@ -441,7 +690,7 @@ impl Connection { let stream = Stream { id, send_window: self.initial_send, - recv_window: 65535, + recv_window: self.recv_initial, unreleased: 0, credited: false, local_end: false, @@ -914,10 +1163,29 @@ impl Connection { if !f.payload.is_empty() { return Err(frame_error()); } - if !self.settings_awaiting_ack { + let Some(acked) = self.settings_pending.pop_front() else { return Err(protocol("unsolicited SETTINGS ack")); + }; + for (id, value) in acked.iter() { + match id { + Settings::MAX_CONCURRENT_STREAMS => { + self.acked_limits.streams = value as usize + } + Settings::INITIAL_WINDOW_SIZE => self.acked_recv_initial = value as i64, + Settings::MAX_FRAME_SIZE => { + self.acked_limits.frame_size = value as usize + } + Settings::MAX_HEADER_LIST_SIZE => { + self.acked_limits.header_list = value as usize + } + _ => {} + } } - self.settings_awaiting_ack = false; + self.local_settings.merge(acked.iter()); + self.enforce_settings(); + // The deadline was for this frame; the host arms the next. + self.settings_deadline = None; + step.event = Some(Event::SettingsAck(acked)); } else { if f.payload.len() % 6 != 0 { return Err(frame_error()); @@ -959,8 +1227,10 @@ impl Connection { } } self.settings_received = true; + let settings = SettingsFrame { payload: f.payload }; + self.remote_settings.merge(settings.iter()); self.frame(4, 1, 0, &[])?; - step.event = Some(Event::Settings); + step.event = Some(Event::Settings(settings)); } } 5 => return Err(protocol("server push disabled")), @@ -989,6 +1259,7 @@ impl Connection { step.event = Some(Event::Goaway { last_stream: u32be(f.payload) & 0x7fffffff, code: u32be(&f.payload[4..]), + debug: Some(&f.payload[8..]).filter(|d| !d.is_empty()), }); } 8 => { diff --git a/protocols/turnloop-http/src/lib.rs b/protocols/turnloop-http/src/lib.rs index d574932..a598903 100644 --- a/protocols/turnloop-http/src/lib.rs +++ b/protocols/turnloop-http/src/lib.rs @@ -23,7 +23,7 @@ //! so neither is optional: //! //! * **`consumed > 0` with no event** is progress with nothing to hand up: the -//! HTTP/2 client preface, a SETTINGS acknowledgement, a PRIORITY frame, an +//! HTTP/2 client preface, a PRIORITY frame, an //! unknown frame type, or a frame the peer had in flight for a stream that is //! already gone; an HTTP/1 chunk-size line, chunk CRLF or empty trailer //! block. A loop that continues only while an event came back stalls here, and diff --git a/protocols/turnloop-http/tests/codecs.rs b/protocols/turnloop-http/tests/codecs.rs index 583cc78..e61bea2 100644 --- a/protocols/turnloop-http/tests/codecs.rs +++ b/protocols/turnloop-http/tests/codecs.rs @@ -757,7 +757,8 @@ fn drive(to: &mut http2::Connection, input: &mut Vec) -> Result, let progressed = consumed > 0 || step.event.is_some(); if let Some(event) = step.event { seen.push(match event { - http2::Event::Settings => "Settings".to_string(), + http2::Event::Settings(_) => "Settings".to_string(), + http2::Event::SettingsAck(_) => "SettingsAck".to_string(), http2::Event::Headers { stream, .. } => format!("Headers s={stream}"), http2::Event::Data { stream, bytes, .. } => { format!("Data s={stream} n={}", bytes.len()) @@ -1163,7 +1164,7 @@ fn h2_step_has_two_independent_zero_cases() { vec![ (24, false), // consumed > 0, event == None: the preface. KEEP GOING. (9, true), // SETTINGS - (9, false), // consumed > 0, event == None: the SETTINGS ack. + (9, true), // SettingsAck: the ack of our initial SETTINGS. (14, false), // consumed > 0, event == None: PRIORITY. (0, false), // consumed == 0, event == None: exhausted. STOP. ] diff --git a/protocols/turnloop-http/tests/http2_settings.rs b/protocols/turnloop-http/tests/http2_settings.rs new file mode 100644 index 0000000..6daff43 --- /dev/null +++ b/protocols/turnloop-http/tests/http2_settings.rs @@ -0,0 +1,420 @@ +//! SETTINGS as a live, two-sided negotiation (PerryTS/turnloop#87), and the +//! GOAWAY opaque data that rides beside it. +use std::time::{Duration, Instant}; +use turnloop_http::{ + http1::Header, + http2::{self, Connection, Event, Limits, Role, Settings}, +}; + +/// One observed event, owned so it outlives the input it borrowed from. +#[derive(Debug, PartialEq)] +enum Seen { + Settings(Vec<(u16, u32)>), + SettingsAck(Vec<(u16, u32)>), + Headers(u32), + Data(u32, usize), + Reset(u32, u32), + Goaway(u32, Option>), + Other, +} + +/// Drive `to` with the documented loop condition until it stops. +fn drive(to: &mut Connection, input: &mut Vec) -> Result, &'static str> { + let mut seen = Vec::new(); + loop { + let step = to.receive(input).map_err(|e| e.code)?; + let consumed = step.consumed; + let progressed = consumed > 0 || step.event.is_some(); + if let Some(event) = step.event { + seen.push(match event { + Event::Settings(frame) => Seen::Settings(frame.iter().collect()), + Event::SettingsAck(settings) => Seen::SettingsAck(settings.iter().collect()), + Event::Headers { stream, .. } => Seen::Headers(stream), + Event::Data { stream, bytes, .. } => Seen::Data(stream, bytes.len()), + Event::Reset { stream, code } => Seen::Reset(stream, code), + Event::Goaway { code, debug, .. } => Seen::Goaway(code, debug.map(<[u8]>::to_vec)), + _ => Seen::Other, + }); + } + input.drain(..consumed); + if !progressed { + return Ok(seen); + } + } +} +fn ship(from: &mut Connection) -> Vec { + let wire = from.output().to_vec(); + from.consume_output(wire.len()).unwrap(); + wire +} +fn request() -> Vec
{ + vec![ + Header::new(":method", "POST"), + Header::new(":scheme", "http"), + Header::new(":path", "/"), + Header::new(":authority", "localhost"), + ] +} +/// Default-limits client and server, both initial SETTINGS exchanged and acked. +fn handshake() -> (Connection, Connection) { + let mut client = Connection::new(Role::Client, Limits::default()).unwrap(); + let mut server = Connection::new(Role::Server, Limits::default()).unwrap(); + drive(&mut server, &mut ship(&mut client)).unwrap(); + drive(&mut client, &mut ship(&mut server)).unwrap(); + drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!( + (client.pending_settings(), server.pending_settings()), + (0, 0) + ); + (client, server) +} +const DEFAULT_ADVERTISED: [(u16, u32); 3] = [(3, 100), (5, 16384), (6, 32768)]; + +/// Item 1. A second SETTINGS while the first is outstanding is ordinary: RFC +/// 9113 section 6.5.3 acknowledges them in order. With a single "awaiting ack" +/// flag the second ack was "unsolicited" and failed the connection. +#[test] +fn settings_acks_are_an_ordered_queue() { + let mut client = Connection::new(Role::Client, Limits::default()).unwrap(); + let mut server = Connection::new(Role::Server, Limits::default()).unwrap(); + let mut second = Settings::new(); + second.set(Settings::MAX_CONCURRENT_STREAMS, 7); + let mut third = Settings::new(); + third.set(Settings::MAX_HEADER_LIST_SIZE, 65536); + // Two more before the peer has seen even the first. + server.settings(&second).unwrap(); + server.settings(&third).unwrap(); + assert_eq!(server.pending_settings(), 3); + + // The server reads the client's preface and SETTINGS and acks them. + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!( + seen, + vec![Seen::Settings(vec![ + (2, 0), + (3, 100), + (5, 16384), + (6, 32768) + ])] + ); + let seen = drive(&mut client, &mut ship(&mut server)).unwrap(); + assert_eq!( + seen, + vec![ + Seen::Settings(DEFAULT_ADVERTISED.to_vec()), + Seen::Settings(vec![(3, 7)]), + Seen::Settings(vec![(6, 65536)]), + Seen::SettingsAck(vec![(2, 0), (3, 100), (5, 16384), (6, 32768)]), + ] + ); + // The client acked all three; the server pairs them with what it sent. + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!( + seen, + vec![ + Seen::SettingsAck(DEFAULT_ADVERTISED.to_vec()), + Seen::SettingsAck(vec![(3, 7)]), + Seen::SettingsAck(vec![(6, 65536)]), + ] + ); + assert_eq!(server.pending_settings(), 0); + assert_eq!( + server.local_settings().iter().collect::>(), + vec![(3, 7), (5, 16384), (6, 65536)] + ); + + // A fourth ack has nothing to acknowledge, and that is still an error. + let mut wire = Vec::new(); + http2::encode_frame(4, 1, 0, &[], &mut wire).unwrap(); + assert_eq!(drive(&mut server, &mut wire), Err("PROTOCOL_ERROR")); +} + +/// Item 2. The acknowledgement is an event carrying the settings it +/// acknowledges, and it clears the host's deadline for that frame - so a host +/// can resolve `session.settings(obj, callback)` and time the round trip. +#[test] +fn settings_ack_is_an_event_and_clears_its_deadline() { + let (mut client, mut server) = handshake(); + assert_eq!(server.next_timeout(), None); + let mut change = Settings::new(); + change.set(Settings::INITIAL_WINDOW_SIZE, 100_000); + server.settings(&change).unwrap(); + let deadline = Instant::now() + Duration::from_secs(5); + server.set_settings_deadline(Some(deadline)); + assert_eq!(server.next_timeout(), Some(deadline)); + + let seen = drive(&mut client, &mut ship(&mut server)).unwrap(); + assert_eq!(seen, vec![Seen::Settings(vec![(4, 100_000)])]); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![(4, 100_000)])]); + assert_eq!(server.next_timeout(), None); + assert_eq!( + server.handle_timeout(deadline + Duration::from_secs(1)), + None + ); +} + +/// Item 3. `Event::Settings` carries the peer's frame as sent - order, +/// duplicates and unknown identifiers included - and `remote_settings` +/// accumulates what is in force. +#[test] +fn settings_event_carries_the_peers_parameters() { + let mut server = Connection::new(Role::Server, Limits::default()).unwrap(); + let mut wire = http2::PREFACE.to_vec(); + let mut payload = Vec::new(); + for (id, value) in [(4u16, 1000u32), (0x7a, 9), (4, 2000), (5, 20000)] { + payload.extend_from_slice(&id.to_be_bytes()); + payload.extend_from_slice(&value.to_be_bytes()); + } + http2::encode_frame(4, 0, 0, &payload, &mut wire).unwrap(); + let step = server.receive(&wire).unwrap(); + wire.drain(..step.consumed); + let step = server.receive(&wire).unwrap(); + let Some(Event::Settings(frame)) = step.event else { + panic!("expected Settings, got {:?}", step.event); + }; + assert_eq!( + frame.iter().collect::>(), + vec![(4, 1000), (0x7a, 9), (4, 2000), (5, 20000)] + ); + assert_eq!( + frame.get(4), + Some(2000), + "the last value is the one in force" + ); + assert_eq!(frame.get(1), None); + assert_eq!( + server.remote_settings().iter().collect::>(), + vec![(4, 2000), (5, 20000), (0x7a, 9)] + ); + + // An empty SETTINGS frame is an event too, with nothing in it. + let mut wire = Vec::new(); + http2::encode_frame(4, 0, 0, &[], &mut wire).unwrap(); + let step = server.receive(&wire).unwrap(); + assert!(matches!(step.event, Some(Event::Settings(frame)) if frame.is_empty())); +} + +/// Item 4. GOAWAY's opaque data reaches the host, and "none" stays +/// distinguishable from data - Node reports `undefined`, not an empty buffer. +#[test] +fn goaway_carries_its_debug_data_and_absent_stays_absent() { + let (mut client, mut server) = handshake(); + server.goaway(11, 0, b"enhance").unwrap(); + assert_eq!( + drive(&mut client, &mut ship(&mut server)).unwrap(), + vec![Seen::Goaway(11, Some(b"enhance".to_vec()))] + ); + let (mut client, mut server) = handshake(); + server.shutdown().unwrap(); + assert_eq!( + drive(&mut client, &mut ship(&mut server)).unwrap(), + vec![Seen::Goaway(0, None)] + ); +} + +/// Item 5. The initial SETTINGS is the host's choice - empty, as Node's server +/// sends it, or any parameters, serialised in ascending identifier order. +#[test] +fn initial_settings_are_host_chosen_and_ordered() { + let mut server = + Connection::with_settings(Role::Server, Limits::default(), &Settings::new()).unwrap(); + assert_eq!(server.output(), [0, 0, 0, 4, 0, 0, 0, 0, 0]); + + // Set out of order; written in ascending order. + let mut chosen = Settings::new(); + chosen + .set(Settings::MAX_HEADER_LIST_SIZE, 1 << 16) + .set(Settings::HEADER_TABLE_SIZE, 4096) + .set(Settings::ENABLE_CONNECT_PROTOCOL, 1); + let ordered = Connection::with_settings(Role::Server, Limits::default(), &chosen).unwrap(); + let mut expected = Vec::new(); + http2::encode_frame( + 4, + 0, + 0, + &[ + 0, 1, 0, 0, 0x10, 0, // HEADER_TABLE_SIZE 4096 + 0, 6, 0, 1, 0, 0, // MAX_HEADER_LIST_SIZE 65536 + 0, 8, 0, 0, 0, 1, // ENABLE_CONNECT_PROTOCOL 1 + ], + &mut expected, + ) + .unwrap(); + assert_eq!(ordered.output(), expected); + + // `new` still advertises its limits, now in ascending order for a client + // too: ENABLE_PUSH (2) first rather than appended. + let client = Connection::new(Role::Client, Limits::default()).unwrap(); + let settings = &client.output()[http2::PREFACE.len()..]; + let frame = http2::decode_frame(settings, 16384).unwrap().unwrap(); + let ids: Vec = frame.payload.chunks(6).map(|s| s[1]).collect(); + assert_eq!(ids, [2, 3, 5, 6]); + + // A client that names no ENABLE_PUSH still refuses push: it defaults to on. + let client = + Connection::with_settings(Role::Client, Limits::default(), &Settings::new()).unwrap(); + assert_eq!( + &client.output()[http2::PREFACE.len()..], + [0, 0, 6, 4, 0, 0, 0, 0, 0, 0, 2, 0, 0, 0, 0] + ); + + // An empty-SETTINGS server still serves a request end to end. + let mut client = Connection::new(Role::Client, Limits::default()).unwrap(); + drive(&mut server, &mut ship(&mut client)).unwrap(); + drive(&mut client, &mut ship(&mut server)).unwrap(); + let id = client.open(&request(), true).unwrap(); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![]), Seen::Headers(id)]); +} + +/// Values this endpoint could not honour are refused before a byte is sent. +#[test] +fn settings_it_cannot_honour_are_refused() { + let (_, mut server) = handshake(); + for (id, value) in [ + (Settings::ENABLE_PUSH, 1), + (Settings::HEADER_TABLE_SIZE, 8192), + (Settings::INITIAL_WINDOW_SIZE, 0x8000_0000), + (Settings::MAX_FRAME_SIZE, 16383), + (Settings::MAX_FRAME_SIZE, 0x0100_0000), + (Settings::ENABLE_CONNECT_PROTOCOL, 2), + ] { + let mut bad = Settings::new(); + bad.set(id, value); + assert!(server.settings(&bad).is_err(), "{id}={value}"); + assert!( + Connection::with_settings(Role::Server, Limits::default(), &bad).is_err(), + "{id}={value}" + ); + } + assert!(server.output().is_empty()); + assert_eq!(server.pending_settings(), 0); +} + +/// Item 6. A live change applies to the running connection: a loosening at +/// once (the peer may use it before its ack arrives), a tightening on the ack. +#[test] +fn live_settings_change_applies_loosening_now_and_tightening_on_ack() { + let (mut client, mut server) = handshake(); + let big = vec![0u8; 20000]; + + // Loosen MAX_FRAME_SIZE. The client acts on it as soon as it reads the + // frame, and its 20000-byte DATA may arrive before its ack does. + let mut larger = Settings::new(); + larger.set(Settings::MAX_FRAME_SIZE, 32768); + server.settings(&larger).unwrap(); + assert_eq!( + server.limits().frame_size, + 32768, + "a loosening applies at once" + ); + drive(&mut client, &mut ship(&mut server)).unwrap(); + let ack_and_more = ship(&mut client); + let id = client.open(&request(), false).unwrap(); + assert_eq!(client.send_data(id, &big, false).unwrap(), 20000); + // Deliver the DATA first, then the ack. + let mut data_first = ship(&mut client); + let seen = drive(&mut server, &mut data_first).unwrap(); + assert_eq!(seen, vec![Seen::Headers(id), Seen::Data(id, 20000)]); + let seen = drive(&mut server, &mut ack_and_more.clone()).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![(5, 32768)])]); + server.release_capacity(id, 20000).unwrap(); + drive(&mut client, &mut ship(&mut server)).unwrap(); + + // Tighten it again. Until the ack, a large frame is still legal - the + // client may not have read the change yet. + let mut smaller = Settings::new(); + smaller.set(Settings::MAX_FRAME_SIZE, 16384); + server.settings(&smaller).unwrap(); + assert_eq!( + server.limits().frame_size, + 32768, + "a tightening waits for the ack" + ); + let pending = ship(&mut server); + assert_eq!(client.send_data(id, &big, false).unwrap(), 20000); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::Data(id, 20000)]); + server.release_capacity(id, 20000).unwrap(); + drive(&mut client, &mut ship(&mut server)).unwrap(); + // Now the client reads it, acks, and splits its DATA to 16384. + drive(&mut client, &mut pending.clone()).unwrap(); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![(5, 16384)])]); + assert_eq!(server.limits().frame_size, 16384); + assert_eq!(client.send_data(id, &big, false).unwrap(), 16384); + // A peer that ignores the acknowledged limit is answered FRAME_SIZE_ERROR. + let mut oversized = Vec::new(); + http2::encode_frame(0, 0, id, &big, &mut oversized).unwrap(); + assert_eq!(drive(&mut server, &mut oversized), Err("FRAME_SIZE_ERROR")); +} + +/// Item 6, streams: tightening MAX_CONCURRENT_STREAMS applies on the ack. +/// Streams the peer opened before it read the change are accepted; one past +/// the new limit afterwards is refused as a stream error. +#[test] +fn live_stream_limit_applies_on_ack() { + let (mut client, mut server) = handshake(); + let mut one = Settings::new(); + one.set(Settings::MAX_CONCURRENT_STREAMS, 1); + server.settings(&one).unwrap(); + assert_eq!( + server.limits().streams, + 100, + "a tightening waits for the ack" + ); + let pending = ship(&mut server); + let a = client.open(&request(), false).unwrap(); + let b = client.open(&request(), false).unwrap(); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::Headers(a), Seen::Headers(b)]); + + drive(&mut client, &mut pending.clone()).unwrap(); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![(3, 1)])]); + assert_eq!(server.limits().streams, 1); + // The client now respects the limit itself... + assert!(client.open(&request(), true).is_err()); + // ...and a peer that does not is refused: REFUSED_STREAM, not a + // connection error. The block is static-table only (:method POST, + // :scheme http, :path /), so it decodes against any HPACK state. + let mut wire = Vec::new(); + http2::encode_frame(1, 5, 5, &[0x83, 0x86, 0x84], &mut wire).unwrap(); + assert_eq!( + drive(&mut server, &mut wire).unwrap(), + vec![Seen::Reset(5, 7)] + ); +} + +/// Item 6, windows: INITIAL_WINDOW_SIZE moves every open stream's receive +/// window by the difference (RFC 9113 section 6.9.2) once acknowledged, and +/// new streams start at the new size. +#[test] +fn live_initial_window_change_adjusts_streams() { + let (mut client, mut server) = handshake(); + let a = client.open(&request(), false).unwrap(); + drive(&mut server, &mut ship(&mut client)).unwrap(); + let mut small = Settings::new(); + small.set(Settings::INITIAL_WINDOW_SIZE, 1000); + server.settings(&small).unwrap(); + let pending = ship(&mut server); + // Before the client reads it, 2000 bytes on `a` are within the old window. + assert_eq!(client.send_data(a, &[0; 2000], false).unwrap(), 2000); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::Data(a, 2000)]); + + drive(&mut client, &mut pending.clone()).unwrap(); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::SettingsAck(vec![(4, 1000)])]); + // A new stream starts at 1000, on both sides. + let b = client.open(&request(), false).unwrap(); + assert_eq!(client.send_data(b, &[0; 2000], false).unwrap(), 1000); + let seen = drive(&mut server, &mut ship(&mut client)).unwrap(); + assert_eq!(seen, vec![Seen::Headers(b), Seen::Data(b, 1000)]); + // `a` is now at 1000 - 65535 + (65535 - 2000) = -1000: one more byte on it + // overruns the window the peer was told about. + let mut wire = Vec::new(); + http2::encode_frame(0, 0, a, &[0], &mut wire).unwrap(); + assert_eq!(drive(&mut server, &mut wire), Err("FLOW_CONTROL_ERROR")); +} diff --git a/protocols/turnloop-websocket/src/lib.rs b/protocols/turnloop-websocket/src/lib.rs index 41964fa..5d03470 100644 --- a/protocols/turnloop-websocket/src/lib.rs +++ b/protocols/turnloop-websocket/src/lib.rs @@ -60,6 +60,56 @@ pub struct Connection { terminal: bool, peer_close: Option, } +/// One [`Connection::receive`] step: how much of `input` was taken, and the +/// message it completed, if any. +/// +/// **`consumed` and `message` are independent**, as they are for +/// `turnloop_http`'s `Step`. `receive` reads input into the connection's own +/// buffer and parses from that buffer before it reads again, so the bytes a +/// message is made of may have been consumed by an earlier call: +/// +/// | `consumed` | `message` | meaning | +/// |---|---|---| +/// | `0` | `None` | **stop.** Nothing complete is buffered and the input is empty; read more from the transport before calling again. | +/// | `> 0` | `None` | **keep going.** The input was taken in, all of it, but completes no message yet: a partial frame, or a non-final fragment. The next call returns the stop shape unless more input arrived. | +/// | `> 0` | `Some` | a message, possibly with input left over. | +/// | `0` | `Some` | a message completed from bytes an earlier call consumed - the second of two frames that arrived in one read, say. Normal, and easy to drop. | +/// +/// So the loop condition is the same one-liner, `consumed > 0 || +/// message.is_some()`: +/// +/// ``` +/// # use turnloop_websocket::*; +/// # let mut connection = Connection::new(Role::Server, WebSocketConfig::default()); +/// # let (mut input, mut replies) = (Vec::new(), Vec::new()); +/// loop { +/// let step = connection.receive(&input, &mut replies)?; +/// input.drain(..step.consumed); +/// let progressed = step.consumed > 0 || step.message.is_some(); +/// if let Some(message) = step.message { +/// let _ = message; +/// } +/// connection.flush(&mut replies)?; // automatic pong/close replies +/// if !progressed { +/// break; // read more bytes from the transport, then continue +/// } +/// } +/// # Ok::<(), Error>(()) +/// ``` +/// +/// Looping while `consumed > 0` **drops messages**: the `0`/`Some` step reads +/// as "stop" and its message is never looked at. Reading the transport +/// whenever `input` is empty **stalls** on the same step: no input is left, +/// yet a message is waiting, and a peer waiting for its answer never sends +/// the bytes that would wake the host. Looping while a message came back is, +/// for this type, safe - a step with no message has always taken in all of +/// its input, so nothing is stranded - but it is not the documented +/// condition, and the one above works for every step type in the workspace. +/// +/// `receive` is not idempotent. The `consumed` bytes are in the connection's +/// buffer once it returns, so input that is fed again without being drained +/// is parsed again: a complete message in it is delivered twice, with no +/// error to notice. pub struct Received { pub consumed: usize, pub message: Option, @@ -73,6 +123,8 @@ impl Connection { peer_close: None, } } + /// One decode step. `consumed` and `message` are independent, and + /// `consumed == 0` does not by itself mean "stop": see [`Received`]. pub fn receive(&mut self, input: &[u8], output: &mut Vec) -> Result { if self.terminal { return Err(Error::AlreadyClosed); diff --git a/protocols/turnloop-websocket/tests/websocket.rs b/protocols/turnloop-websocket/tests/websocket.rs index a12a35d..7f6d0f6 100644 --- a/protocols/turnloop-websocket/tests/websocket.rs +++ b/protocols/turnloop-websocket/tests/websocket.rs @@ -31,6 +31,159 @@ fn handshakes_subprotocol_and_bad_key() { assert!(client.verify(&response).is_err()); assert!(ClientHandshake::new("localhost", "/", [0; 16], vec!["a".into(), "a".into()]).is_err()); } +/// Three masked client text frames, the last cut after 5 of its 11 bytes: the +/// shape a transport read that splits a frame produces. +fn three_messages_last_split() -> (Vec, Vec) { + let mut client = Connection::new(Role::Client, config()); + let mut wire = Vec::new(); + for text in ["one", "two", "three"] { + client.send(Message::text(text), &mut wire).unwrap(); + } + assert_eq!(wire.len(), 9 + 9 + 11); + let rest = wire.split_off(9 + 9 + 5); + (wire, rest) +} +/// Drive with the documented condition, `consumed > 0 || message.is_some()`. +fn drain(server: &mut Connection, input: &mut Vec, seen: &mut Vec<(usize, Option)>) { + let mut reply = Vec::new(); + loop { + let step = server.receive(input, &mut reply).unwrap(); + let progressed = step.consumed > 0 || step.message.is_some(); + input.drain(..step.consumed); + seen.push((step.consumed, step.message)); + if !progressed { + return; + } + } +} +/// `Received`'s shapes (PerryTS/turnloop#86). tungstenite reads input into its +/// own buffer and parses from that buffer before it reads again, so a message +/// can complete from bytes an earlier call consumed: `consumed == 0` with a +/// message is normal, and `consumed == 0` alone does not mean "stop". +#[test] +fn received_consumed_and_message_are_independent() { + let (mut input, rest) = three_messages_last_split(); + let mut server = Connection::new(Role::Server, config()); + let mut seen = Vec::new(); + drain(&mut server, &mut input, &mut seen); + assert_eq!( + seen, + vec![ + // Everything was read in; the first message came out of it. + (23, Some(Message::text("one"))), + // consumed == 0 with a message: parsed from the buffer. KEEP GOING. + (0, Some(Message::text("two"))), + // consumed == 0, no message: 5 bytes of a frame header. STOP. + (0, None), + ] + ); + assert!(input.is_empty()); + + let mut input = rest; + let mut seen = Vec::new(); + drain(&mut server, &mut input, &mut seen); + assert_eq!(seen, vec![(6, Some(Message::text("three"))), (0, None)]); + + // consumed > 0 with no message: a partial frame, all of it taken in. + // Calling again is harmless and returns the stop shape. + let (mut input, _) = three_messages_last_split(); + input.truncate(5); + let mut server = Connection::new(Role::Server, config()); + let mut seen = Vec::new(); + drain(&mut server, &mut input, &mut seen); + assert_eq!(seen, vec![(5, None), (0, None)]); +} +/// The conditions the step contract rules out, and the one it does not. +#[test] +fn received_rejects_the_wrong_loop_conditions() { + // "Loop while bytes were consumed" loses "two": it arrives in a step that + // consumed nothing, which this loop reads as the stop signal. + let (mut input, rest) = three_messages_last_split(); + let mut server = Connection::new(Role::Server, config()); + let mut reply = Vec::new(); + let mut messages = Vec::new(); + loop { + let step = server.receive(&input, &mut reply).unwrap(); + if step.consumed == 0 { + break; + } + input.drain(..step.consumed); + messages.extend(step.message); + } + input.extend_from_slice(&rest); + loop { + let step = server.receive(&input, &mut reply).unwrap(); + if step.consumed == 0 { + break; + } + input.drain(..step.consumed); + messages.extend(step.message); + } + assert_eq!( + messages, + vec![Message::text("one"), Message::text("three")], + "\"two\" came back with consumed == 0 and was dropped" + ); + + // "Read the transport whenever the input is empty" stalls: after "one" + // there is no input left, but "two" is complete in the connection's own + // buffer. A host that waits for the peer here waits forever if the peer + // is waiting for an answer to "two". + let (mut input, _) = three_messages_last_split(); + let mut server = Connection::new(Role::Server, config()); + let step = server.receive(&input, &mut reply).unwrap(); + input.drain(..step.consumed); + assert_eq!(step.message, Some(Message::text("one"))); + assert!(input.is_empty(), "nothing left to hand in..."); + let step = server.receive(&input, &mut reply).unwrap(); + assert_eq!( + (step.consumed, step.message), + (0, Some(Message::text("two"))), + "...and yet a message was waiting" + ); + + // "Loop while a message came back" is safe for this type, unlike + // http2::Step: a step without a message has always taken in all of its + // input, so nothing is stranded when such a host goes back to the + // transport. Pinned here because it rests on how tungstenite reads. + let (mut input, rest) = three_messages_last_split(); + let mut server = Connection::new(Role::Server, config()); + let mut messages = Vec::new(); + for chunk in [None, Some(rest)] { + input.extend(chunk.into_iter().flatten()); + loop { + let step = server.receive(&input, &mut reply).unwrap(); + input.drain(..step.consumed); + let Some(message) = step.message else { break }; + messages.push(message); + } + assert!(input.is_empty()); + } + assert_eq!( + messages, + ["one", "two", "three"].map(Message::text).to_vec() + ); +} +/// `receive` is not idempotent. It takes `consumed` bytes into its own buffer, +/// so a host that does not drain them and feeds them again gets every message +/// in them twice - silently, with no error to notice. +#[test] +fn received_is_not_idempotent() { + let (input, _) = three_messages_last_split(); + let one = &input[..9]; + let mut server = Connection::new(Role::Server, config()); + let mut reply = Vec::new(); + let first = server.receive(one, &mut reply).unwrap(); + assert_eq!( + (first.consumed, first.message), + (9, Some(Message::text("one"))) + ); + let again = server.receive(one, &mut reply).unwrap(); + assert_eq!( + (again.consumed, again.message), + (9, Some(Message::text("one"))) + ); +} #[test] fn masking_fragmentation_ping_pong_close_and_limits() { let mut client = Connection::new(Role::Client, config()); @@ -131,6 +284,7 @@ fn serve(mut socket: TcpStream) { loop { let event = ws.receive(&input, &mut output).unwrap(); input.drain(..event.consumed); + let progressed = event.consumed > 0 || event.message.is_some(); if let Some(message) = event.message { match message { Message::Text(text) => { @@ -154,7 +308,9 @@ fn serve(mut socket: TcpStream) { } socket.write_all(&output).unwrap(); output.clear(); - if input.is_empty() { + // Not `input.is_empty()`: a message can be complete in the + // connection's own buffer with no input left to hand in. + if !progressed { let mut b = [0; 1024]; let n = socket.read(&mut b).unwrap(); assert!(n > 0);