The document channel
The built-in channel
Most apps never write a channel. yrby-rails includes Y::DocumentChannel,
much like turbo-rails includes Turbo::StreamsChannel. The browser subscribes
with the signed token that collaborative_document_tag rendered. The channel
looks up the record from that token and saves every change before confirming
it. It rejects a token that's missing, tampered with, expired, signed for a
different attribute, or points at a deleted record.
Getting started shows the setup.
By default, any valid token lets the browser subscribe. To also check the
connected user's current permissions, register a block in config.to_prepare:
# config/initializers/yrby.rb, inside Rails.application.config.to_prepare
Y::DocumentChannel.authorize_document do |record, name|
current_user.present? && record.editable_by?(current_user, attribute: name)
end
The block runs inside the channel, so current_user works. editable_by?
stands for your own permission check. yrby doesn't define it. If the block
returns false or nil, the channel rejects the subscription before sending
anything.
The rest of this page covers writing your own channel with the same concern. You'd do that for documents keyed by room with no record behind them, to use a different store, or to handle authorization yourself.
Writing your own channel
Your channel includes the Y::ActionCable concern from yrby-rails. The
concern implements the y-websocket protocol, which covers document sync and
presence, and it works on both Action Cable and AnyCable. Each document has a
key. Use whatever scheme fits your app, such as
one per record and attribute (post/42/body) or one per room.
# app/channels/document_channel.rb
class DocumentChannel < ApplicationCable::Channel
include Y::ActionCable
def subscribed = sync_subscribed(params[:id])
def receive(data) = sync_receive(data, params[:id])
private
# Rejects everyone until you replace it with your app's check.
# sync_subscribed rejects the subscription unless this returns true.
def authorized?(_document_key) = false
end
Storage defaults to the gem's Y::Document models. Declare the two hooks to
point it somewhere else:
on_load { |key| Y::Document.load_state(key) } # read the document from storage
on_change { |key, update| Y::Document.append(key, update) } # save the change, then broadcast
bin/rails generate yrby:install --channel generates this channel for you,
with both hooks written out. Without --channel, the generator adds only the
storage migration.
The two hooks
When the yrby-rails models are installed and a channel declares neither hook,
both hooks use Y::Document. To use a different store, declare both. If you
declare only one, subscribing raises an error, so reads and writes can't end
up in different stores. A subclass can override one hook and inherit the other,
as long as the parent declares both. Outside a Rails app with yrby-rails there
are no defaults.
on_load receives a key and returns a binary Y.js update, or nil for a new
document. on_change receives a key and the CRDT delta, and it saves that
delta. Both run in the channel instance through instance_exec, so they can
call params, current_user, and any other channel method.
The concern loads the document through on_load whenever a client syncs, and
saves each change through on_change before broadcasting it. It keeps no
document in memory, so AnyCable RPC workers, Puma workers, and separate dynos
can all handle the same document, as long as they share the store and the
cable adapter.
Pass the key on every action: sync_receive(data, params[:id]). AnyCable
builds a new channel instance for each RPC command, so an instance variable
set in subscribed is gone by the time receive runs.
Authorization
sync_subscribed calls authorized?(key) before it opens a stream or sends
any state. The concern's default returns false, so a channel that doesn't
define the method rejects every subscriber. The log line for each rejection
says how to fix it.
Define the method with your app's check. It runs in the channel, so the connection's identity is available:
private
def authorized?(key)
current_user&.can_edit?(key)
end
For a public document, write def authorized?(_key) = true, so anyone reading
the code can see it's public on purpose.
authorized? runs once, when the client subscribes. The subscription is
authorized until it ends. To cut off access during a session, stop the
subscription yourself. Short token lifetimes also help, because the server
checks the token again on every new subscription.
For documents that belong to a record, use Y::Collaborative, which the engine
adds to every Active Record model. The page renders a signed GlobalID for one
attribute, and the channel looks up the record from it. The browser never gets
to pick a document.
<%# the view picks the document and signs it %>
<%= tag.div data: { grant: post.collaborative_sgid(:body) } %>
# the channel looks up the record from the token
def authorized?(_key) = record.present? && record.editable_by?(current_user)
# Not memoized: AnyCable builds a new channel instance for every command.
def record = Y::Collaborative.locate(params[:grant], :body)
A token signed for :body verifies only under the "yrby/body" purpose.
locate returns nil for a tampered, expired, or wrong-attribute token.
The same idea works without records. This site's demo rooms don't exist until
someone opens one, so there's no record to sign. Instead, the page signs the
room name with Rails.application.message_verifier, and authorized? only
accepts the key that token decodes to. This works well when subscribers are
anonymous.
Save before broadcast
The concern calls on_change with every change before broadcasting it. Your
block should save it durably.
on_change do |key, update|
# Write synchronously and durably. `update` is the CRDT delta.
AuditLog.append!(key, update) # raising here rejects the change
end
If the block raises, the server rejects the change and doesn't send it to anyone. The cost is one synchronous write per change. The gem doesn't lock documents, so concurrent writes can save the same update twice. That's harmless, because applying a CRDT update twice has no effect.
Delivery guarantees
These guarantees hold whether you run one process or hundreds across many servers.
- Every copy of the document ends up the same. Updates can arrive out of order, twice, or at the same time, and the result doesn't change.
- Once the server confirms an update, it's saved, even if it arrived before an update it depends on. That early update waits in the document. The client that sent the missing one hasn't had it confirmed, so it keeps resending it until the server saves it.
on_changeruns at least once for every update, before the server confirms or broadcasts it. Replaying what it saved rebuilds the document. If you need exactly-once behavior, makeon_changeidempotent. The CRDT handles duplicates either way.- If
on_changeraises, the server drops the update without replying. The client keeps it and resends it on a timer and on reconnect. That's what you want when the store was down for a moment. But if the block raises every time for the same edit, the client retries forever, because nothing tells it to stop. Put permanent rejections inauthorized?, which runs when the client subscribes. - Messages larger than
max_frame_bytes(8 MiB by default) are dropped the same way, before the server parses them. This limits how much work one client can cause. Typing never gets close, but a big paste, an embedded image, or a large first sync can. A real edit over the limit gets retried forever, like one that makeson_changeraise. The server logs each drop with the document key and update id. To add a user or connection id to that line, overridesync_log_contexton the channel.
Frame validation
Before the server uses or forwards a message, it checks that it's one complete, well-formed protocol message. It drops anything malformed, truncated, oversized, of an unknown type, or holding more than one message. If the Rust code panics, Ruby raises an exception and the process keeps running. So one client can't send a message that breaks the room for everyone else.
Validation doesn't limit how many messages a client sends. A client sending
valid messages as fast as it can still needs a rate limit. This site's
channels put token buckets in front of sync_receive for that, and the site's
README describes
them.
Reliable delivery (acks)
The server acknowledges document updates. Each browser update includes an
"id", and the server replies { "ack": <id> } after on_change returns.
client -> server { "update": "<base64 update>", "id": 42 }
server -> client { "ack": 42 } # saved; the client can drop update 42
yrby-client's ActionCableProvider does this for you. It queues local
updates until they're confirmed and sends them merged into one update. The id
is the highest sequence number in the batch, so one ack confirms all of them.
If the server gets an update it already has, applying it changes nothing, and
the server acks it again. Presence updates aren't acked.
Causal gaps
An update can reach the server before the update it depends on. That's normal, and the server saves and confirms it like any other. Yjs holds it in the document as a pending struct and applies it once the missing update arrives. It costs the same as any other update, because the server appends it, forwards it, and acks it without rebuilding the document.
Other clients get pending structs too, because handle_sync_message answers
with the full state. They hold the same pending struct and apply it the same
way. Nothing special has to happen to close the gap. The client that sent the
missing update hasn't had it confirmed, so it keeps resending it. Compaction
is the exception. It leaves pending structs out of the snapshot, because a
pending struct folded into the snapshot could never be applied.
Gaps are easy to miss, because the pending edit doesn't show up in the
document until its dependency arrives. A gap is worth an alert when no connected
client can fill it. Use the on_gap hook for that. Whenever the server
loads a document to send its state and finds a gap, it calls on_gap with the
document key.
class DocumentChannel < ApplicationCable::Channel
include Y::ActionCable
on_gap { |key| StatsD.increment("yrby.gap", tags: ["doc:#{key}"]) }
end
The server also logs open gaps at info. It catches and logs errors raised in
the hook, so a broken metrics call can't break frame handling.
Sync flow
Browser Server
| |
|--- subscribe -------------------------->| authorized?
|<-- SyncStep1 (server's state vector) ---| read through on_load
|--- SyncStep2 (what the server lacks) -->| saved through on_change
| |
|--- SyncStep1 + awareness -------------->|
|<-- SyncStep2 (full state) --------------| the browser is synced
| |
|--- update { id } ---------------------->| saved, then broadcast
|<-- { ack: id } -------------------------|
|<-- updates from other clients ----------|
Message type constants
Y::MSG_SYNC # 0 - document sync message
Y::MSG_AWARENESS # 1 - presence update
Y::MSG_SYNC_STEP1 # 0 - state vector request
Y::MSG_SYNC_STEP2 # 1 - update response
Y::MSG_SYNC_UPDATE # 2 - incremental update
Protocol codec
These are module functions on Y, because classifying and unwrapping a
message doesn't need any state.
Y.message_kind(frame) # => 0 drop / 1 step1 / 2 update / 3 awareness / 4 query
Y.update_from_message(frame) # => the document delta in a frame, or nil
Y.wrap_update(update_bytes) # => a raw document update wrapped as a sync Update frame
This page is copied from the repo README. If the two disagree, go by the README.