macula (macula v13.2.2)

View Source

Macula SDK — Public API for mesh applications.

This is the main entry point for applications using the Macula SDK.

Apps connect via connect/2, which returns a macula_client pool that internally wraps N peering links to N stations. publish/4,5, subscribe/4,5, unsubscribe/2, call/5,6, providers/3,4, advertise/5, unadvertise/3, call_stream/5, advertise_stream/5, and unadvertise_stream/3 route through the pool with realm-per-call semantics. See macula_pubsub for the slice module of the publish/subscribe surface.

LOCAL streaming (call_stream/2,3, open_stream/3,4, advertise_stream/2,3, unadvertise_stream/1) dispatches in-process via macula_stream_local — for unit tests and same-BEAM pairs.

Erlang distribution over the mesh ships via join_mesh/1 (V2 pool carrier) or join_dist_relay/1 (dedicated dist relay). See macula_dist_pool / macula_dist_relay_client.

Summary

Functions

Abort the stream with an error frame.

Register a procedure handler on a V2 pool and advertise it: the pool resolves its own D25 provider authorization — the realm-signed org_directory and the org-signed procedure_delegation that names the pool's node id, both fetched from the DHT — verifies it against the pinned realm key, and hands every link the advertisement spec: for each link the pool signs an advertisement naming the station that link is connected to as serving_station, sent as an ADVERTISE frame (signed again on reconnect and on link respawn, never past the chain's earlier expiry). The pool renews the chain at a third of its remaining life (macula_client:advertise/8, D32): a delegation lives 30 minutes, so a provider stays advertised while the realm keeps reissuing, and drops out once it stops. The frame registers the procedure at that station; it does not put a record in the DHT, which is what a direct-dialling caller resolves (macula_response:advertise_direct/6,7 publishes that one). The procedure must carry an org namespace, and the pool must run a provisioned identity whose delegation the org has published, and pin the realm's key (realm_trust at connect); a missing piece fails fast with {error, {provider_authorization, _}}.

Advertise a LOCAL in-process streaming procedure (default: server_stream).

Advertise a LOCAL in-process streaming procedure with mode.

Register a streaming procedure handler on a V2 pool. Fans out to every healthy link and stores in pool state for replay on link respawn. A caller reaches this provider only through a procedure_advertisement record that names it; registering the handler publishes none. See macula_client:advertise_stream/5.

As advertise_stream/5, with Opts. auth sets the streaming procedure's policy, the same set advertise/5 takes: open (default), {ucan_required, IssuerNodeId} or {realm_member_required, RealmKeyId, RequiredCan}. A consumer presents its token with call_stream/5's ucan_token opt. stations registers it on the links to those stations only, as advertise/5 does.

Wait for the terminal reply (client-stream / bidi).

Call Procedure in Realm at the provider that serves it: resolve the procedure's verified advertisements, authorized against the realm key the pool pinned for Realm, reach the station a candidate names directly, and call its provider there. See macula_direct_dial:call/5. A procedure the pool's linked stations serve themselves, such as _dht.*, goes through macula_client:call_linked_station/5.

As call/5, with Opts. #{provider => NodeId} calls THAT provider of Procedure and no other: the procedure is resolved and its advertisements trust-checked exactly as call/5 does, and only the named provider's are tried. A provider with no trusted advertisement by the deadline is {error, {unresolved, provider_not_advertised}}. To have every provider of a procedure answer, list them with providers/3,4 and make one call per provider. See macula_direct_dial:call/6.

Issue a CALL to Target, a provider's node_id, at ONE specific station, dialing it directly even if it is not in the pool's seed set. Station is a seed URL (e.g. <<"quic://[::1]:4433">>). The pool reuses an existing link or dials and monitors a new one, waits for the handshake, and calls there. This is the direct-dial data path: resolve a provider's serving_station and its endpoint, then reach it in one hop. See macula_client:call_station/11. With no options it carries no signed state to decide on sealing from, so it is refused as {error, {confidentiality, no_signed_state}} (macula 13.0.0): use call_station/8 with advertisement, confidential => required or confidential => off (see call_seal/5).

