Packages

macula

0.32.5
7.1.0 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_relay_handler.erl
Raw

src/macula_relay_handler.erl

%%%-------------------------------------------------------------------
%%% @doc Relay handler — one process per connected node.
%%%
%%% Manages a single QUIC stream to a node. Routes messages using:
%%% - pg groups for pub/sub (subscribe = join, publish = broadcast)
%%% - gproc for RPC (register = reg, call = lookup + forward)
%%%
%%% Cross-relay routing uses dual pg groups per topic:
%%% - `{relay_topic, Topic}' — ALL handlers (client + peer) join.
%%% Used for local PUBLISH broadcasting. Peer handlers forward
%%% messages to their remote relay.
%%% - `{relay_local, Topic}' — only CLIENT handlers join.
%%% Used for cross-relay delivery (loop-safe). When a message
%%% arrives from a peer relay, it is delivered ONLY to local
%%% client handlers via this group.
%%%
%%% Peer vs client is determined by the CONNECT handshake: peers
%%% send `type => <<"relay">>'. This prevents pub/sub loops where
%%% a message would bounce between relays indefinitely.
%%%
%%% When this process dies (node disconnects), pg and gproc automatically
%%% remove all subscriptions and procedure registrations.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_relay_handler).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([start_link/2]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
conn :: reference(),
stream :: reference() | undefined,
node_id :: binary(),
node_name :: binary(), %% stable identity (e.g. "hecate@beam00.lab")
is_peer :: boolean(), %% true if this is a relay-to-relay peering connection
identified :: boolean(), %% true after CONNECT handshake processed
pending_msgs :: [term()], %% messages received before CONNECT (replayed after)
recv_buffer :: binary(),
pending_calls :: #{binary() => pid()} %% call_id => caller handler pid
}).
%%====================================================================
%% API
%%====================================================================
start_link(Conn, Stream) ->
gen_server:start_link(?MODULE, {Conn, Stream}, []).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init({Conn, Stream}) ->
%% Relay transfers ownership to us after start_link returns.
%% We wait for the signal before setting active mode.
?LOG_INFO("[relay_handler] Waiting for stream ownership"),
{ok, #state{conn = Conn, stream = Stream, node_id = <<>>,
node_name = <<>>, is_peer = false, identified = false,
pending_msgs = [], recv_buffer = <<>>, pending_calls = #{}}}.
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%%====================================================================
%% QUIC data received
%%====================================================================
handle_info(ownership_transferred, #state{stream = Stream} = State) ->
quicer:setopt(Stream, active, true),
?LOG_INFO("[relay_handler] Ownership received, stream active"),
{noreply, State};
handle_info({quic, Data, Stream, _Flags}, #state{stream = Stream} = State)
when is_binary(Data) ->
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
{NewBuffer, State2} = process_buffer(Buffer, State),
{noreply, State2#state{recv_buffer = NewBuffer}};
%% New stream from this connection (e.g., for additional data)
handle_info({quic, Data, NewStream, _Flags}, State) when is_binary(Data) ->
quicer:setopt(NewStream, active, true),
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
{NewBuffer, State2} = process_buffer(Buffer, State),
{noreply, State2#state{recv_buffer = NewBuffer, stream = NewStream}};
%%====================================================================
%% Inter-handler messages (from other relay_handler processes)
%%====================================================================
%% Pub/sub: another handler published to a topic we're subscribed to
handle_info({relay_publish, Topic, Payload}, State) ->
send_to_node(publish, #{
<<"topic">> => Topic,
<<"payload">> => Payload,
<<"qos">> => 0,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
}, State),
{noreply, State};
%% RPC: another handler wants us to forward a call to our node
handle_info({relay_call, CallerPid, CallId, Procedure, Args}, State) ->
%% Track: when our node replies, forward to CallerPid
PendingCalls = maps:put(CallId, CallerPid, State#state.pending_calls),
send_to_node(call, #{
<<"call_id">> => CallId,
<<"procedure">> => Procedure,
<<"args">> => Args
}, State),
{noreply, State#state{pending_calls = PendingCalls}};
%% RPC reply from caller side: send reply back to the calling node
handle_info({relay_reply, CallId, Result}, State) ->
send_to_node(reply, #{
<<"call_id">> => CallId,
<<"result">> => Result
}, State),
{noreply, State};
%%====================================================================
%% QUIC lifecycle
%%====================================================================
handle_info({quic, peer_send_shutdown, _Stream, _}, State) ->
?LOG_INFO("[relay_handler] Peer send shutdown, closing"),
{stop, normal, State};
handle_info({quic, shutdown, _Ref, _Reason}, State) ->
{stop, normal, State};
handle_info({quic, closed, _Ref, _Reason}, State) ->
{stop, normal, State};
handle_info({quic, streams_available, _Conn, _Info}, State) ->
{noreply, State};
handle_info({quic, peer_needs_streams, _Conn, _Info}, State) ->
{noreply, State};
handle_info(Info, State) ->
?LOG_DEBUG("[relay_handler] Unhandled: ~p", [Info]),
{noreply, State}.
terminate(Reason, _State) ->
?LOG_INFO("[relay_handler] Terminated: ~p (pg/gproc auto-cleanup)", [Reason]),
%% pg and gproc automatically remove this process from all groups/registrations
ok.
%%====================================================================
%% Message processing
%%====================================================================
%% Process complete messages from the receive buffer.
%% Returns {RemainingBuffer, UpdatedState}.
process_buffer(Buffer, State) when byte_size(Buffer) < 8 ->
{Buffer, State};
process_buffer(<<_Version:8, _TypeId:8, _Flags:8, _Reserved:8,
PayloadLen:32/big-unsigned, Rest/binary>> = Buffer, State) ->
case byte_size(Rest) >= PayloadLen of
true ->
MsgBytes = binary:part(Buffer, 0, 8 + PayloadLen),
Remaining = binary:part(Buffer, 8 + PayloadLen, byte_size(Buffer) - 8 - PayloadLen),
State2 = handle_message(macula_protocol_decoder:decode(MsgBytes), State),
process_buffer(Remaining, State2);
false ->
{Buffer, State}
end.
%%====================================================================
%% Protocol message handlers
%%====================================================================
%% CONNECT — node identifies itself
handle_message({ok, {connect, Msg}}, State) ->
NodeId = maps:get(<<"node_id">>, Msg, <<>>),
Realm = maps:get(<<"realm_id">>, Msg, <<>>),
IsPeer = maps:get(<<"type">>, Msg, <<"node">>) =:= <<"relay">>,
NodeName = maps:get(<<"node_name">>, Msg, binary:encode_hex(NodeId)),
Identity = maps:get(<<"identity">>, Msg, #{}),
TargetRelay = maps:get(<<"target_relay">>, Msg, <<>>),
Label = case IsPeer of true -> <<"Peer relay">>; false -> <<"Node">> end,
?LOG_INFO("[relay_handler] ~s connected: realm=~s name=~s target=~s",
[Label, Realm, NodeName, TargetRelay]),
register_client_node(IsPeer, NodeName, Identity, TargetRelay),
send_to_node(pong, #{
<<"timestamp">> => erlang:system_time(millisecond),
<<"server_time">> => erlang:system_time(millisecond)
}, State),
State2 = State#state{node_id = NodeId, node_name = NodeName,
is_peer = IsPeer, identified = true},
replay_pending(State2);
%% Buffer messages received before CONNECT handshake.
%% Replayed after CONNECT sets is_peer correctly, preventing
%% the race where SUBSCRIBE joins wrong pg groups.
handle_message({ok, {Type, _}} = Msg, #state{identified = false, pending_msgs = Pending} = State)
when Type =/= connect ->
State#state{pending_msgs = [Msg | Pending]};
%% SUBSCRIBE — join pg groups for topics
%% Client handlers join both {relay_topic, T} and {relay_local, T}.
%% Peer handlers join only {relay_topic, T} (prevents cross-relay loops).
handle_message({ok, {subscribe, Msg}}, #state{is_peer = IsPeer} = State) ->
Topics = maps:get(<<"topics">>, Msg, []),
lists:foreach(fun(Topic) -> join_topic(Topic, IsPeer) end, Topics),
State;
%% UNSUBSCRIBE — leave pg groups
handle_message({ok, {unsubscribe, Msg}}, State) ->
Topics = maps:get(<<"topics">>, Msg, []),
lists:foreach(fun(Topic) ->
pg:leave(pg, {relay_topic, Topic}, self()),
pg:leave(pg, {relay_local, Topic}, self())
end, Topics),
State;
%% PUBLISH on _relay.graph — merge into relay graph (QUIC-based graph propagation)
handle_message({ok, {publish, #{<<"topic">> := <<"_relay.graph">>, <<"payload">> := Payload}}}, State) ->
case catch json:decode(Payload) of
Entries when is_list(Entries) ->
?LOG_DEBUG("[relay_handler] Graph update: ~b entries", [length(Entries)]),
catch macula_relay_peering:merge_remote_graph(Entries);
_ ->
?LOG_WARNING("[relay_handler] Invalid graph update payload")
end,
State;
%% PUBLISH — broadcast to all subscribers via pg
handle_message({ok, {publish, Msg}}, State) ->
Topic = maps:get(<<"topic">>, Msg, <<>>),
Payload = maps:get(<<"payload">>, Msg, <<>>),
Members = try pg:get_members(pg, {relay_topic, Topic}) catch _:_ -> [] end,
?LOG_INFO("[relay_handler] PUBLISH ~s → ~p subscriber(s)", [Topic, length(Members)]),
lists:foreach(fun(Pid) ->
case Pid =/= self() of
true -> Pid ! {relay_publish, Topic, Payload};
false -> ok
end
end, Members),
State;
%% REGISTER_PROCEDURE — register in gproc for RPC routing
handle_message({ok, {register_procedure, Msg}}, State) ->
Procedure = maps:get(<<"procedure">>, Msg, <<>>),
%% Use property (not name) so multiple nodes can register same procedure
try
gproc:reg({p, l, {relay_rpc, Procedure}}),
?LOG_INFO("[relay_handler] Registered RPC procedure: ~s", [Procedure])
catch
_:_ -> ?LOG_DEBUG("[relay_handler] Procedure ~s already registered", [Procedure])
end,
State;
%% CALL — find provider via gproc, forward to its handler.
%% If not found locally, try peer relays via the peering module.
handle_message({ok, {call, Msg}}, State) ->
Procedure = maps:get(<<"procedure">>, Msg, <<>>),
CallId = maps:get(<<"call_id">>, Msg, <<>>),
Args = maps:get(<<"args">>, Msg, #{}),
case gproc:lookup_pids({p, l, {relay_rpc, Procedure}}) of
[] ->
?LOG_INFO("[relay_handler] RPC ~s not local, trying peer relays", [Procedure]),
forward_rpc_to_peers(Procedure, CallId, Args, State);
[ProviderPid | _] ->
?LOG_INFO("[relay_handler] Forwarding RPC ~s to provider ~p", [Procedure, ProviderPid]),
ProviderPid ! {relay_call, self(), CallId, Procedure, Args}
end,
State;
%% REPLY — route back to the caller handler
handle_message({ok, {reply, Msg}}, State) ->
CallId = maps:get(<<"call_id">>, Msg, <<>>),
Result = maps:get(<<"result">>, Msg, maps:get(<<"error">>, Msg, #{})),
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined ->
?LOG_DEBUG("[relay_handler] Reply for unknown call_id ~s", [CallId]);
CallerPid ->
CallerPid ! {relay_reply, CallId, Result}
end,
State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)};
%% PING — respond with PONG
handle_message({ok, {ping, Msg}}, State) ->
Timestamp = maps:get(<<"timestamp">>, Msg, erlang:system_time(millisecond)),
send_to_node(pong, #{
<<"timestamp">> => Timestamp,
<<"server_time">> => erlang:system_time(millisecond)
}, State),
State;
%% Catch-all for messages we don't handle
handle_message({ok, {Type, _Msg}}, State) ->
?LOG_DEBUG("[relay_handler] Ignoring message type: ~p", [Type]),
State;
handle_message({error, Reason}, State) ->
?LOG_WARNING("[relay_handler] Decode error: ~p", [Reason]),
State.
%%====================================================================
%% Send helper
%%====================================================================
%% Replay messages that arrived before CONNECT.
replay_pending(#state{pending_msgs = []} = State) -> State;
replay_pending(#state{pending_msgs = Msgs} = State) ->
?LOG_INFO("[relay_handler] Replaying ~b buffered message(s)", [length(Msgs)]),
lists:foldl(fun(Msg, S) -> handle_message(Msg, S) end,
State#state{pending_msgs = []},
lists:reverse(Msgs)).
%% Register this handler as a connected client node (not peers).
%% Joins both the global group AND the target-relay-specific group
%% for multi-tenant relay support.
register_client_node(true, _NodeName, _Identity, _TargetRelay) -> ok;
register_client_node(false, NodeName, Identity, TargetRelay) ->
put(relay_node_name, NodeName),
put(relay_node_identity, Identity),
put(relay_target, TargetRelay),
pg:join(pg, relay_connected_nodes, self()),
%% Also join target-specific group for multi-tenant node counting
case TargetRelay of
<<>> -> ok;
_ -> pg:join(pg, {relay_connected_nodes, TargetRelay}, self())
end.
%% Join topic pg groups. Clients join both groups; peers only relay_topic.
join_topic(Topic, false) ->
pg:join(pg, {relay_topic, Topic}, self()),
pg:join(pg, {relay_local, Topic}, self()),
notify_peering(topic_subscribed, Topic),
?LOG_INFO("[relay_handler] Subscribed to ~s (peer=false)", [Topic]);
join_topic(Topic, true) ->
pg:join(pg, {relay_topic, Topic}, self()),
?LOG_INFO("[relay_handler] Subscribed to ~s (peer=true)", [Topic]).
%% Forward RPC call to peer relays when procedure not found locally.
%% Spawns async — tries each peer, sends first success back to caller.
forward_rpc_to_peers(Procedure, CallId, Args, _State) ->
HandlerPid = self(),
spawn(fun() ->
case erlang:whereis(macula_relay_peering) of
undefined ->
HandlerPid ! {relay_reply, CallId,
#{<<"error">> => #{<<"code">> => <<"procedure_not_found">>,
<<"message">> => <<"No node has registered this procedure">>}}};
PeeringPid ->
Clients = gen_server:call(PeeringPid, peer_clients, 2000),
Result = try_rpc_on_peers(Procedure, Args, Clients),
case Result of
{ok, Response} ->
HandlerPid ! {relay_reply, CallId, Response};
{error, _} ->
HandlerPid ! {relay_reply, CallId,
#{<<"error">> => #{<<"code">> => <<"procedure_not_found">>,
<<"message">> => <<"No provider on any relay">>}}}
end
end
end).
try_rpc_on_peers(_Procedure, _Args, []) ->
{error, not_found};
try_rpc_on_peers(Procedure, Args, [{_Url, ClientPid} | Rest]) ->
case catch macula_relay_client:call(ClientPid, Procedure, Args, 4000) of
{ok, Response} -> {ok, Response};
_ -> try_rpc_on_peers(Procedure, Args, Rest)
end.
%% Notify relay peering module (if running) about subscription changes.
notify_peering(Event, Topic) ->
case erlang:whereis(macula_relay_peering) of
undefined -> ok;
Pid -> gen_server:cast(Pid, {Event, Topic})
end.
send_to_node(Type, Msg, #state{stream = Stream}) ->
Binary = macula_protocol_encoder:encode(Type, Msg),
case macula_quic:async_send(Stream, Binary) of
ok -> ok;
{error, Reason} ->
?LOG_WARNING("[relay_handler] Send ~p failed: ~p", [Type, Reason])
end.