Packages
macula
0.4.2
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_client_client.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Macula SDK client connection manager.
%%%
%%% Manages individual HTTP/3 (QUIC) connections to Macula mesh nodes.
%%% Handles connection lifecycle, authentication, message sending/receiving,
%%% and subscription management.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_client_client).
-behaviour(gen_server).
%% API
-export([
start_link/2,
stop/1,
publish/3,
publish/4,
subscribe/3,
unsubscribe/2,
call/3,
call/4
]).
%% 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").
-record(state, {
url :: binary(),
opts :: map(),
realm :: binary(),
node_id :: binary(),
%% QUIC connection and stream
connection :: undefined | pid(),
stream :: undefined | pid(),
%% Connection state
status :: connecting | connected | disconnected | error,
%% Subscriptions: #{SubscriptionRef => {Topic, Callback}}
subscriptions :: #{reference() => {binary(), fun((map()) -> ok)}},
%% Pending RPC calls: #{CallId => {From, Timer}}
pending_calls :: #{binary() => {term(), reference()}},
%% Message ID counter
msg_id_counter :: non_neg_integer(),
%% Receive buffer for partial messages
recv_buffer :: binary()
}).
-define(DEFAULT_TIMEOUT, 5000).
-define(CALL_TIMEOUT, 30000).
-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.
-spec publish(pid(), binary(), map() | binary(), map()) -> ok | {error, term()}.
publish(Client, Topic, Data, Opts) ->
gen_server:call(Client, {publish, Topic, Data, Opts}, ?DEFAULT_TIMEOUT).
%% @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 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).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Url, Opts}) ->
%% Parse URL to extract host and port
{Host, Port} = parse_url(Url),
%% Get realm (required)
Realm = case maps:get(realm, Opts, undefined) of
undefined -> error({missing_required_option, realm});
R when is_binary(R) -> R;
R when is_list(R) -> list_to_binary(R);
R when is_atom(R) -> atom_to_binary(R)
end,
%% Generate or get node ID
NodeId = maps:get(node_id, Opts, generate_node_id()),
State = #state{
url = Url,
opts = Opts#{host => Host, port => Port},
realm = Realm,
node_id = NodeId,
connection = undefined,
stream = undefined,
status = connecting,
subscriptions = #{},
pending_calls = #{},
msg_id_counter = 0,
recv_buffer = <<>>
},
%% Initiate connection asynchronously
self() ! connect,
{ok, State}.
%% @private
handle_call({publish, Topic, Data, Opts}, _From, #state{status = connected} = State) ->
%% Build publish message
Qos = maps:get(qos, Opts, 0),
Retain = maps:get(retain, Opts, false),
{MsgId, State2} = next_message_id(State),
%% Encode data to binary if it's a map
Payload = case Data of
D when is_binary(D) -> D;
D when is_map(D) -> encode_json(D);
D when is_list(D) -> list_to_binary(D)
end,
PublishMsg = #{
topic => ensure_binary(Topic),
payload => Payload,
qos => Qos,
retain => Retain,
message_id => MsgId
},
%% Encode and send
case send_message(publish, PublishMsg, State2) of
ok ->
{reply, ok, State2};
{error, Reason} ->
{reply, {error, Reason}, State2}
end;
handle_call({publish, _Topic, _Data, _Opts}, _From, State) ->
{reply, {error, not_connected}, State};
handle_call({subscribe, Topic, Callback}, _From, #state{status = connected} = State) ->
%% Generate subscription reference
SubRef = make_ref(),
%% Send subscribe message
SubscribeMsg = #{
topics => [ensure_binary(Topic)],
qos => 0
},
case send_message(subscribe, SubscribeMsg, State) of
ok ->
%% Store subscription
Subscriptions = maps:put(SubRef, {ensure_binary(Topic), Callback},
State#state.subscriptions),
State2 = State#state{subscriptions = Subscriptions},
{reply, {ok, SubRef}, State2};
{error, Reason} ->
{reply, {error, Reason}, State}
end;
handle_call({subscribe, _Topic, _Callback}, _From, State) ->
{reply, {error, not_connected}, State};
handle_call({unsubscribe, SubRef}, _From, #state{status = connected} = State) ->
case maps:get(SubRef, State#state.subscriptions, undefined) of
undefined ->
{reply, {error, not_subscribed}, State};
{Topic, _Callback} ->
%% Send unsubscribe message
UnsubscribeMsg = #{
topics => [Topic]
},
case send_message(unsubscribe, UnsubscribeMsg, State) of
ok ->
%% Remove subscription
Subscriptions = maps:remove(SubRef, State#state.subscriptions),
State2 = State#state{subscriptions = Subscriptions},
{reply, ok, State2};
{error, Reason} ->
{reply, {error, Reason}, State}
end
end;
handle_call({unsubscribe, _SubRef}, _From, State) ->
{reply, {error, not_connected}, State};
handle_call({call, Procedure, Args, Opts}, From, #state{status = connected} = State) ->
%% Generate call ID
{CallId, State2} = next_message_id(State),
%% Encode args
EncodedArgs = case Args of
A when is_binary(A) -> A;
A when is_map(A) -> encode_json(A);
A when is_list(A) -> encode_json(A)
end,
%% Build call message (RPC call)
CallMsg = #{
procedure => ensure_binary(Procedure),
args => EncodedArgs,
call_id => CallId
},
case send_message(call, CallMsg, State2) of
ok ->
%% Set up timeout timer
Timeout = maps:get(timeout, Opts, ?CALL_TIMEOUT),
Timer = erlang:send_after(Timeout, self(), {call_timeout, CallId}),
%% Store pending call
PendingCalls = maps:put(CallId, {From, Timer}, State2#state.pending_calls),
State3 = State2#state{pending_calls = PendingCalls},
{noreply, State3};
{error, Reason} ->
{reply, {error, Reason}, State2}
end;
handle_call({call, _Procedure, _Args, _Opts}, _From, State) ->
{reply, {error, not_connected}, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info(connect, State) ->
case do_connect(State) of
{ok, State2} ->
{noreply, State2};
{error, Reason} ->
?LOG_ERROR("Connection failed: ~p, retrying in ~p ms",
[Reason, ?CONNECT_RETRY_DELAY]),
erlang:send_after(?CONNECT_RETRY_DELAY, self(), connect),
{noreply, State#state{status = error}}
end;
handle_info({quic, Data, Stream, _Props}, #state{stream = Stream} = State) ->
%% Received data from QUIC stream
handle_received_data(Data, State);
handle_info({call_timeout, CallId}, State) ->
%% RPC call timed out
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined ->
{noreply, State};
{From, _Timer} ->
gen_server:reply(From, {error, timeout}),
PendingCalls = maps:remove(CallId, State#state.pending_calls),
{noreply, State#state{pending_calls = PendingCalls}}
end;
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, #state{stream = Stream, connection = Conn}) ->
%% Clean up QUIC resources
catch macula_quic:close(Stream),
catch macula_quic:close(Conn),
ok.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Establish QUIC connection and perform handshake.
-spec do_connect(#state{}) -> {ok, #state{}} | {error, term()}.
do_connect(State) ->
#{host := Host, port := Port} = State#state.opts,
%% Connect via QUIC
QuicOpts = [
{alpn, ["macula"]},
{verify, none} %% TODO: Add proper TLS verification
],
ConnectResult = try
macula_quic:connect(Host, Port, QuicOpts, ?DEFAULT_TIMEOUT)
catch
_:Error ->
{error, Error}
end,
case ConnectResult of
{ok, Conn} ->
%% Open bidirectional stream
case macula_quic:open_stream(Conn) of
{ok, Stream} ->
%% Send CONNECT message
ConnectMsg = #{
version => <<"1.0">>,
node_id => State#state.node_id,
realm_id => State#state.realm,
capabilities => [pubsub, rpc]
},
case send_message_raw(connect, ConnectMsg, Stream) of
ok ->
?LOG_INFO("Connected to Macula mesh: ~s:~p", [Host, Port]),
{ok, State#state{
connection = Conn,
stream = Stream,
status = connected
}};
{error, Reason} ->
macula_quic:close(Stream),
macula_quic:close(Conn),
{error, {handshake_failed, Reason}}
end;
{error, Reason} ->
macula_quic:close(Conn),
{error, {stream_open_failed, Reason}}
end;
{error, Reason} ->
{error, {connection_failed, Reason}};
{error, Type, Details} ->
{error, {connection_failed, {Type, Details}}};
Other ->
{error, {connection_failed, Other}}
end.
%% @doc Send a protocol message through the established stream.
-spec send_message(atom(), map(), #state{}) -> ok | {error, term()}.
send_message(Type, Msg, #state{stream = Stream}) ->
send_message_raw(Type, Msg, Stream).
%% @doc Send a protocol message through a stream (raw).
-spec send_message_raw(atom(), map(), pid()) -> ok | {error, term()}.
send_message_raw(Type, Msg, Stream) ->
try
Binary = macula_protocol_encoder:encode(Type, Msg),
macula_quic:send(Stream, Binary)
catch
error:Reason ->
{error, {encode_error, Reason}}
end.
%% @doc Handle received data from the stream.
-spec handle_received_data(binary(), #state{}) -> {noreply, #state{}}.
handle_received_data(Data, State) ->
%% Append to receive buffer
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
%% Try to decode messages
{Messages, RemainingBuffer} = decode_messages(Buffer, []),
%% Process each message
State2 = lists:foldl(fun process_message/2, State, Messages),
{noreply, State2#state{recv_buffer = RemainingBuffer}}.
%% @doc Decode all complete messages from buffer.
-spec decode_messages(binary(), list()) -> {list(), binary()}.
decode_messages(Buffer, Acc) when byte_size(Buffer) < 8 ->
%% Not enough for header
{lists:reverse(Acc), Buffer};
decode_messages(<<_Version:8, _TypeId:8, _Flags:8, _Reserved:8,
PayloadLen:32/big-unsigned, Rest/binary>> = Buffer, Acc) ->
case byte_size(Rest) of
ActualLen when ActualLen >= PayloadLen ->
%% We have a complete message
case macula_protocol_decoder:decode(Buffer) of
{ok, {Type, Msg}} ->
%% Skip this message and continue
<<_:8/binary, _Payload:PayloadLen/binary, Remaining/binary>> = Buffer,
decode_messages(Remaining, [{Type, Msg} | Acc]);
{error, _Reason} ->
%% Decode error, skip this message
<<_:8/binary, _Payload:PayloadLen/binary, Remaining/binary>> = Buffer,
decode_messages(Remaining, Acc)
end;
_ ->
%% Incomplete message, wait for more data
{lists:reverse(Acc), Buffer}
end.
%% @doc Process a received message.
-spec process_message({atom(), map()}, #state{}) -> #state{}.
process_message({publish, Msg}, State) ->
%% Handle incoming publish (for subscriptions)
#{topic := Topic, payload := Payload} = Msg,
%% Find matching subscriptions and invoke callbacks
maps:fold(fun(_SubRef, {SubTopic, Callback}, St) ->
case SubTopic of
Topic ->
%% Decode payload
Event = decode_json(Payload),
%% Invoke callback (spawn to avoid blocking)
spawn(fun() -> Callback(Event) end);
_ ->
ok
end,
St
end, State, State#state.subscriptions);
process_message({reply, Msg}, State) ->
%% Handle RPC reply
#{call_id := CallId, result := Result} = Msg,
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined ->
State;
{From, Timer} ->
erlang:cancel_timer(Timer),
DecodedResult = decode_json(Result),
gen_server:reply(From, {ok, DecodedResult}),
State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)}
end;
process_message({pong, _Msg}, State) ->
%% TODO: Handle ping/pong for keepalive
State;
process_message(_OtherMsg, State) ->
%% Ignore unknown messages
State.
%% @doc Parse URL to extract host and port.
-spec parse_url(binary()) -> {string(), inet:port_number()}.
parse_url(Url) when is_binary(Url) ->
parse_url(binary_to_list(Url));
parse_url("https://" ++ Rest) ->
parse_host_port(Rest, 443);
parse_url("http://" ++ Rest) ->
parse_host_port(Rest, 80);
parse_url(Rest) ->
parse_host_port(Rest, 443).
parse_host_port(HostPort, DefaultPort) ->
case string:split(HostPort, ":") of
[Host] ->
{Host, DefaultPort};
[Host, PortStr] ->
Port = list_to_integer(string:trim(PortStr, trailing, "/")),
{Host, Port}
end.
%% @doc Generate a random node ID.
-spec generate_node_id() -> binary().
generate_node_id() ->
%% 32-byte random ID
crypto:strong_rand_bytes(32).
%% @doc Get next message ID.
-spec next_message_id(#state{}) -> {binary(), #state{}}.
next_message_id(State) ->
Counter = State#state.msg_id_counter + 1,
%% 16-byte message ID from counter
MsgId = <<Counter:128>>,
{MsgId, State#state{msg_id_counter = Counter}}.
%% @doc Ensure value is binary.
-spec ensure_binary(binary() | list() | atom()) -> binary().
ensure_binary(B) when is_binary(B) -> B;
ensure_binary(L) when is_list(L) -> list_to_binary(L);
ensure_binary(A) when is_atom(A) -> atom_to_binary(A).
%% @doc Encode map/list to JSON binary.
-spec encode_json(map() | list()) -> binary().
encode_json(Data) ->
json:encode(Data).
%% @doc Decode JSON binary to map/list.
-spec decode_json(binary()) -> map() | list().
decode_json(Binary) ->
json:decode(Binary).