As call_station/7, presenting a capability token to a gated provider via Opts (#{ucan_token => Token}). Empty/absent = none. Slice 7b dual-trust. Opts also carries the station this dial must prove, expected_node_id (see macula_client:call_station/11), and may set dial_timeout_ms, how much of TimeoutMs the wait for a fresh link's handshake may take (default: all of it). pin_tls_cert => true is REFUSED with {error, {refused, {pin_tls_cert, no_pin_primitive_for_mldsa87_identity}}}, and verify in any value with {error, {refused, {verify, one_verification_mode}}} (see trust_options_checked/1).

Open a LOCAL in-process server-stream call. Used for unit tests and same-BEAM dispatch via macula_stream_local.

Open a LOCAL in-process server-stream call with options.

Open a streaming RPC to Procedure's provider in Realm: resolve the provider through its procedure_advertisement and open the stream at its serving station, naming the provider as the target, as call/5 does for a single-reply call. Same as macula_direct_dial:call_stream/5. The returned stream is bound to that station's link (errors with peer_down if the link dies; caller re-opens). Optsucan_token presents a UCAN to a streaming procedure advertised with an auth policy (see advertise_stream/6). An open whose signed STREAM_OPEN would be longer than max_stream_open_bytes (1 MiB by default) returns {error, {open_too_large, Limit}} without sending anything. The stream is sealed as call/6 seals a call: confidential => required fails closed without a key, and confidential => off is refused as {error, {confidentiality, off_needs_explicit_target}} (see call_stream_station/7).

Open a streaming RPC to Target, a provider's node_id, by DIALING a specific station directly (direct-dial): the streaming analogue of call_station/7. Compose it with DHT resolution (find_records -> read_procedure_advertisement -> station_endpoint) to reach a stream provider in one hop, exactly as a unary caller does. Opts may set dial_timeout_ms (default 10_000) and a mode. Opts also names the station this dial must prove, expected_node_id (see macula_client:call_station/11); pin_tls_cert => true and verify are REFUSED as call_station/8 describes.

OTP child spec to drop a V2 pool into a caller's supervision tree. Give node_identity as a loader {Module, Function, Args} that returns {ok, Key}. Its Args say where the key is and never hold the key, because a supervisor logs them when a start fails. See macula_client:child_spec/3.

Stop a V2 pool. Every subscriber receives a final {macula_event_gone, SubRef, pool_closed} message.

Half-close the write side; recv still drains.

Close a V1 stream (both sides). Renamed from close/1 in 3.11.0 because close/1 now refers to the V2 pool surface.

Connect to the Macula relay mesh and return a pool handle.

The dist relay client that join_dist_relay/1 started, if it is running.

Ensure this node is running in distributed mode.

A pool link to Station, for a consumer that must address one station through one link, as an overlay does: a live link the pool already holds to that station, else one it dials. The station is pinned by the expected_node_id in its seed map or in Opts, which a link requires (D5, D16); pin_tls_cert => true and verify are refused as call_station/8 describes. Answers {ok, Link} once the link's handshake has completed, within TimeoutMs, or {error, not_connected}. The pool owns and monitors the link: it respawns it and ends it with the pool, so a consumer holds the pid and never stops it. Drive it with overlay_subscribe/3 and send_overlay_frame/2,3.

A field of a map a peer supplied, or undefined when it is absent. See field/3.

A field of a map a peer supplied (D26), or Default when it is absent. A map from the codec carries its text keys as {text, Bin}; a map handed over in process may carry atom or binary keys. The lookup tries {text, Name}, then the atom, then the binary, so a handler reads both kinds of map the same way. Looking up a binary name never creates an atom.

Fetch a record from the mesh DHT by its macula_record:storage_key/1.

As find_record/2, waiting at most TimeoutMs for the reply, for a caller that bounds its work by a deadline of its own.

Fetch EVERY record stored at Key — the full multi-value set, e.g. every procedure_advertisement under one procedure's storage key. Where find_record/2 returns the first record (or not_found), this returns the whole list, empty when none.

As find_records/2, waiting at most TimeoutMs for the reply, for a caller that bounds its work by a deadline of its own.

Return every record of a given type currently visible from the pool's connected stations, each verified under the node's crypto profile; a record that does not verify is dropped.

Fetch MCID in Realm, as get_content/4 with no options.

Fetch MCID from a node that shares it in Realm (D27): find its announcements, open a stream to a sharer through the station it announced, and verify every block and the manifest against their content ids (macula_content_fetch). A sharer that fails moves the fetch to the next; {error, {unavailable, [{Sharer, Reason}]}} when all fail, {error, not_shared} when none announces it, {error, invalid_mcid} for an id that is not tag 2. Opts: max_bytes (256 MiB), max_chunks, chunk_timeout_ms, parallel. No realm key is needed: content verifies itself.

Enable Erlang distribution over a dedicated dist relay (macula-io/macula-dist-relay).

Join the Macula relay mesh with Erlang distribution.

Per-link snapshot of a V2 pool — one entry per spawned link with its peer station node_id (pubkey), dial host, pid, and connected flag. Use this to resolve a specific station (by pubkey or hostname) to its link for targeted, per-station operations. See macula_client:links/1 for the link_info() shape.

Subscribe to node up/down events.

Open a LOCAL in-process client-stream or bidi call. Used for unit tests and same-BEAM dispatch via macula_stream_local.

Open a LOCAL in-process stream with explicit mode.

Take the overlay frames for Realm arriving on Link (from ensure_station_link/4), delivered to Subscriber as {macula_overlay_frame, SubRef, Frame, Meta}, where Meta names the frame's authenticated sender and, for a relayed frame, the station it came via. A frame from the connected peer, or a GOSSIP, for a realm nobody on this link subscribed to is dropped.

Drop an overlay subscription. Idempotent.

Parse a MACULA_STATIONS seed list: comma-separated <node id>@<host>:<port> entries, the node id as 64 lowercase hex characters, the host a DNS name, an IPv4 address or a bracketed IPv6 address, the port 1 to 65535. Returns the seeds in order, each pinned to its node id. A refusal names the entry's position and what is wrong with it, and carries neither the value, nor a node id, nor a host. See macula_stations:parse/1.

Resolve this pool's own D25 provider authorization for an org-namespaced Procedure under Realm — the realm-signed org_directory and the org-signed procedure_delegation naming the pool's node id, both fetched from the DHT and verified against the realm key the pool pinned for Realm at connect — as the authorization opt macula_direct_dial:publish_advertisement/5 and macula_response:advertise_direct/6,7 take: a map of the two records' encoded wire forms. advertise/5 resolves this same chain itself for the wire frame; a caller that publishes the direct-dial record needs it explicitly, since the station refuses an org-namespaced record without one (no_authorization).

As provider_authorization/3, reading the resolution's DHT calls from the provider_io/0 seam entries in Opts (defaults to the facade's own functions) — the seam advertise/5 accepts too.

Who provides Procedure in Realm: each provider whose advertisement passes the same trust check call/5 applies, with the station it serves from, in one DHT lookup bounded by TimeoutMs. A provider whose record has not replicated yet is not listed; ask again for a fresher answer. See macula_direct_dial:providers/4.

Publish to (Realm, Topic) on Pool. Equivalent to publish/5 with empty opts.

Publish to (Realm, Topic) on Pool with options, group to seal it under a sealed group's current epoch. See macula_pubsub:publish/5 for honored opts.

Store a signed record in the mesh DHT via a V2 pool.

Receive the next chunk (blocks).

Send a binary chunk on the stream.

Send a chunk with explicit encoding.

Send a pre-built, pre-signed overlay frame to the peer at the other end of Link. {error, not_connected} until its handshake completes.

Send a pre-built, pre-signed overlay frame to TargetPeer (its node_id), relayed by the station Link is connected to. {error, not_connected} until the link's handshake completes.

Server-side: emit the terminal reply value.

Share Bytes in Realm, as share_content/4 with no options.

Share Bytes in Realm from this node (D27): the pool's sharer keeps them, serves them on the node's own content procedure and announces their root in the DHT, naming the realm, the station the node is reachable through and that procedure (macula_content_sharer). Returns the root's content id: bytes of at most 256 KiB are one raw block, larger bytes a manifest over their chunks. The content is available while this node is online. Opts: name, and org, the org whose delegation serves the content (<org>/content_v1_<node id>); without org the node serves in its own namespace (~<node id>/content_v1).

