Packages

macula

0.38.8
7.0.0 6.0.0 5.2.2 5.2.1 5.2.0 5.1.0 5.0.0 4.8.0 4.7.1 4.7.0 4.6.0 4.5.0 4.4.10 4.4.9 4.4.8 4.4.7 4.4.6 4.4.5 4.4.4 4.4.3 4.4.2 4.4.1 4.4.0 4.3.1 4.3.0 4.2.9 4.2.8 4.2.7 4.2.6 4.2.5 4.2.4 4.2.3 4.2.2 4.2.1 4.2.0 4.1.1 4.1.0 4.0.0 3.16.0 3.15.3 3.15.2 3.15.1 3.14.0 3.13.0 3.12.1 3.12.0 3.11.1 3.11.0 3.10.3 3.10.2 3.10.1 3.9.0 3.8.0 3.7.0 3.5.0 3.4.0 3.3.0 3.2.0 3.1.0 3.0.0 2.1.1 2.1.0 2.0.0 1.5.2 1.5.1 1.4.30 1.4.29 1.4.28 1.4.27 1.4.26 1.4.25 1.4.24 1.4.23 1.4.22 1.4.21 1.4.20 1.4.19 1.4.18 1.4.17 1.4.16 1.4.15 1.4.14 1.4.13 1.4.11 1.4.10 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.1 1.3.0 1.2.0 1.1.0 1.0.10 1.0.9 1.0.8 1.0.7 1.0.6 1.0.5 1.0.4 1.0.3 1.0.2 1.0.1 1.0.0 0.48.6 0.48.5 0.48.4 0.48.3 0.48.2 0.48.1 0.48.0 0.47.1 0.47.0 0.46.3 0.46.1 0.46.0 0.45.3 0.45.2 0.45.1 0.45.0 0.44.2 0.44.1 0.44.0 0.43.3 0.43.2 0.43.1 0.43.0 0.42.9 0.42.8 0.42.7 0.42.6 0.42.5 0.42.4 0.42.3 0.42.2 0.42.1 0.42.0 0.41.1 0.41.0 0.40.1 0.40.0 0.39.9 0.39.8 0.39.7 0.39.6 0.39.5 0.39.4 0.39.3 0.39.2 0.39.1 0.39.0 0.38.8 0.38.7 0.38.6 0.38.5 0.38.4 0.38.3 0.38.2 0.38.1 0.38.0 0.37.7 0.37.6 0.37.5 0.37.4 0.37.3 0.37.2 0.37.1 0.37.0 0.36.6 0.36.5 0.36.4 0.36.3 0.36.2 0.36.1 0.36.0 0.35.4 0.35.3 0.35.2 0.35.1 0.35.0 0.34.1 0.34.0 0.33.1 0.33.0 0.32.5 0.32.4 0.32.3 0.32.2 0.32.1 0.32.0 0.31.9 0.31.8 0.31.7 0.31.6 0.31.5 0.31.4 0.31.3 0.31.2 0.31.1 0.31.0 0.30.10 0.30.9 0.30.8 0.30.7 0.30.6 0.30.5 0.30.4 0.30.3 0.30.2 0.30.1 0.30.0 0.29.0 0.28.3 0.28.2 0.28.1 0.28.0 0.27.1 0.27.0 0.26.1 0.26.0 0.25.6 0.25.5 0.25.4 0.25.3 0.25.2 0.25.1 0.25.0 0.24.6 0.24.5 0.24.4 0.24.3 0.24.2 0.24.1 0.24.0 0.23.3 0.23.2 0.23.1 0.23.0 0.22.12 0.22.11 0.22.10 0.22.9 0.22.8 0.22.7 0.22.6 0.22.5 0.22.4 0.22.3 0.22.2 0.22.1 0.22.0 0.21.7 0.21.6 0.21.5 0.21.4 0.21.2 0.21.1 0.21.0 0.20.25 0.20.24 0.20.23 0.20.22 0.20.21 0.20.20 0.20.19 0.20.18 0.20.17 0.20.16 0.20.15 0.20.14 0.20.13 0.20.12 0.20.11 0.20.10 0.20.9 0.20.8 0.20.7 0.20.6 0.20.5 0.20.3 0.20.2 0.20.1 0.20.0 0.19.2 0.19.1 0.19.0 0.18.1 0.18.0 0.17.4 0.17.3 0.17.2 0.17.1 0.17.0 0.16.6 0.16.5 0.16.4 0.16.3 0.16.2 0.16.1 0.16.0 0.15.1 0.15.0 0.14.3 0.14.2 0.14.1 0.14.0 0.12.6 0.12.5 0.12.3 0.11.3 0.10.2 0.10.1 0.10.0 0.9.2 0.9.1 0.9.0 0.8.25 0.8.24 0.8.23 0.8.22 0.8.21 0.8.20 0.8.19 0.8.18 0.8.17 0.8.16 0.8.15 0.8.14 0.8.13 0.8.12 0.8.11 0.8.10 0.8.9 0.8.8 0.8.7 0.8.6 0.8.5 0.8.4 0.8.3 0.8.2 0.8.1 0.8.0 0.7.30 0.7.29 0.7.28 0.7.27 0.7.26 0.7.25 0.7.24 0.7.23 0.7.22 0.7.21 0.7.20 0.7.19 0.7.18 0.7.17 0.7.16 0.7.15 0.7.14 0.7.13 0.7.12 0.7.11 0.7.10 0.7.9 0.7.8 0.7.7 0.7.6 0.7.5 0.7.4 0.7.3 0.7.2 0.7.1 0.7.0 0.6.7 0.6.6 0.6.5 0.6.4 0.6.3 0.6.2 0.6.1 0.6.0 0.5.0 0.4.4 0.4.3 0.4.2 0.4.1 0.4.0 0.3.4 0.3.3 0.3.2 0.3.1

Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise

