Step 4: Show a statement
Alice wants a statement: every deposit and withdrawal, with the balance after it. The account cannot answer that. It decides commands and keeps only what its rules need, which is the owner and the balance. This step adds the other half of the split from the overview: a read model, a table shaped for one screen, which a projection keeps up to date from the journal.
The same feature in a CRUD application
A CRUD application writes the statement row in the same transaction as the balance, and the statement query reads the table that the write code changed:
BEGIN;
UPDATE accounts SET balance = balance - 30 WHERE id = 'alice';
INSERT INTO transactions (account, entry, amount)
VALUES ('alice', 'Withdrawal', -30);
COMMIT;
-- The statement reads the rows the write code inserted.
SELECT entry, amount FROM transactions WHERE account = 'alice';
Every code path that changes a balance must also insert the transaction row. In FCQRS, the account stores only its event. One projection writes the statement from the journal, and because the journal holds every event, the statement can be rebuilt from it.
Create the 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
module Statement =
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open FCQRS.Common
open Account
Fcqrs_250-tutorial_004-show-a-statement.md_page.AccountFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
Fcqrs_250-tutorial_004-show-a-statement.md_page.Account.AccountCommandOpenowner: stringstringAn abbreviation for the CLI type . Basic Types
Depositamount: decimaldecimalAn abbreviation for the CLI type . Basic Types
WithdrawFcqrs_250-tutorial_004-show-a-statement.md_page.Account.AccountEventOpenedDepositedWithdrawnRejectedreason: stringFcqrs_250-tutorial_004-show-a-statement.md_page.Account.AccountStateOwner: string optionoptionThe type of optional values. When used from other CLI languages the empty option is the null value. Use the constructors Some and None to create values of this type. Use the values in the Option module to manipulate values of this type, or pattern match against the values directly. 'None' values will appear as the value null to other CLI languages. Instance methods on this type will appear as static methods to other CLI languages due to the use of null as a value representation. Options
Balance: decimalinitial: AccountStateNoneThe representation of "No value"
decide: Command<AccountCommand> -> AccountState -> EventAction<AccountEvent>command: Command<AccountCommand>FCQRS.Common.Command`1Represents a command to be processed by an aggregate actor. <typeparam name="'CommandDetails">The specific type of the command payload.</typeparam>
state: AccountStateCommandDetails: 'CommandDetailsThe specific details or payload of the command.
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
DeferEventPublish and fold the event in the live actor without storing it or incrementing the persisted version. A deferred event does not start a saga; running sagas still receive it.
PersistEventPersist the event to the journal. The actor's state will be updated using the event handler *after* persistence succeeds.
(<=): 'T -> 'T -> boolStructural less-than-or-equal comparison The first parameter. The second parameter. The result of the comparison. 5 <= 1 // Evaluates to false 5 <= 5 // Evaluates to true [1; 5] <= [1; 6] // Evaluates to true
(>): 'T -> 'T -> boolStructural greater-than The first parameter. The second parameter. The result of the comparison. 5 > 1 // Evaluates to true 5 > 5 // Evaluates to false (1, "a") > (1, "z") // Evaluates to false
fold: Event<AccountEvent> -> AccountState -> AccountStateevent: Event<AccountEvent>FCQRS.Common.Event`1Represents an event generated by an aggregate actor as a result of processing a command. <typeparam name="'EventDetails">The specific type of the event payload.</typeparam>
EventDetails: 'EventDetailsThe specific details or payload of the event.
Fcqrs_250-tutorial_004-show-a-statement.md_page.StatementSystemDataCommonThreadingTasksAkkaPersistenceQueryDapper// The read model: one row per stored event, with the balance after it.
let createTable (connection: DbConnection) =
connection.Execute
"CREATE TABLE IF NOT EXISTS statement (
account TEXT NOT NULL,
version INTEGER NOT NULL,
entry TEXT NOT NULL,
amount NUMERIC NOT NULL,
balance NUMERIC NOT NULL,
PRIMARY KEY (account, version))"
|> ignore
createTable: DbConnection -> unitconnection: DbConnectionSystem.Data.Common.DbConnectionDefines the core behavior of database connections and provides a base class for database-specific connections.
Execute: string * obj * System.Data.IDbTransaction * System.Nullable<int> * System.Nullable<System.Data.CommandType> -> intExecute parameterized SQL. The connection to query on. The SQL to execute for this query. The parameters to use for this query. The transaction to use for this query. Number of seconds before command execution timeout. Is it a stored proc or a batch? The number of rows affected.
(|>): 'T1 -> ('T1 -> 'U) -> 'UApply a function to a value, the value being on the left, the function on the right The argument. The function. The function result. let doubleIt x = x * 2 3 |> doubleIt // Evaluates to 6
ignore: 'T -> unitIgnore the passed value. This is often used to throw away results of a computation. The value to ignore. ignore 55555 // Evaluates to ()
// The read model: one row per stored event, with the balance after it.
public static void CreateTable(DbConnection connection) =>
connection.Execute(
"""
CREATE TABLE IF NOT EXISTS statement (
account TEXT NOT NULL,
version INTEGER NOT NULL,
entry TEXT NOT NULL,
amount NUMERIC NOT NULL,
balance NUMERIC NOT NULL,
PRIMARY KEY (account, version))
""");
The statement is an ordinary table. version is the version of the event that produced the row, so
each event adds at most one row.
Write the projection
// Adds a row whose balance continues from the account's previous row.
let private addRow =
"INSERT INTO statement (account, version, entry, amount, balance)
SELECT @Account, @Version, @Entry, @Amount,
COALESCE((SELECT balance FROM statement WHERE account = @Account
ORDER BY version DESC LIMIT 1), 0) + @Amount"
// FCQRS calls this for each stored event, inside a transaction it commits.
let handle (connection: DbConnection) (transaction: DbTransaction)
(envelope: EventEnvelope) =
task {
match envelope.Event with
// Only account events go on a statement; Sender is the account's ID.
| :? Event<AccountEvent> as stored ->
let add (entry: string) (amount: decimal) =
let row =
{| Account = string stored.Sender.Value
Version = envelope.SequenceNr
Entry = entry
Amount = amount |}
connection.ExecuteAsync(addRow, row, transaction) :> Task
match stored.EventDetails with
| Opened owner -> do! add $"Opened for {owner}" 0m
| Deposited amount -> do! add "Deposit" amount
| Withdrawn amount -> do! add "Withdrawal" -amount
// A rejection is a reply only; the journal never holds one.
| Rejected _ -> ()
| _ -> ()
}
:> Task
addRow: stringhandle: DbConnection -> DbTransaction -> EventEnvelope -> Taskconnection: DbConnectionSystem.Data.Common.DbConnectionDefines the core behavior of database connections and provides a base class for database-specific connections.
transaction: DbTransactionSystem.Data.Common.DbTransactionDefines the core behavior of database transactions and provides a base class for database-specific transactions.
envelope: EventEnvelopeAkka.Persistence.Query.EventEnvelopeEvent wrapper adding meta data for the events in the result stream of query, or similar queries. The is the time the event was stored, in ticks. The value of this property represents the number of 100-nanosecond intervals that have elapsed since 12:00:00 midnight, January 1, 0001 in the Gregorian calendar (same as `DateTime.Now.Ticks`).
task: TaskBuilderBuilds a task using computation expression syntax.
Event: objFCQRS.Common.Event`1Represents an event generated by an aggregate actor as a result of processing a command. <typeparam name="'EventDetails">The specific type of the event payload.</typeparam>
Fcqrs_250-tutorial_004-show-a-statement.md_page.Account.AccountEventstored: Event<AccountEvent>add: string -> decimal -> Taskentry: stringstringAn abbreviation for the CLI type . Basic Types
amount: decimaldecimalAn abbreviation for the CLI type . Basic Types
row: {| Account: string; Amount: decimal; Entry: string; Version: int64 |}Account: stringstring: 'T -> stringConverts the argument to a string using ToString. For standard integer and floating point values and any type that implements IFormattable, ToString conversion uses CultureInfo.InvariantCulture. The input value. The converted string. string 'A' // evaluates to "A" string 0xff // evaluates to "255" string -10 // evaluates to "-10"
Value: FCQRS.Model.Data.AggregateIdGet the value of a 'Some' option. A NullReferenceException is raised if the option is 'None'.
Sender: FCQRS.Model.Data.AggregateId optionAn optional identifier for the actor that generated the event.
Version: int64SequenceNr: int64Entry: stringAmount: decimalExecuteAsync: string * obj * System.Data.IDbTransaction * System.Nullable<int> * System.Nullable<System.Data.CommandType> -> Task<int>Execute a command asynchronously using Task. The connection to query on. The SQL to execute for this query. The parameters to use for this query. The transaction to use for this query. Number of seconds before command execution timeout. Is it a stored proc or a batch? The number of rows affected.
System.Threading.Tasks.TaskRepresents an asynchronous operation.
EventDetails: 'EventDetailsThe specific details or payload of the event.
Openedowner: stringDepositedWithdrawn(~-): ^T -> ^TOverloaded unary negation. The value to negate. The result of the operation.
Rejected// Adds a row whose balance continues from the account's previous row.
const string AddRow =
"""
INSERT INTO statement (account, version, entry, amount, balance)
SELECT @Account, @Version, @Entry, @Amount,
COALESCE((SELECT balance FROM statement WHERE account = @Account
ORDER BY version DESC LIMIT 1), 0) + @Amount
""";
// FCQRS calls this for each stored event, inside a transaction it commits.
public static async Task Handle(
DbConnection connection, DbTransaction transaction, EventEnvelope envelope)
{
// Only account events go on a statement; Sender is the account's ID.
if (envelope.Event is not Event<AccountEvent> { Sender: { } sender } stored)
return;
Task Add(string entry, decimal amount)
{
var row = new
{
Account = sender.Value.ToString(),
Version = envelope.SequenceNr,
Entry = entry,
Amount = amount
};
return connection.ExecuteAsync(AddRow, row, transaction);
}
await (stored.EventDetails switch
{
Opened opened => Add($"Opened for {opened.Owner}", 0m),
Deposited deposited => Add("Deposit", deposited.Amount),
Withdrawn withdrawn => Add("Withdrawal", -withdrawn.Amount),
// A rejection is a reply only; the journal never holds one.
Rejected => Task.CompletedTask
});
}
FCQRS calls the handler for each stored event and passes it a connection and a transaction. It commits the handler's writes together with the projection's progress. If the handler fails, neither is committed, and the event is handled again when the projection restarts.
The handler receives events from every aggregate, so it keeps only account events. envelope.Event
is the stored event, Sender is the ID of the account that stored it, and envelope.SequenceNr is
its version. Rejected never arrives, because the journal never holds it.
Each row's balance continues from the account's previous row. That is correct because a projection receives one account's events in version order. Across accounts it follows the journal's numbers, which on PostgreSQL can differ from the order writes commit, so a handler that combines several accounts must not depend on which account's event arrives first.
Register the projection
Shared setup
open System
open System.Data.Common
open System.IO
open Microsoft.Data.Sqlite
open Microsoft.Extensions.Configuration
open Microsoft.Extensions.Logging
open FCQRS.Actor
open FCQRS.Common
open FCQRS.FSharp
open FCQRS.ProjectionStorage
open FCQRS.Projections
open Account
let database = Path.Combine(AppContext.BaseDirectory, "accounts.db")
let connectionString = $"Data Source={database}"
// Startup and the account registration are the same as in step 2.
let logging = LoggerFactory.Create(fun _ -> ())
let configuration = ConfigurationBuilder().Build()
let connection = Fcqrs.connect DBType.Sqlite connectionString
let api = Fcqrs.actor configuration logging (Some connection) "accounts"
let accounts =
Fcqrs.aggregate api
{ Name = "Account"
Initial = initial
Decide = decide
Fold = fold
Snapshots = Default
Passivation = PassivationPolicy.Default }
Fcqrs.wireSagaStarters api []
SystemDataCommonIOMicrosoftSqliteExtensionsConfigurationLoggingFCQRSFCQRS.ActorFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
FCQRS.FSharpIdiomatic-F# functional facade for FCQRS. Gives F# consumers the same one-call ergonomics the C# host-builder (HostExtensions.fs) gives C#, but with F# idioms: records-of-functions for the definitions, typed handles for the results, an explicit wiring pipeline, and plain helpers for saga side effects. It is a *pure addition* that wraps only the existing primitives (IActor.InitializeActor / SagaBuilder.initSimple / Projections.startTracked / InitializeSagaStarter / CreateCommandSubscription / Actor.api) and changes nothing in the C# interop layer or the core. open FCQRS.FSharp let api = Fcqrs.actor config loggerFactory (Some (Fcqrs.connect DBType.Sqlite conn)) "Cluster" let documents = Fcqrs.aggregate api { Name="Document"; Initial=...; Decide=...; Fold=... } let slugs = Fcqrs.aggregate api { Name="Slug"; Initial=...; Decide=...; Fold=... } let publication = Fcqrs.saga api (publicationDef documents.Factory slugs.Factory) Fcqrs.wireSagaStarters api [ publication ] let subs = Fcqrs.projection api (Projection.single 0 updateReadModel) // (Projection.multi when you must control which notifications publish) // send a command and await the matching aggregate reply: let! ev = documents.Send (Fcqrs.newCid()) (Fcqrs.aggregateId id) cmd (fun e -> ...)
FCQRS.ProjectionStorageSQL storage for journal-wide, transactional projection catch-up.
FCQRS.ProjectionsTransactional projections with a journal-wide, durable catch-up boundary.
Fcqrs_250-tutorial_004-show-a-statement.md_page.Accountdatabase: stringSystem.IO.PathPerforms operations on instances that contain file or directory path information. These operations are performed in a cross-platform manner.
Combine: string * string -> stringCombines two strings into a path. The first path to combine. The second path to combine. .NET Framework and .NET Core versions older than 2.1: or contains one or more of the invalid characters defined in . or is . The combined paths. If one of the specified paths is a zero-length string, this method returns the other path. If contains an absolute path, this method returns .
System.AppContextProvides members for setting and retrieving data about an application's context.
BaseDirectory: stringGets the file path of the base directory that the assembly resolver uses to probe for assemblies. The file path of the base directory that the assembly resolver uses to probe for assemblies.
connectionString: stringlogging: ILoggerFactoryMicrosoft.Extensions.Logging.LoggerFactoryProduces instances of classes based on the given providers.
Create: Action<ILoggingBuilder> -> ILoggerFactoryCreates new instance of configured using provided delegate. A delegate to configure the . The that was created.
configuration: IConfigurationRoot``.ctor``: unit -> unitBuild: unit -> IConfigurationRootBuilds an with keys and values from the set of providers registered in . An with keys and values from the registered providers.
connection: ConnectionFCQRS.FSharp.Fcqrsconnect: DBType -> string -> ConnectionBuild a SQLite/etc. Connection from a raw connection string of any length.
FCQRS.Actor.DBTypeRepresents the type of database connection
SqliteSQLite using Microsoft.Data.Sqlite provider
api: IActoractor: IConfiguration -> ILoggerFactory -> Connection option -> string -> IActorCreate the actor system from plain values (cluster name as a string).
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
accounts: AggregateHandle<AccountCommand,AccountEvent>aggregate: IActor -> Aggregate<'State,'Command,'Event> -> AggregateHandle<'Command,'Event>Register an aggregate and return its typed handle. Calling this IS the registration (it initializes the sharding region).
Name: stringInitial: 'Stateinitial: AccountStateDecide: Command<'Command> -> 'State -> EventAction<'Event>handleCommand (decide): command + current state -> what to do.
decide: Command<AccountCommand> -> AccountState -> EventAction<AccountEvent>Fold: Event<'Event> -> 'State -> 'StateapplyEvent (fold): event + current state -> next state (pure).
fold: Event<AccountEvent> -> AccountState -> AccountStateSnapshots: SnapshotPolicySnapshot cadence: Default (config / 30), NoSnapshots, or Every n.
DefaultUse the global config (config:akka:persistence:snapshot-version-count), or 30.
Passivation: PassivationPolicyIdle passivation: PassivationPolicy.Default (configuration, then Akka's 120s), After an idle period, or Never.
FCQRS.Common.PassivationPolicyIdle passivation for an aggregate type, set per entity at registration. Passivation stops an idle actor and releases its in-memory state; the next command recovers it from the journal. Only messages routed through cluster sharding count as activity. Sagas ignore this: their shard regions remember entities, which disables idle passivation in Akka.NET. A saga stops at StopSaga or abort instead.
DefaultUse configuration: `akka.cluster.sharding.<EntityName>.passivate-idle-entity-after`, then `akka.cluster.sharding.passivate-idle-entity-after`, then Akka.NET's 120s.
wireSagaStarters: IActor -> SagaHandle list -> unitWire every registered saga into one saga-starter (or the empty starter if none). Call after the aggregates + sagas are registered.
// The statement table lives in the same SQLite file as the journal.
do
use connection = new SqliteConnection(connectionString)
Statement.createTable connection
// FCQRS passes each new journal event to Statement.handle, and commits the
// handler's rows and the projection's progress in one transaction.
let store =
SqlProjectionStore(
ProjectionSqlDialect.Sqlite,
Func<DbConnection>(fun () -> new SqliteConnection(connectionString)))
let options = TransactionalProjectionOptions("Statement", store)
let statement = Fcqrs.transactionalProjection api options Statement.handle
connection: SqliteConnectionMicrosoft.Data.Sqlite.SqliteConnectionRepresents a connection to a SQLite database. Connection Strings Async Limitations
connectionString: stringFcqrs_250-tutorial_004-show-a-statement.md_page.StatementcreateTable: DbConnection -> unitstore: SqlProjectionStore``.ctor``: ProjectionSqlDialect * Func<DbConnection> -> SqlProjectionStoreUses one database for both the journal and the transactional read model.
FCQRS.ProjectionStorage.ProjectionSqlDialectDatabase SQL syntax supported by the transactional projection store.
Sqlite: ProjectionSqlDialectSQLite with a provider such as Microsoft.Data.Sqlite.
System.Func`1Encapsulates a method that has no parameters and returns a value of the type specified by the parameter. The type of the return value of the method that this delegate encapsulates. The return value of the method that this delegate encapsulates.
System.Data.Common.DbConnectionDefines the core behavior of database connections and provides a base class for database-specific connections.
options: TransactionalProjectionOptions``.ctor``: string * SqlProjectionStore -> TransactionalProjectionOptionsstatement: IProjectionFCQRS.FSharp.FcqrstransactionalProjection: IActor -> TransactionalProjectionOptions -> (DbConnection -> DbTransaction -> Akka.Persistence.Query.EventEnvelope -> Threading.Tasks.Task) -> IProjectionStarts a transactional projection with journal-wide CatchUpAsync support. Write the read model through the supplied connection and transaction; FCQRS commits those updates and contiguous per-persistence-ID progress together. Retain unprocessed journal history. Each pass applies events in journal order: per persistence ID in sequence order, and across persistence IDs in write order on SQLite; on PostgreSQL an event that commits late is applied in a later pass. See TransactionalProjectionOptions.
api: IActorhandle: DbConnection -> DbTransaction -> Akka.Persistence.Query.EventEnvelope -> Threading.Tasks.Task// The statement table lives in the same SQLite file as the journal.
using (var connection = new SqliteConnection(connectionString))
Statement.CreateTable(connection);
// FCQRS passes each new journal event to Statement.Handle, and commits the
// handler's rows and the projection's progress in one transaction.
var store = new SqlProjectionStore(
ProjectionSqlDialect.Sqlite, () => new SqliteConnection(connectionString));
var options = new TransactionalProjectionOptions("Statement", store);
var builder = Host.CreateApplicationBuilder();
builder.Logging.ClearProviders();
builder.Services.AddFcqrs(connectionString, "accounts")
.AddAggregate<Account>()
.AddTransactionalProjection(options, Statement.Handle);
using var host = builder.Build();
await host.StartAsync();
"Statement" names the projection's progress. FCQRS records, for each account, the version of the
last event this projection committed, and resumes from there after a restart. Give each read model its
own name. In C#, AddTransactionalProjection registers the projection with the host, and the program
gets it as IProjection.
Send commands and wait for the statement
Shared setup
let describe event =
match event with
| Opened owner -> $"Opened for {owner}"
| Deposited amount -> $"Deposited {amount}"
| Withdrawn amount -> $"Withdrew {amount}"
| Rejected reason -> $"Rejected: {reason}"
describe: AccountEvent -> stringevent: AccountEventOpenedowner: stringDepositedamount: decimalWithdrawnRejectedreason: stringlet alice = Fcqrs.aggregateId "alice"
// Send a command and wait until the statement includes the event it stored.
let send command =
let reply =
Fcqrs.sendAwaiting
statement accounts (Fcqrs.newCid ()) alice command (fun _ -> true)
|> Async.RunSynchronously
let description = describe reply.EventDetails
printfn $"{description} (version {reply.Version})"
send (Open "Alice")
send (Deposit 100m)
send (Withdraw 30m)
send (Deposit 50m)
send (Withdraw 500m)
alice: FCQRS.Model.Data.AggregateIdFCQRS.FSharp.FcqrsaggregateId: string -> FCQRS.Model.Data.AggregateIdAn aggregate id from a string (e.g. a document/user key). Any non-blank id works: the shard names entity actors Uri.EscapeDataString(entityId), so characters Akka actor names would reject directly (spaces, %) are escaped before they reach an actor path.
send: AccountCommand -> unitcommand: AccountCommandreply: Event<AccountEvent>sendAwaiting: FCQRS.Query.ISubscribe<FCQRS.Model.Data.IMessageWithCID> -> AggregateHandle<'Command,'Event> -> FCQRS.Model.Data.CID -> FCQRS.Model.Data.AggregateId -> 'Command -> ('Event -> bool) -> Async<Event<'Event>>Read-your-writes in one call: subscribe on the CID BEFORE sending, send, then await the projection ONLY if the delivered ack was journaled. A deferred (rejection-style) ack never reaches the journal, so a naive await would hang until timeout. The subscribe-before-send ordering is what makes the wait race-free; owning it here means callers cannot get it backwards. Awaits exactly ONE projected event: a batch persist (PersistAllEvents) caller should Subscribe with an explicit take instead. Envelopes without the delivery stamp (pre-stamp FCQRS) are treated as journaled. The projection wait is bounded by `akka.fcqrs.command-timeout` (default 30s, bare number = seconds): a projection that suppresses or filters out the matching notification raises TimeoutException instead of hanging the caller forever.
statement: IProjectionaccounts: AggregateHandle<AccountCommand,AccountEvent>newCid: unit -> FCQRS.Model.Data.CIDA fresh correlation id (UUID v7).
(|>): 'T1 -> ('T1 -> 'U) -> 'UApply a function to a value, the value being on the left, the function on the right The argument. The function. The function result. let doubleIt x = x * 2 3 |> doubleIt // Evaluates to 6
Microsoft.FSharp.Control.FSharpAsyncHolds static members for creating and manipulating asynchronous computations. See also F# Language Guide - Async Workflows. Async Programming
RunSynchronously: Async<'T> * int option * Threading.CancellationToken option -> 'TRuns the asynchronous computation and await its result. If an exception occurs in the asynchronous computation then an exception is re-raised by this function. If no cancellation token is provided then the default cancellation token is used. The computation is started on the current thread if is null, has of true, and no timeout is specified. Otherwise the computation is started by queueing a new work item in the thread pool, and the current thread is blocked awaiting the completion of the computation. The timeout parameter is given in milliseconds. A value of -1 is equivalent to . The computation to run. The amount of time in milliseconds to wait for the result of the computation before raising a . If no value is provided for timeout then a default of -1 is used to correspond to . The cancellation token to be associated with the computation. If one is not supplied, the default cancellation token is used. The result of the computation. Starting Async Computations printfn "A" let result = async { printfn "B" do! Async.Sleep(1000) printfn "C" 17 } |> Async.RunSynchronously printfn "D" Prints "A", "B" immediately, then "C", "D" in 1 second. result is set to 17.
description: stringdescribe: AccountEvent -> stringEventDetails: 'EventDetailsThe specific details or payload of the event.
printfn: Printf.TextWriterFormat<'T> -> 'TPrint to stdout using the given format, and add a newline. The formatter. The formatted result. See Printf.printfn (link: ) for examples.
OpenDepositWithdraw// Sends commands to accounts and returns the event each one replied with.
var accounts = host.Services
.GetRequiredService<Handler<AccountCommand, AccountEvent>>();
// Publishes a notification after it commits each event.
var statement = host.Services.GetRequiredService<IProjection>();
var alice = Values.CreateAggregateId("alice");
// Send a command and wait until the statement includes the event it stored.
async Task Send(AccountCommand command)
{
var cid = Values.NewCID();
// Subscribe first: a notification sent before the subscription is lost.
using var projected = statement.SubscribeForFirst(cid);
var reply = await accounts(_ => true, cid, alice, command);
// A rejection is not stored, so no notification comes for it.
if (reply.Journaled?.Value != false)
await projected.Task.WaitAsync(TimeSpan.FromSeconds(30));
var description = Describe(reply.EventDetails);
Console.WriteLine($"{description} (version {reply.Version})");
}
await Send(new Open("Alice"));
await Send(new Deposit(100m));
await Send(new Withdraw(30m));
await Send(new Deposit(50m));
await Send(new Withdraw(500m));
The account replies as soon as it stores the event. The projection writes the statement row afterwards, in its own transaction, so a query right after the reply can miss the newest row.
A correlation ID labels one request. FCQRS copies it to the events the command causes, and the
projection publishes a notification carrying it after it commits each event. Fcqrs.sendAwaiting
subscribes to the command's correlation ID, sends the command, and waits for that notification. The C#
Send spells out the same steps.
- Subscribe before sending. A notification published before the subscription exists is missed, and the wait would last until its timeout.
- A rejection has no notification. It is not stored, so the projection never sees it.
sendAwaitingreturns right away when the reply was not stored; the C# code checksJournaled. - The wait is bounded.
sendAwaitingraisesTimeoutExceptionafter the command timeout, 30 seconds by default. The C# code usesWaitAsyncfor the same bound.
The last argument, fun _ -> true (_ => true in C#), chooses which reply completes the send. Each
command here produces one reply, so it accepts any.
Waiting covers this projection only. Another read model, or a system that reads the journal, can still
be behind. Subscriptions live in memory for the duration of the request; they are not a queue a client
can reconnect to. To wait for every event stored so far instead of one request's events, call
CatchUpAsync on the projection, as Catch up projections
shows.
Query the statement
// The statement is an ordinary table: read it with SQL.
let printStatement () =
use connection = new SqliteConnection(connectionString)
connection.Open()
use query = connection.CreateCommand()
query.CommandText <-
"SELECT version, entry, amount, balance FROM statement
WHERE account = 'alice' ORDER BY version"
use rows = query.ExecuteReader()
printfn "\nStatement for alice:"
printfn " version entry amount balance"
while rows.Read() do
let version, entry = rows.GetInt64 0, rows.GetString 1
let amount, balance = rows.GetDecimal 2, rows.GetDecimal 3
printfn $" {version,7} {entry,-18}{amount,7}{balance,9}"
printStatement ()
printStatement: unit -> unitconnection: SqliteConnectionMicrosoft.Data.Sqlite.SqliteConnectionRepresents a connection to a SQLite database. Connection Strings Async Limitations
connectionString: stringOpen: unit -> unitOpens a connection to the database using the value of . If Mode=ReadWriteCreate is used (the default) the file is created, if it doesn't already exist. A SQLite error occurs while opening the connection.
query: SqliteCommandCreateCommand: unit -> SqliteCommandCreates a new command associated with the connection. The new command. The command's property will also be set to the current transaction.
CommandText: stringGets or sets the SQL to execute against the database. The SQL to execute against the database. Batching
rows: SqliteDataReaderExecuteReader: unit -> SqliteDataReaderExecutes the against the database and returns a data reader. The data reader. A SQLite error occurs during execution. Database Errors Batching
printfn: Printf.TextWriterFormat<'T> -> 'TPrint to stdout using the given format, and add a newline. The formatter. The formatted result. See Printf.printfn (link: ) for examples.
Read: unit -> boolAdvances to the next row in the result set. if there are more rows; otherwise, .
version: int64entry: stringGetInt64: int -> int64Gets the value of the specified column as a . The zero-based column ordinal. The value of the column.
GetString: int -> stringGets the value of the specified column as a . The zero-based column ordinal. The value of the column.
amount: decimalbalance: decimalGetDecimal: int -> decimalGets the value of the specified column as a . The zero-based column ordinal. The value of the column.
// The statement is an ordinary table: read it with SQL.
void PrintStatement()
{
using var connection = new SqliteConnection(connectionString);
connection.Open();
using var query = connection.CreateCommand();
query.CommandText =
"""
SELECT version, entry, amount, balance FROM statement
WHERE account = 'alice' ORDER BY version
""";
using var rows = query.ExecuteReader();
Console.WriteLine();
Console.WriteLine("Statement for alice:");
Console.WriteLine(" version entry amount balance");
while (rows.Read())
{
var (version, entry) = (rows.GetInt64(0), rows.GetString(1));
var (amount, balance) = (rows.GetDecimal(2), rows.GetDecimal(3));
Console.WriteLine($" {version,7} {entry,-18}{amount,7}{balance,9}");
}
}
PrintStatement();
Reading a read model is a plain SQL query. The table already has the shape the screen needs, so the query does not fold events or join tables.
Run it
From samples/accounts:
dotnet run --project 4-show-a-statement/fsharp
dotnet run --project 4-show-a-statement/csharp
Both programs print:
Opened for Alice (version 1)
Deposited 100 (version 2)
Withdrew 30 (version 3)
Deposited 50 (version 4)
Rejected: Insufficient funds: 120 available (version 4)
Statement for alice:
version entry amount balance
1 Opened for Alice 0 0
2 Deposit 100 100
3 Withdrawal -30 70
4 Deposit 50 120
The rejected withdrawal has no row: the projection only sees stored events.
Run it again
Rejected: The account is already open (version 4)
Deposited 100 (version 5)
Withdrew 30 (version 6)
Deposited 50 (version 7)
Rejected: Insufficient funds: 240 available (version 7)
Statement for alice:
version entry amount balance
1 Opened for Alice 0 0
2 Deposit 100 100
3 Withdrawal -30 70
4 Deposit 50 120
5 Deposit 100 220
6 Withdrawal -30 190
7 Deposit 50 240
The projection resumed after version 4, the last event it had committed for Alice, so rows 1 to 4 were not written again. The primary key would reject a repeated row.
Rebuild the statement
The statement holds nothing that the journal does not, so it can be rebuilt. Stop the program, delete the statement's rows and the projection's progress, and run it again:
DELETE FROM statement;
DELETE FROM fcqrs_projection_progress
WHERE projection_name = 'Statement';
The projection starts again from each account's first event and writes every row, then continues with the new commands. Rebuild this way after fixing a bug in the handler or changing the table. Rebuild a read model covers rebuilding beside a model that must stay available.
Next
Step 5: Transfer money moves money from Alice to Bob. A transfer changes two accounts, and each account decides only with its own state, so the transfer needs a workflow that coordinates them: a saga.