Skip to content

Replay notifications

A replay re-delivers history to an observer, and some reactors and reducers need to know when that starts and when it is over. A reactor that feeds a cache or an external index might pause its writes for the duration; a reducer that keeps in-memory state might reset it.

Implement ICanBeNotifiedWhenReplay on the reactor or reducer class and Chronicle calls BeginReplay() when a replay of the observer begins and EndReplay() when it ends. The interface lives in the Cratis.Chronicle.Observation namespace.

using Cratis.Chronicle.Events;
using Cratis.Chronicle.Observation;
using Cratis.Chronicle.Reactors;
[EventType]
public record ReplayNotifiedOrderPlaced(string OrderId);
public interface IReplayNotificationGate
{
void Pause();
void Resume();
}
public class ReplayNotifiedOrderReactor(IReplayNotificationGate gate) : IReactor, ICanBeNotifiedWhenReplay
{
public Task BeginReplay()
{
gate.Pause();
return Task.CompletedTask;
}
public Task EndReplay()
{
gate.Resume();
return Task.CompletedTask;
}
public Task OrderPlaced(ReplayNotifiedOrderPlaced @event)
{
// Handles the event both live and during the replay.
return Task.CompletedTask;
}
}

View C# snippet source on GitHub

Reducers use the same interface:

using Cratis.Chronicle.Events;
using Cratis.Chronicle.Observation;
using Cratis.Chronicle.Reducers;
[EventType]
public record ReplayNotifiedItemAdded(string OrderId);
public record ReplayNotifiedOrderTotals(int Items);
public class ReplayNotifiedOrderTotalsReducer : IReducerFor<ReplayNotifiedOrderTotals>, ICanBeNotifiedWhenReplay
{
public Task BeginReplay() => Task.CompletedTask;
public Task EndReplay() => Task.CompletedTask;
public ReplayNotifiedOrderTotals ItemAdded(ReplayNotifiedItemAdded @event, ReplayNotifiedOrderTotals? current) =>
new((current?.Items ?? 0) + 1);
}

View C# snippet source on GitHub

The notifications are calls of their own. They are not events and they do not go through the event handlers, so a [Replay] handler does not see them and an event handler is not called for them. The events of the replay are delivered between the two calls, as they always are.

Chronicle creates the reactor or reducer for each notification, from a new dependency injection scope, and disposes it afterwards. BeginReplay() and EndReplay() therefore run on different instances, and neither is the instance that handles the events. Anything they need to share must live outside the class, in a service registered as a singleton, such as the one injected in the reactor example above. A scoped or transient service is created anew for each notification, so it could not carry state from BeginReplay() to EndReplay().

Replaying one partition, for example with the observers replay-partition command of the Cratis CLI, has its own interface, ICanBeNotifiedWhenPartitionReplayed. It receives the partition being replayed:

using Cratis.Chronicle.Events;
using Cratis.Chronicle.Observation;
using Cratis.Chronicle.Reactors;
[EventType]
public record PartitionNotifiedOrderPlaced(string OrderId);
public class PartitionNotifiedOrderReactor : IReactor, ICanBeNotifiedWhenPartitionReplayed
{
public Task BeginReplayPartition(Partition partition)
{
// The partition is the event source id being replayed.
return Task.CompletedTask;
}
public Task EndReplayPartition(Partition partition)
{
return Task.CompletedTask;
}
public Task OrderPlaced(PartitionNotifiedOrderPlaced @event)
{
return Task.CompletedTask;
}
}

View C# snippet source on GitHub

The two interfaces are independent. Implement either or both, depending on which replays you need to know about.

  • A reactor that has no class to implement the interface on: one defined fluently with the reactor builder, or registered at runtime with a handler delegate. Those are never notified.
  • Other clients. These interfaces are part of the .NET client.