Skip to main content

stream

Function stream 

Source
pub fn stream<T, F, Fut>(
    state: T,
    step: F,
) -> Sse<impl Stream<Item = Result<Event, Infallible>> + Send + 'static>
where T: Send + 'static, F: FnMut(T) -> Fut + Send + 'static, Fut: Future<Output = Option<(Event, T)>> + Send + 'static,
Expand description

An event stream: step(state) gives the next event and the next state, or None to end the stream.

step runs once per event, only while the client reads the response; the first call happens after the handler returned (so the response headers leave at once). state carries what the next step needs (a counter, a cursor, an id); it must be Send like every Ocre value. See the module documentation for costs and limits.

ยงExamples

use axum::response::IntoResponse;
use ocre::sse::{self, Event};

let response = sse::stream(1, |n: u32| async move {
    (n <= 2).then(|| (Event::default().id(n.to_string()).data(format!("step {n}")), n + 1))
})
.into_response();
assert_eq!(response.headers()["content-type"], "text/event-stream");
let body = pollster::block_on(axum::body::to_bytes(response.into_body(), usize::MAX)).unwrap();
assert_eq!(body, "id: 1\ndata: step 1\n\nid: 2\ndata: step 2\n\n");