Scalable WebSockets
Scalable WebSocket APIs using the Celerity runtime
Spec Version: v2026-02-27-draft
Ecosystem Compatibility: v0 (Current) / v1 (Preview)
Ecosystem Compatibility
- v0 (Current): Local development environments, and Amazon ECS deployments, which is where a WebSocket API can be deployed as a cluster of nodes. A local environment usually runs a single node, which behaves the same way and simply never finds another node in its group.
- v1 (Preview): Multi-cloud deployments of WebSocket APIs as containerised applications using the Celerity runtime.
Read more about Celerity versions here.
Overview
The Celerity runtime supports deployment as a horizontally scalable cluster of WebSocket servers. The approach is to allow for multiple WebSocket servers to be deployed in a way in which it does not matter which server a client connects to. This is made possible by using a shared message broker that is used to publish messages to other nodes in the cluster. Each node will then filter messages received from other nodes based on the target connection ID for the message and the clients that are connected to the node.
The Celerity runtime uses Redis OSS1 pub/sub as the message broker for publishing messages to nodes in the cluster. This allows for a scalable and efficient way to handle WebSocket connections and messages. For improvements in reliability, this may be extended in the future to support message brokers with more robust delivery guarantees.
The runtime also uses Redis OSS1 to store mappings of connection IDs to node groups, allowing for the runtime to select the channel to publish messages to in an efficient manner.

