#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Framing {
Sse,
Ndjson,
Whole,
}
impl Framing {
pub fn split(&self, bytes: &[u8]) -> Vec<WireFrame> {
match self {
Self::Sse => SseFramer::new()
.push(bytes)
.filter(|event| !event.data.trim().is_empty())
.map(|event| WireFrame::Text(event.data))
.collect(),
Self::Ndjson => {
let mut framer = NdjsonFramer::new();
let mut lines: Vec<Vec<u8>> = framer.push(bytes).collect();
lines.extend(framer.finish());
lines.into_iter().map(WireFrame::from_bytes).collect()
}
Self::Whole if bytes.is_empty() => Vec::new(),
Self::Whole => vec![WireFrame::from_bytes(bytes.to_vec())],
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum WireFrame {
Text(String),
Bytes(Vec<u8>),
}
impl WireFrame {
pub fn as_str(&self) -> std::borrow::Cow<'_, str> {
match self {
Self::Text(text) => std::borrow::Cow::Borrowed(text),
Self::Bytes(bytes) => String::from_utf8_lossy(bytes),
}
}
pub fn from_bytes(bytes: Vec<u8>) -> Self {
match String::from_utf8(bytes) {
Ok(text) => Self::Text(text),
Err(error) => Self::Bytes(error.into_bytes()),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SseEvent {
pub event: String,
pub data: String,
pub id: String,
pub retry: Option<u64>,
}
const BOM: [u8; 3] = [0xef, 0xbb, 0xbf];
#[derive(Debug, Default)]
pub struct SseFramer {
buffer: Vec<u8>,
ready: Vec<SseEvent>,
event_type: String,
data: String,
last_event_id: String,
retry: Option<u64>,
bom_prefix: usize,
bom_done: bool,
since_dispatch: usize,
}
impl SseFramer {
pub fn new() -> Self {
Self::default()
}
pub fn push(&mut self, chunk: &[u8]) -> std::vec::Drain<'_, SseEvent> {
let chunk = self.strip_bom(chunk);
self.buffer.extend_from_slice(chunk);
while let Some((line_len, consumed)) = terminated_line(&self.buffer) {
let line = String::from_utf8_lossy(self.buffer.get(..line_len).unwrap_or_default())
.into_owned();
self.buffer.drain(..consumed);
self.since_dispatch = self.since_dispatch.saturating_add(consumed);
self.line(&line);
}
self.ready.drain(..)
}
pub fn pending(&self) -> usize {
self.since_dispatch.saturating_add(self.buffer.len())
}
fn strip_bom<'a>(&mut self, chunk: &'a [u8]) -> &'a [u8] {
if self.bom_done {
return chunk;
}
let mut rest = chunk;
while let Some(&byte) = rest.first() {
if BOM.get(self.bom_prefix) != Some(&byte) {
let matched = std::mem::take(&mut self.bom_prefix);
self.bom_done = true;
self.buffer
.extend_from_slice(BOM.get(..matched).unwrap_or_default());
return rest;
}
self.bom_prefix += 1;
rest = rest.get(1..).unwrap_or_default();
if self.bom_prefix == BOM.len() {
self.bom_prefix = 0;
self.bom_done = true;
return rest;
}
}
rest
}
fn line(&mut self, line: &str) {
if line.is_empty() {
self.dispatch();
return;
}
if let Some(rest) = line.strip_prefix(':') {
let _ = rest;
return;
}
let (field, value) = match line.split_once(':') {
Some((field, value)) => (field, value.strip_prefix(' ').unwrap_or(value)),
None => (line, ""),
};
match field {
"event" => {
self.event_type.clear();
self.event_type.push_str(value);
}
"data" => {
self.data.push_str(value);
self.data.push('\n');
}
"id" if !value.contains('\0') => {
self.last_event_id.clear();
self.last_event_id.push_str(value);
}
"retry" if !value.is_empty() && value.bytes().all(|b| b.is_ascii_digit()) => {
self.retry = value.parse().ok();
}
_ => {}
}
}
fn dispatch(&mut self) {
self.since_dispatch = 0;
if self.data.is_empty() {
self.event_type.clear();
return;
}
self.data.pop();
let event = if self.event_type.is_empty() {
"message".to_owned()
} else {
std::mem::take(&mut self.event_type)
};
self.ready.push(SseEvent {
event,
data: std::mem::take(&mut self.data),
id: self.last_event_id.clone(),
retry: self.retry.take(),
});
self.event_type.clear();
}
}
#[derive(Debug, Default)]
pub struct NdjsonFramer {
buffer: Vec<u8>,
ready: Vec<Vec<u8>>,
}
impl NdjsonFramer {
pub fn new() -> Self {
Self::default()
}
pub fn push(&mut self, chunk: &[u8]) -> std::vec::Drain<'_, Vec<u8>> {
self.buffer.extend_from_slice(chunk);
while let Some(pos) = self.buffer.iter().position(|byte| *byte == b'\n') {
let mut line: Vec<u8> = self.buffer.drain(..=pos).collect();
line.pop();
if line.last() == Some(&b'\r') {
line.pop();
}
if !line.is_empty() {
self.ready.push(line);
}
}
self.ready.drain(..)
}
pub fn finish(&mut self) -> Option<Vec<u8>> {
let line = std::mem::take(&mut self.buffer);
(!line.is_empty()).then_some(line)
}
pub fn pending(&self) -> usize {
self.buffer.len()
}
}
fn terminated_line(buffer: &[u8]) -> Option<(usize, usize)> {
let pos = buffer
.iter()
.position(|byte| *byte == b'\n' || *byte == b'\r')?;
match buffer.get(pos) {
Some(b'\n') => Some((pos, pos + 1)),
Some(b'\r') => match buffer.get(pos + 1) {
Some(b'\n') => Some((pos, pos + 2)),
Some(_) => Some((pos, pos + 1)),
None => None,
},
_ => None,
}
}
#[cfg(test)]
mod tests;