Packages
macula
6.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
Current section
Files
src/client/macula_client.erl
%% @doc `macula_client' — the canonical pool client.
%%
%% Holds N peering links to N stations and routes ops with
%% replication, subscription replay, and inbound-event dedup. Apps
%% don't manage individual `macula_station_link' workers; they call
%% `macula_client' (or the `macula' facade, which re-exports the
%% public surface).
%%
%% Per `PLAN_V2_PARITY' Q2 §1: pool is the canonical client handle.
%% A single-station link is an internal worker only.
%%
%% == Lifecycle ==
%%
%% <pre>
%% {ok, Pool} = macula_client:connect(Seeds, Opts).
%% ok = macula_client:publish(Pool, Realm, Topic, Payload, #{}).
%% {ok, Sub} = macula_client:subscribe(Pool, Realm, Topic, self(), #{}).
%% receive {macula_event, Sub, Topic, Payload, Meta} -> ... end.
%% ok = macula_client:unsubscribe(Pool, Sub).
%% ok = macula_client:close(Pool).
%% </pre>
%%
%% == Replication ==
%%
%% `publish/5' fans the PUBLISH frame to `replication_factor' (default
%% 1) currently-spawned links. **Partial success counts as success**
%% per `PLAN_V2_PARITY' §5.1.1: the call returns `ok' as soon as one
%% link accepts the frame; the others are best-effort. When zero
%% links are spawned the call returns
%% `{error, {transient, no_healthy_station}}'.
%%
%% `subscribe/5' applies to every spawned link. The pool delivers a
%% deduped event stream to the consumer regardless of which link
%% relayed any given EVENT.
%%
%% == Dedup ==
%%
%% Inbound EVENT frames are keyed by `(Realm, Publisher, Seq)' in an
%% ETS table owned by the pool. The table is swept every
%% `dedup_sweep_ms' (default 30s) for entries older than
%% `dedup_window_ms' (default 60s).
%%
%% == Replay ==
%%
%% When a link's process dies the pool monitor fires; the pool
%% schedules a respawn after ?LINK_RESPAWN_DELAY_MS (1s). On respawn,
%% the pool re-issues every currently-tracked (Realm, Topic)
%% subscription against the new link via the internal
%% macula_client_replay helper.
-module(macula_client).
-behaviour(gen_server).
-export([connect/2, close/1, child_spec/3, status/1, links/1]).
%% Internal API — called by `macula_pubsub' (and future surfaces).
-export([publish/5, subscribe/5, unsubscribe/2]).
%% RPC fan-out (since 3.16.0) — called by the `macula' facade.
-export([call/5, advertise/4, unadvertise/3]).
%% Streaming RPC (since 3.17.0) — called by the `macula' facade.
-export([call_stream/5, advertise_stream/5, unadvertise_stream/3]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2,
terminate/2, code_change/3]).
-export_type([pool/0, opts/0, seed/0, status/0, link_info/0, handler/0, stream_handler/0]).
-type pool() :: pid().
%% RPC handler — accepted by `advertise/4'. Either a 1-arg fun called
%% with the inbound payload, or `{Module, Function}' invoked as
%% `Module:Function(Payload)'. Re-exported here so consumers do not
%% have to reach into the private `macula_station_link' module.
-type handler() :: fun((term()) -> term())
| {module(), atom()}.
%% Streaming RPC handler — accepted by `advertise_stream/5'. A 2-arg
%% fun invoked as `Handler(StreamPid, Args)' where `StreamPid' is
%% the local `macula_stream' bound to the inbound STREAM_OPEN.
-type stream_handler() :: fun((pid(), term()) -> any()).
%% Aggregate health snapshot of a pool. See `status/1'.
-type status() :: #{
seeds := [seed()],
healthy_links := non_neg_integer(),
failed_links := non_neg_integer(),
self_node_id := macula_identity:pubkey(),
subscriptions := non_neg_integer()
}.
%% Per-link view returned by `links/1'. One entry per configured seed
%% that currently has a spawned link worker. `node_id' is the peer
%% station's pubkey (`undefined' until CONNECT/HELLO completes);
%% `host' is the dial host parsed from the seed.
-type link_info() :: #{
seed := seed(),
host := binary() | undefined,
pid := pid(),
connected := boolean(),
node_id := macula_identity:pubkey() | undefined
}.
-type seed() :: binary() | string()
| #{host := binary() | string(),
port := inet:port_number()}.
-type opts() :: #{
%% Shared Ed25519 keypair for every link in the pool. Stations see
%% the pool as a single peer (one pubkey across N links).
%% Auto-generated when absent.
identity => macula_identity:key_pair(),
%% How many of the pool's currently-connected links accept a
%% single PUBLISH frame. Partial success counts as success
%% (`PLAN_V2_PARITY' §5.1.1). Default 1.
replication_factor => pos_integer(),
%% Per-link capability bitfield, forwarded to every
%% `macula_station_link'. Default 0. Reserved for future use.
capabilities => non_neg_integer(),
%% ALPN identifiers offered to the QUIC handshake. Default
%% `[<<"macula">>]'.
alpn => [binary()],
%% Per-link CONNECT/HELLO deadline in milliseconds. Default 30_000.
%% Applies to each link independently — total pool readiness
%% wallclock can be up to N×timeout for sequential dial fallback.
connect_timeout_ms => pos_integer(),
%% Inbound-EVENT dedup window in milliseconds. The pool keys
%% inbound events on `(Realm, Publisher, Seq)' so duplicate
%% deliveries from multiple subscribed links collapse to one
%% emission per consumer. Default 60_000.
dedup_window_ms => non_neg_integer(),
%% How often the dedup table is swept for entries older than
%% `dedup_window_ms'. Default 30_000.
dedup_sweep_ms => pos_integer()
}.
%% V1 multi_relay options that have NO V2 equivalent. Callers passing
%% these from a V1 migration get a one-shot warning and the opt is
%% silently ignored. Keeping the names listed here so the warning
%% can name them helpfully.
-define(V1_LEGACY_OPTS, [relays, realm, site, connections]).
-define(DEFAULT_REPLICATION, 1).
-define(DEFAULT_DEDUP_WINDOW_MS, 60_000).
-define(DEFAULT_DEDUP_SWEEP_MS, 30_000).
-define(LINK_RESPAWN_DELAY_MS, 1_000).
-record(link_state, {
seed :: seed(),
pid :: pid() | undefined,
mon :: reference() | undefined
}).
-record(sub_spec, {
realm :: <<_:256>>,
topic :: binary(),
subscriber :: pid(),
mon :: reference()
}).
-record(state, {
seeds :: [seed()],
identity :: macula_identity:key_pair(),
link_opts :: map(),
replication :: pos_integer(),
dedup_window :: non_neg_integer(),
dedup_sweep :: pos_integer(),
%% Pool-owned monotonic publish sequence. Stamped onto every
%% outbound PUBLISH (via `macula_station_link:publish/5') so the
%% station-side `(publisher, seq)' dedup stays stable across link
%% respawns — the publisher pubkey is the pool's, shared by all
%% links, so the seq must be owned by the pool, not the link.
%% Seeded from wall-clock µs at init so a pool restart does not
%% re-issue seqs that collide with the pre-restart tail still in a
%% station's dedup window (see
%% macula-station/plans/PLAN_PUBSUB_E2E_SIGNED_EVENTS.md).
publish_seq :: non_neg_integer(),
%% seed → link_state
links = #{} :: #{seed() => #link_state{}},
%% pool-owned SubRef → sub_spec
subs = #{} :: #{reference() => #sub_spec{}},
%% {realm, topic} → set of pool-owned SubRefs
topic_index = #{} :: #{{<<_:256>>, binary()} => sets:set(reference())},
%% Advertised procedures — pool replays these on link respawn.
%% {realm, procedure} → handler
procs = #{} :: #{{<<_:256>>, binary()} => handler()},
%% Advertised streaming procedures — replayed on link respawn
%% alongside `procs'. {realm, procedure} → {mode, handler}
stream_procs = #{} :: #{{<<_:256>>, binary()} =>
{macula_frame:stream_mode(),
stream_handler()}},
dedup_tab :: ets:tid()
}).
%%====================================================================
%% Public API
%%====================================================================
%% @doc Spawn a pool with one link per seed. Returns immediately;
%% link handshakes complete asynchronously. Publish/subscribe block
%% until at least one link is connected (or fail with
%% `{error, {transient, no_healthy_station}}' on the publish path).
-spec connect([seed()], opts()) -> {ok, pool()} | {error, term()}.
connect(Seeds, Opts) when is_list(Seeds), is_map(Opts) ->
gen_server:start_link(?MODULE, {Seeds, Opts}, []).
%% @doc Stop the pool. Every subscriber receives a final
%% `{macula_event_gone, SubRef, pool_closed}' message; every link
%% terminates with the pool.
-spec close(pool()) -> ok.
close(Pool) ->
gen_server:stop(Pool, normal, 5_000).
%% @doc OTP child spec — drop the pool into a caller's supervision
%% tree. `Id' is the supervisor child id.
-spec child_spec(term(), [seed()], opts()) -> supervisor:child_spec().
child_spec(Id, Seeds, Opts) ->
#{id => Id,
start => {?MODULE, connect, [Seeds, Opts]},
restart => permanent,
shutdown => 5_000,
type => worker,
modules => [?MODULE]}.
%% @doc Issue a CALL frame against the pool. Tries each healthy link
%% in turn and returns the first non-error reply. Returns
%% `{error, no_healthy_station}' when no link has completed its
%% CONNECT/HELLO handshake.
%%
%% Realm is per-call (32 bytes). Different realms can share a single
%% pool with no extra plumbing.
-spec call(pool(), <<_:256>>, binary(), term(), pos_integer()) ->
{ok, term()} | {error, term()}.
call(Pool, Realm, Procedure, Payload, TimeoutMs)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure),
is_integer(TimeoutMs), TimeoutMs > 0 ->
gen_server:call(Pool, {rpc_call, Realm, Procedure, Payload, TimeoutMs},
TimeoutMs + 1_000).
%% @doc Advertise a procedure handler on every healthy link. Stored
%% in pool state so links respawned later replay the advertisement.
%% Returns `ok' when at least one link accepted the registration.
-spec advertise(pool(), <<_:256>>, binary(), handler()) ->
ok | {error, term()}.
advertise(Pool, Realm, Procedure, Handler)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure),
(is_function(Handler, 1) orelse
(is_tuple(Handler) andalso tuple_size(Handler) =:= 2)) ->
gen_server:call(Pool, {advertise, Realm, Procedure, Handler}, 5_000).
%% @doc Drop a previously-advertised procedure on every healthy link
%% and remove it from the pool's replay state. Idempotent.
-spec unadvertise(pool(), <<_:256>>, binary()) -> ok.
unadvertise(Pool, Realm, Procedure)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure) ->
gen_server:call(Pool, {unadvertise, Realm, Procedure}, 5_000).
%% @doc Open a streaming RPC against the pool. Picks the first
%% currently-healthy link and opens the stream there; the returned
%% stream pid is sticky — if the underlying link dies, the stream
%% errors with `{error, peer_down}' and the caller must re-open.
%%
%% Returns `{error, no_healthy_station}' when no link has completed
%% its CONNECT/HELLO handshake. `Realm' (32 bytes) and `Procedure'
%% name the remote endpoint. `Args' is the opening payload; `Opts'
%% accepts `mode' (default `server_stream'), `owner' (default the
%% calling pid), and `deadline_ms'.
-spec call_stream(pool(), <<_:256>>, binary(), term(), map()) ->
{ok, pid()} | {error, term()}.
call_stream(Pool, Realm, Procedure, Args, Opts)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure),
is_map(Opts) ->
gen_server:call(Pool,
{rpc_call_stream, Realm, Procedure, Args,
Opts#{owner => maps:get(owner, Opts, self())}},
5_000).
%% @doc Advertise a streaming procedure handler on every healthy
%% link. Stored in pool state so links respawned later replay the
%% advertisement. Returns `ok' when at least one link accepted the
%% registration.
-spec advertise_stream(pool(), <<_:256>>, binary(),
macula_frame:stream_mode(),
stream_handler()) ->
ok | {error, term()}.
advertise_stream(Pool, Realm, Procedure, Mode, Handler)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure),
(Mode =:= server_stream orelse Mode =:= client_stream
orelse Mode =:= bidi),
is_function(Handler, 2) ->
gen_server:call(Pool,
{advertise_stream, Realm, Procedure, Mode, Handler},
5_000).
%% @doc Drop a streaming procedure on every healthy link and remove
%% it from the pool's replay state. Idempotent.
-spec unadvertise_stream(pool(), <<_:256>>, binary()) -> ok.
unadvertise_stream(Pool, Realm, Procedure)
when is_pid(Pool),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure) ->
gen_server:call(Pool, {unadvertise_stream, Realm, Procedure}, 5_000).
%% @doc Aggregate health snapshot of the pool. Single round-trip to
%% the pool's gen_server plus one `is_connected' probe per spawned
%% link (each capped at 1s). Suitable for `/health' or
%% `/status' endpoints; not for hot-loop polling.
%%
%% Counts:
%% <ul>
%% <li>`healthy_links' — links whose worker pid is alive and whose
%% CONNECT/HELLO handshake has completed.</li>
%% <li>`failed_links' — every other configured seed (link not yet
%% spawned, dead, or still handshaking).</li>
%% </ul>
-spec status(pool()) -> {ok, status()}.
status(Pool) when is_pid(Pool) ->
gen_server:call(Pool, status, 5_000).
%% @doc Per-link snapshot of the pool — one `link_info()' per
%% configured seed that currently has a spawned link worker. Unlike
%% `status/1' (which only aggregates counts), this exposes each link's
%% `node_id' (peer station pubkey), dial `host', `pid', and
%% `connected' flag, so a caller can resolve a specific station (by
%% pubkey or hostname) to its link and address it directly.
%%
%% One `is_connected/1' + `peer_node_id/1' probe per spawned link
%% (each capped at 1s). Not for hot-loop polling.
-spec links(pool()) -> {ok, [link_info()]}.
links(Pool) when is_pid(Pool) ->
gen_server:call(Pool, links, 5_000).
%% @doc Publish a frame to `replication_factor' currently-spawned
%% links. Partial success = success. Realm is per-call (32 bytes).
%%
%% The payload is checked for wire admissibility HERE, in the caller's
%% process, before the pool is touched. Downstream the send is a
%% `gen_statem:cast' into a shared peering connection that encodes
%% without a try/catch, so an unrepresentable term would kill that
%% connection and every other producer's in-flight traffic with it,
%% asynchronously, after this function had already answered `ok'.
%% Checking first is what makes the `ok' falsifiable.
%%
%% Returns `{error, {unsupported_payload_type, Type, Path}}' naming the
%% offending value and where it sits in the term. Floats are the common
%% case: scale them to integers (micro-units) or send binary strings.
-spec publish(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) ->
publish_checked(macula_frame:check_payload(Payload),
Pool, Realm, Topic, Payload, Opts).
publish_checked(ok, Pool, Realm, Topic, Payload, Opts) ->
Timeout = maps:get(timeout_ms, Opts, 5_000),
gen_server:call(Pool, {publish, Realm, Topic, Payload, Opts},
Timeout + 500);
publish_checked({error, _} = Rejected, _Pool, _Realm, _Topic, _Payload, _Opts) ->
Rejected.
%% @doc Subscribe `Subscriber' to `(Realm, Topic)'. The pool
%% subscribes every currently-spawned link and dedupes inbound
%% events before fan-out. Returns `{ok, SubRef}'; `Subscriber'
%% receives `{macula_event, SubRef, Topic, Payload, Meta}' for each
%% delivered event and `{macula_event_gone, SubRef, Reason}' once
%% when the pool closes or the subscriber pid dies.
-spec subscribe(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) ->
gen_server:call(Pool, {subscribe, Realm, Topic, Subscriber, Opts},
5_000).
%% @doc Drop a subscription. Idempotent — unknown `SubRef' is a
%% no-op. The wire-level link subscription persists for the pool's
%% lifetime (one wire sub per `(Realm, Topic)' multiplexed across
%% local consumers); Phase 4 will tighten this.
-spec unsubscribe(pool(), reference()) -> ok.
unsubscribe(Pool, SubRef) when is_pid(Pool), is_reference(SubRef) ->
gen_server:call(Pool, {unsubscribe, SubRef}, 5_000).
%%====================================================================
%% gen_server
%%====================================================================
init({Seeds, Opts}) ->
process_flag(trap_exit, true),
warn_legacy_opts(Opts),
Identity = maps:get(identity, Opts, macula_identity:generate()),
LinkOpts = #{
identity => Identity,
capabilities => maps:get(capabilities, Opts, 0),
alpn => maps:get(alpn, Opts, [<<"macula">>]),
connect_timeout_ms => maps:get(connect_timeout_ms, Opts, 30_000)
},
DedupWindow = maps:get(dedup_window_ms, Opts, ?DEFAULT_DEDUP_WINDOW_MS),
DedupSweep = maps:get(dedup_sweep_ms, Opts, ?DEFAULT_DEDUP_SWEEP_MS),
Replication = maps:get(replication_factor, Opts, ?DEFAULT_REPLICATION),
DedupTab = macula_client_dedup:new(),
State0 = #state{seeds = Seeds, identity = Identity,
link_opts = LinkOpts, replication = Replication,
dedup_window = DedupWindow, dedup_sweep = DedupSweep,
dedup_tab = DedupTab,
publish_seq = erlang:system_time(microsecond)},
State1 = lists:foldl(fun start_link_for_seed/2, State0, Seeds),
erlang:send_after(DedupSweep, self(), dedup_sweep),
{ok, State1}.
handle_call({publish, Realm, Topic, Payload, _Opts}, From, S) ->
%% Publish only to links that have completed CONNECT/HELLO. A
%% frame sent to a still-handshaking link is dropped on the floor
%% — unlike ADVERTISE, which the link replays on connect — so
%% selecting the first `replication' *spawned* links could report
%% `{error, not_connected}' while other links are healthy. RPC and
%% streams already filter by `is_connected/1'; publish must too.
%%
%% Dispatch via a one-shot worker so concurrent publishes don't
%% serialise through this gen_server (the per-link `publish/4'
%% calls are sync 5s timeouts; under load they pile up at the pool
%% and the harness's `multi_publisher_pubsub' case fails with
%% empty receives).
Targets = connected_link_pids(S),
N = min(length(Targets), S#state.replication),
Selected = lists:sublist(Targets, N),
AllTargets = Targets,
%% One pool-monotone seq per fact, reused across every replicated
%% link so `{publisher, seq}' identifies the fact regardless of
%% which station relayed it.
Seq = S#state.publish_seq,
_ = spawn(fun() ->
Results = [macula_station_link:publish(P, Realm, Topic, Payload, Seq)
|| P <- Selected],
gen_server:reply(From, summarize_publish(Results, AllTargets))
end),
{noreply, S#state{publish_seq = Seq + 1}};
handle_call({subscribe, Realm, Topic, Subscriber, _Opts}, _From, S) ->
SubRef = make_ref(),
Mon = erlang:monitor(process, Subscriber),
Spec = #sub_spec{realm = Realm, topic = Topic,
subscriber = Subscriber, mon = Mon},
Key = {Realm, Topic},
AlreadyTracked = maps:is_key(Key, S#state.topic_index),
NewS = register_sub(SubRef, Spec, S),
issue_wire_subs(AlreadyTracked, Realm, Topic, NewS),
{reply, {ok, SubRef}, NewS};
handle_call({unsubscribe, SubRef}, _From, S) ->
{reply, ok, drop_sub(SubRef, S)};
handle_call({rpc_call, Realm, Procedure, Payload, TimeoutMs}, From, S) ->
%% Worker-spawn so concurrent CALLs don't serialise through the
%% pool gen_server. Each per-link `macula_station_link:call/5'
%% is a sync gen_server:call to the link; with the old
%% `{reply, ..., S}' shape every caller blocked the pool until
%% the link replied, capping concurrent CALL throughput at 1.
Pids = spawned_link_pids(S),
_ = spawn(fun() ->
Reply = call_first_success(Pids, Realm, Procedure, Payload,
TimeoutMs),
gen_server:reply(From, Reply)
end),
{noreply, S};
handle_call({advertise, Realm, Procedure, Handler}, _From,
#state{procs = P} = S) ->
Pids = spawned_link_pids(S),
Reply = fanout_advertise(Pids, Realm, Procedure, Handler),
{reply, Reply, S#state{procs = P#{{Realm, Procedure} => Handler}}};
handle_call({unadvertise, Realm, Procedure}, _From,
#state{procs = P} = S) ->
_ = fanout_unadvertise(spawned_link_pids(S), Realm, Procedure),
{reply, ok, S#state{procs = maps:remove({Realm, Procedure}, P)}};
handle_call({rpc_call_stream, Realm, Procedure, Args, Opts}, From, S) ->
%% Worker-spawn for the same reason as `rpc_call' — the harness's
%% `many_concurrent_streams' fires N parallel `call_stream/4' from
%% separate caller processes; without this each one queued behind
%% the pool gen_server.
Pids = spawned_link_pids(S),
_ = spawn(fun() ->
Reply = stream_first_healthy(Pids, Realm, Procedure, Args, Opts),
gen_server:reply(From, Reply)
end),
{noreply, S};
handle_call({advertise_stream, Realm, Procedure, Mode, Handler}, _From,
#state{stream_procs = SP} = S) ->
Pids = spawned_link_pids(S),
Reply = fanout_advertise_stream(Pids, Realm, Procedure, Mode, Handler),
{reply, Reply,
S#state{stream_procs = SP#{{Realm, Procedure} => {Mode, Handler}}}};
handle_call({unadvertise_stream, Realm, Procedure}, _From,
#state{stream_procs = SP} = S) ->
_ = fanout_unadvertise_stream(spawned_link_pids(S), Realm, Procedure),
{reply, ok,
S#state{stream_procs = maps:remove({Realm, Procedure}, SP)}};
handle_call(status, _From,
#state{seeds = Seeds, links = Links, subs = Subs,
identity = Identity} = S) ->
{Healthy, Failed} = count_link_health(Seeds, Links),
Status = #{
seeds => Seeds,
healthy_links => Healthy,
failed_links => Failed,
self_node_id => macula_identity:public(Identity),
subscriptions => map_size(Subs)
},
{reply, {ok, Status}, S};
handle_call(links, _From, #state{links = Links} = S) ->
{reply, {ok, link_infos(Links)}, S};
handle_call(_Req, _From, S) ->
{reply, {error, unknown_call}, S}.
handle_cast(_Msg, S) -> {noreply, S}.
handle_info({macula_event, _LinkSubRef, Topic, Payload, Meta}, S) ->
Realm = maps:get(realm, Meta, <<0:256>>),
Publisher = maps:get(publisher, Meta),
Seq = maps:get(seq, Meta),
on_inbound_event(macula_client_dedup:check(S#state.dedup_tab,
Realm, Publisher, Seq),
Realm, Topic, Payload, Meta, S);
handle_info({macula_event_gone, _LinkSubRef, _Reason}, S) ->
%% A link torn down its subscription end. Pool will respawn the
%% link via the DOWN handler and replay subs. Don't propagate to
%% local consumers — they see a continuous stream.
{noreply, S};
handle_info({'DOWN', Mon, process, Pid, Reason}, S) ->
on_down(Mon, Pid, Reason, S);
handle_info({respawn_link, Seed}, S) ->
{noreply, on_respawn_link(Seed, S)};
handle_info(dedup_sweep, S) ->
_ = macula_client_dedup:sweep(S#state.dedup_tab, S#state.dedup_window),
erlang:send_after(S#state.dedup_sweep, self(), dedup_sweep),
{noreply, S};
handle_info({'EXIT', _Pid, _Reason}, S) ->
%% Links are linked to us via gen_server:start_link in
%% start_link_for_seed (we trap_exit). The DOWN monitor fires
%% alongside; that path handles cleanup. Drop the EXIT.
{noreply, S};
handle_info(_Other, S) ->
{noreply, S}.
terminate(_Reason, #state{subs = Subs}) ->
%% Notify every subscriber that the pool is gone.
maps:foreach(
fun(SubRef, #sub_spec{subscriber = Pid, mon = Mon}) ->
erlang:demonitor(Mon, [flush]),
Pid ! {macula_event_gone, SubRef, pool_closed}
end, Subs),
ok.
code_change(_OldVsn, S, _Extra) -> {ok, S}.
%%====================================================================
%% Internals — link lifecycle
%%====================================================================
start_link_for_seed(Seed, S) ->
LinkOpts = (S#state.link_opts)#{seed => Seed},
after_link_start(macula_station_link:start_link(LinkOpts), Seed, S).
after_link_start({ok, Pid}, Seed, S) ->
Mon = erlang:monitor(process, Pid),
LinkState = #link_state{seed = Seed, pid = Pid, mon = Mon},
S#state{links = (S#state.links)#{Seed => LinkState}};
after_link_start({error, Reason}, Seed, S) ->
macula_diagnostics:event(<<"_macula.client.link_start_failed">>,
#{seed => Seed, reason => Reason}),
erlang:send_after(?LINK_RESPAWN_DELAY_MS, self(), {respawn_link, Seed}),
Empty = #link_state{seed = Seed, pid = undefined, mon = undefined},
S#state{links = (S#state.links)#{Seed => Empty}}.
spawned_link_pids(#state{links = Links}) ->
[P || #link_state{pid = P} <- maps:values(Links), is_pid(P)].
%% Live links that have completed CONNECT/HELLO. Used by publish,
%% which (unlike advertise) gains nothing from dispatching to a
%% mid-handshake link.
connected_link_pids(#state{} = S) ->
[P || P <- spawned_link_pids(S),
is_process_alive(P),
macula_station_link:is_connected(P)].
%% Surface a one-shot warning when a caller passes V1 multi_relay
%% options that have no V2 equivalent. The opts are silently dropped
%% (V2's `init/1' simply doesn't read them) but the warning gives a
%% caller migrating from V1 a chance to spot the no-op.
warn_legacy_opts(Opts) ->
Stale = [K || K <- ?V1_LEGACY_OPTS, maps:is_key(K, Opts)],
notify_legacy(Stale).
notify_legacy([]) -> ok;
notify_legacy(Keys) ->
logger:notice(
"[macula_client] ignoring V1-only opts ~p — V2 is realm-per-call "
"and one-link-per-seed. See macula:connect/2 docs.", [Keys]),
ok.
%% First-success across the pool's healthy links. Tries each link in
%% turn; the first non-error reply wins. Falls through on
%% per-link errors (timeout, not_connected) so a single dead link does
%% not block the call.
call_first_success([], _Realm, _Proc, _Payload, _Tmo) ->
{error, no_healthy_station};
call_first_success([Pid | Rest], Realm, Proc, Payload, Tmo) ->
next_or_first(macula_station_link:is_connected(Pid),
Pid, Rest, Realm, Proc, Payload, Tmo).
next_or_first(false, _Pid, Rest, Realm, Proc, Payload, Tmo) ->
call_first_success(Rest, Realm, Proc, Payload, Tmo);
next_or_first(true, Pid, Rest, Realm, Proc, Payload, Tmo) ->
keep_or_next(macula_station_link:call(Pid, Realm, Proc, Payload, Tmo),
Rest, Realm, Proc, Payload, Tmo).
keep_or_next({ok, _} = R, _Rest, _Realm, _Proc, _Payload, _Tmo) -> R;
keep_or_next({error, _} = E, [], _Realm, _Proc, _Payload, _Tmo) -> E;
keep_or_next({error, _}, Rest, Realm, Proc, Payload, Tmo) ->
call_first_success(Rest, Realm, Proc, Payload, Tmo).
%% Fan-out advertise: register on every live link. Returns ok if at
%% least one link accepted; per-link errors are logged and discarded.
%%
%% Pre-handshake links MUST receive the call too — `advertise/4' on
%% the link gen_server updates its local `procedures' map regardless
%% of connection state, and `drain_pending_advertises/1' replays that
%% map on the next handshake. Filtering by `is_connected/1' here
%% leaves the link's map out of sync with the pool's intent: a later
%% `unadvertise' that *also* gets filtered (still pre-handshake)
%% never clears the link's map, and the link will silently re-ADVERTISE
%% the dead procedure when it eventually handshakes — the station
%% re-registers a stale entry that nothing in the SDK will ever
%% withdraw.
fanout_advertise([], _Realm, _Proc, _Handler) ->
{error, no_healthy_station};
fanout_advertise(Pids, Realm, Proc, Handler) ->
Results = [safe_link_advertise(P, Realm, Proc, Handler)
|| P <- Pids, is_process_alive(P)],
summarize_advertise([R || R <- Results, R =/= skipped]).
safe_link_advertise(Pid, Realm, Proc, Handler) ->
try macula_station_link:advertise(Pid, Realm, Proc, Handler)
catch _:_ -> skipped
end.
summarize_advertise([]) ->
{error, no_healthy_station};
summarize_advertise(Results) ->
Ok = lists:any(fun(ok) -> true; (_) -> false end, Results),
case Ok of
true -> ok;
false -> {error, all_stations_failed}
end.
%% Fan-out unadvertise: best-effort; ignored errors. The local pool
%% state is dropped regardless so subsequent CALLs surface
%% `unknown_next_peer' from the station.
%%
%% MUST dispatch to every LIVE link (not just connected ones): the
%% link gen_server's `unadvertise' handler clears its local
%% `procedures' map unconditionally, and the wire UNADVERTISE is
%% best-effort inside `maybe_send_unadvertise' (no-op when
%% pre-handshake). Filtering by `is_connected/1' here leaks: a
%% link that was disconnected at unadvertise time keeps the proc in
%% its local map, and on the next handshake `drain_pending_advertises'
%% replays a now-dead ADVERTISE — the station re-registers an entry
%% that the pool already considers withdrawn.
fanout_unadvertise(Pids, Realm, Proc) ->
[_ = safe_link_unadvertise(P, Realm, Proc)
|| P <- Pids, is_process_alive(P)],
ok.
safe_link_unadvertise(Pid, Realm, Proc) ->
try macula_station_link:unadvertise(Pid, Realm, Proc)
catch _:_ -> skipped
end.
%% Sticky-to-link selection for streams. Walk the healthy links in
%% order; the first one that opens cleanly wins. The returned stream
%% pid is bound to that link's `{remote_via_link, _, _}' peer; if
%% the link dies, the stream errors and the caller re-opens.
%% Per-link `{error, not_connected}' (handshake not done) falls
%% through; any other error short-circuits and is returned to the
%% caller, since it likely indicates a real problem (deadline,
%% protocol mismatch) the next link would also hit.
stream_first_healthy([], _Realm, _Proc, _Args, _Opts) ->
{error, no_healthy_station};
stream_first_healthy([Pid | Rest], Realm, Proc, Args, Opts) ->
on_stream_link(macula_station_link:is_connected(Pid),
Pid, Rest, Realm, Proc, Args, Opts).
on_stream_link(false, _Pid, Rest, Realm, Proc, Args, Opts) ->
stream_first_healthy(Rest, Realm, Proc, Args, Opts);
on_stream_link(true, Pid, Rest, Realm, Proc, Args, Opts) ->
keep_or_next_stream(macula_station_link:call_stream(
Pid, Realm, Proc, Args, Opts),
Rest, Realm, Proc, Args, Opts).
keep_or_next_stream({ok, _Stream} = R, _Rest, _Realm, _Proc, _Args, _Opts) ->
R;
keep_or_next_stream({error, not_connected}, Rest, Realm, Proc, Args, Opts) ->
stream_first_healthy(Rest, Realm, Proc, Args, Opts);
keep_or_next_stream({error, _} = E, _Rest, _Realm, _Proc, _Args, _Opts) ->
E.
%% Fan-out streaming advertise across every live link. Same shape
%% as `fanout_advertise/4' for unary; partial success counts. Same
%% rationale for dispatching to pre-handshake links — see the
%% comment on `fanout_advertise/4'.
fanout_advertise_stream([], _Realm, _Proc, _Mode, _Handler) ->
{error, no_healthy_station};
fanout_advertise_stream(Pids, Realm, Proc, Mode, Handler) ->
Results = [safe_link_advertise_stream(P, Realm, Proc, Mode, Handler)
|| P <- Pids, is_process_alive(P)],
summarize_advertise([R || R <- Results, R =/= skipped]).
safe_link_advertise_stream(Pid, Realm, Proc, Mode, Handler) ->
try macula_station_link:advertise_stream(Pid, Realm, Proc, Mode, Handler)
catch _:_ -> skipped
end.
fanout_unadvertise_stream(Pids, Realm, Proc) ->
[_ = safe_link_unadvertise_stream(P, Realm, Proc)
|| P <- Pids, is_process_alive(P)],
ok.
safe_link_unadvertise_stream(Pid, Realm, Proc) ->
try macula_station_link:unadvertise_stream(Pid, Realm, Proc)
catch _:_ -> skipped
end.
%% Count `(healthy, failed)' links across configured seeds. A seed is
%% healthy when its worker pid is alive AND its station_link reports
%% `is_connected'. Anything else (no pid yet, dead pid, mid-handshake)
%% counts as failed. Probes are sequential; cap at 1s per probe via
%% `is_connected/1' so a hung station can't stall the whole
%% `status/1' call past one second per stuck seed.
count_link_health(Seeds, Links) ->
lists:foldl(fun(Seed, Acc) -> tally_seed(maps:find(Seed, Links), Acc) end,
{0, 0}, Seeds).
tally_seed({ok, #link_state{pid = Pid}}, {H, F}) when is_pid(Pid) ->
bump(link_healthy(Pid), H, F);
tally_seed(_, {H, F}) ->
{H, F + 1}.
bump(true, H, F) -> {H + 1, F};
bump(false, H, F) -> {H, F + 1}.
link_healthy(Pid) ->
is_process_alive(Pid) andalso macula_station_link:is_connected(Pid).
%% Build one `link_info()' per spawned link. Skips seeds whose link
%% worker is not (yet) a live pid — those have no addressable station.
link_infos(Links) ->
[link_info(Seed, Pid)
|| {Seed, #link_state{pid = Pid}} <- maps:to_list(Links),
is_pid(Pid)].
link_info(Seed, Pid) ->
Connected = link_healthy(Pid),
#{seed => Seed,
host => seed_host(Seed),
pid => Pid,
connected => Connected,
node_id => link_node_id(Pid, Connected)}.
%% Only probe the peer pubkey on a connected link; a mid-handshake
%% link answers `{error, not_connected}'.
link_node_id(Pid, true) ->
case macula_station_link:peer_node_id(Pid) of
{ok, NodeId} -> NodeId;
{error, not_connected} -> undefined
end;
link_node_id(_Pid, false) ->
undefined.
%% Dial host parsed from a seed. Mirrors `macula_station_link:parse_seed/1'
%% host extraction without re-dialing — URL form or pre-parsed map.
seed_host(#{host := H}) when is_binary(H) -> H;
seed_host(#{host := H}) when is_list(H) -> list_to_binary(H);
seed_host(Url) when is_binary(Url) -> seed_host(binary_to_list(Url));
seed_host(Url) when is_list(Url) ->
case uri_string:parse(Url) of
#{host := H} when H =/= "" -> list_to_binary(H);
_ -> undefined
end;
seed_host(_) ->
undefined.
on_respawn_link(Seed, S) ->
NewS = start_link_for_seed(Seed, S),
replay_to_seed(maps:get(Seed, NewS#state.links, undefined), NewS).
replay_to_seed(#link_state{pid = Pid}, S) when is_pid(Pid) ->
macula_client_replay:subs_to(Pid, S#state.topic_index),
macula_client_replay:advs_to(Pid, S#state.procs),
macula_client_replay:stream_advs_to(Pid, S#state.stream_procs),
S;
replay_to_seed(_, S) ->
S.
%%====================================================================
%% Internals — DOWN routing (link vs subscriber)
%%====================================================================
on_down(Mon, Pid, Reason, S) ->
on_down_routed(find_link_by_mon(Mon, S), Mon, Pid, Reason, S).
on_down_routed({ok, Seed}, _Mon, Pid, Reason, S) ->
macula_diagnostics:event(<<"_macula.client.link_down">>,
#{seed => Seed, pid => Pid, reason => Reason}),
erlang:send_after(?LINK_RESPAWN_DELAY_MS, self(), {respawn_link, Seed}),
{noreply, S#state{links = maps:remove(Seed, S#state.links)}};
on_down_routed(error, Mon, _Pid, _Reason, S) ->
{noreply, on_subscriber_down(Mon, S)}.
find_link_by_mon(Mon, #state{links = Links}) ->
case [Seed || {Seed, #link_state{mon = M}} <- maps:to_list(Links),
M =:= Mon] of
[Seed | _] -> {ok, Seed};
[] -> error
end.
on_subscriber_down(Mon, #state{subs = Subs} = S) ->
Found = [SubRef || {SubRef, #sub_spec{mon = M}}
<- maps:to_list(Subs), M =:= Mon],
lists:foldl(fun drop_sub/2, S, Found).
%%====================================================================
%% Internals — subscription bookkeeping
%%====================================================================
register_sub(SubRef, #sub_spec{realm = R, topic = T} = Spec,
#state{subs = Subs, topic_index = Idx} = S) ->
Key = {R, T},
Set = maps:get(Key, Idx, sets:new()),
NewIdx = Idx#{Key => sets:add_element(SubRef, Set)},
NewSubs = Subs#{SubRef => Spec},
S#state{subs = NewSubs, topic_index = NewIdx}.
drop_sub(SubRef, #state{subs = Subs} = S) ->
drop_sub_take(maps:take(SubRef, Subs), SubRef, S).
drop_sub_take(error, _SubRef, S) ->
S;
drop_sub_take({#sub_spec{realm = R, topic = T, mon = Mon}, NewSubs},
SubRef, #state{topic_index = Idx} = S) ->
erlang:demonitor(Mon, [flush]),
Key = {R, T},
NewSet = sets:del_element(SubRef, maps:get(Key, Idx, sets:new())),
NewIdx = on_index_after_drop(sets:is_empty(NewSet), Key, NewSet, Idx),
S#state{subs = NewSubs, topic_index = NewIdx}.
on_index_after_drop(true, Key, _Set, Idx) -> maps:remove(Key, Idx);
on_index_after_drop(false, Key, Set, Idx) -> Idx#{Key => Set}.
issue_wire_subs(true, _Realm, _Topic, _S) ->
%% A sibling consumer already triggered the wire-level subscribe;
%% the pool fans out to every local SubRef on inbound EVENT.
ok;
issue_wire_subs(false, Realm, Topic, S) ->
PoolPid = self(),
[_ = macula_station_link:subscribe(P, Realm, Topic, PoolPid)
|| P <- spawned_link_pids(S)],
ok.
%%====================================================================
%% Internals — inbound event fan-out
%%====================================================================
on_inbound_event(duplicate, _Realm, _Topic, _Payload, _Meta, S) ->
{noreply, S};
on_inbound_event(new, Realm, Topic, Payload, Meta, S) ->
fan_to_local(Realm, Topic, Payload, Meta, S),
{noreply, S}.
fan_to_local(Realm, Topic, Payload, Meta,
#state{topic_index = Idx, subs = Subs}) ->
fan_to_set(maps:find({Realm, Topic}, Idx), Topic, Payload, Meta, Subs).
fan_to_set(error, _Topic, _Payload, _Meta, _Subs) ->
ok;
fan_to_set({ok, Set}, Topic, Payload, Meta, Subs) ->
sets:fold(fun(SubRef, _) ->
deliver_one(SubRef, Topic, Payload, Meta, Subs)
end, ok, Set).
deliver_one(SubRef, Topic, Payload, Meta, Subs) ->
deliver_to(maps:find(SubRef, Subs), SubRef, Topic, Payload, Meta).
deliver_to(error, _SubRef, _Topic, _Payload, _Meta) ->
ok;
deliver_to({ok, #sub_spec{subscriber = Pid}}, SubRef, Topic, Payload, Meta) ->
Pid ! {macula_event, SubRef, Topic, Payload, Meta},
ok.
%%====================================================================
%% Internals — publish summary
%%====================================================================
summarize_publish([], []) ->
{error, {transient, no_healthy_station}};
summarize_publish([], _NotEmpty) ->
%% Replication factor capped at 0 by config; treat as no-op ok.
ok;
summarize_publish(Results, _Targets) ->
on_publish_results(lists:any(fun(R) -> R =:= ok end, Results), Results).
on_publish_results(true, _Results) -> ok;
on_publish_results(false, [First | _]) -> First;
on_publish_results(false, []) -> {error, no_publish_attempts}.