Sign a domain record (tags 0x20 to 0xFF) as this node, with the pool's node identity key, in the pool's own process. See macula_client:sign_domain_record/2.

Sign a record this node signs about itself with the pool's node identity key, in the pool's own process: a node record, a procedure advertisement or a content announcement that names this node. See macula_client:sign_node_record/2 for the record checks and refusals. Optsnot_after bounds the signed record's expiry: the pool judges it on its own clock, refuses a bound already passed as {error, not_after_passed}, ends the record at the bound when the bound comes before the record's lifetime runs out, and keeps the built lifetime otherwise.

Aggregate health snapshot of a V2 pool. Suitable for /health or /status endpoints; not for hot-loop polling. See macula_client:status/1 for the full shape.

The caller's seal report of a stream it opened (DESIGN_E2E_SEAL_REPORT): sealed 1 with the id of the KEM key the stream is sealed to, or 0 for a clear stream, and provider, the stream's target. It states that the mechanism ran on this exchange, nothing more. It settles on the provider's first chunk or reply opened under the stream's key (on a clear stream, its first chunk, reply or end); after a reseal it names the reseal's key. Before it settles, on a stream that ended first, an error included, and on a local in-process stream it is {error, not_settled}; on a served stream, {error, not_a_caller}. It asks the stream process, so it answers while that process lives: until its owner ends, or the stream is closed.

Subscribe Subscriber to (Realm, Topic) on Pool. Equivalent to subscribe/5 with empty opts.

Subscribe Subscriber to (Realm, Topic) on Pool with options. The delivery option chooses how a single publisher's out-of-order arrivals are handled

Subscribe with a callback function. The SDK spawns a small receiver process internally and invokes the callback once per inbound event. See macula_pubsub:subscribe_callback/4.

Subscribe to live record-stored events filtered by type.

The binary of a text value a peer supplied (D26): the binary of {text, Bin}, a binary unchanged, and badarg for anything else.

Stop advertising a procedure on a V2 pool.

Stop advertising a LOCAL streaming procedure.

Stop advertising a streaming procedure on a V2 pool.

Unsubscribe from node up/down events.

Stop sharing the root MCID in Realm: its announcement is withdrawn and its bytes dropped. Idempotent.

Drop a pool subscription. Idempotent.

Withdraw a record this node signed — a node record, a procedure advertisement, a content announcement or a domain record — with a tombstone signed by the pool's node identity key, in the pool's own process. See macula_client:withdraw_node_record/3.

Types

m_record/0

-type m_record() :: macula_record:m_record().

mcid/0

-type mcid() :: <<_:400>>.

pool/0

-type pool() :: macula_client:pool().

procedure/0

-type procedure() :: binary().

provider_io/0

-type provider_io() ::
          #{status => fun((pool()) -> {ok, macula_client:status()} | {error, term()}),
            find_record => fun((pool(), record_key()) -> {ok, m_record()} | {error, term()}),
            sign_node_record =>
                fun((pool(), m_record(), #{not_after := integer()}) ->
                        {ok, m_record()} | {error, term()}),
            realm_key => fun((pool(), realm()) -> {ok, binary()} | none)}.

realm/0

-type realm() :: <<_:256>>.

32-byte realm tag.

record_key/0

-type record_key() :: <<_:256>>.

DHT storage key — macula_record:storage_key/1 output.

record_type/0

-type record_type() :: macula_record:type_tag().

stream/0

-type stream() :: pid().

stream_handler/0

-type stream_handler() :: fun((stream(), term()) -> any()).

stream_mode/0

-type stream_mode() :: server_stream | client_stream | bidi.

topic/0

-type topic() :: binary().

Functions

abort(Stream, Code, Message)

-spec abort(stream(), binary(), binary()) -> ok | {error, {text_too_long | invalid_text, code}}.

Abort the stream with an error frame.

advertise(Pool, Realm, Procedure, Handler, Opts)

-spec advertise(pool(), realm(), procedure(), macula_client:handler(), map()) -> ok | {error, term()}.

Register a procedure handler on a V2 pool and advertise it: the pool resolves its own D25 provider authorization — the realm-signed org_directory and the org-signed procedure_delegation that names the pool's node id, both fetched from the DHT — verifies it against the pinned realm key, and hands every link the advertisement spec: for each link the pool signs an advertisement naming the station that link is connected to as serving_station, sent as an ADVERTISE frame (signed again on reconnect and on link respawn, never past the chain's earlier expiry). The pool renews the chain at a third of its remaining life (macula_client:advertise/8, D32): a delegation lives 30 minutes, so a provider stays advertised while the realm keeps reissuing, and drops out once it stops. The frame registers the procedure at that station; it does not put a record in the DHT, which is what a direct-dialling caller resolves (macula_response:advertise_direct/6,7 publishes that one). The procedure must carry an org namespace, and the pool must run a provisioned identity whose delegation the org has published, and pin the realm's key (realm_trust at connect); a missing piece fails fast with {error, {provider_authorization, _}}.

Optsauth sets the procedure's policy: open (default, serve any identified caller), {ucan_required, IssuerNodeId} (gated to tokens from one known node), or {realm_member_required, RealmKeyId, RequiredCan} (gated to realm membership at a specific tier) -- see macula_client:auth_policy() for the full set.

stations => [StationNodeId] registers the procedure on the links to those stations only, each named by the node_id its link pins, and a link respawned later registers it again only when its station is one of them. A station the pool holds no link to is refused as {error, {station_not_linked, StationNodeId}}, an empty list as {error, {stations, empty}} and anything but a list of 32-byte node ids as {error, {stations, malformed}}; nothing is registered or kept. Without it every link registers the procedure.

Opts may also carry advertise, an arity-6 override for the pool fan-out, which then decides where the procedure registers (see macula_response:advertise_opts()), and the provider_io/0 seam entries the resolution reads its DHT calls from.

await_reply(Stream)

-spec await_reply(stream()) -> {ok, term()} | {error, term()}.

Wait for the terminal reply (client-stream / bidi).

await_reply(Stream, Timeout)

