Packages

macula

0.8.17
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).
%% 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]).
%% 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(map(), pid()) -> {error, not_supported}.
connect(_Opts, _EventHandler) ->
{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(), pid()) -> {ok, reference()} | {error, term()}.
subscribe(Pid, Topic, HandlerPid) ->
gen_server:call(Pid, {subscribe, Topic, HandlerPid}).
%% @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).
%%==============================================================================
%% gen_server callbacks
%%==============================================================================
init(Opts) ->
Realm = maps:get(realm, Opts, <<"default">>),
EventHandler = maps:get(event_handler, Opts, self()),
io:format("[LocalClient] Initializing local client for realm ~s~n", [Realm]),
%% Find the local gateway process
case find_gateway() of
{ok, GatewayPid} ->
io:format("[LocalClient] Connected to local gateway: ~p~n", [GatewayPid]),
monitor(process, GatewayPid),
State = #state{
realm = Realm,
gateway_pid = GatewayPid,
event_handler = EventHandler
},
{ok, State};
{error, Reason} ->
io:format("[LocalClient] Failed to find gateway: ~p~n", [Reason]),
{stop, {gateway_not_found, Reason}}
end.
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}.
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,
%% Find all callbacks for subscriptions matching this topic
%% (since we store SubRef -> {Topic, Callback})
maps:foreach(fun(_SubRef, {SubTopic, Callback}) ->
case SubTopic of
Topic ->
%% Topic matches exactly, invoke callback
try
Callback(Payload)
catch
Class:Reason:Stacktrace ->
io:format("[LocalClient] Callback error for topic ~s: ~p:~p~n~p~n",
[Topic, Class, Reason, Stacktrace])
end;
_ ->
%% Different topic, skip
ok
end
end, Subs),
{noreply, State};
handle_info({'DOWN', _Ref, process, GatewayPid, Reason}, #state{gateway_pid = GatewayPid} = State) ->
io:format("[LocalClient] Gateway down: ~p. Attempting reconnect...~n", [Reason]),
%% Try to reconnect to gateway
case find_gateway() of
{ok, NewGateway} ->
io:format("[LocalClient] Reconnected to gateway: ~p~n", [NewGateway]),
monitor(process, NewGateway),
{noreply, State#state{gateway_pid = NewGateway}};
{error, _} ->
io:format("[LocalClient] Gateway not available, stopping~n"),
{stop, gateway_unavailable, State}
end;
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() ->
case whereis(macula_gateway) of
Pid when is_pid(Pid) ->
{ok, Pid};
undefined ->
{error, not_found}
end.