use crate::canonical::CanonicalError;
use super::frame::{Decoder, Frame, Framing};
impl Framing {
pub fn decoder(self) -> Box<dyn Decoder> {
match self {
Framing::Sse => Box::new(SseDecoder::default()),
Framing::Ndjson => Box::new(NdjsonDecoder::default()),
Framing::Identity => Box::new(IdentityDecoder),
}
}
}
const BOM: &[u8] = &[0xEF, 0xBB, 0xBF];
#[derive(Default)]
pub struct SseDecoder {
buf: Vec<u8>,
scan: usize,
bom_stripped: bool,
}
impl Decoder for SseDecoder {
fn push(&mut self, chunk: Vec<u8>) -> Result<Vec<Frame>, CanonicalError> {
self.buf.extend_from_slice(&chunk);
if !self.strip_bom() {
return Ok(Vec::new()); }
let mut frames = Vec::new();
while let Some(rel) = find_frame_end(&self.buf[self.scan..]) {
let end = self.scan + rel;
let block: Vec<u8> = self.buf.drain(..end).collect();
self.scan = 0; if let Some(frame) = parse_block(&block) {
frames.push(frame);
}
}
self.scan = self.buf.len().saturating_sub(3);
Ok(frames)
}
fn finish(&mut self) -> Result<Vec<Frame>, CanonicalError> {
let block = std::mem::take(&mut self.buf);
Ok(parse_block(&block).into_iter().collect())
}
}
impl SseDecoder {
fn strip_bom(&mut self) -> bool {
if self.bom_stripped {
return true;
}
if self.buf.len() < BOM.len() && BOM.starts_with(&self.buf) {
return false; }
if self.buf.starts_with(BOM) {
self.buf.drain(..BOM.len());
}
self.bom_stripped = true;
true
}
}
fn find_frame_end(buf: &[u8]) -> Option<usize> {
(0..buf.len()).find_map(|i| {
if buf[i..].starts_with(b"\n\n") {
Some(i + 2)
} else if buf[i..].starts_with(b"\r\n\r\n") {
Some(i + 4)
} else {
None
}
})
}
fn parse_block(block: &[u8]) -> Option<Frame> {
let mut event: Option<String> = None;
let mut data: Vec<Vec<u8>> = Vec::new();
for raw_line in block.split(|&b| b == b'\n') {
let line = raw_line.strip_suffix(b"\r").unwrap_or(raw_line);
let Some(colon) = line.iter().position(|&b| b == b':') else {
continue;
};
let field = &line[..colon];
let value = &line[colon + 1..];
let value = value.strip_prefix(b" ").unwrap_or(value);
match field {
b"event" => event = Some(String::from_utf8_lossy(value).into_owned()),
b"data" => data.push(value.to_vec()),
_ => {} }
}
if event.is_none() && data.is_empty() {
None
} else {
Some(Frame {
event,
data: data.join(&b'\n'),
status: None,
})
}
}
#[derive(Default)]
pub struct NdjsonDecoder {
buf: Vec<u8>,
}
impl Decoder for NdjsonDecoder {
fn push(&mut self, chunk: Vec<u8>) -> Result<Vec<Frame>, CanonicalError> {
self.buf.extend_from_slice(&chunk);
let mut frames = Vec::new();
while let Some(nl) = self.buf.iter().position(|&b| b == b'\n') {
let mut line: Vec<u8> = self.buf.drain(..=nl).collect();
line.pop(); if let Some(frame) = line_frame(line) {
frames.push(frame);
}
}
Ok(frames) }
fn finish(&mut self) -> Result<Vec<Frame>, CanonicalError> {
let line = std::mem::take(&mut self.buf);
Ok(line_frame(line).into_iter().collect())
}
}
fn line_frame(mut line: Vec<u8>) -> Option<Frame> {
if line.last() == Some(&b'\r') {
line.pop();
}
if line.is_empty() {
None
} else {
Some(Frame {
event: None,
data: line,
status: None,
})
}
}
pub struct IdentityDecoder;
impl Decoder for IdentityDecoder {
fn push(&mut self, chunk: Vec<u8>) -> Result<Vec<Frame>, CanonicalError> {
Ok(vec![Frame {
event: None,
data: chunk,
status: None,
}])
}
fn finish(&mut self) -> Result<Vec<Frame>, CanonicalError> {
Ok(vec![]) }
}