Packages

macula

4.0.0
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 pubsub macula_pubsub.erl
Raw

src/pubsub/macula_pubsub.erl

%% @doc Pubsub surface for the V2 SDK.
%%
%% Thin delegation over `macula_client' (the pool). The pool owns
%% the link state machine, replication, replay, and dedup; this
%% module is the named public entry point that consumers reach for
%% (or, more often, the `macula' facade re-exports of the same
%% functions).
%%
%% == Realm-per-call ==
%%
%% Per `PLAN_V2_PARITY' Q2 ยง2: every call carries its own 32-byte
%% realm tag. There is no connect-time default realm. A single pool
%% can multiplex any number of realms with no extra plumbing.
%%
%% == Quick start ==
%%
%% ```
%% {ok, Pool} = macula:connect(Seeds, ConnectOpts),
%% ok = macula_pubsub:publish(Pool, Realm, Topic, Payload),
%% {ok, Sub} = macula_pubsub:subscribe(Pool, Realm, Topic, self()),
%% receive
%% {macula_event, Sub, Topic, Payload, Meta} -> ok
%% end,
%% ok = macula_pubsub:unsubscribe(Pool, Sub).
%% '''
%%
%% See `docs/guides/PUBSUB_GUIDE.md' for a full guide and
%% `docs/migrations/V1_TO_V2_PUBSUB.md' for the breaking changes
%% from the pre-3.11.0 facade.
-module(macula_pubsub).
-export([publish/4, publish/5,
subscribe/4, subscribe/5,
subscribe_callback/4,
unsubscribe/2]).
-export_type([callback/0]).
%% Callback shape accepted by `subscribe_callback/4'. Invoked once
%% per inbound event in a separate receiver process so a slow
%% callback does not back-pressure the pool.
-type callback() :: fun((Topic :: binary(),
Payload :: term(),
Meta :: map()) -> any()).
%% @doc Publish to `(Realm, Topic)' on `Pool'. Equivalent to
%% `publish/5' with empty opts.
-spec publish(macula_client:pool(), <<_:256>>, binary(), term()) ->
ok | {error, term()}.
publish(Pool, Realm, Topic, Payload) ->
publish(Pool, Realm, Topic, Payload, #{}).
%% @doc Publish to `(Realm, Topic)' on `Pool'.
%%
%% `Opts' currently honored:
%% <ul>
%% <li>`timeout_ms' โ€” gen_server call timeout (default 5_000).
%% Most apps leave this as default.</li>
%% </ul>
%%
%% Returns `ok' as soon as one configured station accepts the
%% PUBLISH frame (partial success = success, per
%% `PLAN_V2_PARITY' ยง5.1.1). Returns
%% `{error, {transient, no_healthy_station}}' when the pool has no
%% spawned links; the caller may retry.
-spec publish(macula_client:pool(), <<_:256>>, binary(), term(), map()) ->
ok | {error, term()}.
publish(Pool, Realm, Topic, Payload, Opts)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Topic),
is_map(Opts) ->
macula_client:publish(Pool, Realm, Topic, Payload, Opts).
%% @doc Subscribe `Subscriber' to `(Realm, Topic)' via `Pool'.
%% Equivalent to `subscribe/5' with empty opts.
-spec subscribe(macula_client:pool(), <<_:256>>, binary(), pid()) ->
{ok, reference()}.
subscribe(Pool, Realm, Topic, Subscriber) ->
subscribe(Pool, Realm, Topic, Subscriber, #{}).
%% @doc Subscribe `Subscriber' to `(Realm, Topic)' via `Pool'.
%%
%% Returns `{ok, SubRef}'. `Subscriber' subsequently receives
%% `{macula_event, SubRef, Topic, Payload, Meta}' for each delivered
%% event; `Meta' is a map carrying `realm', `publisher', `seq', and
%% `delivered_via'. Stores receive `{macula_event_gone, SubRef,
%% Reason}' once when the subscription terminates (pool close,
%% subscriber pid death).
%%
%% `Opts' is a forward-compatible map; Phase 1 honors no
%% subscribe-time options. Future phases (history replay, server-
%% side filters) will add named keys.
-spec subscribe(macula_client:pool(), <<_:256>>, binary(), pid(), map()) ->
{ok, reference()}.
subscribe(Pool, Realm, Topic, Subscriber, Opts)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Topic),
is_pid(Subscriber),
is_map(Opts) ->
macula_client:subscribe(Pool, Realm, Topic, Subscriber, Opts).
%% @doc Drop a subscription. Idempotent โ€” unknown `SubRef' is a
%% no-op.
-spec unsubscribe(macula_client:pool(), reference()) -> ok.
unsubscribe(Pool, SubRef) when is_pid(Pool), is_reference(SubRef) ->
macula_client:unsubscribe(Pool, SubRef).
%% @doc Subscribe with a callback function instead of a receiver pid.
%% Spawns a small receiver process internally that drives the
%% callback for every inbound event. The receiver monitors the
%% caller; if the caller dies, the receiver follows and the
%% subscription is cleaned up by the pool's standard subscriber-DOWN
%% path.
%%
%% A crashing callback does NOT kill the receiver โ€” the exception is
%% logged and the next event is delivered. This is intentional: a
%% transient bug in event handler N should not lose events N+1..M.
%%
%% Caller cleanup: invoke `unsubscribe(Pool, SubRef)' with the
%% returned ref. The receiver shuts down on the resulting
%% `macula_event_gone' message.
-spec subscribe_callback(macula_client:pool(), <<_:256>>, binary(),
callback()) ->
{ok, reference()} | {error, term()}.
subscribe_callback(Pool, Realm, Topic, Callback)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Topic),
is_function(Callback, 3) ->
Caller = self(),
{Receiver, Mon} = spawn_monitor(
fun() -> receiver_init(Caller, Pool, Realm, Topic, Callback) end),
await_init(Receiver, Mon).
await_init(Receiver, Mon) ->
receive
{?MODULE, started, Receiver, SubRef} ->
erlang:demonitor(Mon, [flush]),
{ok, SubRef};
{?MODULE, failed, Receiver, Err} ->
erlang:demonitor(Mon, [flush]),
Err;
{'DOWN', Mon, process, Receiver, Reason} ->
{error, {receiver_died, Reason}}
after 5_000 ->
exit(Receiver, init_timeout),
erlang:demonitor(Mon, [flush]),
{error, init_timeout}
end.
receiver_init(Caller, Pool, Realm, Topic, Callback) ->
CallerMon = erlang:monitor(process, Caller),
on_subscribe(macula_client:subscribe(Pool, Realm, Topic, self(), #{}),
Caller, CallerMon, Callback).
on_subscribe({ok, SubRef}, Caller, CallerMon, Callback) ->
Caller ! {?MODULE, started, self(), SubRef},
receiver_loop(SubRef, CallerMon, Callback);
on_subscribe({error, _} = E, Caller, CallerMon, _Callback) ->
Caller ! {?MODULE, failed, self(), E},
erlang:demonitor(CallerMon, [flush]).
receiver_loop(SubRef, CallerMon, Callback) ->
receive
{macula_event, SubRef, Topic, Payload, Meta} ->
invoke(Callback, Topic, Payload, Meta),
receiver_loop(SubRef, CallerMon, Callback);
{macula_event_gone, SubRef, _Reason} ->
ok;
{'DOWN', CallerMon, process, _, _} ->
ok
end.
%% Guard the user callback so a transient handler bug does not
%% wedge the entire subscription stream. This is the rare place where
%% try/catch is the right tool: the SDK is owning a long-lived
%% receiver on behalf of an opaque consumer fun.
invoke(Callback, Topic, Payload, Meta) ->
try Callback(Topic, Payload, Meta) of
_ -> ok
catch
Class:Reason:Stack ->
logger:warning("[macula_pubsub] callback crashed: ~p:~p~n~p",
[Class, Reason, Stack])
end.