Packages

macula

0.7.23
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_clients.erl
Raw

src/macula_gateway_clients.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Clients Worker GenServer - tracks connected clients.
%%%
%%% Responsibilities:
%%% - Track connected clients with metadata (BOUNDED POOL)
%%% - Enforce max_clients limit with backpressure
%%% - Monitor client processes for automatic cleanup
%%% - Store bidirectional streams for client communication
%%% - Provide client query APIs
%%%
%%% Pattern: Bounded client pool with backpressure
%%% - Tracks clients with max_clients limit (default: 10,000)
%%% - Rejects new clients when pool is full (backpressure)
%%% - Allows updates to existing clients even when pool is full
%%%
%%% Configuration:
%%% - max_clients: Maximum concurrent clients (default: 10,000)
%%%
%%% Extracted from macula_gateway.erl (Phase 2)
%%% Renamed from macula_gateway_client_manager (Phase 2 QUIC refactoring)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_clients).
-behaviour(gen_server).
%% API
-export([
start_link/1,
stop/1,
client_connected/3,
client_disconnected/2,
get_client_info/2,
get_all_clients/1,
is_client_alive/2,
store_client_stream/3,
get_client_stream/2,
get_stream_by_endpoint/2,
get_all_node_ids/1
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-type client_info() :: #{
realm := binary(),
node_id := binary(),
endpoint => binary(),
capabilities => [atom()]
}.
-record(state, {
opts :: map(),
max_clients :: integer(), % Maximum clients allowed
clients :: #{pid() => client_info()}, % client_pid => client_info
monitors :: #{reference() => pid()}, % monitor_ref => client_pid
client_streams :: #{binary() => pid()}, % node_id => stream_pid
endpoint_to_stream :: #{binary() => pid()} % endpoint => stream_pid
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the client manager with options.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Stop the client manager.
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Register a connected client with metadata.
%% Monitors the client process for automatic cleanup on death.
-spec client_connected(pid(), pid(), client_info()) -> ok.
client_connected(Pid, ClientPid, ClientInfo) ->
gen_server:call(Pid, {client_connected, ClientPid, ClientInfo}).
%% @doc Unregister a disconnected client.
-spec client_disconnected(pid(), pid()) -> ok.
client_disconnected(Pid, ClientPid) ->
gen_server:call(Pid, {client_disconnected, ClientPid}).
%% @doc Get information about a specific client.
-spec get_client_info(pid(), pid()) -> {ok, client_info()} | not_found.
get_client_info(Pid, ClientPid) ->
gen_server:call(Pid, {get_client_info, ClientPid}).
%% @doc Get all connected clients.
-spec get_all_clients(pid()) -> {ok, [{pid(), client_info()}]}.
get_all_clients(Pid) ->
gen_server:call(Pid, get_all_clients).
%% @doc Check if a client is alive (process still running).
-spec is_client_alive(pid(), pid()) -> boolean().
is_client_alive(_Pid, ClientPid) ->
erlang:is_process_alive(ClientPid).
%% @doc Store a bidirectional stream for a client node (legacy 3-arg version).
-spec store_client_stream(pid(), binary(), pid()) -> ok.
store_client_stream(Pid, NodeId, StreamPid) ->
store_client_stream(Pid, NodeId, StreamPid, <<>>).
%% @doc Store a bidirectional stream for a client node with endpoint tracking.
-spec store_client_stream(pid(), binary(), pid(), binary()) -> ok.
store_client_stream(Pid, NodeId, StreamPid, Endpoint) when is_binary(Endpoint) ->
gen_server:call(Pid, {store_client_stream, NodeId, StreamPid, Endpoint}).
%% @doc Get the stored stream for a client node.
-spec get_client_stream(pid(), binary()) -> {ok, pid()} | not_found.
get_client_stream(Pid, NodeId) ->
gen_server:call(Pid, {get_client_stream, NodeId}).
%% @doc Get the stream PID for a given endpoint URL.
%% Used for routing pub/sub messages to remote subscribers.
-spec get_stream_by_endpoint(pid(), binary()) -> {ok, pid()} | {error, not_found}.
get_stream_by_endpoint(Pid, Endpoint) ->
gen_server:call(Pid, {get_stream_by_endpoint, Endpoint}).
%% @doc Get all node IDs with stored client streams (for debugging).
-spec get_all_node_ids(pid()) -> [binary()].
get_all_node_ids(Pid) ->
gen_server:call(Pid, get_all_node_ids).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
io:format("[Clients] Initializing client manager~n"),
%% Get max clients from opts or use default (10,000)
MaxClients = maps:get(max_clients, Opts, 10000),
io:format("[Clients] Max clients: ~p~n", [MaxClients]),
State = #state{
opts = Opts,
max_clients = MaxClients,
clients = #{},
monitors = #{},
client_streams = #{},
endpoint_to_stream = #{}
},
io:format("[Clients] Client manager initialized~n"),
{ok, State}.
handle_call({client_connected, ClientPid, ClientInfo}, _From,
#state{clients = Clients, max_clients = MaxClients} = State) ->
%% Check if client already connected (update case)
case maps:is_key(ClientPid, Clients) of
true ->
%% Update existing client info (allowed even when pool is full)
NewClients = maps:put(ClientPid, ClientInfo, Clients),
NewState = State#state{clients = NewClients},
{reply, ok, NewState};
false ->
%% New client - check if pool is full
case maps:size(Clients) >= MaxClients of
true ->
%% Pool full - reject new client (backpressure)
io:format("[ClientManager] Client pool full (~p clients), rejecting new client~n",
[MaxClients]),
{reply, {error, max_clients_reached}, State};
false ->
%% Pool has space - monitor and store new client
MonitorRef = erlang:monitor(process, ClientPid),
NewClients = maps:put(ClientPid, ClientInfo, Clients),
NewMonitors = maps:put(MonitorRef, ClientPid, State#state.monitors),
NewState = State#state{
clients = NewClients,
monitors = NewMonitors
},
{reply, ok, NewState}
end
end;
handle_call({client_disconnected, ClientPid}, _From, State) ->
NewState = remove_client(ClientPid, State),
{reply, ok, NewState};
handle_call({get_client_info, ClientPid}, _From, State) ->
Result = case maps:get(ClientPid, State#state.clients, undefined) of
undefined -> not_found;
Info -> {ok, Info}
end,
{reply, Result, State};
handle_call(get_all_clients, _From, State) ->
Clients = maps:to_list(State#state.clients),
{reply, {ok, Clients}, State};
%% Legacy 3-arg version (no endpoint)
handle_call({store_client_stream, NodeId, StreamPid, <<>>}, _From, State) ->
ClientStreams = maps:put(NodeId, StreamPid, State#state.client_streams),
NewState = State#state{client_streams = ClientStreams},
{reply, ok, NewState};
%% New 4-arg version with endpoint tracking
handle_call({store_client_stream, NodeId, StreamPid, Endpoint}, _From, State) when is_binary(Endpoint), byte_size(Endpoint) > 0 ->
%% Store both node_id → stream and endpoint → stream mappings
ClientStreams = maps:put(NodeId, StreamPid, State#state.client_streams),
EndpointToStream = maps:put(Endpoint, StreamPid, State#state.endpoint_to_stream),
NewState = State#state{
client_streams = ClientStreams,
endpoint_to_stream = EndpointToStream
},
io:format("[ClientManager] Tracking endpoint → stream: ~s → ~p~n", [Endpoint, StreamPid]),
{reply, ok, NewState};
handle_call({get_client_stream, NodeId}, _From, State) ->
Result = case maps:get(NodeId, State#state.client_streams, undefined) of
undefined -> not_found;
StreamPid -> {ok, StreamPid}
end,
{reply, Result, State};
handle_call({get_stream_by_endpoint, Endpoint}, _From, State) ->
Result = case maps:get(Endpoint, State#state.endpoint_to_stream, undefined) of
undefined -> {error, not_found};
StreamPid -> {ok, StreamPid}
end,
{reply, Result, State};
handle_call(get_all_node_ids, _From, State) ->
NodeIds = maps:keys(State#state.client_streams),
{reply, NodeIds, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle client process death - automatic cleanup.
handle_info({'DOWN', MonitorRef, process, ClientPid, _Reason}, State) ->
%% Remove monitor reference
Monitors = maps:remove(MonitorRef, State#state.monitors),
%% Remove client from registry
NewState = remove_client(ClientPid, State#state{monitors = Monitors}),
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Remove a client from the registry.
%% Does not demonitor - that's handled separately in handle_info.
%% Also removes associated client stream from client_streams and endpoint_to_stream maps.
-spec remove_client(pid(), #state{}) -> #state{}.
remove_client(ClientPid, State) ->
%% Get client info to extract node_id and endpoint before removing
case maps:get(ClientPid, State#state.clients, undefined) of
undefined ->
%% Client not found, just return state
State;
ClientInfo ->
%% Extract node_id and remove from client_streams
NodeId = maps:get(node_id, ClientInfo),
NewClientStreams = maps:remove(NodeId, State#state.client_streams),
%% Extract endpoint (if present) and remove from endpoint_to_stream
Endpoint = maps:get(endpoint, ClientInfo, undefined),
NewEndpointToStream = case Endpoint of
undefined -> State#state.endpoint_to_stream;
<<>> -> State#state.endpoint_to_stream;
_ ->
io:format("[ClientManager] Cleaning up endpoint mapping: ~s~n", [Endpoint]),
maps:remove(Endpoint, State#state.endpoint_to_stream)
end,
%% Remove from clients map
NewClients = maps:remove(ClientPid, State#state.clients),
State#state{
clients = NewClients,
client_streams = NewClientStreams,
endpoint_to_stream = NewEndpointToStream
}
end.