1use crate::{
4 NetError,
5 line::{DEFAULT_MAX_LINE_BYTES, LineDecoder},
6};
7
8#[derive(Clone, Debug, PartialEq, Eq)]
10pub struct SseEvent {
11 pub event: Option<String>,
13 pub data: String,
15}
16
17#[derive(Debug)]
34pub struct SseDecoder {
35 lines: LineDecoder,
36 event: Option<String>,
37 data: Vec<String>,
38 data_bytes: usize,
39 max_line_bytes: Option<usize>,
40 have_record: bool,
41}
42
43impl Default for SseDecoder {
44 fn default() -> Self {
45 Self::new()
46 }
47}
48
49impl SseDecoder {
50 pub fn new() -> Self {
52 Self::with_max_line_bytes(DEFAULT_MAX_LINE_BYTES)
53 }
54
55 pub fn unbounded() -> Self {
59 Self {
60 lines: LineDecoder::unbounded(),
61 event: None,
62 data: Vec::new(),
63 data_bytes: 0,
64 max_line_bytes: None,
65 have_record: false,
66 }
67 }
68
69 pub fn with_max_line_bytes(max_line_bytes: usize) -> Self {
72 Self {
73 lines: LineDecoder::with_max_line_bytes(max_line_bytes),
74 event: None,
75 data: Vec::new(),
76 data_bytes: 0,
77 max_line_bytes: Some(max_line_bytes),
78 have_record: false,
79 }
80 }
81
82 pub fn push(&mut self, bytes: &[u8]) -> Result<Vec<SseEvent>, NetError> {
84 self.push_checked(bytes)
85 }
86
87 pub fn push_checked(&mut self, bytes: &[u8]) -> Result<Vec<SseEvent>, NetError> {
89 let mut events = Vec::new();
90 for line in self.lines.push_checked(bytes)? {
91 let text = String::from_utf8_lossy(&line);
94 if let Some(event) = self.push_line_checked(&text)? {
95 events.push(event);
96 }
97 }
98 Ok(events)
99 }
100
101 pub fn push_line(&mut self, line: &str) -> Option<SseEvent> {
105 self.push_line_checked(line)
106 .expect("sse decoder cap exceeded; use push_line_checked for fallible handling")
107 }
108
109 pub fn push_line_checked(&mut self, line: &str) -> Result<Option<SseEvent>, NetError> {
112 self.check_line(line.len())?;
113 if line.is_empty() {
114 return Ok(self.dispatch());
115 }
116 if line.starts_with(':') {
117 return Ok(None);
120 }
121 let (field, value) = match line.split_once(':') {
122 Some((field, rest)) => (field, rest.strip_prefix(' ').unwrap_or(rest)),
123 None => (line, ""),
124 };
125 match field {
126 "event" => {
127 self.event = Some(value.to_owned());
128 self.have_record = true;
129 }
130 "data" => {
131 self.add_data(value)?;
132 self.data.push(value.to_owned());
133 self.have_record = true;
134 }
135 _ => {}
136 }
137 Ok(None)
138 }
139
140 fn dispatch(&mut self) -> Option<SseEvent> {
141 if !self.have_record {
142 return None;
143 }
144 let event = self.event.take();
145 let data = std::mem::take(&mut self.data).join("\n");
146 self.data_bytes = 0;
147 self.have_record = false;
148 Some(SseEvent { event, data })
149 }
150
151 pub fn flush(&mut self) -> Option<SseEvent> {
154 self.flush_checked()
155 .expect("sse decoder cap exceeded; use flush_checked for fallible handling")
156 }
157
158 pub fn flush_checked(&mut self) -> Result<Option<SseEvent>, NetError> {
161 if let Some(line) = self.lines.flush_checked()? {
162 let text = String::from_utf8_lossy(&line);
163 if let Some(event) = self.push_line_checked(&text)? {
164 return Ok(Some(event));
165 }
166 }
167 Ok(self.dispatch())
168 }
169
170 fn add_data(&mut self, value: &str) -> Result<(), NetError> {
171 let separator = usize::from(!self.data.is_empty());
172 let next_len = self
173 .data_bytes
174 .saturating_add(separator)
175 .saturating_add(value.len());
176 self.check_line(next_len)?;
177 self.data_bytes = next_len;
178 Ok(())
179 }
180
181 fn check_line(&self, len: usize) -> Result<(), NetError> {
182 if let Some(max) = self.max_line_bytes
183 && len > max
184 {
185 return Err(NetError::LineTooLong { max });
186 }
187 Ok(())
188 }
189}