Packages
macula
0.25.6
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
Current section
Files
src/macula_peer.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Macula Peer - Mesh Participant API (v0.7.0+).
%%%
%%% This module provides the high-level API for mesh participants.
%%% Use this module to connect to a Macula mesh and communicate via pub/sub or RPC.
%%%
%%% == Quick Start ==
%%%
%%% ```
%%% %% 1. Connect to a gateway
%%% {ok, Peer} = macula_peer:start_link(<<"https://gateway.example.com:9443">>, #{
%%% realm => <<"com.example.app">>
%%% }).
%%%
%%% %% 2. Subscribe to events
%%% ok = macula_peer:subscribe(Peer, <<"sensor.temperature">>, self()).
%%%
%%% %% 3. Publish an event
%%% ok = macula_peer:publish(Peer, <<"sensor.temperature">>, #{
%%% device_id => <<"sensor-001">>,
%%% celsius => 21.5
%%% }).
%%%
%%% %% 4. Call a remote service
%%% {ok, Result} = macula_peer:call(Peer, <<"calculator.add">>, #{a => 5, b => 3}).
%%%
%%% %% 5. Advertise a service
%%% ok = macula_peer:advertise(Peer, <<"calculator.add">>, fun(#{a := A, b := B}) ->
%%% #{result => A + B}
%%% end, #{ttl => 300}).
%%% '''
%%%
%%% == Architecture ==
%%%
%%% The peer acts as a facade/coordinator, delegating to specialized child processes:
%%% - `macula_connection': QUIC transport layer (send/receive, encoding/decoding)
%%% - `macula_pubsub_handler': Pub/sub message routing
%%% - `macula_rpc_handler': RPC call/response handling
%%% - `macula_advertisement_manager': DHT service advertisements
%%%
%%% Renamed from macula_connection in v0.7.0 for clarity:
%%% - `macula_peer' = mesh participant (this module)
%%% - `macula_connection' = QUIC transport (low-level)
%%%
%%% == Multi-Tenancy via Realms ==
%%%
%%% Realms provide logical isolation for different applications:
%%%
%%% ```
%%% %% App 1
%%% {ok, Peer1} = macula_peer:start_link(GatewayUrl, #{realm => <<"com.app1">>}).
%%%
%%% %% App 2 (completely isolated from App 1)
%%% {ok, Peer2} = macula_peer:start_link(GatewayUrl, #{realm => <<"com.app2">>}).
%%% '''
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_peer).
-behaviour(gen_server).
%% API
-export([
start_link/2,
stop/1,
publish/3,
publish/4,
subscribe/3,
unsubscribe/2,
discover_subscribers/2,
call/3,
call/4,
call_to/5,
advertise/4,
unadvertise/2,
get_node_id/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2,
code_change/3
]).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
-record(state, {
url :: binary(),
realm :: binary(),
node_id :: binary(),
%% Supervision tree child PIDs
supervisor_pid :: pid(),
connection_manager_pid :: pid(),
pubsub_handler_pid :: pid(),
rpc_handler_pid :: pid(),
advertisement_manager_pid :: pid()
}).
-define(CONNECT_RETRY_DELAY, 1000).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start a client connection to a Macula mesh.
-spec start_link(binary(), map()) -> {ok, pid()} | {error, term()}.
start_link(Url, Opts) ->
gen_server:start_link(?MODULE, {Url, Opts}, []).
%% @doc Stop the client connection.
-spec stop(pid()) -> ok.
stop(Client) ->
gen_server:stop(Client).
%% @doc Publish an event through this client (no options).
-spec publish(pid(), binary(), map() | binary()) -> ok | {error, term()}.
publish(Client, Topic, Data) ->
publish(Client, Topic, Data, #{}).
%% @doc Publish an event through this client with options.
%% This is fire-and-forget - returns ok immediately without blocking.
%% Use QoS 1 in Opts if you need delivery confirmation.
-spec publish(pid(), binary(), map() | binary(), map()) -> ok.
publish(Client, Topic, Data, Opts) ->
gen_server:cast(Client, {publish, Topic, Data, Opts}),
ok.
%% @doc Subscribe to a topic through this client.
-spec subscribe(pid(), binary(), fun((map()) -> ok)) ->
{ok, reference()} | {error, term()}.
subscribe(Client, Topic, Callback) ->
gen_server:call(Client, {subscribe, Topic, Callback}, ?DEFAULT_TIMEOUT).
%% @doc Unsubscribe from a topic.
-spec unsubscribe(pid(), reference()) -> ok | {error, term()}.
unsubscribe(Client, SubRef) ->
gen_server:call(Client, {unsubscribe, SubRef}, ?DEFAULT_TIMEOUT).
%% @doc Discover subscribers to a topic via DHT query.
-spec discover_subscribers(pid(), binary()) ->
{ok, [#{node_id := binary(), endpoint := binary()}]} | {error, term()}.
discover_subscribers(Client, Topic) ->
gen_server:call(Client, {discover_subscribers, Topic}, ?DEFAULT_TIMEOUT).
%% @doc Get the node ID of this peer.
-spec get_node_id(pid()) -> {ok, binary()} | {error, term()}.
get_node_id(Client) ->
gen_server:call(Client, get_node_id, ?DEFAULT_TIMEOUT).
%% @doc Make an RPC call through this client (default timeout).
-spec call(pid(), binary(), map() | list()) -> {ok, term()} | {error, term()}.
call(Client, Procedure, Args) ->
call(Client, Procedure, Args, #{}).
%% @doc Make an RPC call through this client with options.
-spec call(pid(), binary(), map() | list(), map()) ->
{ok, term()} | {error, term()}.
call(Client, Procedure, Args, Opts) ->
Timeout = maps:get(timeout, Opts, ?CALL_TIMEOUT),
gen_server:call(Client, {call, Procedure, Args, Opts}, Timeout + 1000).
%% @doc Make an RPC call to a specific target node.
%%
%% Unlike `call/4' which discovers any provider via DHT, this function
%% sends the RPC directly to the specified target node.
-spec call_to(pid(), binary(), binary(), map() | list(), map()) ->
{ok, term()} | {error, term()}.
call_to(Client, TargetNodeId, Procedure, Args, Opts) ->
Timeout = maps:get(timeout, Opts, ?CALL_TIMEOUT),
gen_server:call(Client, {call_to, TargetNodeId, Procedure, Args, Opts}, Timeout + 1000).
%% @doc Advertise a service handler for a procedure.
%%
%% This makes the local handler available to other mesh nodes via DHT.
%% The handler will be periodically re-advertised based on TTL.
-spec advertise(pid(), binary(), fun((map()) -> {ok, term()} | {error, term()}), map()) ->
ok | {error, term()}.
advertise(Client, Procedure, Handler, Opts) ->
gen_server:call(Client, {advertise, Procedure, Handler, Opts}, ?DEFAULT_TIMEOUT).
%% @doc Stop advertising a service.
%%
%% Removes the local handler and stops advertising to the DHT.
-spec unadvertise(pid(), binary()) -> ok | {error, term()}.
unadvertise(Client, Procedure) ->
gen_server:call(Client, {unadvertise, Procedure}, ?DEFAULT_TIMEOUT).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Url, Opts}) ->
%% Parse URL to extract host and port
{Host, Port} = macula_utils:parse_url(Url),
%% Get realm (required)
Realm = get_realm_from_opts(Opts),
%% Generate or get node ID
NodeId = maps:get(node_id, Opts, macula_utils:generate_node_id()),
?LOG_INFO("[Connection Facade] Starting supervision tree for ~s", [Url]),
%% Prepare opts for handlers (includes node_id, realm, url)
HandlerOpts = Opts#{
node_id => NodeId,
realm => Realm,
url => Url,
host => Host,
port => Port
},
%% Start supervision tree
{ok, SupPid} = macula_peer_system:start_link(Url, HandlerOpts),
%% Look up child PIDs from supervisor
Children = supervisor:which_children(SupPid),
ConnMgrPid = find_child_pid(Children, connection_manager),
PubSubPid = find_child_pid(Children, pubsub_handler),
RpcPid = find_child_pid(Children, rpc_handler),
AdvMgrPid = find_child_pid(Children, advertisement_manager),
?LOG_INFO("[Connection Facade] Supervision tree started - ConnMgr: ~p, PubSub: ~p, RPC: ~p, AdvMgr: ~p",
[ConnMgrPid, PubSubPid, RpcPid, AdvMgrPid]),
%% Wait for QUIC connection to be established before returning
%% This prevents race conditions where subscribe is called before connection is ready
%% Note: wait_for_connection returns ok even on timeout - connection will establish
%% asynchronously via macula_connection's exponential backoff retry loop
wait_for_connection(ConnMgrPid, 10000),
%% Send connection_manager_pid to children that need it
gen_server:cast(PubSubPid, {set_connection_manager_pid, ConnMgrPid}),
gen_server:cast(RpcPid, {set_connection_manager_pid, ConnMgrPid}),
gen_server:cast(AdvMgrPid, {set_connection_manager_pid, ConnMgrPid}),
State = #state{
url = Url,
realm = Realm,
node_id = NodeId,
supervisor_pid = SupPid,
connection_manager_pid = ConnMgrPid,
pubsub_handler_pid = PubSubPid,
rpc_handler_pid = RpcPid,
advertisement_manager_pid = AdvMgrPid
},
%% Defer health advertisement — gateway may be busy at boot
self() ! {deferred_advertise_health, AdvMgrPid, NodeId},
{ok, State}.
%% @private Advertise the built-in _peer.health RPC procedure.
%% Returns node identity, uptime, and version. Used by LAN scanners
%% and mesh-level health monitoring.
advertise_node_health(AdvMgrPid, NodeId) ->
StartTime = erlang:system_time(second),
Handler = fun(_Args) ->
{ok, #{
node_id => binary:encode_hex(NodeId),
node_name => atom_to_binary(node()),
uptime_s => erlang:system_time(second) - StartTime,
otp_release => list_to_binary(erlang:system_info(otp_release)),
version => app_version(macula)
}}
end,
case macula_advertisement_manager:advertise_service(AdvMgrPid, <<"_peer.health">>, Handler, #{}) of
{ok, _Ref} ->
?LOG_INFO("[Peer] Advertised _peer.health");
{error, Reason} ->
?LOG_WARNING("[Peer] Failed to advertise _peer.health: ~p", [Reason])
end.
%% @private Get application version from .app.src at runtime.
app_version(App) ->
case application:get_key(App, vsn) of
{ok, Vsn} -> list_to_binary(Vsn);
_ -> <<"unknown">>
end.
%% @private
%% NOTE: publish is now handled via handle_cast (fire-and-forget semantics)
%% Delegate to pubsub_handler
handle_call({subscribe, Topic, Callback}, _From, State) ->
Result = macula_pubsub_handler:subscribe(State#state.pubsub_handler_pid, Topic, Callback),
{reply, Result, State};
%% Delegate to pubsub_handler
handle_call({unsubscribe, SubRef}, _From, State) ->
Result = macula_pubsub_handler:unsubscribe(State#state.pubsub_handler_pid, SubRef),
{reply, Result, State};
%% Discover subscribers via DHT.
%% Queries routing server with find_value (local + network lookup).
%% Catches gen_server timeout to prevent crashing the peer process.
handle_call({discover_subscribers, Topic}, _From, State) ->
TopicKey = crypto:hash(sha256, Topic),
Result = safe_discover(TopicKey, Topic, State#state.node_id),
{reply, Result, State};
%% Get node ID
handle_call(get_node_id, _From, State) ->
{reply, {ok, State#state.node_id}, State};
%% Delegate to rpc_handler
handle_call({call, Procedure, Args, Opts}, _From, State) ->
Result = macula_rpc_handler:call(State#state.rpc_handler_pid, Procedure, Args, Opts),
{reply, Result, State};
%% Delegate to rpc_handler (targeted call)
handle_call({call_to, TargetNodeId, Procedure, Args, Opts}, _From, State) ->
Result = macula_rpc_handler:call_to(State#state.rpc_handler_pid, TargetNodeId, Procedure, Args, Opts),
{reply, Result, State};
%% Delegate to advertisement_manager
handle_call({advertise, Procedure, Handler, Opts}, _From, State) ->
Result = macula_advertisement_manager:advertise_service(State#state.advertisement_manager_pid, Procedure, Handler, Opts),
{reply, Result, State};
%% Delegate to advertisement_manager
handle_call({unadvertise, Procedure}, _From, State) ->
Result = macula_advertisement_manager:unadvertise_service(State#state.advertisement_manager_pid, Procedure),
{reply, Result, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
%% Async publish - fire-and-forget semantics
%% Handle both {publish, ...} (from macula_peer:publish/4) and
%% {publish_async, ...} (from macula:publish/4 facade)
handle_cast({publish, Topic, Data, Opts}, State) ->
do_publish(Topic, Data, Opts, State);
handle_cast({publish_async, Topic, Data, Opts}, State) ->
do_publish(Topic, Data, Opts, State);
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info({deferred_advertise_health, AdvMgrPid, NodeId}, State) ->
spawn(fun() -> advertise_node_health(AdvMgrPid, NodeId) end),
{noreply, State};
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, #state{supervisor_pid = SupPid}) ->
?LOG_INFO("[Connection Facade] Terminating, stopping supervision tree"),
%% Stop the supervisor (will stop all children)
case SupPid of
Pid when is_pid(Pid) ->
macula_peer_system:stop(Pid);
_ ->
ok
end,
ok.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private Safely discover subscribers, catching gen_server timeouts.
-spec safe_discover(binary(), binary(), binary()) -> {ok, list()} | {error, term()}.
safe_discover(TopicKey, Topic, NodeId) ->
case whereis(macula_routing_server) of
undefined ->
logger:warning("[~s] Routing server not running, cannot discover subscribers",
[NodeId]),
{error, routing_server_not_running};
Pid ->
try macula_routing_server:find_value(Pid, TopicKey, 20) of
{ok, Subscribers} when is_list(Subscribers) ->
logger:debug("[~s] Found ~p subscriber(s) for topic ~s",
[NodeId, length(Subscribers), Topic]),
{ok, Subscribers};
{error, not_found} ->
{ok, []};
{error, Reason} ->
logger:warning("[~s] Failed to discover subscribers for ~s: ~p",
[NodeId, Topic, Reason]),
{error, Reason}
catch
exit:{timeout, _} ->
logger:warning("[~s] Discover subscribers timed out for ~s", [NodeId, Topic]),
{error, timeout};
exit:{noproc, _} ->
{error, routing_server_not_running}
end
end.
%% @doc Extract and normalize realm from options (pattern matching on type).
-spec get_realm_from_opts(map()) -> binary().
get_realm_from_opts(Opts) ->
normalize_realm(maps:get(realm, Opts, undefined)).
normalize_realm(undefined) ->
error({missing_required_option, realm});
normalize_realm(Realm) when is_binary(Realm) ->
Realm;
normalize_realm(Realm) when is_list(Realm) ->
list_to_binary(Realm);
normalize_realm(Realm) when is_atom(Realm) ->
atom_to_binary(Realm).
%% @doc Execute publish operation - shared by both {publish, ...} and {publish_async, ...}
-spec do_publish(binary(), map() | binary(), map(), #state{}) -> {noreply, #state{}}.
do_publish(Topic, Data, Opts, State) ->
?LOG_DEBUG("[Peer] publish received: topic=~s, pubsub_pid=~p",
[Topic, State#state.pubsub_handler_pid]),
%% Delegate to pubsub_handler (which is also async)
macula_pubsub_handler:publish(State#state.pubsub_handler_pid, Topic, Data, Opts),
{noreply, State}.
%% @doc Find child PID from supervisor children list.
-spec find_child_pid(list(), atom()) -> pid().
find_child_pid(Children, ChildId) ->
extract_child_pid(lists:keyfind(ChildId, 1, Children), ChildId).
extract_child_pid({_, Pid, _Type, _Modules}, _ChildId) when is_pid(Pid) ->
Pid;
extract_child_pid({_, undefined, _Type, _Modules}, ChildId) ->
error({child_not_started, ChildId});
extract_child_pid(false, ChildId) ->
error({child_not_found, ChildId}).
%% @doc Wait for connection manager to reach connected status.
%% Polls the connection status at 100ms intervals until connected or timeout.
%% Returns ok even on timeout to not block callers - messages may be queued.
-spec wait_for_connection(pid(), pos_integer()) -> ok.
wait_for_connection(ConnMgrPid, Timeout) when Timeout > 0 ->
wait_for_connection_status(is_process_alive(ConnMgrPid), ConnMgrPid, Timeout);
wait_for_connection(_ConnMgrPid, _Timeout) ->
%% Timeout reached - continue anyway, messages will be queued or dropped
?LOG_WARNING("[Connection Facade] Connection wait timed out, continuing anyway"),
ok.
%% @private Process not alive - sleep and retry
wait_for_connection_status(false, ConnMgrPid, Timeout) ->
timer:sleep(100),
wait_for_connection(ConnMgrPid, Timeout - 100);
%% @private Process alive - check status
wait_for_connection_status(true, ConnMgrPid, Timeout) ->
handle_connection_status(macula_connection:get_status(ConnMgrPid), ConnMgrPid, Timeout).
%% @private Connected - done
handle_connection_status(connected, _ConnMgrPid, _Timeout) ->
ok;
%% @private Not connected - sleep and retry
handle_connection_status(_Other, ConnMgrPid, Timeout) ->
timer:sleep(100),
wait_for_connection(ConnMgrPid, Timeout - 100).