Event Processing
Understanding how reducers process events is crucial for building effective read models. This guide covers the event processing model, method patterns, and advanced techniques.
Event Processing Model
Section titled “Event Processing Model”Event Method Discovery
Section titled “Event Method Discovery”Reducers use convention-based method discovery. Chronicle automatically finds and invokes methods that:
- Accept a known event type as the first parameter — the method name itself can be anything; discovery matches on the parameter type, not the name (a descriptive name such as
OrderCreatedfor anOrderCreatedevent is good practice, not a requirement) - Accept the current read model state (nullable) as the second parameter
- Optionally accept
EventContextas the third parameter
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingOrderCreated(Guid OrderId);
[EventType]public record EventProcessingItemAdded(decimal Price);
public record EventProcessingOrderSummary(Guid OrderId, decimal Total, DateTimeOffset LastUpdated);
public class EventProcessingOrderSummaryReducer : IReducerFor<EventProcessingOrderSummary>{ public EventProcessingOrderSummary Created(EventProcessingOrderCreated @event, EventProcessingOrderSummary? current, EventContext context) { return new EventProcessingOrderSummary(@event.OrderId, 0m, context.Occurred); }
public EventProcessingOrderSummary? ItemAdded(EventProcessingItemAdded @event, EventProcessingOrderSummary? current, EventContext context) { if (current is null) return null; // Skip if no order exists
return current with { Total = current.Total + @event.Price, LastUpdated = context.Occurred }; }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.util.UUID
@EventTypedata class EventProcessingMethodDiscoveryOrderCreated(val orderId: UUID)
@EventTypedata class EventProcessingMethodDiscoveryItemAdded(val price: Double)
@ReadModeldata class EventProcessingMethodDiscoveryOrderSummary(val orderId: UUID = UUID(0, 0), val total: Double = 0.0)
@Reducerclass EventProcessingMethodDiscoveryOrderSummaryReducer { // Discovery matches on the first parameter's type, not the method name - a descriptive name // such as `created` for an OrderCreated event is good practice, not a requirement. fun created(event: EventProcessingMethodDiscoveryOrderCreated, current: EventProcessingMethodDiscoveryOrderSummary?, context: EventContext) = EventProcessingMethodDiscoveryOrderSummary(event.orderId, 0.0)
fun itemAdded(event: EventProcessingMethodDiscoveryItemAdded, current: EventProcessingMethodDiscoveryOrderSummary?, context: EventContext): EventProcessingMethodDiscoveryOrderSummary? { if (current == null) return null
return current.copy(total = current.total + event.price) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.util.UUID;
@EventTyperecord EventProcessingMethodDiscoveryOrderCreated(UUID orderId) {}
@EventTyperecord EventProcessingMethodDiscoveryItemAdded(double price) {}
@ReadModelrecord EventProcessingMethodDiscoveryOrderSummary(UUID orderId, double total) { EventProcessingMethodDiscoveryOrderSummary() { this(new UUID(0, 0), 0.0); }}
@Reducerclass EventProcessingMethodDiscoveryOrderSummaryReducer { // Discovery matches on the first parameter's type, not the method name - a descriptive name // such as `created` for an OrderCreated event is good practice, not a requirement. EventProcessingMethodDiscoveryOrderSummary created( EventProcessingMethodDiscoveryOrderCreated event, EventProcessingMethodDiscoveryOrderSummary current, EventContext context) { return new EventProcessingMethodDiscoveryOrderSummary(event.orderId(), 0.0); }
EventProcessingMethodDiscoveryOrderSummary itemAdded( EventProcessingMethodDiscoveryItemAdded event, EventProcessingMethodDiscoveryOrderSummary current, EventContext context) { if (current == null) return null;
return new EventProcessingMethodDiscoveryOrderSummary(current.orderId(), current.total() + event.price()); }}defmodule EventProcessingOrderCreated do use Chronicle.Events.EventType, id: "event-processing-order-created"
defstruct [:order_id]end
defmodule EventProcessingItemAdded do use Chronicle.Events.EventType, id: "event-processing-item-added"
defstruct [:price]end
defmodule EventProcessingOrderSummary do defstruct [:order_id, :total, :last_updated]end
defmodule EventProcessingOrderSummaryReducer do use Chronicle.Reducers.Reducer, model: EventProcessingOrderSummary
alias EventProcessingOrderCreated alias EventProcessingItemAdded
@handles EventProcessingOrderCreated @handles EventProcessingItemAdded
# Dispatch matches on the event struct's type, not a method name — # Chronicle picks the reduce/3 clause whose first-argument pattern matches # the incoming event. @impl true def reduce(%EventProcessingOrderCreated{} = event, _current, context) do %EventProcessingOrderSummary{ order_id: event.order_id, total: 0, last_updated: Map.get(context, :occurred) } end
# Skip if no order exists yet def reduce(%EventProcessingItemAdded{}, nil, _context), do: nil
def reduce(%EventProcessingItemAdded{} = event, current, context) do %{ current | total: current.total + event.price, last_updated: Map.get(context, :occurred) } endendimport { EventContext, eventType, Guid, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingOrderCreated { constructor(readonly orderId: Guid) {}}
@eventType()class EventProcessingItemAdded { constructor(readonly price: number) {}}
class EventProcessingOrderSummary { orderId: Guid = Guid.empty; total = 0; lastUpdated = new Date();}
// Method names must be the exact camelCase of the event's class name -// Chronicle discovers handlers by name, not by parameter type.@reducer('', undefined, EventProcessingOrderSummary)class EventProcessingOrderSummaryReducer { eventProcessingOrderCreated( event: EventProcessingOrderCreated, current: EventProcessingOrderSummary | undefined, context: EventContext ): EventProcessingOrderSummary { return { orderId: event.orderId, total: 0, lastUpdated: context.occurred }; }
eventProcessingItemAdded( event: EventProcessingItemAdded, current: EventProcessingOrderSummary | undefined, context: EventContext ): EventProcessingOrderSummary | undefined { if (!current) return undefined; // Skip if no order exists
return { ...current, total: current.total + event.price, lastUpdated: context.occurred }; }}When more than one method matches
Section titled “When more than one method matches”Discovery matches on the parameter types, so a helper you extract out of a reducer method keeps the same shape the reducer method has and matches the same event type. When that happens, Chronicle picks one method per event type, in this order:
- A public method beats a non-public one. Reducer methods are public by convention; anything else is an implementation detail, so a private helper never displaces the method that calls it.
- The richest signature wins. Between two methods of the same accessibility, the one taking the most parameters asked for the most context — so the overload that also takes the
EventContextis preferred. - Method name, ordinal, so a genuine tie resolves the same way on every run rather than being left to reflection order.
A private method is still a valid reducer method when nothing else handles that event type — the precedence only decides between candidates, it does not disqualify any of them. If a helper is not meant to take part in dispatch at all, change its parameters so it no longer has the reducer shape.
Event Source Isolation
Section titled “Event Source Isolation”Each reducer method is called for a single event source:
- Events for the same event source (e.g., the same order) are processed sequentially
- The
currentparameter contains the current state for that specific event source - Each event source has its own independent state
Sequential Processing
Section titled “Sequential Processing”Events are guaranteed to be processed in sequence order:
- Events are ordered by their sequence number
- Each method is called once per event
- The return value becomes the
currentparameter for the next event
Method Signatures
Section titled “Method Signatures”Basic Synchronous Pattern
Section titled “Basic Synchronous Pattern”The simplest pattern accepts the event and current state:
public interface IEventProcessingBasicSyncPattern<TReadModel, TEvent> where TReadModel : class{ // Process event and return new state TReadModel Process(TEvent @event, TReadModel? current);}import io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingBasicSyncRecorded(val value: Double)
@ReadModeldata class EventProcessingBasicSyncStats(val total: Double = 0.0)
@Reducerclass EventProcessingBasicSyncStatsReducer { // Process event and return new state fun recorded(event: EventProcessingBasicSyncRecorded, current: EventProcessingBasicSyncStats?): EventProcessingBasicSyncStats = EventProcessingBasicSyncStats((current?.total ?: 0.0) + event.value)}import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
@EventTyperecord EventProcessingBasicSyncRecorded(double value) {}
@ReadModelrecord EventProcessingBasicSyncStats(double total) { EventProcessingBasicSyncStats() { this(0.0); }}
@Reducerclass EventProcessingBasicSyncStatsReducer { // Process event and return new state EventProcessingBasicSyncStats recorded(EventProcessingBasicSyncRecorded event, EventProcessingBasicSyncStats current) { double total = current == null ? 0.0 : current.total(); return new EventProcessingBasicSyncStats(total + event.value()); }}defmodule EventProcessingBasicSyncPatternMetricRecorded do use Chronicle.Events.EventType, id: "event-processing-basic-sync-pattern-metric-recorded"
defstruct [:value]end
defmodule EventProcessingBasicSyncPatternStats do defstruct total: 0end
defmodule EventProcessingBasicSyncPatternReducer do use Chronicle.Reducers.Reducer, model: EventProcessingBasicSyncPatternStats
alias EventProcessingBasicSyncPatternMetricRecorded
@handles EventProcessingBasicSyncPatternMetricRecorded
# The simplest pattern: the event and the current state. Elixir's reduce/3 # always receives a third context argument too — ignore it with `_` when # you don't need it. @impl true def reduce(%EventProcessingBasicSyncPatternMetricRecorded{} = event, current, _context) do total = if current, do: current.total, else: 0 %EventProcessingBasicSyncPatternStats{total: total + event.value} endendinterface EventProcessingBasicSyncPattern<TReadModel, TEvent> { // Process event and return new state process(event: TEvent, current: TReadModel | undefined): TReadModel;}Pattern with Event Context
Section titled “Pattern with Event Context”Access event metadata by adding EventContext parameter:
using Cratis.Chronicle.Events;
public interface IEventProcessingPatternWithContext<TReadModel, TEvent> where TReadModel : class{ // Access occurred time, correlation ID, etc. TReadModel Process(TEvent @event, TReadModel? current, EventContext context);}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingPatternWithContextRecorded(val value: Double)
@ReadModeldata class EventProcessingPatternWithContextStats(val total: Double = 0.0, val lastUpdated: java.time.Instant = java.time.Instant.EPOCH)
@Reducerclass EventProcessingPatternWithContextStatsReducer { // Access occurred time, correlation id, etc. via the third parameter fun recorded( event: EventProcessingPatternWithContextRecorded, current: EventProcessingPatternWithContextStats?, context: EventContext ): EventProcessingPatternWithContextStats = EventProcessingPatternWithContextStats((current?.total ?: 0.0) + event.value, context.occurred)}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.Instant;
@EventTyperecord EventProcessingPatternWithContextRecorded(double value) {}
@ReadModelrecord EventProcessingPatternWithContextStats(double total, Instant lastUpdated) { EventProcessingPatternWithContextStats() { this(0.0, Instant.EPOCH); }}
@Reducerclass EventProcessingPatternWithContextStatsReducer { // Access occurred time, correlation id, etc. via the third parameter EventProcessingPatternWithContextStats recorded( EventProcessingPatternWithContextRecorded event, EventProcessingPatternWithContextStats current, EventContext context) { double total = current == null ? 0.0 : current.total(); return new EventProcessingPatternWithContextStats(total + event.value(), context.getOccurred()); }}defmodule EventProcessingPatternWithContextMetricRecorded do use Chronicle.Events.EventType, id: "event-processing-pattern-with-context-metric-recorded"
defstruct [:value]end
defmodule EventProcessingPatternWithContextStats do defstruct [:total, :last_updated]end
defmodule EventProcessingPatternWithContextReducer do use Chronicle.Reducers.Reducer, model: EventProcessingPatternWithContextStats
alias EventProcessingPatternWithContextMetricRecorded
@handles EventProcessingPatternWithContextMetricRecorded
# Access occurred time, sequence number, and other metadata through the # third `context` argument. @impl true def reduce(%EventProcessingPatternWithContextMetricRecorded{} = event, current, context) do total = if current, do: current.total, else: 0
%EventProcessingPatternWithContextStats{ total: total + event.value, last_updated: Map.get(context, :occurred) } endendimport { EventContext } from '@cratis/chronicle';
interface EventProcessingPatternWithContext<TReadModel, TEvent> { // Access occurred time, correlation ID, etc. process(event: TEvent, current: TReadModel | undefined, context: EventContext): TReadModel;}Async Patterns
Section titled “Async Patterns”Both patterns support async methods:
using Cratis.Chronicle.Events;
public interface IEventProcessingAsyncPatterns<TReadModel, TEvent> where TReadModel : class{ // Async without context Task<TReadModel> ProcessAsync(TEvent @event, TReadModel? current);
// Async with context Task<TReadModel> ProcessWithContextAsync(TEvent @event, TReadModel? current, EventContext context);}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingAsyncRecorded(val value: Double)
@EventTypedata class EventProcessingAsyncRecordedWithContext(val value: Double)
@ReadModeldata class EventProcessingAsyncStats(val total: Double = 0.0)
@Reducerclass EventProcessingAsyncStatsReducer { // Async without context suspend fun recorded(event: EventProcessingAsyncRecorded, current: EventProcessingAsyncStats?): EventProcessingAsyncStats = EventProcessingAsyncStats((current?.total ?: 0.0) + event.value)
// Async with context suspend fun recordedWithContext( event: EventProcessingAsyncRecordedWithContext, current: EventProcessingAsyncStats?, context: EventContext ): EventProcessingAsyncStats = EventProcessingAsyncStats((current?.total ?: 0.0) + event.value)}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
@EventTyperecord EventProcessingAsyncRecorded(double value) {}
@EventTyperecord EventProcessingAsyncRecordedWithContext(double value) {}
@ReadModelrecord EventProcessingAsyncStats(double total) { EventProcessingAsyncStats() { this(0.0); }}
// Java reducer methods are never suspending - they run on the observation thread directly, so// there is no separate "async" shape to opt into. Await I/O with whatever blocking or// CompletableFuture-based client your dependency exposes.@Reducerclass EventProcessingAsyncStatsReducer { EventProcessingAsyncStats recorded(EventProcessingAsyncRecorded event, EventProcessingAsyncStats current) { double total = current == null ? 0.0 : current.total(); return new EventProcessingAsyncStats(total + event.value()); }
EventProcessingAsyncStats recordedWithContext( EventProcessingAsyncRecordedWithContext event, EventProcessingAsyncStats current, EventContext context) { double total = current == null ? 0.0 : current.total(); return new EventProcessingAsyncStats(total + event.value()); }}Elixir does not support this workflow yet.import { EventContext } from '@cratis/chronicle';
interface EventProcessingAsyncPatterns<TReadModel, TEvent> { // Async without context process(event: TEvent, current: TReadModel | undefined): Promise<TReadModel>;
// Async with context processWithContext(event: TEvent, current: TReadModel | undefined, context: EventContext): Promise<TReadModel>;}Event Context
Section titled “Event Context”The EventContext provides metadata about the event:
using Cratis.Chronicle.Auditing;using Cratis.Chronicle.Events;using Cratis.Chronicle.Identities;
// Illustrative subset of Cratis.Chronicle.Events.EventContext's real shapepublic record EventProcessingEventContextShape( EventSequenceNumber SequenceNumber, EventSourceId EventSourceId, EventType EventType, DateTimeOffset Occurred, CorrelationId CorrelationId, IEnumerable<Causation> Causation, Identity CausedBy);// ... and more — see EventContext for the full member listimport io.cratis.chronicle.auditing.Causationimport io.cratis.chronicle.events.EventTypeDescriptorimport io.cratis.chronicle.identity.Identityimport java.time.Instantimport java.util.UUID
// Illustrative subset of io.cratis.chronicle.events.EventContext's real shapedata class EventProcessingEventContextShape( val sequenceNumber: Long, val eventSourceId: String, val eventType: EventTypeDescriptor, val occurred: Instant, val correlationId: UUID, val causedBy: Identity, val causation: List<Causation> = emptyList())// ... and more - see EventContext for the full member listimport io.cratis.chronicle.auditing.Causation;import io.cratis.chronicle.events.EventTypeDescriptor;import io.cratis.chronicle.identity.Identity;
import java.time.Instant;import java.util.List;import java.util.UUID;
// Illustrative subset of io.cratis.chronicle.events.EventContext's real shaperecord EventProcessingEventContextShape( long sequenceNumber, String eventSourceId, EventTypeDescriptor eventType, Instant occurred, UUID correlationId, Identity causedBy, List<Causation> causation) {}// ... and more - see EventContext for the full member listdefmodule EventProcessingEventContextShape do @moduledoc """ Illustrative subset of the context map reduce/3 receives as its third argument — see Chronicle.Reducers.Reducer for the authoritative list. """
# %{event_source_id: ..., sequence_number: ..., occurred: ..., observation_state: ...} def describe(%{ event_source_id: event_source_id, sequence_number: sequence_number, occurred: occurred, observation_state: observation_state }) do {event_source_id, sequence_number, occurred, observation_state} endendimport { CausationEntry, EventType } from '@cratis/chronicle';
// Illustrative subset of the real EventContext shape from '@cratis/chronicle'interface EventProcessingEventContextShape { readonly sequenceNumber: bigint; readonly eventSourceId: string; readonly eventType: EventType; readonly occurred: Date; readonly correlationId: string; readonly causation: ReadonlyArray<CausationEntry>;}// ... see EventContext for the authoritative member listUsing Event Context
Section titled “Using Event Context”using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingContextOrderPlaced(Guid OrderId, decimal Amount);
public record EventProcessingOrderSummaryWithContext( Guid OrderId, decimal Total, DateTimeOffset PlacedAt, string PlacedBy, CorrelationId CorrelationId);
public class EventProcessingOrderSummaryWithContextReducer : IReducerFor<EventProcessingOrderSummaryWithContext>{ public EventProcessingOrderSummaryWithContext Placed(EventProcessingContextOrderPlaced @event, EventProcessingOrderSummaryWithContext? current, EventContext context) => new( OrderId: @event.OrderId, Total: @event.Amount, PlacedAt: context.Occurred, PlacedBy: context.CausedBy.ToString()!, CorrelationId: context.CorrelationId);}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.time.Instantimport java.util.UUID
@EventTypedata class EventProcessingContextOrderPlaced(val orderId: UUID, val amount: Double)
@ReadModeldata class EventProcessingOrderSummaryWithContext( val orderId: UUID = UUID(0, 0), val total: Double = 0.0, val placedAt: Instant = Instant.EPOCH, val placedBy: String = "", val correlationId: UUID = UUID(0, 0))
@Reducerclass EventProcessingOrderSummaryWithContextReducer { fun placed(event: EventProcessingContextOrderPlaced, current: EventProcessingOrderSummaryWithContext?, context: EventContext) = EventProcessingOrderSummaryWithContext( orderId = event.orderId, total = event.amount, placedAt = context.occurred, placedBy = context.causedBy.subject, correlationId = context.correlationId )}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.Instant;import java.util.UUID;
@EventTyperecord EventProcessingContextOrderPlaced(UUID orderId, double amount) {}
@ReadModelrecord EventProcessingOrderSummaryWithContext( UUID orderId, double total, Instant placedAt, String placedBy, UUID correlationId) { EventProcessingOrderSummaryWithContext() { this(new UUID(0, 0), 0.0, Instant.EPOCH, "", new UUID(0, 0)); }}
@Reducerclass EventProcessingOrderSummaryWithContextReducer { EventProcessingOrderSummaryWithContext placed(EventProcessingContextOrderPlaced event, EventProcessingOrderSummaryWithContext current, EventContext context) { return new EventProcessingOrderSummaryWithContext( event.orderId(), event.amount(), context.getOccurred(), context.getCausedBy().getSubject(), context.getCorrelationId()); }}defmodule EventProcessingContextOrderPlaced do use Chronicle.Events.EventType, id: "event-processing-context-order-placed"
defstruct [:order_id, :amount]end
defmodule EventProcessingOrderSummaryWithContext do defstruct [:order_id, :total, :placed_at, :sequence_number]end
defmodule EventProcessingOrderSummaryWithContextReducer do use Chronicle.Reducers.Reducer, model: EventProcessingOrderSummaryWithContext
alias EventProcessingContextOrderPlaced
@handles EventProcessingContextOrderPlaced
@impl true def reduce(%EventProcessingContextOrderPlaced{} = event, _current, context) do %EventProcessingOrderSummaryWithContext{ order_id: event.order_id, total: event.amount, placed_at: Map.get(context, :occurred), sequence_number: Map.get(context, :sequence_number) } endendimport { EventContext, eventType, Guid, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingContextOrderPlaced { constructor(readonly orderId: Guid, readonly amount: number) {}}
class EventProcessingOrderSummaryWithContext { orderId: Guid = Guid.empty; total = 0; placedAt = new Date(); correlationId = '';}
@reducer('', undefined, EventProcessingOrderSummaryWithContext)class EventProcessingOrderSummaryWithContextReducer { eventProcessingContextOrderPlaced( event: EventProcessingContextOrderPlaced, current: EventProcessingOrderSummaryWithContext | undefined, context: EventContext ): EventProcessingOrderSummaryWithContext { return { orderId: event.orderId, total: event.amount, placedAt: context.occurred, correlationId: context.correlationId }; }}Current State Parameter
Section titled “Current State Parameter”The current parameter represents the previously computed state for this event source.
First Event
Section titled “First Event”When processing the first event for an event source:
currentisnull- Initialize your read model with appropriate values
- You can return
null(declare the handler’s return type asTReadModel?) to skip creating state for certain events
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingDataRecorded(decimal Value);
public record EventProcessingAnalytics(int EventCount, DateTimeOffset FirstEventTime, DateTimeOffset LastEventTime, decimal TotalValue);
public class EventProcessingAnalyticsReducer : IReducerFor<EventProcessingAnalytics>{ public EventProcessingAnalytics Recorded(EventProcessingDataRecorded @event, EventProcessingAnalytics? current, EventContext context) { if (current is null) { // First event - initialize state return new EventProcessingAnalytics( EventCount: 1, FirstEventTime: context.Occurred, LastEventTime: context.Occurred, TotalValue: @event.Value); }
// Update existing state return current with { EventCount = current.EventCount + 1, LastEventTime = context.Occurred, TotalValue = current.TotalValue + @event.Value }; }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.util.UUID
@EventTypedata class EventProcessingOrderCreated(val orderId: UUID)
@EventTypedata class EventProcessingItemAdded(val price: Double)
@ReadModeldata class EventProcessingOrderSummary(val orderId: UUID = UUID(0, 0), val total: Double = 0.0, val lastUpdated: java.time.Instant = java.time.Instant.EPOCH)
@Reducerclass EventProcessingOrderSummaryReducer { fun created(event: EventProcessingOrderCreated, current: EventProcessingOrderSummary?, context: EventContext) = EventProcessingOrderSummary(event.orderId, 0.0, context.occurred)
fun itemAdded(event: EventProcessingItemAdded, current: EventProcessingOrderSummary?, context: EventContext): EventProcessingOrderSummary? { if (current == null) return null // Skip if no order exists
return current.copy(total = current.total + event.price, lastUpdated = context.occurred) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.Instant;import java.util.UUID;
@EventTyperecord EventProcessingOrderCreated(UUID orderId) {}
@EventTyperecord EventProcessingItemAdded(double price) {}
@ReadModelrecord EventProcessingOrderSummary(UUID orderId, double total, Instant lastUpdated) { EventProcessingOrderSummary() { this(new UUID(0, 0), 0.0, Instant.EPOCH); }}
@Reducerclass EventProcessingOrderSummaryReducer { EventProcessingOrderSummary created(EventProcessingOrderCreated event, EventProcessingOrderSummary current, EventContext context) { return new EventProcessingOrderSummary(event.orderId(), 0.0, context.getOccurred()); }
EventProcessingOrderSummary itemAdded(EventProcessingItemAdded event, EventProcessingOrderSummary current, EventContext context) { if (current == null) return null; // Skip if no order exists
return new EventProcessingOrderSummary(current.orderId(), current.total() + event.price(), context.getOccurred()); }}defmodule EventProcessingDataRecorded do use Chronicle.Events.EventType, id: "event-processing-data-recorded"
defstruct [:value]end
defmodule EventProcessingAnalytics do defstruct [:event_count, :first_event_time, :last_event_time, :total_value]end
defmodule EventProcessingAnalyticsReducer do use Chronicle.Reducers.Reducer, model: EventProcessingAnalytics
alias EventProcessingDataRecorded
@handles EventProcessingDataRecorded
@impl true def reduce(%EventProcessingDataRecorded{} = event, nil, context) do occurred = Map.get(context, :occurred)
%EventProcessingAnalytics{ event_count: 1, first_event_time: occurred, last_event_time: occurred, total_value: event.value } end
def reduce(%EventProcessingDataRecorded{} = event, current, context) do %{ current | event_count: current.event_count + 1, last_event_time: Map.get(context, :occurred), total_value: current.total_value + event.value } endendimport { EventContext, eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingDataRecorded { constructor(readonly value: number) {}}
class EventProcessingAnalytics { eventCount = 0; firstEventTime = new Date(); lastEventTime = new Date(); totalValue = 0;}
@reducer('', undefined, EventProcessingAnalytics)class EventProcessingAnalyticsReducer { eventProcessingDataRecorded( event: EventProcessingDataRecorded, current: EventProcessingAnalytics | undefined, context: EventContext ): EventProcessingAnalytics { if (!current) { // First event - initialize state return { eventCount: 1, firstEventTime: context.occurred, lastEventTime: context.occurred, totalValue: event.value }; }
// Update existing state return { ...current, eventCount: current.eventCount + 1, lastEventTime: context.occurred, totalValue: current.totalValue + event.value }; }}Subsequent Events
Section titled “Subsequent Events”For subsequent events:
currentcontains the state from the previous event- Use record’s
withexpression to create modified copies - Return the new state
Processing Patterns
Section titled “Processing Patterns”Pattern 1: Accumulation
Section titled “Pattern 1: Accumulation”Accumulate values across events:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingMetricRecorded(decimal Value);
public record EventProcessingStatistics(decimal Sum, int Count, decimal Average);
public class EventProcessingStatisticsReducer : IReducerFor<EventProcessingStatistics>{ public EventProcessingStatistics Recorded(EventProcessingMetricRecorded @event, EventProcessingStatistics? current) { var sum = (current?.Sum ?? 0) + @event.Value; var count = (current?.Count ?? 0) + 1;
return new EventProcessingStatistics(sum, count, sum / count); }}import io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingMetricRecorded(val value: Double)
@ReadModeldata class EventProcessingStatistics(val sum: Double = 0.0, val count: Int = 0, val average: Double = 0.0)
@Reducerclass EventProcessingStatisticsReducer { fun recorded(event: EventProcessingMetricRecorded, current: EventProcessingStatistics?): EventProcessingStatistics { val sum = (current?.sum ?: 0.0) + event.value val count = (current?.count ?: 0) + 1 return EventProcessingStatistics(sum, count, sum / count) }}import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
@EventTyperecord EventProcessingMetricRecorded(double value) {}
@ReadModelrecord EventProcessingStatistics(double sum, int count, double average) { EventProcessingStatistics() { this(0.0, 0, 0.0); }}
@Reducerclass EventProcessingStatisticsReducer { EventProcessingStatistics recorded(EventProcessingMetricRecorded event, EventProcessingStatistics current) { double sum = (current == null ? 0.0 : current.sum()) + event.value(); int count = (current == null ? 0 : current.count()) + 1; return new EventProcessingStatistics(sum, count, sum / count); }}defmodule EventProcessingMetricRecorded do use Chronicle.Events.EventType, id: "event-processing-metric-recorded"
defstruct [:value]end
defmodule EventProcessingStatistics do defstruct sum: 0, count: 0, average: 0end
defmodule EventProcessingStatisticsReducer do use Chronicle.Reducers.Reducer, model: EventProcessingStatistics
alias EventProcessingMetricRecorded
@handles EventProcessingMetricRecorded
@impl true def reduce(%EventProcessingMetricRecorded{} = event, current, _context) do sum = (if current, do: current.sum, else: 0) + event.value count = (if current, do: current.count, else: 0) + 1
%EventProcessingStatistics{sum: sum, count: count, average: sum / count} endendimport { eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingMetricRecorded { constructor(readonly value: number) {}}
class EventProcessingStatistics { sum = 0; count = 0; average = 0;}
@reducer('', undefined, EventProcessingStatistics)class EventProcessingStatisticsReducer { eventProcessingMetricRecorded(event: EventProcessingMetricRecorded, current: EventProcessingStatistics | undefined): EventProcessingStatistics { const sum = (current?.sum ?? 0) + event.value; const count = (current?.count ?? 0) + 1;
return { sum, count, average: sum / count }; }}Pattern 2: State Transitions
Section titled “Pattern 2: State Transitions”Track state changes through events:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingOrderCreatedForStatus(Guid OrderId);
[EventType]public record EventProcessingOrderPaid(Guid OrderId);
[EventType]public record EventProcessingOrderShipped(Guid OrderId);
[EventType]public record EventProcessingOrderDelivered(Guid OrderId);
[EventType]public record EventProcessingOrderCancelled(Guid OrderId);
public record EventProcessingOrderStatus(string State, DateTimeOffset LastUpdated);
public class EventProcessingOrderStatusReducer : IReducerFor<EventProcessingOrderStatus>{ public EventProcessingOrderStatus Created(EventProcessingOrderCreatedForStatus @event, EventProcessingOrderStatus? current, EventContext context) => new EventProcessingOrderStatus("Created", context.Occurred);
public EventProcessingOrderStatus Paid(EventProcessingOrderPaid @event, EventProcessingOrderStatus? current, EventContext context) => new EventProcessingOrderStatus("Paid", context.Occurred);
public EventProcessingOrderStatus Shipped(EventProcessingOrderShipped @event, EventProcessingOrderStatus? current, EventContext context) => new EventProcessingOrderStatus("Shipped", context.Occurred);
public EventProcessingOrderStatus Delivered(EventProcessingOrderDelivered @event, EventProcessingOrderStatus? current, EventContext context) => new EventProcessingOrderStatus("Delivered", context.Occurred);
public EventProcessingOrderStatus Cancelled(EventProcessingOrderCancelled @event, EventProcessingOrderStatus? current, EventContext context) => new EventProcessingOrderStatus("Cancelled", context.Occurred);}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.time.Instantimport java.util.UUID
@EventTypedata class EventProcessingOrderCreatedForStatus(val orderId: UUID)
@EventTypedata class EventProcessingOrderPaid(val orderId: UUID)
@EventTypedata class EventProcessingOrderShipped(val orderId: UUID)
@EventTypedata class EventProcessingOrderDelivered(val orderId: UUID)
@EventTypedata class EventProcessingOrderCancelled(val orderId: UUID)
@ReadModeldata class EventProcessingOrderStatus(val state: String = "", val lastUpdated: Instant = Instant.EPOCH)
@Reducerclass EventProcessingOrderStatusReducer { fun created(event: EventProcessingOrderCreatedForStatus, current: EventProcessingOrderStatus?, context: EventContext) = EventProcessingOrderStatus("Created", context.occurred)
fun paid(event: EventProcessingOrderPaid, current: EventProcessingOrderStatus?, context: EventContext) = EventProcessingOrderStatus("Paid", context.occurred)
fun shipped(event: EventProcessingOrderShipped, current: EventProcessingOrderStatus?, context: EventContext) = EventProcessingOrderStatus("Shipped", context.occurred)
fun delivered(event: EventProcessingOrderDelivered, current: EventProcessingOrderStatus?, context: EventContext) = EventProcessingOrderStatus("Delivered", context.occurred)
fun cancelled(event: EventProcessingOrderCancelled, current: EventProcessingOrderStatus?, context: EventContext) = EventProcessingOrderStatus("Cancelled", context.occurred)}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.Instant;import java.util.UUID;
@EventTyperecord EventProcessingOrderCreatedForStatus(UUID orderId) {}
@EventTyperecord EventProcessingOrderPaid(UUID orderId) {}
@EventTyperecord EventProcessingOrderShipped(UUID orderId) {}
@EventTyperecord EventProcessingOrderDelivered(UUID orderId) {}
@EventTyperecord EventProcessingOrderCancelled(UUID orderId) {}
@ReadModelrecord EventProcessingOrderStatus(String state, Instant lastUpdated) { EventProcessingOrderStatus() { this("", Instant.EPOCH); }}
@Reducerclass EventProcessingOrderStatusReducer { EventProcessingOrderStatus created(EventProcessingOrderCreatedForStatus event, EventProcessingOrderStatus current, EventContext context) { return new EventProcessingOrderStatus("Created", context.getOccurred()); }
EventProcessingOrderStatus paid(EventProcessingOrderPaid event, EventProcessingOrderStatus current, EventContext context) { return new EventProcessingOrderStatus("Paid", context.getOccurred()); }
EventProcessingOrderStatus shipped(EventProcessingOrderShipped event, EventProcessingOrderStatus current, EventContext context) { return new EventProcessingOrderStatus("Shipped", context.getOccurred()); }
EventProcessingOrderStatus delivered(EventProcessingOrderDelivered event, EventProcessingOrderStatus current, EventContext context) { return new EventProcessingOrderStatus("Delivered", context.getOccurred()); }
EventProcessingOrderStatus cancelled(EventProcessingOrderCancelled event, EventProcessingOrderStatus current, EventContext context) { return new EventProcessingOrderStatus("Cancelled", context.getOccurred()); }}defmodule EventProcessingOrderCreatedForStatus do use Chronicle.Events.EventType, id: "event-processing-order-created-for-status"
defstruct [:order_id]end
defmodule EventProcessingOrderPaid do use Chronicle.Events.EventType, id: "event-processing-order-paid"
defstruct [:order_id]end
defmodule EventProcessingOrderShipped do use Chronicle.Events.EventType, id: "event-processing-order-shipped"
defstruct [:order_id]end
defmodule EventProcessingOrderDelivered do use Chronicle.Events.EventType, id: "event-processing-order-delivered"
defstruct [:order_id]end
defmodule EventProcessingOrderCancelled do use Chronicle.Events.EventType, id: "event-processing-order-cancelled"
defstruct [:order_id]end
defmodule EventProcessingOrderStatus do defstruct [:state, :last_updated]end
defmodule EventProcessingOrderStatusReducer do use Chronicle.Reducers.Reducer, model: EventProcessingOrderStatus
alias EventProcessingOrderCreatedForStatus alias EventProcessingOrderPaid alias EventProcessingOrderShipped alias EventProcessingOrderDelivered alias EventProcessingOrderCancelled
@handles EventProcessingOrderCreatedForStatus @handles EventProcessingOrderPaid @handles EventProcessingOrderShipped @handles EventProcessingOrderDelivered @handles EventProcessingOrderCancelled
@impl true def reduce(%EventProcessingOrderCreatedForStatus{}, _current, context), do: %EventProcessingOrderStatus{state: "created", last_updated: Map.get(context, :occurred)}
def reduce(%EventProcessingOrderPaid{}, _current, context), do: %EventProcessingOrderStatus{state: "paid", last_updated: Map.get(context, :occurred)}
def reduce(%EventProcessingOrderShipped{}, _current, context), do: %EventProcessingOrderStatus{state: "shipped", last_updated: Map.get(context, :occurred)}
def reduce(%EventProcessingOrderDelivered{}, _current, context), do: %EventProcessingOrderStatus{state: "delivered", last_updated: Map.get(context, :occurred)}
def reduce(%EventProcessingOrderCancelled{}, _current, context), do: %EventProcessingOrderStatus{state: "cancelled", last_updated: Map.get(context, :occurred)}endimport { EventContext, eventType, Guid, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingOrderCreatedForStatus { constructor(readonly orderId: Guid) {}}
@eventType()class EventProcessingOrderPaid { constructor(readonly orderId: Guid) {}}
@eventType()class EventProcessingOrderShipped { constructor(readonly orderId: Guid) {}}
@eventType()class EventProcessingOrderDelivered { constructor(readonly orderId: Guid) {}}
@eventType()class EventProcessingOrderCancelled { constructor(readonly orderId: Guid) {}}
class EventProcessingOrderStatus { state = ''; lastUpdated = new Date();}
@reducer('', undefined, EventProcessingOrderStatus)class EventProcessingOrderStatusReducer { eventProcessingOrderCreatedForStatus(event: EventProcessingOrderCreatedForStatus, current: EventProcessingOrderStatus | undefined, context: EventContext): EventProcessingOrderStatus { return { state: 'Created', lastUpdated: context.occurred }; }
eventProcessingOrderPaid(event: EventProcessingOrderPaid, current: EventProcessingOrderStatus | undefined, context: EventContext): EventProcessingOrderStatus { return { state: 'Paid', lastUpdated: context.occurred }; }
eventProcessingOrderShipped(event: EventProcessingOrderShipped, current: EventProcessingOrderStatus | undefined, context: EventContext): EventProcessingOrderStatus { return { state: 'Shipped', lastUpdated: context.occurred }; }
eventProcessingOrderDelivered(event: EventProcessingOrderDelivered, current: EventProcessingOrderStatus | undefined, context: EventContext): EventProcessingOrderStatus { return { state: 'Delivered', lastUpdated: context.occurred }; }
eventProcessingOrderCancelled(event: EventProcessingOrderCancelled, current: EventProcessingOrderStatus | undefined, context: EventContext): EventProcessingOrderStatus { return { state: 'Cancelled', lastUpdated: context.occurred }; }}Pattern 3: Collection Building
Section titled “Pattern 3: Collection Building”Build collections from events:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingCustomerAction(string Type, string Description);
public record EventProcessingActivity(string Type, DateTimeOffset Timestamp, string Description);public record EventProcessingCustomerActivityLog(List<EventProcessingActivity> Activities);
public class EventProcessingCustomerActivityLogReducer : IReducerFor<EventProcessingCustomerActivityLog>{ public EventProcessingCustomerActivityLog Recorded(EventProcessingCustomerAction @event, EventProcessingCustomerActivityLog? current, EventContext context) { // Copy rather than mutate — current.Activities may still be referenced by a held snapshot var activities = new List<EventProcessingActivity>(current?.Activities ?? []);
activities.Add(new EventProcessingActivity( @event.Type, context.Occurred, @event.Description));
return new EventProcessingCustomerActivityLog(activities); }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.time.Instant
@EventTypedata class EventProcessingCustomerAction(val type: String, val description: String)
data class EventProcessingActivity(val type: String, val timestamp: Instant, val description: String)
@ReadModeldata class EventProcessingCustomerActivityLog(val activities: List<EventProcessingActivity> = emptyList())
@Reducerclass EventProcessingCustomerActivityLogReducer { fun recorded( event: EventProcessingCustomerAction, current: EventProcessingCustomerActivityLog?, context: EventContext ): EventProcessingCustomerActivityLog { // Copy rather than mutate - current.activities may still be referenced by a held snapshot val activities = (current?.activities ?: emptyList()) + EventProcessingActivity(event.type, context.occurred, event.description)
return EventProcessingCustomerActivityLog(activities) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.Instant;import java.util.ArrayList;import java.util.List;
@EventTyperecord EventProcessingCustomerAction(String type, String description) {}
record EventProcessingActivity(String type, Instant timestamp, String description) {}
@ReadModelrecord EventProcessingCustomerActivityLog(List<EventProcessingActivity> activities) { EventProcessingCustomerActivityLog() { this(List.of()); }}
@Reducerclass EventProcessingCustomerActivityLogReducer { EventProcessingCustomerActivityLog recorded( EventProcessingCustomerAction event, EventProcessingCustomerActivityLog current, EventContext context) { // Copy rather than mutate - current.activities() may still be referenced by a held snapshot List<EventProcessingActivity> activities = new ArrayList<>(current == null ? List.of() : current.activities()); activities.add(new EventProcessingActivity(event.type(), context.getOccurred(), event.description()));
return new EventProcessingCustomerActivityLog(activities); }}defmodule EventProcessingCustomerAction do use Chronicle.Events.EventType, id: "event-processing-customer-action"
defstruct [:type, :description]end
defmodule EventProcessingActivity do defstruct [:type, :timestamp, :description]end
defmodule EventProcessingCustomerActivityLog do defstruct activities: []end
defmodule EventProcessingCustomerActivityLogReducer do use Chronicle.Reducers.Reducer, model: EventProcessingCustomerActivityLog
alias EventProcessingCustomerAction alias EventProcessingActivity
@handles EventProcessingCustomerAction
@impl true def reduce(%EventProcessingCustomerAction{} = event, current, context) do activities = if current, do: current.activities, else: []
activity = %EventProcessingActivity{ type: event.type, timestamp: Map.get(context, :occurred), description: event.description }
%EventProcessingCustomerActivityLog{activities: activities ++ [activity]} endendimport { EventContext, eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingCustomerAction { constructor(readonly type: string, readonly description: string) {}}
class EventProcessingActivity { type = ''; timestamp = new Date(); description = '';}
class EventProcessingCustomerActivityLog { activities: EventProcessingActivity[] = [];}
@reducer('', undefined, EventProcessingCustomerActivityLog)class EventProcessingCustomerActivityLogReducer { eventProcessingCustomerAction( event: EventProcessingCustomerAction, current: EventProcessingCustomerActivityLog | undefined, context: EventContext ): EventProcessingCustomerActivityLog { // Copy rather than mutate — current.activities may still be referenced by a held snapshot const activities = [...(current?.activities ?? [])];
activities.push({ type: event.type, timestamp: context.occurred, description: event.description });
return { activities }; }}Pattern 4: Time-Based Aggregation
Section titled “Pattern 4: Time-Based Aggregation”Aggregate events within time windows:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingHourlyMetricRecorded(decimal Value);
public record EventProcessingHourlyMetrics(Dictionary<int, decimal> MetricsByHour);
public class EventProcessingHourlyMetricsReducer : IReducerFor<EventProcessingHourlyMetrics>{ public EventProcessingHourlyMetrics Recorded(EventProcessingHourlyMetricRecorded @event, EventProcessingHourlyMetrics? current, EventContext context) { var metricsByHour = new Dictionary<int, decimal>(current?.MetricsByHour ?? []); var hour = context.Occurred.Hour;
if (!metricsByHour.ContainsKey(hour)) metricsByHour[hour] = 0;
metricsByHour[hour] += @event.Value;
return new EventProcessingHourlyMetrics(metricsByHour); }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingHourlyMetricRecorded(val value: Double)
@ReadModeldata class EventProcessingHourlyMetrics(val metricsByHour: Map<Int, Double> = emptyMap())
@Reducerclass EventProcessingHourlyMetricsReducer { fun recorded(event: EventProcessingHourlyMetricRecorded, current: EventProcessingHourlyMetrics?, context: EventContext): EventProcessingHourlyMetrics { val hour = context.occurred.atZone(java.time.ZoneOffset.UTC).hour val metricsByHour = (current?.metricsByHour ?: emptyMap()) + (hour to (current?.metricsByHour?.get(hour) ?: 0.0) + event.value)
return EventProcessingHourlyMetrics(metricsByHour) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.time.ZoneOffset;import java.util.HashMap;import java.util.Map;
@EventTyperecord EventProcessingHourlyMetricRecorded(double value) {}
@ReadModelrecord EventProcessingHourlyMetrics(Map<Integer, Double> metricsByHour) { EventProcessingHourlyMetrics() { this(Map.of()); }}
@Reducerclass EventProcessingHourlyMetricsReducer { EventProcessingHourlyMetrics recorded(EventProcessingHourlyMetricRecorded event, EventProcessingHourlyMetrics current, EventContext context) { Map<Integer, Double> metricsByHour = new HashMap<>(current == null ? Map.of() : current.metricsByHour()); int hour = context.getOccurred().atZone(ZoneOffset.UTC).getHour();
metricsByHour.merge(hour, event.value(), Double::sum);
return new EventProcessingHourlyMetrics(metricsByHour); }}defmodule EventProcessingHourlyMetricRecorded do use Chronicle.Events.EventType, id: "event-processing-hourly-metric-recorded"
defstruct [:value]end
defmodule EventProcessingHourlyMetrics do defstruct metrics_by_hour: %{}end
defmodule EventProcessingHourlyMetricsReducer do use Chronicle.Reducers.Reducer, model: EventProcessingHourlyMetrics
alias EventProcessingHourlyMetricRecorded
@handles EventProcessingHourlyMetricRecorded
@impl true def reduce(%EventProcessingHourlyMetricRecorded{} = event, current, context) do metrics_by_hour = if current, do: current.metrics_by_hour, else: %{} {:ok, occurred, _} = DateTime.from_iso8601(Map.get(context, :occurred))
updated = Map.update(metrics_by_hour, occurred.hour, event.value, &(&1 + event.value))
%EventProcessingHourlyMetrics{metrics_by_hour: updated} endendimport { EventContext, eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingHourlyMetricRecorded { constructor(readonly value: number) {}}
class EventProcessingHourlyMetrics { metricsByHour: Record<number, number> = {};}
@reducer('', undefined, EventProcessingHourlyMetrics)class EventProcessingHourlyMetricsReducer { eventProcessingHourlyMetricRecorded( event: EventProcessingHourlyMetricRecorded, current: EventProcessingHourlyMetrics | undefined, context: EventContext ): EventProcessingHourlyMetrics { const metricsByHour = { ...(current?.metricsByHour ?? {}) }; const hour = context.occurred.getHours();
metricsByHour[hour] = (metricsByHour[hour] ?? 0) + event.value;
return { metricsByHour }; }}Pattern 5: Conditional Processing
Section titled “Pattern 5: Conditional Processing”Skip processing based on conditions:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingAccountOpened(Guid AccountId);
[EventType]public record EventProcessingDepositMade(decimal Amount);
[EventType]public record EventProcessingAccountClosed;
public record EventProcessingAccount(Guid AccountId, decimal Balance, bool IsActive);
public class EventProcessingAccountReducer : IReducerFor<EventProcessingAccount>{ public EventProcessingAccount Opened(EventProcessingAccountOpened @event, EventProcessingAccount? current, EventContext context) { return new EventProcessingAccount(@event.AccountId, 0m, true); }
public EventProcessingAccount? DepositMade(EventProcessingDepositMade @event, EventProcessingAccount? current, EventContext context) { // Skip if account doesn't exist or is not active if (current is null || !current.IsActive) return current;
return current with { Balance = current.Balance + @event.Amount }; }
public EventProcessingAccount? Closed(EventProcessingAccountClosed @event, EventProcessingAccount? current, EventContext context) { if (current is null) return null;
return current with { IsActive = false }; }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.util.UUID
@EventTypedata class EventProcessingAccountOpened(val accountId: UUID)
@EventTypedata class EventProcessingDepositMade(val amount: Double)
@EventTypeclass EventProcessingAccountClosed
@ReadModeldata class EventProcessingAccount(val accountId: UUID = UUID(0, 0), val balance: Double = 0.0, val isActive: Boolean = false)
@Reducerclass EventProcessingAccountReducer { fun opened(event: EventProcessingAccountOpened, current: EventProcessingAccount?, context: EventContext) = EventProcessingAccount(event.accountId, 0.0, true)
fun depositMade(event: EventProcessingDepositMade, current: EventProcessingAccount?, context: EventContext): EventProcessingAccount? { // Skip if account doesn't exist or is not active if (current == null || !current.isActive) return current
return current.copy(balance = current.balance + event.amount) }
fun closed(event: EventProcessingAccountClosed, current: EventProcessingAccount?, context: EventContext): EventProcessingAccount? { if (current == null) return null
return current.copy(isActive = false) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.util.UUID;
@EventTyperecord EventProcessingAccountOpened(UUID accountId) {}
@EventTyperecord EventProcessingDepositMade(double amount) {}
@EventTyperecord EventProcessingAccountClosed() {}
@ReadModelrecord EventProcessingAccount(UUID accountId, double balance, boolean isActive) { EventProcessingAccount() { this(new UUID(0, 0), 0.0, false); }}
@Reducerclass EventProcessingAccountReducer { EventProcessingAccount opened(EventProcessingAccountOpened event, EventProcessingAccount current, EventContext context) { return new EventProcessingAccount(event.accountId(), 0.0, true); }
EventProcessingAccount depositMade(EventProcessingDepositMade event, EventProcessingAccount current, EventContext context) { // Skip if account doesn't exist or is not active if (current == null || !current.isActive()) return current;
return new EventProcessingAccount(current.accountId(), current.balance() + event.amount(), current.isActive()); }
EventProcessingAccount closed(EventProcessingAccountClosed event, EventProcessingAccount current, EventContext context) { if (current == null) return null;
return new EventProcessingAccount(current.accountId(), current.balance(), false); }}defmodule EventProcessingAccountOpened do use Chronicle.Events.EventType, id: "event-processing-account-opened"
defstruct [:account_id]end
defmodule EventProcessingDepositMade do use Chronicle.Events.EventType, id: "event-processing-deposit-made"
defstruct [:amount]end
defmodule EventProcessingAccountClosed do use Chronicle.Events.EventType, id: "event-processing-account-closed"
defstruct []end
defmodule EventProcessingAccount do defstruct [:account_id, :balance, :is_active]end
defmodule EventProcessingAccountReducer do use Chronicle.Reducers.Reducer, model: EventProcessingAccount
alias EventProcessingAccountOpened alias EventProcessingDepositMade alias EventProcessingAccountClosed
@handles EventProcessingAccountOpened @handles EventProcessingDepositMade @handles EventProcessingAccountClosed
@impl true def reduce(%EventProcessingAccountOpened{} = event, _current, _context) do %EventProcessingAccount{account_id: event.account_id, balance: 0, is_active: true} end
# Skip if the account doesn't exist or is no longer active def reduce(%EventProcessingDepositMade{}, nil, _context), do: nil def reduce(%EventProcessingDepositMade{}, %{is_active: false} = current, _context), do: current
def reduce(%EventProcessingDepositMade{} = event, current, _context) do %{current | balance: current.balance + event.amount} end
def reduce(%EventProcessingAccountClosed{}, nil, _context), do: nil def reduce(%EventProcessingAccountClosed{}, current, _context), do: %{current | is_active: false}endimport { EventContext, eventType, Guid, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingAccountOpened { constructor(readonly accountId: Guid) {}}
@eventType()class EventProcessingDepositMade { constructor(readonly amount: number) {}}
@eventType()class EventProcessingAccountClosed {}
class EventProcessingAccount { accountId: Guid = Guid.empty; balance = 0; isActive = false;}
@reducer('', undefined, EventProcessingAccount)class EventProcessingAccountReducer { eventProcessingAccountOpened(event: EventProcessingAccountOpened, current: EventProcessingAccount | undefined): EventProcessingAccount { return { accountId: event.accountId, balance: 0, isActive: true }; }
eventProcessingDepositMade( event: EventProcessingDepositMade, current: EventProcessingAccount | undefined, context: EventContext ): EventProcessingAccount | undefined { // Skip if account doesn't exist or is not active if (!current || !current.isActive) return current;
return { ...current, balance: current.balance + event.amount }; }
eventProcessingAccountClosed( event: EventProcessingAccountClosed, current: EventProcessingAccount | undefined ): EventProcessingAccount | undefined { if (!current) return undefined;
return { ...current, isActive: false }; }}Error Handling
Section titled “Error Handling”Handling Invalid State
Section titled “Handling Invalid State”When there’s no state to build on yet, return null from a nullable-returning handler (TReadModel? Handle(...)) to leave the read model unset:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingSkipItemAdded(decimal Price);
public record EventProcessingSkipOrderSummary(decimal Total);
public class EventProcessingSkipOrderSummaryReducer : IReducerFor<EventProcessingSkipOrderSummary>{ public EventProcessingSkipOrderSummary? ItemAdded(EventProcessingSkipItemAdded @event, EventProcessingSkipOrderSummary? current, EventContext context) { // Can't add items if order doesn't exist if (current is null) return null;
return current with { Total = current.Total + @event.Price }; }}import io.cratis.chronicle.events.EventContextimport io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingSkipItemAdded(val price: Double)
@ReadModeldata class EventProcessingSkipOrderSummary(val total: Double = 0.0)
@Reducerclass EventProcessingSkipOrderSummaryReducer { fun itemAdded(event: EventProcessingSkipItemAdded, current: EventProcessingSkipOrderSummary?, context: EventContext): EventProcessingSkipOrderSummary? { // Can't add items if order doesn't exist if (current == null) return null
return current.copy(total = current.total + event.price) }}import io.cratis.chronicle.events.EventContext;import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
@EventTyperecord EventProcessingSkipItemAdded(double price) {}
@ReadModelrecord EventProcessingSkipOrderSummary(double total) { EventProcessingSkipOrderSummary() { this(0.0); }}
@Reducerclass EventProcessingSkipOrderSummaryReducer { EventProcessingSkipOrderSummary itemAdded(EventProcessingSkipItemAdded event, EventProcessingSkipOrderSummary current, EventContext context) { // Can't add items if order doesn't exist if (current == null) return null;
return new EventProcessingSkipOrderSummary(current.total() + event.price()); }}defmodule EventProcessingSkipItemAdded do use Chronicle.Events.EventType, id: "event-processing-skip-item-added"
defstruct [:price]end
defmodule EventProcessingSkipOrderSummary do defstruct total: 0end
defmodule EventProcessingSkipOrderSummaryReducer do use Chronicle.Reducers.Reducer, model: EventProcessingSkipOrderSummary
alias EventProcessingSkipItemAdded
@handles EventProcessingSkipItemAdded
# Can't add items if the order doesn't exist yet — returning nil when # current is already nil is a no-op. @impl true def reduce(%EventProcessingSkipItemAdded{}, nil, _context), do: nil
def reduce(%EventProcessingSkipItemAdded{} = event, current, _context) do %{current | total: current.total + event.price} endendimport { EventContext, eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingSkipItemAdded { constructor(readonly price: number) {}}
class EventProcessingSkipOrderSummary { total = 0;}
@reducer('', undefined, EventProcessingSkipOrderSummary)class EventProcessingSkipOrderSummaryReducer { eventProcessingSkipItemAdded( event: EventProcessingSkipItemAdded, current: EventProcessingSkipOrderSummary | undefined, context: EventContext ): EventProcessingSkipOrderSummary | undefined { // Can't add items if order doesn't exist if (!current) return undefined;
return { total: current.total + event.price }; }}Note the difference between skipping an event and removing a read model: returning null when current was already null is a no-op, but returning null when current holds real state deletes that read model. To ignore an event without discarding existing state, return current unchanged instead.
Recording Errors in State
Section titled “Recording Errors in State”Include error information in your read model:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingInvalidDataDetected(string Reason);
public record EventProcessingValidationResult(bool IsValid, List<string> Errors);
public class EventProcessingValidationResultReducer : IReducerFor<EventProcessingValidationResult>{ public EventProcessingValidationResult Detected(EventProcessingInvalidDataDetected @event, EventProcessingValidationResult? current) { var errors = new List<string>(current?.Errors ?? []) { @event.Reason };
return new EventProcessingValidationResult(false, errors); }}import io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingInvalidDataDetected(val reason: String)
@ReadModeldata class EventProcessingValidationResult(val isValid: Boolean = true, val errors: List<String> = emptyList())
@Reducerclass EventProcessingValidationResultReducer { fun detected(event: EventProcessingInvalidDataDetected, current: EventProcessingValidationResult?): EventProcessingValidationResult { val errors = (current?.errors ?: emptyList()) + event.reason
return EventProcessingValidationResult(false, errors) }}import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.util.ArrayList;import java.util.List;
@EventTyperecord EventProcessingInvalidDataDetected(String reason) {}
@ReadModelrecord EventProcessingValidationResult(boolean isValid, List<String> errors) { EventProcessingValidationResult() { this(true, List.of()); }}
@Reducerclass EventProcessingValidationResultReducer { EventProcessingValidationResult detected(EventProcessingInvalidDataDetected event, EventProcessingValidationResult current) { List<String> errors = new ArrayList<>(current == null ? List.of() : current.errors()); errors.add(event.reason());
return new EventProcessingValidationResult(false, errors); }}defmodule EventProcessingInvalidDataDetected do use Chronicle.Events.EventType, id: "event-processing-invalid-data-detected"
defstruct [:reason]end
defmodule EventProcessingValidationResult do defstruct is_valid: true, errors: []end
defmodule EventProcessingValidationResultReducer do use Chronicle.Reducers.Reducer, model: EventProcessingValidationResult
alias EventProcessingInvalidDataDetected
@handles EventProcessingInvalidDataDetected
@impl true def reduce(%EventProcessingInvalidDataDetected{} = event, current, _context) do errors = if current, do: current.errors, else: [] %EventProcessingValidationResult{is_valid: false, errors: errors ++ [event.reason]} endendimport { eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingInvalidDataDetected { constructor(readonly reason: string) {}}
class EventProcessingValidationResult { isValid = true; errors: string[] = [];}
@reducer('', undefined, EventProcessingValidationResult)class EventProcessingValidationResultReducer { eventProcessingInvalidDataDetected( event: EventProcessingInvalidDataDetected, current: EventProcessingValidationResult | undefined ): EventProcessingValidationResult { const errors = [...(current?.errors ?? []), event.reason];
return { isValid: false, errors }; }}Performance Optimization
Section titled “Performance Optimization”Minimize Object Creation
Section titled “Minimize Object Creation”Leverage record types with with expressions:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingMinimalMetricRecorded(decimal Value);
public record EventProcessingMinimalStats(int Count, decimal Sum);
public class EventProcessingMinimalStatsReducer : IReducerFor<EventProcessingMinimalStats>{ // Efficient - only creates a new object when needed public EventProcessingMinimalStats Recorded(EventProcessingMinimalMetricRecorded @event, EventProcessingMinimalStats? current) { if (current is null) return new EventProcessingMinimalStats(Count: 1, Sum: @event.Value);
return current with { Count = current.Count + 1, Sum = current.Sum + @event.Value }; }}import io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModel
@EventTypedata class EventProcessingMinimalMetricRecorded(val value: Double)
@ReadModeldata class EventProcessingMinimalStats(val count: Int = 0, val sum: Double = 0.0)
@Reducerclass EventProcessingMinimalStatsReducer { // Efficient - only creates a new object when needed fun recorded(event: EventProcessingMinimalMetricRecorded, current: EventProcessingMinimalStats?): EventProcessingMinimalStats { if (current == null) return EventProcessingMinimalStats(count = 1, sum = event.value)
return current.copy(count = current.count + 1, sum = current.sum + event.value) }}import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
@EventTyperecord EventProcessingMinimalMetricRecorded(double value) {}
@ReadModelrecord EventProcessingMinimalStats(int count, double sum) { EventProcessingMinimalStats() { this(0, 0.0); }}
@Reducerclass EventProcessingMinimalStatsReducer { // Efficient - only creates a new object when needed EventProcessingMinimalStats recorded(EventProcessingMinimalMetricRecorded event, EventProcessingMinimalStats current) { if (current == null) return new EventProcessingMinimalStats(1, event.value());
return new EventProcessingMinimalStats(current.count() + 1, current.sum() + event.value()); }}defmodule EventProcessingMinimalMetricRecorded do use Chronicle.Events.EventType, id: "event-processing-minimal-metric-recorded"
defstruct [:value]end
defmodule EventProcessingMinimalStats do defstruct count: 0, sum: 0end
defmodule EventProcessingMinimalStatsReducer do use Chronicle.Reducers.Reducer, model: EventProcessingMinimalStats
alias EventProcessingMinimalMetricRecorded
@handles EventProcessingMinimalMetricRecorded
@impl true def reduce(%EventProcessingMinimalMetricRecorded{} = event, nil, _context) do %EventProcessingMinimalStats{count: 1, sum: event.value} end
def reduce(%EventProcessingMinimalMetricRecorded{} = event, current, _context) do %{current | count: current.count + 1, sum: current.sum + event.value} endendimport { eventType, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingMinimalMetricRecorded { constructor(readonly value: number) {}}
class EventProcessingMinimalStats { count = 0; sum = 0;}
@reducer('', undefined, EventProcessingMinimalStats)class EventProcessingMinimalStatsReducer { // Efficient - only creates a new object when needed eventProcessingMinimalMetricRecorded( event: EventProcessingMinimalMetricRecorded, current: EventProcessingMinimalStats | undefined ): EventProcessingMinimalStats { if (!current) { return { count: 1, sum: event.value }; }
return { count: current.count + 1, sum: current.sum + event.value }; }}Reuse Collections
Section titled “Reuse Collections”Copy a collection before mutating it — a previously-returned state instance may still be held by a snapshot or another reader, so mutating it in place can corrupt state you don’t own:
using Cratis.Chronicle.Events;using Cratis.Chronicle.Reducers;
[EventType]public record EventProcessingReuseItemAdded(Guid ItemId, string Name);
public record EventProcessingItem(Guid ItemId, string Name);public record EventProcessingItemList(List<EventProcessingItem> Items);
public class EventProcessingItemListReducer : IReducerFor<EventProcessingItemList>{ public EventProcessingItemList ItemAdded(EventProcessingReuseItemAdded @event, EventProcessingItemList? current) { // Copy rather than mutate current.Items directly — a held snapshot may still reference it var items = new List<EventProcessingItem>(current?.Items ?? []) { new(@event.ItemId, @event.Name) };
return new EventProcessingItemList(items); }}import io.cratis.chronicle.events.EventTypeimport io.cratis.chronicle.observation.Reducerimport io.cratis.chronicle.readModels.ReadModelimport java.util.UUID
@EventTypedata class EventProcessingReuseItemAdded(val itemId: UUID, val name: String)
data class EventProcessingItem(val itemId: UUID, val name: String)
@ReadModeldata class EventProcessingItemList(val items: List<EventProcessingItem> = emptyList())
@Reducerclass EventProcessingItemListReducer { fun itemAdded(event: EventProcessingReuseItemAdded, current: EventProcessingItemList?): EventProcessingItemList { // Build a new list rather than mutate current.items directly - a held snapshot may still // reference it val items = (current?.items ?: emptyList()) + EventProcessingItem(event.itemId, event.name)
return EventProcessingItemList(items) }}import io.cratis.chronicle.events.EventType;import io.cratis.chronicle.observation.Reducer;import io.cratis.chronicle.readModels.ReadModel;
import java.util.ArrayList;import java.util.List;import java.util.UUID;
@EventTyperecord EventProcessingReuseItemAdded(UUID itemId, String name) {}
record EventProcessingItem(UUID itemId, String name) {}
@ReadModelrecord EventProcessingItemList(List<EventProcessingItem> items) { EventProcessingItemList() { this(List.of()); }}
@Reducerclass EventProcessingItemListReducer { EventProcessingItemList itemAdded(EventProcessingReuseItemAdded event, EventProcessingItemList current) { // Copy rather than mutate current.items() directly - a held snapshot may still reference it List<EventProcessingItem> items = new ArrayList<>(current == null ? List.of() : current.items()); items.add(new EventProcessingItem(event.itemId(), event.name()));
return new EventProcessingItemList(items); }}defmodule EventProcessingReuseItemAdded do use Chronicle.Events.EventType, id: "event-processing-reuse-item-added"
defstruct [:item_id, :name]end
defmodule EventProcessingItem do defstruct [:item_id, :name]end
defmodule EventProcessingItemList do defstruct items: []end
defmodule EventProcessingItemListReducer do use Chronicle.Reducers.Reducer, model: EventProcessingItemList
alias EventProcessingReuseItemAdded alias EventProcessingItem
@handles EventProcessingReuseItemAdded
# Elixir lists are immutable, so `++` always returns a new list rather than # mutating a shared one — there is no in-place mutation hazard here. @impl true def reduce(%EventProcessingReuseItemAdded{} = event, current, _context) do items = if current, do: current.items, else: [] item = %EventProcessingItem{item_id: event.item_id, name: event.name}
%EventProcessingItemList{items: items ++ [item]} endendimport { eventType, Guid, reducer } from '@cratis/chronicle';
@eventType()class EventProcessingReuseItemAdded { constructor(readonly itemId: Guid, readonly name: string) {}}
class EventProcessingItem { itemId: Guid = Guid.empty; name = '';}
class EventProcessingItemList { items: EventProcessingItem[] = [];}
@reducer('', undefined, EventProcessingItemList)class EventProcessingItemListReducer { eventProcessingReuseItemAdded( event: EventProcessingReuseItemAdded, current: EventProcessingItemList | undefined ): EventProcessingItemList { // Copy rather than mutate current.items directly — a held snapshot may still reference it const items = [...(current?.items ?? []), { itemId: event.itemId, name: event.name }];
return { items }; }}Best Practices
Section titled “Best Practices”- Use record types - Prefer immutable record types for read models
- Keep logic pure - Avoid side effects; only compute state from events
- Handle null safely - Always check
currentfor null on first event - Use with expressions - Leverage record’s
withfor clean state updates - Return null to skip or remove - Declare the handler’s return type as
TReadModel?and returnnull— but only when you mean to clear the read model, not merely to ignore the event (returncurrentunchanged for that) - Access context when needed - Use
EventContextfor metadata like timestamps - Name methods clearly - Use descriptive method names that match event types
- Test thoroughly - Unit test with various event sequences and edge cases