Packages

macula

0.7.25
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_pubsub_delivery.erl
Raw

src/macula_pubsub_delivery.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Message routing and delivery to local and remote subscribers.
%%% Combines local registry and remote discovery for full fan-out.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_delivery).
%% API
-export([
deliver_local/2,
deliver_remote/3,
publish/4,
get_matching_patterns/2
]).
%% Types
-type message() :: #{
topic := binary(),
payload := term(),
timestamp := integer()
}.
-type delivery_result() :: ok | {ok, term()} | {error, term()}.
-type discovery_fun() :: fun((binary()) -> {ok, [map()]} | {error, term()}).
-type send_fun() :: fun((message(), macula_pubsub_discovery:address()) -> ok | {error, term()}).
-export_type([message/0, delivery_result/0, discovery_fun/0, send_fun/0]).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Deliver message to all matching local subscribers.
%% Crashes if subscriber callback fails - indicates dead subscriber process.
-spec deliver_local(message(), macula_pubsub_registry:registry()) -> [delivery_result()].
deliver_local(Message, Registry) ->
Topic = maps:get(topic, Message),
%% Find matching subscriptions
Subscriptions = macula_pubsub_registry:match(Registry, Topic),
%% Deliver to each callback (let it crash on dead subscribers)
lists:map(
fun(Sub) ->
Callback = maps:get(callback, Sub),
Callback ! Message,
{ok, maps:get(subscriber_id, Sub)}
end,
Subscriptions
).
%% @doc Deliver message to remote subscribers via QUIC.
-spec deliver_remote(message(), [macula_pubsub_discovery:subscriber()], send_fun()) ->
[delivery_result()].
deliver_remote(Message, RemoteSubscribers, SendFun) ->
lists:map(
fun(RemoteSub) ->
Address = maps:get(address, RemoteSub),
SendFun(Message, Address)
end,
RemoteSubscribers
).
%% @doc Publish message to both local and remote subscribers.
%% Returns {LocalResults, RemoteResults}.
-spec publish(message(), macula_pubsub_registry:registry(), discovery_fun(), send_fun()) ->
{[delivery_result()], [delivery_result()]}.
publish(Message, Registry, DiscoveryFun, SendFun) ->
Topic = maps:get(topic, Message),
%% Deliver to local subscribers
LocalResults = deliver_local(Message, Registry),
%% Find matching patterns for remote discovery
MatchingPatterns = get_matching_patterns(Topic, Registry),
%% Discover remote subscribers for each matching pattern
RemoteSubscribers = discover_remote_subscribers(MatchingPatterns, DiscoveryFun),
%% Deliver to remote subscribers
RemoteResults = deliver_remote(Message, RemoteSubscribers, SendFun),
{LocalResults, RemoteResults}.
%% @doc Get all unique patterns that match the topic.
%% Used for remote subscriber discovery.
-spec get_matching_patterns(binary(), macula_pubsub_registry:registry()) -> [binary()].
get_matching_patterns(Topic, Registry) ->
%% Get all patterns from registry
AllPatterns = macula_pubsub_registry:list_patterns(Registry),
%% Filter to matching patterns
MatchingPatterns = lists:filter(
fun(Pattern) ->
macula_pubsub_topic:matches(Topic, Pattern)
end,
AllPatterns
),
%% Return unique patterns
lists:usort(MatchingPatterns).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Discover remote subscribers for multiple patterns.
-spec discover_remote_subscribers([binary()], discovery_fun()) ->
[macula_pubsub_discovery:subscriber()].
discover_remote_subscribers(Patterns, DiscoveryFun) ->
%% Query discovery for each pattern
AllSubscribers = lists:flatmap(
fun(Pattern) ->
case DiscoveryFun(Pattern) of
{ok, Subscribers} -> Subscribers;
{error, _Reason} -> [] % Ignore discovery errors
end
end,
Patterns
),
%% Deduplicate by node_id
deduplicate_by_node_id(AllSubscribers).
%% @doc Remove duplicate subscribers by node_id.
-spec deduplicate_by_node_id([macula_pubsub_discovery:subscriber()]) ->
[macula_pubsub_discovery:subscriber()].
deduplicate_by_node_id(Subscribers) ->
%% Build map: node_id -> subscriber
SubscriberMap = lists:foldl(
fun(Sub, Acc) ->
NodeId = maps:get(node_id, Sub),
Acc#{NodeId => Sub}
end,
#{},
Subscribers
),
%% Return unique subscribers
maps:values(SubscriberMap).