Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / crates / npro-test / tests / ws_client_replay.rs
Author[]Andy Green <andy@warmcat.com> 2026-10-05 13:43 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-10-05 16:11 UTC
Tree396443a6367471c81245a9b620b36c92052af2f3   Raw Patch
 
npro-ws: the server's ws handshake and connection
npro-ws: the server's ws handshake and connection

The server half of phase 1e: a new no_std crate, the port of C's ws role
(lib/sansio/ws), sans-IO as the h1 sides are.

 - handshake::server() is C's lws_process_ws_upgrade() over the h1
   server's request, its checks in C's order, each refusal C's status: a
   GET; "upgrade" among the Connection tokens; a key under 128 bytes and
   a Host; version 13, a 400 with none and a 426 saying
   "sec-websocket-version: 13" for another; and the first subprotocol of
   the request's list the server has, or its default for none.  The
   accept is base64 of the SHA-1 of the key and RFC 6455's GUID, and
   response_101() writes C's 101 header for header.

 - conn::Ws is C's lws_ws_rx_sm(), its writeable handling and its close
   states.  rx() takes a thing at a time and unmasks a payload where it
   lies in the input, handing the application each piece of a message as
   it arrives, with whether it starts and ends it.  Control frames are
   gathered: a ping is answered with C's one pending pong, a pong given
   to the application, and the peer's close given to it and answered with
   its own payload, its code made 1002 where C's answer_peer_close makes
   it.  What C refuses is refused with C's code and reason ("frag ctl",
   "bad opc", "bad cont", "rsv bits", "client unmasked", "ctl len", "bad
   len", 1002; "huge frame", 1009; "bad utf8", "partial utf8", 1007), and
   nothing is read after either close.  tx() writes in C's order and
   pulls the application's payload; a pong owed when we begin a close is
   forgotten, as C forgets it; close_when_flushed() is C's close with no
   close frame once the last message has gone.

The workspace names npro-core and npro-h1 with a version beside the
path, since a published crate's in-tree dependencies need one; deny.toml
admits npro-ws, and ci.sh and sai.sh hold it to the no_std builds.

