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 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 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#[derive(Debug)]
147enum BomHeaderState {
148 Parsing(Vec<u8>),
149 Consumed,
150}
151
152const BOM_HEADER: &[u8] = b"\xEF\xBB\xBF";
153
154fn 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 complete_lines: VecDeque<Vec<u8>>,
177 incomplete_line: Option<Vec<u8>>,
180 last_char_was_cr: bool,
183 event_data: Option<EventData>,
185 last_event_id: Option<String>,
187 sse: VecDeque<SSE>,
188 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 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 Bytes::from_owner(owned_rest)
230 } else {
231 return Ok(());
232 }
233 } else {
234 bytes
235 };
236
237 self.decode_and_buffer_lines(bytes_to_process);
247 self.parse_complete_lines_into_event()?;
248
249 Ok(())
250 }
251
252 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 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 fn decode_and_buffer_lines(&mut self, chunk: Bytes) {
344 let mut lines = chunk.split_inclusive(|&b| b == b'\n' || b == b'\r');
345 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 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 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 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.complete_lines
407 .push_back(line[..line.len() - 1].to_vec());
408 } else if line.is_empty() {
409 trace!("chunk ended with a line terminator");
411 } else if lines.peek().is_some() {
412 self.complete_lines.push_back(line.to_vec());
415 } else {
416 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 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 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 assert!(parser
828 .process_bytes(Bytes::from(b"\xEF\xBB".as_slice()))
829 .is_ok());
830 assert!(parser.get_event().is_none());
831 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 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}