Reference

WebSockets

A WebSocket hub — live connections grouped into rooms, with broadcast fan-out that stays correct across replicas via the framework Coordinator.

WebSockets — the hub

github.com/redelay/go-modules/ws gives a Redelay app live, bidirectional connections. Clients open a socket, subscribe to rooms, and receive every message broadcast to a room they are in. It is the realtime counterpart to the notifications module's one-way SSE: use SSE for server→client notifications, the WebSocket hub when clients need to join and leave arbitrary channels and receive targeted fan-out.

shell
go get github.com/redelay/go-modules/ws
go
import _ "github.com/redelay/go-modules/ws"

It mounts at /ws — outside the ROUTE_PREFIX — matching how realtime endpoints are conventionally addressed.

What the hub is

The hub is the piece every realtime feature needs and none should rebuild: connection bookkeeping, room membership, and the one genuinely hard part — making a broadcast reach clients connected to a different process.

ConcernHandled by
Accept a socket, track it, tear it down cleanlythe upgrade handler + per-connection read/write pumps
Group connections into rooms; join/leaveHub.join / Hub.leave, or client subscribe frames
Deliver a message to a room's membersHub.Broadcast / Hub.BroadcastData
Deliver across replicasthe Coordinator bridge (see below)
Stop one slow client stalling the roomslow-consumer eviction

Endpoints

MethodPathAuthPurpose
GET/ws/connectquery tokenThe WebSocket upgrade
GET/ws/healthpublicLiveness + connection/room counts
GET/ws/infobearerHub summary + the caller's own connection count
POST/ws/broadcast/{session_id}bearerFan an arbitrary JSON message out to a room

Only the three control-plane endpoints appear in /openapi.json. The /ws/connect upgrade does not — a socket has no request/response body to describe.

Authenticating the upgrade

A browser's WebSocket API cannot set an Authorization header, so the upgrade reads a token from the query string:

js
const ws = new WebSocket(`wss://api.example.com/ws/connect?token=${accessToken}&session_id=room-42`)

?session_id= (or ?room=) pre-joins a room on connect, so a client that only ever watches one channel need not send a subscribe frame. Set WS_REQUIRE_AUTH=false to allow anonymous sockets (a public live feed); the control-plane POST endpoints are always authenticated regardless.

The client protocol

After connecting, a client manages its own subscriptions by sending JSON frames. It cannot broadcast — that goes through the authenticated POST endpoint, so a connected browser cannot push to other users just by holding a socket open.

json
{ "action": "subscribe",   "room": "room-42" }
{ "action": "unsubscribe", "room": "room-42" }
{ "action": "ping" }

Every message the server sends is one envelope:

json
{ "type": "biometric", "room": "room-42", "data": { "hr": 145 }, "ts": "2026-07-21T..." }

type names the kind of update so the client can route on it without parsing data. Control acks (subscribed, unsubscribed, pong) use the same shape.

Broadcasting from Go

Any module can reach the hub and push a message:

go
hub.BroadcastData(ctx, "room-42", "biometric", map[string]any{"hr": 145})

BroadcastData marshals the value into the envelope's data; Broadcast takes a pre-built ws.Message. Both deliver to every subscriber of the room — on this replica and all others.

Cross-replica fan-out

This is the property a naïve hub silently lacks. A hub that only fanned out in-process would pass every single-process test and then drop half its messages the moment a second replica came up: a client connected to replica A never sees a broadcast that originated on replica B.

When a Coordinator is wired (COORD=nats or COORD=redis) the hub bridges broadcasts across replicas over a single pub/sub subject. Each replica tags its published copies with a per-instance origin id and ignores its own echo, so a broadcast is delivered exactly once to each local subscriber.

text
COORD unset / noop   →  in-process fan-out only  (correct for one replica)
COORD=nats|redis     →  bridged across every replica

GET /ws/health reports "distributed": true when the bridge is active. The local fan-out and the cross-replica publish are independent: a message still reaches local subscribers even if the Coordinator is momentarily unavailable.

Slow-consumer eviction

A hub cannot let one stalled client back up fan-out to everyone else. Each connection has a bounded send buffer (WS_SEND_BUFFER, default 32); a client that fills it has fallen too far behind to catch up and is closed rather than allowed to block the room. Dropping a straggler is the right trade — a live feed that waits for the slowest tab is worse than a reconnect for that tab.

Adding your own control endpoints

Domain-specific realtime endpoints attach to the same /ws mount and the same hub through the ws.RouteContributor seam — so two modules never contend for the mount, and a domain module never assembles auth middleware itself.

go
// A module that pushes workout sensor data to a session's subscribers.
func (m *Module) WSRoutes(reg *ws.RouteRegistry) {
    m.hub = reg.Hub
    reg.Authed.Post("/biometric", m.handleBiometric)   // → POST /ws/biometric
}

func (m *Module) handleBiometric(w http.ResponseWriter, r *http.Request) {
    sessionID := r.URL.Query().Get("session_id")
    var data map[string]any
    _ = httputil.ReadJSON(r, &data)
    _ = m.hub.BroadcastData(r.Context(), sessionID, "biometric", data)
    httputil.WriteJSON(w, http.StatusOK, map[string]any{"session_id": sessionID})
}

RouteRegistry hands the contributor the shared *ws.Hub plus two router groups — Authed (requires a valid active user) and Public. To have the endpoint appear in /openapi.json at its top-level path, implement DynamicRoutesProvider and set AbsolutePath: true on the descriptor (see below).

Environment variables

VariableDefaultDescription
WS_MOUNT_PATH/wsWhere the hub mounts, outside the API prefix
WS_PING_INTERVAL_SECONDS25Keepalive ping interval (stays inside typical 60s idle timeouts)
WS_WRITE_TIMEOUT_SECONDS10Per-message write timeout to one client
WS_SEND_BUFFER32Queued messages per client before a slow consumer is dropped
WS_ALLOWED_ORIGINS(empty)CORS allow-list for the handshake; empty = same-origin only
WS_REQUIRE_AUTHtrueRequire a token on the upgrade; false allows anonymous sockets

AbsolutePath — top-level OpenAPI paths

The hub relies on a small framework capability worth knowing on its own: modules.DynamicRouteDescriptor.AbsolutePath.

By default the OpenAPI spec generator prefixes every dynamic route with ROUTE_PREFIX, so a route declared at /ws/health would be advertised at /api/v1/ws/health while its handler lives at /ws/health. Setting AbsolutePath: true tells the generator the path is already complete and must be emitted verbatim. Use it for any route mounted at a top-level path outside the API prefix that still needs to appear in /openapi.json.

go
func (m *Module) DynamicRoutes() []modules.DynamicRouteDescriptor {
    return []modules.DynamicRouteDescriptor{{
        Method: "GET", Path: "/ws/health",
        Source: "ws", Summary: "WebSocket hub health",
        SuccessStatus: 200, AbsolutePath: true,
    }}
}