KubeMQ
Client SDKsElixirTutorials

Persistent Pub/Sub

Publish and subscribe to persistent KubeMQ Events Store messages using the Elixir SDK pub/sub API.

Overview

This tutorial builds a publisher and subscriber on a KubeMQ Events Store channel — reach for this pattern when a subscriber can't guarantee it's listening the instant a message is published. Plain events are fire-and-forget: publish with no one subscribed and the message is gone. Events Store persists every event to a durable, ordered log, so a subscriber connecting seconds or a full restart later still catches up — useful for anything needing a complete history, like an audit trail or event-sourced state.

The two calls involved: send_event_store/2 publishes and returns {:ok, %EventStoreResult{}} confirming storage plus a broker-assigned sequence number, and subscribe_to_events_store/3 requires a start_at option telling the broker where to start — new events only (:start_new_only, used here), from the first stored event, or a given sequence or time. Production subscribers usually resume from a saved checkpoint instead of starting fresh.

Gotchas: starting from new events means anything published earlier is silently skipped — this sample papers over that race with a fixed Process.sleep instead of a ready signal, fine for a demo but not production. Replaying from the first event on every restart replays the whole log, which gets costly on a busy channel. Persistence isn't consumer coordination: each independent subscriber gets its own full replay unless grouped with a consumer group.

Prerequisites

  • KubeMQ server running on localhost:50000
  • Elixir SDK installed ({:kubemq, "~> 1.0"} in mix.exs)

Code

main.exs
channel = "elixir-events-store.persistent-pubsub"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "elixir-es-basic")

event = KubeMQ.EventStore.new(channel: channel, body: "Persistent event", metadata: "audit")

case KubeMQ.Client.send_event_store(client, event) do
  {:ok, result} ->
    IO.puts("Event stored! Sent: #{result.sent}")

  {:error, err} ->
    IO.puts("Store failed: #{err.message}")
end

parent = self()

{:ok, sub} =
  KubeMQ.Client.subscribe_to_events_store(client, channel,
    start_at: :start_new_only,
    on_event: fn event ->
      IO.puts("Received stored event: #{event.body}")
      send(parent, :got_it)
    end
  )

Process.sleep(500)

event2 = KubeMQ.EventStore.new(channel: channel, body: "Second persistent event")
{:ok, _} = KubeMQ.Client.send_event_store(client, event2)

receive do
  :got_it -> IO.puts("Event delivery confirmed!")
after
  5_000 -> IO.puts("Timeout")
end

KubeMQ.Subscription.cancel(sub)
KubeMQ.Client.close(client)

How It Works

  • send_event_store/2 persists the event and returns {:ok, %EventStoreResult{}} with confirmation
  • subscribe_to_events_store/3 requires a start_at position to control which events are delivered
  • With :start_new_only, only events published after subscribing are received

Was this page helpful?

On this page