Packages

macula

0.8.15
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 = #{} :: #{binary() => reference()},
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, HandlerPid}, _From, State) ->
#state{gateway_pid = Gateway, realm = Realm, subscriptions = Subs} = State,
%% Subscribe via local gateway
case gen_server:call(Gateway, {local_subscribe, Realm, Topic, HandlerPid}) of
{ok, SubRef} ->
NewSubs = maps:put(Topic, SubRef, 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 tracking
NewSubs = maps:filter(fun(_, Ref) -> Ref =/= SubRef end, 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}.
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
maps:foreach(fun(_, SubRef) ->
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.