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_pubsub_discovery.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% DHT integration for finding remote subscribers.
%%% Uses Kademlia DHT to publish and discover subscriptions.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_discovery).
%% API
-export([
find_subscribers/2,
find_with_cache/3,
find_with_cache/4,
announce/4,
unannounce/3,
refresh_cache/2
]).
%% Types
-type pattern() :: binary().
-type node_id() :: binary().
-type address() :: {inet:ip_address(), inet:port_number()}.
-type subscriber() :: #{node_id := node_id(), address := address()}.
-type dht_lookup_fun() :: fun((pattern()) -> {ok, [subscriber()]} | {error, term()}).
-type dht_publish_fun() :: fun((pattern(), node_id(), address()) -> ok | {error, term()}).
-type dht_unpublish_fun() :: fun((pattern(), node_id()) -> ok | {error, term()}).
-export_type([
pattern/0,
node_id/0,
address/0,
subscriber/0,
dht_lookup_fun/0,
dht_publish_fun/0,
dht_unpublish_fun/0
]).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Find remote subscribers for a pattern via DHT.
-spec find_subscribers(pattern(), dht_lookup_fun()) -> {ok, [subscriber()]} | {error, term()}.
find_subscribers(Pattern, DhtLookupFun) ->
DhtLookupFun(Pattern).
%% @doc Find subscribers with cache (default TTL: 300 seconds).
-spec find_with_cache(pattern(), macula_pubsub_cache:cache(), dht_lookup_fun()) ->
{ok, [subscriber()], macula_pubsub_cache:cache()}.
find_with_cache(Pattern, Cache, DhtLookupFun) ->
find_with_cache(Pattern, Cache, DhtLookupFun, 300).
%% @doc Find subscribers with cache and custom TTL.
-spec find_with_cache(pattern(), macula_pubsub_cache:cache(), dht_lookup_fun(), pos_integer()) ->
{ok, [subscriber()], macula_pubsub_cache:cache()}.
find_with_cache(Pattern, Cache, DhtLookupFun, TTL) ->
%% Check if cached entry is fresh
case macula_pubsub_cache:is_expired(Cache, Pattern, TTL) of
true ->
%% Cache miss or expired, query DHT
case DhtLookupFun(Pattern) of
{ok, Subscribers} ->
%% Cache the result
Cache2 = macula_pubsub_cache:put(Cache, Pattern, Subscribers),
{ok, Subscribers, Cache2};
{error, Reason} ->
%% Return error, don't update cache
{error, Reason, Cache}
end;
false ->
%% Cache hit with fresh entry
{ok, Subscribers, Cache2} = macula_pubsub_cache:get(Cache, Pattern),
{ok, Subscribers, Cache2}
end.
%% @doc Announce local subscription to DHT.
-spec announce(pattern(), node_id(), address(), dht_publish_fun()) -> ok | {error, term()}.
announce(Pattern, LocalNodeId, LocalAddress, DhtPublishFun) ->
DhtPublishFun(Pattern, LocalNodeId, LocalAddress).
%% @doc Remove local subscription from DHT.
-spec unannounce(pattern(), node_id(), dht_unpublish_fun()) -> ok | {error, term()}.
unannounce(Pattern, LocalNodeId, DhtUnpublishFun) ->
DhtUnpublishFun(Pattern, LocalNodeId).
%% @doc Refresh cache by removing expired entries.
-spec refresh_cache(macula_pubsub_cache:cache(), pos_integer()) -> macula_pubsub_cache:cache().
refresh_cache(Cache, TTL) ->
%% Get all patterns in cache
Patterns = get_all_patterns(Cache),
%% Invalidate expired patterns
lists:foldl(
fun(Pattern, AccCache) ->
case macula_pubsub_cache:is_expired(AccCache, Pattern, TTL) of
true -> macula_pubsub_cache:invalidate(AccCache, Pattern);
false -> AccCache
end
end,
Cache,
Patterns
).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Extract all patterns from cache.
%% NOTE: This is a workaround since macula_pubsub_cache doesn't expose list_patterns.
%% We reconstruct it from the cache entries.
-spec get_all_patterns(macula_pubsub_cache:cache()) -> [pattern()].
get_all_patterns(#{entries := Entries}) ->
lists:map(
fun(Entry) -> maps:get(pattern, Entry) end,
Entries
).