use eventsource_stream::Eventsource;
use futures::stream::{Stream, StreamExt};
use std::pin::Pin;
use crate::error::{Error, Result};
use crate::utils::parse_json_value_strict_str;
pub fn parse_sse_stream(
response: reqwest::Response,
) -> Pin<Box<dyn Stream<Item = Result<serde_json::Value>> + Send>> {
let event_stream = response.bytes_stream().eventsource();
let json_stream = event_stream.filter_map(|event_result| async move {
match event_result {
Ok(event) => {
let data = event.data;
if data == "[DONE]" {
return None;
}
match parse_json_value_strict_str(&data) {
Ok(json) => Some(Ok(json)),
Err(e) => Some(Err(Error::Json(e))),
}
}
Err(e) => Some(Err(Error::Inference(format!("SSE error: {}", e)))),
}
});
Box::pin(json_stream)
}