Packages
Elixir client for SpacetimeDB — BSATN binary protocol, WebSocket subscriptions, reducer calls, live ETS table mirrors
Current section
Files
Jump to
Current section
Files
lib/spacetimedb/protocol/bsatn.ex
defmodule SpacetimeDB.Protocol.BSATN do
@moduledoc """
Encodes and decodes SpacetimeDB WebSocket messages using the
`v1.bsatn.spacetimedb` binary subprotocol.
## Wire format
Every WebSocket binary frame carries exactly one message, encoded as a
BSATN sum type:
<<tag::u8, ...fields>>
### Client → server message tags
| Tag | Message |
|-----|---------|
| 0 | `CallReducer` |
| 1 | `Subscribe` |
| 2 | `OneOffQuery` |
| 3 | `SubscribeSingle` |
| 4 | `Unsubscribe` |
| 5 | `SubscribeMulti` |
### Server → client message tags
| Tag | Message |
|-----|---------|
| 0 | `InitialSubscription` |
| 1 | `TransactionUpdate` |
| 2 | `TransactionUpdateLight` |
| 3 | `IdentityToken` |
| 4 | `OneOffQueryResponse` |
| 5 | `SubscribeApplied` |
| 6 | `UnsubscribeApplied` |
| 7 | `SubscriptionError` |
| 8 | `SubscribeMultiApplied` |
Row data within `TableUpdate` messages is raw BSATN — each row is encoded
as the product type matching that table's schema. Row bytes are surfaced as
opaque `binary()` values in `SpacetimeDB.Types.TableUpdate` fields; use
`SpacetimeDB.BSATN.Schema` to define your table schemas and decode them.
"""
alias SpacetimeDB.{BSATN, Types}
@subprotocol "v1.bsatn.spacetimedb"
@doc "The WebSocket subprotocol header value."
def subprotocol, do: @subprotocol
# ---------------------------------------------------------------------------
# Client → Server encoding
# ---------------------------------------------------------------------------
@doc "Encode a `CallReducer` message (tag 0)."
def encode_call_reducer(reducer, args_bsatn, request_id, flags \\ 0)
when is_binary(args_bsatn) do
<<0::8,
BSATN.encode_string(reducer)::binary,
BSATN.encode_bytes(args_bsatn)::binary,
BSATN.encode_u32(request_id)::binary,
BSATN.encode_u8(flags)::binary>>
end
@doc "Encode a `Subscribe` message (tag 1)."
def encode_subscribe(query_strings, request_id) when is_list(query_strings) do
queries = BSATN.encode_array(query_strings, &BSATN.encode_string/1)
<<1::8, queries::binary, BSATN.encode_u32(request_id)::binary>>
end
@doc "Encode a `OneOffQuery` message (tag 2)."
def encode_one_off_query(query_string, message_id) when is_binary(message_id) do
<<2::8,
BSATN.encode_bytes(message_id)::binary,
BSATN.encode_string(query_string)::binary>>
end
@doc "Encode a `SubscribeSingle` message (tag 3)."
def encode_subscribe_single(query, request_id, query_id) do
<<3::8,
BSATN.encode_string(query)::binary,
BSATN.encode_u32(request_id)::binary,
BSATN.encode_u64(query_id)::binary>>
end
@doc "Encode an `Unsubscribe` message (tag 4)."
def encode_unsubscribe(request_id, query_id) do
<<4::8,
BSATN.encode_u32(request_id)::binary,
BSATN.encode_u64(query_id)::binary>>
end
@doc "Encode a `SubscribeMulti` message (tag 5)."
def encode_subscribe_multi(query_strings, request_id, query_id) when is_list(query_strings) do
queries = BSATN.encode_array(query_strings, &BSATN.encode_string/1)
<<5::8,
queries::binary,
BSATN.encode_u32(request_id)::binary,
BSATN.encode_u64(query_id)::binary>>
end
# ---------------------------------------------------------------------------
# Server → Client decoding
# ---------------------------------------------------------------------------
@doc "Decode a BSATN binary frame from the server into a typed struct."
@spec decode(binary()) :: {:ok, term()} | {:error, term()}
def decode(<<tag::8, rest::binary>>), do: decode_tag(tag, rest)
def decode(_), do: {:error, :empty_frame}
# tag 0: InitialSubscription
defp decode_tag(0, bin) do
with {:ok, tables, bin} <- decode_table_updates(bin),
{:ok, request_id, bin} <- BSATN.decode_u32(bin),
{:ok, exec_time, _rest} <- BSATN.decode_u64(bin) do
{:ok,
%Types.InitialSubscription{
request_id: request_id,
tables: tables,
execution_time_micros: exec_time
}}
end
end
# tag 1: TransactionUpdate
defp decode_tag(1, bin) do
with {:ok, status, tables, bin} <- decode_update_status(bin),
{:ok, timestamp_us, bin} <- BSATN.decode_u64(bin),
{:ok, caller_identity, bin} <- BSATN.decode_bytes(bin),
{:ok, caller_conn_id, bin} <- BSATN.decode_bytes(bin),
{:ok, reducer_call, bin} <- decode_reducer_call(bin),
{:ok, energy, _rest} <- BSATN.decode_u128(bin) do
{:ok,
%Types.TransactionUpdate{
status: status,
tables: tables,
timestamp: %Types.Timestamp{microseconds_since_epoch: timestamp_us},
caller_identity: Base.encode16(caller_identity, case: :lower),
caller_connection_id: Base.encode16(caller_conn_id, case: :lower),
reducer_call: reducer_call,
energy_consumed: energy
}}
end
end
# tag 2: TransactionUpdateLight
defp decode_tag(2, bin) do
with {:ok, tables, _rest} <- decode_table_updates(bin) do
{:ok, %Types.TransactionUpdate{status: :committed, tables: tables}}
end
end
# tag 3: IdentityToken
defp decode_tag(3, bin) do
with {:ok, identity, bin} <- BSATN.decode_bytes(bin),
{:ok, token, bin} <- BSATN.decode_string(bin),
{:ok, conn_id, _rest} <- BSATN.decode_bytes(bin) do
{:ok,
%Types.IdentityToken{
identity: Base.encode16(identity, case: :lower),
token: token,
connection_id: Base.encode16(conn_id, case: :lower)
}}
end
end
# tag 4: OneOffQueryResponse
defp decode_tag(4, bin) do
with {:ok, message_id, bin} <- BSATN.decode_bytes(bin),
{:ok, error, bin} <- BSATN.decode_option(bin, &BSATN.decode_string/1),
{:ok, tables, bin} <- decode_rows_in_tables(bin),
{:ok, exec_time, _rest} <- BSATN.decode_u64(bin) do
{:ok,
%Types.OneOffQueryResponse{
message_id: message_id,
error: error,
tables: tables,
execution_time_micros: exec_time
}}
end
end
# tag 5: SubscribeApplied
defp decode_tag(5, bin) do
with {:ok, request_id, bin} <- BSATN.decode_u32(bin),
{:ok, query_id, bin} <- BSATN.decode_u64(bin),
{:ok, tables, _rest} <- decode_rows_in_tables(bin) do
{:ok, %Types.SubscribeApplied{request_id: request_id, query_id: query_id, tables: tables}}
end
end
# tag 6: UnsubscribeApplied
defp decode_tag(6, bin) do
with {:ok, request_id, bin} <- BSATN.decode_u32(bin),
{:ok, query_id, bin} <- BSATN.decode_u64(bin),
{:ok, tables, _rest} <- decode_rows_in_tables(bin) do
{:ok,
%Types.UnsubscribeApplied{request_id: request_id, query_id: query_id, tables: tables}}
end
end
# tag 7: SubscriptionError
defp decode_tag(7, bin) do
with {:ok, request_id, bin} <- BSATN.decode_u32(bin),
{:ok, query_id, bin} <- BSATN.decode_u64(bin),
{:ok, error, _rest} <- BSATN.decode_string(bin) do
{:ok,
%Types.SubscriptionError{request_id: request_id, query_id: query_id, error: error}}
end
end
# tag 8: SubscribeMultiApplied — same shape as SubscribeApplied
defp decode_tag(8, bin), do: decode_tag(5, bin)
defp decode_tag(tag, _bin), do: {:ok, {:unknown_tag, tag}}
# ---------------------------------------------------------------------------
# Helpers — update status
# ---------------------------------------------------------------------------
# UpdateStatus is a sum type: 0=Committed, 1=Failed, 2=OutOfEnergy
defp decode_update_status(<<0::8, rest::binary>>) do
with {:ok, tables, bin} <- decode_table_updates(rest) do
{:ok, :committed, tables, bin}
end
end
defp decode_update_status(<<1::8, rest::binary>>) do
with {:ok, reason, bin} <- BSATN.decode_string(rest) do
{:ok, {:failed, reason}, [], bin}
end
end
defp decode_update_status(<<2::8, rest::binary>>), do: {:ok, :out_of_energy, [], rest}
defp decode_update_status(_), do: {:error, :invalid_update_status}
# ---------------------------------------------------------------------------
# Helpers — table updates
# ---------------------------------------------------------------------------
defp decode_table_updates(bin) do
BSATN.decode_array(bin, &decode_table_update/1)
end
defp decode_table_update(bin) do
with {:ok, table_id, bin} <- BSATN.decode_u32(bin),
{:ok, table_name, bin} <- BSATN.decode_string(bin),
{:ok, inserts, bin} <- BSATN.decode_array(bin, &BSATN.decode_bytes/1),
{:ok, deletes, rest} <- BSATN.decode_array(bin, &BSATN.decode_bytes/1) do
{:ok,
%Types.TableUpdate{
table_id: table_id,
table_name: table_name,
inserts: inserts,
deletes: deletes
}, rest}
end
end
# For SubscribeApplied / UnsubscribeApplied / OneOffQueryResponse the table
# format is slightly different (table_name + inserts only, no table_id)
defp decode_rows_in_tables(bin) do
BSATN.decode_array(bin, &decode_rows_table/1)
end
defp decode_rows_table(bin) do
with {:ok, table_name, bin} <- BSATN.decode_string(bin),
{:ok, inserts, bin} <- BSATN.decode_array(bin, &BSATN.decode_bytes/1),
{:ok, deletes, rest} <- BSATN.decode_array(bin, &BSATN.decode_bytes/1) do
{:ok,
%Types.TableUpdate{
table_name: table_name,
inserts: inserts,
deletes: deletes
}, rest}
end
end
# ---------------------------------------------------------------------------
# Helpers — reducer call
# ---------------------------------------------------------------------------
defp decode_reducer_call(<<0::8, rest::binary>>), do: {:ok, nil, rest}
defp decode_reducer_call(<<1::8, rest::binary>>) do
with {:ok, reducer_name, bin} <- BSATN.decode_string(rest),
{:ok, request_id, bin} <- BSATN.decode_u32(bin),
{:ok, status, rest} <- decode_call_status(bin) do
{:ok,
%Types.ReducerCallInfo{
reducer_name: reducer_name,
request_id: request_id,
status: status
}, rest}
end
end
defp decode_reducer_call(_), do: {:error, :invalid_reducer_call}
defp decode_call_status(<<0::8, rest::binary>>), do: {:ok, :committed, rest}
defp decode_call_status(<<1::8, rest::binary>>) do
with {:ok, msg, bin} <- BSATN.decode_string(rest), do: {:ok, {:failed, msg}, bin}
end
defp decode_call_status(<<2::8, rest::binary>>), do: {:ok, :out_of_energy, rest}
defp decode_call_status(_), do: {:error, :invalid_call_status}
end