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

SlotSet withBuilt-in short namesDefaultBehaviour
AuthAUTH_BACKENDstatic, httpstatickraken_auth
BrokerBROKER_BACKENDsyn, mqttsynkraken_broker
StoreSTORE_BACKENDets, noopetskraken_store
ControlCONTROL_BACKENDnoop, httpnoopkraken_control
Presence storepresence_store_backend in sys.confignoopnoopkraken_presence_store
Wakewake_backend in sys.confignoopnoopkraken_wake
Delivery storedelivery_store_backend in sys.config, plus durable_delivery set to truenoop, etsnoopkraken_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>:

RequestBody
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-failuresOne 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 advertise persistent presence. With the default noop, 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 with durable_delivery set to true. The built-in ets backend replays from the ets message store, so it also needs RECORD_MESSAGES on.

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.