Write a saga
This recipe coordinates a money transfer across two accounts. Alice's account stores TransferSent,
which takes the money out of her balance. Bob's account decides whether to accept it. A saga delivers
the money to the target account and, if the target rejects it, refunds the source.
Transfer money builds and runs this saga step by step. This page
lists what every saga needs, using the same code.
Read Sagas first if you need the ground-up explanation of transitions,
SagaStartingEvent, the starter handshake, and recovery re-drive.
Motivation: Use a saga here because the transfer crosses two independent owners and the conversation must survive a process restart.
Write the state table first
| Current state | Incoming event | Next state | Command issued |
|---|---|---|---|
| not started | TransferSent from the source |
Delivering |
ReceiveTransfer to the target |
Delivering |
TransferReceived from the target |
Completed |
none; stop saga |
Delivering |
Rejected from the target |
Refunding |
RefundTransfer to the source |
Refunding |
TransferRefunded from the source |
Completed |
none; stop saga |
Delivering or Refunding |
no reply before the deadline | the same state | the same command |
The implementation has one function for the first three columns and another for the last column.
Motivation: Writing the table first exposes missing outcomes and accidental loops before routing, persistence, or language syntax can hide them.
Map incoming events to persisted states
Each state carries what its command needs:
Shared setup
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 }
// 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
module Transfer =
open System
open FCQRS.Common
open FCQRS.FSharp
open Account
Fcqrs_450-how-to_007-write-a-saga.md_page.AccountFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
Fcqrs_450-how-to_007-write-a-saga.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_450-how-to_007-write-a-saga.md_page.Account.AccountEventOpenedDepositedWithdrawnTransferSentTransferReceivedTransferRefundedRejectedreason: stringFcqrs_450-how-to_007-write-a-saga.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 [ ].
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
id: 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
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.
(-): ^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"
Fcqrs_450-how-to_007-write-a-saga.md_page.TransferSystemFCQRS.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_450-how-to_007-write-a-saga.md_page.Transfer.TransferDetailsId: stringstringAn abbreviation for the CLI type . Basic Types
Source: stringTarget: stringAmount: decimaldecimalAn abbreviation for the CLI type . Basic Types
Fcqrs_450-how-to_007-write-a-saga.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;
handleEvent receives events from every participant as obj. Match the typed envelope and current
saga state together. The state is None when the starting event first reaches user code. Sender
identifies the account that stored an event, which separates the target's rejection from any other.
// 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_450-how-to_007-write-a-saga.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_450-how-to_007-write-a-saga.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()
};
StateChangedEvent next stores the next saga state. UnhandledEvent rejects an event that does not
belong in the current state. Do not issue commands from this function; state persistence must complete
first. The ExpectationExhausted case answers a missed deadline, described
below.
In C#, the two cases of the state option are two methods. Start receives a message before the saga
has a state, normally the event that started it, and returns the first state. HandleEvent receives
the later messages with the stored state. Both take object deliberately: a saga also receives other
aggregates' reply events and ToSelf timeout payloads.
Map persisted states to commands
applySideEffects runs after the state is stored and again after recovery. It returns commands plus
the transition FCQRS should make after issuing them. In C#, ApplySideEffects runs only once a user
state exists, so it receives the state directly.
// 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_450-how-to_007-write-a-saga.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);
The command helpers select a target. C# has the same helpers on SagaCommands:
toOriginator factory command: the exact aggregate instance whose event started the saga;toAggregate factory id command: another aggregate instance selected by id;toActor actorRef command: an arbitrary actor reference;toSelf command: a message back to the saga itself, which arrives inhandleEvent;toOriginatorAfter factory delayMs taskName command,toAggregateAfter, andtoSelfAfter: delayed variants for timeout or retry behaviour. AtoSelfAfterreminder is a hand-rolled saga timeout: enter a state, schedule a wake-up, and lethandleEventdecide whether it still matters.
The returned saga transition means:
Stay: keep waiting in the current state after sending commands;StayExpecting expectation: keep waiting, with a deadline and retry schedule for the wait, asexpectingreturns (next section);NextState next: persist another state immediately and run its side effects;StopSaga: send any returned commands, then complete and passivate. Delayed commands returned withStopSagaare still delivered, excepttoSelfAfterones. Those are cancelled with a warning, so a completed saga cannot resurrect itself.
Give every wait a deadline
A waiting state can hang forever for two reasons: the command produced no reply event (the target
decided IgnoreEvent), or a message was lost between nodes. expecting in F#, or Expect with
Expectations.Create in C#, declares what the wait expects and what happens when it does not arrive.
The transfer's deliver helper sends its command on state entry, sends exactly that command again every
5 seconds while no state transition is persisted, and after 30 seconds delivers an
ExpectationExhausted message to handleEvent. The handler must answer it with a transition.
The rules that make this safe:
- The deadline is measured from the persisted state-entry time, not from a timer. A restart or crash loop re-arms the schedule from the journal and cannot postpone the deadline.
- Re-sent commands must be retry-safe. This is the same contract recovery re-drives already impose; the expectation adds no new obligation on the target aggregate.
- A timeout means the outcome is unknown, not failed. The reply may still arrive after the saga
escalated, so the escalated state's
handleEventshould decide what a late success means instead of ignoring it. - An unhandled
ExpectationExhausted(no matching case, or a handler exception) is logged as an error and re-delivered one deadline period later. The framework never invents a terminal state. DeadlineandRetryEveryare explicit; there are no defaults. Size the deadline above worst-case shard handoff plus journal latency, and preferBackoff(which applies jitter) when many sagas can wait on the same aggregate.- One expectation per state. A state waiting on several aggregates with different deadlines should be split into one state per wait.
The hand-rolled equivalent, toSelfAfter with an attempt counter in the state, remains valid and
shows exactly what the framework automates.
Exhaustion has two answers: escalate or renew
Escalating to a failure or compensation state is correct while the workflow can still change its mind. Some waits cannot fail. Once a saga has persisted a decision that other aggregates may already have acted on, the only correct behaviour is to keep delivering that decision until every participant has confirmed it.
The transfer is such a wait. The money left Alice's account when her account stored TransferSent. A
timeout in Delivering does not show whether Bob's account stored TransferReceived: the command may
have arrived and only the reply been lost. Refunding on a timeout could then pay the money twice, once
to Bob and once back to Alice. The saga therefore answers exhaustion by entering the same state again,
as the ExpectationExhausted case in handleEvent does.
A self-transition persists a new state entry, which re-anchors the deadline and re-arms the schedule: the saga retries forever, but in journaled cycles. Each renewal is a durable event operators can alert on, so a participant that never recovers shows up as a growing trail of renewals instead of a silent hang. This is the deliberate blocking behaviour of a decided workflow made observable, not a bug.
Choose per state:
- Escalate when a timeout can still resolve the workflow: report a rejection, compensate, or release a hold. A wait before any money moves, such as a fraud check before the source account sends the transfer, belongs here.
- Renew when the state represents a decision already made. Never abort after the decision; alert on repeated renewals and fix the participant instead.
Declare the start event
StartOn answers “which originator event creates one new instance of this saga?” Match only the event
that begins the workflow.
// 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_450-how-to_007-write-a-saga.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;
Originator supplies the aggregate factory used by the starting handshake and by toOriginator.
InitialData supplies fixed data available to the saga functions. Current workflow progress belongs in
the state-machine cases; use unit when no additional fixed data is needed. In C#, the saga class
declares SagaName, InitialData, Originator, and StartsOn, and AddSaga registers the class.
Do not construct SagaStartingEvent yourself. FCQRS creates and stores that runtime envelope from the
event accepted by StartOn.
A saga starts once per correlation ID and originator aggregate. A command that would start it again
with the same correlation ID stores nothing and fails with SagaAlreadyStartedException, so send each
command that can start a saga with a new correlation ID.
Correlation IDs
explains when that happens.
Motivation: The start rule lets FCQRS subscribe the saga before the originator publishes the one event that begins the workflow. Without that handshake, the new saga could miss its first event.
Register the saga and starter rules
Register participant aggregates before constructing the saga, then wire every saga start rule once. Here the accounts are both the originator and the target, so the saga needs one factory:
Shared setup
open Microsoft.Extensions.Configuration
open Microsoft.Extensions.Logging
open FCQRS.Actor
open FCQRS.Common
open FCQRS.FSharp
open Account
let connectionString = "Data Source=accounts.db"
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 }
MicrosoftExtensionsConfigurationLoggingFCQRSFCQRS.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_450-how-to_007-write-a-saga.md_page.AccountconnectionString: stringlogging: ILoggerFactoryMicrosoft.Extensions.Logging.LoggerFactoryProduces instances of classes based on the given providers.
Create: System.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_450-how-to_007-write-a-saga.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);
wireSagaStarters is not optional. It installs the predicates and the safe-start handshake that
subscribes a new saga before the originator publishes its starting event. The C# host builder wires
the saga starter from all registered sagas at startup, so C# has no counterpart to call.
Make recovery commands safe
After reconstructing saga state, FCQRS invokes applySideEffects with recovering = true. Delivery of
the previous command is uncertain, so each waiting state must do one of the following:
- resend an idempotent command;
- query an external operation by a stable idempotency key;
- issue a recovery-specific reconciliation command;
- move to an explicit failed or manual-resolution path.
Motivation: Recovery repeats the next intended action because the journal can prove the stored state, but it cannot prove whether an outgoing message crossed the process boundary before failure.
In this example, each account records the IDs of the transfers it has received and refunded. A
repeated ReceiveTransfer or RefundTransfer for a known ID returns the first answer as a deferred
reply and moves no money. The normal commands are therefore safe to issue again, after recovery or on
each retry. Transfer money sends a repeated delivery and shows the
result.
Do not return no command merely because recovering is true. If the process stopped before delivery,
that leaves the workflow waiting forever. Add a timeout for every event that may never arrive.
Use Test your domain to test the event-to-state and state-to-command functions independently.