use core::fmt;
const UTF8_BOM: [u8; 3] = [0xEF, 0xBB, 0xBF];
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SseLimits {
line_bytes: usize,
event_bytes: usize,
keepalive_lines: usize,
}
impl SseLimits {
#[must_use]
pub const fn new(
max_line_bytes: usize,
max_event_bytes: usize,
max_keepalive_lines: usize,
) -> Option<Self> {
if max_line_bytes == 0 || max_event_bytes == 0 || max_keepalive_lines == 0 {
return None;
}
Some(Self {
line_bytes: max_line_bytes,
event_bytes: max_event_bytes,
keepalive_lines: max_keepalive_lines,
})
}
#[must_use]
pub const fn max_line_bytes(&self) -> usize {
self.line_bytes
}
#[must_use]
pub const fn max_event_bytes(&self) -> usize {
self.event_bytes
}
#[must_use]
pub const fn max_keepalive_lines(&self) -> usize {
self.keepalive_lines
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SseEndOfStream {
pub discarded_pending_event: bool,
pub discarded_partial_line: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SseParseError {
LineTooLong {
limit_bytes: usize,
},
EventTooLarge {
limit_bytes: usize,
},
KeepaliveFlood {
limit_lines: usize,
},
Poisoned,
}
#[derive(Debug)]
pub(crate) enum SsePushError<E> {
Parse(SseParseError),
Consumer(E),
}
impl fmt::Display for SseParseError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::LineTooLong { limit_bytes } => {
write!(formatter, "SSE line exceeds {limit_bytes} bytes")
}
Self::EventTooLarge { limit_bytes } => {
write!(formatter, "SSE event exceeds {limit_bytes} bytes")
}
Self::KeepaliveFlood { limit_lines } => {
write!(
formatter,
"SSE stream exceeds {limit_lines} consecutive non-dispatching lines"
)
}
Self::Poisoned => formatter.write_str("SSE parser already refused earlier input"),
}
}
}
impl std::error::Error for SseParseError {}
#[derive(Debug)]
pub struct BoundedSseParser {
limits: SseLimits,
raw_line: Vec<u8>,
pending_cr: bool,
bom_window_open: bool,
data: String,
event_raw_bytes: usize,
keepalive_lines: usize,
poisoned: bool,
}
impl BoundedSseParser {
#[must_use]
pub fn new(limits: SseLimits) -> Self {
Self {
limits,
raw_line: Vec::new(),
pending_cr: false,
bom_window_open: true,
data: String::new(),
event_raw_bytes: 0,
keepalive_lines: 0,
poisoned: false,
}
}
#[must_use]
pub fn buffered_bytes(&self) -> usize {
self.raw_line.len() + self.data.len()
}
pub fn push(&mut self, chunk: &[u8]) -> Result<Vec<String>, SseParseError> {
let mut dispatched = Vec::new();
self.push_with(chunk, |payload| {
dispatched.push(payload);
Ok::<_, ()>(())
})
.map_err(|error| match error {
SsePushError::Parse(error) => error,
SsePushError::Consumer(()) => unreachable!("infallible SSE event collector refused"),
})?;
Ok(dispatched)
}
pub(crate) fn push_with<E>(
&mut self,
chunk: &[u8],
mut accept: impl FnMut(String) -> Result<(), E>,
) -> Result<(), SsePushError<E>> {
if self.poisoned {
return Err(SsePushError::Parse(SseParseError::Poisoned));
}
for &byte in chunk {
if self.pending_cr {
self.pending_cr = false;
if byte == b'\n' {
continue;
}
}
match byte {
b'\r' => {
self.complete_line_with(&mut accept)?;
self.pending_cr = true;
}
b'\n' => self.complete_line_with(&mut accept)?,
other => {
if self.raw_line.len() >= self.limits.line_bytes {
return Err(SsePushError::Parse(self.poison(
SseParseError::LineTooLong {
limit_bytes: self.limits.line_bytes,
},
)));
}
self.raw_line.push(other);
}
}
}
Ok(())
}
pub fn finish(self) -> Result<SseEndOfStream, SseParseError> {
if self.poisoned {
return Err(SseParseError::Poisoned);
}
Ok(SseEndOfStream {
discarded_pending_event: !self.data.is_empty(),
discarded_partial_line: !self.raw_line.is_empty(),
})
}
fn poison(&mut self, error: SseParseError) -> SseParseError {
self.raw_line = Vec::new();
self.data = String::new();
self.event_raw_bytes = 0;
self.keepalive_lines = 0;
self.poisoned = true;
error
}
fn complete_line_with<E>(
&mut self,
accept: &mut impl FnMut(String) -> Result<(), E>,
) -> Result<(), SsePushError<E>> {
let mut raw = core::mem::take(&mut self.raw_line);
if self.bom_window_open {
self.bom_window_open = false;
if raw.starts_with(&UTF8_BOM) {
raw.drain(..UTF8_BOM.len());
}
}
let raw_len = raw.len();
let line = String::from_utf8_lossy(&raw);
if line.len() > self.limits.line_bytes {
return Err(SsePushError::Parse(self.poison(
SseParseError::LineTooLong {
limit_bytes: self.limits.line_bytes,
},
)));
}
match self.process_line_with(&line, raw_len, accept) {
Ok(()) => Ok(()),
Err(SsePushError::Parse(error)) => Err(SsePushError::Parse(self.poison(error))),
Err(SsePushError::Consumer(error)) => Err(SsePushError::Consumer(error)),
}
}
fn process_line_with<E>(
&mut self,
line: &str,
raw_len: usize,
accept: &mut impl FnMut(String) -> Result<(), E>,
) -> Result<(), SsePushError<E>> {
if line.is_empty() {
if self.data.is_empty() {
return self
.count_non_dispatching_line()
.map_err(SsePushError::Parse);
}
let mut payload = core::mem::take(&mut self.data);
if payload.ends_with('\n') {
payload.pop();
}
self.event_raw_bytes = 0;
self.keepalive_lines = 0;
return accept(payload).map_err(|error| {
self.poison(SseParseError::Poisoned);
SsePushError::Consumer(error)
});
}
if line.starts_with(':') {
return self
.count_non_dispatching_line()
.map_err(SsePushError::Parse);
}
let (field, value) = match line.split_once(':') {
Some((field, value)) => (field, value.strip_prefix(' ').unwrap_or(value)),
None => (line, ""),
};
if field == "data" {
let decoded_after = self
.data
.len()
.saturating_add(value.len())
.saturating_add(1);
let raw_after = self.event_raw_bytes.saturating_add(raw_len);
if decoded_after > self.limits.event_bytes || raw_after > self.limits.event_bytes {
return Err(SsePushError::Parse(self.poison(
SseParseError::EventTooLarge {
limit_bytes: self.limits.event_bytes,
},
)));
}
self.data.push_str(value);
self.data.push('\n');
self.event_raw_bytes = raw_after;
self.keepalive_lines = 0;
return Ok(());
}
self.count_non_dispatching_line()
.map_err(SsePushError::Parse)
}
fn count_non_dispatching_line(&mut self) -> Result<(), SseParseError> {
self.keepalive_lines = self.keepalive_lines.saturating_add(1);
if self.keepalive_lines > self.limits.keepalive_lines {
return Err(SseParseError::KeepaliveFlood {
limit_lines: self.limits.keepalive_lines,
});
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::{BoundedSseParser, SseEndOfStream, SseLimits, SseParseError, SsePushError};
fn generous_limits() -> SseLimits {
SseLimits::new(4_096, 65_536, 64).expect("nonzero test limits")
}
fn parse_all(
input: &[u8],
limits: SseLimits,
) -> Result<(Vec<String>, SseEndOfStream), SseParseError> {
let mut parser = BoundedSseParser::new(limits);
let events = parser.push(input)?;
let end = parser.finish()?;
Ok((events, end))
}
fn assert_chunk_invariant(input: &[u8], limits: SseLimits) -> (Vec<String>, SseEndOfStream) {
let reference = parse_all(input, limits).expect("reference parse succeeds");
let mut byte_parser = BoundedSseParser::new(limits);
let mut byte_events = Vec::new();
for byte in input {
byte_events.extend(
byte_parser
.push(core::slice::from_ref(byte))
.expect("byte-by-byte parse succeeds"),
);
}
let byte_end = byte_parser.finish().expect("byte-by-byte finish succeeds");
assert_eq!(
(&byte_events, byte_end),
(&reference.0, reference.1),
"byte-by-byte feed must match whole-buffer feed"
);
for split in 0..=input.len() {
let mut parser = BoundedSseParser::new(limits);
let mut events = parser
.push(&input[..split])
.expect("first split half parses");
events.extend(
parser
.push(&input[split..])
.expect("second split half parses"),
);
let end = parser.finish().expect("split finish succeeds");
assert_eq!(
(&events, end),
(&reference.0, reference.1),
"split at byte {split} must match whole-buffer feed"
);
}
reference
}
#[test]
fn dispatches_single_data_event() {
let (events, end) = assert_chunk_invariant(b"data: hello\n\n", generous_limits());
assert_eq!(events, ["hello"]);
assert_eq!(
end,
SseEndOfStream {
discarded_pending_event: false,
discarded_partial_line: false
}
);
}
#[test]
fn data_field_name_variants_match_the_standard() {
let (events, _) = assert_chunk_invariant(b"data\n\n", generous_limits());
assert_eq!(events, [""]);
let (events, _) = assert_chunk_invariant(b"data:\n\n", generous_limits());
assert_eq!(events, [""]);
let (events, _) = assert_chunk_invariant(b"data: \n\n", generous_limits());
assert_eq!(events, [""]);
let (events, _) = assert_chunk_invariant(b"data: two spaces\n\n", generous_limits());
assert_eq!(events, [" two spaces"]);
}
#[test]
fn multiple_data_lines_join_with_inserted_newlines() {
let (events, _) =
assert_chunk_invariant(b"data: a\ndata: b\ndata: c\n\n", generous_limits());
assert_eq!(events, ["a\nb\nc"]);
}
#[test]
fn exactly_one_trailing_inserted_newline_is_removed() {
let (events, _) = assert_chunk_invariant(b"data: a\ndata:\n\n", generous_limits());
assert_eq!(events, ["a\n"]);
}
#[test]
fn blank_line_without_data_produces_no_event() {
let (events, _) = assert_chunk_invariant(b"\n\n\n", generous_limits());
assert_eq!(events, Vec::<String>::new());
}
#[test]
fn empty_data_event_is_distinct_from_no_data_event() {
let (events, _) = assert_chunk_invariant(b"\ndata:\n\n\n", generous_limits());
assert_eq!(events, [""]);
}
#[test]
fn comment_lines_produce_no_event() {
let (events, _) =
assert_chunk_invariant(b": keepalive\n: another\ndata: x\n\n", generous_limits());
assert_eq!(events, ["x"]);
}
#[test]
fn inert_fields_are_parsed_but_carry_no_state() {
let input: &[u8] = b"event: custom\nid: 7\nretry: 50\nunknown: y\ndata: x\n\nid: 8\n\n";
let (events, _) = assert_chunk_invariant(input, generous_limits());
assert_eq!(events, ["x"]);
}
#[test]
fn field_without_colon_is_a_name_with_empty_value() {
let (events, _) = assert_chunk_invariant(b"event\ndata: x\n\n", generous_limits());
assert_eq!(events, ["x"]);
}
#[test]
fn crlf_bare_lf_and_bare_cr_terminate_lines_identically() {
for input in [
b"data: a\r\n\r\n".as_slice(),
b"data: a\n\n".as_slice(),
b"data: a\r\r".as_slice(),
b"data: a\r\n\n".as_slice(),
b"data: a\n\r\n".as_slice(),
] {
let (events, _) = assert_chunk_invariant(input, generous_limits());
assert_eq!(events, ["a"], "input {input:?}");
}
}
#[test]
fn leading_bom_is_stripped_exactly_once() {
let mut input = vec![0xEF, 0xBB, 0xBF];
input.extend_from_slice(b"data: x\n\n");
let (events, _) = assert_chunk_invariant(&input, generous_limits());
assert_eq!(events, ["x"]);
}
#[test]
fn midstream_bom_is_ordinary_content() {
let mut input = b"data: ".to_vec();
input.extend_from_slice(&[0xEF, 0xBB, 0xBF]);
input.extend_from_slice(b"x\n\n");
let (events, _) = assert_chunk_invariant(&input, generous_limits());
assert_eq!(events, ["\u{FEFF}x"]);
let mut input = b"data: a\n".to_vec();
input.extend_from_slice(&[0xEF, 0xBB, 0xBF]);
input.extend_from_slice(b"data: b\n\n");
let (events, _) = assert_chunk_invariant(&input, generous_limits());
assert_eq!(events, ["a"]);
}
#[test]
fn malformed_utf8_is_replaced_never_fatal() {
let (events, _) = assert_chunk_invariant(b"data: \xFF\n\n", generous_limits());
assert_eq!(events, ["\u{FFFD}"]);
let (events, _) = assert_chunk_invariant(b"data: \xE2\x82ok\n\n", generous_limits());
assert_eq!(events, ["\u{FFFD}ok"]);
}
#[test]
fn multibyte_sequences_split_across_chunks_survive() {
let (events, _) = assert_chunk_invariant("data: €\n\n".as_bytes(), generous_limits());
assert_eq!(events, ["€"]);
}
#[test]
fn sequence_interrupted_by_terminator_is_replaced_within_its_line() {
let (events, _) = assert_chunk_invariant(b"data: \xE2\x82\ndata: x\n\n", generous_limits());
assert_eq!(events, ["\u{FFFD}\nx"]);
}
#[test]
fn nul_bytes_pass_through_data_values() {
let (events, _) = assert_chunk_invariant(b"data: a\x00b\n\n", generous_limits());
assert_eq!(events, ["a\u{0000}b"]);
}
#[test]
fn eof_discards_unterminated_pending_event() {
let limits = generous_limits();
let mut parser = BoundedSseParser::new(limits);
let events = parser.push(b"data: never dispatched\n").expect("parses");
assert_eq!(events, Vec::<String>::new());
let end = parser.finish().expect("finish succeeds");
assert!(end.discarded_pending_event);
assert!(!end.discarded_partial_line);
}
#[test]
fn eof_discards_unterminated_partial_line() {
let limits = generous_limits();
let mut parser = BoundedSseParser::new(limits);
parser.push(b"data: x\n\ndata: partial").expect("parses");
let end = parser.finish().expect("finish succeeds");
assert!(!end.discarded_pending_event);
assert!(end.discarded_partial_line);
}
#[test]
fn raw_line_bound_is_exact() {
let limits = SseLimits::new(8, 65_536, 64).expect("limits");
let (events, _) = assert_chunk_invariant(b"data: xx\n\n", limits);
assert_eq!(events, ["xx"]);
let mut parser = BoundedSseParser::new(limits);
let error = parser.push(b"data: xxx\n\n").expect_err("line too long");
assert_eq!(error, SseParseError::LineTooLong { limit_bytes: 8 });
assert_eq!(parser.buffered_bytes(), 0, "refusal releases buffers");
assert_eq!(
parser.push(b"data: x\n\n"),
Err(SseParseError::Poisoned),
"a refused stream accepts nothing further"
);
}
#[test]
fn replacement_expansion_counts_against_the_decoded_line_bound() {
let limits = SseLimits::new(12, 65_536, 64).expect("limits");
let mut parser = BoundedSseParser::new(limits);
let error = parser
.push(b"data: \xFF\xFF\xFF\n\n")
.expect_err("decoded expansion exceeds the line ceiling");
assert_eq!(error, SseParseError::LineTooLong { limit_bytes: 12 });
}
#[test]
fn event_bound_is_exact_across_data_lines() {
let limits = SseLimits::new(4_096, 30, 64).expect("limits");
let (events, _) = assert_chunk_invariant(b"data: abcd\ndata: abcd\ndata: abcd\n\n", limits);
assert_eq!(events, ["abcd\nabcd\nabcd"]);
let mut parser = BoundedSseParser::new(limits);
let error = parser
.push(b"data: abcd\ndata: abcd\ndata: abcd\ndata: abcd\n\n")
.expect_err("event too large");
assert_eq!(error, SseParseError::EventTooLarge { limit_bytes: 30 });
assert_eq!(parser.buffered_bytes(), 0, "refusal releases buffers");
}
#[test]
fn keepalive_flood_is_bounded_and_reset_by_data() {
let limits = SseLimits::new(4_096, 65_536, 3).expect("limits");
let (events, _) = assert_chunk_invariant(
b": a\n: b\n: c\ndata: x\n\n: a\n: b\n: c\ndata: y\n\n",
limits,
);
assert_eq!(events, ["x", "y"]);
let mut parser = BoundedSseParser::new(limits);
let error = parser
.push(b": a\n: b\n: c\n: d\n")
.expect_err("comment flood");
assert_eq!(error, SseParseError::KeepaliveFlood { limit_lines: 3 });
}
#[test]
fn dispatch_resets_buffered_bytes() {
let mut parser = BoundedSseParser::new(generous_limits());
parser.push(b"data: hello").expect("parses");
assert!(parser.buffered_bytes() > 0);
let events = parser.push(b"\n\n").expect("dispatches");
assert_eq!(events, ["hello"]);
assert_eq!(parser.buffered_bytes(), 0);
}
#[test]
fn chunk_packed_consumer_overflow_stops_before_materializing_the_tail() {
let accepted = 3_usize;
let chunk = b"data: x\n\n".repeat(accepted + 2);
let mut parser = BoundedSseParser::new(generous_limits());
let mut observed = 0_usize;
let error = parser
.push_with(&chunk, |_| {
observed += 1;
if observed > accepted { Err(()) } else { Ok(()) }
})
.expect_err("the first event beyond the one-variable count limit refuses");
assert!(matches!(error, SsePushError::Consumer(())));
assert_eq!(observed, accepted + 1);
assert_eq!(
parser.buffered_bytes(),
0,
"refusal releases parser buffers"
);
assert_eq!(
parser.push(b"data: later\n\n"),
Err(SseParseError::Poisoned),
"a consumer-refused chunk cannot materialize a later tail"
);
}
#[test]
fn zero_limits_are_rejected_at_construction() {
assert!(SseLimits::new(0, 1, 1).is_none());
assert!(SseLimits::new(1, 0, 1).is_none());
assert!(SseLimits::new(1, 1, 0).is_none());
}
#[test]
fn poisoned_finish_reports_poisoned_not_a_discard_summary() {
let limits = SseLimits::new(8, 65_536, 64).expect("limits");
let mut parser = BoundedSseParser::new(limits);
parser
.push(b"data: way too long for the line bound\n\n")
.expect_err("line too long");
assert_eq!(parser.finish(), Err(SseParseError::Poisoned));
}
#[test]
fn interleaved_events_dispatch_in_stream_order() {
let input: &[u8] =
b"data: first\n\n: keepalive\nevent: progress\ndata: second\n\ndata: third\n\n";
let (events, end) = assert_chunk_invariant(input, generous_limits());
assert_eq!(events, ["first", "second", "third"]);
assert!(!end.discarded_pending_event);
}
}