Using WebSockets with platformOS
In this guide we will show you how to integrate WebSockets into your application based on the chat functionality we implemented for the pOS Marketplace Template.
WebSocket is a protocol for real-time communication between a client and a server. It allows to build an interactive UI where data can be updated from the server without the need of doing http requests. Each client can subscribe to multiple channels and rooms. Messages sent to a certain room are then broadcasted to all subscribers of that room. The pOS implementation of Websockets uses ActionCable.
A client that only needs to receive broadcasts - a notification feed, a live counter, a dashboard - can subscribe over Server-Sent Events instead, with no WebSocket library at all. See Following a room with Server-Sent Events (SSE) below.
Requirements
To understand this topic, you should be familiar with:
Steps
Step 1: Setup
Make sure you have added the @rails/actioncable package "@rails/actioncable": "^6.0.3-2" to your package.json file. This package handles communication between the server and the client.
Step 2: Create a connection
Consumers require an instance of the connection on their side. This can be established using the code below.
// modules/chat/src/js/consumer.js
import { createConsumer } from "@rails/actioncable"
export default createConsumer('/websocket')
Step 3: Subscribe to a channel
A consumer becomes a subscriber by creating a subscription to a given channel:
// modules/chat/src/js/chat.js
import consumer from "./consumer";
const chat = function(){
// cache 'this' value not to be overwritten later
const module = this;
module.createSubscription = () => {
module.channel = consumer.subscriptions.create(
{
channel: 'conversate', //channel name
room_id: module.conversationId, //room ID
}
)
}
}
Step 4: Handle callbacks
There are 4 callbacks handled by the subscription that allow you to react on certain events. You can add handlers for them when you create the subscription object:
// modules/chat/src/js/chat.js
import consumer from "./consumer";
const chat = function(){
// cache 'this' value not to be overwritten later
const module = this;
module.createSubscription = () => {
module.channel = consumer.subscriptions.create(
{
channel: 'conversate',
room_id: module.conversationId
},
{
received: function(data){
// function responsible for handling new message
module.showMessage(data);
},
connected: function(){
// on connect we want to enable the message input
module.settings.messageInput.disabled = false;
module.settings.messageInput.focus();
},
rejected: function(){
// when connection has been rejected we want to notify the user about the error
module.blocked();
},
disconnected: function(){
// when consumer has been disconnected we want to notify the user about the error
module.blocked();
}
}
);
};
}
Handling rejected subscriptions
The rejected() callback fires when the server refuses a subscription, but ActionCable's reject_subscription frame cannot carry a reason - so rejected() receives no arguments. To find out why a subscription was refused, platformOS sends a structured subscription_error message immediately before the rejection. Read it in your received(data) handler (correlated automatically to this subscription) and branch on its type so it isn't mistaken for an application message:
received: function(data) {
if (data.type === 'subscription_error') {
// { type: 'subscription_error', code, message, retryable }
console.warn(`Subscription refused (${data.code}): ${data.message}`);
if (data.retryable) {
// a transient server error - resubscribing later may succeed
}
return;
}
// otherwise it's an application message
module.showMessage(data);
}
The code field is a stable, machine-readable value you can switch on:
code |
Meaning | retryable |
|---|---|---|
unauthorized |
The subscribed partial did not return true, or the channel has no subscribed partial and the connection is not allowed the open default (cross-origin, or websockets_require_subscribed_partial is enabled). |
false |
instance_not_found |
The connection could not be matched to a platformOS instance. | false |
subscribed_partial_error |
Your subscribed partial raised a Liquid error. The detail is written to the marketplace error log under the WebSocketSubscribeError code (visible in the admin), never sent over the socket. |
false |
internal |
An unexpected server error. Retrying may help. | true |
Messages sent to the client are always generic and never leak internal detail. If you receive neither a successful connection nor a subscription_error (only silence), treat it as a transport problem (dropped socket) and reconnect - not as a permission denial.
Step 5: Secure room subscription
To make sure that only authorized users can access certain rooms (for example, rooms for conversations where they are one of the pariticpants) you can use a partial that will be called when the user tries to subscribe to a certain room.
This file should be stored in views/partials/channels/:channel_name/subscribed.liquid. In this example, we use code from the pOS Marketplace Template and the channel is called conversate. This partial needs to return a true string if the current user is authorized to subscribe to the room.
{%- comment -%}
modules/chat/public/views/partials/channels/conversate/subscribed.liquid
{%- endcomment -%}
{% liquid
function current_profile = 'lib/current_organization_profile', user_id: context.current_user.id
assign room_id = context.params.room_id
function conversation = 'modules/chat/lib/queries/conversations/find_by_participant', id: room_id, participant_id: current_profile.id
if conversation
echo 'true'
else
echo 'false'
endif
%}
In the pOS Marketplace Template, each user has a profile and the profile is linked with a conversation. If the conversation is found then we return true to confirm that the user can subscribe to this room.
Because the websockets_require_subscribed_partial flag in app/config.yml defaults to true, a channel without a channels/:channel_name/subscribed.liquid partial rejects every subscription with an unauthorized verdict - so adding this partial is how you authorize subscribers. You can restore the legacy open-by-default behaviour (any same-origin subscriber is admitted) by setting websockets_require_subscribed_partial: false, but this is not recommended. Regardless of the flag, a cross-origin connection is always rejected when the channel has no subscribed partial.
Step 6: Send messages
To send a message, you need to call the send method on the subscription object you have created. This method takes Object as an argument.
// modules/chat/src/js/chat.js
module.sendMessage = (message) => {
let messageData = {
message: encodeHtml(message),
author_id: module.settings.currentUserId,
sender_name: module.settings.messageInput.getAttribute('data-from-name'),
created_at: new Date(),
create: true
};
module.channel.send(messageData);
};
Step 7: Set server callback on received message
You may want to trigger some code when the server receives a message from a client. To do so, you need to create a partial that will be called when the message is received. This file should be stored in views/partials/channels/:channel_name/receive.liquid. In this example, we use code from the pOS Marketplace Template and the channel is called conversate.
context.params contains all fields that have been sent in the previous step.
{%- comment -%}
modules/chat/public/views/partials/channels/conversate/receive.liquid
{%- endcomment -%}
{% liquid
function current_profile = 'lib/current_organization_profile', user_id: context.current_user.id
assign room_id = context.params.room_id
function conversation = 'modules/chat/lib/queries/conversations/find_by_participant', id: room_id, participant_id: current_profile.id
if conversation
assign message_safe = context.params.message | raw_escape_string
assign object = { "conversation_id": conversation.id, "autor_id": current_profile.id, "message": message_safe }
function message = 'modules/chat/lib/commands/messages/create', object: object
if message.valid != true
log message, 'ERROR receive message'
echo "false"
else
function res = 'modules/chat/lib/commands/conversations/mark_unread', conversation: conversation, current_profile: current_profile
endif
else
echo "false"
endif
%}
In this example, we are retrieving the conversation object and creating a message object to persist data in the database.
To avoid broadcasting (for example if the user is not authorized to add a message in a conversation or when an validation error is triggered) ensure the partial returns "false" (whitespace characters do not matter).
Employing Cross-Site Request Forgery token validation
If you've enabled websockets_require_csrf_token (default) in your app/config.yml, the CSRF token must be passed as below:
const token = csrfToken();
createConsumer(`/websocket?authenticity_token=${token}`)
The csrfToken function can retrieve the token from the meta tag for example:
const csrfToken = () => {
const meta = document.querySelector('meta[name="csrf-token"]');
if (!meta) { throw new Error('Unable to find CSRF token meta'); }
const token = meta.getAttribute('content');
if (!token) { throw new Error('Unable to get CSRF token value'); }
return token;
}
Following a room with Server-Sent Events (SSE)
A browser does not have to open a WebSocket to follow a room. The same channel is also reachable over Server-Sent Events at /anycable-events, through the browser's built-in EventSource - no @rails/actioncable package, no consumer, no build step.
Nothing changes on the platformOS side. An SSE subscriber goes through the same channels/conversate/subscribed.liquid partial from Step 5: Secure room subscription, obeys the same websockets_require_subscribed_partial and websockets_require_csrf_token flags, and receives the same broadcasts - including the ones sent from the backend with the channel_send_message mutation described in the next section.
The trade-off is that SSE is server-to-client only: it can read a room, it cannot write to one. Use it for notification feeds, live counters, dashboards and any other broadcast-only stream, and keep the WebSocket transport for chat, where the client also sends.
// modules/chat/src/js/notifications-stream.js
const notificationsStream = function(){
// cache 'this' value not to be overwritten later
const module = this;
module.createStream = () => {
// The *whole* identifier, as URL-encoded JSON - see the warning below.
const identifier = JSON.stringify({
channel: 'conversate', // channel name
room_id: module.conversationId // room ID
});
const url = `/anycable-events?identifier=${encodeURIComponent(identifier)}`
+ `&authenticity_token=${encodeURIComponent(csrfToken())}`;
module.source = new EventSource(url);
module.source.addEventListener('confirm_subscription', () => {
// the room is open - the equivalent of the connected() callback
module.settings.messageList.classList.remove('is-loading');
});
module.source.onmessage = (event) => {
// event.data is the broadcast payload itself
module.showMessage(JSON.parse(event.data));
};
module.source.onerror = () => {
if (module.source.readyState === EventSource.CLOSED) {
// the browser has given up - the subscription was refused, or the token
// is no longer valid. Retrying this URL cannot help, see below.
module.blocked();
}
// otherwise readyState is CONNECTING and the browser is already retrying
};
};
module.closeStream = () => module.source.close();
}
csrfToken() is the helper from the previous section, and showMessage / blocked are the same handlers the WebSocket subscription uses in Step 4.
Always send the full identifier
The endpoint also accepts a ?channel=conversate shorthand, but that form carries only the channel name. ?channel=conversate&room_id=12 reaches the server as {"channel":"conversate"} - room_id is dropped on the way - so context.params.room_id in your subscribed partial is blank and every client using the shorthand ends up in one shared room instead of its own. Always build the full identifier JSON, as in the example above.
Callback equivalents
subscriptions.create callback |
EventSource |
|---|---|
connected() |
the confirm_subscription event (the stream opens with a welcome event first) |
received(data) |
onmessage, where event.data is the payload itself - there is no Action Cable envelope to unwrap |
rejected() |
onerror with readyState === EventSource.CLOSED - the server answers the request with 401 |
disconnected() |
onerror with readyState === EventSource.CONNECTING - the browser is already reconnecting on its own |
Do not count on receiving the structured subscription_error message described in Handling rejected subscriptions here: on this transport the refusal arrives as the HTTP 401 above, which is why the handler branches on readyState rather than on a payload. The reason is still recorded in your Instance logs exactly as before - a subscribed partial that raises is logged under the WebSocketSubscribeError code.
The CSRF token, and why a reconnect needs a fresh one
EventSource cannot set request headers, so when websockets_require_csrf_token is enabled the token goes in the query string, exactly as it does in the WebSocket URL. Two things follow from that:
- An anonymous visitor - the usual audience for a read-only stream - is never given the
CSRF-TOKENcookie, so the token has to be templated into the page for thecsrfToken()helper to find. Render it into the meta tag with{{ context.authenticity_token }}. - The browser's own reconnect replays the URL the
EventSourcewas constructed with, so it can never pick up a newer token. Once the session behind that token is gone - a logout, an expiry - the stream is refused for good andreadyStatestaysCLOSED. Recovering means constructing a newEventSourcewith a freshly read token, not waiting for the built-in retry.
What SSE cannot do
Step 6: Send messages and Step 7: Set server callback on received message have no SSE equivalent: an EventSource has no send(), and channels/conversate/receive.liquid is never invoked for an SSE client. A channel that defines a receive partial is still perfectly subscribable this way - a conversation with two-way WebSocket chat can be read over SSE in parallel - it is only the sending half that is missing.
A note on Origin
Browsers send no Origin header at all on a same-origin EventSource request - unlike a WebSocket handshake, which always carries one - and add it only when the request is cross-origin. platformOS therefore reads a missing Origin on this transport as same-origin, and evaluates a cross-origin EventSource exactly as it evaluates a WebSocket. This only matters for a channel with no subscribed partial running with websockets_require_subscribed_partial: false, where being same-origin is what admits the subscriber; with a subscribed partial in place - the recommended setup - your own authorization decides, on every transport.
GraphQL mutation to broadcast to a channel and a room
You can broadcast a message from the backend, to a specific channel and room, using the mutation channel_send_message.
The mutation takes the channel_name, the room_id, and the JSON payload as mandatory arguments to send.
For example:
mutation broadcast_message {
channel_send_message(
channel_name: "conversate"
room_id: "12"
payload: {
message: "Some message",
author_id: 2,
sender_name: "Some Sender",
created_at: "2030-07-01 10:30:00",
create: true
}
)
}
Other examples
In the pOS Marketplace Template, we have also implemented a notifications functionality that uses WebSockets to inform user about unread messages. You can find it in the file modules/chat/src/js/notifications.js.
Q&A
- Socket timeouts — how does it handle reconnect? Is that the job of the client or can the server reconnect in a permanent listen mode?
The JS client @rails/actioncable is responsible for reconnecting.
- Can the websocket service connect to anything else (a notification service to receive events)?
Yes, if that service handles websocket connections.
- Timeouts
There is 5 seconds timeout for handling subscribed.liquid and receive.liquid partials.
- How do you define a message?
Message is a JSON object.
- Parallel processing, assuming just one connection, but how does the server keep up if there are streams of events one after another,
does it split off thread / processes in the background to cover this or queue this somehow?
Each client has their connection and each connection is handled by a thread on the server side.
If there are multiple incoming messages from a single connection then they are queued.