Packages
macula
0.20.21
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_gateway_system/macula_direct_routing.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Direct Routing Table - Bypass bootstrap for known subscriber endpoints.
%%%
%%% Problem: Every publish goes through bootstrap even for known subscribers.
%%% Solution: Cache {NodeId, Endpoint} mappings and route directly via QUIC.
%%%
%%% Design:
%%% - ETS table for O(1) lookup by NodeId
%%% - TTL-based expiration (default 5 minutes)
%%% - Automatic cleanup of stale entries
%%% - Thread-safe concurrent reads
%%%
%%% Expected improvement:
%%% - 3-5x latency reduction for messages to known subscribers
%%% - Reduced load on bootstrap gateway
%%%
%%% Configuration Options:
%%% - ttl_ms: Route entry TTL (default: 300000ms = 5 minutes)
%%% - cleanup_interval_ms: How often to clean stale entries (default: 60000ms)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_direct_routing).
-behaviour(gen_server).
%% API
-export([
start_link/0,
start_link/1,
lookup/1,
store/2,
store_from_subscriber/1,
remove/1,
clear_all/0,
stats/0
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-define(SERVER, ?MODULE).
-define(TABLE, macula_direct_routing_table).
-define(DEFAULT_TTL_MS, 300000). % 5 minutes default TTL
-define(DEFAULT_CLEANUP_INTERVAL_MS, 60000). % Cleanup every 1 minute
-record(state, {
table :: ets:tid(),
ttl_ms :: pos_integer(),
cleanup_interval_ms :: pos_integer(),
hits :: non_neg_integer(),
misses :: non_neg_integer(),
stores :: non_neg_integer()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the direct routing table with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start the direct routing table with options.
%% Options:
%% - ttl_ms: Route entry TTL in milliseconds (default: 300000)
%% - cleanup_interval_ms: Cleanup interval (default: 60000)
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Look up endpoint for a node ID.
%% Returns {ok, Endpoint} on hit, or miss on cache miss/expired.
-spec lookup(binary()) -> {ok, binary()} | miss.
lookup(NodeId) when is_binary(NodeId) ->
case ets:lookup(?TABLE, NodeId) of
[{NodeId, Endpoint, ExpiresAt}] ->
Now = erlang:system_time(millisecond),
case ExpiresAt > Now of
true ->
%% Cache hit - increment counter asynchronously
gen_server:cast(?SERVER, cache_hit),
{ok, Endpoint};
false ->
%% Expired - treat as miss
ets:delete(?TABLE, NodeId),
gen_server:cast(?SERVER, cache_miss),
miss
end;
[] ->
gen_server:cast(?SERVER, cache_miss),
miss
end.
%% @doc Store endpoint for a node ID.
%% Entry will expire after TTL.
-spec store(binary(), binary()) -> ok.
store(NodeId, Endpoint) when is_binary(NodeId), is_binary(Endpoint) ->
gen_server:cast(?SERVER, {store, NodeId, Endpoint}).
%% @doc Store routing info from a subscriber map (from DHT lookup).
%% Extracts node_id and endpoint from subscriber info.
-spec store_from_subscriber(map()) -> ok.
store_from_subscriber(#{node_id := NodeId, endpoint := Endpoint}) ->
store(NodeId, Endpoint);
store_from_subscriber(#{<<"node_id">> := NodeId, <<"endpoint">> := Endpoint}) ->
store(NodeId, Endpoint);
store_from_subscriber(_) ->
%% Missing required fields - skip silently
ok.
%% @doc Remove routing entry for a node ID.
-spec remove(binary()) -> ok.
remove(NodeId) when is_binary(NodeId) ->
gen_server:cast(?SERVER, {remove, NodeId}).
%% @doc Clear all routing entries.
-spec clear_all() -> ok.
clear_all() ->
gen_server:call(?SERVER, clear_all).
%% @doc Get routing table statistics.
-spec stats() -> map().
stats() ->
gen_server:call(?SERVER, stats).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
TtlMs = maps:get(ttl_ms, Opts, ?DEFAULT_TTL_MS),
CleanupIntervalMs = maps:get(cleanup_interval_ms, Opts, ?DEFAULT_CLEANUP_INTERVAL_MS),
Table = ets:new(?TABLE, [
named_table,
set,
public, % Allow direct reads from any process
{read_concurrency, true}
]),
%% Schedule periodic cleanup
erlang:send_after(CleanupIntervalMs, self(), cleanup),
{ok, #state{
table = Table,
ttl_ms = TtlMs,
cleanup_interval_ms = CleanupIntervalMs,
hits = 0,
misses = 0,
stores = 0
}}.
handle_call(stats, _From, #state{hits = Hits, misses = Misses, stores = Stores} = State) ->
Total = Hits + Misses,
HitRate = case Total of
0 -> 0.0;
_ -> Hits / Total * 100
end,
TableSize = ets:info(?TABLE, size),
Stats = #{
hits => Hits,
misses => Misses,
stores => Stores,
total_lookups => Total,
hit_rate => HitRate,
table_size => TableSize
},
{reply, Stats, State};
handle_call(clear_all, _From, State) ->
ets:delete_all_objects(?TABLE),
{reply, ok, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(cache_hit, #state{hits = Hits} = State) ->
{noreply, State#state{hits = Hits + 1}};
handle_cast(cache_miss, #state{misses = Misses} = State) ->
{noreply, State#state{misses = Misses + 1}};
handle_cast({store, NodeId, Endpoint}, #state{ttl_ms = TtlMs, stores = Stores} = State) ->
ExpiresAt = erlang:system_time(millisecond) + TtlMs,
ets:insert(?TABLE, {NodeId, Endpoint, ExpiresAt}),
{noreply, State#state{stores = Stores + 1}};
handle_cast({remove, NodeId}, State) ->
ets:delete(?TABLE, NodeId),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(cleanup, #state{cleanup_interval_ms = CleanupIntervalMs} = State) ->
Now = erlang:system_time(millisecond),
%% Remove expired entries
MatchSpec = [{{'$1', '_', '$2'}, [{'<', '$2', Now}], [true]}],
_ = ets:select_delete(?TABLE, MatchSpec),
%% Schedule next cleanup
erlang:send_after(CleanupIntervalMs, self(), cleanup),
{noreply, State};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.