ocre/sse.rs
1//! Server-Sent Events: a response that sends events as they happen (Rails' `ActionController::Live` with `SSE`).
2//!
3//! A handler returns [`stream()`](fn@stream)`(state, step)`: Ocre calls `step` for each
4//! event, and the Worker streams every event to the client as soon as it is
5//! produced, until `step` returns `None`. The response is axum's [`Sse`]
6//! (`Content-Type: text/event-stream`, `Cache-Control: no-cache`), which
7//! browsers read with `new EventSource(url)` or htmx's SSE extension.
8//! Pace events with [`sleep`](crate::sleep).
9//!
10//! Streams fit work that one request follows from start to end: the
11//! progress of an import, tokens from an AI model, a countdown. To push
12//! changes to many open pages (a new comment), use [`realtime`] channels,
13//! which do not keep a request open per page.
14//!
15//! # Free plan
16//!
17//! - **CPU**: only the work of each `step` counts toward the 10 ms per
18//! request; waiting in [`sleep`](crate::sleep) or on I/O does not.
19//! - **Duration**: a response may stream for as long as the client stays
20//! connected; when the client goes away, the stream stops at the next event.
21//! - **One invocation**: the whole stream is one request. Its binding calls
22//! count toward the per-request limits (50 subrequests on the free plan,
23//! D1 queries included), so a stream cannot poll D1 every second for
24//! minutes; end it and let the client reconnect, or use [`realtime`].
25//! - **Reconnects**: `EventSource` reconnects about 3 s after a stream ends,
26//! and each reconnect is a new request (100,000 a day on the free plan).
27//! Send a last event (e.g. `event: done`) on which the page calls
28//! `source.close()`, or set [`Event::retry`] to reconnect less often.
29//!
30//! [`realtime`]: https://docs.rs/ocre/latest/ocre/realtime/index.html
31//!
32//! # Examples
33//!
34//! ```no_run
35//! use std::time::Duration;
36//!
37//! use axum::{Router, routing::get};
38//! use ocre::{Ctx, sse::{self, Event}};
39//!
40//! pub fn routes() -> Router<Ctx> {
41//! Router::new().route("/countdown", get(countdown))
42//! }
43//!
44//! /// `data: 3`, `data: 2`, `data: 1` one second apart, then `event: done`.
45//! async fn countdown() -> impl axum::response::IntoResponse {
46//! sse::stream(Some(3), |left: Option<u32>| async move {
47//! let left = left?;
48//! ocre::sleep(Duration::from_secs(1)).await;
49//! Some(match left {
50//! 0 => (Event::default().event("done").data(""), None),
51//! n => (Event::default().data(n.to_string()), Some(n - 1)),
52//! })
53//! })
54//! }
55//! ```
56//!
57//! In the page: `const source = new EventSource("/countdown");
58//! source.addEventListener("done", () => source.close());`.
59
60use std::{convert::Infallible, future::Future};
61
62pub use axum::response::sse::{Event, Sse};
63use futures_util::{Stream, StreamExt as _, stream};
64
65/// An event stream: `step(state)` gives the next event and the next state, or `None` to end the stream.
66///
67/// `step` runs once per event, only while the client reads the response;
68/// the first call happens after the handler returned (so the response
69/// headers leave at once). `state` carries what the next step needs (a
70/// counter, a cursor, an id); it must be `Send` like every Ocre value. See
71/// the [module documentation](self) for costs and limits.
72///
73/// # Examples
74///
75/// ```
76/// use axum::response::IntoResponse;
77/// use ocre::sse::{self, Event};
78///
79/// let response = sse::stream(1, |n: u32| async move {
80/// (n <= 2).then(|| (Event::default().id(n.to_string()).data(format!("step {n}")), n + 1))
81/// })
82/// .into_response();
83/// assert_eq!(response.headers()["content-type"], "text/event-stream");
84/// let body = pollster::block_on(axum::body::to_bytes(response.into_body(), usize::MAX)).unwrap();
85/// assert_eq!(body, "id: 1\ndata: step 1\n\nid: 2\ndata: step 2\n\n");
86/// ```
87pub fn stream<T, F, Fut>(state: T, step: F) -> Sse<impl Stream<Item = Result<Event, Infallible>> + Send + 'static>
88where
89 T: Send + 'static,
90 F: FnMut(T) -> Fut + Send + 'static,
91 Fut: Future<Output = Option<(Event, T)>> + Send + 'static,
92{
93 Sse::new(stream::unfold(state, step).map(Ok))
94}
95
96#[cfg(test)]
97#[path = "../tests/sse.rs"]
98mod tests;