KubeMQ
Client SDKsRubyReference

Events Store

KubeMQ Ruby SDK API reference for persistent pub/sub with replay capabilities.

EventStoreMessage

Outbound durable event message. Unlike EventMessage, events store messages are persisted and can be replayed.

msg = KubeMQ::PubSub::EventStoreMessage.new(
  channel: "orders.created",
  metadata: "order-event",
  body: '{"order_id": 100}',
  tags: { "region" => "us-east" }
)

Attributes

AttributeTypeDefaultDescription
idStringAuto-generated UUIDUnique message identifier
channelStringRequiredTarget channel name
metadataStringnilArbitrary metadata string
bodyStringnilMessage payload (binary-safe)
tagsHash{String => String}{}User-defined key-value tags

EventStoreStartPosition

Constants controlling where the broker begins replaying persisted events.

ConstantValueDescription
START_NEW_ONLY1Only events published after the subscription starts
START_FROM_FIRST2Replay from the first event ever stored
START_FROM_LAST3Replay from the most recent event
START_AT_SEQUENCE4Replay from a specific sequence number
START_AT_TIME5Replay from a specific Unix timestamp
START_AT_TIME_DELTA6Replay events stored within the last N seconds

EventsStoreSubscription

Configuration for subscribing to durable events store channels.

sub = KubeMQ::PubSub::EventsStoreSubscription.new(
  channel: "orders.created",
  start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST,
  start_position_value: 0,
  group: nil
)

Attributes

AttributeTypeDefaultDescription
channelStringRequiredChannel name
start_positionIntegerRequiredOne of the EventStoreStartPosition constants
start_position_valueInteger0Sequence number, timestamp, or delta
groupStringnilConsumer group for load-balanced delivery

PubSubClient Methods

send_event_store(message)

Sends a single event to a durable events store channel.

result = client.send_event_store(msg)
puts "Stored: #{result.id}" if result.sent

Returns: EventStoreResult

Raises: ValidationError, ClientClosedError, ConnectionError

subscribe_to_events_store(subscription, cancellation_token:, on_error:, &block)

Subscribes to durable events store messages. On reconnect, playback resumes from the last received sequence number.

sub = KubeMQ::PubSub::EventsStoreSubscription.new(
  channel: "orders",
  start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST
)
token = KubeMQ::CancellationToken.new

client.subscribe_to_events_store(sub, cancellation_token: token) do |event|
  puts "seq=#{event.sequence}: #{event.body}"
end

Yields: EventStoreReceived with channel, metadata, body, tags, id, and sequence attributes.

Returns: Subscription handle.

create_events_store_sender

Creates a streaming events store sender for high-throughput publishing with per-message confirmation.

sender = client.create_events_store_sender
result = sender.publish(msg)
puts "Sent: #{result.id}, confirmed: #{result.sent}"
sender.close

Returns: PubSub::EventStoreSender

EventStoreReceived

Inbound event received from an events store subscription.

AttributeTypeDescription
idStringMessage identifier
channelStringChannel the event arrived on
metadataStringMetadata string
bodyStringMessage payload
tagsHashKey-value tags
sequenceIntegerBroker-assigned sequence number

EventStoreResult

Result returned from send_event_store.

AttributeTypeDescription
idStringMessage identifier
sentBooleanWhether the event was persisted
errorStringError message (if any)

Was this page helpful?

On this page