File size: 1,558 Bytes
9ebf6d4 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 | use codex_http_client::ByteStream;
use codex_http_client::StreamError;
use eventsource_stream::Eventsource;
use futures::StreamExt;
use tokio::sync::mpsc;
use tokio::time::Duration;
use tokio::time::timeout;
/// Minimal SSE helper that forwards raw `data:` frames as UTF-8 strings.
///
/// Errors and idle timeouts are sent as `Err(StreamError)` before the task exits.
pub fn sse_stream(
stream: ByteStream,
idle_timeout: Duration,
tx: mpsc::Sender<Result<String, StreamError>>,
) {
tokio::spawn(async move {
let mut stream = stream
.map(|res| res.map_err(|e| StreamError::Stream(e.to_string())))
.eventsource();
loop {
match timeout(idle_timeout, stream.next()).await {
Ok(Some(Ok(ev))) => {
if tx.send(Ok(ev.data.clone())).await.is_err() {
return;
}
}
Ok(Some(Err(e))) => {
let _ = tx.send(Err(StreamError::Stream(e.to_string()))).await;
return;
}
Ok(None) => {
let _ = tx
.send(Err(StreamError::Stream(
"stream closed before completion".into(),
)))
.await;
return;
}
Err(_) => {
let _ = tx.send(Err(StreamError::Timeout)).await;
return;
}
}
}
});
}
|