#[derive(Debug, Clone)]
pub struct SseEvent {
pub event_type: Option<String>,
pub data: String,
}
pub struct SseParser {
buffer: String,
}
impl SseParser {
pub fn new() -> Self {
Self { buffer: String::new() }
}
pub fn feed(&mut self, chunk: &str) -> Vec<SseEvent> {
self.buffer.push_str(chunk);
let mut events = Vec::new();
loop {
let unix = self.buffer.find("\n\n").map(|pos| (pos, 2));
let windows = self.buffer.find("\r\n\r\n").map(|pos| (pos, 4));
let split_pos = match (unix, windows) {
(Some(unix), Some(windows)) => Some(if unix.0 < windows.0 { unix } else { windows }),
(Some(delimiter), None) | (None, Some(delimiter)) => Some(delimiter),
(None, None) => None,
};
match split_pos {
Some((pos, delim_len)) => {
let event_str: String = self.buffer.drain(..pos + delim_len).collect();
if let Some(event) = Self::parse_event(&event_str) {
events.push(event);
}
}
None => {
break;
}
}
}
events
}
fn parse_event(event_str: &str) -> Option<SseEvent> {
let mut event_type = None;
let mut data_parts = Vec::new();
for line in event_str.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
if let Some(value) = line.strip_prefix("event:") {
event_type = Some(value.trim().to_string());
} else if let Some(value) = line.strip_prefix("data:") {
data_parts.push(value.trim().to_string());
}
}
if data_parts.is_empty() {
return None;
}
let data = data_parts.join("\n");
Some(SseEvent { event_type, data })
}
pub fn has_remaining(&self) -> bool {
!self.buffer.is_empty()
}
#[allow(dead_code)]
pub fn remaining(&self) -> &str {
&self.buffer
}
}
impl Default for SseParser {
fn default() -> Self {
Self::new()
}
}