-spec await_reply(stream(), timeout()) -> {ok, term()} | {error, term()}.

call(Pool, Realm, Procedure, Payload, TimeoutMs)

-spec call(pool(), realm(), procedure(), term(), 1..600000) -> {ok, term()} | {error, term()}.

Call Procedure in Realm at the provider that serves it: resolve the procedure's verified advertisements, authorized against the realm key the pool pinned for Realm, reach the station a candidate names directly, and call its provider there. See macula_direct_dial:call/5. A procedure the pool's linked stations serve themselves, such as _dht.*, goes through macula_client:call_linked_station/5.

call(Pool, Realm, Procedure, Payload, TimeoutMs, Opts)

-spec call(pool(),
           realm(),
           procedure(),
           term(),
           1..600000,
           #{provider => <<_:256>>,
             confidential => preferred | required,
             report => boolean(),
             ucan_token => binary()}) ->
              {ok, term()} | {ok, term(), macula_station_link:report()} | {error, term()}.

As call/5, with Opts. #{provider => NodeId} calls THAT provider of Procedure and no other: the procedure is resolved and its advertisements trust-checked exactly as call/5 does, and only the named provider's are tried. A provider with no trusted advertisement by the deadline is {error, {unresolved, provider_not_advertised}}. To have every provider of a procedure answer, list them with providers/3,4 and make one call per provider. See macula_direct_dial:call/6.

The call is sealed to the KEM key of the advertisement it resolves, when that advertisement names one (E2E design §8.1). #{confidential => required} fails the call closed when none does. A lookup never downgrades a call, so confidential => off is refused as {error, {confidentiality, off_needs_explicit_target}}: a clear call is call_station/8's, to a target the application names itself.

report => true asks for the call's seal report (DESIGN_E2E_SEAL_REPORT): a result then comes back as {ok, Result, Report}, Report saying whether the request that produced it was sealed and to which key (sealed 1 with seal_key_id, or 0), and the provider it was addressed to. An error comes back as it is. Any value but a boolean is {error, {invalid_option, report}}.

ucan_token presents a UCAN to a gated provider, as call_station/8 does, on every station call the call makes; anything but bytes is {error, {invalid_option, ucan_token}} before anything is looked up.

call_station(Pool, Station, Target, Realm, Procedure, Payload, TimeoutMs)

-spec call_station(pool(), macula_client:seed(), <<_:256>>, realm(), procedure(), term(), 1..600000) ->
                      {ok, term()} | {error, term()}.

Issue a CALL to Target, a provider's node_id, at ONE specific station, dialing it directly even if it is not in the pool's seed set. Station is a seed URL (e.g. <<"quic://[::1]:4433">>). The pool reuses an existing link or dials and monitors a new one, waits for the handshake, and calls there. This is the direct-dial data path: resolve a provider's serving_station and its endpoint, then reach it in one hop. See macula_client:call_station/11. With no options it carries no signed state to decide on sealing from, so it is refused as {error, {confidentiality, no_signed_state}} (macula 13.0.0): use call_station/8 with advertisement, confidential => required or confidential => off (see call_seal/5).

call_station(Pool, Station, Target, Realm, Procedure, Payload, TimeoutMs, Opts)

-spec call_station(pool(),
                   macula_client:seed(),
                   <<_:256>>,
                   realm(),
                   procedure(),
                   term(),
                   1..600000,
                   map()) ->
                      {ok, term()} | {ok, term(), macula_station_link:report()} | {error, term()}.

