Skip to main content

rig_http/http_client/
framing.rs

1//! Incremental SSE and newline-delimited framing of response bytes.
2//! SSE dispatch requires a blank line; incomplete events are not flushed at EOF.
3//! NDJSON permits a final unterminated line through [`NdjsonFramer::finish`].
4//!
5//! ```
6//! use rig_http::http_client::framing::NdjsonFramer;
7//!
8//! let mut framer = NdjsonFramer::new();
9//! assert_eq!(framer.push(b"{}\n").collect::<Vec<_>>(), vec![b"{}".to_vec()]);
10//! ```
11
12/// How a reply's bytes are split into wire frames.
13#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub enum Framing {
15    /// `text/event-stream`: WHATWG event-stream grammar ([`SseFramer`]).
16    Sse,
17    /// Newline-delimited JSON ([`NdjsonFramer`]).
18    Ndjson,
19    /// The whole body is one frame.
20    Whole,
21}
22
23impl Framing {
24    /// Split a whole reply body into the frames a decoder reads: one per SSE
25    /// event with a non-blank `data` field, one per NDJSON line, or the
26    /// whole body.
27    ///
28    /// ```
29    /// use rig_http::http_client::framing::{Framing, WireFrame};
30    ///
31    /// let frames = Framing::Sse.split(b"data: {}\n\ndata: [DONE]\n\n");
32    /// assert_eq!(frames.len(), 2);
33    /// assert_eq!(frames[1].as_str(), "[DONE]");
34    /// ```
35    pub fn split(&self, bytes: &[u8]) -> Vec<WireFrame> {
36        match self {
37            Self::Sse => SseFramer::new()
38                .push(bytes)
39                .filter(|event| !event.data.trim().is_empty())
40                .map(|event| WireFrame::Text(event.data))
41                .collect(),
42            Self::Ndjson => {
43                let mut framer = NdjsonFramer::new();
44                let mut lines: Vec<Vec<u8>> = framer.push(bytes).collect();
45                lines.extend(framer.finish());
46                lines.into_iter().map(WireFrame::from_bytes).collect()
47            }
48            Self::Whole if bytes.is_empty() => Vec::new(),
49            Self::Whole => vec![WireFrame::from_bytes(bytes.to_vec())],
50        }
51    }
52}
53
54/// One transport frame, after framing but before decoding.
55///
56/// The transport layer (SSE framer, NDJSON splitter, websocket reader) owns
57/// byte splitting and yields these; a decoder never splits bytes.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum WireFrame {
60    /// Text payload, such as an SSE `data:` field or WebSocket message.
61    Text(String),
62    /// Byte payload, such as an NDJSON line or binary frame.
63    Bytes(Vec<u8>),
64}
65
66impl WireFrame {
67    /// The frame payload as text (lossy for byte frames).
68    pub fn as_str(&self) -> std::borrow::Cow<'_, str> {
69        match self {
70            Self::Text(text) => std::borrow::Cow::Borrowed(text),
71            Self::Bytes(bytes) => String::from_utf8_lossy(bytes),
72        }
73    }
74
75    /// A frame of `bytes`: text when they are UTF-8.
76    pub fn from_bytes(bytes: Vec<u8>) -> Self {
77        match String::from_utf8(bytes) {
78            Ok(text) => Self::Text(text),
79            Err(error) => Self::Bytes(error.into_bytes()),
80        }
81    }
82}
83
84/// One dispatched `text/event-stream` event.
85#[derive(Debug, Clone, Default, PartialEq, Eq)]
86pub struct SseEvent {
87    /// The `event:` field, `"message"` when the stream named none.
88    pub event: String,
89    /// The joined `data:` fields, without the terminating newline.
90    pub data: String,
91    /// The stream's last event id, empty when none was ever set.
92    pub id: String,
93    /// Most recent parsed `retry:` value since the previous emitted event.
94    /// This value does not trigger automatic reconnection.
95    pub retry: Option<u64>,
96}
97
98/// The optional leading byte-order mark, which the grammar ignores.
99const BOM: [u8; 3] = [0xef, 0xbb, 0xbf];
100
101/// A push parser for the WHATWG `text/event-stream` grammar.
102#[derive(Debug, Default)]
103pub struct SseFramer {
104    /// Bytes not yet consumed as a complete line.
105    buffer: Vec<u8>,
106    /// Events completed by the current `push`.
107    ready: Vec<SseEvent>,
108    event_type: String,
109    data: String,
110    last_event_id: String,
111    retry: Option<u64>,
112    /// How many BOM bytes have matched so far, while the prefix is undecided.
113    bom_prefix: usize,
114    bom_done: bool,
115    /// Bytes of complete lines consumed since the last blank line.
116    since_dispatch: usize,
117}
118
119impl SseFramer {
120    /// A framer at the start of a stream.
121    pub fn new() -> Self {
122        Self::default()
123    }
124
125    /// Feed one transport chunk; yields every event it completed.
126    pub fn push(&mut self, chunk: &[u8]) -> std::vec::Drain<'_, SseEvent> {
127        let chunk = self.strip_bom(chunk);
128        self.buffer.extend_from_slice(chunk);
129        while let Some((line_len, consumed)) = terminated_line(&self.buffer) {
130            let line = String::from_utf8_lossy(self.buffer.get(..line_len).unwrap_or_default())
131                .into_owned();
132            self.buffer.drain(..consumed);
133            self.since_dispatch = self.since_dispatch.saturating_add(consumed);
134            self.line(&line);
135        }
136        self.ready.drain(..)
137    }
138
139    /// Saturating count of buffered bytes and complete lines since the last
140    /// blank line. Excludes any undecided leading BOM prefix; exposes no frame data.
141    pub fn pending(&self) -> usize {
142        self.since_dispatch.saturating_add(self.buffer.len())
143    }
144
145    /// Consume the stream's optional leading BOM, which may itself be split
146    /// across chunks. Returns the chunk with any BOM bytes removed.
147    fn strip_bom<'a>(&mut self, chunk: &'a [u8]) -> &'a [u8] {
148        if self.bom_done {
149            return chunk;
150        }
151        let mut rest = chunk;
152        while let Some(&byte) = rest.first() {
153            if BOM.get(self.bom_prefix) != Some(&byte) {
154                // Not a BOM after all: replay the bytes that matched so far.
155                let matched = std::mem::take(&mut self.bom_prefix);
156                self.bom_done = true;
157                self.buffer
158                    .extend_from_slice(BOM.get(..matched).unwrap_or_default());
159                return rest;
160            }
161            self.bom_prefix += 1;
162            rest = rest.get(1..).unwrap_or_default();
163            if self.bom_prefix == BOM.len() {
164                self.bom_prefix = 0;
165                self.bom_done = true;
166                return rest;
167            }
168        }
169        rest
170    }
171
172    /// Process one complete line of the grammar.
173    fn line(&mut self, line: &str) {
174        if line.is_empty() {
175            self.dispatch();
176            return;
177        }
178        // A line starting with a colon is a comment.
179        if let Some(rest) = line.strip_prefix(':') {
180            let _ = rest;
181            return;
182        }
183        let (field, value) = match line.split_once(':') {
184            Some((field, value)) => (field, value.strip_prefix(' ').unwrap_or(value)),
185            None => (line, ""),
186        };
187        match field {
188            "event" => {
189                self.event_type.clear();
190                self.event_type.push_str(value);
191            }
192            "data" => {
193                self.data.push_str(value);
194                self.data.push('\n');
195            }
196            // A NUL in an id is ignored outright, per the grammar.
197            "id" if !value.contains('\0') => {
198                self.last_event_id.clear();
199                self.last_event_id.push_str(value);
200            }
201            "retry" if !value.is_empty() && value.bytes().all(|b| b.is_ascii_digit()) => {
202                self.retry = value.parse().ok();
203            }
204            _ => {}
205        }
206    }
207
208    /// A blank line ends an event. An event with no data is not dispatched:
209    /// it only resets the event type, as the grammar says.
210    fn dispatch(&mut self) {
211        self.since_dispatch = 0;
212        if self.data.is_empty() {
213            self.event_type.clear();
214            return;
215        }
216        self.data.pop();
217        let event = if self.event_type.is_empty() {
218            "message".to_owned()
219        } else {
220            std::mem::take(&mut self.event_type)
221        };
222        self.ready.push(SseEvent {
223            event,
224            data: std::mem::take(&mut self.data),
225            id: self.last_event_id.clone(),
226            retry: self.retry.take(),
227        });
228        self.event_type.clear();
229    }
230}
231
232/// Splits nonempty newline-delimited payloads without validating JSON.
233/// A nonempty unterminated last line is available through [`Self::finish`].
234#[derive(Debug, Default)]
235pub struct NdjsonFramer {
236    buffer: Vec<u8>,
237    ready: Vec<Vec<u8>>,
238}
239
240impl NdjsonFramer {
241    /// A framer at the start of a stream.
242    pub fn new() -> Self {
243        Self::default()
244    }
245
246    /// Feed one transport chunk; yields every line it completed. Blank lines
247    /// are not frames.
248    pub fn push(&mut self, chunk: &[u8]) -> std::vec::Drain<'_, Vec<u8>> {
249        self.buffer.extend_from_slice(chunk);
250        while let Some(pos) = self.buffer.iter().position(|byte| *byte == b'\n') {
251            let mut line: Vec<u8> = self.buffer.drain(..=pos).collect();
252            line.pop();
253            if line.last() == Some(&b'\r') {
254                line.pop();
255            }
256            if !line.is_empty() {
257                self.ready.push(line);
258            }
259        }
260        self.ready.drain(..)
261    }
262
263    /// The unterminated last line, if the stream ended with one.
264    pub fn finish(&mut self) -> Option<Vec<u8>> {
265        let line = std::mem::take(&mut self.buffer);
266        (!line.is_empty()).then_some(line)
267    }
268
269    /// Bytes of an unterminated trailing line, for truncation diagnostics.
270    pub fn pending(&self) -> usize {
271        self.buffer.len()
272    }
273}
274
275/// The first complete line in `buffer` as `(line length, bytes consumed)`.
276///
277/// Defers a trailing CR until another byte arrives so a split CRLF is consumed
278/// as one terminator.
279fn terminated_line(buffer: &[u8]) -> Option<(usize, usize)> {
280    let pos = buffer
281        .iter()
282        .position(|byte| *byte == b'\n' || *byte == b'\r')?;
283    match buffer.get(pos) {
284        Some(b'\n') => Some((pos, pos + 1)),
285        Some(b'\r') => match buffer.get(pos + 1) {
286            Some(b'\n') => Some((pos, pos + 2)),
287            Some(_) => Some((pos, pos + 1)),
288            None => None,
289        },
290        _ => None,
291    }
292}
293
294#[cfg(test)]
295mod tests;