Not yet: the client's side, close deadlines and keepalive, a configured
frame maximum and a message maximum.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019kg5Eemy68ZaqDBcUJQG6J
diff --git a/Cargo.lock b/Cargo.lock index 9be86bd..ccaae9d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -30,3 +30,11 @@ dependencies = [ "npro-core", "npro-h1", ] + +[[package]] +name = "npro-ws" +version = "0.0.2" +dependencies = [ + "npro-core", + "npro-h1", +] diff --git a/Cargo.toml b/Cargo.toml index 60b0372..bab3ec9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,6 +10,12 @@ license = "MIT" homepage = "https://npro.rs" repository = "https://libwebsockets.org/git/npro" +# The crates here that a published crate depends on: a published crate's +# dependency needs a version beside its path, which crates.io resolves +[workspace.dependencies] +npro-core = { path = "crates/npro-core", version = "0.0.2" } +npro-h1 = { path = "crates/npro-h1", version = "0.0.2" } + [workspace.lints.rust] unsafe_code = "forbid" missing_docs = "deny" diff --git a/crates/npro-ws/Cargo.toml b/crates/npro-ws/Cargo.toml new file mode 100644 index 0000000..5a83646 --- /dev/null +++ b/crates/npro-ws/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "npro-ws" +description = "The sans-IO websockets of npro: the handshake, framing and close" +readme = "../../README.md" +keywords = ["websocket", "sans-io", "no-std"] +categories = ["network-programming", "web-programming::websocket", "no-std"] +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +homepage.workspace = true +repository.workspace = true + +[dependencies] +npro-core.workspace = true +npro-h1.workspace = true + +[lints] +workspace = true diff --git a/crates/npro-ws/src/conn.rs b/crates/npro-ws/src/conn.rs new file mode 100644 index 0000000..7de04e1 --- /dev/null +++ b/crates/npro-ws/src/conn.rs @@ -0,0 +1,867 @@ +//! A ws connection: C's `lws_ws_rx_sm()`, its writeable handling and its +//! close, sans-IO. +//! +//! [`Ws::rx`] takes the peer's bytes, at most one thing each call, and +//! unmasks a frame's payload where it lies in the input, so a message's +//! data is handed to the application without a copy: [`Event::Message`], +//! a piece of a message as it arrives, with whether it starts and ends it, +//! as C's `lws_is_first_fragment()` and `lws_is_final_fragment()` say. +//! Control frames, at most 125 bytes, are gathered here: a ping is answered +//! with a pong, a pong is given to the application, and the peer's close +//! is given to it and answered with the peer's own payload, its code made +//! 1002 if it is one no peer may send. +//! +//! What C refuses, this refuses, with C's close code and reason: a +//! fragmented control frame ("frag ctl"), a reserved opcode ("bad opc"), a +//! continuation out of place ("bad cont"), RSV bits ("rsv bits"), a server's +//! unmasked frame ("client unmasked"), a long control frame ("ctl len"), a +//! length with its top bit ("bad len"), all 1002; a frame longer than C's +//! 256MiB, 1009 "huge frame"; text that is not UTF-8, 1007 "bad utf8" or +//! "partial utf8". After its close, the connection reads nothing more. +//! +//! [`Ws::tx`] writes in C's order: what is in flight first (the 101, a frame +//! begun), then our own close, then the pong, then the answer to the peer's +//! close, then the application's next frame, whose payload it pulls. A +//! pong still owed when we begin a close is forgotten, as C forgets it. + +use npro_core::utf8::Utf8Validator; +use npro_h1::server::TxSource; + +/// The longest frame C takes: `LWS_WS_MAX_RX_FRAME_LEN`. +pub const MAX_FRAME: u64 = 0x1000_0000; + +/// The longest a control frame's payload may be. +const MAX_CTL: usize = 125; + +/// Which end of the connection this is. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Side { + /// A server: the client's frames must be masked, ours are not. + Server, +} + +/// What a message is. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Kind { + /// Text, which must be UTF-8. + Text, + /// Binary. + Binary, +} + +/// What the peer's bytes were. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Event<'a> { + /// A piece of a message: C's `LWS_CALLBACK_RECEIVE`. + Message { + /// What the message is. + kind: Kind, + /// The piece. + data: &'a [u8], + /// It starts the message. + first: bool, + /// It ends the message. + last: bool, + }, + /// A pong, with its payload: C's `LWS_CALLBACK_RECEIVE_PONG`. + Pong(&'a [u8]), + /// The peer's close, with its payload: C's + /// `LWS_CALLBACK_WS_PEER_INITIATED_CLOSE`. It is answered. + PeerClose(&'a [u8]), +} + +/// What [`Ws::rx`] took. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct Rx<'a> { + /// How many bytes it took. The rest are the caller's to hand in again. + pub consumed: usize, + /// What they were, if they came to something. + pub event: Option<Event<'a>>, +} + +/// What the connection asks of its carrier, once it is done with. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Close { + /// Stop sending once what was written has gone. + Shutdown, + /// Release it: the close handshake is over. + Release, +} + +/// Why a frame cannot be sent. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum SendError { + /// A frame is still going, or the connection is closing. + Busy, +} + +impl core::fmt::Display for SendError { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.write_str("a frame is still going, or the connection is closing") + } +} + +impl core::error::Error for SendError {} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Op { + Continuation, + Text, + Binary, + Close, + Ping, + Pong, +} + +impl Op { + const fn control(self) -> bool { + matches!(self, Self::Close | Self::Ping | Self::Pong) + } + + const fn code(self) -> u8 { + match self { + Self::Continuation => 0, + Self::Text => 1, + Self::Binary => 2, + Self::Close => 8, + Self::Ping => 9, + Self::Pong => 10, + } + } +} + +/// A frame's header as it has come so far. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct Frame { + op: Op, + fin: bool, + masked: bool, + len: u64, + mask: [u8; 4], +} + +/// Where the frame parser is: C's `lws_rx_parse_state`. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Parse { + /// The first byte of a frame. + First, + /// The length byte. + Len(Frame), + /// This many more bytes of an extended length. + LenMore(Frame, u8), + /// This many more bytes of the mask. + Mask(Frame, u8), + /// This much payload is still to come; `at` of it came so far. + Payload(Frame, u64), + /// Nothing more is read. + Stopped, +} + +/// Where a message is, between frames. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Msg { + /// Between messages. + Idle, + /// A message is under way, and its first piece has been given or not. + Open { kind: Kind, given: Given }, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Given { + Nothing, + Some, +} + +/// A control frame, its payload gathered or to be sent. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct Ctl { + buf: [u8; MAX_CTL], + len: u8, +} + +impl Ctl { + const fn new() -> Self { + Self { + buf: [0; MAX_CTL], + len: 0, + } + } + + fn payload(&self) -> &[u8] { + self.buf.get(..usize::from(self.len)).unwrap_or_default() + } + + fn push(&mut self, b: &[u8]) -> bool { + let at = usize::from(self.len); + let Some(end) = at.checked_add(b.len()).filter(|e| *e <= MAX_CTL) else { + return false; + }; + if let Some(d) = self.buf.get_mut(at..end) { + d.copy_from_slice(b); + } + self.len = u8::try_from(end).unwrap_or(0); + true + } + + fn close(code: u16, reason: &[u8]) -> Self { + let mut c = Self::new(); + let _ = c.push(&code.to_be_bytes()) && c.push(reason); + c + } +} + +/// Where the close is: C's close states, `lwsi_close()`. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Closing { + /// Open. + None, + /// Our close is to go: `LCS_WAITING_TO_SEND_CLOSE`. + WaitingToSend(Ctl), + /// Our close has gone: `LCS_AWAITING_CLOSE_ACK`. + AwaitingAck, + /// The peer's close is to be answered: `LCS_RETURNED_CLOSE`. + Returned(Ctl), + /// The application closes once what it sent has gone: + /// `LCS_FLUSHING_BEFORE_CLOSE`. + Flushing, + /// Done with. + Closed(Close), +} + +/// Bytes being written, and how many of them have gone. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct Out { + buf: [u8; 256], + len: usize, + sent: usize, +} + +impl Out { + const fn new() -> Self { + Self { + buf: [0; 256], + len: 0, + sent: 0, + } + } + + const fn pending(&self) -> bool { + self.sent < self.len + } + + fn set(&mut self, b: &[u8]) { + *self = Self::new(); + if let Some(d) = self.buf.get_mut(..b.len()) { + d.copy_from_slice(b); + self.len = b.len(); + } + } + + fn drain(&mut self, out: &mut [u8]) -> usize { + let rest = self.buf.get(self.sent..self.len).unwrap_or_default(); + let n = rest.len().min(out.len()); + if let (Some(d), Some(s)) = (out.get_mut(..n), rest.get(..n)) { + d.copy_from_slice(s); + } + self.sent = self.sent.saturating_add(n); + n + } +} + +/// The application's frame, its header written or not, its payload owed. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum App { + Idle, + Sending { owed: u64 }, +} + +/// A frame's header, unmasked, as a server writes it. +fn frame_header(op: Op, len: u64, out: &mut Out) { + // the length is 7 bits, or 126 and 16 bits, or 127 and 64 bits: the + // last of its big endian bytes, after the marker + let be = len.to_be_bytes(); + let (marker, extra) = match u8::try_from(len) { + Ok(short @ 0..126) => (short, 0), + Ok(_) | Err(_) if u16::try_from(len).is_ok() => (126, 2), + Ok(_) | Err(_) => (127, 8), + }; + let mut header = [0x80 | op.code(), marker, 0, 0, 0, 0, 0, 0, 0, 0]; + let used = 2usize.saturating_add(extra); + if let (Some(d), Some(s)) = ( + header.get_mut(2..used), + be.get(be.len().saturating_sub(extra)..), + ) { + d.copy_from_slice(s); + } + out.set(header.get(..used).unwrap_or_default()); +} + +/// A control frame, header and payload. +fn control_frame(op: Op, c: &Ctl, out: &mut Out) { + let p = c.payload(); + let mut b = [0u8; 2 + MAX_CTL]; + if let Some(h) = b.get_mut(..2) { + h.copy_from_slice(&[0x80 | op.code(), c.len]); + } + if let Some(d) = b.get_mut(2..2usize.saturating_add(p.len())) { + d.copy_from_slice(p); + } + out.set(b.get(..2usize.saturating_add(p.len())).unwrap_or_default()); +} + +/// One ws connection. +/// +/// ``` +/// use npro_ws::conn::{Event, Kind, Ws}; +/// +/// let mut ws = Ws::server(b""); +/// // a masked "Hi", with a zero mask +/// let mut frame = *b"\x81\x82\0\0\0\0Hi"; +/// let rx = ws.rx(&mut frame); +/// assert_eq!( +/// rx.event, +/// Some(Event::Message { kind: Kind::Text, data: b"Hi", first: true, last: true }) +/// ); +/// ``` +#[derive(Clone, Debug)] +pub struct Ws { + parse: Parse, + msg: Msg, + utf8: Utf8Validator, + /// A control frame's payload, gathered. + ctl: Ctl, + /// The pong owed, if any: C's one pending pong. + pong: Option<Ctl>, + closing: Closing, + out: Out, + app: App, +} + +impl Ws { + /// A server's connection, `first` being what goes before its frames: + /// the 101. + #[must_use] + pub fn server(first: &[u8]) -> Self { + let mut out = Out::new(); + out.set(first); + Self { + parse: Parse::First, + msg: Msg::Idle, + utf8: Utf8Validator::new(), + ctl: Ctl::new(), + pong: None, + closing: Closing::None, + out, + app: App::Idle, + } + } + + /// What the connection asks of its carrier, once it is done with. + #[must_use] + pub const fn close(&self) -> Option<Close> { + match self.closing { + Closing::Closed(c) => Some(c), + Closing::None + | Closing::WaitingToSend(_) + | Closing::AwaitingAck + | Closing::Returned(_) + | Closing::Flushing => None, + } + } + + /// Whether the connection has something of its own to write. + #[must_use] + pub const fn wants_write(&self) -> bool { + self.out.pending() + || self.pong.is_some() + || matches!( + self.closing, + Closing::WaitingToSend(_) | Closing::Returned(_) + ) + } + + /// Fails the connection with our close: C's `lws_close_reason()` and + /// `LWS_HPI_RET_PLEASE_CLOSE_ME`. + fn refuse<'a>(&mut self, code: u16, reason: &[u8]) -> Rx<'a> { + self.parse = Parse::Stopped; + if matches!(self.closing, Closing::None) { + self.closing = Closing::WaitingToSend(Ctl::close(code, reason)); + } + Rx { + consumed: 0, + event: None, + } + } + + /// Takes bytes from the peer: see [`Event`]. A frame's payload is + /// unmasked where it lies in `input`. + pub fn rx<'a>(&'a mut self, input: &'a mut [u8]) -> Rx<'a> { + let mut used = 0usize; + loop { + match self.parse { + Parse::Stopped => { + // after its close, nothing more is read + return Rx { + consumed: input.len(), + event: None, + }; + } + Parse::Payload(f, left) => return self.payload(f, left, input, used), + Parse::First | Parse::Len(_) | Parse::LenMore(..) | Parse::Mask(..) => {} + } + let Some(&c) = input.get(used) else { + return Rx { + consumed: used, + event: None, + }; + }; + used = used.saturating_add(1); + if let Some((code, reason)) = self.header(c) { + let mut r = self.refuse(code, reason); + r.consumed = used; + return r; + } + } + } + + /// One byte of a frame's header; `Some` refuses the frame. + fn header(&mut self, c: u8) -> Option<(u16, &'static [u8])> { + match self.parse { + Parse::First => { + let fin = c & 0x80 != 0; + let op = match c & 0x0f { + 0 => Op::Continuation, + 1 => Op::Text, + 2 => Op::Binary, + 8 => Op::Close, + 9 => Op::Ping, + 10 => Op::Pong, + _ => { + if c & 0x08 != 0 && !fin { + return Some((1002, b"frag ctl")); + } + return Some((1002, b"bad opc")); + } + }; + if op.control() && !fin { + return Some((1002, b"frag ctl")); + } + match (op, self.msg) { + (Op::Text | Op::Binary, Msg::Open { .. }) | (Op::Continuation, Msg::Idle) => { + return Some((1002, b"bad cont")); + } + (Op::Text, Msg::Idle) => { + self.utf8 = Utf8Validator::new(); + self.msg = Msg::Open { + kind: Kind::Text, + given: Given::Nothing, + }; + } + (Op::Binary, Msg::Idle) => { + self.msg = Msg::Open { + kind: Kind::Binary, + given: Given::Nothing, + }; + } + (Op::Continuation, Msg::Open { .. }) + | (Op::Close | Op::Ping | Op::Pong, Msg::Idle | Msg::Open { .. }) => {} + } + if c & 0x70 != 0 { + return Some((1002, b"rsv bits")); + } + self.parse = Parse::Len(Frame { + op, + fin, + masked: false, + len: 0, + mask: [0; 4], + }); + } + Parse::Len(mut f) => { + f.masked = c & 0x80 != 0; + if !f.masked { + return Some((1002, b"client unmasked")); + } + match c & 0x7f { + 126 | 127 if f.op.control() => return Some((1002, b"ctl len")), + 126 => self.parse = Parse::LenMore(f, 2), + 127 => self.parse = Parse::LenMore(f, 8), + n => { + f.len = u64::from(n); + self.parse = Parse::Mask(f, 4); + } + } + } + Parse::LenMore(mut f, left) => { + if left == 8 && c & 0x80 != 0 { + return Some((1002, b"bad len")); + } + f.len = (f.len << 8) | u64::from(c); + let left = left.saturating_sub(1); + if left > 0 { + self.parse = Parse::LenMore(f, left); + } else if f.len > MAX_FRAME { + return Some((1009, b"huge frame")); + } else { + self.parse = Parse::Mask(f, 4); + } + } + Parse::Mask(mut f, left) => { + let i = usize::from(4u8.saturating_sub(left)); + if let Some(m) = f.mask.get_mut(i) { + *m = c; + } + let left = left.saturating_sub(1); + self.parse = if left > 0 { + Parse::Mask(f, left) + } else { + self.ctl = Ctl::new(); + Parse::Payload(f, f.len) + }; + } + Parse::Payload(..) | Parse::Stopped => {} + } + None + } + + /// The payload of `f`, `left` of it still to come. + fn payload<'a>(&'a mut self, f: Frame, left: u64, input: &'a mut [u8], used: usize) -> Rx<'a> { + let rest = input.get_mut(used..).unwrap_or_default(); + let n = usize::try_from(left).unwrap_or(usize::MAX).min(rest.len()); + let at = f.len.saturating_sub(left); + let piece = rest.get_mut(..n).unwrap_or_default(); + for (i, b) in piece.iter_mut().enumerate() { + let k = at.wrapping_add(u64::try_from(i).unwrap_or(0)) % 4; + *b ^= f + .mask + .get(usize::try_from(k).unwrap_or(0)) + .copied() + .unwrap_or(0); + } + let left = left.saturating_sub(u64::try_from(n).unwrap_or(left)); + let consumed = used.saturating_add(n); + if f.op.control() { + let _ = self.ctl.push(piece); + if left > 0 { + self.parse = Parse::Payload(f, left); + return Rx { + consumed, + event: None, + }; + } + self.parse = Parse::First; + let event = self.control(f.op); + return Rx { consumed, event }; + } + self.parse = if left > 0 { + Parse::Payload(f, left) + } else { + Parse::First + }; + // a piece of a message: given at the frame's end, or as it comes + if n == 0 && left > 0 { + return Rx { + consumed, + event: None, + }; + } + let Msg::Open { kind, given } = self.msg else { + return Rx { + consumed, + event: None, + }; + }; + let last = f.fin && left == 0; + if kind == Kind::Text { + if self.utf8.feed(piece).is_err() { + let mut r = self.refuse(1007, b"bad utf8"); + r.consumed = consumed; + return r; + } + if last && !self.utf8.at_boundary() { + let mut r = self.refuse(1007, b"partial utf8"); + r.consumed = consumed; + return r; + } + } + self.msg = if last { + Msg::Idle + } else { + Msg::Open { + kind, + given: Given::Some, + } + }; + // nothing for the app once a close is under way + if !matches!(self.closing, Closing::None) { + return Rx { + consumed, + event: None, + }; + } + Rx { + consumed, + event: Some(Event::Message { + kind, + data: piece, + first: given == Given::Nothing, + last, + }), + } + } + + /// A whole control frame has come. + fn control(&mut self, op: Op) -> Option<Event<'_>> { + match op { + Op::Ping => { + // one pong owed at a time: a second ping is dropped + if self.pong.is_none() { + self.pong = Some(self.ctl); + } + None + } + Op::Pong => (self.ctl.len > 0).then(|| Event::Pong(self.ctl.payload())), + Op::Close => self.peer_close(), + Op::Continuation | Op::Text | Op::Binary => None, + } + } + + /// The peer's close: C's handling of `LWSWSOPC_CLOSE`. + fn peer_close(&mut self) -> Option<Event<'_>> { + match self.closing { + // a second close changes nothing; nor is one answered while + // the app's last goes + Closing::Returned(_) | Closing::Flushing | Closing::Closed(_) => None, + // the answer to ours: done + Closing::AwaitingAck | Closing::WaitingToSend(_) => { + self.closing = Closing::Closed(Close::Release); + self.parse = Parse::Stopped; + None + } + Closing::None => { + if self.ctl.len >= 2 { + let code = u16::from_be_bytes([ + self.ctl.buf.first().copied().unwrap_or(0), + self.ctl.buf.get(1).copied().unwrap_or(0), + ]); + // a code no peer may send is answered as a protocol error + if code < 1000 + || matches!(code, 1004..=1006 | 1012..=1015) + || (1016..3000).contains(&code) + { + if let Some(b) = self.ctl.buf.get_mut(..2) { + b.copy_from_slice(&1002u16.to_be_bytes()); + } + } + } + self.closing = Closing::Returned(self.ctl); + // after the peer's close, nothing more is read + self.parse = Parse::Stopped; + Some(Event::PeerClose(self.ctl.payload())) + } + } + } + + /// Commits a whole message, `len` bytes, whose payload [`Ws::tx`] pulls: + /// C's `lws_write()` of a final frame. + /// + /// # Errors + /// + /// [`SendError::Busy`] while a frame is still going, or the connection + /// is closing. + pub fn send(&mut self, kind: Kind, len: u64) -> Result<(), SendError> { + if self.app != App::Idle || self.out.pending() || !matches!(self.closing, Closing::None) { + return Err(SendError::Busy); + } + let op = match kind { + Kind::Text => Op::Text, + Kind::Binary => Op::Binary, + }; + frame_header(op, len, &mut self.out); + self.app = App::Sending { owed: len }; + Ok(()) + } + + /// The application is done: the connection closes once what it sent + /// has gone, without a close frame, as C's + /// `lws_raw_transaction_completed()`. + pub const fn close_when_flushed(&mut self) { + if matches!(self.closing, Closing::None) { + self.closing = Closing::Flushing; + } + } + + /// Writes what is owed the peer into `out`: see the module's + /// description. + pub fn tx(&mut self, out: &mut [u8], src: &mut dyn TxSource) -> usize { + let mut written = 0usize; + loop { + let room = out.get_mut(written..).unwrap_or_default(); + if room.is_empty() { + return written; + } + // what is in flight goes first + if self.out.pending() { + written = written.saturating_add(self.out.drain(room)); + continue; + } + if let App::Sending { owed } = self.app { + let cap = usize::try_from(owed).unwrap_or(usize::MAX).min(room.len()); + let n = src.fill(room.get_mut(..cap).unwrap_or_default()).min(cap); + written = written.saturating_add(n); + let owed = owed.saturating_sub(u64::try_from(n).unwrap_or(owed)); + self.app = if owed == 0 { + App::Idle + } else { + App::Sending { owed } + }; + if owed > 0 { + // the rest of the payload is not here yet + return written; + } + continue; + } + match self.closing { + Closing::WaitingToSend(c) => { + control_frame(Op::Close, &c, &mut self.out); + self.closing = Closing::AwaitingAck; + continue; + } + Closing::Flushing => { + self.closing = Closing::Closed(Close::Shutdown); + return written; + } + Closing::None + | Closing::AwaitingAck + | Closing::Returned(_) + | Closing::Closed(_) => {} + } + // the pong goes while open, or ahead of the answer to the + // peer's close, its ping having come first (RFC 6455 5.5.2); if + // we began the close, it is forgotten + if let Some(p) = self.pong.take() { + match self.closing { + Closing::None | Closing::Returned(_) => { + control_frame(Op::Pong, &p, &mut self.out); + } + Closing::WaitingToSend(_) + | Closing::AwaitingAck + | Closing::Flushing + | Closing::Closed(_) => {} + } + continue; + } + if let Closing::Returned(c) = self.closing { + control_frame(Op::Close, &c, &mut self.out); + self.closing = Closing::Closed(Close::Shutdown); + continue; + } + return written; + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + struct Nothing; + impl TxSource for Nothing { + fn fill(&mut self, _: &mut [u8]) -> usize { + 0 + } + } + + /// What the server writes after taking `frames`. + fn answers(frames: &[u8]) -> ([u8; 64], usize) { + let mut ws = Ws::server(b""); + let mut input = frames.to_vec(); + let mut at = 0; + while at < input.len() { + let rx = ws.rx(&mut input[at..]); + if rx.consumed == 0 { + break; + } + at = at.checked_add(rx.consumed).unwrap(); + } + let mut out = [0u8; 64]; + let n = ws.tx(&mut out, &mut Nothing); + (out, n) + } + + #[test] + fn refusals_close_with_cs_codes() { + for (frames, close) in [ + (&b"\x83\x80\0\0\0\0"[..], &b"\x88\x09\x03\xeabad opc"[..]), + (b"\x09\x80\0\0\0\0", b"\x88\x0a\x03\xeafrag ctl"), + (b"\x80\x80\0\0\0\0", b"\x88\x0a\x03\xeabad cont"), + (b"\xc1\x80\0\0\0\0", b"\x88\x0a\x03\xearsv bits"), + (b"\x81\x00", b"\x88\x11\x03\xeaclient unmasked"), + (b"\x89\xfe\0\x7e", b"\x88\x09\x03\xeactl len"), + (b"\x82\xff\x80", b"\x88\x09\x03\xeabad len"), + (b"\x81\x81\0\0\0\0\xff", b"\x88\x0a\x03\xefbad utf8"), + (b"\x81\x81\0\0\0\0\xc3", b"\x88\x0e\x03\xefpartial utf8"), + ] { + let (out, n) = answers(frames); + assert_eq!(&out[..n], close, "{}", frames.escape_ascii()); + } + } + + #[test] + fn a_reserved_close_code_is_answered_as_1002() { + let (out, n) = answers(b"\x88\x82\0\0\0\0\x03\xed"); + assert_eq!(&out[..n], b"\x88\x02\x03\xea"); + } + + #[test] + fn a_message_in_pieces_says_its_first_and_last() { + let mut ws = Ws::server(b""); + let mut a = *b"\x01\x81\0\0\0\0a"; + assert_eq!( + ws.rx(&mut a).event, + Some(Event::Message { + kind: Kind::Text, + data: b"a", + first: true, + last: false + }) + ); + let mut b = *b"\x80\x81\0\0\0\0b"; + assert_eq!( + ws.rx(&mut b).event, + Some(Event::Message { + kind: Kind::Text, + data: b"b", + first: false, + last: true + }) + ); + } + + #[test] + fn a_pong_owed_is_forgotten_once_we_close() { + // a ping, then a reserved opcode we refuse + let (out, n) = answers(b"\x89\x81\0\0\0\0a\x83\x80\0\0\0\0"); + assert_eq!(&out[..n], b"\x88\x09\x03\xeabad opc"); + } + + #[test] + fn a_header_takes_the_shortest_length_form() { + for (len, want) in [ + (125, &b"\x82\x7d"[..]), + (126, b"\x82\x7e\x00\x7e"), + (0xffff, b"\x82\x7e\xff\xff"), + (0x1_0000, b"\x82\x7f\0\0\0\0\0\x01\0\0"), + ] { + let mut out = Out::new(); + frame_header(Op::Binary, len, &mut out); + assert_eq!(&out.buf[..out.len], want, "{len}"); + } + } + + #[test] + fn a_second_ping_while_a_pong_is_owed_is_dropped() { + let (out, n) = answers(b"\x89\x81\0\0\0\0a\x89\x81\0\0\0\0b"); + assert_eq!(&out[..n], b"\x8a\x01a"); + } +} diff --git a/crates/npro-ws/src/handshake.rs b/crates/npro-ws/src/handshake.rs new file mode 100644 index 0000000..70066ef --- /dev/null +++ b/crates/npro-ws/src/handshake.rs @@ -0,0 +1,267 @@ +//! A server's side of the ws handshake: C's `lws_process_ws_upgrade()` and +//! `handshake_0405()`. +//! +//! C's checks, in C's order, each refusal C's status: an upgrade is a GET; +//! its `Connection` names the token `upgrade`; it has a key, of less than +//! 128 bytes, and a Host; its version is `13`, a 400 without one and a 426 +//! (saying `sec-websocket-version: 13`) for another; and it asks for a +//! subprotocol the server has, the first of its list that it has, or with +//! no list, the server's default. + +use npro_core::base64; +use npro_core::sha1::Sha1; +use npro_h1::table::HeaderTable; +use npro_h1::token::Token; + +/// The GUID RFC 6455 4.2.2 appends to a key. +const GUID: &[u8] = b"258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; + +/// The longest key C takes, and one: C's `MAX_WEBSOCKET_04_KEY_LEN`. +const MAX_KEY: usize = 128; + +/// The length of an accept value: base64 of a SHA-1. +pub const ACCEPT_LEN: usize = 28; + +/// An upgrade the server takes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct Accepted { + /// The index of the subprotocol in the server's list. + pub protocol: usize, + /// Whether the request named it, so the 101 says it. + pub named: bool, + accept: [u8; ACCEPT_LEN], +} + +impl Accepted { + /// The `Sec-WebSocket-Accept` value. + #[must_use] + pub const fn accept(&self) -> &[u8; ACCEPT_LEN] { + &self.accept + } +} + +/// Why an upgrade is refused, as C says why. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Refusal { + /// Not a GET (400). + NotGet, + /// No `upgrade` in `Connection` (400). + NoConnectionUpgrade, + /// No key, too long a key, or no Host (400). + KeyOrHost, + /// No version (400). + NoVersion, + /// A version that is not 13 (426). + Version, + /// A protocol list that is not one (400). + ProtocolList, + /// No protocol the server has (400). + NoProtocol, +} + +impl Refusal { + /// The status C answers with. + #[must_use] + pub const fn status(self) -> u16 { + match self { + Self::Version => 426, + Self::NotGet + | Self::NoConnectionUpgrade + | Self::KeyOrHost + | Self::NoVersion + | Self::ProtocolList + | Self::NoProtocol => 400, + } + } + + /// The header C's 426 adds to its status page. + #[must_use] + pub const fn header(self) -> Option<(&'static [u8], &'static [u8])> { + match self { + Self::Version => Some((b"sec-websocket-version", b"13")), + Self::NotGet + | Self::NoConnectionUpgrade + | Self::KeyOrHost + | Self::NoVersion + | Self::ProtocolList + | Self::NoProtocol => None, + } + } +} + +/// Whether `c` is a token's byte, RFC 9110 5.6.2's tchar. +const fn tchar(c: u8) -> bool { + c.is_ascii_alphanumeric() + || matches!( + c, + b'!' | b'#' + | b'$' + | b'%' + | b'&' + | b'\'' + | b'*' + | b'+' + | b'-' + | b'.' + | b'^' + | b'_' + | b'`' + | b'|' + | b'~' + ) +} + +/// The elements of a comma separated list of tokens, each `Some` token, or +/// `None` for one that is not a token: what C's `lws_tokenize()` takes for +/// these lists. +fn tokens(v: &[u8]) -> impl Iterator<Item = Option<&[u8]>> { + v.split(|c| *c == b',').filter_map(|e| { + let blank = |c: &u8| *c == b' ' || *c == b'\t'; + let start = e.iter().position(|c| !blank(c))?; + let end = e + .iter() + .rposition(|c| !blank(c)) + .map_or(0, |n| n.saturating_add(1)); + let t = e.get(start..end)?; + Some(t.iter().all(|c| tchar(*c)).then_some(t)) + }) +} + +/// Checks an upgrade request against a server having the subprotocols +/// `protocols`, and `default` for a request naming none. +/// +/// ``` +/// use npro_h1::head::{Config, Head, Side}; +/// use npro_ws::handshake::server; +/// +/// let mut h = Head::new([0u8; 1024], Side::Server, Config::new())?; +/// h.rx(b"GET /chat HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\n\ +/// Connection: Upgrade\r\nSec-WebSocket-Version: 13\r\n\ +/// Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n\r\n")?; +/// let a = server(h.table(), &[b"chat"], Some(0)).unwrap(); +/// assert_eq!(a.accept(), b"s3pPLMBiTxaQ9kYGzzhZRbK+xOo="); +/// # Ok::<(), Box<dyn core::error::Error>>(()) +/// ``` +/// +/// # Errors +/// +/// The [`Refusal`], whose status and header C answers with. +pub fn server<S: AsRef<[u8]> + AsMut<[u8]>>( + t: &HeaderTable<S>, + protocols: &[&[u8]], + default: Option<usize>, +) -> Result<Accepted, Refusal> { + if !t.is_present(Token::GetUri) { + return Err(Refusal::NotGet); + } + let mut buf = [0u8; MAX_KEY]; + let conn = t + .copy( + Token::Connection, + buf.get_mut(..MAX_KEY - 1).unwrap_or_default(), + ) + .ok() + .filter(|n| *n > 0) + .and_then(|n| buf.get(..n)) + .ok_or(Refusal::NoConnectionUpgrade)?; + let mut upgrade = false; + for tok in tokens(conn) { + match tok { + Some(name) if name.eq_ignore_ascii_case(b"upgrade") => { + upgrade = true; + break; + } + Some(_) => {} + None => return Err(Refusal::NoConnectionUpgrade), + } + } + if !upgrade { + return Err(Refusal::NoConnectionUpgrade); + } + + let key_len = t.total_len(Token::WsKey); + if key_len == 0 || key_len >= MAX_KEY || !t.is_present(Token::Host) { + return Err(Refusal::KeyOrHost); + } + let mut key = [0u8; MAX_KEY]; + let key = t + .copy(Token::WsKey, &mut key) + .ok() + .and_then(|n| key.get(..n)) + .ok_or(Refusal::KeyOrHost)?; + + match t.copy(Token::WsVersion, &mut buf) { + Ok(0) => return Err(Refusal::NoVersion), + Ok(2) if buf.get(..2) == Some(b"13".as_slice()) => {} + Ok(_) | Err(_) => return Err(Refusal::Version), + } + + let mut list = [0u8; MAX_KEY]; + let list = t + .copy( + Token::WsProtocol, + list.get_mut(..MAX_KEY - 1).unwrap_or_default(), + ) + .ok() + .and_then(|n| list.get(..n)) + .ok_or(Refusal::ProtocolList)?; + let (protocol, named) = if list.is_empty() { + let d = default + .filter(|d| *d < protocols.len()) + .ok_or(Refusal::NoProtocol)?; + (d, false) + } else { + let mut found = None; + for tok in tokens(list) { + let name = tok.ok_or(Refusal::ProtocolList)?; + if name.len() >= 64 { + return Err(Refusal::ProtocolList); + } + if let Some(i) = protocols.iter().position(|p| *p == name) { + found = Some(i); + break; + } + } + (found.ok_or(Refusal::NoProtocol)?, true) + }; + + let mut h = Sha1::new(); + h.update(key); + h.update(GUID); + let mut accept = [0u8; ACCEPT_LEN]; + base64::encode(&h.finish(), &mut accept).map_err(|_| Refusal::KeyOrHost)?; + Ok(Accepted { + protocol, + named, + accept, + }) +} + +/// The most a 101 C writes may have here. +pub const MAX_101: usize = 256; + +/// C's 101 for `a`, `name` being its subprotocol's: written into `out`, +/// returning how much of it. +/// +/// # Errors +/// +/// `None` if `out` is too small. +#[must_use] +pub fn response_101(a: &Accepted, name: &[u8], out: &mut [u8]) -> Option<usize> { + let mut at = 0usize; + let mut put = |b: &[u8]| -> Option<()> { + let end = at.checked_add(b.len())?; + out.get_mut(at..end)?.copy_from_slice(b); + at = end; + Some(()) + }; + put(b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: WebSocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: ")?; + put(&a.accept)?; + // the protocol is said only if the request named one, and it has a name + if a.named && !name.is_empty() { + put(b"\r\nSec-WebSocket-Protocol: ")?; + put(name)?; + } + put(b"\r\n\r\n")?; + Some(at) +} diff --git a/crates/npro-ws/src/lib.rs b/crates/npro-ws/src/lib.rs new file mode 100644 index 0000000..483678d --- /dev/null +++ b/crates/npro-ws/src/lib.rs @@ -0,0 +1,14 @@ +//! npro-ws: websockets for npro, sans-IO. +//! +//! The port of C libwebsockets' ws role (`lib/sansio/ws`): +//! +//! - [`handshake`]: a server's checks of an upgrade request, and its 101; +//! - [`conn`]: a ws connection, its frames in and out, and its close. +//! +//! The client's handshake comes next, in phase 1e of the port plan. + +#![no_std] +#![forbid(unsafe_code)] + +pub mod conn; +pub mod handshake; diff --git a/deny.toml b/deny.toml index d19e6b0..e3e32e3 100644 --- a/deny.toml +++ b/deny.toml @@ -39,6 +39,7 @@ allow = [ "npro", "npro-core", "npro-h1", + "npro-ws", "npro-fuzz", "npro-test", diff --git a/scripts/ci.sh b/scripts/ci.sh index 6210e53..5ea589d 100755 --- a/scripts/ci.sh +++ b/scripts/ci.sh @@ -43,7 +43,7 @@ cargo "+$msrv" check --workspace --all-targets --all-features --locked echo "== no_std" # the sans-IO crates must build for a target with no std at all -for c in npro-core npro-h1; do +for c in npro-core npro-h1 npro-ws; do cargo build -p "$c" --all-features --target thumbv7em-none-eabihf --locked done diff --git a/scripts/sai.sh b/scripts/sai.sh index a5d38b8..6d0d362 100755 --- a/scripts/sai.sh +++ b/scripts/sai.sh @@ -32,7 +32,7 @@ export PATH jobs="${SAI_PARALLEL:-4}" # the sans-IO crates, which must build with no std at all -nostd_crates="npro-core npro-h1" +nostd_crates="npro-core npro-h1 npro-ws" nostd_targets="thumbv6m-none-eabi thumbv7em-none-eabihf riscv32imc-unknown-none-elf" . scripts/require.sh
Page fetched 0s ago, creation time: 5ms (vhost etag hits: 0%, cache hits: 0%)