Catch up projections

A month-end report can require every change already saved across the journal, including changes to other accounts. Call CatchUpAsync after the aggregate reply, then query the transactional projection's read model:

Shared setup
module Account =
    open FCQRS.Common

    // What a caller can ask an account to do.
    type AccountCommand =
        | Open of owner: string
        | Deposit of amount: decimal
        | Withdraw of amount: decimal

    // What the account replies. Rejected is a reply only: it is never stored.
    type AccountEvent =
        | Opened of owner: string
        | Deposited of amount: decimal
        | Withdrawn of amount: decimal
        | Rejected of reason: string

    // What the account knows now, rebuilt from its events.
    type AccountState = { Owner: string option; Balance: decimal }

    // The state before the account's first event.
    let initial = { Owner = None; Balance = 0m }

    // Chooses what to do with a command, based on the current state.
    let decide (command: Command<AccountCommand>) (state: AccountState) =
        match command.CommandDetails, state.Owner with
        | Open _, Some _ -> DeferEvent(Rejected "The account is already open")
        | Open owner, None -> PersistEvent(Opened owner)
        | _, None -> DeferEvent(Rejected "The account is not open")
        | (Deposit amount | Withdraw amount), _ when amount <= 0m ->
            DeferEvent(Rejected "The amount must be positive")
        | Deposit amount, _ -> PersistEvent(Deposited amount)
        | Withdraw amount, _ when amount > state.Balance ->
            DeferEvent(Rejected $"Insufficient funds: {state.Balance} available")
        | Withdraw amount, _ -> PersistEvent(Withdrawn amount)

    // Applies one event to the state. A rejection changes nothing.
    let fold (event: Event<AccountEvent>) (state: AccountState) =
        match event.EventDetails with
        | Opened owner -> { state with Owner = Some owner }
        | Deposited amount -> { state with Balance = state.Balance + amount }
        | Withdrawn amount -> { state with Balance = state.Balance - amount }
        | Rejected _ -> state
open System
open FCQRS.Model.Data
open FCQRS.Common
open FCQRS.FSharp
open Account
open System.Threading
open FCQRS.Projections
let depositAndCatchUp (accounts: AggregateHandle<AccountCommand, AccountEvent>)
                      (projection: IProjection) (cid: CID) (alice: AggregateId)
                      (cancellationToken: CancellationToken) = async {
let! reply = accounts.Send cid alice (Deposit 100m) (fun _ -> true)
do! projection.CatchUpAsync(cancellationToken) |> Async.AwaitTask
// Inspect the reply, then query the balances of every account.
var reply = await accounts(_ => true, cid, alice, new Deposit(100m));
await projection.CatchUpAsync(cancellationToken);
// Inspect the reply, then query the balances of every account.

These transactional projection APIs are available in FCQRS 6.4.0 and later.

CatchUpAsync captures one journal snapshot: the highest committed sequence number for each persistence identity visible in one database read. A persistence identity identifies one actor's journal history. The call succeeds after this projection has committed every event through those captured sequence numbers. The snapshot includes aggregate events and saga journal records. It excludes Akka's cluster-sharding records, whose persistence identities start with /sharding/. Akka deletes their early history itself, and they are not application events.

The order matters. Await the aggregate reply before calling CatchUpAsync so the snapshot includes that event when the reply was journaled. It also includes every other event committed before the snapshot, across all application persistence identities in this journal. Writes arriving after the snapshot do not extend this call's target.

To stay cheap on a large journal, most snapshots read only writes numbered after what the projection handled LateWriteWindow earlier, 30 seconds by default. A write can take its number some time before it commits. One that took longer than the window is handled by the next full scan, within FullScanInterval, and a CatchUpAsync call in between can return before it. No event is skipped; only that call's promise narrows for such a write.

The checkpoint is durable, so this wait does not require a correlation subscription before sending. It also works after a deferred reply, although a deferred reply itself adds no journal event. Check the command outcome separately: projection completion does not turn a rejected command into a successful one.

The examples use the account from the tutorial's withdraw money step. The tutorial's statement is a transactional projection too.

Register a transactional projection

Use Fcqrs.transactionalProjection in F# or AddTransactionalProjection in C#. These are separate from the registration in Add a projection, whose handler writes outside the transaction. The transactional runner stores one checkpoint per persistence identity under a stable projection name.

The application supplies a SQL connection factory and a handler. FCQRS opens a transaction, passes its connection and transaction to the handler, records progress, and commits them together. Return from the handler only when its writes have completed. Use the supplied transaction for every read-model change covered by this projection.

Use the Balances table from Add a projection. The following SQLite registration uses one database for the journal, read model, and FCQRS-owned progress tables:

Shared setup
}
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open Microsoft.Data.Sqlite
open FCQRS.ProjectionStorage
let connectionString = "Data Source=app.db;"
let configuration = Microsoft.Extensions.Configuration.ConfigurationBuilder().Build()
let loggerFactory = Microsoft.Extensions.Logging.LoggerFactory.Create(fun _ -> ())
module Sqlite =
    let api = Fcqrs.actor configuration loggerFactory
                  (Some(Fcqrs.connect FCQRS.Actor.DBType.Sqlite connectionString)) "accounts"
