const DEFAULT_MAX_EVENT_BYTES: usize = 1 << 20;
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
#[error("an SSE line or event exceeded the {limit}-byte limit")]
pub struct DecodeError {
limit: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Event {
pub name: String,
pub data: String,
}
#[derive(Debug)]
pub struct Decoder {
partial: String,
name: String,
data: String,
started: bool,
max_event_bytes: usize,
}
impl Default for Decoder {
fn default() -> Self {
Self {
partial: String::new(),
name: String::new(),
data: String::new(),
started: false,
max_event_bytes: DEFAULT_MAX_EVENT_BYTES,
}
}
}
impl Decoder {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[cfg(test)]
fn with_max_event_bytes(max_event_bytes: usize) -> Self {
Self {
max_event_bytes,
..Self::default()
}
}
pub fn push(&mut self, chunk: &[u8]) -> Result<Vec<Event>, DecodeError> {
self.partial.push_str(&String::from_utf8_lossy(chunk));
let mut out = Vec::new();
while let Some((line, rest)) = split_line(&self.partial) {
if line.len() > self.max_event_bytes {
return Err(self.too_large());
}
if rest.len() >= self.partial.len() {
break;
}
let line = line.to_owned();
self.partial = rest.to_owned();
if let Some(event) = self.line(&line) {
out.push(event);
}
if self.name.len().saturating_add(self.data.len()) > self.max_event_bytes {
return Err(self.too_large());
}
}
if self.partial.len() > self.max_event_bytes {
return Err(self.too_large());
}
Ok(out)
}
fn too_large(&self) -> DecodeError {
DecodeError {
limit: self.max_event_bytes,
}
}
fn line(&mut self, line: &str) -> Option<Event> {
if line.is_empty() {
if !self.started {
return None;
}
let name = std::mem::take(&mut self.name);
let data = std::mem::take(&mut self.data);
self.started = false;
return Some(Event { name, data });
}
if line.starts_with(':') {
return None;
}
let (field, value) = match line.split_once(':') {
Some((f, v)) => (f, v.strip_prefix(' ').unwrap_or(v)),
None => (line, ""),
};
match field {
"event" => {
value.clone_into(&mut self.name);
self.started = true;
}
"data" => {
if !self.data.is_empty() {
self.data.push('\n');
}
self.data.push_str(value);
self.started = true;
}
_ => {}
}
None
}
}
fn split_line(buf: &str) -> Option<(&str, &str)> {
let bytes = buf.as_bytes();
let idx = bytes.iter().position(|&b| b == b'\n' || b == b'\r')?;
match bytes[idx] {
b'\n' => Some((&buf[..idx], &buf[idx + 1..])),
_ if idx + 1 == bytes.len() => None,
_ if bytes[idx + 1] == b'\n' => Some((&buf[..idx], &buf[idx + 2..])),
_ => Some((&buf[..idx], &buf[idx + 1..])),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn all(chunks: &[&str]) -> Vec<Event> {
let mut d = Decoder::new();
let mut out = Vec::new();
for c in chunks {
out.extend(d.push(c.as_bytes()).expect("valid SSE"));
}
out
}
#[test]
fn a_named_event_with_data_dispatches_on_the_blank_line() {
assert_eq!(
all(&["event: message_start\ndata: {\"a\":1}\n\n"]),
vec![Event {
name: "message_start".to_owned(),
data: "{\"a\":1}".to_owned(),
}]
);
}
#[test]
fn an_event_split_across_chunks_is_reassembled() {
assert_eq!(
all(&["event: mess", "age_start\nda", "ta: {\"a\":", "1}\n", "\n"]),
vec![Event {
name: "message_start".to_owned(),
data: "{\"a\":1}".to_owned(),
}]
);
}
#[test]
fn a_boundary_inside_the_terminator_still_dispatches_once() {
let events = all(&["data: x\n", "\ndata: y\n\n"]);
assert_eq!(events.len(), 2);
assert_eq!(events[0].data, "x");
assert_eq!(events[1].data, "y");
}
#[test]
fn multi_line_data_is_joined_with_newlines() {
assert_eq!(all(&["data: one\ndata: two\n\n"])[0].data, "one\ntwo");
}
#[test]
fn comments_are_ignored() {
let events = all(&[": keep-alive\ndata: x\n\n"]);
assert_eq!(events.len(), 1);
assert_eq!(events[0].data, "x");
}
#[test]
fn a_comment_alone_dispatches_nothing() {
assert!(all(&[": ping\n\n"]).is_empty());
}
#[test]
fn crlf_and_bare_cr_terminate_lines() {
assert_eq!(all(&["data: a\r\n\r\n"])[0].data, "a");
assert_eq!(all(&["data: b\r\rx"])[0].data, "b");
}
#[test]
fn crlf_is_one_terminator_not_two() {
let events = all(&["data: a\r\ndata: b\r\n\r\n"]);
assert_eq!(
events.len(),
1,
"the CRLF was read as two terminators, so the blank line it \
manufactured dispatched an event early: {events:?}"
);
assert_eq!(events[0].data, "a\nb");
}
#[test]
fn a_crlf_split_across_chunks_is_still_one_terminator() {
let events = all(&["data: a\r", "\ndata: b\r\n\r\n"]);
assert_eq!(events.len(), 1, "{events:?}");
assert_eq!(events[0].data, "a\nb");
}
#[test]
fn a_stream_ending_on_a_bare_cr_holds_it_back() {
assert!(all(&["data: b\r"]).is_empty());
}
#[test]
fn a_split_crlf_is_not_read_as_two_terminators() {
let events = all(&["data: a\r", "\n\r\n"]);
assert_eq!(events.len(), 1, "{events:?}");
assert_eq!(events[0].data, "a");
}
#[test]
fn only_one_leading_space_is_stripped() {
assert_eq!(all(&["data: x\n\n"])[0].data, " x");
assert_eq!(all(&["data:x\n\n"])[0].data, "x");
}
#[test]
fn a_field_with_no_colon_has_an_empty_value() {
let events = all(&["data\n\n"]);
assert_eq!(events.len(), 1);
assert_eq!(events[0].data, "");
}
#[test]
fn reconnection_fields_are_ignored_but_do_not_break_the_event() {
let events = all(&["id: 7\nretry: 3000\nevent: e\ndata: d\n\n"]);
assert_eq!(events.len(), 1);
assert_eq!(events[0].name, "e");
assert_eq!(events[0].data, "d");
}
#[test]
fn an_incomplete_trailing_line_is_not_dispatched() {
assert!(all(&["event: message_start\ndata: {\"a\""]).is_empty());
}
#[test]
fn an_unnamed_event_still_carries_its_data() {
let events = all(&["data: x\n\n"]);
assert_eq!(events[0].name, "");
assert_eq!(events[0].data, "x");
}
#[test]
fn invalid_utf8_degrades_rather_than_aborting() {
let mut d = Decoder::new();
let events = d.push(b"data: \xff\xfe\n\n").expect("valid SSE");
assert_eq!(events.len(), 1, "a bad byte must not lose the event");
}
#[test]
fn an_unterminated_line_cannot_grow_without_bound() {
let mut d = Decoder::with_max_event_bytes(8);
let err = d
.push(b"123456789")
.expect_err("an endless line must be bounded");
assert_eq!(err.limit, 8);
}
#[test]
fn multi_line_event_data_cannot_grow_without_bound() {
let mut d = Decoder::with_max_event_bytes(8);
d.push(b"data: 1\n").expect("first line fits");
d.push(b"data: 2\n").expect("second line fits");
d.push(b"data: 3\n").expect("third line fits");
d.push(b"data: 4\n").expect("fourth line fits");
d.push(b"data: 5\n")
.expect_err("the assembled event exceeds the limit");
}
}