Packages

macula

0.37.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_system macula_subscriber_cache.erl
Raw

src/macula_gateway_system/macula_subscriber_cache.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% DHT Subscriber Cache - Caches topic→subscribers mappings for fast pub/sub.
%%%
%%% Problem: DHT lookups on every PUBLISH cause 50-200ms latency per message.
%%% Solution: Cache topic→subscribers mappings with TTL-based expiration.
%%%
%%% Design:
%%% - ETS table for O(1) lookup
%%% - TTL-based expiration (default 5 seconds)
%%% - Automatic cache invalidation on subscribe/unsubscribe
%%% - Thread-safe concurrent reads
%%% - Adaptive rate-limiting to prevent discovery storms (default 2 seconds)
%%%
%%% Rate-Limiting:
%%% When cache expires and many publishes occur, rate-limiting prevents
%%% flooding the DHT with lookup queries. Only one DHT query per topic
%%% is allowed within the min_discovery_interval_ms window.
%%%
%%% Expected improvement:
%%% - 5-10x latency reduction for high-frequency topics (caching)
%%% - 2-3x improvement during high-frequency publishing (rate-limiting)
%%%
%%% Configuration Options:
%%% - ttl_ms: Cache entry TTL (default: 5000ms)
%%% - min_discovery_interval_ms: Minimum time between DHT queries per topic (default: 2000ms)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_subscriber_cache).
-behaviour(gen_server).
%% API
-export([
start_link/0,
start_link/1,
lookup/1,
store/2,
invalidate/1,
invalidate_all/0,
stats/0,
%% Rate-limiting API
should_query_dht/1,
record_dht_query/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-define(SERVER, ?MODULE).
-define(TABLE, macula_subscriber_cache_table).
-define(RATE_LIMIT_TABLE, macula_discovery_rate_limit_table).
-define(DEFAULT_TTL_MS, 5000). % 5 seconds default TTL
-define(DEFAULT_MIN_DISCOVERY_INTERVAL_MS, 2000). % 2 seconds between DHT queries per topic
-define(CLEANUP_INTERVAL_MS, 1000). % Cleanup every 1 second
-record(state, {
table :: ets:tid(),
rate_limit_table :: ets:tid(),
ttl_ms :: pos_integer(),
min_discovery_interval_ms :: pos_integer(),
hits :: non_neg_integer(),
misses :: non_neg_integer(),
rate_limited :: non_neg_integer()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the subscriber cache with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start the subscriber cache with options.
%% Options:
%% - ttl_ms: Cache entry TTL in milliseconds (default: 5000)
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Look up subscribers for a topic from cache.
%% Returns {ok, Subscribers} on cache hit, or {miss, TopicKey} on cache miss.
%% TopicKey is the SHA256 hash of the topic binary.
-spec lookup(binary()) -> {ok, list()} | {miss, binary()}.
lookup(Topic) when is_binary(Topic) ->
TopicKey = crypto:hash(sha256, Topic),
lookup_by_key(TopicKey).
%% @doc Store subscribers for a topic in cache.
%% Subscribers will expire after TTL.
-spec store(binary(), list()) -> ok.
store(Topic, Subscribers) when is_binary(Topic), is_list(Subscribers) ->
TopicKey = crypto:hash(sha256, Topic),
store_by_key(TopicKey, Subscribers).
%% @doc Invalidate cache entry for a topic.
%% Should be called when subscription changes occur.
-spec invalidate(binary()) -> ok.
invalidate(Topic) when is_binary(Topic) ->
TopicKey = crypto:hash(sha256, Topic),
invalidate_by_key(TopicKey).
%% @doc Invalidate all cache entries.
%% Useful for testing or when major topology changes occur.
-spec invalidate_all() -> ok.
invalidate_all() ->
gen_server:call(?SERVER, invalidate_all).
%% @doc Get cache statistics.
-spec stats() -> map().
stats() ->
gen_server:call(?SERVER, stats).
%% @doc Check if we should query DHT for a topic (rate-limiting check).
%% Returns true if enough time has passed since last DHT query.
%% Returns false if we recently queried DHT (rate-limited).
-spec should_query_dht(binary()) -> boolean().
should_query_dht(Topic) when is_binary(Topic) ->
TopicKey = crypto:hash(sha256, Topic),
should_query_dht_by_key(TopicKey).
%% @doc Record that we just performed a DHT query for a topic.
%% Call this after a successful DHT lookup to reset the rate-limit timer.
-spec record_dht_query(binary()) -> ok.
record_dht_query(Topic) when is_binary(Topic) ->
TopicKey = crypto:hash(sha256, Topic),
record_dht_query_by_key(TopicKey).
%%%===================================================================
%%% Internal Lookup Functions (direct ETS access for speed)
%%%===================================================================
lookup_by_key(TopicKey) ->
case ets:lookup(?TABLE, TopicKey) of
[{TopicKey, Subscribers, ExpiresAt}] ->
Now = erlang:system_time(millisecond),
case ExpiresAt > Now of
true ->
%% Cache hit - increment counter asynchronously
gen_server:cast(?SERVER, cache_hit),
{ok, Subscribers};
false ->
%% Expired - treat as miss
ets:delete(?TABLE, TopicKey),
gen_server:cast(?SERVER, cache_miss),
{miss, TopicKey}
end;
[] ->
gen_server:cast(?SERVER, cache_miss),
{miss, TopicKey}
end.
store_by_key(TopicKey, Subscribers) ->
gen_server:cast(?SERVER, {store, TopicKey, Subscribers}).
invalidate_by_key(TopicKey) ->
gen_server:cast(?SERVER, {invalidate, TopicKey}).
%% @private
%% @doc Check rate-limit table to see if we can query DHT.
%% Direct ETS access for performance.
should_query_dht_by_key(TopicKey) ->
Now = erlang:system_time(millisecond),
case ets:lookup(?RATE_LIMIT_TABLE, TopicKey) of
[{TopicKey, LastQueryTime, MinInterval}] ->
TimeSinceLastQuery = Now - LastQueryTime,
case TimeSinceLastQuery >= MinInterval of
true ->
%% Enough time has passed - allow query
true;
false ->
%% Rate-limited - too soon
gen_server:cast(?SERVER, rate_limited),
false
end;
[] ->
%% No previous query - allow
true
end.
%% @private
%% @doc Record DHT query in rate-limit table.
record_dht_query_by_key(TopicKey) ->
gen_server:cast(?SERVER, {record_dht_query, TopicKey}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
TtlMs = maps:get(ttl_ms, Opts, ?DEFAULT_TTL_MS),
MinDiscoveryIntervalMs = maps:get(min_discovery_interval_ms, Opts, ?DEFAULT_MIN_DISCOVERY_INTERVAL_MS),
Table = ets:new(?TABLE, [
named_table,
set,
public, % Allow direct reads from any process
{read_concurrency, true}
]),
RateLimitTable = ets:new(?RATE_LIMIT_TABLE, [
named_table,
set,
public, % Allow direct reads from any process
{read_concurrency, true}
]),
%% Schedule periodic cleanup
erlang:send_after(?CLEANUP_INTERVAL_MS, self(), cleanup),
{ok, #state{
table = Table,
rate_limit_table = RateLimitTable,
ttl_ms = TtlMs,
min_discovery_interval_ms = MinDiscoveryIntervalMs,
hits = 0,
misses = 0,
rate_limited = 0
}}.
handle_call(stats, _From, #state{hits = Hits, misses = Misses, rate_limited = RateLimited} = State) ->
Total = Hits + Misses,
HitRate = case Total of
0 -> 0.0;
_ -> Hits / Total * 100
end,
TableSize = ets:info(?TABLE, size),
RateLimitTableSize = ets:info(?RATE_LIMIT_TABLE, size),
Stats = #{
hits => Hits,
misses => Misses,
total => Total,
hit_rate => HitRate,
table_size => TableSize,
rate_limited => RateLimited,
rate_limit_table_size => RateLimitTableSize
},
{reply, Stats, State};
handle_call(invalidate_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, TopicKey, Subscribers}, #state{ttl_ms = TtlMs} = State) ->
ExpiresAt = erlang:system_time(millisecond) + TtlMs,
ets:insert(?TABLE, {TopicKey, Subscribers, ExpiresAt}),
{noreply, State};
handle_cast({invalidate, TopicKey}, State) ->
ets:delete(?TABLE, TopicKey),
{noreply, State};
handle_cast(rate_limited, #state{rate_limited = RateLimited} = State) ->
{noreply, State#state{rate_limited = RateLimited + 1}};
handle_cast({record_dht_query, TopicKey}, #state{min_discovery_interval_ms = MinInterval} = State) ->
Now = erlang:system_time(millisecond),
ets:insert(?RATE_LIMIT_TABLE, {TopicKey, Now, MinInterval}),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(cleanup, #state{min_discovery_interval_ms = MinInterval} = State) ->
Now = erlang:system_time(millisecond),
%% Remove expired cache entries
CacheMatchSpec = [{{'$1', '_', '$2'}, [{'<', '$2', Now}], [true]}],
_NumCacheDeleted = ets:select_delete(?TABLE, CacheMatchSpec),
%% Remove stale rate-limit entries (older than 2x min interval)
%% This prevents the rate-limit table from growing unbounded
MaxRateLimitAge = MinInterval * 2,
CutoffTime = Now - MaxRateLimitAge,
RateLimitMatchSpec = [{{'$1', '$2', '_'}, [{'<', '$2', CutoffTime}], [true]}],
_NumRateLimitDeleted = ets:select_delete(?RATE_LIMIT_TABLE, RateLimitMatchSpec),
%% Schedule next cleanup
erlang:send_after(?CLEANUP_INTERVAL_MS, self(), cleanup),
{noreply, State};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.