uptrakit_openapi_client/
sse.rs1#[derive(Debug, Clone, PartialEq, Eq)]
11pub struct RawSseEvent {
12 pub event_type: String,
14 pub data: String,
16 pub id: Option<String>,
18}
19
20#[derive(Debug, thiserror::Error)]
22pub enum SseError {
23 #[error("stream read error: {0}")]
24 Transport(#[from] reqwest::Error),
25}
26
27#[expect(
35 clippy::string_slice,
36 reason = "pos is sourced from find_event_boundary which searches for ASCII byte sequences; all slice boundaries are ASCII-safe"
37)]
38pub fn parse_sse_stream(
39 response: reqwest::Response,
40) -> impl futures_util::Stream<Item = Result<RawSseEvent, SseError>> {
41 futures_util::stream::unfold(
42 (response, String::new()),
43 |(mut response, mut buffer)| async move {
44 loop {
45 if let Some(pos) = find_event_boundary(&buffer) {
47 let event_text = buffer[..pos].to_string();
48 let skip = if buffer[pos..].starts_with("\r\n\r\n") {
50 4
51 } else if buffer[pos..].starts_with("\n\n") {
52 2
53 } else {
54 2
56 };
57 buffer = buffer[pos + skip..].to_string();
58
59 if let Some(event) = parse_event(&event_text) {
60 return Some((Ok(event), (response, buffer)));
61 }
62 continue;
64 }
65
66 match response.chunk().await {
68 Ok(Some(bytes)) => {
69 let text = String::from_utf8_lossy(&bytes);
70 buffer.push_str(&text);
71 }
72 Ok(None) => {
73 if !buffer.trim().is_empty() {
75 let event_text = std::mem::take(&mut buffer);
76 if let Some(event) = parse_event(&event_text) {
77 return Some((Ok(event), (response, buffer)));
78 }
79 }
80 return None;
81 }
82 Err(e) => {
83 return Some((Err(SseError::Transport(e)), (response, buffer)));
84 }
85 }
86 }
87 },
88 )
89}
90
91fn find_event_boundary(s: &str) -> Option<usize> {
93 if let Some(pos) = s.find("\r\n\r\n") {
95 if let Some(nn_pos) = s.find("\n\n") {
97 return Some(nn_pos.min(pos));
98 }
99 return Some(pos);
100 }
101 if let Some(pos) = s.find("\n\n") {
102 return Some(pos);
103 }
104 s.find("\r\r")
105}
106
107fn parse_event(text: &str) -> Option<RawSseEvent> {
110 let mut event_type = None;
111 let mut data_parts: Vec<&str> = Vec::new();
112 let mut id = None;
113
114 for line in text.lines() {
115 if line.starts_with(':') {
116 continue;
118 }
119
120 if let Some(value) = line.strip_prefix("event:") {
121 event_type = Some(value.trim().to_string());
122 } else if let Some(value) = line.strip_prefix("data:") {
123 data_parts.push(value.strip_prefix(' ').unwrap_or(value));
124 } else if let Some(value) = line.strip_prefix("id:") {
125 id = Some(value.trim().to_string());
126 }
127 }
129
130 if data_parts.is_empty() {
131 return None;
132 }
133
134 Some(RawSseEvent {
135 event_type: event_type.unwrap_or_else(|| "message".to_string()),
136 data: data_parts.join("\n"),
137 id,
138 })
139}
140
141#[cfg(test)]
142mod tests {
143 use super::*;
144
145 #[test]
146 fn parse_single_event() {
147 let text = "event: output\ndata: hello world";
148 let event = parse_event(text).expect("should parse");
149 assert_eq!(event.event_type, "output");
150 assert_eq!(event.data, "hello world");
151 assert_eq!(event.id, None);
152 }
153
154 #[test]
155 fn parse_event_default_type() {
156 let text = "data: just data";
157 let event = parse_event(text).expect("should parse");
158 assert_eq!(event.event_type, "message");
159 assert_eq!(event.data, "just data");
160 }
161
162 #[test]
163 fn parse_event_with_id() {
164 let text = "event: completed\ndata: {\"status\":\"done\"}\nid: 42";
165 let event = parse_event(text).expect("should parse");
166 assert_eq!(event.event_type, "completed");
167 assert_eq!(event.data, "{\"status\":\"done\"}");
168 assert_eq!(event.id.as_deref(), Some("42"));
169 }
170
171 #[test]
172 fn parse_event_multi_data_lines() {
173 let text = "data: line1\ndata: line2\ndata: line3";
174 let event = parse_event(text).expect("should parse");
175 assert_eq!(event.data, "line1\nline2\nline3");
176 }
177
178 #[test]
179 fn parse_event_skips_comments() {
180 let text = ": this is a comment\nevent: output\ndata: payload";
181 let event = parse_event(text).expect("should parse");
182 assert_eq!(event.event_type, "output");
183 assert_eq!(event.data, "payload");
184 }
185
186 #[test]
187 fn parse_event_no_data_returns_none() {
188 let text = "event: output\n: comment only";
189 assert!(parse_event(text).is_none());
190 }
191
192 #[test]
193 fn find_boundary_double_newline() {
194 assert_eq!(find_event_boundary("data: hello\n\ndata: world"), Some(11));
195 }
196
197 #[test]
198 fn find_boundary_crlf() {
199 assert_eq!(
200 find_event_boundary("data: hello\r\n\r\ndata: world"),
201 Some(11)
202 );
203 }
204
205 #[test]
206 fn find_boundary_none() {
207 assert_eq!(find_event_boundary("data: hello\n"), None);
208 }
209}