Skip to content

Define a reactor fluently

A reactor is usually a class whose handler methods Chronicle discovers. When the handlers belong to something else, such as a host that builds behavior from configuration or a framework that wires reactors for its own components, you can declare the reactor fluently instead. Give it a stable identifier and a handler for each event type.

On<TEvent> subscribes the reactor to that event type and hands the handler the deserialized event. Each overload can take the event alone or the event together with its event context, and each can be synchronous or return a Task:

using Cratis.Chronicle.Events;
[EventType]
public record InvoiceIssued(string Customer, decimal Amount);
[EventType]
public record InvoicePaid(decimal Amount);
public interface IInvoiceNotifications
{
Task Issued(string customer, decimal amount);
Task Paid(EventSourceId invoice, decimal amount);
}
public class InvoiceReactorRegistration
{
public async Task Register(IEventStore eventStore, IInvoiceNotifications notifications)
{
await eventStore.Reactors.Register(
"invoice-notifications",
reactor => reactor
.On<InvoiceIssued>(@event => notifications.Issued(@event.Customer, @event.Amount))
.On<InvoicePaid>((@event, context) => notifications.Paid(context.EventSourceId, @event.Amount)));
}
}

View C# snippet source on GitHub

The reactor observes only the event types it has handlers for. A handler receives the event in the generation it subscribed to, which is the latest generation of its type, even when the event was appended in another generation.

Subscribe declares a catch-all handler. It receives each event as an object together with its context, and it makes the reactor an all-events subscription:

using Cratis.Chronicle.Events;
using Cratis.Chronicle.Reactors;
public interface IAuditTrail
{
Task Record(EventSourceId eventSourceId, object @event);
}
public class AuditReactorRegistration
{
public async Task<IReactorHandler> Register(IEventStore eventStore, IAuditTrail auditTrail)
{
var definition = eventStore.Reactors.Define(
"audit-trail",
reactor => reactor.Subscribe((@event, context) => auditTrail.Record(context.EventSourceId, @event)));
// The definition is inspectable before it is registered.
Console.WriteLine($"{definition.Id} subscribes to all events: {definition.SubscribesToAllEvents}");
Console.WriteLine($"Event types: {string.Join(", ", definition.EventTypes.Select(_ => _.Id))}");
return await eventStore.Reactors.Register(definition);
}
}

View C# snippet source on GitHub

A catch-all handler receives objects, and only an event type the client knows can become one. So an all-events subscription covers every event type the client knows when the reactor is defined, which is every discovered event type, each in its latest generation. An event type the client does not know is not part of the subscription. Neither is an event type that only becomes known after the reactor was defined. EventTypes on the definition lists the exact set Chronicle is asked to deliver.

A reactor can combine typed and catch-all handlers. For each event the typed handlers for its type run first, in the order they were declared, and then the catch-all handlers run.

Reactors.Define builds the definition without registering it. Its Id, EventTypes, ClrTypes, SubscribesToAllEvents, EventSequenceId and IsReplayable show what the reactor will observe. Reactors.Register(definition) registers it. Reactors.Register(id, define) does both in one call.

A fluently defined reactor is an ordinary Chronicle reactor. It has its own observer and keeps its position under its identifier. It has state, failed partitions and replay, and GetHandlerById looks it up later. Handlers are awaited. If one throws or its task fails, the partition fails without advancing past that event, just as it does for a failing method on a reactor class. Use OnEventSequence to observe another event sequence, and NotReplayable when repeating the side effects is unsafe.

Keep the identifier stable across restarts, and make side effects idempotent so they tolerate replay, retries and reconnects.

Registering on an event store sends the reactor to Chronicle immediately, and again whenever the connection is re-established. To declare a reactor before the client exists, from hosting or dependency injection setup, add it to the client options. Every event store the client creates defines it the first time it discovers its artifacts, and registers it with the rest of its reactors from then on:

using Microsoft.Extensions.Hosting;
public static class FluentReactorAtStartup
{
public static void Configure(string[] args)
{
var builder = Host.CreateApplicationBuilder(args);
builder.AddCratisChronicle(configureOptions: options => options.ExplicitArtifacts
.RegisterReactor(
"invoice-log",
reactor => reactor
.On<InvoiceIssued>(@event => Console.WriteLine($"Issued to {@event.Customer}"))
.NotReplayable()));
}
}

View C# snippet source on GitHub

The definition callback runs once per event store, at that first discovery, so an all-events subscription registered this way covers the event types known at that point.

For reactors that handle raw JSON for event types known only by identifier, see Register at Runtime.