This diagram is a high-level overview of the flow for a WebSocket API cluster that doesn't include details about acknowledgements and message deduplication.
Key Namespacing
Every key and channel described here is named under a prefix that identifies the application, in the form ${prefix}:node-groups, ${prefix}:conn:${connectionId} and so on. Nodes sharing a prefix are one cluster. Two applications sharing a Redis deployment therefore need prefixes of their own, or each will publish messages to the other's nodes for connections they have never heard of, which will not produce an error or deliver a message.
The prefix must be identical across an application's nodes, so it should be derived from something that identifies the application rather than anything belonging to a node, such as the service name.
The keys that describe group membership should additionally share a hash tag, so that a Redis Cluster deployment places them in one slot. A script that reads across slots is refused, and these keys are few and small. Connection entries must stay outside that tag, since there is one per connected client and they can be spread across the cluster.
Startup Process
When a node starts up, it will connect to the message broker and will join any node group with capacity or start a new node group. Membership is recorded rather than inferred: each group has a member set at ${prefix}:node-group-members:${id}, and the groups themselves are listed in a set at ${prefix}:node-groups. To select a group, a node reads each group's member set, discards the members that are no longer running, and joins the group with the fewest live members that still has capacity. If no group has capacity, the node starts a new one. The ID created for a new node group should be a random, unique and compact ID such as a NanoID. This will have a final form of ${prefix}:node-group:${id}.
Membership is recorded, not counted from subscribers
An earlier revision of this document selected a group by counting subscribers to its channel. A subscriber count is answered only by the Redis node that receives the question, so under cluster mode it reports one shard's view of the cluster. It also cannot distinguish a node that has died from one that was never there, because a dead node's subscription and a missing subscription look identical. Reading a member set and checking each member's liveness key answers both, and reads the same from any Redis node.
Discarding, counting and joining must happen atomically, as a single script evaluated by the shared store. Nodes start together far more often than they start alone, a rolling deployment or a scale-out event starts many at once, and if each reads a group as having room before any of them has taken it, they all join the same group. The bound is then ineffective as a dozen nodes can land in a group whose capacity is two.
The node group pattern is ${prefix}:node-group:* and the node group ack pattern is ${prefix}:node-group-ack:*.
The node group ack channel is a mirror of the node group channel that is only used for acknowledgements.
Important
The same configurable capacity for node groups needs to be shared across all nodes in the cluster, this defaults to 5 nodes.
The process of joining a node group is done by subscribing to the node group and node group ack channels and keeping the node group ID in memory to be used in connection management and message handling.
Node Liveness
A node says it is still running by writing a key at ${prefix}:node:${nodeName}, whose value is the ID of the group it belongs to and which carries an expiry. The node refreshes it several times inside that expiry. Three is a reasonable choice, so that two refreshes can be lost to a stalled task or a slow round trip before anything concludes the node has gone. The expiry defaults to 30 seconds.
A node whose key has expired is dropped from its group's member set by the next node to read it, freeing the place it held. Its value also serves routing where an acknowledgement is addressed to a node, and reading that node's key gives the group whose ack channel should carry it, which avoids having to remember where each message came from.
When a node shuts down cleanly it should remove itself from its member set and delete its own key rather than waiting out the expiry, during which its group looks fuller than it is and messages are still published to it for connections it no longer holds.
Changing Node Group
A node that is dropped from its group while it was slow to refresh takes a place again on its next attempt, in whichever group has room, which may not be the group it was in. Nothing keeps a place for it, since a place kept for a node that never returns is capacity permanently lost.
Where the group changes, the node must subscribe to the new group's channels before it stops listening to the old ones, and keep both for a grace period, defaulting to 5 seconds. A sender reads where a connection is and publishes a moment later, so there are always messages in flight addressed to the group a node has just left. Giving up the old subscription immediately would lose those messages.
Connection Management
When a client connects to a node, the node must add the connection ID to the node group that the node is a member of. This is done by writing the connection ID to group mapping to the shared store. The key should be in the form of ${prefix}:conn:${connectionId} with the value being the node group ID.
When a client disconnects from a node, the node must remove the connection ID from the node group that the node is a member of. This is done by removing the connection ID to group mapping from the shared store.
Connection entries carry the same expiry as the node liveness key, and a node refreshes its own entries alongside it, in a single batched request. A node that dies without shutting down cleanly therefore leaves neither a membership nor a set of connection entries that would continue receiving messages for clients that died with it. Refreshing writes the entry rather than only extending it, so an entry that expired while a node was too busy to say otherwise is restored.
Failing to write a connection entry should not cost the client its connection. The connection works whatever the shared store knows; what is lost is the other nodes being able to find it, which the next refresh repairs.
Publishing Messages
When a node receives a message from a WebSocket client, it will first check if the connection ID is for a connection that is connected to the node. If it is, the message will be processed by the node. If it is not, the message will be published via the message broker to a subset of nodes that are more likely to be connected to the target client.
To determine the channel to send the message to, the node will look up the node group by the connection ID. Node groups are used to prevent the need to broadcast every message to every node in the cluster. Groups are used to reduce the amount of channels that need to be managed by the message broker while still providing the benefits of reducing the number of messages that need to be processed by each node.
The node will then publish the message to the channel for the node group and listen for an acknowledgement for the message, see the next section for more details on acknowledgements.
Where no group is recorded for a connection ID, the message must be published to every group rather than dropped. A client that connected a moment ago looks exactly like a client that does not exist, and the two cannot be told apart from another node. Every group hearing a message that only one of them can act on costs a message; dropping it loses a message for a client that is there.
Acknowledgements
Every time a message is sent to a node group in the cluster, the unique identifier of the source node will need to be included to allow for listening to acknowledgements that will be published by the upstream node that is connected to the target connection ID for the message.
The sender node will be listening for acknowledgements on the ack channel for its node group. It will then filter the acknowledgements based on the source node ID of the message and update the status of the message for the provided message ID. If an acknowledgement is not received by a configurable timeout, the sender node will publish the message again to the node group. This will continue until the acknowledgement is received or a maximum number of retries is reached. If the maximum number of retries is reached, the sender node will mark the message as lost and ensure that clients that should be notified of the lost message are informed.
See the Lost Messages section for more details on how lost messages are handled.
Handling Duplicates
Each node will keep track of messages that it has received and has been able to successfully forward to the client. This is done by storing the message ID in the shared store. The key should be in the form of ${prefix}:msg:${messageId} with the value being 1. These entries must expire after a configurable timeout (defaulting to 5 minutes) to prevent the shared store from growing indefinitely.
Checking and recording a message ID must be one operation rather than a read followed by a write, which the shared store's conditional set with an expiry provides. Deciding whether a message has been seen cannot be deferred, since the node has to decide before it forwards; so a separate write buys nothing and only opens a window in which two copies both read the ID as absent. One conditional set costs the same single round trip that a read alone would, and where several messages arrive together their operations can be pipelined.
A node that recognises a message as one it has already forwarded must still acknowledge it, and only skip the delivery. Skipping the acknowledgement as well would leave the sender resending until its attempts run out, and then telling the clients waiting on it that the message was lost, when in fact it was delivered on the first attempt. Reporting a delivered message as lost is worse than delivering it twice, which is what this section exists to avoid in the first place.
The timeout should be chosen carefully based on the acknowledgement timeout and the maximum number of retries, the message processed entry TTL should be set for a value that is greater than the acknowledgement timeout multiplied by the maximum number of retries.
Using a shared store allows for other nodes to be able to detect duplicates for messages that are being resent due to an acknowledgement timeout. This will not protect against the same message being sent multiple times with a different message ID, it is the responsibility of the application layer to handle content-based deduplication.
Failed requests to shared store
If recording a message ID in the shared store fails, the node will not be able to detect duplicates for messages that are being resent due to an acknowledgement timeout. This will result in the message being sent more than once to the client.
If the check fails, the message will be forwarded, where message delivery is prioritised over deduplication.
Footnotes
Last updated on