macula_pubsub (macula v13.4.0)

View Source

Pubsub surface for the V2 SDK.

Thin delegation over macula_client (the pool). The pool owns the link state machine, replication, replay, and dedup; this module is the named public entry point that consumers reach for (or, more often, the macula facade re-exports of the same functions).

Realm-per-call

Per PLAN_V2_PARITY Q2 §2: every call carries its own 32-byte realm tag. There is no connect-time default realm. A single pool can multiplex any number of realms with no extra plumbing.

Quick start

  {ok, Pool} = macula:connect(Seeds, ConnectOpts),
  ok          = macula_pubsub:publish(Pool, Realm, Topic, Payload),
  {ok, Sub}   = macula_pubsub:subscribe(Pool, Realm, Topic, self()),
  receive
      {macula_event, Sub, Topic, Payload, Meta} -> ok
  end,
  ok          = macula_pubsub:unsubscribe(Pool, Sub).

See docs/guides/pubsub/PUBSUB_GUIDE.md for a full guide.

Summary

Functions

The meta a subscriber receives with an event, built from a publication that verified and the EVENT's delivered_via. A link calls this for every event it delivers, so the meta has one producer.

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

Publish to (Realm, Topic) on Pool.

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

Subscribe Subscriber to (Realm, Topic) via Pool.

Subscribe with a callback function instead of a receiver pid. Spawns a small receiver process internally that drives the callback for every inbound event. The receiver monitors the caller; if the caller dies, the receiver follows and the subscription is cleaned up by the pool's standard subscriber-DOWN path. It monitors the pool too, and ends when the pool dies: a pool that is killed sends no macula_event_gone.

Drop a subscription. Idempotent — unknown SubRef is a no-op.

Types

callback/0

-type callback() :: fun((Topic :: binary(), Payload :: term(), Meta :: event_meta()) -> any()).

event_meta/0

-type event_meta() ::
          #{realm := <<_:256>>,
            publisher := <<_:256>>,
            seq := non_neg_integer(),
            published_at := non_neg_integer(),
            delivered_via := macula_frame:delivery_channel(),
            publication_hash := <<_:384>>,
            expires_at := non_neg_integer(),
            sealed => 0 | 1,
            seal_key_id => <<_:64>>}.

Functions

event_meta(_, DeliveredVia)

The meta a subscriber receives with an event, built from a publication that verified and the EVENT's delivered_via. A link calls this for every event it delivers, so the meta has one producer.

publish(Pool, Realm, Topic, Payload)

-spec publish(macula_client:pool(), <<_:256>>, binary(), 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(macula_client:pool(), <<_:256>>, binary(), term(), map()) -> ok | {error, term()}.

Publish to (Realm, Topic) on Pool.

Opts currently honored:

  • timeout_ms — gen_server call timeout (default 5_000). Most apps leave this as default.
  • group — a sealed group's prefix, which Topic must be under (plans/DESIGN_E2E_SEALED_PUBSUB.md): the payload is sealed under the group's current epoch, its key pulled from the org's distributor first. A refusal fails the publish closed as {error, {group, Reason}}; a topic outside the prefix, or a prefix with no org segment, is {error, {invalid_option, group}}. ucan_token carries the org's grant to the distributor, and distributor pins its node_id.

Without group, a topic under a group this node holds is refused as {error, {confidentiality, {group_held, Prefix}}} rather than sent in the clear.

Returns ok as soon as one configured station accepts the PUBLISH frame (partial success = success, per PLAN_V2_PARITY §5.1.1). Returns {error, {transient, no_healthy_station}} when the pool has no spawned links; the caller may retry.

subscribe(Pool, Realm, Topic, Subscriber)

-spec subscribe(macula_client:pool(), <<_:256>>, binary(), pid()) ->
                   {ok, reference()} | {error, {text_too_long | invalid_text, topic}}.

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

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

-spec subscribe(macula_client:pool(), <<_:256>>, binary(), 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) via Pool.

Returns {ok, SubRef}. Subscriber subsequently receives {macula_event, SubRef, Topic, Payload, Meta} for each delivered event, where Meta is an event_meta(). Only a publication that verified is delivered. Subscriber also receives {macula_event_gone, SubRef, Reason} once when the subscription terminates (pool close, subscriber pid death).

Opts honors delivery (see macula:subscribe/5) and group, a sealed group's prefix Topic (or every topic a pattern matches) must be under, with ucan_token and distributor as for publish/5. The group is joined before the subscription is made, and a refusal fails it closed as {error, {group, Reason}}. Under a group, a sealed event arrives opened, its meta saying sealed => 1 and seal_key_id; one that cannot be opened arrives once as {macula_event_unopened, SubRef, Topic, #{publisher, seal_key_id, reason}}, and nothing of its payload. reason is one of unknown_epoch, epoch_expired, not_a_member, membership_unknown, no_distributor, no_group (a sealed event on a subscription that named no group) and tag_invalid. A clear event carries sealed => 0, and one under a prefix this node holds as required is refused, counted and logged naming its publisher.

subscribe_callback(Pool, Realm, Topic, Callback)

-spec subscribe_callback(macula_client:pool(), <<_:256>>, binary(), callback()) ->
                            {ok, reference()} | {error, term()}.

Subscribe with a callback function instead of a receiver pid. Spawns a small receiver process internally that drives the callback for every inbound event. The receiver monitors the caller; if the caller dies, the receiver follows and the subscription is cleaned up by the pool's standard subscriber-DOWN path. It monitors the pool too, and ends when the pool dies: a pool that is killed sends no macula_event_gone.

A crashing callback does NOT kill the receiver — the exception is logged and the next event is delivered. This is intentional: a transient bug in event handler N should not lose events N+1..M.

Caller cleanup: invoke unsubscribe(Pool, SubRef) with the returned ref. The receiver shuts down on the resulting macula_event_gone message.

unsubscribe(Pool, SubRef)

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

Drop a subscription. Idempotent — unknown SubRef is a no-op.