// NuGet: Dapper, Microsoft.Data.Sqlite
open System
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open Microsoft.Data.Sqlite
open FCQRS.Common
open FCQRS.FSharp
open FCQRS.ProjectionStorage
open FCQRS.Projections
open Account

let openAccount =
    "insert into Balances (Id, Owner, Balance) values (@Id, @Owner, 0)"
let changeBalance =
    "update Balances set Balance = Balance + @Amount where Id = @Id"

let handle (connection: DbConnection) (transaction: DbTransaction)
           (envelope: EventEnvelope) : Task =
    task {
        match envelope.Event with
        // Sender is the ID of the account that stored the event.
        | :? Event<AccountEvent> as stored ->
            let id = string stored.Sender.Value
            let add (amount: decimal) =
                task {
                    let row = {| Id = id; Amount = amount |}
                    let! rows =
                        connection.ExecuteAsync(changeBalance, row, transaction)
                    if rows <> 1 then
                        failwith "A balance change needs an earlier Opened"
                }
            match stored.EventDetails with
            | Opened owner ->
                let row = {| Id = id; Owner = owner |}
                let! _ = connection.ExecuteAsync(openAccount, row, transaction)
                ()
            | Deposited amount -> do! add amount
            | Withdrawn amount -> do! add -amount
            | Rejected _ -> ()
        | _ -> ()
    } :> Task

let store =
    SqlProjectionStore(
        ProjectionSqlDialect.Sqlite,
        Func<DbConnection>(fun () -> new SqliteConnection(connectionString)))
let options = TransactionalProjectionOptions("Balances", store)
let projection = Fcqrs.transactionalProjection api options handle
// NuGet: Dapper, Microsoft.Data.Sqlite
using System.Data.Common;
using Akka.Persistence.Query;
using Dapper;
using Microsoft.Data.Sqlite;
using FCQRS;
using static FCQRS.Common;
using static FCQRS.ProjectionStorage;
using static FCQRS.Projections;

const string OpenAccount =
    "insert into Balances (Id, Owner, Balance) values (@Id, @Owner, 0)";
const string ChangeBalance =
    "update Balances set Balance = Balance + @Amount where Id = @Id";

static async Task Handle(
    DbConnection connection, DbTransaction transaction, EventEnvelope envelope)
{
    // Sender is the ID of the account that stored the event.
    if (envelope.Event is not Event<AccountEvent> { Sender: { } sender } stored)
        return;
    var id = sender.Value.ToString();

    async Task Add(decimal amount)
    {
        var row = new { Id = id, Amount = amount };
        var rows =
            await connection.ExecuteAsync(ChangeBalance, row, transaction);
        if (rows != 1)
            throw new InvalidOperationException(
                "A balance change needs an earlier Opened");
    }

    switch (stored.EventDetails)
    {
        case Opened opened:
            var owner = new { Id = id, opened.Owner };
            await connection.ExecuteAsync(OpenAccount, owner, transaction);
            break;
        case Deposited deposited: await Add(deposited.Amount); break;
        case Withdrawn withdrawn: await Add(-withdrawn.Amount); break;
    }
}

var store = new SqlProjectionStore(
    ProjectionSqlDialect.Sqlite, () => new SqliteConnection(connectionString));
var options = new TransactionalProjectionOptions("Balances", store);

builder.Services.AddFcqrs(connectionString, "accounts")
    .AddAggregate<Account>()
    .AddTransactionalProjection(options, Handle);

The C# registration adds IProjection to dependency injection. Inject it into the caller that waits for catch-up. The F# facade returns the same interface. It also implements ISubscribe, so existing correlation subscriptions remain available, with aggregate notifications published after commit.

The runner applies each read of the journal in journal order, but an event whose write commits late is applied in a later read, so a saga's follow-up event can commit before the event that caused it. Both carry the same correlation ID. A subscription for that correlation ID, such as Subscribe(cid, ...) or sendAwaiting, therefore receives its notifications after the whole snapshot commits. A waiter woken by the follow-up can read the event that caused it. It can also wait longer than the commit of its own event, and it is cancelled if the projection fails before the snapshot commits. A subscription without a correlation ID, including a filter on the correlation ID, receives each notification when its event commits.

Subscriptions belong to the local worker. If several workers share the same projection name and store, a worker can observe progress committed by another worker without publishing that other worker's notifications. Use CatchUpAsync for the durable completion boundary across those workers.

The factory must return a new, unopened connection. For separate journal and read-model databases, use the constructor taking journalConnectionFactory and projectionConnectionFactory. Both databases must use the selected dialect. The journal factory must read the authoritative database; the projection factory must open the store updated by the handler. Journal schema, table, and persistence-ID and sequence-number column overrides must match the Akka journal configuration.

The runner validates the effective HOCON settings for both SQL journal readers and writers. DataOptionsSetup and MultiDataOptionsSetup overrides are unsupported because they can replace those settings. If akka.persistence.query.journal.sql.write-plugin is set, it must identify the active write journal. Configured Akka event adapters are also unsupported by this reader.

