Step 5: Transfer money
Alice sends 30 to Bob. Each account decides only with its own state (step 2), so no account can move money into another. A saga coordinates the transfer: it reacts to an event from one account, stores how far the transfer has got, and sends commands to the accounts involved. When the target account cannot take the money, the saga gives it back to the sender.
The same feature in a CRUD application
A CRUD application moves money between two rows in one database transaction:
BEGIN;
UPDATE accounts SET balance = balance - 30 WHERE id = 'alice';
UPDATE accounts SET balance = balance + 30 WHERE id = 'bob';
COMMIT;
That works only while both rows live in one database. When the target account belongs to another service or another bank, the debit and the credit are separate operations, and something has to finish or undo the transfer when the second half fails, including after a restart. In FCQRS, every account is its own aggregate and stores its own events, so the same holds within one application. The saga is the part that finishes or undoes the transfer, and it stores its progress so a restart does not lose a transfer halfway.
Add transfer commands and events
module Account
open FCQRS.Common
// What a caller, or the transfer saga, can ask an account to do.
type AccountCommand =
| Open of owner: string
| Deposit of amount: decimal
| Withdraw of amount: decimal
| SendTransfer of transferId: string * target: string * amount: decimal
| ReceiveTransfer of transferId: string * source: string * amount: decimal
| RefundTransfer of transferId: string * target: string * 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
| TransferSent of transferId: string * target: string * amount: decimal
| TransferReceived of transferId: string * source: string * amount: decimal
| TransferRefunded of transferId: string * target: string * amount: decimal
| Rejected of reason: string
// What the account knows now. The sets hold the IDs of transfers it has
// handled, so a repeated transfer command moves no money twice.
type AccountState =
{ Owner: string option
Balance: decimal
Sent: Set<string>
Received: Set<string>
Refunded: Set<string> }
// The state before the account's first event.
let initial =
{ Owner = None
Balance = 0m
Sent = Set.empty
Received = Set.empty
Refunded = Set.empty }
Fcqrs_250-tutorial_005-transfer-money.md_page.AccountFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountCommandOpenowner: stringstringAn abbreviation for the CLI type . Basic Types
Depositamount: decimaldecimalAn abbreviation for the CLI type . Basic Types
WithdrawSendTransfertransferId: stringtarget: stringReceiveTransfersource: stringRefundTransferFcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventOpenedDepositedWithdrawnTransferSentTransferReceivedTransferRefundedRejectedreason: stringFcqrs_250-tutorial_005-transfer-money.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: decimalSent: Set<string>Microsoft.FSharp.Collections.FSharpSet`1Immutable sets based on binary trees, where elements are ordered by F# generic comparison. By default comparison is the F# structural comparison function or uses implementations of the IComparable interface on element values. See the module for further operations on sets. All members of this class are thread-safe and may be used concurrently from multiple threads.
Received: Set<string>Refunded: Set<string>initial: AccountStateNoneThe representation of "No value"
Microsoft.FSharp.Collections.SetModuleContains operations for working with values of type .
empty: Set<'T>The empty set for the type 'T. Set.empty<int> Evaluates to set [ ].
// What a caller, or the transfer saga, can ask an account to do.
public union AccountCommand(
Open, Deposit, Withdraw, SendTransfer, ReceiveTransfer, RefundTransfer);
public sealed record Open(string Owner);
public sealed record Deposit(decimal Amount);
public sealed record Withdraw(decimal Amount);
public sealed record SendTransfer(
string TransferId, string Target, decimal Amount);
public sealed record ReceiveTransfer(
string TransferId, string Source, decimal Amount);
public sealed record RefundTransfer(
string TransferId, string Target, decimal Amount);
// What the account replies. Rejected is a reply only: it is never stored.
public union AccountEvent(
Opened, Deposited, Withdrawn,
TransferSent, TransferReceived, TransferRefunded, Rejected);
public sealed record Opened(string Owner);
public sealed record Deposited(decimal Amount);
public sealed record Withdrawn(decimal Amount);
public sealed record TransferSent(
string TransferId, string Target, decimal Amount);
public sealed record TransferReceived(
string TransferId, string Source, decimal Amount);
public sealed record TransferRefunded(
string TransferId, string Target, decimal Amount);
public sealed record Rejected(string Reason);
// What the account knows now. The sets hold the IDs of transfers it has
// handled, so a repeated transfer command moves no money twice.
public sealed record AccountState(
string? Owner,
decimal Balance,
ImmutableHashSet<string> Sent,
ImmutableHashSet<string> Received,
ImmutableHashSet<string> Refunded);
A transfer has three messages, each named with a transfer ID that the caller chooses:
SendTransferasks the source account to send money. It storesTransferSentand debits itself.ReceiveTransferasks the target account to take the money. It storesTransferReceived.RefundTransferasks the source account to take back money the target did not accept. It storesTransferRefunded.
The state now remembers which transfer IDs each account has sent, received, and refunded.
Make each transfer command safe to repeat
// 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 | SendTransfer(_, _, amount)), _
when amount <= 0m -> DeferEvent(Rejected "The amount must be positive")
| Deposit amount, _ -> PersistEvent(Deposited amount)
| (Withdraw amount | SendTransfer(_, _, amount)), _
when amount > state.Balance ->
DeferEvent(Rejected $"Insufficient funds: {state.Balance} available")
| Withdraw amount, _ -> PersistEvent(Withdrawn amount)
| SendTransfer(id, _, _), _ when state.Sent.Contains id ->
DeferEvent(Rejected $"Transfer {id} was already sent")
| SendTransfer(id, target, amount), _ ->
PersistEvent(TransferSent(id, target, amount))
// A repeated delivery gets the first answer; no money moves.
| ReceiveTransfer(id, source, amount), _ when state.Received.Contains id ->
DeferEvent(TransferReceived(id, source, amount))
| ReceiveTransfer(id, source, amount), _ ->
PersistEvent(TransferReceived(id, source, amount))
| RefundTransfer(id, target, amount), _ when state.Refunded.Contains id ->
DeferEvent(TransferRefunded(id, target, amount))
| RefundTransfer(id, target, amount), _ ->
PersistEvent(TransferRefunded(id, target, amount))
// Applies one event. A rejection or a repeated reply 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 }
| TransferSent(id, _, amount) ->
{ state with
Balance = state.Balance - amount
Sent = state.Sent.Add id }
// FCQRS folds a repeated reply too; a known ID changes nothing.
| TransferReceived(id, _, _) when state.Received.Contains id -> state
| TransferReceived(id, _, amount) ->
{ state with
Balance = state.Balance + amount
Received = state.Received.Add id }
| TransferRefunded(id, _, _) when state.Refunded.Contains id -> state
| TransferRefunded(id, _, amount) ->
{ state with
Balance = state.Balance + amount
Refunded = state.Refunded.Add id }
| Rejected _ -> state
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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountCommandstate: AccountStateFcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountStateCommandDetails: 'CommandDetailsThe specific details or payload of the command.
Owner: string optionOpenSomeThe 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.
Rejectedowner: stringNoneThe representation of "No value"
PersistEventPersist the event to the journal. The actor's state will be updated using the event handler *after* persistence succeeds.
OpenedDepositamount: decimalWithdrawSendTransfer(<=): '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
Deposited(>): '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
Balance: decimalWithdrawnid: stringContains: string -> boolA useful shortcut for Set.contains. See the Set module for further operations on sets. The value to check. True if the set contains value. let set = Set.empty.Add(2).Add(3) printfn $"Does the set contain 1? {set.Contains(1)}" The sample evaluates to the following output: Does the set contain 1? false
Sent: Set<string>target: stringTransferSentReceiveTransfersource: stringReceived: Set<string>TransferReceivedRefundTransferRefunded: Set<string>TransferRefundedfold: 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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventEventDetails: 'EventDetailsThe specific details or payload of the event.
(-): ^T1 -> ^T2 -> ^T3Overloaded subtraction operator The first parameter. The second parameter. The result of the operation. 10 - 2 // Evaluates to 8
Add: string -> Set<string>A useful shortcut for Set.add. Note this operation produces a new set and does not mutate the original set. The new set will share many storage nodes with the original. See the Set module for further operations on sets. The value to add to the set. The result set. let set = Set.empty.Add(1).Add(1).Add(2) printfn $"The new set is: {set}" The sample evaluates to the following output: The new set is: set [1; 2]
(+): ^T1 -> ^T2 -> ^T3Overloaded addition operator The first parameter. The second parameter. The result of the operation. 2 + 2 // Evaluates to 4 "Hello " + "World" // Evaluates to "Hello World"
public sealed class Account
: Aggregate<AccountState, AccountCommand, AccountEvent>
{
// The name stored with every event of this aggregate.
public override string EntityName => "Account";
// The state before the account's first event.
public override AccountState InitialState => new(null, 0m, [], [], []);
// Chooses what to do with a command, based on the current state.
public override EventAction<AccountEvent> HandleCommand(
Command<AccountCommand> command, AccountState state) =>
(command.CommandDetails, state.Owner) switch
{
(Open, not null) => Reject("The account is already open"),
(Open open, null) => Store(new Opened(open.Owner)),
(_, null) => Reject("The account is not open"),
(Deposit { Amount: <= 0m } or Withdraw { Amount: <= 0m }
or SendTransfer { Amount: <= 0m }, _) =>
Reject("The amount must be positive"),
(Deposit deposit, _) => Store(new Deposited(deposit.Amount)),
(Withdraw withdraw, _) when withdraw.Amount > state.Balance =>
Reject($"Insufficient funds: {state.Balance} available"),
(SendTransfer send, _) when send.Amount > state.Balance =>
Reject($"Insufficient funds: {state.Balance} available"),
(Withdraw withdraw, _) => Store(new Withdrawn(withdraw.Amount)),
(SendTransfer send, _) when state.Sent.Contains(send.TransferId) =>
Reject($"Transfer {send.TransferId} was already sent"),
(SendTransfer send, _) => Store(
new TransferSent(send.TransferId, send.Target, send.Amount)),
// A repeated delivery gets the first answer; no money moves.
(ReceiveTransfer receive, _)
when state.Received.Contains(receive.TransferId) =>
Repeat(new TransferReceived(
receive.TransferId, receive.Source, receive.Amount)),
(ReceiveTransfer receive, _) =>
Store(new TransferReceived(
receive.TransferId, receive.Source, receive.Amount)),
(RefundTransfer refund, _)
when state.Refunded.Contains(refund.TransferId) =>
Repeat(new TransferRefunded(
refund.TransferId, refund.Target, refund.Amount)),
(RefundTransfer refund, _) =>
Store(new TransferRefunded(
refund.TransferId, refund.Target, refund.Amount))
};
// Applies one event. A rejection or a repeated reply changes nothing.
public override AccountState ApplyEvent(
Event<AccountEvent> stored, AccountState state) =>
stored.EventDetails switch
{
Opened opened => state with { Owner = opened.Owner },
Deposited deposited =>
state with { Balance = state.Balance + deposited.Amount },
Withdrawn withdrawn =>
state with { Balance = state.Balance - withdrawn.Amount },
TransferSent sent => state with
{
Balance = state.Balance - sent.Amount,
Sent = state.Sent.Add(sent.TransferId)
},
// FCQRS folds a repeated reply too; a known ID changes nothing.
TransferReceived received
when state.Received.Contains(received.TransferId) => state,
TransferReceived received => state with
{
Balance = state.Balance + received.Amount,
Received = state.Received.Add(received.TransferId)
},
TransferRefunded refunded
when state.Refunded.Contains(refunded.TransferId) => state,
TransferRefunded refunded => state with
{
Balance = state.Balance + refunded.Amount,
Refunded = state.Refunded.Add(refunded.TransferId)
},
Rejected => state
};
// Stores the event and replies with it.
static EventAction<AccountEvent> Store(AccountEvent @event) =>
EventActions.Persist(@event);
// Replies without storing anything.
static EventAction<AccountEvent> Reject(string reason) =>
EventActions.Defer<AccountEvent>(new Rejected(reason));
// Replies with an earlier answer again, without storing it.
static EventAction<AccountEvent> Repeat(AccountEvent @event) =>
EventActions.Defer(@event);
}
The rules use the transfer IDs in two ways:
- A second
SendTransferwith the same ID is rejected. A caller that retries a request, for example after a timeout, cannot send the money twice. - A second
ReceiveTransferorRefundTransferwith the same ID gets the first answer again, as a reply that is not stored. The saga can deliver these commands more than once, as the next sections show, and the money must move only once.
FCQRS applies a repeated reply with fold too, as it does every deferred reply, so fold checks the
ID and leaves the state unchanged. An account that is not open rejects ReceiveTransfer; that
rejection is what sends a transfer back.
Describe where a transfer can be
Shared setup
module Statement =
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open FCQRS.Common
open Account
// 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
// 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
| TransferSent(id, target, amount) -> do! add $"Transfer {id} to {target}" -amount
| TransferReceived(id, source, amount) -> do! add $"Transfer {id} from {source}" amount
| TransferRefunded(id, _, amount) -> do! add $"Refund of transfer {id}" amount
// A rejection is a reply only; the journal never holds one.
| Rejected _ -> ()
| _ -> ()
}
:> Task
module Transfer =
open System
open FCQRS.Common
open FCQRS.FSharp
open Account
Fcqrs_250-tutorial_005-transfer-money.md_page.StatementSystemDataCommonThreadingTasksAkkaPersistenceQueryDapperFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
Fcqrs_250-tutorial_005-transfer-money.md_page.AccountcreateTable: 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 ()
addRow: stringhandle: DbConnection -> DbTransaction -> EventEnvelope -> Tasktransaction: 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_005-transfer-money.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.
TransferSentid: stringtarget: stringTransferReceivedsource: stringTransferRefundedRejectedFcqrs_250-tutorial_005-transfer-money.md_page.TransferFCQRS.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 -> ...)
// What a transfer moves, and between which accounts.
type TransferDetails =
{ Id: string
Source: string
Target: string
Amount: decimal }
// Where a transfer is.
type TransferState =
// Waiting for the target account to take the money.
| Delivering of TransferDetails
// The target turned it down: waiting for the source account's refund.
| Refunding of TransferDetails
| Completed
Fcqrs_250-tutorial_005-transfer-money.md_page.Transfer.TransferDetailsId: stringstringAn abbreviation for the CLI type . Basic Types
Source: stringTarget: stringAmount: decimaldecimalAn abbreviation for the CLI type . Basic Types
Fcqrs_250-tutorial_005-transfer-money.md_page.Transfer.TransferStateDeliveringRefundingCompleted// What a transfer moves, and between which accounts.
public sealed record TransferDetails(
string Id, string Source, string Target, decimal Amount);
// Where a transfer is.
public union TransferState(Delivering, Refunding, Completed);
// Waiting for the target account to take the money.
public sealed record Delivering(TransferDetails Transfer);
// The target turned it down: waiting for the source account's refund.
public sealed record Refunding(TransferDetails Transfer);
public sealed record Completed;
// The saga needs no fixed data besides its state.
public sealed record TransferData;
The saga's state is what it has stored about one transfer. TransferDetails holds what the transfer
moves; both waiting states carry it, because each needs it for its next command. Write the states as
a table before the code:
| State | Event | Next state | Command sent |
|---|---|---|---|
| none | TransferSent |
Delivering |
ReceiveTransfer to the target |
Delivering |
TransferReceived |
Completed |
none; the saga stops |
Delivering |
Rejected from the target |
Refunding |
RefundTransfer to the source |
Refunding |
TransferRefunded |
Completed |
none; the saga stops |
React to events
// Turns an event into the next state to store.
let handleEvent (message: obj) (saga: SagaState<unit, TransferState option>) =
match message, saga.State with
| (:? Event<AccountEvent> as event), state ->
// Sender is the ID of the account that stored the event.
let sender = string event.Sender.Value
match event.EventDetails, state with
| TransferSent(id, target, amount), None ->
let transfer =
{ Id = id; Source = sender; Target = target; Amount = amount }
StateChangedEvent(Delivering transfer)
| TransferReceived(id, _, _), Some(Delivering transfer)
when id = transfer.Id -> StateChangedEvent Completed
| Rejected _, Some(Delivering transfer) when sender = transfer.Target ->
StateChangedEvent(Refunding transfer)
| TransferRefunded(id, _, _), Some(Refunding transfer)
when id = transfer.Id -> StateChangedEvent Completed
| _ -> UnhandledEvent
// No answer in time. The outcome is unknown, so enter the same state again,
// which sends the command again.
| :? ExpectationExhausted, Some(Delivering _ | Refunding _ as waiting) ->
StateChangedEvent waiting
| _ -> UnhandledEvent
handleEvent: obj -> SagaState<unit,TransferState option> -> EventAction<TransferState>message: objobjAn abbreviation for the CLI type . Basic Types
saga: SagaState<unit,TransferState option>FCQRS.Common.SagaState`2Represents the state of a saga instance. <typeparam name="'SagaData">The type of the custom data held by the saga.</typeparam> <typeparam name="'State">The type representing the saga's current state machine state (e.g., an enum or DU).</typeparam>
unitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
Fcqrs_250-tutorial_005-transfer-money.md_page.Transfer.TransferStateoptionThe 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
State: 'StateThe current state machine state of the saga.
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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventevent: Event<AccountEvent>state: TransferState optionsender: 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.
EventDetails: 'EventDetailsThe specific details or payload of the event.
TransferSentid: stringtarget: stringamount: decimalNoneThe representation of "No value"
transfer: TransferDetailsId: stringSource: stringTarget: stringAmount: decimalStateChangedEventIndicate that the state of a saga has changed (used internally by sagas for persistence).
DeliveringTransferReceivedSomeThe representation of "Value of type 'T" The input value. An option representing the value.
(=): 'T -> 'T -> boolStructural equality The first parameter. The second parameter. The result of the comparison. 5 = 5 // Evaluates to true 5 = 6 // Evaluates to false [1; 2] = [1; 2] // Evaluates to true (1, 5) = (1, 6) // Evaluates to false
CompletedRejectedRefundingTransferRefundedUnhandledEventIndicate that the command or event could not be handled in the current state.
FCQRS.Common.ExpectationExhaustedDelivered to the saga's event handler when an expectation's deadline has passed without a state transition. The handler must match this type and answer with a state change (typically to a domain failure or compensation state). An unhandled exhaustion is logged as an error and re-delivered one deadline period later; the framework never invents a terminal state. The original reply may still arrive after exhaustion — the escalated state's handler should decide what a late success means.
waiting: TransferState// Turns the event that started the transfer into its first state.
public override EventAction<TransferState> Start(
object message, TransferData data) =>
message is Event<AccountEvent>
{
EventDetails: TransferSent sent, Sender: { } sender
}
? StateChanged(new Delivering(
new(sent.TransferId, sender.Value.ToString(), sent.Target,
sent.Amount)))
: Unhandled();
// Turns a later event into the next state to store.
public override EventAction<TransferState> HandleEvent(
object message, SagaState<TransferData, TransferState> saga) =>
(message, saga.State) switch
{
// Sender is the ID of the account that stored the event.
(Event<AccountEvent> { Sender: { } sender } stored, var state) =>
React(stored.EventDetails, sender.Value.ToString(), state),
// No answer in time. The outcome is unknown, so enter the same
// state again, which sends the command again.
(ExpectationExhausted, Delivering or Refunding) =>
StateChanged(saga.State),
_ => Unhandled()
};
static EventAction<TransferState> React(
AccountEvent @event, string sender, TransferState state) =>
(@event, state) switch
{
(TransferReceived received, Delivering delivering)
when received.TransferId == delivering.Transfer.Id =>
StateChanged(new Completed()),
(Rejected, Delivering delivering)
when sender == delivering.Transfer.Target =>
StateChanged(new Refunding(delivering.Transfer)),
(TransferRefunded refunded, Refunding refunding)
when refunded.TransferId == refunding.Transfer.Id =>
StateChanged(new Completed()),
_ => Unhandled()
};
handleEvent receives events from every account that takes part in the transfer, as obj,
together with the stored state. Before the first state is stored, that state is None. It returns the
next state to store, or UnhandledEvent for an event that does not belong in the current state. It
sends no commands: FCQRS stores the state first.
C# splits the two cases. Start receives the event that started the transfer, before the saga has a
state, and returns the first state. HandleEvent receives the later events with the stored state.
The events reach this saga through the transfer's correlation ID (step 4).
The saga's commands carry the correlation ID of the SendTransfer request, and so do the replies the
accounts send back. Sender identifies the account that replied, so a rejection counts only when it
comes from the target.
Send commands for each state
// Returns the commands for a stored state. It runs again after recovery; the
// last argument says whether FCQRS is recovering, and this saga ignores it.
let applySideEffects accounts (saga: SagaState<unit, TransferState>) _ =
// Send now and every 5 seconds; after 30 seconds, tell handleEvent.
let deliver command =
let retry = FixedInterval(TimeSpan.FromSeconds 5.)
expecting (TimeSpan.FromSeconds 30.) retry [ command ], []
match saga.State with
| Delivering transfer ->
let command =
ReceiveTransfer(transfer.Id, transfer.Source, transfer.Amount)
deliver (toAggregate accounts transfer.Target command)
| Refunding transfer ->
let command =
RefundTransfer(transfer.Id, transfer.Target, transfer.Amount)
deliver (toOriginator accounts command)
| Completed -> StopSaga, []
applySideEffects: AggregateFactory -> SagaState<unit,TransferState> -> 'a -> SagaTransition<'b> * 'c listaccounts: AggregateFactorysaga: SagaState<unit,TransferState>FCQRS.Common.SagaState`2Represents the state of a saga instance. <typeparam name="'SagaData">The type of the custom data held by the saga.</typeparam> <typeparam name="'State">The type representing the saga's current state machine state (e.g., an enum or DU).</typeparam>
unitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
Fcqrs_250-tutorial_005-transfer-money.md_page.Transfer.TransferStatedeliver: ExecuteCommand -> SagaTransition<'d> * 'e listcommand: ExecuteCommandretry: RetryScheduleFixedIntervalRe-send at a fixed interval.
System.TimeSpanRepresents a time interval.
FromSeconds: float -> TimeSpanReturns a that represents a specified number of seconds, where the specification is accurate to the nearest millisecond. A number of seconds, accurate to the nearest millisecond. is less than TimeSpan.MinValue or greater than TimeSpan.MaxValue. -or- is . -or- is . is equal to . An object that represents .
expecting: TimeSpan -> RetrySchedule -> ExecuteCommand list -> SagaTransition<'State>Declare a saga expectation: stay in this state, send `resend` now, re-send exactly those commands on `retryEvery`, and once `deadline` (measured from the persisted state-entry time, so restarts cannot postpone it) has passed without a state transition, deliver an ExpectationExhausted message to HandleEvent. The handler must answer it with a transition, typically to a failure or compensation state. Resend commands must be retry-safe and must not carry their own DelayInMs.
State: 'StateThe current state machine state of the saga.
Deliveringtransfer: TransferDetailscommand: AccountCommandReceiveTransferId: stringSource: stringAmount: decimaltoAggregate: AggregateFactory -> string -> obj -> ExecuteCommandSend a command to a specific aggregate instance by id (cross-aggregate).
Target: stringRefundingRefundTransfertoOriginator: AggregateFactory -> obj -> ExecuteCommandSend a command back to the saga's originator aggregate.
CompletedStopSagaThe saga should stop and terminate
// Returns the commands for a stored state. It runs again after recovery;
// `recovering` says whether FCQRS is recovering, and this saga ignores it.
public override SagaSideEffectResult<TransferState> ApplySideEffects(
SagaState<TransferData, TransferState> saga, bool recovering) =>
saga.State switch
{
Delivering(var (id, source, target, amount)) =>
Deliver(ToAccount(target, new ReceiveTransfer(id, source, amount))),
Refunding(var (id, _, target, amount)) =>
Deliver(ToSource(new RefundTransfer(id, target, amount))),
Completed => new() { Transition = StopSaga(), Commands = [] }
};
// Send now and every 5 seconds; after 30 seconds, tell HandleEvent.
static SagaSideEffectResult<TransferState> Deliver(ExecuteCommand command) =>
new()
{
Transition = Stay(),
Expect = Expectations.Create(
[command],
TimeSpan.FromSeconds(30),
RetrySchedules.Fixed(TimeSpan.FromSeconds(5)))
};
// SagaCommands takes any object. These helpers take an AccountCommand, so the
// compiler rejects anything else sent to an account.
ExecuteCommand ToAccount(string id, AccountCommand command) =>
SagaCommands.ToAggregate(accounts, id, command);
ExecuteCommand ToSource(AccountCommand command) =>
SagaCommands.ToOriginator(accounts, command);
applySideEffects (ApplySideEffects in C#) runs after FCQRS stores a state. toAggregate sends a
command to an account by ID; toOriginator sends one to the account whose event started the saga.
Completed returns StopSaga, which ends the saga.
The function runs again when FCQRS recovers the saga after a restart. The journal shows the stored
state, but not whether the command left the process before it stopped, so the saga sends it again.
This is why ReceiveTransfer and RefundTransfer must be safe to repeat.
expecting (Expectations.Create in C#) gives each wait a deadline. FCQRS sends the command, sends
it again every 5 seconds until a new state is stored, and after 30 seconds passes
ExpectationExhausted to handleEvent. The deadline counts from the moment the state was stored, so
a restart does not extend it.
ExpectationExhausted means that no answer came, not that the transfer failed: Bob may have taken the
money. The money is in flight, so the transfer must not give up. handleEvent enters the same state
again, which stores a new state and starts another 30 seconds. Every round is stored, so a stuck
transfer is visible in the saga's journal.
In C#, SagaCommands takes a command as object. The helpers take an AccountCommand, so the
compiler rejects anything that is not an account command.
Register the saga
// A transfer starts when an account stores TransferSent.
let startsOn (event: Event<AccountEvent>) =
match event.EventDetails with
| TransferSent _ -> true
| _ -> false
let definition accounts =
{ Name = "Transfer"
InitialData = ()
Originator = accounts
HandleEvent = handleEvent
ApplySideEffects = applySideEffects accounts
StartOn = startsOn
Snapshots = Default }
startsOn: Event<AccountEvent> -> boolevent: 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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventEventDetails: 'EventDetailsThe specific details or payload of the event.
TransferSentdefinition: AggregateFactory -> Saga<unit,TransferState,AccountEvent>accounts: AggregateFactoryName: stringInitialData: 'DataOriginator: AggregateFactoryThe aggregate the saga starts from (its commands' Originator target).
HandleEvent: obj -> SagaState<'Data,'State option> -> EventAction<'State>handleEvent: obj -> SagaState<unit,TransferState option> -> EventAction<TransferState>ApplySideEffects: SagaState<'Data,'State> -> bool -> SagaTransition<'State> * ExecuteCommand listapplySideEffects: AggregateFactory -> SagaState<unit,TransferState> -> 'a -> SagaTransition<'b> * 'c listStartOn: Event<'OriginatorEvent> -> boolWhich originator events spawn an instance of this saga. Typed to the originator's event so 'OriginatorEvent is inferred from the definition; there is no type argument to remember (or to get silently wrong).
Snapshots: SnapshotPolicySnapshot cadence: Default (config / 30), NoSnapshots, or Every n.
DefaultUse the global config (config:akka:persistence:snapshot-version-count), or 30.
// A transfer starts when an account stores TransferSent.
public override bool StartsOn(Event<AccountEvent> stored) =>
stored.EventDetails is TransferSent;
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.Model.Data
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 }
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 -> ...)
ModelFCQRS.Model.DataFCQRS.ProjectionStorageSQL storage for journal-wide, transactional projection catch-up.
FCQRS.ProjectionsTransactional projections with a journal-wide, durable catch-up boundary.
Fcqrs_250-tutorial_005-transfer-money.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.
// Register the saga with the accounts it sends commands to, then install its
// start rule. From now on, each stored TransferSent starts one transfer.
let transfers = Fcqrs.saga api (Transfer.definition accounts.Factory)
Fcqrs.wireSagaStarters api [ transfers ]
transfers: SagaHandleFCQRS.FSharp.Fcqrssaga: IActor -> Saga<'Data,'State,'OriginatorEvent> -> SagaHandleRegister a saga and return its handle. The originator-event type is inferred from `def.StartOn`, so there are no type arguments to supply: `Fcqrs.saga api def`.
api: IActorFcqrs_250-tutorial_005-transfer-money.md_page.Transferdefinition: AggregateFactory -> Saga<unit,Transfer.TransferState,AccountEvent>accounts: AggregateHandle<AccountCommand,AccountEvent>Factory: AggregateFactoryEntity-ref factory (DEFAULT_SHARD applied). Hand this to a saga to target it.
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.
// Register the saga with the accounts it sends commands to. The host installs
// its start rule: from now on, each stored TransferSent starts one transfer.
builder.Services.AddFcqrs(connectionString, "accounts")
.AddAggregate<Account>()
.AddSaga(services => new Transfer(services.AggregateFactory<Account>()))
.AddTransactionalProjection(options, Statement.Handle);
startsOn (StartsOn in C#) selects the event that begins a transfer. FCQRS starts one saga for each
stored TransferSent. A rejected SendTransfer is not stored, so it starts nothing. Transfer is
the saga's stored name, like an aggregate's EntityName.
Fcqrs.wireSagaStarters installs the start rules, and must run after every saga is registered. It
also makes Alice's account wait, before it publishes TransferSent, until the new saga listens for the
transfer's correlation ID, so the saga cannot miss its first event. The C# host installs the rules
itself when it starts.
Send transfers and wait for them
Shared setup
// The statement from step 4, with rows for transfers.
do
use connection = new SqliteConnection(connectionString)
Statement.createTable connection
let store =
SqlProjectionStore(
ProjectionSqlDialect.Sqlite,
Func<DbConnection>(fun () -> new SqliteConnection(connectionString)))
let options = TransactionalProjectionOptions("Statement", store)
let statement = Fcqrs.transactionalProjection api options Statement.handle
let describe event =
match event with
| Opened owner -> $"Opened for {owner}"
| Deposited amount -> $"Deposited {amount}"
| Withdrawn amount -> $"Withdrew {amount}"
| TransferSent(id, target, amount) -> $"Sent {amount} to {target} ({id})"
| TransferReceived(id, source, amount) -> $"Received {amount} from {source} ({id})"
| TransferRefunded(id, _, amount) -> $"Refunded {amount} ({id})"
| Rejected reason -> $"Rejected: {reason}"
let show (reply: Event<AccountEvent>) =
let stored = if reply.Journaled = Some true then "stored" else "not stored"
let description = describe reply.EventDetails
printfn $"{description} (version {reply.Version}, {stored})"
let alice = Fcqrs.aggregateId "alice"
let bob = Fcqrs.aggregateId "bob"
// Send a command and wait until the statement includes the event it stored.
let send account command =
Fcqrs.sendAwaiting statement accounts (Fcqrs.newCid ()) account command (fun _ -> true)
|> Async.RunSynchronously
|> show
connection: SqliteConnectionMicrosoft.Data.Sqlite.SqliteConnectionRepresents a connection to a SQLite database. Connection Strings Async Limitations
connectionString: stringFcqrs_250-tutorial_005-transfer-money.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.Taskdescribe: AccountEvent -> stringevent: AccountEventOpenedowner: stringDepositedamount: decimalWithdrawnTransferSentid: stringtarget: stringTransferReceivedsource: stringTransferRefundedRejectedreason: stringshow: Event<AccountEvent> -> unitreply: 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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventstored: stringJournaled: bool optionWhether this envelope's event was journaled, read from the delivery stamp: Some true (a projection event will follow), Some false (a deferred/publish-only reply — nothing to await), or None (an envelope that never passed through aggregate delivery, e.g. read back from the journal, or produced by a pre-stamp FCQRS).
(=): 'T -> 'T -> boolStructural equality The first parameter. The second parameter. The result of the comparison. 5 = 5 // Evaluates to true 5 = 6 // Evaluates to false [1; 2] = [1; 2] // Evaluates to true (1, 5) = (1, 6) // Evaluates to false
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
description: 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.
alice: AggregateIdaggregateId: string -> 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.
bob: AggregateIdsend: AggregateId -> AccountCommand -> unitaccount: AggregateIdcommand: AccountCommandsendAwaiting: FCQRS.Query.ISubscribe<IMessageWithCID> -> AggregateHandle<'Command,'Event> -> CID -> 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.
accounts: AggregateHandle<AccountCommand,AccountEvent>newCid: unit -> 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.
// A transfer ends when the target stores the money or the source gets it back.
let finished (message: IMessageWithCID) =
match message with
| :? Event<AccountEvent> as event ->
match event.EventDetails with
| TransferReceived _ | TransferRefunded _ -> true
| _ -> false
| _ -> false
// Ask Alice's account to send money, and wait until the saga has finished.
let transfer id target amount =
let cid = Fcqrs.newCid ()
// Subscribe first: the saga can finish before the reply arrives.
use outcome = statement.Subscribe(cid, finished, 1)
let command = SendTransfer(id, target, amount)
let reply =
accounts.Send cid alice command (fun _ -> true)
|> Async.RunSynchronously
show reply
if reply.Journaled = Some true then
outcome.Task.WaitAsync(TimeSpan.FromSeconds 30.).Wait()
transfer "t1" "bob" 30m
// Carol has no account, so this transfer comes back.
transfer "t2" "carol" 20m
finished: IMessageWithCID -> boolmessage: IMessageWithCIDFCQRS.Model.Data.IMessageWithCIDInterface for messages that carry a Correlation ID (CID).
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>
Fcqrs_250-tutorial_005-transfer-money.md_page.Account.AccountEventevent: Event<AccountEvent>EventDetails: 'EventDetailsThe specific details or payload of the event.
TransferReceivedTransferRefundedtransfer: string -> string -> decimal -> unitid: stringtarget: stringamount: decimalcid: CIDFCQRS.FSharp.FcqrsnewCid: unit -> CIDA fresh correlation id (UUID v7).
outcome: FCQRS.Query.IAwaitableDisposablestatement: IProjectionSubscribe: CID * (IMessageWithCID -> bool) * int * (IMessageWithCID -> unit) option * Threading.CancellationToken option -> FCQRS.Query.IAwaitableDisposableSubscribes to events matching a specific correlation ID and an additional filter. The correlation ID to match. Additional predicate to filter events after CID matching. Maximum number of events to process. Optional callback function to handle the event. An optional cancellation token to cancel the subscription.
command: AccountCommandSendTransferreply: Event<AccountEvent>accounts: AggregateHandle<AccountCommand,AccountEvent>Send: CID -> AggregateId -> 'Command -> ('Event -> bool) -> Async<Event<'Event>>Send a command and await the first matching aggregate event. Fails with SagaAlreadyStartedException, with nothing stored, when the command would start a saga that its correlation ID already started on this aggregate.
alice: AggregateId(|>): '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.
show: Event<AccountEvent> -> unitJournaled: bool optionWhether this envelope's event was journaled, read from the delivery stamp: Some true (a projection event will follow), Some false (a deferred/publish-only reply — nothing to await), or None (an envelope that never passed through aggregate delivery, e.g. read back from the journal, or produced by a pre-stamp FCQRS).
(=): 'T -> 'T -> boolStructural equality The first parameter. The second parameter. The result of the comparison. 5 = 5 // Evaluates to true 5 = 6 // Evaluates to false [1; 2] = [1; 2] // Evaluates to true (1, 5) = (1, 6) // Evaluates to false
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
WaitAsync: TimeSpan -> Threading.Tasks.TaskGets a that will complete when this completes or when the specified timeout expires. The timeout after which the should be faulted with a if it hasn't otherwise completed. The representing the asynchronous wait. It may or may not be the same instance as the current instance.
Task: Threading.Tasks.TaskWait: unit -> unitWaits for the to complete execution. The has been disposed. The task was canceled. The collection contains a object. -or- An exception was thrown during the execution of the task. The collection contains information about the exception or exceptions.
System.TimeSpanRepresents a time interval.
FromSeconds: float -> TimeSpanReturns a that represents a specified number of seconds, where the specification is accurate to the nearest millisecond. A number of seconds, accurate to the nearest millisecond. is less than TimeSpan.MinValue or greater than TimeSpan.MaxValue. -or- is . -or- is . is equal to . An object that represents .
// A transfer ends when the target stores the money or the source gets it back.
static bool Finished(Data.IMessageWithCID message) =>
message is Event<AccountEvent>
{
EventDetails: TransferReceived or TransferRefunded
};
// Ask Alice's account to send money, and wait until the saga has finished.
async Task SendTransfer(string id, string target, decimal amount)
{
var cid = Values.NewCID();
// Subscribe first: the saga can finish before the reply arrives.
using var outcome = statement.SubscribeForFirst(cid, Finished);
var reply = await accounts(
_ => true, cid, alice, new SendTransfer(id, target, amount));
Show(reply);
if (reply.Journaled?.Value == true)
await outcome.Task.WaitAsync(TimeSpan.FromSeconds(30));
}
await SendTransfer("t1", "bob", 30m);
// Carol has no account, so this transfer comes back.
await SendTransfer("t2", "carol", 20m);
The reply to SendTransfer arrives when Alice's account has stored TransferSent, which is before
the transfer has finished. The program waits for the transfer's last event instead: TransferReceived
from the target, or TransferRefunded from Alice. It subscribes to the correlation ID on the
statement projection from step 4, with a filter for those two events, before it sends.
Deliver a transfer twice
// After a restart, a saga sends its last command again. Do the same by hand:
send bob (ReceiveTransfer("t1", "alice", 30m))
send: AggregateId -> AccountCommand -> unitbob: AggregateIdReceiveTransfer// After a restart, a saga sends its last command again. Do the same by hand:
await Send(bob, new ReceiveTransfer("t1", "alice", 30m));
Run it
From samples/accounts:
dotnet run --project 5-transfer-money/fsharp
dotnet run --project 5-transfer-money/csharp
Both programs print:
Opened for Alice (version 1, stored)
Deposited 100 (version 2, stored)
Opened for Bob (version 1, stored)
Sent 30 to bob (t1) (version 3, stored)
Sent 20 to carol (t2) (version 4, stored)
Received 30 from alice (t1) (version 2, not stored)
Statement for alice:
version entry amount balance
1 Opened for Alice 0 0
2 Deposit 100 100
3 Transfer t1 to bob -30 70
4 Transfer t2 to carol -20 50
5 Refund of transfer t2 20 70
Statement for bob:
version entry amount balance
1 Opened for Bob 0 0
2 Transfer t1 from alice 30 30
Transfer t1 moved 30 from Alice to Bob. Transfer t2 came back: Carol has no account, so her
account rejected ReceiveTransfer, and the saga refunded Alice, which is row 5 of her statement. The
second delivery of t1 got the first answer again without storing anything, and Bob's statement has
one row for it.
Between rows 4 and 5, Alice's statement shows 20 that has left her account and reached no one. A saga is not a transaction: a query during a transfer sees the money in flight.
Run it again
Rejected: The account is already open (version 5, not stored)
Deposited 100 (version 6, stored)
Rejected: The account is already open (version 2, not stored)
Rejected: Transfer t1 was already sent (version 6, not stored)
Rejected: Transfer t2 was already sent (version 6, not stored)
Received 30 from alice (t1) (version 2, not stored)
The program sends the same requests again. Only the deposit is new: the transfer IDs stop t1 and
t2 from sending money a second time, and the statements gain one row, Alice's deposit.
What remains the application's job
- Repeatable commands. Every command a saga sends can arrive more than once. Give it an ID the receiver remembers, as the transfer ID does here.
- A deadline for every wait. A state that waits for an answer needs
expecting, or its own timer, so a lost message cannot leave it waiting forever. - External systems. Stored saga states do not make a call to another system happen exactly once. A call to a payment provider needs its own idempotency key, timeout, retry policy, and a way to undo or escalate. Sagas covers recovery in detail, and Write a saga covers the other command targets and timers.
Next
Step 6: Add a memo adds a memo to transfers. TransferSent events are already
stored without one, so the step shows how to change an event that the journal already holds.