rig_http/http_client/
framing.rs1#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub enum Framing {
15 Sse,
17 Ndjson,
19 Whole,
21}
22
23impl Framing {
24 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#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum WireFrame {
60 Text(String),
62 Bytes(Vec<u8>),
64}
65
66impl WireFrame {
67 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 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#[derive(Debug, Clone, Default, PartialEq, Eq)]
86pub struct SseEvent {
87 pub event: String,
89 pub data: String,
91 pub id: String,
93 pub retry: Option<u64>,
96}
97
98const BOM: [u8; 3] = [0xef, 0xbb, 0xbf];
100
101#[derive(Debug, Default)]
103pub struct SseFramer {
104 buffer: Vec<u8>,
106 ready: Vec<SseEvent>,
108 event_type: String,
109 data: String,
110 last_event_id: String,
111 retry: Option<u64>,
112 bom_prefix: usize,
114 bom_done: bool,
115 since_dispatch: usize,
117}
118
119impl SseFramer {
120 pub fn new() -> Self {
122 Self::default()
123 }
124
125 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 pub fn pending(&self) -> usize {
142 self.since_dispatch.saturating_add(self.buffer.len())
143 }
144
145 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 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 fn line(&mut self, line: &str) {
174 if line.is_empty() {
175 self.dispatch();
176 return;
177 }
178 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 "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 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#[derive(Debug, Default)]
235pub struct NdjsonFramer {
236 buffer: Vec<u8>,
237 ready: Vec<Vec<u8>>,
238}
239
240impl NdjsonFramer {
241 pub fn new() -> Self {
243 Self::default()
244 }
245
246 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 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 pub fn pending(&self) -> usize {
271 self.buffer.len()
272 }
273}
274
275fn 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;