Supervised buffered connection event emitter modeled after Baileys'
makeEventBuffer.
Each emitter owns its buffering server and callback task supervisor. Internal taps remain ordered, while each public subscriber has an independent serial, bounded delivery lane. Slow or reentrant subscribers therefore do not block event ingestion or internal protocol handling.
Summary
Functions
Enter buffering mode.
Return whether buffering is active.
Wrap work in a nested buffering context.
Emit one event, buffering it when the current buffer policy requires it.
Flush the active buffer, retaining it when the dispatch queue is full.
Register an event-map subscriber and return its unsubscribe function.
Add values used by conditional buffered-event evaluation.
Start a supervised event emitter runtime.
Register a pre-dispatch tap and return its unsubscribe function.
Types
@type emit_error() :: :dispatch_queue_full
@type event() :: atom()
Functions
@spec buffer(GenServer.server()) :: :ok
Enter buffering mode.
@spec buffering?(GenServer.server()) :: boolean()
Return whether buffering is active.
@spec create_buffered_function(GenServer.server(), (-> term())) :: (-> term())
Wrap work in a nested buffering context.
@spec emit(GenServer.server(), event(), term()) :: :ok | {:error, emit_error()}
Emit one event, buffering it when the current buffer policy requires it.
@spec flush(GenServer.server()) :: boolean() | {:error, emit_error()}
Flush the active buffer, retaining it when the dispatch queue is full.
@spec process(GenServer.server(), (map() -> term())) :: (-> :ok)
Register an event-map subscriber and return its unsubscribe function.
@spec seed(GenServer.server(), map()) :: :ok
Add values used by conditional buffered-event evaluation.
@spec start_link(keyword()) :: Supervisor.on_start()
Start a supervised event emitter runtime.
@spec tap(GenServer.server(), (map() -> term())) :: (-> :ok)
Register a pre-dispatch tap and return its unsubscribe function.