Skip to main content

eventsource_client/
event_parser.rs

1use std::{collections::VecDeque, convert::TryFrom, str::from_utf8};
2
3use bytes::Bytes;
4use log::{debug, log_enabled, trace};
5use pin_project::pin_project;
6
7use crate::response::Response;
8
9use super::error::{Error, Result};
10
11#[derive(Default, PartialEq)]
12struct EventData {
13    pub event_type: String,
14    pub data: String,
15    pub id: Option<String>,
16    pub retry: Option<u64>,
17}
18
19impl EventData {
20    fn new() -> Self {
21        Self::default()
22    }
23
24    pub fn append_data(&mut self, value: &str) {
25        self.data.push_str(value);
26        self.data.push('\n');
27    }
28
29    pub fn with_id(mut self, value: Option<String>) -> Self {
30        self.id = value;
31        self
32    }
33}
34
35#[derive(Debug, Eq, PartialEq)]
36pub enum SSE {
37    Connected(ConnectionDetails),
38    Event(Event),
39    Comment(String),
40}
41
42impl TryFrom<EventData> for Option<SSE> {
43    type Error = Error;
44
45    fn try_from(event_data: EventData) -> std::result::Result<Self, Self::Error> {
46        if event_data == EventData::default() {
47            return Err(Error::InvalidEvent);
48        }
49
50        if event_data.data.is_empty() {
51            return Ok(None);
52        }
53
54        let event_type = if event_data.event_type.is_empty() {
55            String::from("message")
56        } else {
57            event_data.event_type
58        };
59
60        let mut data = event_data.data.clone();
61        data.truncate(data.len() - 1);
62
63        let id = event_data.id.clone();
64
65        let retry = event_data.retry;
66
67        Ok(Some(SSE::Event(Event {
68            event_type,
69            data,
70            id,
71            retry,
72        })))
73    }
74}
75
76#[derive(Clone, Debug, Eq, PartialEq)]
77pub struct ConnectionDetails {
78    response: Response,
79}
80
81impl ConnectionDetails {
82    pub(crate) fn new(response: Response) -> Self {
83        Self { response }
84    }
85
86    /// Returns information describing the response at the time of connection.
87    pub fn response(&self) -> &Response {
88        &self.response
89    }
90}
91
92#[derive(Clone, Debug, Eq, PartialEq)]
93pub struct Event {
94    pub event_type: String,
95    pub data: String,
96    pub id: Option<String>,
97    pub retry: Option<u64>,
98}
99
100const LOGIFY_MAX_CHARS: usize = 100;
101fn logify(bytes: &[u8]) -> String {
102    let stringified = from_utf8(bytes).unwrap_or("<bad utf8>");
103    stringified.chars().take(LOGIFY_MAX_CHARS).collect()
104}
105
106fn parse_field(line: &[u8]) -> Result<Option<(&str, &str)>> {
107    if line.is_empty() {
108        return Err(Error::InvalidLine(
109            "should never try to parse an empty line (probably a bug)".into(),
110        ));
111    }
112
113    match line.iter().position(|&b| b':' == b) {
114        Some(0) => {
115            let value = &line[1..];
116            debug!("comment: {}", logify(value));
117            Ok(Some(("comment", parse_value(value)?)))
118        }
119        Some(colon_pos) => {
120            let key = &line[0..colon_pos];
121            let key = parse_key(key)?;
122
123            let mut value = &line[colon_pos + 1..];
124            // remove the first initial space character if any (but remove no other whitespace)
125            if value.starts_with(b" ") {
126                value = &value[1..];
127            }
128
129            debug!("key: {}, value: {}", key, logify(value));
130
131            Ok(Some((key, parse_value(value)?)))
132        }
133        None => Ok(Some((parse_key(line)?, ""))),
134    }
135}
136
137fn parse_key(key: &[u8]) -> Result<&str> {
138    from_utf8(key).map_err(|e| Error::InvalidLine(format!("malformed key: {e:?}")))
139}
140
141fn parse_value(value: &[u8]) -> Result<&str> {
142    from_utf8(value).map_err(|e| Error::InvalidLine(format!("malformed value: {e:?}")))
143}
144
145// A state machine for handling the BOM header.
146#[derive(Debug)]
147enum BomHeaderState {
148    Parsing(Vec<u8>),
149    Consumed,
150}
151
152const BOM_HEADER: &[u8] = b"\xEF\xBB\xBF";
153
154// Try to consume the BOM header from the given bytes.
155// If the BOM header is found, return the remaining bytes, otherwise return the origin buffer.
156// Return `None` if we cannot determine whether the BOM header is present.
157fn try_consume_bom_header(buf: &[u8]) -> Option<&[u8]> {
158    if buf.len() < BOM_HEADER.len() {
159        if BOM_HEADER.starts_with(buf) {
160            None
161        } else {
162            Some(buf)
163        }
164    } else if buf.starts_with(BOM_HEADER) {
165        Some(&buf[BOM_HEADER.len()..])
166    } else {
167        Some(buf)
168    }
169}
170
171#[pin_project]
172#[must_use = "streams do nothing unless polled"]
173pub struct EventParser {
174    /// buffer for lines we know are complete (terminated) but not yet parsed into event fields, in
175    /// the order received
176    complete_lines: VecDeque<Vec<u8>>,
177    /// buffer for the most-recently received line, pending completion (by a newline terminator) or
178    /// extension (by more non-newline bytes)
179    incomplete_line: Option<Vec<u8>>,
180    /// flagged if the last character processed as a carriage return; used to help process CRLF
181    /// pairs
182    last_char_was_cr: bool,
183    /// the event data currently being decoded
184    event_data: Option<EventData>,
185    /// the last-seen event ID; events without an ID will take on this value until it is updated.
186    last_event_id: Option<String>,
187    sse: VecDeque<SSE>,
188    /// state machine for handling the BOM header
189    bom_header_state: BomHeaderState,
190}
191
192impl EventParser {
193    pub fn new() -> Self {
194        Self {
195            complete_lines: VecDeque::with_capacity(10),
196            incomplete_line: None,
197            last_char_was_cr: false,
198            event_data: None,
199            last_event_id: None,
200            sse: VecDeque::with_capacity(3),
201            bom_header_state: BomHeaderState::Parsing(Vec::new()),
202        }
203    }
204
205    pub fn was_processing(&self) -> bool {
206        if self.incomplete_line.is_some() || !self.complete_lines.is_empty() {
207            true
208        } else {
209            !self.sse.is_empty()
210        }
211    }
212
213    pub fn get_event(&mut self) -> Option<SSE> {
214        self.sse.pop_front()
215    }
216
217    pub fn process_bytes(&mut self, bytes: Bytes) -> Result<()> {
218        trace!("Parsing bytes {bytes:?}");
219
220        // According to the SSE spec, a BOM header may be present at the beginning of the stream,
221        // which must be stripped before the message processing.
222        let bytes_to_process =
223            if let BomHeaderState::Parsing(header_buf) = &mut self.bom_header_state {
224                header_buf.extend_from_slice(&bytes);
225                if let Some(rest) = try_consume_bom_header(header_buf) {
226                    let owned_rest = rest.to_vec();
227                    self.bom_header_state = BomHeaderState::Consumed;
228                    // Once the BOM header is consumed, we can process the rest of the bytes.
229                    Bytes::from_owner(owned_rest)
230                } else {
231                    return Ok(());
232                }
233            } else {
234                bytes
235            };
236
237        // We get bytes from the underlying stream in chunks.  Decoding a chunk has two phases:
238        // decode the chunk into lines, and decode the lines into events.
239        //
240        // We counterintuitively do these two phases in reverse order. Because both lines and
241        // events may be split across chunks, we need to ensure we have a complete
242        // (newline-terminated) line before parsing it, and a complete event
243        // (empty-line-terminated) before returning it. So we buffer lines between poll()
244        // invocations, and begin by processing any incomplete events from previous invocations,
245        // before requesting new input from the underlying stream and processing that.
246        self.decode_and_buffer_lines(bytes_to_process);
247        self.parse_complete_lines_into_event()?;
248
249        Ok(())
250    }
251
252    // Populate the event fields from the complete lines already seen, until we either encounter an
253    // empty line - indicating we've decoded a complete event - or we run out of complete lines to
254    // process.
255    //
256    // Returns the event for dispatch if it is complete.
257    fn parse_complete_lines_into_event(&mut self) -> Result<()> {
258        loop {
259            let mut seen_empty_line = false;
260
261            while let Some(line) = self.complete_lines.pop_front() {
262                if line.is_empty() && self.event_data.is_some() {
263                    seen_empty_line = true;
264                    break;
265                } else if line.is_empty() {
266                    continue;
267                }
268
269                if let Some((key, value)) = parse_field(&line)? {
270                    if key == "comment" {
271                        self.sse.push_back(SSE::Comment(value.to_string()));
272                        continue;
273                    }
274
275                    if !matches!(key, "event" | "data" | "id" | "retry") {
276                        continue;
277                    }
278
279                    let id = &self.last_event_id;
280                    let event_data = self
281                        .event_data
282                        .get_or_insert_with(|| EventData::new().with_id(id.clone()));
283
284                    if key == "event" {
285                        event_data.event_type = value.to_string()
286                    } else if key == "data" {
287                        event_data.append_data(value);
288                    } else if key == "id" {
289                        // If id contains a null byte, it is a non-fatal error and the rest of
290                        // the event should be parsed if possible.
291                        if value.chars().any(|c| c == '\0') {
292                            debug!("Ignoring event ID containing null byte");
293                            continue;
294                        }
295
296                        if value.is_empty() {
297                            self.last_event_id = Some("".to_string());
298                        } else {
299                            self.last_event_id = Some(value.to_string());
300                        }
301
302                        event_data.id.clone_from(&self.last_event_id)
303                    } else if key == "retry" {
304                        match value.parse::<u64>() {
305                            Ok(retry) => {
306                                event_data.retry = Some(retry);
307                            }
308                            _ => debug!("Failed to parse {value:?} into retry value"),
309                        };
310                    }
311                }
312            }
313
314            if seen_empty_line {
315                let event_data = self.event_data.take();
316
317                trace!(
318                    "seen empty line, event_data is {:?})",
319                    event_data.as_ref().map(|event_data| &event_data.event_type)
320                );
321
322                if let Some(event_data) = event_data {
323                    match Option::<SSE>::try_from(event_data) {
324                        Err(e) => return Err(e),
325                        Ok(None) => (),
326                        Ok(Some(event)) => self.sse.push_back(event),
327                    };
328                }
329
330                continue;
331            } else {
332                trace!("processed all complete lines but event_data not yet complete");
333            }
334
335            break;
336        }
337
338        Ok(())
339    }
340
341    // Decode a chunk into lines and buffer them for subsequent parsing, taking account of
342    // incomplete lines from previous chunks.
343    fn decode_and_buffer_lines(&mut self, chunk: Bytes) {
344        let mut lines = chunk.split_inclusive(|&b| b == b'\n' || b == b'\r');
345        // The first and last elements in this split are special. The spec requires lines to be
346        // terminated. But lines may span chunks, so:
347        //  * the last line, if non-empty (i.e. if chunk didn't end with a line terminator),
348        //    should be buffered as an incomplete line
349        //  * the first line should be appended to the incomplete line, if any
350
351        if let Some(incomplete_line) = self.incomplete_line.as_mut() {
352            if let Some(line) = lines.next() {
353                trace!(
354                    "extending line from previous chunk: {:?}+{:?}",
355                    logify(incomplete_line),
356                    logify(line)
357                );
358
359                self.last_char_was_cr = false;
360                if !line.is_empty() {
361                    // Checking the last character handles lines where the last character is a
362                    // terminator, but also where the entire line is a terminator.
363                    match line.last().unwrap() {
364                        b'\r' => {
365                            incomplete_line.extend_from_slice(&line[..line.len() - 1]);
366                            let il = self.incomplete_line.take();
367                            self.complete_lines.push_back(il.unwrap());
368                            self.last_char_was_cr = true;
369                        }
370                        b'\n' => {
371                            incomplete_line.extend_from_slice(&line[..line.len() - 1]);
372                            let il = self.incomplete_line.take();
373                            self.complete_lines.push_back(il.unwrap());
374                        }
375                        _ => incomplete_line.extend_from_slice(line),
376                    };
377                }
378            }
379        }
380
381        let mut lines = lines.peekable();
382        while let Some(line) = lines.next() {
383            if let Some(actually_complete_line) = self.incomplete_line.take() {
384                // we saw the next line, so the previous one must have been complete after all
385                trace!(
386                    "previous line was complete: {:?}",
387                    logify(&actually_complete_line)
388                );
389                self.complete_lines.push_back(actually_complete_line);
390            }
391
392            if self.last_char_was_cr && line == *b"\n" {
393                // This is a continuation of a \r\n pair, so we can ignore this line. We do need to
394                // reset our flag though.
395                self.last_char_was_cr = false;
396                continue;
397            }
398
399            self.last_char_was_cr = false;
400            if line.ends_with(b"\r") {
401                self.complete_lines
402                    .push_back(line[..line.len() - 1].to_vec());
403                self.last_char_was_cr = true;
404            } else if line.ends_with(b"\n") {
405                // self isn't a continuation, but rather a line ending with a LF terminator.
406                self.complete_lines
407                    .push_back(line[..line.len() - 1].to_vec());
408            } else if line.is_empty() {
409                // this is the last line and it's empty, no need to buffer it
410                trace!("chunk ended with a line terminator");
411            } else if lines.peek().is_some() {
412                // this line isn't the last and we know from previous checks it doesn't end in a
413                // terminator, so we can consider it complete
414                self.complete_lines.push_back(line.to_vec());
415            } else {
416                // last line needs to be buffered as it may be incomplete
417                trace!("buffering incomplete line: {:?}", logify(line));
418                self.incomplete_line = Some(line.to_vec());
419            }
420        }
421
422        if log_enabled!(log::Level::Trace) {
423            for line in &self.complete_lines {
424                trace!("complete line: {:?}", logify(line));
425            }
426            if let Some(line) = &self.incomplete_line {
427                trace!("incomplete line: {:?}", logify(line));
428            }
429        }
430    }
431}
432
433#[cfg(test)]
434mod tests {
435    use std::str::FromStr;
436
437    use super::{Error::*, *};
438    use proptest::proptest;
439    use test_case::test_case;
440
441    fn field<'a>(key: &'a str, value: &'a str) -> Result<Option<(&'a str, &'a str)>> {
442        Ok(Some((key, value)))
443    }
444
445    /// Requires an event to be popped from the given parser.
446    /// Event properties can be asserted using a closure.
447    fn require_pop_event<F>(parser: &mut EventParser, f: F)
448    where
449        F: FnOnce(Event),
450    {
451        if let Some(SSE::Event(event)) = parser.get_event() {
452            f(event)
453        } else {
454            panic!("Event should have been received")
455        }
456    }
457
458    #[test]
459    fn test_logify_handles_code_point_boundaries() {
460        let phase = String::from_str(
461            "这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。",
462        )
463        .expect("Invalid sample string");
464
465        let input: &[u8] = phase.as_bytes();
466        let result = logify(input);
467
468        assert!(result == "这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如此。这是一条很长的消息,最初导致我们的代码出现恐慌。我希望情况不再如");
469    }
470
471    #[test]
472    fn test_parse_field_invalid() {
473        assert!(parse_field(b"").is_err());
474
475        match parse_field(b"\x80: invalid UTF-8") {
476            Err(InvalidLine(msg)) => assert!(msg.contains("Utf8Error")),
477            res => panic!("expected InvalidLine error, got {res:?}"),
478        }
479    }
480
481    #[test]
482    fn test_event_id_error_if_invalid_utf8() {
483        let mut bytes = Vec::from("id: ");
484        let mut invalid = vec![b'\xf0', b'\x28', b'\x8c', b'\xbc'];
485        bytes.append(&mut invalid);
486        bytes.push(b'\n');
487        let mut parser = EventParser::new();
488        assert!(parser.process_bytes(Bytes::from(bytes)).is_err());
489    }
490
491    #[test]
492    fn test_parse_field_comments() {
493        assert_eq!(parse_field(b":"), field("comment", ""));
494        assert_eq!(
495            parse_field(b":hello \0 world"),
496            field("comment", "hello \0 world")
497        );
498        assert_eq!(parse_field(b":event: foo"), field("comment", "event: foo"));
499    }
500
501    #[test]
502    fn test_parse_field_valid() {
503        assert_eq!(parse_field(b"event:foo"), field("event", "foo"));
504        assert_eq!(parse_field(b"event: foo"), field("event", "foo"));
505        assert_eq!(parse_field(b"event:  foo"), field("event", " foo"));
506        assert_eq!(parse_field(b"event:\tfoo"), field("event", "\tfoo"));
507        assert_eq!(parse_field(b"event: foo "), field("event", "foo "));
508
509        assert_eq!(parse_field(b"disconnect:"), field("disconnect", ""));
510        assert_eq!(parse_field(b"disconnect: "), field("disconnect", ""));
511        assert_eq!(parse_field(b"disconnect:  "), field("disconnect", " "));
512        assert_eq!(parse_field(b"disconnect:\t"), field("disconnect", "\t"));
513
514        assert_eq!(parse_field(b"disconnect"), field("disconnect", ""));
515
516        assert_eq!(parse_field(b" : foo"), field(" ", "foo"));
517        assert_eq!(parse_field(b"\xe2\x98\x83: foo"), field("☃", "foo"));
518    }
519
520    fn event(typ: &str, data: &str) -> SSE {
521        SSE::Event(Event {
522            data: data.to_string(),
523            id: None,
524            event_type: typ.to_string(),
525            retry: None,
526        })
527    }
528
529    fn event_with_id(typ: &str, data: &str, id: &str) -> SSE {
530        SSE::Event(Event {
531            data: data.to_string(),
532            id: Some(id.to_string()),
533            event_type: typ.to_string(),
534            retry: None,
535        })
536    }
537
538    #[test]
539    fn test_event_without_data_yields_no_event() {
540        let mut parser = EventParser::new();
541        assert!(parser.process_bytes(Bytes::from("id: abc\n\n")).is_ok());
542        assert!(parser.get_event().is_none());
543    }
544
545    #[test_case("unknown: value"; "unknown_with_value")]
546    #[test_case("unknown"; "unknown_without_colon")]
547    #[test_case("unknown:"; "unknown_without_value")]
548    #[test_case("DATA: value"; "uppercase_data")]
549    #[test_case("EVENT: update"; "uppercase_event")]
550    #[test_case("ID: cursor"; "uppercase_id")]
551    #[test_case("RETRY: 42"; "uppercase_retry")]
552    #[test_case("☃: value"; "unicode_name")]
553    #[test_case(" : value"; "space_name")]
554    fn test_unknown_field_blocks_are_ignored(field: &str) {
555        for line_ending in ["\n", "\r", "\r\n"] {
556            let input =
557                format!("{field}{line_ending}{line_ending}data: hello{line_ending}{line_ending}");
558
559            for split in 0..=input.len() {
560                let mut parser = EventParser::new();
561                let (first, second) = input.as_bytes().split_at(split);
562
563                assert!(parser.process_bytes(Bytes::copy_from_slice(first)).is_ok());
564                assert!(parser.process_bytes(Bytes::copy_from_slice(second)).is_ok());
565                assert_eq!(parser.get_event(), Some(event("message", "hello")));
566                assert!(parser.get_event().is_none());
567            }
568        }
569    }
570
571    #[test]
572    fn test_unknown_fields_preserve_event_data() {
573        let mut parser = EventParser::new();
574        assert!(parser
575            .process_bytes(Bytes::from(
576                "unknown: before\nid: cursor\nevent: update\ndata: hello\n\
577                 unknown\nretry: 42\nDATA: ignored\ndata: world\nunknown: after\n\n"
578            ))
579            .is_ok());
580        assert_eq!(
581            parser.get_event(),
582            Some(SSE::Event(Event {
583                event_type: "update".into(),
584                data: "hello\nworld".into(),
585                id: Some("cursor".into()),
586                retry: Some(42),
587            }))
588        );
589        assert!(parser.get_event().is_none());
590    }
591
592    #[test]
593    fn test_ignore_id_containing_null() {
594        let mut parser = EventParser::new();
595        assert!(parser
596            .process_bytes(Bytes::from("id: a\x00bc\nevent: add\ndata: abc\n\n"))
597            .is_ok());
598
599        if let Some(SSE::Event(event)) = parser.get_event() {
600            assert!(event.id.is_none());
601        } else {
602            panic!("Event should have been received");
603        }
604    }
605
606    #[test_case("event: add\ndata: hello\n\n", "add".into())]
607    #[test_case("data: hello\n\n", "message".into())]
608    fn test_event_can_parse_type_correctly(chunk: &'static str, event_type: String) {
609        let mut parser = EventParser::new();
610
611        assert!(parser.process_bytes(Bytes::from(chunk)).is_ok());
612
613        require_pop_event(&mut parser, |e| assert_eq!(event_type, e.event_type));
614    }
615
616    #[test_case("data: hello\n\n", event("message", "hello"); "parses event body with LF")]
617    #[test_case("data: hello\n\r", event("message", "hello"); "parses event body with LF and trailing CR")]
618    #[test_case("data: hello\r\n\n", event("message", "hello"); "parses event body with CRLF")]
619    #[test_case("data: hello\r\n\r", event("message", "hello"); "parses event body with CRLF and trailing CR")]
620    #[test_case("data: hello\r\r", event("message", "hello"); "parses event body with CR")]
621    #[test_case("data: hello\r\r\n", event("message", "hello"); "parses event body with CR and trailing CRLF")]
622    #[test_case("id: 1\ndata: hello\n\n", event_with_id("message", "hello", "1"))]
623    #[test_case("id: 😀\ndata: hello\n\n", event_with_id("message", "hello", "😀"))]
624    fn test_decode_chunks_simple(chunk: &'static str, event: SSE) {
625        let mut parser = EventParser::new();
626        assert!(parser.process_bytes(Bytes::from(chunk)).is_ok());
627        assert_eq!(parser.get_event().unwrap(), event);
628        assert!(parser.get_event().is_none());
629    }
630
631    #[test_case("persistent-event-id.sse"; "persistent-event-id.sse")]
632    fn test_last_id_persists_if_not_overridden(file: &str) {
633        let contents = read_contents_from_file(file);
634        let mut parser = EventParser::new();
635        assert!(parser.process_bytes(Bytes::from(contents)).is_ok());
636
637        require_pop_event(&mut parser, |e| assert_eq!(e.id, Some("1".into())));
638        require_pop_event(&mut parser, |e| assert_eq!(e.id, Some("1".into())));
639        require_pop_event(&mut parser, |e| assert_eq!(e.id, Some("3".into())));
640        require_pop_event(&mut parser, |e| assert_eq!(e.id, Some("3".into())));
641    }
642
643    #[test_case(b":hello\n"; "with LF")]
644    #[test_case(b":hello\r"; "with CR")]
645    #[test_case(b":hello\r\n"; "with CRLF")]
646    fn test_decode_chunks_comments_are_generated(chunk: &'static [u8]) {
647        let mut parser = EventParser::new();
648        assert!(parser.process_bytes(Bytes::from(chunk)).is_ok());
649        assert!(parser.get_event().is_some());
650    }
651
652    #[test]
653    fn test_comment_is_separate_from_event() {
654        let mut parser = EventParser::new();
655        let result = parser.process_bytes(Bytes::from(":comment\ndata:hello\n\n"));
656        assert!(result.is_ok());
657
658        let comment = parser.get_event();
659        assert!(matches!(comment, Some(SSE::Comment(_))));
660
661        let event = parser.get_event();
662        assert!(matches!(event, Some(SSE::Event(_))));
663
664        assert!(parser.get_event().is_none());
665    }
666
667    #[test]
668    fn test_comment_with_trailing_blank_line() {
669        let mut parser = EventParser::new();
670        let result = parser.process_bytes(Bytes::from(":comment\n\r\n\r"));
671        assert!(result.is_ok());
672
673        let comment = parser.get_event();
674        assert!(matches!(comment, Some(SSE::Comment(_))));
675
676        assert!(parser.get_event().is_none());
677    }
678
679    #[test_case(&["data:", "hello\n\n"], event("message", "hello"); "data split")]
680    #[test_case(&["data:hell", "o\n\n"], event("message", "hello"); "data truncated")]
681    fn test_decode_message_split_across_chunks(chunks: &[&'static str], event: SSE) {
682        let mut parser = EventParser::new();
683
684        if let Some((last, chunks)) = chunks.split_last() {
685            for chunk in chunks {
686                assert!(parser.process_bytes(Bytes::from(*chunk)).is_ok());
687                assert!(parser.get_event().is_none());
688            }
689
690            assert!(parser.process_bytes(Bytes::from(*last)).is_ok());
691            assert_eq!(parser.get_event(), Some(event));
692            assert!(parser.get_event().is_none());
693        } else {
694            panic!("Failed to split last");
695        }
696    }
697
698    #[test_case(&["data:hell", "o\n\ndata:", "world\n\n"], &[event("message", "hello"), event("message", "world")]; "with lf")]
699    #[test_case(&["data:hell", "o\r\rdata:", "world\r\r"], &[event("message", "hello"), event("message", "world")]; "with cr")]
700    #[test_case(&["data:hell", "o\r\n\ndata:", "world\r\n\n"], &[event("message", "hello"), event("message", "world")]; "with crlf")]
701    fn test_decode_multiple_messages_split_across_chunks(chunks: &[&'static str], events: &[SSE]) {
702        let mut parser = EventParser::new();
703
704        for chunk in chunks {
705            assert!(parser.process_bytes(Bytes::from(*chunk)).is_ok());
706        }
707
708        for event in events {
709            assert_eq!(parser.get_event().unwrap(), *event);
710        }
711
712        assert!(parser.get_event().is_none());
713    }
714
715    #[test]
716    fn test_decode_line_split_across_chunks() {
717        let mut parser = EventParser::new();
718        assert!(parser.process_bytes(Bytes::from("data:foo")).is_ok());
719        assert!(parser.process_bytes(Bytes::from("")).is_ok());
720        assert!(parser.process_bytes(Bytes::from("baz\n\n")).is_ok());
721        assert_eq!(parser.get_event(), Some(event("message", "foobaz")));
722        assert!(parser.get_event().is_none());
723
724        assert!(parser.process_bytes(Bytes::from("data:foo")).is_ok());
725        assert!(parser.process_bytes(Bytes::from("bar")).is_ok());
726        assert!(parser.process_bytes(Bytes::from("baz\n\n")).is_ok());
727        assert_eq!(parser.get_event(), Some(event("message", "foobarbaz")));
728        assert!(parser.get_event().is_none());
729    }
730
731    #[test]
732    fn test_decode_concatenates_multiple_values_for_same_field() {
733        let mut parser = EventParser::new();
734        assert!(parser.process_bytes(Bytes::from("data:hello\n")).is_ok());
735        assert!(parser.process_bytes(Bytes::from("data:world\n\n")).is_ok());
736        assert_eq!(parser.get_event(), Some(event("message", "hello\nworld")));
737        assert!(parser.get_event().is_none());
738    }
739
740    #[test_case("\n\n\n\n" ; "all LFs")]
741    #[test_case("\r\r\r\r" ; "all CRs")]
742    #[test_case("\r\n\r\n\r\n\r\n" ; "all CRLFs")]
743    fn test_decode_repeated_terminators(chunk: &'static str) {
744        let mut parser = EventParser::new();
745        assert!(parser.process_bytes(Bytes::from(chunk)).is_ok());
746
747        // spec seems unclear on whether this should actually dispatch empty events, but that seems
748        // unhelpful for all practical purposes
749        assert!(parser.get_event().is_none());
750    }
751
752    #[test]
753    fn test_decode_extra_terminators_between_events() {
754        let mut parser = EventParser::new();
755        assert!(parser
756            .process_bytes(Bytes::from("data: abc\n\n\ndata: def\n\n"))
757            .is_ok());
758
759        assert_eq!(parser.get_event(), Some(event("message", "abc")));
760        assert_eq!(parser.get_event(), Some(event("message", "def")));
761        assert!(parser.get_event().is_none());
762    }
763
764    #[test_case("one-event.sse"; "one-event.sse")]
765    #[test_case("one-event-crlf.sse"; "one-event-crlf.sse")]
766    fn test_decode_one_event(file: &str) {
767        let contents = read_contents_from_file(file);
768        let mut parser = EventParser::new();
769        assert!(parser.process_bytes(Bytes::from(contents)).is_ok());
770
771        require_pop_event(&mut parser, |e| {
772            assert_eq!(e.event_type, "patch");
773            assert!(e
774                .data
775                .contains(r#"path":"/flags/goals.02.featureWithGoals"#));
776        });
777    }
778
779    #[test_case("two-events.sse"; "two-events.sse")]
780    #[test_case("two-events-crlf.sse"; "two-events-crlf.sse")]
781    fn test_decode_two_events(file: &str) {
782        let contents = read_contents_from_file(file);
783        let mut parser = EventParser::new();
784        assert!(parser.process_bytes(Bytes::from(contents)).is_ok());
785
786        require_pop_event(&mut parser, |e| {
787            assert_eq!(e.event_type, "one");
788            assert_eq!(e.data, "One");
789        });
790
791        require_pop_event(&mut parser, |e| {
792            assert_eq!(e.event_type, "two");
793            assert_eq!(e.data, "Two");
794        });
795    }
796
797    #[test_case("big-event-followed-by-another.sse"; "big-event-followed-by-another.sse")]
798    #[test_case("big-event-followed-by-another-crlf.sse"; "big-event-followed-by-another-crlf.sse")]
799    fn test_decode_big_event_followed_by_another(file: &str) {
800        let contents = read_contents_from_file(file);
801        let mut parser = EventParser::new();
802        assert!(parser.process_bytes(Bytes::from(contents)).is_ok());
803
804        require_pop_event(&mut parser, |e| {
805            assert_eq!(e.event_type, "patch");
806            assert!(e.data.len() > 10_000);
807            assert!(e.data.contains(r#"path":"/flags/big.00.bigFeatureKey"#));
808        });
809
810        require_pop_event(&mut parser, |e| {
811            assert_eq!(e.event_type, "patch");
812            assert!(e
813                .data
814                .contains(r#"path":"/flags/goals.02.featureWithGoals"#));
815        });
816    }
817
818    fn read_contents_from_file(name: &str) -> Vec<u8> {
819        std::fs::read(format!("test-data/{name}"))
820            .unwrap_or_else(|_| panic!("couldn't read {name}"))
821    }
822
823    #[test]
824    fn test_event_parser_with_bom_header_split_across_chunks() {
825        let mut parser = EventParser::new();
826        // First chunk: partial BOM
827        assert!(parser
828            .process_bytes(Bytes::from(b"\xEF\xBB".as_slice()))
829            .is_ok());
830        assert!(parser.get_event().is_none());
831        // Second chunk: rest of BOM + data
832        assert!(parser
833            .process_bytes(Bytes::from(b"\xBFdata: hello\n\n".as_slice()))
834            .is_ok());
835        assert_eq!(parser.get_event(), Some(event("message", "hello")));
836        assert!(parser.get_event().is_none());
837    }
838
839    #[test]
840    fn test_event_parser_only_strips_initial_bom() {
841        let mut parser = EventParser::new();
842        assert!(parser
843            .process_bytes(Bytes::from(b"\xEF\xBB\xBFdata: first\n\n".as_slice()))
844            .is_ok());
845        assert_eq!(parser.get_event(), Some(event("message", "first")));
846
847        // Only the initial BOM is stripped, so this field name is not "data".
848        assert!(parser
849            .process_bytes(Bytes::from(b"\xEF\xBB\xBFdata: second\n\n".as_slice()))
850            .is_ok());
851        assert!(parser.get_event().is_none());
852        assert!(parser.process_bytes(Bytes::from("data: third\n\n")).is_ok());
853        assert_eq!(parser.get_event(), Some(event("message", "third")));
854        assert!(parser.get_event().is_none());
855    }
856
857    proptest! {
858        #[test]
859        fn test_decode_and_buffer_lines_does_not_crash(next in "(\r\n|\r|\n)*event: [^\n\r:]*(\r\n|\r|\n)", previous in "(\r\n|\r|\n)*event: [^\n\r:]*(\r\n|\r|\n)") {
860            let mut parser = EventParser::new();
861            parser.incomplete_line = Some(previous.as_bytes().to_vec());
862            parser.decode_and_buffer_lines(Bytes::from(next));
863        }
864    }
865}