As call_station/7, presenting a capability token to a gated provider via Opts (#{ucan_token => Token}). Empty/absent = none. Slice 7b dual-trust. Opts also carries the station this dial must prove, expected_node_id (see macula_client:call_station/11), and may set dial_timeout_ms, how much of TimeoutMs the wait for a fresh link's handshake may take (default: all of it). pin_tls_cert => true is REFUSED with {error, {refused, {pin_tls_cert, no_pin_primitive_for_mldsa87_identity}}}, and verify in any value with {error, {refused, {verify, one_verification_mode}}} (see trust_options_checked/1).

report => true asks for the call's seal report, as call/6 does: a result then comes back as {ok, Result, Report} (see macula_station_link:call/9). It is honoured here because the pool's own direct dial calls through this function. Any value but a boolean is {error, {invalid_option, report}}, before anything is sent.

call_stream(Procedure, Args)

-spec call_stream(procedure(), term()) -> {ok, stream()} | {error, term()}.

Open a LOCAL in-process server-stream call. Used for unit tests and same-BEAM dispatch via macula_stream_local.

call_stream(Procedure, Args, Opts)

-spec call_stream(procedure(), term(), map()) -> {ok, stream()} | {error, term()}.

Open a LOCAL in-process server-stream call with options.

call_stream(Pool, Realm, Procedure, Args, Opts)

-spec call_stream(pool(), realm(), procedure(), term(), map()) -> {ok, stream()} | {error, term()}.

Open a streaming RPC to Procedure's provider in Realm: resolve the provider through its procedure_advertisement and open the stream at its serving station, naming the provider as the target, as call/5 does for a single-reply call. Same as macula_direct_dial:call_stream/5. The returned stream is bound to that station's link (errors with peer_down if the link dies; caller re-opens). Optsucan_token presents a UCAN to a streaming procedure advertised with an auth policy (see advertise_stream/6). An open whose signed STREAM_OPEN would be longer than max_stream_open_bytes (1 MiB by default) returns {error, {open_too_large, Limit}} without sending anything. The stream is sealed as call/6 seals a call: confidential => required fails closed without a key, and confidential => off is refused as {error, {confidentiality, off_needs_explicit_target}} (see call_stream_station/7).

call_stream_station(Pool, Station, Target, Realm, Procedure, Args, Opts)

-spec call_stream_station(pool(), macula_client:seed(), <<_:256>>, realm(), procedure(), term(), map()) ->
                             {ok, stream()} | {error, term()}.

Open a streaming RPC to Target, a provider's node_id, by DIALING a specific station directly (direct-dial): the streaming analogue of call_station/7. Compose it with DHT resolution (find_records -> read_procedure_advertisement -> station_endpoint) to reach a stream provider in one hop, exactly as a unary caller does. Opts may set dial_timeout_ms (default 10_000) and a mode. Opts also names the station this dial must prove, expected_node_id (see macula_client:call_station/11); pin_tls_cert => true and verify are REFUSED as call_station/8 describes.

Whether the STREAM_OPEN and every later frame is sealed is decided from signed state only, as for call_station/8 (call_seal/5, E2E design §8.1): advertisement, confidential => off or confidential => required. With none of them the open is refused {error, {confidentiality, no_signed_state}} before anything is sent.

child_spec(Id, Seeds, Opts)

OTP child spec to drop a V2 pool into a caller's supervision tree. Give node_identity as a loader {Module, Function, Args} that returns {ok, Key}. Its Args say where the key is and never hold the key, because a supervisor logs them when a start fails. See macula_client:child_spec/3.

close(Pool)

-spec close(pool()) -> ok.

Stop a V2 pool. Every subscriber receives a final {macula_event_gone, SubRef, pool_closed} message.

close_send(Stream)

-spec close_send(stream()) -> ok.

Half-close the write side; recv still drains.

close_stream(Stream)

-spec close_stream(stream()) -> ok.

Close a V1 stream (both sides). Renamed from close/1 in 3.11.0 because close/1 now refers to the V2 pool surface.

connect(Seeds, Opts)

-spec connect([macula_client:seed()], macula_client:opts()) -> {ok, pool()} | {error, term()}.

Connect to the Macula relay mesh and return a pool handle.

Seeds is a list of station endpoints: #{host, port, expected_node_id} maps, or URL binaries or strings. Every seed names the node_id it expects, in the seed or in the expected_node_id option; otherwise no pool starts and {error, {seeds, expected_node_id_required}} is returned. The pool spawns one peering link per seed and routes ops with replication, replay, and event dedup. Returns immediately; link handshakes complete asynchronously.

Honored opts (full reference: macula_client:opts()):

  • node_identity: the pool's node identity key, in the node's crypto profile. If absent, the pool uses the node's one stored identity (macula_node_keys:node_identity/1, at the node_identity_path application env): loaded if stored, ground and stored once if not, and refused, never replaced, if the file exists and will not load. Every pool on the node, and every restart, is then the same node.
  • realm_trust: the realm keys the pool pins, one per realm id, as #{RealmId => RealmKey}, each realm's public key as carried. A call trusts an org namespaced advertisement only through the key pinned for its realm. Refused as {error, {realm_trust, invalid}} unless every id is 32 bytes and every key is well formed for the node's crypto profile, and as {error, {realm_trust, profile_mismatch}} for a key of the other profile.
  • replication_factor — links per PUBLISH (default 2, since 10.19.0).
  • capabilities — per-link bitfield (default 0).
  • alpn — QUIC ALPN list (default [<<"macula">>]).
  • connect_timeout_ms — per-link CONNECT/HELLO deadline (default 30_000).
  • dedup_sweep_ms: how often the inbound publication dedup table is swept.
  • admission_sweep_ms: how often the request admission is swept for entries past their deadline plus 5 minutes (default 30_000).
  • renew_backoff_ms / renew_recheck_ms: the retry delay of a failed renewal of an advertised chain (default 5_000, doubling, never past its not_after) and how often a chain past it is asked for again (default 300_000). See advertise/5 and D32.
  • verify — ⚠ REFUSED in any value, here and on every seed in Seeds, with {error, {refused, {verify, one_verification_mode}}}. A link trusts a station in one way only: its handshake signature under the key of its ML-DSA-87 certificate, and its identity through the handshake. There is no chain to check.
  • expected_node_id — the station node_id the handshake must prove, for every link this pool dials.
  • pin_tls_cert — ⚠ true is REFUSED, here and on every seed in Seeds, with {error, {refused, {pin_tls_cert, no_pin_primitive_for_mldsa87_identity}}}. No pin primitive can express an ML-DSA-87 identity, so the option could never be honoured; it pinned nothing at any value. false and an absent key pass and change nothing. See macula#15.

Legacy opts silently dropped (with a one-shot logger:notice): relays (use the Seeds positional argument), realm (V2 is realm-per-call), site (no V2 analog), connections (one link per seed; add more seeds to grow the pool).

See macula_client for the canonical pool implementation and macula_pubsub for the slice module.

dist_relay_client()

-spec dist_relay_client() -> {ok, pid()} | {error, not_joined}.

The dist relay client that join_dist_relay/1 started, if it is running.

The client exits with {relay_closed, Reason} when the relay closes the connection and is not restarted. Monitor the returned pid and call join_dist_relay/1 again after it goes down.

ensure_distributed()

-spec ensure_distributed() -> ok | {error, term()}.

Ensure this node is running in distributed mode.

ensure_station_link(Pool, Station, Opts, TimeoutMs)

-spec ensure_station_link(pool(), macula_client:seed(), map(), pos_integer()) ->
                             {ok, pid()} | {error, term()}.

A pool link to Station, for a consumer that must address one station through one link, as an overlay does: a live link the pool already holds to that station, else one it dials. The station is pinned by the expected_node_id in its seed map or in Opts, which a link requires (D5, D16); pin_tls_cert => true and verify are refused as call_station/8 describes. Answers {ok, Link} once the link's handshake has completed, within TimeoutMs, or {error, not_connected}. The pool owns and monitors the link: it respawns it and ends it with the pool, so a consumer holds the pid and never stops it. Drive it with overlay_subscribe/3 and send_overlay_frame/2,3.

field(Name, Map)

-spec field(atom() | binary(), map()) -> term().

A field of a map a peer supplied, or undefined when it is absent. See field/3.

field(Name, Map, Default)

-spec field(atom() | binary(), map(), term()) -> term().

A field of a map a peer supplied (D26), or Default when it is absent. A map from the codec carries its text keys as {text, Bin}; a map handed over in process may carry atom or binary keys. The lookup tries {text, Name}, then the atom, then the binary, so a handler reads both kinds of map the same way. Looking up a binary name never creates an atom.

find_record(Pool, Key)

-spec find_record(pool(), record_key()) -> {ok, m_record()} | {error, not_found | term()}.

Fetch a record from the mesh DHT by its macula_record:storage_key/1.

Returns {error, not_found} when no record exists at the key. The record is verified under the node's crypto profile before it is returned; one that does not verify returns its refusal from macula_record:verify/2, such as {error, expired}.

find_record(Pool, Key, TimeoutMs)

-spec find_record(pool(), record_key(), pos_integer()) -> {ok, m_record()} | {error, not_found | term()}.

As find_record/2, waiting at most TimeoutMs for the reply, for a caller that bounds its work by a deadline of its own.

find_records(Pool, Key)

-spec find_records(pool(), record_key()) -> {ok, [m_record()]} | {error, term()}.

Fetch EVERY record stored at Key — the full multi-value set, e.g. every procedure_advertisement under one procedure's storage key. Where find_record/2 returns the first record (or not_found), this returns the whole list, empty when none.

The relay's local store is a signer-deduped multiset: one record per signing key at a storage key, so N providers of one procedure return N records. Each record is verified under the node's crypto profile, and a record that does not verify is dropped.

find_records(Pool, Key, TimeoutMs)

-spec find_records(pool(), record_key(), pos_integer()) -> {ok, [m_record()]} | {error, term()}.

As find_records/2, waiting at most TimeoutMs for the reply, for a caller that bounds its work by a deadline of its own.

find_records_by_type(Pool, Type)

-spec find_records_by_type(pool(), record_type()) -> {ok, [m_record()]} | {error, term()}.

Return every record of a given type currently visible from the pool's connected stations, each verified under the node's crypto profile; a record that does not verify is dropped.

Coverage depends on each station's view of the DHT — a single station sees its local replicas plus whatever its peers have gossiped. Aggregating across the full mesh requires querying multiple stations and deduplicating by record key.

get_content(Pool, Realm, MCID)

-spec get_content(pool(), realm(), mcid()) -> {ok, binary()} | {error, term()}.

Fetch MCID in Realm, as get_content/4 with no options.

get_content(Pool, Realm, MCID, Opts)

-spec get_content(pool(), realm(), mcid(), map()) -> {ok, binary()} | {error, term()}.

Fetch MCID from a node that shares it in Realm (D27): find its announcements, open a stream to a sharer through the station it announced, and verify every block and the manifest against their content ids (macula_content_fetch). A sharer that fails moves the fetch to the next; {error, {unavailable, [{Sharer, Reason}]}} when all fail, {error, not_shared} when none announces it, {error, invalid_mcid} for an id that is not tag 2. Opts: max_bytes (256 MiB), max_chunks, chunk_timeout_ms, parallel. No realm key is needed: content verifies itself.

join_dist_relay(Opts)

-spec join_dist_relay(map()) -> ok | {error, term()}.

Enable Erlang distribution over a dedicated dist relay (macula-io/macula-dist-relay).

Different from join_mesh/1: - Connects to a dist relay (port 4434, ALPN macula-dist), NOT the pub/sub station mesh - No mesh_client, no pub/sub subscriptions — only dist traffic - Uses raw QUIC stream routing with no MessagePack overhead

Options: - url (required): &lt;&lt;"quic://relay.example.com:4434"&gt;&gt;

After this returns ok, standard OTP distribution (rpc:call/4, gen_server:call/3 across nodes, pg groups, etc.) works across firewalls via the dist relay.

The relay client runs as a temporary child of the macula application supervisor. It does not reconnect: when the relay closes the connection the client ends and is not restarted. Monitor the pid from dist_relay_client/0 to learn that, then call this function again. Returns {error, macula_not_started} when the macula application is not running.

join_mesh(Opts)

-spec join_mesh(map()) -> ok | {error, term()}.

Join the Macula relay mesh with Erlang distribution.

After calling this, standard OTP distribution works across firewalls. Opts takes:

  • relays (required): the V2 pool's seeds, each a map with host, port and expected_node_id, the relay's 32-byte node_id, which every dial checks (D16). A relay without one refuses the join with {error, {relays, expected_node_id_required}} before any pool starts.
  • node_identity: the V2 pool's node identity key, macula_node_keys:node_key(). Default: the node's one stored identity, as for connect/2.

Internally builds a V2 macula_client:pool() and registers it with macula_dist_pool as the carrier for _dist.tunnel.* traffic. Dist tunnel frames travel under the all-zeros realm (protocol-internal infrastructure, not bound to any user realm).

links(Pool)

-spec links(pool()) -> {ok, [macula_client:link_info()]}.

Per-link snapshot of a V2 pool — one entry per spawned link with its peer station node_id (pubkey), dial host, pid, and connected flag. Use this to resolve a specific station (by pubkey or hostname) to its link for targeted, per-station operations. See macula_client:links/1 for the link_info() shape.

monitor_nodes()

-spec monitor_nodes() -> ok.

Subscribe to node up/down events.

open_stream(Procedure, Args, Opts)

-spec open_stream(procedure(), term(), map()) -> {ok, stream()} | {error, term()}.

Open a LOCAL in-process client-stream or bidi call. Used for unit tests and same-BEAM dispatch via macula_stream_local.

open_stream(Procedure, Args, Opts, Mode)

-spec open_stream(procedure(), term(), map(), stream_mode()) -> {ok, stream()} | {error, term()}.

Open a LOCAL in-process stream with explicit mode.

overlay_subscribe(Link, Realm, Subscriber)

-spec overlay_subscribe(pid(), realm(), pid()) -> {ok, reference()} | {error, term()}.

Take the overlay frames for Realm arriving on Link (from ensure_station_link/4), delivered to Subscriber as {macula_overlay_frame, SubRef, Frame, Meta}, where Meta names the frame's authenticated sender and, for a relayed frame, the station it came via. A frame from the connected peer, or a GOSSIP, for a realm nobody on this link subscribed to is dropped.

overlay_unsubscribe(Link, SubRef)

-spec overlay_unsubscribe(pid(), reference()) -> ok.

Drop an overlay subscription. Idempotent.

parse_stations(Value)

-spec parse_stations(binary()) -> {ok, [macula_stations:seed()]} | {error, term()}.

Parse a MACULA_STATIONS seed list: comma-separated <node id>@<host>:<port> entries, the node id as 64 lowercase hex characters, the host a DNS name, an IPv4 address or a bracketed IPv6 address, the port 1 to 65535. Returns the seeds in order, each pinned to its node id. A refusal names the entry's position and what is wrong with it, and carries neither the value, nor a node id, nor a host. See macula_stations:parse/1.

provider_authorization(Pool, Realm, Procedure)

-spec provider_authorization(pool(), realm(), procedure()) ->
                                {ok,
                                 #{org_directory := binary(), procedure_delegation := binary()} |
                                 undefined} |
                                {error, term()}.

Resolve this pool's own D25 provider authorization for an org-namespaced Procedure under Realm — the realm-signed org_directory and the org-signed procedure_delegation naming the pool's node id, both fetched from the DHT and verified against the realm key the pool pinned for Realm at connect — as the authorization opt macula_direct_dial:publish_advertisement/5 and macula_response:advertise_direct/6,7 take: a map of the two records' encoded wire forms. advertise/5 resolves this same chain itself for the wire frame; a caller that publishes the direct-dial record needs it explicitly, since the station refuses an org-namespaced record without one (no_authorization).

Fails fast with {error, {provider_authorization, _}} when the procedure has no org namespace, a chain piece is missing from the DHT, the pool pinned no key for Realm, or the chain does not verify against that key — the same failures advertise/5 reports.

A procedure in this node's own namespace, ~<own node_id>/<name>, needs no chain (D25 item 6, revised 2026-09-24): the answer is {ok, undefined}, which publish_advertisement/5 reads as no authorization. Another node's namespace is {error, {provider_authorization, not_own_namespace}}.

provider_authorization(Pool, Realm, Procedure, Opts)

-spec provider_authorization(pool(), realm(), procedure(), provider_io()) ->
                                {ok,
                                 #{org_directory := binary(), procedure_delegation := binary()} |
                                 undefined} |
                                {error, term()}.

As provider_authorization/3, reading the resolution's DHT calls from the provider_io/0 seam entries in Opts (defaults to the facade's own functions) — the seam advertise/5 accepts too.

