Skip to main content

ocre/runtime/realtime/
channel.rs

1//! The `OcreChannel` Durable Object class. It sits in its own module because
2//! `#[durable_object]` generates public wasm-bindgen glue that cannot carry
3//! docs; the parent allows `missing_docs` for this module only.
4
5use worker::{
6    DurableObject, Env, Method, Request, State, WebSocket, WebSocketIncomingMessage, WebSocketPair, durable_object,
7};
8
9use crate::realtime::{SUBSCRIBER_HEADER, Subscriber, close_code};
10
11/// The Durable Object class behind each realtime channel, exported by Ocre under the name `OcreChannel`.
12///
13/// One instance per channel name holds that channel's WebSockets. It
14/// accepts them with the WebSocket Hibernation API, so it is evicted from
15/// memory between broadcasts while browsers stay connected (hibernated
16/// sockets cost no duration), and it stores nothing but each socket's
17/// attachment (its identity and whether it may publish). What subscribers
18/// send is ignored, or relayed to the others as JSON for sockets connected
19/// with [`WebSocketUpgrade::rebroadcast`](crate::realtime::WebSocketUpgrade::rebroadcast);
20/// it completes the closing handshakes browsers start. Apps never
21/// call it directly: they use [`broadcast`](crate::realtime::broadcast) and
22/// [`WebSocketUpgrade::connect`](crate::realtime::WebSocketUpgrade::connect).
23///
24/// Free plan: Durable Objects must be SQLite-backed (`storage: "sqlite"`);
25/// one object accepts up to 32,768 WebSockets. cloudflare.config.ts declares it
26/// (`ocre g scaffold ... --realtime` adds this, with `demo` as the app name):
27///
28/// ```ts
29/// // worker.env
30/// CHANNELS: bindings.durableObject({ worker: "demo", exportName: "OcreChannel" }),
31/// // worker.exports
32/// OcreChannel: exports.durableObject({ storage: "sqlite" }),
33/// ```
34#[durable_object(websocket)]
35pub struct OcreChannel {
36    state: State,
37}
38
39impl DurableObject for OcreChannel {
40    fn new(state: State, _env: Env) -> Self {
41        Self { state }
42    }
43
44    /// `POST`: send the body to every socket, answer how many got it.
45    /// Anything else: a WebSocket handshake forwarded by `connect`.
46    async fn fetch(&self, mut req: Request) -> worker::Result<worker::Response> {
47        if req.method() == Method::Post {
48            let message = req.text().await?;
49            let sent = self.state.get_websockets().iter().filter(|ws| ws.send_with_str(&message).is_ok()).count();
50            return worker::Response::ok(sent.to_string());
51        }
52        let subscriber = Subscriber::from_header(req.headers().get(SUBSCRIBER_HEADER)?.as_deref());
53        let pair = WebSocketPair::new()?;
54        if let Some(joined) = subscriber.presence("joined") {
55            self.send_to_others(&pair.server, &joined);
56        }
57        self.state.accept_web_socket(&pair.server);
58        pair.server.serialize_attachment(&subscriber)?;
59        worker::Response::from_websocket(pair.client)
60    }
61
62    /// Relays text from sockets allowed to publish to every other socket; ignores the rest.
63    async fn websocket_message(&self, ws: WebSocket, message: WebSocketIncomingMessage) -> worker::Result<()> {
64        let WebSocketIncomingMessage::String(text) = message else {
65            return Ok(());
66        };
67        let subscriber: Subscriber = ws.deserialize_attachment()?.unwrap_or_default();
68        let Some(relayed) = subscriber.relay(&text) else {
69            return Ok(());
70        };
71        self.send_to_others(&ws, &relayed);
72        Ok(())
73    }
74
75    /// Completes the closing handshake the browser started.
76    async fn websocket_close(&self, ws: WebSocket, code: usize, reason: String, _clean: bool) -> worker::Result<()> {
77        self.announce_departure(&ws);
78        // Already closed when the runtime auto-replied (compatibility date 2026-04-07 or later).
79        let _ = ws.close(Some(close_code(code)), Some(reason));
80        Ok(())
81    }
82
83    async fn websocket_error(&self, ws: WebSocket, _error: worker::Error) -> worker::Result<()> {
84        self.announce_departure(&ws);
85        Ok(())
86    }
87}
88
89impl OcreChannel {
90    /// Sends `text` to every socket of the channel but `sender`.
91    fn send_to_others(&self, sender: &WebSocket, text: &str) {
92        let sender: &worker::web_sys::WebSocket = sender.as_ref();
93        for other in self.state.get_websockets() {
94            let target: &worker::web_sys::WebSocket = other.as_ref();
95            if target != sender {
96                let _ = other.send_with_str(text);
97            }
98        }
99    }
100
101    /// Tells the others that an identified publisher left.
102    fn announce_departure(&self, ws: &WebSocket) {
103        let subscriber: Subscriber = ws.deserialize_attachment().ok().flatten().unwrap_or_default();
104        if let Some(left) = subscriber.presence("left") {
105            self.send_to_others(ws, &left);
106        }
107    }
108}