Packages
macula
0.30.5
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)
%%%
%%% 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(),
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}) ->
%% Stream already accepted by relay via async_accept_stream
quicer:setopt(Stream, active, true),
?LOG_INFO("[relay_handler] Started with stream"),
{ok, #state{conn = Conn, stream = Stream, node_id = <<>>,
recv_buffer = <<>>, pending_calls = #{}}}.
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%%====================================================================
%% QUIC data received
%%====================================================================
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, <<>>),
?LOG_INFO("[relay_handler] Node connected: realm=~s node=~s",
[Realm, binary:encode_hex(NodeId)]),
%% Send PONG as handshake acknowledgment
send_to_node(pong, #{
<<"timestamp">> => erlang:system_time(millisecond),
<<"server_time">> => erlang:system_time(millisecond)
}, State),
State#state{node_id = NodeId};
%% SUBSCRIBE — join pg groups for topics
handle_message({ok, {subscribe, Msg}}, State) ->
Topics = maps:get(<<"topics">>, Msg, []),
lists:foreach(fun(Topic) ->
pg:join(pg, {relay_topic, Topic}, self()),
?LOG_INFO("[relay_handler] Subscribed to ~s", [Topic])
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())
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.