Current section

Files

Jump to
macula src macula_local_client.erl
Raw

src/macula_local_client.erl

%% @doc Local client for in-VM workloads to connect to macula_gateway
%%
%% This module provides process-to-process communication between workloads
%% running in the same BEAM VM as the Macula platform and the local gateway.
%% Unlike macula_peer which creates QUIC connections, this connects directly
%% to the local macula_gateway process.
%%
%% Architecture:
%% Phoenix/Elixir App → macula_local_client → macula_gateway
%% ↓ (QUIC)
%% Other Peers
%%
%% @end
-module(macula_local_client).
-behaviour(gen_server).
-behaviour(macula_client_behaviour).
-include_lib("kernel/include/logger.hrl").
%% API - Connection management
-export([connect/2, connect_local/1, disconnect/1]).
%% API - Pub/Sub
-export([publish/3, publish/4, subscribe/3, unsubscribe/2, discover_subscribers/2]).
%% API - RPC
-export([call/3, call/4, advertise/3, advertise/4, unadvertise/2]).
%% API - Utility
-export([get_node_id/1]).
%% API - Platform Layer (v0.10.0+)
-export([
register_workload/2,
get_leader/1,
subscribe_leader_changes/2,
propose_crdt_update/3,
propose_crdt_update/4,
read_crdt/2
]).
%% Legacy API (kept for backward compatibility)
-export([start_link/1, stop/1, register_procedure/3, unregister_procedure/2]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
realm :: binary(),
gateway_pid :: pid() | undefined,
subscriptions = #{} :: #{reference() => {binary(), fun()}}, % SubRef -> {Topic, Callback}
registrations = #{} :: #{binary() => fun()},
event_handler :: pid() | undefined
}).
%%==============================================================================
%% API
%%==============================================================================
%%------------------------------------------------------------------------------
%% Connection Management
%%------------------------------------------------------------------------------
%% @doc Connect to remote gateway (not supported for local client)
%% For compatibility with macula_client_behaviour
-spec connect(binary() | string(), map()) -> {error, not_supported}.
connect(_Url, _Opts) ->
{error, not_supported}.
%% @doc Create a local client connection to the gateway
-spec connect_local(map()) -> {ok, pid()} | {error, term()}.
connect_local(Opts) ->
start_link(Opts).
%% @doc Disconnect the client
-spec disconnect(pid()) -> ok.
disconnect(Pid) ->
stop(Pid).
%% @doc Start a local client connection to the gateway (legacy API)
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Stop the local client (legacy API)
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%%------------------------------------------------------------------------------
%% Pub/Sub Operations
%%------------------------------------------------------------------------------
%% @doc Publish an event to a topic
-spec publish(pid(), binary(), map()) -> ok | {error, term()}.
publish(Pid, Topic, Payload) ->
publish(Pid, Topic, Payload, #{}).
%% @doc Publish an event to a topic with options
-spec publish(pid(), binary(), map(), map()) -> ok | {error, term()}.
publish(Pid, Topic, Payload, Opts) ->
gen_server:call(Pid, {publish, Topic, Payload, Opts}).
%% @doc Subscribe to a topic
-spec subscribe(pid(), binary(), fun((map()) -> ok)) -> {ok, reference()} | {error, term()}.
subscribe(Pid, Topic, Callback) ->
gen_server:call(Pid, {subscribe, Topic, Callback}).
%% @doc Unsubscribe from a topic
-spec unsubscribe(pid(), reference()) -> ok | {error, term()}.
unsubscribe(Pid, SubRef) ->
gen_server:call(Pid, {unsubscribe, SubRef}).
%% @doc Discover subscribers of a topic via DHT query
-spec discover_subscribers(pid(), binary()) -> {ok, [binary()]} | {error, term()}.
discover_subscribers(Pid, Topic) ->
gen_server:call(Pid, {discover_subscribers, Topic}).
%%------------------------------------------------------------------------------
%% RPC Operations
%%------------------------------------------------------------------------------
%% @doc Call an RPC procedure with default options
-spec call(pid(), binary(), list()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args) ->
call(Pid, Procedure, Args, #{}).
%% @doc Call an RPC procedure
-spec call(pid(), binary(), list(), map()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args, Opts) ->
gen_server:call(Pid, {call, Procedure, Args, Opts}, 30000).
%% @doc Advertise an RPC service with default options
-spec advertise(pid(), binary(), fun()) -> {ok, reference()} | {error, term()}.
advertise(Pid, Procedure, Handler) ->
advertise(Pid, Procedure, Handler, #{}).
%% @doc Advertise an RPC service with options
-spec advertise(pid(), binary(), fun(), map()) -> {ok, reference()} | {error, term()}.
advertise(Pid, Procedure, Handler, Opts) ->
gen_server:call(Pid, {advertise, Procedure, Handler, Opts}).
%% @doc Unadvertise an RPC service
-spec unadvertise(pid(), binary()) -> ok | {error, term()}.
unadvertise(Pid, Procedure) ->
gen_server:call(Pid, {unadvertise, Procedure}).
%% @doc Register an RPC procedure (legacy API, use advertise/3 instead)
-spec register_procedure(pid(), binary(), fun()) -> ok | {error, term()}.
register_procedure(Pid, Procedure, Handler) ->
advertise(Pid, Procedure, Handler, #{}).
%% @doc Unregister an RPC procedure (legacy API, use unadvertise/2 instead)
-spec unregister_procedure(pid(), binary()) -> ok | {error, term()}.
unregister_procedure(Pid, Procedure) ->
unadvertise(Pid, Procedure).
%%------------------------------------------------------------------------------
%% Utility Operations
%%------------------------------------------------------------------------------
%% @doc Get the node ID of the local gateway
-spec get_node_id(pid()) -> {ok, binary()} | {error, term()}.
get_node_id(Pid) ->
gen_server:call(Pid, get_node_id).
%%------------------------------------------------------------------------------
%% Platform Layer Operations (v0.10.0+)
%%------------------------------------------------------------------------------
%% @doc Register this workload with the Platform Layer
-spec register_workload(pid(), map()) -> {ok, map()} | {error, term()}.
register_workload(Pid, Opts) ->
gen_server:call(Pid, {register_workload, Opts}).
%% @doc Get the current Platform Layer leader node ID
-spec get_leader(pid()) -> {ok, binary()} | {error, no_leader | term()}.
get_leader(Pid) ->
gen_server:call(Pid, get_leader).
%% @doc Subscribe to Platform Layer leader change notifications
-spec subscribe_leader_changes(pid(), fun((map()) -> ok)) ->
{ok, reference()} | {error, term()}.
subscribe_leader_changes(Pid, Callback) ->
gen_server:call(Pid, {subscribe_leader_changes, Callback}).
%% @doc Propose a CRDT update with default options (LWW-Register)
-spec propose_crdt_update(pid(), binary(), term()) -> ok | {error, term()}.
propose_crdt_update(Pid, Key, Value) ->
propose_crdt_update(Pid, Key, Value, #{crdt_type => lww_register}).
%% @doc Propose a CRDT update with specific type
-spec propose_crdt_update(pid(), binary(), term(), map()) -> ok | {error, term()}.
propose_crdt_update(Pid, Key, Value, Opts) ->
gen_server:call(Pid, {propose_crdt_update, Key, Value, Opts}).
%% @doc Read the current value of a CRDT-managed state entry
-spec read_crdt(pid(), binary()) -> {ok, term()} | {error, not_found | term()}.
read_crdt(Pid, Key) ->
gen_server:call(Pid, {read_crdt, Key}).
%%==============================================================================
%% gen_server callbacks
%%==============================================================================
init(Opts) ->
Realm = maps:get(realm, Opts, <<"default">>),
EventHandler = maps:get(event_handler, Opts, self()),
?LOG_INFO("Initializing local client for realm ~s", [Realm]),
%% Find and connect to the local gateway process
GatewayResult = find_gateway(),
init_with_gateway(GatewayResult, Realm, EventHandler).
%% @private Gateway found - initialize state
init_with_gateway({ok, GatewayPid}, Realm, EventHandler) ->
?LOG_INFO("Connected to local gateway: ~p", [GatewayPid]),
monitor(process, GatewayPid),
State = #state{
realm = Realm,
gateway_pid = GatewayPid,
event_handler = EventHandler
},
{ok, State};
%% @private Gateway not found - stop with error
init_with_gateway({error, Reason}, _Realm, _EventHandler) ->
?LOG_ERROR("Failed to find gateway: ~p", [Reason]),
{stop, {gateway_not_found, Reason}}.
handle_call({publish, Topic, Payload, Opts}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm} = State,
%% Send publish request to gateway via gen_server:call
%% Note: Opts are currently ignored for local clients
_ = Opts,
Result = case gen_server:call(Gateway, {local_publish, Realm, Topic, Payload}) of
ok -> ok;
{error, _} = Error -> Error
end,
{reply, Result, State};
handle_call({subscribe, Topic, Callback}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm, subscriptions = Subs} = State,
%% Subscribe via local gateway, passing our PID (not the callback)
%% Gateway will send pubsub events to us, and we'll invoke the callback
case gen_server:call(Gateway, {local_subscribe, Realm, Topic, self()}) of
{ok, SubRef} ->
%% Store callback for later invocation
NewSubs = maps:put(SubRef, {Topic, Callback}, Subs),
{reply, {ok, SubRef}, State#state{subscriptions = NewSubs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call({unsubscribe, SubRef}, _From, State) ->
#state{gateway_pid = Gateway, subscriptions = Subs} = State,
case gen_server:call(Gateway, {local_unsubscribe, SubRef}) of
ok ->
%% Remove from local callback tracking
NewSubs = maps:remove(SubRef, Subs),
{reply, ok, State#state{subscriptions = NewSubs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call({discover_subscribers, Topic}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm} = State,
%% Forward DHT query to gateway
Result = gen_server:call(Gateway, {local_discover_subscribers, Realm, Topic}),
{reply, Result, State};
handle_call({call, Procedure, Args, Opts}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm} = State,
%% Route RPC call through local gateway
Result = gen_server:call(Gateway, {local_rpc_call, Realm, Procedure, Args, Opts}, 30000),
{reply, Result, State};
handle_call({advertise, Procedure, Handler, Opts}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm, registrations = Regs} = State,
%% Forward advertise request to gateway
case gen_server:call(Gateway, {local_advertise, Realm, Procedure, Handler, Opts}) of
{ok, Ref} ->
NewRegs = maps:put(Procedure, Handler, Regs),
{reply, {ok, Ref}, State#state{registrations = NewRegs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call({unadvertise, Procedure}, _From, State) ->
#state{gateway_pid = Gateway, registrations = Regs} = State,
%% Forward unadvertise request to gateway
case gen_server:call(Gateway, {local_unadvertise, Procedure}) of
ok ->
NewRegs = maps:remove(Procedure, Regs),
{reply, ok, State#state{registrations = NewRegs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call({register, Procedure, Handler}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm, registrations = Regs} = State,
case gen_server:call(Gateway, {local_register_procedure, Realm, Procedure, Handler}) of
ok ->
NewRegs = maps:put(Procedure, Handler, Regs),
{reply, ok, State#state{registrations = NewRegs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call({unregister, Procedure}, _From, State) ->
#state{gateway_pid = Gateway, registrations = Regs} = State,
case gen_server:call(Gateway, {local_unregister_procedure, Procedure}) of
ok ->
NewRegs = maps:remove(Procedure, Regs),
{reply, ok, State#state{registrations = NewRegs}};
{error, _} = Error ->
{reply, Error, State}
end;
handle_call(get_node_id, _From, State) ->
#state{gateway_pid = Gateway} = State,
Result = gen_server:call(Gateway, local_get_node_id),
{reply, Result, State};
%%------------------------------------------------------------------------------
%% Platform Layer Handlers (v0.14.0+ - Masterless/CRDT-based)
%%------------------------------------------------------------------------------
handle_call({register_workload, Opts}, _From, State) ->
%% v0.14.0: Masterless architecture - no leader election
%% Workloads register locally; state is CRDT-replicated
WorkloadName = maps:get(workload_name, Opts, <<"unknown">>),
Capabilities = maps:get(capabilities, Opts, []),
?LOG_INFO("Registered workload ~s with capabilities ~p (masterless mode)",
[WorkloadName, Capabilities]),
%% Return platform info (no leader in masterless design)
Info = #{
leader_node => undefined, % No leader - masterless design
cluster_size => 0, % Cluster size via DHT peer count
platform_version => <<"0.14.0">>,
architecture => masterless
},
{reply, {ok, Info}, State};
handle_call(get_leader, _From, _State) ->
%% v0.14.0: No leader in masterless design
%% Return undefined - applications should use CRDTs for coordination
{reply, {ok, undefined}, _State};
handle_call({subscribe_leader_changes, _Callback}, _From, State) ->
%% v0.14.0: No leader changes in masterless design
%% Return error - applications should use CRDT-based state instead
{reply, {error, not_supported_in_masterless_architecture}, State};
handle_call({propose_crdt_update, Key, Value, Opts}, _From, State) ->
CrdtType = maps:get(crdt_type, Opts, lww_register),
Reply = do_propose_crdt_update(CrdtType, Key, Value),
{reply, Reply, State};
handle_call({read_crdt, Key}, _From, State) ->
Reply = do_read_crdt(Key),
{reply, Reply, State}.
%% Async publish - fire-and-forget from caller's perspective
handle_cast({publish_async, Topic, Payload, Opts}, State) ->
#state{gateway_pid = Gateway, realm = Realm} = State,
%% Send publish request to gateway via gen_server:cast (async)
%% Note: Opts are currently ignored for local clients
_ = Opts,
gen_server:cast(Gateway, {local_publish_async, Realm, Topic, Payload}),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle pubsub events from gateway
%% Gateway sends messages in format: {publish, Topic, Payload}
handle_info({publish, Topic, Payload}, State) ->
#state{subscriptions = Subs} = State,
?LOG_DEBUG("Received publish message for topic ~s", [Topic]),
?LOG_DEBUG("Current subscriptions: ~p", [Subs]),
%% Find all callbacks for subscriptions matching this topic
%% (since we store SubRef -> {Topic, Callback})
MatchCount = maps:fold(fun(SubRef, {SubTopic, Callback}, Acc) ->
?LOG_DEBUG("Checking SubRef=~p, SubTopic=~p against Topic=~p",
[SubRef, SubTopic, Topic]),
case SubTopic of
Topic ->
%% Topic matches exactly, invoke callback
?LOG_DEBUG("Match found, invoking callback for topic ~s", [Topic]),
invoke_callback_safe(Callback, Payload, Topic),
Acc + 1;
_ ->
%% Different topic, skip
?LOG_DEBUG("No match: ~p =/= ~p", [SubTopic, Topic]),
Acc
end
end, 0, Subs),
?LOG_DEBUG("Total callbacks invoked for topic ~s: ~p", [Topic, MatchCount]),
{noreply, State};
handle_info({'DOWN', _Ref, process, GatewayPid, Reason}, #state{gateway_pid = GatewayPid} = State) ->
?LOG_WARNING("Gateway down: ~p. Attempting reconnect...", [Reason]),
GatewayResult = find_gateway(),
handle_gateway_reconnect(GatewayResult, State);
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, #state{subscriptions = Subs, registrations = Regs, gateway_pid = Gateway}) ->
%% Clean up subscriptions (keys are SubRefs now)
maps:foreach(fun(SubRef, _TopicCallback) ->
gen_server:call(Gateway, {local_unsubscribe, SubRef})
end, Subs),
%% Clean up registrations (unadvertise all procedures)
maps:foreach(fun(Procedure, _) ->
gen_server:call(Gateway, {local_unadvertise, Procedure})
end, Regs),
ok.
%%==============================================================================
%% Internal functions
%%==============================================================================
%% @doc Find the local gateway process via whereis
-spec find_gateway() -> {ok, pid()} | {error, not_found}.
find_gateway() ->
GatewayPid = whereis(macula_gateway),
gateway_lookup_result(GatewayPid).
%% @private Gateway process found
gateway_lookup_result(Pid) when is_pid(Pid) ->
{ok, Pid};
%% @private Gateway not registered
gateway_lookup_result(undefined) ->
{error, not_found}.
%% NOTE: get_leader_reply/1 removed in v0.14.0 (masterless architecture)
%% Leader election is no longer used - applications use CRDTs for coordination
%% @private Gateway reconnect succeeded
handle_gateway_reconnect({ok, NewGateway}, State) ->
?LOG_INFO("Reconnected to gateway: ~p", [NewGateway]),
monitor(process, NewGateway),
{noreply, State#state{gateway_pid = NewGateway}};
%% @private Gateway reconnect failed
handle_gateway_reconnect({error, _}, State) ->
?LOG_ERROR("Gateway not available, stopping"),
{stop, gateway_unavailable, State}.
%% @private LWW-Register CRDT update
do_propose_crdt_update(lww_register, Key, Value) ->
Timestamp = erlang:system_time(millisecond),
TableName = macula_crdt_storage,
ensure_crdt_table_exists(TableName),
ets:insert(TableName, {Key, {Value, Timestamp}}),
ok;
%% @private Unsupported CRDT type
do_propose_crdt_update(CrdtType, _Key, _Value) ->
{error, {unsupported_crdt_type, CrdtType}}.
%% @private Ensure ETS table exists (handles race condition)
ensure_crdt_table_exists(TableName) ->
TableRef = ets:whereis(TableName),
do_ensure_table(TableRef, TableName).
%% @private Table doesn't exist - create it
do_ensure_table(undefined, TableName) ->
handle_table_creation(catch ets:new(TableName, [named_table, public, set]));
%% @private Table exists - nothing to do
do_ensure_table(_TableRef, _TableName) ->
ok.
%% @private Read from CRDT storage
do_read_crdt(Key) ->
TableName = macula_crdt_storage,
TableRef = ets:whereis(TableName),
do_crdt_lookup(TableRef, TableName, Key).
%% @private Table doesn't exist
do_crdt_lookup(undefined, _TableName, _Key) ->
{error, not_found};
%% @private Table exists - perform lookup
do_crdt_lookup(_TableRef, TableName, Key) ->
LookupResult = ets:lookup(TableName, Key),
crdt_lookup_result(LookupResult).
%% @private Key found in CRDT table
crdt_lookup_result([{_Key, {Value, _Timestamp}}]) ->
{ok, Value};
%% @private Key not found
crdt_lookup_result([]) ->
{error, not_found}.
%% @private Handle ETS table creation result
handle_table_creation({'EXIT', {badarg, _}}) ->
ok; %% Table already exists (race condition)
handle_table_creation(_TableRef) ->
ok.
%% @private Invoke callback safely with catch expression
invoke_callback_safe(Callback, Payload, Topic) ->
handle_callback_result(catch Callback(Payload), Topic).
%% @private Handle callback result
handle_callback_result({'EXIT', {Reason, Stacktrace}}, Topic) ->
?LOG_ERROR("Callback error for topic ~s: ~p~nStacktrace: ~p",
[Topic, Reason, Stacktrace]);
handle_callback_result({'EXIT', Reason}, Topic) ->
?LOG_ERROR("Callback error for topic ~s: ~p", [Topic, Reason]);
handle_callback_result(_Result, Topic) ->
?LOG_DEBUG("Callback completed successfully for topic ~s", [Topic]).