Packages
macula
0.31.3
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(),
is_peer :: boolean(), %% true if this is a relay-to-relay peering connection
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 = <<>>,
is_peer = false, 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">>,
?LOG_INFO("[relay_handler] ~s connected: realm=~s node=~s",
[case IsPeer of true -> <<"Peer relay">>; false -> <<"Node">> end,
Realm, binary:encode_hex(NodeId)]),
send_to_node(pong, #{
<<"timestamp">> => erlang:system_time(millisecond),
<<"server_time">> => erlang:system_time(millisecond)
}, State),
State#state{node_id = NodeId, is_peer = IsPeer};
%% 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) ->
pg:join(pg, {relay_topic, Topic}, self()),
case IsPeer of
false -> pg:join(pg, {relay_local, Topic}, self());
true -> ok
end,
?LOG_INFO("[relay_handler] Subscribed to ~s (peer=~p)", [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
%%====================================================================
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.