Packages
macula
3.11.1
7.1.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_station_link.erl
%% @private
%% @doc Per-station link — internal to `macula_client' (the pool).
%%
%% A `macula_station_link' is a `gen_server' that owns one
%% `macula_peering' connection to a single station endpoint. The
%% pool spawns one link per healthy seed and routes operations
%% through them transparently. **Application code should not call
%% `macula_station_link' directly** — use `macula_client' (or the
%% `macula' facade), which handles failover, replication, dedup,
%% and subscription replay.
%%
%% This module is kept publicly accessible for diagnostics and
%% special-case use (e.g. probing a specific station). It is
%% marked `@private' so it does not appear in user-facing
%% documentation indices.
%%
%% Per `PLAN_V2_PARITY' Q6: the per-station worker name is
%% `macula_station_link' (not `macula_station_client') — a station
%% is an identity bound to one IPv6:port; one relay box hosts many
%% stations; a "client" name is taken by the pool above.
%%
%% It drives the CONNECT/HELLO handshake as the client side, then
%% exposes three surfaces over the same peering pipe:
%%
%% <ul>
%% <li><strong>Request/response</strong> — `call/5' sends a CALL
%% frame and matches inbound RESULT/ERROR frames against
%% pending callers using the 16-byte CALL id. Convenience
%% wrappers cover `_dht.put_record', `_dht.find_record', and
%% `_dht.find_records_by_type'.</li>
%% <li><strong>Streaming subscribe</strong> — `subscribe/4' sends
%% a SUBSCRIBE frame and registers a delivery pid. Inbound
%% EVENT frames matching the (realm, topic) fan out to
%% subscribers as
%% `{macula_event, SubRef, Topic, Payload, Meta}'. On
%% disconnect each subscriber receives a single
%% `{macula_event_gone, SubRef, Reason}'.</li>
%% <li><strong>Publish</strong> — `publish/4' sends a PUBLISH
%% frame fire-and-forget. Per-link monotonic `seq' counter
%% stamps each frame for downstream dedup.</li>
%% </ul>
%%
%% == Realm-per-call ==
%%
%% Per `PLAN_V2_PARITY' Q2 sub-decision §2: realm is **per-call**, not
%% connect-time. Stations are realm-agnostic infrastructure; every
%% wire frame carries its own 32-byte `realm' tag. The link advertises
%% an empty realms list in CONNECT and stamps the realm passed to each
%% public op onto the outbound frame.
%%
%% == Lifecycle ==
%%
%% <ol>
%% <li>`start_link/1' — spawn worker, schedule connect.</li>
%% <li>`connect_now/1' (cast) — build connect opts, call
%% `macula_peering:connect/1', store the worker pid.</li>
%% <li>Peering handshake completes → `{macula_peering, connected,
%% Pid, PeerNodeId}' arrives → state moves to `connected'.</li>
%% <li>`call/5' from caller → build CALL frame, sign happens inside
%% peering, store `{from, deadline_timer}` keyed by CALL id, send
%% frame via `macula_peering:send_frame/2'.</li>
%% <li>RESULT or ERROR arrives as `{macula_peering, frame, Pid, Frame}'
%% → look up `call_id', cancel timer, reply to caller.</li>
%% <li>`{macula_peering, disconnected, Pid, Reason}' → fail all
%% pending calls with `{error, {disconnected, Reason}}', notify
%% all subscribers via `macula_event_gone', stop the client
%% (caller is responsible for restart / reconnect).</li>
%% </ol>
%%
%% == Call reply taxonomy ==
%%
%% <table>
%% <tr><th>Inbound frame</th><th>`call/5' returns</th></tr>
%% <tr><td>RESULT(payload=`{error, Reason}')</td><td>`{ok, {error, Reason}}'</td></tr>
%% <tr><td>RESULT(payload=Value)</td><td>`{ok, Value}'</td></tr>
%% <tr><td>ERROR(code=C, name=N)</td><td>`{error, {call_error, C, N}}'</td></tr>
%% <tr><td>(deadline elapses)</td><td>`{error, timeout}'</td></tr>
%% <tr><td>(connection drops)</td><td>`{error, {disconnected, Reason}}'</td></tr>
%% </table>
-module(macula_station_link).
-behaviour(gen_server).
-export([
start_link/1,
stop/1,
call/5,
publish/4,
put_record/2, put_record/3,
find_record/2, find_record/3,
find_records_by_type/2, find_records_by_type/3,
subscribe/4,
unsubscribe/2,
is_connected/1,
peer_node_id/1
]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2,
terminate/2, code_change/3]).
-export_type([opts/0]).
-type url() :: binary() | string().
-type opts() :: #{
%% Endpoint to dial. Either a URL (https://host:port) or a
%% pre-parsed #{host, port} map.
seed := url() | #{host := binary() | string(),
port := inet:port_number()},
%% Local Ed25519 keypair used to sign the CONNECT frame and any
%% subsequent application frames. Auto-generated when absent.
identity => macula_identity:key_pair(),
%% Capability bitfield announced in CONNECT (default 0).
capabilities => non_neg_integer(),
%% ALPN list passed through to QUIC (default [<<"macula">>]).
alpn => [binary()],
%% Connect timeout in ms (default 30_000).
connect_timeout_ms => non_neg_integer()
}.
-define(DHT_REALM, <<0:256>>).
-define(DEFAULT_DEADLINE_MS, 5_000).
-define(CONNECT_RETRY_BACKOFF_MS, 1_000).
-record(state, {
seed :: #{host := binary() | string(),
port := inet:port_number()},
identity :: macula_identity:key_pair(),
capabilities :: non_neg_integer(),
alpn :: [binary()],
connect_timeout_ms :: non_neg_integer(),
%% peering worker pid (`macula_peering_conn`). undefined while
%% disconnected.
peer_pid :: pid() | undefined,
%% peer's node id, set on `connected'.
peer_node_id :: macula_identity:pubkey() | undefined,
%% map of CALL id (16 bytes) -> {From, TimerRef}.
pending = #{} :: #{<<_:128>> => {gen_server:from(), reference()}},
%% Active topic subscriptions keyed by SubRef returned to the
%% subscriber. The reverse `topic_index' lets inbound EVENT
%% frames fan out to all SubRefs subscribed to a given
%% (realm, topic) without scanning the whole subscriptions map.
subscriptions = #{} :: #{reference() => subscription()},
topic_index = #{} :: #{{<<_:256>>, binary()} => sets:set(reference())},
%% Monotonic per-link publish sequence (stamps outbound PUBLISH
%% frames). Resets on link respawn — pool dedup absorbs the gap.
publish_seq = 0 :: non_neg_integer()
}).
-type subscription() :: {Realm :: <<_:256>>,
Topic :: binary(),
Subscriber :: pid(),
Mon :: reference()}.
%%====================================================================
%% Public API
%%====================================================================
%% @doc Start a station-client connected to `seed'.
%% Returns once the gen_server is alive; the QUIC handshake completes
%% asynchronously. Use `is_connected/1' to poll readiness or just
%% issue `call/5' (which blocks the caller until ready or until its
%% timeout elapses).
-spec start_link(opts()) -> {ok, pid()} | {error, term()}.
start_link(Opts) when is_map(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Issue a CALL frame and block until the station replies, the
%% deadline elapses, or the connection drops.
%%
%% `Realm' is the 32-byte realm id stamped on the outbound CALL frame.
%% Stations are realm-agnostic infrastructure; the realm is carried
%% per-frame so a single link can multiplex many realms.
%%
%% `Procedure' is the V2 procedure name, e.g.
%% `<<"_dht.find_records_by_type">>'. `Payload' is any term that
%% `macula_frame:call/1' accepts (typically a map).
-spec call(pid(), <<_:256>>, binary(), term(), pos_integer()) ->
{ok, term()} | {error, term()}.
call(Pid, Realm, Procedure, Payload, TimeoutMs)
when is_pid(Pid),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Procedure),
is_integer(TimeoutMs), TimeoutMs > 0 ->
%% gen_server timeout = TimeoutMs + 500 to give the server time to
%% report a clean `{error, timeout}' rather than the caller seeing
%% a hard `exit({timeout, ...})'.
GenTimeout = TimeoutMs + 500,
try
gen_server:call(Pid, {call, Realm, Procedure, Payload, TimeoutMs},
GenTimeout)
catch
%% try/catch retained: collapses the three distinct gen_server
%% exit signals into the SDK's call-result taxonomy. Without
%% it the caller sees `exit({timeout, _})' instead of
%% `{error, timeout}', breaking the contract documented above.
exit:{timeout, _} -> {error, timeout};
exit:{noproc, _} -> {error, noproc};
exit:{normal, _} -> {error, gone}
end.
%% @doc Send a PUBLISH frame fire-and-forget. The link stamps a
%% monotonic per-link `seq' onto the frame and the local
%% `published_at_ms' clock; the station relays it to subscribers.
%%
%% Returns `ok' once the frame is on the wire, `{error, not_connected}'
%% when the link has not yet completed the QUIC handshake. Publishes
%% are NOT queued during disconnect — they would arrive at the wrong
%% wall-clock and fight pool-level dedup. The pool retries on a peer
%% link instead.
-spec publish(pid(), <<_:256>>, binary(), term()) ->
ok | {error, not_connected | term()}.
publish(Pid, Realm, Topic, Payload)
when is_pid(Pid),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Topic) ->
gen_server:call(Pid, {publish, Realm, Topic, Payload}, 5_000).
%% @doc Convenience wrapper for `_dht.put_record'. The record must be
%% a fully-signed `macula_record:record()' map (build via
%% `macula_record:envelope/3,4' + `macula_record:sign/2'). Returns
%% `ok' on success, `{error, Reason}' on RPC failure or unexpected
%% reply.
%%
%% Stations replicate the put across the K-nearest peers in their
%% Kademlia routing table, so a single `put_record/2' call against
%% any one connected station propagates to the rest of the DHT.
%%
%% DHT-internal procedures travel under the all-zeros realm tag —
%% they are protocol-internal, not bound to any business realm.
-spec put_record(pid(), map()) -> ok | {error, term()}.
put_record(Pid, Record) ->
put_record(Pid, Record, ?DEFAULT_DEADLINE_MS).
-spec put_record(pid(), map(), pos_integer()) -> ok | {error, term()}.
put_record(Pid, Record, TimeoutMs) when is_pid(Pid), is_map(Record) ->
classify_put(call(Pid, ?DHT_REALM, <<"_dht.put_record">>,
Record, TimeoutMs)).
classify_put({ok, ok}) -> ok;
classify_put({ok, Other}) -> {error, {unexpected_reply, Other}};
classify_put({error, _} = E) -> E.
%% @doc Convenience wrapper for `_dht.find_record'. Looks up a record
%% by its `macula_record:storage_key/1' (32-byte BLAKE3 digest).
%% Returns `{error, not_found}' when no record exists at the key.
%% Callers SHOULD verify the returned record's signature with
%% `macula_record:verify/1' before trusting its payload.
-spec find_record(pid(), <<_:256>>) ->
{ok, map()} | {error, not_found | term()}.
find_record(Pid, Key) ->
find_record(Pid, Key, ?DEFAULT_DEADLINE_MS).
-spec find_record(pid(), <<_:256>>, pos_integer()) ->
{ok, map()} | {error, not_found | term()}.
find_record(Pid, Key, TimeoutMs)
when is_pid(Pid), is_binary(Key), byte_size(Key) =:= 32 ->
classify_find(call(Pid, ?DHT_REALM, <<"_dht.find_record">>,
#{key => Key}, TimeoutMs)).
classify_find({ok, #{type := _, payload := _, sig := _} = R}) -> {ok, R};
classify_find({ok, not_found}) -> {error, not_found};
classify_find({ok, Other}) -> {error, {unexpected_reply, Other}};
classify_find({error, _} = E) -> E.
%% @doc Convenience wrapper for `_dht.find_records_by_type'. Returns
%% the decoded list of signed records (CBOR-decoded maps as produced
%% by `macula_record').
-spec find_records_by_type(pid(), 0..255) ->
{ok, [map()]} | {error, term()}.
find_records_by_type(Pid, Type) ->
find_records_by_type(Pid, Type, ?DEFAULT_DEADLINE_MS).
-spec find_records_by_type(pid(), 0..255, pos_integer()) ->
{ok, [map()]} | {error, term()}.
find_records_by_type(Pid, Type, TimeoutMs)
when is_integer(Type), Type >= 0, Type =< 255 ->
classify_records(call(Pid, ?DHT_REALM, <<"_dht.find_records_by_type">>,
#{type => Type}, TimeoutMs)).
classify_records({ok, Records}) when is_list(Records) -> {ok, Records};
classify_records({ok, Other}) -> {error, {unexpected_reply, Other}};
classify_records({error, _} = E) -> E.
%% @doc Subscribe to a peering pubsub topic in `Realm'. Sends a
%% SUBSCRIBE frame to the connected station and registers
%% `Subscriber' as the delivery pid for inbound EVENT frames matching
%% `(Realm, Topic)'.
%%
%% Returns `{ok, SubRef}' once the SUBSCRIBE frame is sent (or queued
%% if the peering handshake has not yet completed — drained on
%% `connected'). Stations do not acknowledge SUBSCRIBE — the contract
%% is best-effort, mirroring the existing peering pubsub semantics.
%%
%% Subscriber receives one of:
%%
%% <ul>
%% <li>`{macula_event, SubRef, Topic, Payload, Meta}' — every time
%% an EVENT frame arrives for `(Realm, Topic)'. `Meta' is a map
%% with `realm', `publisher', `seq', and `delivered_via'
%% fields.</li>
%% <li>`{macula_event_gone, SubRef, Reason}' — once, when the
%% connection drops or the client stops. The subscription map
%% is cleared on the same transition.</li>
%% </ul>
%%
%% The client monitors `Subscriber'; if it dies the subscription is
%% torn down (best-effort UNSUBSCRIBE on the wire).
-spec subscribe(pid(), <<_:256>>, binary(), pid()) ->
{ok, reference()}.
subscribe(Client, Realm, Topic, Subscriber)
when is_pid(Client),
is_binary(Realm), byte_size(Realm) =:= 32,
is_binary(Topic), is_pid(Subscriber) ->
gen_server:call(Client, {subscribe, Realm, Topic, Subscriber}, 5_000).
%% @doc Drop a subscription. Sends a best-effort UNSUBSCRIBE frame
%% to the station and clears local bookkeeping. Always returns `ok',
%% even when `SubRef' is unknown — unsubscribe is idempotent.
-spec unsubscribe(pid(), reference()) -> ok.
unsubscribe(Client, SubRef)
when is_pid(Client), is_reference(SubRef) ->
gen_server:call(Client, {unsubscribe, SubRef}, 5_000).
-spec is_connected(pid()) -> boolean().
is_connected(Pid) ->
case gen_server:call(Pid, is_connected, 1_000) of
true -> true;
false -> false
end.
-spec peer_node_id(pid()) -> {ok, macula_identity:pubkey()} | {error, not_connected}.
peer_node_id(Pid) ->
gen_server:call(Pid, peer_node_id, 1_000).
%%====================================================================
%% gen_server
%%====================================================================
init(Opts) ->
Seed = parse_seed(maps:get(seed, Opts)),
Identity = maps:get(identity, Opts, macula_identity:generate()),
Caps = maps:get(capabilities, Opts, 0),
Alpn = maps:get(alpn, Opts, [<<"macula">>]),
Tmo = maps:get(connect_timeout_ms, Opts, 30_000),
State = #state{seed = Seed, identity = Identity,
capabilities = Caps, alpn = Alpn,
connect_timeout_ms = Tmo},
process_flag(trap_exit, true),
self() ! attempt_connect,
{ok, State}.
handle_call({call, _Realm, _Proc, _Payload, _Tmo}, _From,
#state{peer_pid = undefined} = S) ->
{reply, {error, not_connected}, S};
handle_call({call, Realm, Proc, Payload, Tmo}, From,
#state{peer_pid = Pid, identity = Id, pending = P} = S) ->
CallId = crypto:strong_rand_bytes(16),
Caller = macula_identity:public(Id),
DeadlineMs = erlang:system_time(millisecond) + Tmo,
Frame = macula_frame:call(#{
call_id => CallId,
procedure => Proc,
realm => Realm,
payload => Payload,
deadline_ms => DeadlineMs,
caller => Caller
}),
ok = macula_peering:send_frame(Pid, Frame),
TRef = erlang:send_after(Tmo, self(), {call_timeout, CallId}),
{noreply, S#state{pending = P#{CallId => {From, TRef}}}};
handle_call({publish, _Realm, _Topic, _Payload}, _From,
#state{peer_node_id = undefined} = S) ->
%% Require the full HELLO handshake before publishing — the
%% peering worker may exist mid-handshake while the wire is not
%% yet ready for application frames. Matches `is_connected/1'.
{reply, {error, not_connected}, S};
handle_call({publish, Realm, Topic, Payload}, _From,
#state{peer_pid = Pid, identity = Id,
publish_seq = Seq} = S) ->
Pub = macula_identity:public(Id),
Frame = macula_frame:publish(#{
topic => Topic,
realm => Realm,
publisher => Pub,
seq => Seq,
payload => Payload,
published_at_ms => erlang:system_time(millisecond)
}),
ok = macula_peering:send_frame(Pid, Frame),
{reply, ok, S#state{publish_seq = Seq + 1}};
handle_call(is_connected, _From, #state{peer_pid = undefined} = S) ->
{reply, false, S};
handle_call(is_connected, _From, #state{peer_node_id = undefined} = S) ->
{reply, false, S};
handle_call(is_connected, _From, S) ->
{reply, true, S};
handle_call(peer_node_id, _From, #state{peer_node_id = undefined} = S) ->
{reply, {error, not_connected}, S};
handle_call(peer_node_id, _From, #state{peer_node_id = Id} = S) ->
{reply, {ok, Id}, S};
handle_call({subscribe, Realm, Topic, Subscriber}, _From,
#state{subscriptions = Subs, topic_index = Idx} = S) ->
SubRef = make_ref(),
Mon = erlang:monitor(process, Subscriber),
NewSubs = Subs#{SubRef => {Realm, Topic, Subscriber, Mon}},
NewIdx = add_topic_sub(Realm, Topic, SubRef, Idx),
%% Send the SUBSCRIBE frame now if peering is up; otherwise the
%% `connected' handler drains every stored subscription on
%% handshake completion. Avoids the race where a consumer calls
%% `subscribe/4' immediately after `start_link/1' before the
%% peering CONNECT/HELLO has finished — the SUBSCRIBE used to
%% return `{error, not_connected}' and silently never land on
%% the wire even though the client became connected milliseconds
%% later.
maybe_send_subscribe(Realm, Topic, S),
{reply, {ok, SubRef}, S#state{subscriptions = NewSubs,
topic_index = NewIdx}};
handle_call({unsubscribe, SubRef}, _From, S) ->
{reply, ok, on_unsubscribe(SubRef, S)};
handle_call(_Req, _From, S) ->
{reply, {error, unknown_call}, S}.
handle_cast(_Msg, S) -> {noreply, S}.
%%-------------------------------------------------------------------
%% Connect
%%-------------------------------------------------------------------
handle_info(attempt_connect, #state{seed = Seed, identity = Id,
capabilities = Caps, alpn = Alpn,
connect_timeout_ms = Tmo} = S) ->
Pub = macula_identity:public(Id),
PeeringOpts = #{
role => client,
target => Seed#{alpn => Alpn, timeout_ms => Tmo},
node_id => Pub,
identity => Id,
%% Realm-agnostic: the link advertises no realm membership.
%% Each frame carries its own realm tag.
realms => [],
capabilities => Caps,
controlling_pid => self()
},
after_connect_request(macula_peering:connect(PeeringOpts), S);
handle_info({macula_peering, connected, Pid, PeerNodeId},
#state{peer_pid = Pid} = S) ->
NewS = S#state{peer_node_id = PeerNodeId},
drain_pending_subscribes(NewS),
{noreply, NewS};
handle_info({macula_peering, frame, Pid, Frame},
#state{peer_pid = Pid} = S) ->
{noreply, on_frame(Frame, S)};
handle_info({macula_peering, disconnected, Pid, Reason},
#state{peer_pid = Pid} = S) ->
NewS = fail_all_pending({disconnected, Reason}, S),
%% Stop normally — the supervisor (or owning gen_server) decides
%% whether to restart us.
{stop, normal, NewS#state{peer_pid = undefined,
peer_node_id = undefined}};
handle_info({call_timeout, CallId}, #state{pending = P} = S) ->
on_timeout(maps:take(CallId, P), S);
handle_info({'EXIT', Pid, Reason}, #state{peer_pid = Pid} = S) ->
NewS = fail_all_pending({peering_exit, Reason}, S),
{stop, normal, NewS#state{peer_pid = undefined,
peer_node_id = undefined}};
handle_info({'DOWN', Mon, process, _Pid, _Reason}, S) ->
{noreply, on_subscriber_down(Mon, S)};
handle_info(_Other, S) ->
{noreply, S}.
terminate(_Reason, #state{peer_pid = Pid}) when is_pid(Pid) ->
catch macula_peering:close(Pid, client_stop),
ok;
terminate(_Reason, _S) ->
ok.
code_change(_OldVsn, S, _Extra) -> {ok, S}.
%%====================================================================
%% Internals
%%====================================================================
after_connect_request({ok, Pid}, S) ->
link(Pid),
{noreply, S#state{peer_pid = Pid}};
after_connect_request({error, Reason}, S) ->
macula_diagnostics:event(<<"_macula.station_link.connect_failed">>, #{
reason => Reason,
seed => S#state.seed
}),
erlang:send_after(?CONNECT_RETRY_BACKOFF_MS, self(), attempt_connect),
{noreply, S}.
%% RESULT
on_frame(#{frame_type := result, call_id := CallId, payload := Payload},
#state{pending = P} = S) ->
deliver_pending(maps:take(CallId, P), {ok, Payload}, S);
%% ERROR
on_frame(#{frame_type := error, call_id := CallId} = Frame,
#state{pending = P} = S) ->
Code = maps:get(code, Frame, 0),
Name = maps:get(name, Frame, undefined),
deliver_pending(maps:take(CallId, P),
{error, {call_error, Code, Name}}, S);
%% EVENT — pubsub delivery. Fan out to every subscriber whose
%% (realm, topic) matches. Stations may push EVENTs without a prior
%% SUBSCRIBE on this connection (e.g. wildcard / catalog channels);
%% silently drop those.
on_frame(#{frame_type := event, topic := Topic, realm := Realm} = Frame, S) ->
deliver_event(Realm, Topic, Frame, S);
%% HyParView / Plumtree / SWIM / content frames pass through here.
%% This client cares only about call/result/error and event; the
%% rest is for dedicated overlay modules.
on_frame(_Frame, S) ->
S.
deliver_pending(error, _Reply, S) ->
%% Unknown call_id (race with timeout, or duplicate reply).
S;
deliver_pending({{From, TRef}, NewP}, Reply, S) ->
_ = erlang:cancel_timer(TRef),
gen_server:reply(From, Reply),
S#state{pending = NewP}.
on_timeout(error, S) ->
{noreply, S};
on_timeout({{From, _OldTRef}, NewP}, S) ->
gen_server:reply(From, {error, timeout}),
{noreply, S#state{pending = NewP}}.
fail_all_pending(Reason, #state{pending = P, subscriptions = Subs} = S) ->
maps:foreach(fun(_CallId, {From, TRef}) ->
_ = erlang:cancel_timer(TRef),
gen_server:reply(From, {error, Reason})
end, P),
maps:foreach(fun(SubRef, {_Realm, _Topic, Subscriber, Mon}) ->
erlang:demonitor(Mon, [flush]),
Subscriber ! {macula_event_gone, SubRef, Reason}
end, Subs),
S#state{pending = #{}, subscriptions = #{}, topic_index = #{}}.
%%-------------------------------------------------------------------
%% Subscription helpers
%%-------------------------------------------------------------------
add_topic_sub(Realm, Topic, SubRef, Idx) ->
Key = {Realm, Topic},
Set = maps:get(Key, Idx, sets:new()),
Idx#{Key => sets:add_element(SubRef, Set)}.
del_topic_sub(Realm, Topic, SubRef, Idx) ->
Key = {Realm, Topic},
on_set_after_del(Key, sets:del_element(SubRef, maps:get(Key, Idx, sets:new())), Idx).
on_set_after_del(Key, Set, Idx) ->
on_empty_set(sets:is_empty(Set), Key, Set, Idx).
on_empty_set(true, Key, _Set, Idx) -> maps:remove(Key, Idx);
on_empty_set(false, Key, Set, Idx) -> Idx#{Key => Set}.
%% Drop a single subscription. Best-effort UNSUBSCRIBE on the wire
%% (drops silently when disconnected — the station prunes stale
%% subscribers eventually). Idempotent: unknown SubRef is a no-op.
on_unsubscribe(SubRef, #state{subscriptions = Subs,
topic_index = Idx,
peer_pid = Pid,
identity = Id} = S) ->
on_unsubscribe_take(maps:take(SubRef, Subs), SubRef, Idx, Pid, Id, S).
on_unsubscribe_take(error, _SubRef, _Idx, _Pid, _Id, S) ->
S;
on_unsubscribe_take({{Realm, Topic, _Subscriber, Mon}, NewSubs},
SubRef, Idx, Pid, Id, S) ->
erlang:demonitor(Mon, [flush]),
NewIdx = del_topic_sub(Realm, Topic, SubRef, Idx),
send_unsubscribe(Pid, Realm, Topic, Id),
S#state{subscriptions = NewSubs, topic_index = NewIdx}.
send_unsubscribe(undefined, _Realm, _Topic, _Id) ->
ok;
send_unsubscribe(Pid, Realm, Topic, Id) ->
SubKey = macula_identity:public(Id),
Frame = macula_frame:unsubscribe(#{topic => Topic,
realm => Realm,
subscriber => SubKey}),
catch macula_peering:send_frame(Pid, Frame),
ok.
%% Send a SUBSCRIBE frame for `(Realm, Topic)' iff peering is connected.
maybe_send_subscribe(_Realm, _Topic, #state{peer_pid = undefined}) ->
ok;
maybe_send_subscribe(Realm, Topic, #state{peer_pid = Pid, identity = Id}) ->
SubKey = macula_identity:public(Id),
Frame = macula_frame:subscribe(#{topic => Topic,
realm => Realm,
subscriber => SubKey}),
catch macula_peering:send_frame(Pid, Frame),
ok.
%% On handshake completion, send a SUBSCRIBE frame for every stored
%% subscription. Subscribers that came in before connect have been
%% sitting in `subscriptions' with no wire frame yet sent — drain
%% them now. De-duplicate by `(Realm, Topic)' since multiple local
%% SubRefs may share the same wire-level subscription (one SUBSCRIBE
%% frame per identity per (realm, topic), reused across consumers).
drain_pending_subscribes(#state{subscriptions = Subs} = S) ->
Pairs = lists:usort(
[{R, T} || {_Ref, {R, T, _Sub, _Mon}} <- maps:to_list(Subs)]),
[maybe_send_subscribe(R, T, S) || {R, T} <- Pairs],
ok.
%% Subscriber pid died — find its SubRef(s) by monitor ref, drop
%% them. A pid can only have one subscription via one monitor, but
%% scan defensively.
on_subscriber_down(Mon, #state{subscriptions = Subs} = S) ->
Found = maps:fold(fun
(SubRef, {_R, _T, _P, M}, Acc) when M =:= Mon -> [SubRef | Acc];
(_, _, Acc) -> Acc
end, [], Subs),
lists:foldl(fun on_unsubscribe/2, S, Found).
%% Fan an EVENT frame out to every subscriber for that (realm, topic).
deliver_event(Realm, Topic, Frame, #state{topic_index = Idx} = S) ->
deliver_event_to(maps:find({Realm, Topic}, Idx), Realm, Topic, Frame, S),
S.
deliver_event_to(error, _Realm, _Topic, _Frame, _S) ->
ok;
deliver_event_to({ok, Set}, Realm, Topic, Frame, #state{subscriptions = Subs}) ->
Payload = maps:get(payload, Frame),
Meta = #{realm => Realm,
publisher => maps:get(publisher, Frame),
seq => maps:get(seq, Frame),
delivered_via => maps:get(delivered_via, Frame, direct)},
sets:fold(fun(SubRef, _) ->
deliver_event_one(SubRef, Topic, Payload, Meta, Subs)
end, ok, Set).
deliver_event_one(SubRef, Topic, Payload, Meta, Subs) ->
fan_event(maps:find(SubRef, Subs), SubRef, Topic, Payload, Meta).
fan_event(error, _SubRef, _Topic, _Payload, _Meta) ->
ok;
fan_event({ok, {_R, _T, Subscriber, _Mon}}, SubRef, Topic, Payload, Meta) ->
Subscriber ! {macula_event, SubRef, Topic, Payload, Meta},
ok.
%%-------------------------------------------------------------------
%% Helpers
%%-------------------------------------------------------------------
parse_seed(#{host := _, port := _} = Map) ->
Map;
parse_seed(Url) when is_binary(Url) ->
parse_seed(binary_to_list(Url));
parse_seed(Url) when is_list(Url) ->
case uri_string:parse(Url) of
#{host := H, port := P} when is_integer(P) ->
#{host => list_to_binary(H), port => P};
#{host := H, scheme := "https"} ->
#{host => list_to_binary(H), port => 4433};
_ ->
error({invalid_seed_url, Url})
end.