Skip to main content
Version: Master

Subscribing to events

To subscribe a handler to events in Rails Event Store you need to use #subscribe method on RailsEventStore::Client

Depending on where you decided to keep the configuration that would usually be in config/application.rb or config/initializers/rails_event_store.rb or one of environment files (config/environments/*.rb).

# config/application.rb
module YourAppName
class Application < Rails::Application
config.to_prepare do
Rails.configuration.event_store = event_store = RailsEventStore::Client.new
event_store.subscribe(OrderNotifier.new, to: [OrderCancelled])
end
end
end

Synchronous handlers​

To subscribe to events publication, you can use #subscribe method. It accepts two arguments:

  • subscriber (an event handler) - which can be a function-like object. That means it needs to respond to the #call method. This way both normal objects and lambda expressions are supported. A block of code can also be passed as a subscriber (&subscriber)
  • to: - which is an array of event types. Your subscriber gets notified only when events of types listed here are be published.

An example usage with the object event handler:

class InvoiceReadModel
def call(event)
# Process an event here.
end
end

subscriber = InvoiceReadModel.new
event_store.subscribe(subscriber, to: [InvoiceCreated, InvoiceUpdated])

You can use Proc objects or lambdas in 3 ways:

event_store.subscribe(to: [InvoicePrinted]) do |event|
# Process an event here...
end
invoice_read_model = ->(event) do
# Process an event here...
end

event_store.subscribe(invoice_read_model, to: [InvoiceCreated, InvoiceUpdated])
send_invoice_email =
Proc.new do |event|
# Process an event here...
end

event_store.subscribe(send_invoice_email, to: [InvoiceAccepted])

Handling exceptions​

If your synchronous handlers raise an exception, it might bubble up and cause problems such as reverting a transaction.

class SyncHandler
def call(event)
# ...
raise StandardError, "ups, something went wrong"
end
end
event_store.subscribe(SyncHandler.new, to: [OrderPlaced])
ActiveRecord::Base.transaction do
event_store.publish(OrderPlaced.new)
# sync handlers executed here
# exception will bubble up
# and rollback the transaction
end

If you don't want your event handlers to cause such potential problems, just swallow the exceptions and send them to your exception tracker.

class SyncHandler
def call(event)
# ...
rescue => e
ExceptionTracker.notify(e)
end
end

Fresh handler state​

If you subscribe an instance of a class (SyncHandler.new), the same object is going to be called with new events.

class SyncHandler
def call(event); end
end
handler = SyncHandler.new
event_store.subscribe(handler, to: [OrderPlaced])
event_store.publish(OrderPlaced.new)
# handler is called

event_store.publish(OrderPlaced.new)
# handler is called again

This can be problematic, especially if you use memoization (the @ivar ||= ... pattern).

class SyncHandler
def call(event)
Rails.logger.warn("Order placed by #{customer_id(event)}")
Stats.increase("orders-#{customer_id(event)}", 1)
end

private

def customer_id(event)
@customer_id ||= event.data.fetch(:customer_id)
end
end

because subsequent events would read the same @customer_id which was memoized when the handler was processing a previous event. To avoid that problem, subscribe a callable that builds a fresh instance for every published event, for example a lambda:

event_store.subscribe(->(event) { SyncHandler.new.call(event) }, to: [OrderPlaced])
event_store.publish(OrderPlaced.new(data: { customer_id: 2 }))
# SyncHandler.new.call is invoked (instance A)

event_store.publish(OrderPlaced.new(data: { customer_id: 3 }))
# SyncHandler.new.call is invoked (instance B)

Alternatively, you can define self.call on the handler class to encapsulate the fresh-instance logic:

class SyncHandler
def self.call(event)
new.call(event)
end

def call(event)
# ...
end
end
event_store.subscribe(SyncHandler, to: [OrderPlaced])

The class responds to call, so it qualifies as a callable subscriber. The tradeoff is that the intent — a fresh instance per event — is less visible at the subscription site than with a lambda.

When are sync handlers executed?​

Those handlers are executed immediately after events are stored in the DB.

ActiveRecord::Base.transaction do
order = Order.new(...).save!
event_store.publish(
OrderPlaced.new(data:{order_id: order.id}),
stream_name: "Order-#{order.id}"
)
# Sync handlers executed here
end

Subscribe for all event types​

You can also subscribe for all event types at once. It is especially useful for logging or debugging events. Use subscribe_to_all_events(subsriber1, &subscriber2) method for that.

class EventsLogger
def initialize(logger)
@logger = logger
end

def call(event)
logger.info("#{event.event_type} published. Data: #{event.data.inspect}")
end

private

attr_reader :logger
end

event_store.subscribe_to_all_events(EventsLogger.new(Rails.logger))
event_store.subscribe_to_all_events { |event| puts event.inspect }

Temporary subscriptions​

Rails Event Store supports temporary (dynamic, one-shot) subscriptions for events. The subscriber gets unsubscribed automatically at the end of the provided block.

class CountImportResults
def initialize()
@ok = 0
@error = 0
end

def call(event)
case event
when ProductImported
@ok += 1
when ProductImportFailed
@error += 1
else
raise ArgumentError
end
end
end

class Import
def run(file)
CSV.parse(file) do |row|
if row_imported(row)
event_store.publish(ProductImported.new(...))
else
event_store.publish(ProductImportFailed.new(...))
end
end
end
end

results = CountImportResults.new
event_types = [ProductImported, ProductImportFailed]
event_store.within do
Import.new.run(file)
end.subscribe(results, to: event_types).call

This can be useful also in controllers:

class OperationsController < ApplicationController
def create
event_store
.within { Operation.new.run(file) }
.subscribe(to: [OperationSucceeded]) { redirect_to results_index_path }
.subscribe(to: [OperationFailed]) { render :new }
.call
end
end

Temporarily subscribing to all events is also supported.

event_store
.within { Import.new.run(file) }
.subscribe_to_all_events(EventsLogger.new(Rails.logger))
.subscribe_to_all_events { |event| puts event.inspect }
.call

You start the temporary subscription by providing a block within which the subscriptions will be active. Then you can chain subscribe and subscribe_to_all_events as many times as you want to register temporary subscribers. When you are ready call call to evaluate the provided block with the temporary subscriptions.

Async handlers​

It's possible to also subscribe asynchronous handlers to events. To enqueue asynchronous handlers as background jobs scheduler class is needed. RailsEventStore provides implementation of a scheduler for ActiveJob library. In that case async handlers are just background jobs implemented as:

class SendOrderEmail < ActiveJob::Base
def perform(payload)
event = event_store.deserialize(serializer: RubyEventStore::Serializers::YAML, **payload.symbolize_keys)
email = event.data.fetch(:customer_email)
OrderMailer.notify_customer(email).deliver_now!
end

private

def event_store
Rails.configuration.event_store
end
end

event_store = RailsEventStore::Client.new
event_store.subscribe(SendOrderEmail, to: [OrderPlaced])

You can also use RailsEventStore::AsyncHandler module that will deserialize the event for you:

class SendOrderEmail < ActiveJob::Base
prepend RailsEventStore::AsyncHandler

def perform(event)
email = event.data.fetch(:customer_email)
OrderMailer.notify_customer(email).deliver_now!
end
end

event_store = RailsEventStore::Client.new
event_store.subscribe(SendOrderEmail, to: [OrderPlaced])

If you'd like to rely solely on Sidekiq, you can use the ruby_event_store-sidekiq_scheduler gem instead.

Custom Scheduler​

Alternatively you could implement your own, custom scheduler. It has to respond to the call and verify methods. The sample CustomScheduler could be implemented as:

class CustomScheduler
# method doing actual schedule; record is a RubyEventStore::Record
def call(klass, record)
klass.perform_async(record.serialize(RubyEventStore::Serializers::YAML).to_h)
end

# method which is checking whether given subscriber is correct for this scheduler
def verify(subscriber)
Class === subscriber && subscriber.respond_to?(:perform_async)
end
end

You can also use our scheduler_lint for more confidence that your scheduler is written correctly.

Then you have to initialize RailsEventStore::Client using asynchronous dispatcher with your custom scheduler:

event_store =
RailsEventStore::Client.new(
message_broker: RubyEventStore::Broker.new(
dispatcher: RailsEventStore::AfterCommitDispatcher.new(scheduler: CustomScheduler.new),
),
)

Often you will want to be able to specify both asynchronous and synchronous dispatchers. In that case, you can use ComposedDispatcher, which accepts arbitrary number of dispatchers and dispatch the event to the first subscriber which is accepted (by verify method) by the dispatcher. This is also our default configuration in RailsEventStore.

event_store =
RailsEventStore::Client.new(
message_broker: RubyEventStore::Broker.new(
dispatcher:
RubyEventStore::ComposedDispatcher.new(
RailsEventStore::AfterCommitDispatcher.new(scheduler: CustomScheduler.new), # our asynchronous dispatcher, which expects that subscriber respond to `perform_async` method
RubyEventStore::SyncScheduler.new, # regular synchronous dispatcher
),
),
)

When are async handlers scheduled?​

The default behaviour and examples above use RailsEventStore::AfterCommitDispatcher, which schedule handlers after the transaction is committed.

class SendOrderEmail < ActiveJob::Base
prepend RailsEventStore::AsyncHandler

def perform(event)
email = event.data.fetch(:customer_email)
OrderMailer.notify_customer(email).deliver_now!
end
end

event_store = RailsEventStore::Client.new
event_store.subscribe(SendOrderEmail, to: [OrderPlaced])

# ...

ActiveRecord::Base.transaction do
order = Order.new(...).save!
event_store.publish(
OrderPlaced.new(data:{order_id: order.id}),
stream_name: "Order-#{order.id}"
)
end
# Async handlers such as SendOrderEmail scheduled here, after transaction is committed

Enqueuing async handlers in bulk​

Each async handler is enqueued with its own perform_later call. When a single action publishes many events, or publishes events with many subscribers, those calls fan out — and with a relational database behind ActiveJob each of them is a separate INSERT.

RailsEventStore::ActiveJobBulkScheduler buffers the handlers scheduled within one database transaction and enqueues all of them with a single ActiveJob.perform_all_later call once the transaction commits. Under the hood perform_all_later uses Sidekiq's bulk push API, or insert_all on backends like Solid Queue and GoodJob.

Opting in means swapping the scheduler — the dispatcher stays the same:

event_store = RailsEventStore::Client.new(
message_broker: RubyEventStore::Broker.new(
dispatcher: RubyEventStore::ComposedDispatcher.new(
RailsEventStore::AfterCommitDispatcher.new(
scheduler: RailsEventStore::ActiveJobBulkScheduler.new(serializer: RubyEventStore::Serializers::YAML)
),
RubyEventStore::SyncScheduler.new
)
)
)

Before you switch, mind the trade-offs:

  • It only works with RailsEventStore::AfterCommitDispatcher, which tells the scheduler when the transaction ends. Paired with RubyEventStore::ImmediateDispatcher, or with a custom dispatcher, the buffered jobs are never enqueued.
  • Jobs enqueued with perform_all_later do not run their before_enqueue, around_enqueue or after_enqueue callbacks. If your subscribers rely on them, stay with RailsEventStore::ActiveJobScheduler.
  • A subscriber registered with .set(...) cannot be enqueued in bulk and still costs one perform_later call.
  • It requires Rails 7.2 or newer, and RailsEventStore::AfterCommitDispatcher refuses to take it below that. RailsEventStore::ActiveJobScheduler keeps working on every Rails version the gem supports.

If your handlers reload the event themselves and only need its id, RailsEventStore::ActiveJobIdOnlyBulkScheduler batches the same way while enqueuing the payload of RailsEventStore::ActiveJobIdOnlyScheduler. It takes no serializer, and since it rejects subscribers registered with .set(...) outright, every handler it accepts is batched:

RailsEventStore::AfterCommitDispatcher.new(scheduler: RailsEventStore::ActiveJobIdOnlyBulkScheduler.new)

Buffering is not limited to these two schedulers. A custom scheduler can batch its own backend by including RailsEventStore::BufferingScheduler, collecting into #buffer and shipping in #ship:

class MyBatchingScheduler
include RailsEventStore::BufferingScheduler

def call(klass, record)
buffer << [klass, record]
end

def verify(subscriber)
true
end

private

def ship(entries)
MyBackend.push_bulk(entries)
end
end

RailsEventStore::AfterCommitDispatcher calls #flush on such a scheduler once the transaction commits. The module is what marks the scheduler as buffering — defining a #flush of your own changes nothing, so a scheduler that already has one for unrelated reasons keeps working untouched.

#flush comes from the module rather than being yours to write, because getting it right is easy to get wrong: it has to do nothing when the buffer is empty, stay idempotent across the one call it gets per scheduled handler, and empty the buffer before shipping it, so a backend that raises does not leave the entries behind for the next transaction to ship a second time. #ship is only reached with a non-empty set of entries.

Such a scheduler requires Rails 7.2 or newer, because the dispatcher signals the end of a transaction through ActiveRecord's Transaction#after_commit, added in that version. Below 7.2 RailsEventStore::AfterCommitDispatcher refuses to take a buffering scheduler at all, rather than letting it collect handlers nothing will ever ship.

Each instance buffers on its own, keyed by identity, so schedulers that compare equal do not drain each other. The buffer is per thread as well, and a flush drops it whole — nothing of the scheduler is left behind on a thread that goes on to serve the next request. What lingers instead is a buffer nobody ever flushes: pair the scheduler with a dispatcher that never signals the commit, and it collects on a long-lived thread and holds on.

The rails_event_store gem ships a lint for the contract, alongside the scheduler lint:

require "rails_event_store/spec/buffering_scheduler_lint"

RSpec.describe MyBatchingScheduler do
it_behaves_like "buffering scheduler", MyBatchingScheduler.new
end

Applications using connects_to​

RailsEventStore::AfterCommitDispatcher joins the transaction that wraps your unit of work — the one you open around saving your records and publishing events. It looks for it on ActiveRecord::Base.

In a default Rails application that is the right place to look. ApplicationRecord without connects_to has no connection of its own, so it resolves to ActiveRecord::Base's pool and ApplicationRecord.transaction is that same transaction.

The same holds when ApplicationRecord is declared with primary_abstract_class, as the Rails generator does, and calls connects_to — for example to add a read replica. Active Record registers the primary abstract class's pool under ActiveRecord::Base, so the transaction is still the one the dispatcher finds.

Any other abstract class calling connects_to gets a connection pool of its own — even when it points at the same database. A transaction opened on that class is then invisible to ActiveRecord::Base, which reports none open, and the dispatcher schedules handlers immediately instead of after the commit. A rollback no longer takes the job back. This is the case for an ApplicationRecord marked only with self.abstract_class = true, common in applications created before Rails 7.0, or for a separate abstract class your unit of work runs on. Tell the dispatcher which class owns the transaction:

class ApplicationRecord < ActiveRecord::Base
self.abstract_class = true # not primary_abstract_class
connects_to database: { writing: :primary, reading: :primary_replica }
end

RailsEventStore::AfterCommitDispatcher.new(
scheduler: RailsEventStore::ActiveJobScheduler.new(serializer: RubyEventStore::Serializers::YAML),
transaction_owner: ApplicationRecord
)

The dispatcher watches one class. If your application opens transactions on several classes with pools of their own, only the one passed as transaction_owner: defers handlers to its commit.

This is independent of where the events themselves are stored. Keeping the event store on its own database changes nothing here — by the time handlers are dispatched the repository has already committed its own transaction, so there is nothing on the event store connection left to join. Pass the class your application opens transactions on, not the one you gave to the repository.

Scheduling async handlers immediately​

You can configure your dispatcher slightly different, to schedule async handlers immediately after events are stored in the database. Note the usage of RubyEventStore::ImmediateDispatcher instead of RailsEventStore::AfterCommitDispatcher.

class SendOrderEmail < ActiveJob::Base
def perform(event)
email = event.data.fetch(:customer_email)
OrderMailer.notify_customer(email).deliver_now!
end
end

event_store = RailsEventStore::Client.new(
message_broker: RubyEventStore::Broker.new(
dispatcher: RubyEventStore::ComposedDispatcher.new(
RubyEventStore::ImmediateDispatcher.new(
scheduler: RailsEventStore::ActiveJobScheduler.new(serializer: RubyEventStore::Serializers::YAML)
),
RubyEventStore::SyncScheduler.new
)
)
)

event_store.subscribe(SendOrderEmail, to: [OrderPlaced])

ActiveRecord::Base.transaction do
order = Order.new(...).save!
event_store.publish(
OrderPlaced.new(data:{order_id: order.id}),
stream_name: "Order-#{order.id}"
)
# Async handlers such as SendOrderEmail scheduled here
end

It means that when your ActiveJob adapter (such as sidekiq or resque) is using non-SQL store your handler might get called before the whole transaction is committed or when the transaction was rolled-back.

Transactional outbox​

Both AfterCommitDispatcher and ImmediateDispatcher enqueue jobs directly to Redis. If the network call to Redis fails after the database transaction commits, the job is lost. To guarantee that jobs are never lost, use the transactional outbox pattern: jobs are written to the same database within the same transaction, and a separate process drains them to Redis.

This is worth keeping in mind with ActiveJobBulkScheduler in particular: a failed enqueue is recorded on the job as enqueue_error rather than raised, and perform_all_later does not report back how many jobs made it.

Removing subscriptions​

When you define a new subscription by subscribe method execution it will return a lambda that allows to remove defined subscription.

Rails.configuration.event_store = event_store = RailsEventStore::Client.new
unsubscribe = event_store.subscribe(OrderNotifier.new, to: [OrderCancelled])
# ...and then when subscription is no longer needed
unsubscribe.call

Unsubscribe lambda will remove all subscriptions defined by subscribe method, when you defined subscription as:

unsubscribe = event_store.subscribe(InvoiceReadModel.new, to: [InvoiceCreated, InvoiceUpdated])

and then execute returned lambda both subscriptions will be removed.

It you need temporary subscription to be defined read more here.