Packages
reckon_db
1.4.4
5.11.0
5.10.4
5.10.3
5.10.1
5.10.0
5.9.1
5.9.0
5.8.3
5.8.2
5.8.1
5.8.0
5.7.0
5.6.1
5.6.0
5.5.5
5.5.4
5.5.3
5.5.2
5.5.1
5.5.0
5.4.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.0.0
3.1.2
3.1.1
3.0.0
2.3.7
2.3.6
2.3.5
2.3.4
2.3.3
2.3.2
2.3.1
2.3.0
2.2.2
2.2.0
2.1.4
2.1.3
2.1.2
2.1.1
2.1.0
2.0.0
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.3
1.3.2
1.3.1
1.3.0
1.2.7
1.2.6
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.1
1.1.0
1.0.3
1.0.2
1.0.1
1.0.0
BEAM-native Event Store built on Khepri/Ra with Raft consensus. Event sourcing, persistent subscriptions, snapshots, and automatic cluster formation via UDP multicast discovery. Ships embedded Rust NIFs for 3-15x acceleration of crypto, hashing, compression, aggregation, filter matching, and grap...
Current section
Files
Jump to
Current section
Files
src/reckon_db_emitter.erl
%% @doc Emitter worker for reckon-db
%%
%% A gen_server that handles event broadcasting for subscriptions.
%% Each subscription has a pool of emitter workers that receive events
%% from Khepri triggers and forward them to subscribers.
%%
%% Message types handled:
%% - {broadcast, Topic, Event}: Broadcast event to all subscribers on topic
%% - {forward_to_local, Topic, Event}: Forward event locally (optimization)
%% - {events, [Event]}: Direct event delivery to subscriber pid
%%
%% @author rgfaber
-module(reckon_db_emitter).
-behaviour(gen_server).
-include("reckon_db.hrl").
-include("reckon_db_telemetry.hrl").
%% API
-export([
start_link/4,
child_spec/4,
update_subscriber/2
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2,
code_change/3
]).
-record(state, {
store_id :: atom(),
subscription_key :: binary(),
subscriber :: pid() | undefined,
topic :: binary(),
active = true :: boolean()
}).
%%====================================================================
%% API
%%====================================================================
%% @doc Start an emitter worker
-spec start_link(atom(), binary(), pid() | undefined, atom()) -> {ok, pid()} | {error, term()}.
start_link(StoreId, SubscriptionKey, Subscriber, EmitterName) ->
gen_server:start_link({local, EmitterName}, ?MODULE,
{StoreId, SubscriptionKey, Subscriber}, []).
%% @doc Create a child spec for the emitter worker
-spec child_spec(atom(), binary(), pid() | undefined, atom()) -> supervisor:child_spec().
child_spec(StoreId, SubscriptionKey, Subscriber, EmitterName) ->
#{
id => EmitterName,
start => {?MODULE, start_link, [StoreId, SubscriptionKey, Subscriber, EmitterName]},
restart => permanent,
shutdown => 5000,
type => worker,
modules => [?MODULE]
}.
%% @doc Update the subscriber pid
-spec update_subscriber(pid() | atom(), pid()) -> ok.
update_subscriber(Emitter, NewSubscriber) ->
gen_server:cast(Emitter, {update_subscriber, NewSubscriber}).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init({StoreId, SubscriptionKey, Subscriber}) ->
process_flag(trap_exit, true),
Topic = reckon_db_emitter_group:topic(StoreId, SubscriptionKey),
%% Join the emitter group
ok = reckon_db_emitter_group:join(StoreId, SubscriptionKey, self()),
logger:info("Emitter worker started: store=~p, subscription=~s, topic=~s",
[StoreId, SubscriptionKey, Topic]),
State = #state{
store_id = StoreId,
subscription_key = SubscriptionKey,
subscriber = Subscriber,
topic = Topic,
active = true
},
{ok, State}.
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast({update_subscriber, NewSubscriber}, State) ->
{noreply, State#state{subscriber = NewSubscriber}};
handle_cast(_Msg, State) ->
{noreply, State}.
%% Handle broadcast messages (from remote emitters)
handle_info({broadcast, Topic, Event}, State) ->
handle_event_delivery(Topic, Event, State),
{noreply, State};
%% Handle forward_to_local messages (from local emitters)
handle_info({forward_to_local, Topic, Event}, State) ->
handle_event_delivery(Topic, Event, State),
{noreply, State};
%% Handle direct events message
handle_info({events, Events}, #state{subscriber = Subscriber} = State) when is_list(Events) ->
maybe_forward_events(Subscriber, Events),
{noreply, State};
%% Handle EXIT from linked processes
handle_info({'EXIT', Pid, Reason}, #state{subscriber = Subscriber} = State) ->
case Subscriber of
Pid ->
logger:info("Subscriber ~p exited with reason: ~p", [Pid, Reason]),
{noreply, State#state{subscriber = undefined}};
_ ->
{noreply, State}
end;
handle_info(_Info, State) ->
{noreply, State}.
terminate(Reason, #state{store_id = StoreId, subscription_key = SubscriptionKey}) ->
%% Leave the emitter group
ok = reckon_db_emitter_group:leave(StoreId, SubscriptionKey, self()),
logger:info("Emitter worker terminating: store=~p, subscription=~s, reason=~p",
[StoreId, SubscriptionKey, Reason]),
ok.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%====================================================================
%% Internal functions
%%====================================================================
%% @private Handle event delivery to subscriber
-spec handle_event_delivery(binary(), event(), #state{}) -> ok.
handle_event_delivery(_Topic, _Event, #state{active = false}) ->
ok;
handle_event_delivery(Topic, Event, #state{
store_id = StoreId,
subscription_key = SubscriptionKey,
subscriber = Subscriber
}) ->
telemetry:execute(
?SUBSCRIPTION_EVENT_DELIVERED,
#{count => 1},
#{store_id => StoreId, subscription_key => SubscriptionKey, topic => Topic}
),
deliver_event(Subscriber, Event, StoreId, SubscriptionKey, Topic).
maybe_forward_events(undefined, _Events) ->
ok;
maybe_forward_events(Pid, Events) when is_pid(Pid) ->
case erlang:is_process_alive(Pid) of
true -> Pid ! {events, Events};
false -> ok
end.
deliver_event(undefined, Event, StoreId, _SubscriptionKey, Topic) ->
broadcast_to_topic(StoreId, Topic, Event);
deliver_event(Pid, Event, StoreId, SubscriptionKey, _Topic) when is_pid(Pid) ->
send_to_subscriber(Pid, Event, StoreId, SubscriptionKey).
%% @private Broadcast event to topic subscribers via pg
-spec broadcast_to_topic(atom(), binary(), event()) -> ok.
broadcast_to_topic(StoreId, Topic, Event) ->
%% Get subscribers for the topic
Group = {StoreId, Topic, subscribers},
Members = pg:get_members(?RECKON_DB_PG_SCOPE, Group),
lists:foreach(
fun(Pid) ->
catch Pid ! {events, [Event]}
end,
Members
),
ok.
%% @private Send event directly to subscriber, stop pool if subscriber is dead
-spec send_to_subscriber(pid(), event(), atom(), binary()) -> ok.
send_to_subscriber(Pid, Event, StoreId, SubscriptionKey) ->
case erlang:is_process_alive(Pid) of
true ->
Pid ! {events, [Event]},
ok;
false ->
logger:warning("Subscriber ~p is dead for subscription ~s in store ~p, "
"stopping emitter pool",
[Pid, SubscriptionKey, StoreId]),
%% Stop the emitter pool asynchronously to avoid blocking the
%% emitter worker during event delivery
spawn(fun() -> reckon_db_emitter_pool:stop(StoreId, SubscriptionKey) end),
ok
end.