For PostgreSQL, add the Npgsql package to the application, choose ProjectionSqlDialect.PostgreSql, and configure the Akka journal for the same database. The handler above uses SQL accepted by both SQLite and PostgreSQL:

Shared setup
module PostgreSql =
    let handle = Sqlite.handle
open Npgsql

let store =
    SqlProjectionStore(
        ProjectionSqlDialect.PostgreSql,
        Func<DbConnection>(fun () -> new NpgsqlConnection(connectionString)))
let options = TransactionalProjectionOptions("Balances", store)

let connection = Fcqrs.connect FCQRS.Actor.DBType.PostgreSQL15 connectionString
let api = Fcqrs.actor configuration loggerFactory (Some connection) "accounts"
let projection = Fcqrs.transactionalProjection api options handle
var store = new SqlProjectionStore(
    ProjectionSqlDialect.PostgreSql,
    () => new Npgsql.NpgsqlConnection(connectionString));
var options = new TransactionalProjectionOptions("Balances", store);

builder.Services.AddFcqrs(
        connectionString, "accounts", FCQRS.Actor.DBType.PostgreSQL15)
    .AddAggregate<Account>()
    .AddTransactionalProjection(options, Handle);

The snapshot and completion contract uses per-identity checkpoints on both databases.

FCQRS owns the transaction and progress tables. Keep the handler's connection and transaction inside the handler, and let FCQRS commit or roll back. Dispose the F# projection handle when its runtime scope ends; the C# host manages the registered projection's lifetime.

Understand completion and ordering

Transactional projections handle events in sequence within each persistence identity. They do not promise a global processing order across identities. A handler combining facts from different aggregates must tolerate their arrival order or explicitly coordinate those dependencies.

The runner requires a complete history through each captured target. It does not treat the largest observed sequence number as proof that missing earlier events were handled. Preserve journal events needed by this projection, and begin with a new read model when starting a new checkpoint history.

The completion boundary covers one projection and the database transaction used by its handler. It does not wait for another projection, an external HTTP call, or work started without awaiting it. A query can observe later committed changes too; completion does not freeze the read model at the captured snapshot.

Catch-up suppresses an ambient TransactionScope so the snapshot sees freshly committed journal data. Projection commits belong to FCQRS transactions independently of the caller's transaction; rolling back the caller's scope does not roll back projection work.

For an account-closing workflow, this wait can establish that the selected projection processed every event committed before the closing event and the subsequent snapshot. Rejecting later commands to the closed account remains an aggregate or application rule. The wait does not drain commands that were sent earlier but have not yet been persisted, and it does not stop future writes.

Configure discovery and waiting

TransactionalProjectionOptions configures the runner:

Property Default Meaning
PollInterval 1 second Background delay before discovering new journal heads after the previous batch finishes
BatchSize 500 Maximum events fetched for one persistence identity per query
CatchUpTimeout 30 seconds Time allowed for the entire catch-up call, including snapshot capture
LateWriteWindow 30 seconds How long after taking its journal number a write may commit and still be found by a query of recent writes
FullScanInterval 5 minutes How often a query reads every history's head, which also finds a write that committed after the window

Background discovery follows new events as the application runs. Each CatchUpAsync call captures its own fixed target once; it does not repeatedly replace the target with a newer journal tail. Most discovery queries read only the journal rows numbered after what the projection handled a LateWriteWindow earlier, so their cost follows recent writes. A full scan, which groups the whole journal by persistence_id, runs when the projection starts, once more a window later, and every FullScanInterval. It costs more as the journal grows. An aggregate on this node that stores an event starts the next query at once.

The catch-up timeout belongs to TransactionalProjectionOptions, separately from akka.fcqrs.command-timeout used by Read your writes.

Recover without losing progress

Read-model updates and their checkpoint share one transaction. If the handler fails before commit, neither becomes durable. After restart, use the same projection name and store to resume from the committed checkpoints. The handler can run again for an event whose transaction did not commit; external side effects therefore need their own retry and idempotency policy.

An error while processing journal events stops the runner and faults IProjection.Completion. Catch-up calls then report the error. Observe this task in the application's worker or health monitoring, correct the cause, and restart the projection. Missing sequence numbers stop the runner instead of silently advancing its checkpoint.

Cancellation or TimeoutException ends the caller's wait. It does not undo the aggregate command or read-model transactions that have already committed. Treat the result as an incomplete confirmation, then retry catch-up or report that the read model has not yet been confirmed current.

The repository's facade tests cover catch-up and recovery:

dotnet run --project test/Facade.Tests/Facade.Tests.fsproj

Set FCQRS_TEST_POSTGRES to a PostgreSQL test-server connection string to include the PostgreSQL integration cases. The test account needs permission to create and drop the tests' isolated databases. CI runs these cases against its PostgreSQL service. Test your domain covers tests for the application's own decisions and replay rules.

For rebuilding query data, continue with Rebuild a read model. For a request that needs one matching notification, use Read your writes.