Skip to main content

WebSocketUpgrade

Struct WebSocketUpgrade 

Source
pub struct WebSocketUpgrade { /* private fields */ }
Expand description

Axum extractor for a WebSocket handshake (Upgrade: websocket), finished with connect.

Check who may listen in the handler, then call connect. Browsers send the session cookie with the handshake, so CurrentUser and Session work in the same handler. serve refuses handshakes from other sites (403), as it does for forms.

Rejection: requests without Upgrade: websocket get Error::BadRequest (400), rendered as an HTML page in full-stack apps and as JSON in API-only apps.

Free plan: each connection (and each reconnection) is one Worker request and one Durable Object request (100,000 a day each); hibernated sockets cost nothing between messages.

§Examples

use axum::{extract::{Path, State}, response::Response};
use ocre::{Ctx, Result, realtime::WebSocketUpgrade};

// GET /realtime/{channel}
async fn connect(State(ctx): State<Ctx>, Path(channel): Path<String>, upgrade: WebSocketUpgrade)
    -> Result<Response> {
    upgrade.connect(&ctx, &channel).await
}

Implementations§

Source§

impl WebSocketUpgrade

Source

pub fn identified_by(self, identity: impl Into<String>) -> Self

Names who is connecting, like Action Cable’s identified_by :current_user.

The identity (a user id, a display name…) stays with the socket in the channel object, even while it hibernates, and is the from of the messages rebroadcast relays, so clients cannot forge it. At most MAX_IDENTITY_LEN bytes; a longer one makes connect fail with Error::Internal.

§Examples
use axum::{extract::{Path, State}, response::Response};
use ocre::{Ctx, Result, realtime::WebSocketUpgrade};

async fn connect(State(ctx): State<Ctx>, Path(channel): Path<String>, upgrade: WebSocketUpgrade)
    -> Result<Response> {
    let user = CurrentUser { id: 7 }; // `user: CurrentUser` as an extractor after `ocre g auth`
    upgrade.identified_by(user.id.to_string()).connect(&ctx, &channel).await
}
Source

pub fn rebroadcast(self) -> Self

Relays what this client sends to the channel’s other subscribers, like a channel that rebroadcasts client data.

Without it, messages from clients are ignored (they only listen). With it, each text message of at most MAX_REBROADCAST_BYTES is sent to every other socket of the channel as JSON: {"from": "<identity>", "data": <message>}, where from is the identified_by identity (null without one) and data is the message parsed as JSON (a string when it is not JSON). An identified publisher’s arrival and departure reach the others as {"event": "joined", "from": "<identity>"} and {"event": "left", ...} (presence, Action Cable’s subscribed and unsubscribed hooks). No app code runs: use it for typing indicators, cursors or chat between JavaScript clients; send anything that must be checked or stored to an ordinary route that saves it and calls broadcast. Relayed messages are never HTML swaps, so a client cannot inject markup into other pages.

Free plan: Cloudflare bills incoming WebSocket messages to a Durable Object at 20 messages per request (100,000 requests a day); relaying is free.

§Examples
use axum::{extract::{Path, State}, response::Response};
use ocre::{Ctx, Result, realtime::WebSocketUpgrade};

// GET /chat/{room}: `ws.send(JSON.stringify({text: "hi"}))` reaches the others as
// {"from":"ada","data":{"text":"hi"}}.
async fn chat(State(ctx): State<Ctx>, Path(room): Path<String>, upgrade: WebSocketUpgrade)
    -> Result<Response> {
    upgrade.identified_by("ada").rebroadcast().connect(&ctx, &format!("chat:{room}")).await
}
Source§

impl WebSocketUpgrade

Source

pub fn connect( self, ctx: &Ctx, channel: &str, ) -> impl Future<Output = Result<Response>> + Send + use<>

Connects the browser to channel, returning the 101 Switching Protocols response to send back.

The handshake is forwarded to the channel’s OcreChannel object, which accepts the WebSocket (hibernating). Check who may listen before calling it: list the allowed channels, or use a channel per record or user (post:12). The returned future is Send.

Free plan: one Durable Object request per connection and reconnection.

§Errors
  • Error::BadRequest (400) when channel is not a valid name (1 to MAX_CHANNEL_LEN ASCII letters, digits, _, -, . or :).
  • Error::Internal (500) when the CHANNELS Durable Object binding is missing (the message names the cloudflare.config.ts entries to add) or the channel object cannot be reached.
§Examples
use axum::{extract::{Path, State}, response::Response};
use ocre::{Ctx, Error, Result, realtime::WebSocketUpgrade};

async fn connect(State(ctx): State<Ctx>, Path(channel): Path<String>, upgrade: WebSocketUpgrade)
    -> Result<Response> {
    match channel.as_str() {
        "posts" => {}
        _ => return Err(Error::NotFound),
    }
    upgrade.connect(&ctx, &channel).await
}

Trait Implementations§

Source§

impl Debug for WebSocketUpgrade

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<S: Sync> FromRequestParts<S> for WebSocketUpgrade

Source§

type Rejection = Error

If the extractor fails it’ll use this “rejection” type. A rejection is a kind of error that can be converted into a response.
Source§

async fn from_request_parts( parts: &mut Parts, _state: &S, ) -> Result<Self, Error>

Perform the extraction.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<S, T> FromRequest<S, ViaParts> for T
where S: Send + Sync, T: FromRequestParts<S>,

§

type Rejection = <T as FromRequestParts<S>>::Rejection

If the extractor fails it’ll use this “rejection” type. A rejection is a kind of error that can be converted into a response.
§

fn from_request( req: Request<Body>, state: &S, ) -> impl Future<Output = Result<T, <T as FromRequest<S, ViaParts>>::Rejection>>

Perform the extraction.
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<S, T> Upcast<T> for S
where T: UpcastFrom<S> + ?Sized, S: ?Sized,

Source§

fn upcast(&self) -> &T
where Self: ErasableGeneric, T: Sized + ErasableGeneric<Repr = Self::Repr>,

Perform a zero-cost type-safe upcast to a wider ref type within the Wasm bindgen generics type system. Read more
Source§

fn upcast_into(self) -> T
where Self: Sized + ErasableGeneric, T: Sized + ErasableGeneric<Repr = Self::Repr>,

Perform a zero-cost type-safe upcast to a wider type within the Wasm bindgen generics type system. Read more
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V