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 }; }}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);}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);}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);}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 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);}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 }; }}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); }}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);}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); }}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); }}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 }; }}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 }; }}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); }}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 }; }}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); }}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