Packages

macula

0.4.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
macula src macula_gateway.erl
Raw

src/macula_gateway.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Macula Gateway - HTTP/3 Message Router
%%%
%%% Main API module for the Macula Gateway.
%%% The gateway can be embedded in applications or run standalone.
%%%
%%% Architecture:
%%% - QUIC Listener: Accepts HTTP/3 connections from SDK clients
%%% - Router: Routes pub/sub messages between clients
%%% - RPC: Handles remote procedure calls
%%% - Realm Manager: Manages multiple realms
%%%
%%% Usage (Embedded):
%%% ```
%%% {ok, Pid} = macula_gateway:start_link([
%%% {port, 9443},
%%% {realm, <<"be.cortexiq.energy">>}
%%% ]).
%%% '''
%%%
%%% Usage (Standalone):
%%% ```
%%% application:start(macula_gateway).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway).
-behaviour(gen_server).
%% API
-export([
start_link/0,
start_link/1,
stop/1,
get_stats/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-record(state, {
port :: inet:port_number(),
realm :: binary(),
listener :: pid() | undefined,
clients :: #{pid() => client_info()},
subscriptions :: #{binary() => [pid()]}, % topic => [client_pids]
registrations :: #{binary() => pid()} % procedure => client_pid
}).
-type client_info() :: #{
realm := binary(),
node_id := binary(),
capabilities := [atom()]
}.
-define(DEFAULT_PORT, 9443).
-define(DEFAULT_REALM, <<"macula.default">>).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start the gateway with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link([]).
%% @doc Start the gateway with custom options.
%% Options:
%% {port, Port} - Listen port (default: 9443)
%% {realm, Realm} - Default realm (default: "macula.default")
-spec start_link(proplists:proplist()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Stop the gateway.
-spec stop(pid()) -> ok.
stop(Gateway) ->
gen_server:stop(Gateway).
%% @doc Get gateway statistics.
-spec get_stats(pid()) -> map().
get_stats(Gateway) ->
gen_server:call(Gateway, get_stats).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
Port = proplists:get_value(port, Opts, ?DEFAULT_PORT),
Realm = proplists:get_value(realm, Opts, ?DEFAULT_REALM),
%% Get TLS certificates
%% Priority: 1. Environment variables (production)
%% 2. Pre-generated certs (development)
{CertFile, KeyFile} = get_tls_certificates(),
io:format("Using TLS certificates:~n"),
io:format(" Cert: ~s~n", [CertFile]),
io:format(" Key: ~s~n", [KeyFile]),
%% Validate certificate files exist
case macula_quic_cert:validate_files(CertFile, KeyFile) of
ok ->
start_quic_listener(Port, Realm, CertFile, KeyFile);
{error, Reason} ->
io:format("Certificate validation failed: ~p~n", [Reason]),
{stop, {cert_validation_failed, Reason}}
end.
%% @private
%% @doc Get TLS certificate paths from environment or use pre-generated ones.
get_tls_certificates() ->
case {os:getenv("TLS_CERT_FILE"), os:getenv("TLS_KEY_FILE")} of
{false, false} ->
%% No env vars, use pre-generated certs
io:format("Using pre-generated TLS certificates~n"),
{"/opt/macula/certs/cert.pem", "/opt/macula/certs/key.pem"};
{CertEnv, KeyEnv} when CertEnv =/= false andalso KeyEnv =/= false ->
%% Use mounted certificates (production)
io:format("Using mounted TLS certificates from environment~n"),
{CertEnv, KeyEnv};
_ ->
%% Partial configuration, log warning and use defaults
io:format("WARNING: Partial TLS environment config, using pre-generated certs~n"),
{"/opt/macula/certs/cert.pem", "/opt/macula/certs/key.pem"}
end.
%% @private
%% @doc Start the QUIC listener with given certificates.
start_quic_listener(Port, Realm, CertFile, KeyFile) ->
%% Start QUIC listener
ListenOpts = [
{cert, CertFile},
{key, KeyFile},
{alpn, ["macula"]},
{peer_unidi_stream_count, 3}
],
case macula_quic:listen(Port, ListenOpts) of
{ok, Listener} ->
io:format("Macula Gateway listening on port ~p (realm: ~s)~n", [Port, Realm]),
%% Mark health server as ready
try
macula_gateway_health:set_ready(true)
catch
_:_ -> ok % Health server might not be running in embedded mode
end,
%% Register diagnostics procedures
try
macula_gateway_diagnostics:register_procedures(self())
catch
_:_ -> ok % Diagnostics service might not be running in embedded mode
end,
%% Start accepting connections
self() ! accept,
State = #state{
port = Port,
realm = Realm,
listener = Listener,
clients = #{},
subscriptions = #{},
registrations = #{}
},
{ok, State};
{error, Reason} ->
io:format("QUIC listen failed: ~p~n", [Reason]),
{stop, {listen_failed, Reason}};
{error, Type, Details} ->
io:format("QUIC listen failed: ~p ~p~n", [Type, Details]),
{stop, {listen_failed, {Type, Details}}};
Other ->
io:format("QUIC listen unexpected result: ~p~n", [Other]),
{stop, {listen_failed, Other}}
end.
handle_call(get_stats, _From, State) ->
Stats = #{
port => State#state.port,
realm => State#state.realm,
clients => maps:size(State#state.clients),
subscriptions => maps:size(State#state.subscriptions),
registrations => maps:size(State#state.registrations)
},
{reply, Stats, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Request, State) ->
{noreply, State}.
handle_info(accept, #state{listener = Listener} = State) ->
%% Accept incoming connection
case macula_quic:accept(Listener, 5000) of
{ok, Conn} ->
%% Spawn handler for this connection
spawn_link(fun() -> handle_client_connection(self(), Conn, State) end),
%% Continue accepting
self() ! accept,
{noreply, State};
{error, timeout} ->
%% No connection, keep accepting
self() ! accept,
{noreply, State};
{error, Reason} ->
io:format("Accept error: ~p~n", [Reason]),
{stop, {accept_error, Reason}, State}
end;
%% Client registered
handle_info({client_connected, ClientPid, ClientInfo}, State) ->
io:format("Client connected: ~p~n", [ClientInfo]),
%% Monitor client
erlang:monitor(process, ClientPid),
Clients = maps:put(ClientPid, ClientInfo, State#state.clients),
{noreply, State#state{clients = Clients}};
%% Client disconnected
handle_info({'DOWN', _Ref, process, ClientPid, _Reason}, State) ->
io:format("Client disconnected: ~p~n", [ClientPid]),
%% Remove client from all subscriptions
Subscriptions = maps:map(fun(_Topic, Subscribers) ->
lists:delete(ClientPid, Subscribers)
end, State#state.subscriptions),
%% Remove client registrations
Registrations = maps:filter(fun(_Proc, Pid) ->
Pid =/= ClientPid
end, State#state.registrations),
%% Remove client
Clients = maps:remove(ClientPid, State#state.clients),
{noreply, State#state{
clients = Clients,
subscriptions = Subscriptions,
registrations = Registrations
}};
%% Publish message (from client)
handle_info({publish, FromPid, Topic, Payload}, State) ->
%% Get all subscribers to this topic
Subscribers = maps:get(Topic, State#state.subscriptions, []),
%% Send to all subscribers (except sender)
[SubPid ! {event, Topic, Payload} || SubPid <- Subscribers, SubPid =/= FromPid],
{noreply, State};
%% Subscribe request
handle_info({subscribe, ClientPid, Topic}, State) ->
%% Add client to topic subscribers
Subscribers = maps:get(Topic, State#state.subscriptions, []),
NewSubscribers = [ClientPid | Subscribers],
Subscriptions = maps:put(Topic, NewSubscribers, State#state.subscriptions),
%% Acknowledge subscription
ClientPid ! {subscribed, Topic},
{noreply, State#state{subscriptions = Subscriptions}};
%% Unsubscribe request
handle_info({unsubscribe, ClientPid, Topic}, State) ->
%% Remove client from topic subscribers
Subscribers = maps:get(Topic, State#state.subscriptions, []),
NewSubscribers = lists:delete(ClientPid, Subscribers),
Subscriptions = maps:put(Topic, NewSubscribers, State#state.subscriptions),
%% Acknowledge unsubscription
ClientPid ! {unsubscribed, Topic},
{noreply, State#state{subscriptions = Subscriptions}};
%% RPC Call request
handle_info({call, FromPid, CallId, Procedure, Args}, State) ->
%% Find registered handler for procedure
case maps:get(Procedure, State#state.registrations, undefined) of
undefined ->
%% No handler registered
FromPid ! {call_error, CallId, <<"wamp.error.no_such_procedure">>};
HandlerPid ->
%% Forward to handler
HandlerPid ! {invoke, FromPid, CallId, Procedure, Args}
end,
{noreply, State};
%% Register procedure
handle_info({register, ClientPid, Procedure}, State) ->
%% Register procedure handler
Registrations = maps:put(Procedure, ClientPid, State#state.registrations),
%% Acknowledge registration
ClientPid ! {registered, Procedure},
{noreply, State#state{registrations = Registrations}};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, #state{listener = Listener}) ->
%% Close listener
case Listener of
undefined -> ok;
_ -> macula_quic:close(Listener)
end,
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Handle a client connection.
handle_client_connection(GatewayPid, Conn, State) ->
accept_client_stream(GatewayPid, Conn, State).
%% @doc Accept stream from client connection.
accept_client_stream(GatewayPid, Conn, State) ->
case macula_quic:accept_stream(Conn, 5000) of
{ok, Stream} ->
receive_connect_message(GatewayPid, Stream, Conn, State);
{error, Reason} ->
io:format("Stream accept error: ~p~n", [Reason]),
macula_quic:close(Conn)
end.
%% @doc Receive and decode CONNECT message from client.
receive_connect_message(GatewayPid, Stream, Conn, State) ->
case macula_quic:recv(Stream, 5000) of
{ok, Data} ->
decode_connect_message(GatewayPid, Stream, Conn, State, Data);
{error, Reason} ->
io:format("Recv error: ~p~n", [Reason]),
macula_quic:close(Stream),
macula_quic:close(Conn)
end.
%% @doc Decode CONNECT message.
decode_connect_message(GatewayPid, Stream, Conn, State, Data) ->
case macula_protocol_decoder:decode(Data) of
{ok, {connect, ConnectMsg}} ->
validate_client_realm(GatewayPid, Stream, Conn, State, ConnectMsg);
{error, Reason} ->
io:format("Decode error: ~p~n", [Reason]),
macula_quic:close(Stream),
macula_quic:close(Conn)
end.
%% @doc Validate client realm and register if valid.
validate_client_realm(GatewayPid, Stream, Conn, State, ConnectMsg) ->
RealmId = maps:get(realm_id, ConnectMsg),
check_realm_match(GatewayPid, Stream, Conn, State, RealmId, ConnectMsg).
%% @doc Check if client realm matches gateway realm.
check_realm_match(GatewayPid, Stream, Conn, State, RealmId, ConnectMsg)
when RealmId =:= State#state.realm ->
%% Register client
ClientInfo = #{
realm => RealmId,
node_id => maps:get(node_id, ConnectMsg),
capabilities => maps:get(capabilities, ConnectMsg, [])
},
GatewayPid ! {client_connected, self(), ClientInfo},
%% Handle client messages
client_message_loop(GatewayPid, Stream, Conn);
check_realm_match(_GatewayPid, Stream, Conn, _State, _OtherRealm, _ConnectMsg) ->
%% Wrong realm, reject
io:format("Client rejected: wrong realm~n"),
macula_quic:close(Stream),
macula_quic:close(Conn).
%% @doc Message loop for connected client.
client_message_loop(GatewayPid, Stream, Conn) ->
receive
%% Event from gateway (subscribed topic)
{event, Topic, Payload} ->
%% Send event to client
EventMsg = #{
topic => Topic,
payload => Payload
},
case macula_protocol_encoder:encode(publish, EventMsg) of
{ok, Data} ->
macula_quic:send(Stream, Data);
{error, Reason} ->
io:format("Encode error: ~p~n", [Reason])
end,
client_message_loop(GatewayPid, Stream, Conn);
%% Subscription confirmed
{subscribed, Topic} ->
io:format("Subscription confirmed: ~s~n", [Topic]),
client_message_loop(GatewayPid, Stream, Conn);
%% Unsubscription confirmed
{unsubscribed, Topic} ->
io:format("Unsubscription confirmed: ~s~n", [Topic]),
client_message_loop(GatewayPid, Stream, Conn);
%% RPC invocation (this client is the handler)
{invoke, CallerPid, CallId, Procedure, Args} ->
%% Send invocation to client
InvokeMsg = #{
call_id => CallId,
procedure => Procedure,
args => Args
},
case macula_protocol_encoder:encode(call, InvokeMsg) of
{ok, Data} ->
macula_quic:send(Stream, Data),
%% Wait for result from client
%% (Simplified - should track pending invocations)
CallerPid ! {call_result, CallId, <<"dummy result">>};
{error, Reason} ->
io:format("Encode error: ~p~n", [Reason]),
CallerPid ! {call_error, CallId, <<"encoding_failed">>}
end,
client_message_loop(GatewayPid, Stream, Conn);
%% Call result
{call_result, _CallId, Result} ->
%% Forward result (simplified)
io:format("Call result: ~p~n", [Result]),
client_message_loop(GatewayPid, Stream, Conn);
%% Call error
{call_error, _CallId, Error} ->
%% Forward error
io:format("Call error: ~p~n", [Error]),
client_message_loop(GatewayPid, Stream, Conn)
after 100 ->
%% Check for incoming messages from client
case macula_quic:recv(Stream, 100) of
{ok, Data} ->
handle_client_message(GatewayPid, Data),
client_message_loop(GatewayPid, Stream, Conn);
{error, timeout} ->
client_message_loop(GatewayPid, Stream, Conn);
{error, Reason} ->
io:format("Client recv error: ~p~n", [Reason]),
macula_quic:close(Stream),
macula_quic:close(Conn)
end
end.
%% @doc Handle message from client.
handle_client_message(GatewayPid, Data) ->
case macula_protocol_decoder:decode(Data) of
{ok, {publish, PubMsg}} ->
Topic = maps:get(topic, PubMsg),
Payload = maps:get(payload, PubMsg),
GatewayPid ! {publish, self(), Topic, Payload};
{ok, {subscribe, SubMsg}} ->
Topics = maps:get(topics, SubMsg, []),
lists:foreach(fun(Topic) ->
GatewayPid ! {subscribe, self(), Topic}
end, Topics);
{ok, {unsubscribe, UnsubMsg}} ->
Topics = maps:get(topics, UnsubMsg, []),
lists:foreach(fun(Topic) ->
GatewayPid ! {unsubscribe, self(), Topic}
end, Topics);
{ok, {call, CallMsg}} ->
CallId = maps:get(call_id, CallMsg),
Procedure = maps:get(procedure, CallMsg),
Args = maps:get(args, CallMsg),
GatewayPid ! {call, self(), CallId, Procedure, Args};
{error, Reason} ->
io:format("Decode error: ~p~n", [Reason])
end.