Packages
macula
0.31.7
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_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)),
Label = case IsPeer of true -> <<"Peer relay">>; false -> <<"Node">> end,
?LOG_INFO("[relay_handler] ~s connected: realm=~s name=~s", [Label, Realm, NodeName]),
register_client_node(IsPeer, NodeName),
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 — 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
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_WARNING("[relay_handler] No provider for RPC ~s", [Procedure]),
send_to_node(reply, #{
<<"call_id">> => CallId,
<<"error">> => #{
<<"code">> => <<"procedure_not_found">>,
<<"message">> => <<"No node has registered this procedure">>
}
}, 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).
register_client_node(true, _NodeName) -> ok;
register_client_node(false, NodeName) ->
put(relay_node_name, NodeName),
pg:join(pg, relay_connected_nodes, self()).
%% 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]).
%% 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.