providers(Pool, Realm, Procedure)

-spec providers(pool(), realm(), procedure()) ->
                   {ok, [#{provider := <<_:256>>, station := <<_:256>>}]} | {error, term()}.

As providers/4, within 5 seconds.

providers(Pool, Realm, Procedure, TimeoutMs)

-spec providers(pool(), realm(), procedure(), 1..600000) ->
                   {ok, [#{provider := <<_:256>>, station := <<_:256>>}]} | {error, term()}.

Who provides Procedure in Realm: each provider whose advertisement passes the same trust check call/5 applies, with the station it serves from, in one DHT lookup bounded by TimeoutMs. A provider whose record has not replicated yet is not listed; ask again for a fresher answer. See macula_direct_dial:providers/4.

publish(Pool, Realm, Topic, Payload)

-spec publish(pool(), realm(), topic(), term()) -> ok | {error, term()}.

Publish to (Realm, Topic) on Pool. Equivalent to publish/5 with empty opts.

publish(Pool, Realm, Topic, Payload, Opts)

-spec publish(pool(), realm(), topic(), term(), map()) -> ok | {error, term()}.

Publish to (Realm, Topic) on Pool with options, group to seal it under a sealed group's current epoch. See macula_pubsub:publish/5 for honored opts.

put_record(Pool, Signed)

-spec put_record(pool(), m_record() | binary()) -> ok | {error, term()}.

Store a signed record in the mesh DHT via a V2 pool.

Build the record via the typed constructors in macula_record (node_record/3,4, content_announcement/3,4, tombstone/2,3, realm_directory/3,4, procedure_advertisement/4,5, etc.) and sign it. A record this node signs about itself is signed by its pool, which holds the node identity key: through macula_client:sign_node_record/2, a domain record (tags 0x20 to 0xFF) through macula_client:sign_domain_record/2, and the tombstone of either through macula_client:withdraw_node_record/3. A realm-, org- or foundation-signed record is signed with macula_record:sign/2 and its signer's key. Pass the signed record or its wire form (macula_record:encode/1). The record travels as its wire form; the station verifies it on receipt and stores it under macula_record:storage_key/1, propagating to the K-nearest peers.

recv(Stream)

-spec recv(stream()) -> {chunk, binary()} | {data, term()} | eof | {error, term()}.

Receive the next chunk (blocks).

recv(Stream, Timeout)

-spec recv(stream(), timeout()) -> {chunk, binary()} | {data, term()} | eof | {error, term()}.

send(Stream, Bin)

-spec send(stream(), binary()) -> ok | {error, term()}.

Send a binary chunk on the stream.

send(Stream, Body, Encoding)

-spec send(stream(), binary() | term(), raw | msgpack) -> ok | {error, term()}.

Send a chunk with explicit encoding.

send_overlay_frame(Link, Frame)

-spec send_overlay_frame(pid(), macula_frame:frame()) -> ok | {error, term()}.

Send a pre-built, pre-signed overlay frame to the peer at the other end of Link. {error, not_connected} until its handshake completes.

send_overlay_frame(Link, TargetPeer, Frame)

-spec send_overlay_frame(pid(), <<_:256>>, macula_frame:frame()) -> ok | {error, term()}.

Send a pre-built, pre-signed overlay frame to TargetPeer (its node_id), relayed by the station Link is connected to. {error, not_connected} until the link's handshake completes.

set_reply(Stream, Result)

-spec set_reply(stream(), term()) -> ok.

Server-side: emit the terminal reply value.

share_content(Pool, Realm, Bytes)

-spec share_content(pool(), realm(), binary()) -> {ok, mcid()} | {error, term()}.

Share Bytes in Realm, as share_content/4 with no options.

share_content(Pool, Realm, Bytes, Opts)

-spec share_content(pool(), realm(), binary(), map()) -> {ok, mcid()} | {error, term()}.

Share Bytes in Realm from this node (D27): the pool's sharer keeps them, serves them on the node's own content procedure and announces their root in the DHT, naming the realm, the station the node is reachable through and that procedure (macula_content_sharer). Returns the root's content id: bytes of at most 256 KiB are one raw block, larger bytes a manifest over their chunks. The content is available while this node is online. Opts: name, and org, the org whose delegation serves the content (<org>/content_v1_<node id>); without org the node serves in its own namespace (~<node id>/content_v1).

sign_domain_record(Pool, Record)

-spec sign_domain_record(pool(), m_record()) -> {ok, m_record()} | {error, term()}.

Sign a domain record (tags 0x20 to 0xFF) as this node, with the pool's node identity key, in the pool's own process. See macula_client:sign_domain_record/2.

sign_node_record(Pool, Record)

-spec sign_node_record(pool(), m_record()) -> {ok, m_record()} | {error, term()}.

Sign a record this node signs about itself with the pool's node identity key, in the pool's own process: a node record, a procedure advertisement or a content announcement that names this node. See macula_client:sign_node_record/2 for the record checks and refusals. Optsnot_after bounds the signed record's expiry: the pool judges it on its own clock, refuses a bound already passed as {error, not_after_passed}, ends the record at the bound when the bound comes before the record's lifetime runs out, and keeps the built lifetime otherwise.

sign_node_record(Pool, Record, Opts)

-spec sign_node_record(pool(), m_record(), #{not_after := integer()}) ->
                          {ok, m_record()} | {error, term()}.

status(Pool)

-spec status(pool()) -> {ok, macula_client:status()}.

Aggregate health snapshot of a V2 pool. Suitable for /health or /status endpoints; not for hot-loop polling. See macula_client:status/1 for the full shape.

stream_report(Stream)

-spec stream_report(stream()) ->
                       {ok, macula_station_link:report()} | {error, not_settled | not_a_caller}.

The caller's seal report of a stream it opened (DESIGN_E2E_SEAL_REPORT): sealed 1 with the id of the KEM key the stream is sealed to, or 0 for a clear stream, and provider, the stream's target. It states that the mechanism ran on this exchange, nothing more. It settles on the provider's first chunk or reply opened under the stream's key (on a clear stream, its first chunk, reply or end); after a reseal it names the reseal's key. Before it settles, on a stream that ended first, an error included, and on a local in-process stream it is {error, not_settled}; on a served stream, {error, not_a_caller}. It asks the stream process, so it answers while that process lives: until its owner ends, or the stream is closed.

subscribe(Pool, Realm, Topic, Subscriber)

-spec subscribe(pool(), realm(), topic(), pid()) ->
                   {ok, reference()} | {error, {text_too_long | invalid_text, topic}}.

Subscribe Subscriber to (Realm, Topic) on Pool. Equivalent to subscribe/5 with empty opts.

subscribe(Pool, Realm, Topic, Subscriber, Opts)

-spec subscribe(pool(), realm(), topic(), pid(), map()) ->
                   {ok, reference()} |
                   {error,
                    {text_too_long | invalid_text, topic} |
                    {invalid_option, group | distributor | ucan_token} |
                    {group, macula_group_keyring:reason()}}.

Subscribe Subscriber to (Realm, Topic) on Pool with options. The delivery option chooses how a single publisher's out-of-order arrivals are handled:

  • ordered (default) — per-publisher FIFO by seq; out-of-order arrivals are buffered and released in order, a genuinely missing seq skipped after order_timeout_ms (a connect/2 option, default 250ms). A new publisher's first facts are held for up to order_timeout_ms, so its order starts at the lowest seq seen.
  • latest_only — deliver only seqs newer than the highest seen for that publisher (drop stale); no buffering, no delay.
  • as_arrives — deliver in raw arrival order; the consumer orders it itself.

Ordering state is kept per publisher, and apart for EVENTs whose publisher signature did not verify. group subscribes under a sealed group, joining it first; its events arrive opened, or once as macula_event_unopened. See macula_pubsub:subscribe/5.

subscribe_callback(Pool, Realm, Topic, Callback)

-spec subscribe_callback(pool(), realm(), topic(), macula_pubsub:callback()) ->
                            {ok, reference()} | {error, term()}.

Subscribe with a callback function. The SDK spawns a small receiver process internally and invokes the callback once per inbound event. See macula_pubsub:subscribe_callback/4.

subscribe_records(Pool, Type, Callback)

-spec subscribe_records(pool(), record_type(), fun((m_record()) -> any())) ->
                           {ok, reference()} | {error, term()}.

Subscribe to live record-stored events filtered by type.

The callback receives each newly-stored record of the given type that verifies under the node's crypto profile. Returns a subscription reference for unsubscribe_records/2. Topic shape is _dht.records.<type>.stored, rendered with the type tag as a decimal integer for log friendliness.

text(Bin)

-spec text({text, binary()} | binary()) -> binary().

The binary of a text value a peer supplied (D26): the binary of {text, Bin}, a binary unchanged, and badarg for anything else.

unadvertise(Pool, Realm, Procedure)

-spec unadvertise(pool(), realm(), procedure()) -> ok.

Stop advertising a procedure on a V2 pool.

unadvertise_stream(Procedure)

-spec unadvertise_stream(procedure()) -> ok.

Stop advertising a LOCAL streaming procedure.

unadvertise_stream(Pool, Realm, Procedure)

-spec unadvertise_stream(pool(), realm(), procedure()) -> ok.

Stop advertising a streaming procedure on a V2 pool.

unmonitor_nodes()

-spec unmonitor_nodes() -> ok.

Unsubscribe from node up/down events.

unshare_content(Pool, Realm, MCID)

-spec unshare_content(pool(), realm(), mcid()) -> ok.

Stop sharing the root MCID in Realm: its announcement is withdrawn and its bytes dropped. Idempotent.

unsubscribe(Pool, SubRef)

-spec unsubscribe(pool(), reference()) -> ok.

Drop a pool subscription. Idempotent.

unsubscribe_records(Pool, Ref)

-spec unsubscribe_records(pool(), reference()) -> ok.

Cancel a subscribe_records/3 subscription.

withdraw_node_record(Pool, Withdrawn, Reason)

-spec withdraw_node_record(pool(), m_record() | binary(), macula_record:reason()) ->
                              {ok, m_record()} | {error, term()}.

Withdraw a record this node signed — a node record, a procedure advertisement, a content announcement or a domain record — with a tombstone signed by the pool's node identity key, in the pool's own process. See macula_client:withdraw_node_record/3.