use anyhow::Result;
use futures_util::StreamExt;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Frame {
pub event: String,
pub data: String,
}
#[derive(Debug, Default)]
pub struct SseFrameParser {
line: Vec<u8>,
event: String,
data: String,
}
impl SseFrameParser {
pub fn push(&mut self, byte: u8) -> Option<Frame> {
if byte != b'\n' {
self.line.push(byte);
return None;
}
if self.line.last() == Some(&b'\r') {
self.line.pop();
}
let line = String::from_utf8_lossy(&self.line).into_owned();
self.line.clear();
self.consume_line(&line)
}
fn consume_line(&mut self, line: &str) -> Option<Frame> {
if line.is_empty() {
let frame = (!self.data.is_empty()).then(|| Frame {
event: std::mem::take(&mut self.event),
data: std::mem::take(&mut self.data),
});
self.event.clear();
self.data.clear();
return frame;
}
if line.starts_with(':') {
return None;
}
if let Some(name) = line.strip_prefix("event:") {
self.event = name.trim().to_string();
} else if let Some(chunk) = line.strip_prefix("data:") {
let chunk = chunk.strip_prefix(' ').unwrap_or(chunk);
if !self.data.is_empty() {
self.data.push('\n');
}
self.data.push_str(chunk);
}
None
}
}
pub async fn stream_events(
endpoint: &str,
query: &str,
on_frame: &mut impl FnMut(Frame),
) -> Result<()> {
let client = reqwest::Client::new();
let response = client
.get(format!("http://{endpoint}/events{query}"))
.header("Accept", "text/event-stream")
.send()
.await?
.error_for_status()?;
let mut bytes = response.bytes_stream();
let mut parser = SseFrameParser::default();
while let Some(chunk) = bytes.next().await {
for byte in chunk? {
if let Some(frame) = parser.push(byte) {
on_frame(frame);
}
}
}
Ok(())
}