Plugins
Everything behind kraken's wire protocol is a plugin slot. Each slot is an Erlang behaviour with built-in implementations, and you choose one per slot by short name, or by naming any module that implements the behaviour.
The seven slots
| Slot | Set with | Built-in short names | Default | Behaviour |
|---|---|---|---|---|
| Auth | AUTH_BACKEND | static, http | static | kraken_auth |
| Broker | BROKER_BACKEND | syn, mqtt | syn | kraken_broker |
| Store | STORE_BACKEND | ets, noop | ets | kraken_store |
| Control | CONTROL_BACKEND | noop, http | noop | kraken_control |
| Presence store | presence_store_backend in sys.config | noop | noop | kraken_presence_store |
| Wake | wake_backend in sys.config | noop | noop | kraken_wake |
| Delivery store | delivery_store_backend in sys.config, plus durable_delivery set to true | noop, ets | noop | kraken_delivery_store |
The first four can be set from the environment; the last three only in sys.config (see Configuration). A value that is not one of the slot's short names is used as a module name, so AUTH_BACKEND=my_auth makes kraken call my_auth:validate_token/1.
kraken also contains kraken_presence_store_ets and kraken_wake_http, which you select by module name. Their source marks both as development and test backends, not for production.
Writing a backend
Auth
-callback validate_token(AccessToken :: binary()) ->
{ok, AuthData :: map()} | {error, Reason :: binary()}.
-callback revalidate_token(ActorTokenId :: binary()) ->
{ok, AuthData :: map()} | {error, Reason :: binary()} | {retry, Reason :: binary()}.
%% optional
-callback check_room_access(ActorTokenId :: binary(), Pattern :: binary()) ->
{ok, AllowedTopics :: list()} | {error, Reason :: binary()}.
Build AuthData with kraken_auth:build_auth_data/1 from a map with binary keys in the same shape as the HTTP backend's client_attrs; see HTTP Auth Contract. An {error, Reason} from validate_token is the reason the client sees. From revalidate_token, {error, Reason} closes the connection with code 4001 and {retry, Reason} tries again at the next heartbeat. kraken's dispatcher keeps the 30 second cache of successful validations for every backend.
Broker
-callback start() -> ok.
-callback connect() -> {ok, Session} | {error, term()}.
-callback connect(AuthData) -> {ok, Session} | {error, term()}.
-callback connect(AuthData, PersistentSession, SessionExpirySeconds) ->
{ok, Session, ClientId} | {error, term()}.
-callback subscribe(Session, MqttTopic, DisplayTopic, WsPid, QoS) -> ok.
-callback unsubscribe(Session, Topic) -> ok.
-callback publish(Session, Topic, Data, Sender | undefined, QoS, Retain) -> ok.
-callback disconnect(Session) -> ok.
-callback format_shared_subscription(BaseTopic, Group) -> binary().
-callback supports_load_balancing() -> boolean().
-callback capabilities() -> map(). %% #{retained, shared_subscriptions, multi_region}
Subscribed connection processes must receive {mqtt_publish, #{topic := Topic, payload := PackedPayload}}, where the payload is MessagePack. When Sender is set, wrap the data as #{<<"data">> => Data, <<"_sender">> => Sender} so the sending connection can drop its own copy. A shared subscription, $share/Group/Topic, must deliver each message to exactly one member of the group.
Store
-callback init() -> ok | {error, term()}.
-callback is_enabled() -> boolean().
-callback log_message(Doc) -> ok | {error, term()}.
-callback log_delivery(Doc) -> ok | {error, term()}.
-callback mark_delivered(MessageId, ActorId, Timestamp) -> ok.
-callback log_event(Doc) -> ok.
-callback ack_delivery(MessageId, ActorId) -> ok | {error, term()}.
-callback batch_ack_deliveries(ActorId, MessageIds) -> ok | {error, term()}.
-callback get_replay_messages(Options) -> {ok, [Msg], Count} | {error, term()}.
-callback get_undelivered_count(ActorId, AppId) -> {ok, Count} | {error, term()}.
-callback terminate() -> ok.
Recording is switched on and off globally by RECORD_MESSAGES, whichever store is configured. Nothing in v0.9.0 reads recorded messages back for clients, and nothing calls get_replay_messages or get_undelivered_count; the only code that reads recorded messages is the ets delivery store's replay for load-balanced subscriptions.
Control
-callback report_usage(Entries) -> ok | {ok, BlockedProjectIds} | {error, term()}.
-callback report_subscription_change(Changes) -> ok | {error, term()}.
-callback report_webhook_failure(Failure) -> ok.
kraken batches the reports itself: usage every 30 seconds, or sooner once a project reaches 100 messages; subscription changes every 500 milliseconds or 50 changes. Returning {ok, BlockedProjectIds} from report_usage replaces the set of blocked projects, and publishes from a blocked project are refused with error 42920, monthly_quota_exceeded.
The built-in http control backend posts these to CONTROL_HTTP_URL, with Authorization: Bearer <BACKEND_SECRET>:
| Request | Body |
|---|---|
POST {CONTROL_HTTP_URL}/usage | { "entries": [{ "projectId", "count", "totalBytes" }] }. Answer 200, optionally with { "blockedProjects": ["..."] }. |
POST {CONTROL_HTTP_URL}/subscriptions | { "changes": [{ "actorTokenId", "topic", "action", "loadBalance", "loadBalanceGroup", "filters" }] }, where action is subscribe or unsubscribe and the last three appear only when set |
POST {CONTROL_HTTP_URL}/webhook-failures | One failed webhook call: organizationId, projectId, appId, roomId, type, webhookUrl, requestHeaders, requestBody, responseStatus, errorMessage, node, timestamp |
The /subscriptions reports are what an auth service needs to restore subscriptions on reconnect.
Presence store, wake and delivery store
These three slots carry features that are off by default:
- Presence store (
upsert/1,offline/1,mark_waking/1,discover/1) keeps a presence record after its connection closes, for clients that advertisepersistentpresence. With the defaultnoop, presence lasts exactly as long as the connection. - Wake (
fire/1) would call out to bring an offline actor back when a message is routed to it. The default does nothing. - Delivery store (
pending/1,claim/1,cursor_get/1,cursor_set/2,ack/1,is_enabled/0) lets a load-balanced worker that reconnects replay, and claim, the messages its group missed while every member was offline. It takes effect only withdurable_deliveryset totrue. The built-inetsbackend replays from theetsmessage store, so it also needsRECORD_MESSAGESon.
The behaviour modules in the kraken source (src/backends/kraken_presence_store.erl, kraken_wake.erl, kraken_delivery_store.erl) document the maps each callback receives.
Embedding kraken
kraken can be a rebar3 dependency of your own Erlang application, which then ships its backend modules alongside kraken and sets the slots in its own sys.config:
%% rebar.config
{deps, [
{kraken, {git, "https://github.com/NoLagApp/kraken.git", {tag, "v0.9.0"}}}
]}.
%% sys.config
[{kraken, [
{auth_backend, my_auth},
{broker_backend, syn},
{store_backend, noop},
{control_backend, noop}
]}].
We compiled an application this way, with kraken v0.9.0 as a dependency and an auth backend of its own, using the same Erlang image kraken's Dockerfile builds with. A release built this way also needs the rest of kraken's settings; copy them from kraken's config/sys.config.src. kraken:stats/0 returns the node name, the cluster's nodes, the connection count and the configured backends, for an embedding application that reports on its own health.