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}