Read your writes
Deposit 100 into Alice's account and immediately show her statement. With an asynchronous projection
the statement can still end before the deposit: the account has stored Deposited 100, but the
statement projection has not applied it yet. Nothing is broken, and a later query would show the
deposit. A response that must include the caller's own change cannot ship "refresh in a moment", so
read-your-writes closes the gap by waiting for the required projection before the query runs.
Motivation: Keep projections asynchronous for throughput and independence, then pay the waiting cost only for a request whose response must include its own change.
Read Correlation IDs and read-your-writes first if you need the mental model behind the sequence, projection boundary, and ephemeral notification.
To wait for every event already committed across the journal, use the projection's CatchUpAsync.
Catch up projections shows the snapshot boundary it waits for. A matching
correlation notification alone does not establish that boundary.
Use the combined F# helper
Fcqrs.sendAwaiting subscribes before sending, sends the command, and waits for one projection
notification when the aggregate reply was journaled:
Shared setup
open System.Threading
open FCQRS.Model.Data
open FCQRS.Common
open FCQRS.FSharp
open FCQRS.Query
module Account =
type AccountCommand =
| Open of owner: string
| Deposit of amount: decimal
| Withdraw of amount: decimal
type AccountEvent =
| Opened of owner: string
| Deposited of amount: decimal
| Withdrawn of amount: decimal
| Rejected of reason: string
open Account
let deposit (api: IActor) (handle: obj -> unit)
(accounts: AggregateHandle<AccountCommand, AccountEvent>)
(cid: CID) (alice: AggregateId) = async {
SystemThreadingFCQRSModelFCQRS.Model.DataFCQRS.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.QueryCorrelation subscriptions and the notifications projections publish to them.
Fcqrs_450-how-to_005-read-your-writes.md_page.AccountFcqrs_450-how-to_005-read-your-writes.md_page.Account.AccountCommandOpenowner: stringstringAn abbreviation for the CLI type . Basic Types
Depositamount: decimaldecimalAn abbreviation for the CLI type . Basic Types
WithdrawFcqrs_450-how-to_005-read-your-writes.md_page.Account.AccountEventOpenedDepositedWithdrawnRejectedreason: stringdeposit: IActor -> (obj -> unit) -> AggregateHandle<AccountCommand,AccountEvent> -> CID -> AggregateId -> Async<Event<AccountEvent>>api: IActorFCQRS.Common.IActorDefines the core functionalities and context provided by the FCQRS environment to actors. This interface provides access to essential Akka.NET services and FCQRS initialization methods.
handle: obj -> unitobjAn abbreviation for the CLI type . Basic Types
unitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
accounts: AggregateHandle<AccountCommand,AccountEvent>FCQRS.FSharp.AggregateHandle`2What you get back after registering an aggregate.
cid: CIDFCQRS.Model.Data.CIDCorrelationID for commands and Sagas. A CID cannot contain '~'. FCQRS joins a CID to an aggregate ID with '~' to name sagas and pub-sub topics, and reads the CID back from the text after the last '~', so a CID containing one never reaches its saga. Creating such a CID throws ArgumentException; the constructor throws instead of returning an error so that FCQRS releases built against the previous signature keep working. Deserialization does not run this check, so IsValid reports it for a CID read from a message.
alice: AggregateIdFCQRS.Model.Data.AggregateIdasync: AsyncBuilderBuilds an asynchronous workflow using computation expression syntax. let sleepExample() = async { printfn "sleeping" do! Async.Sleep 10 printfn "waking up" return 6 } sleepExample() |> Async.RunSynchronously
// The statement projection, which callers can subscribe to.
let statement = Fcqrs.projection api (Projection.single FromStart handle)
// The last argument selects the account's reply: Deposited or Rejected.
let! reply =
Fcqrs.sendAwaiting statement accounts cid alice (Deposit 100m)
(fun _ -> true)
// On return, the statement projection has handled the deposit. Query it now.
statement: FCQRS.Projections.IProjectionFCQRS.FSharp.Fcqrsprojection: IActor -> Projection -> FCQRS.Projections.IProjectionRegister the read-model projection and return the subscription stream. Starts a projection. It follows each aggregate's and saga's own sequence numbers, so it never skips a stored event. Each pass hands the events it found to the handler in journal order, so each aggregate's events arrive in sequence order; on SQLite, events of different aggregates arrive in the order they were written, and on PostgreSQL an event that commits late arrives in a later pass. With `Named` progress, a handler can see an event again after a crash. A handler that throws terminates the process. Requires a SQLite or PostgreSQL journal; an Akka event adapter must turn each journal row into one event.
api: IActorFCQRS.FSharp.ProjectionModuleConstructors for the projection-handler shapes.
single: ProjectionProgress -> (obj -> unit) -> ProjectionSingle-event handler: just update the read model (returns unit); each aggregate event is then published to subscribers as-is. The common case when every event is worth notifying.
FromStartIn memory: the projection reads the whole journal each time it starts, then follows new events. For a read model the process keeps in memory.
handle: obj -> unitreply: Event<AccountEvent>sendAwaiting: 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>cid: CIDalice: AggregateIdDeposit// C# composes the same subscribe-before-send sequence explicitly.
using var projected = statement.SubscribeForFirst(cid);
// The first argument selects the account's reply: Deposited or Rejected.
var reply = await accounts(_ => true, cid, alice, new Deposit(100m));
if (reply.Journaled?.Value != false)
await projected.Task.WaitAsync(TimeSpan.FromSeconds(30));
// The statement projection has now handled the deposit. Query it.
The helper waits for one notification. If a command persists a batch and the projection publishes
several events for the same CID, either filter notifications so only the final required update is
published or compose a subscription with the correct take count.
The F# wait is bounded: if no matching notification arrives within akka.fcqrs.command-timeout,
sendAwaiting raises TimeoutException. The default is 30s. A bare number means seconds, and HOCON
durations such as 500ms also work. A projection that suppresses or filters out the matching event therefore surfaces
as a timeout instead of hanging the request. See
Configuration.
Why "only if journaled"
An aggregate can persist an event or defer a reply. A deferred rejection or idempotent response is returned to the caller but never enters the journal, so no projection will receive it.
FCQRS stamps the delivered envelope with Event.Journaled : bool option:
Some true: the event was stored and can reach a projection;Some false: the reply was deferred, likeRejected, or publish-only and will not reach a projection;None: the envelope predates or bypassed the delivery stamp.
sendAwaiting skips the projection wait for Some false. The C# sequence performs the equivalent
Journaled check explicitly.
Compose the sequence manually
Use the explicit form when waiting for several notifications, adding cancellation, or applying a notification filter:
Shared setup
return reply
}
let sendManually (statement: ISubscribe)
(accounts: AggregateHandle<AccountCommand, AccountEvent>)
(cid: CID) (alice: AggregateId) (command: AccountCommand)
(cancellationToken: CancellationToken) = async {
reply: Event<AccountEvent>sendManually: ISubscribe -> AggregateHandle<AccountCommand,AccountEvent> -> CID -> AggregateId -> AccountCommand -> CancellationToken -> Async<unit>statement: ISubscribeFCQRS.Query.ISubscribeThe canonical subscription stream: a non-generic shorthand for ISubscribe<IMessageWithCID> — the type every FCQRS projection / read-your-writes subscription actually uses (cf. IEnumerable vs IEnumerable<T>). Lets consumers write ISubscribe instead of the closed generic, and inject it by that name.
accounts: AggregateHandle<AccountCommand,AccountEvent>FCQRS.FSharp.AggregateHandle`2What you get back after registering an aggregate.
Fcqrs_450-how-to_005-read-your-writes.md_page.Account.AccountCommandFcqrs_450-how-to_005-read-your-writes.md_page.Account.AccountEventcid: CIDFCQRS.Model.Data.CIDCorrelationID for commands and Sagas. A CID cannot contain '~'. FCQRS joins a CID to an aggregate ID with '~' to name sagas and pub-sub topics, and reads the CID back from the text after the last '~', so a CID containing one never reaches its saga. Creating such a CID throws ArgumentException; the constructor throws instead of returning an error so that FCQRS releases built against the previous signature keep working. Deserialization does not run this check, so IsValid reports it for a CID read from a message.
alice: AggregateIdFCQRS.Model.Data.AggregateIdcommand: AccountCommandcancellationToken: CancellationTokenSystem.Threading.CancellationTokenPropagates notification that operations should be canceled.
async: AsyncBuilderBuilds an asynchronous workflow using computation expression syntax. let sleepExample() = async { printfn "sleeping" do! Async.Sleep 10 printfn "waking up" return 6 } sleepExample() |> Async.RunSynchronously
use awaiter = statement.Subscribe(cid, 1, cancellationToken = cancellationToken)
let! reply = accounts.Send cid alice command (fun _ -> true)
if reply.Journaled <> Some false then
do! awaiter.Task |> Async.AwaitTask
// Query the statement.
awaiter: IAwaitableDisposablestatement: ISubscribeSubscribe: CID * int * (IMessageWithCID -> unit) option * CancellationToken option -> IAwaitableDisposableSubscribes to events matching a specific correlation ID. The correlation ID to match. Maximum number of events to process. Optional callback function to handle the event. An optional cancellation token to cancel the subscription.
cid: CIDcancellationTokencancellationToken: CancellationTokenreply: 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: AggregateIdcommand: AccountCommandJournaled: 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 inequality The first parameter. The second parameter. The result of the comparison. 5 <> 5 // Evaluates to false 5 <> 6 // Evaluates to true [1; 2] <> [1; 2] // Evaluates to false
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
Task: Tasks.Task(|>): '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
AwaitTask: Tasks.Task -> Async<unit>Return an asynchronous computation that will wait for the given task to complete and return its result. The task to await. If an exception occurs in the asynchronous computation then an exception is re-raised by this function. If the task is cancelled then is raised. Note that the task may be governed by a different cancellation token to the overall async computation where the AwaitTask occurs. In practice you should normally start the task with the cancellation token returned by let! ct = Async.CancellationToken, and catch any at the point where the overall async is started. Awaiting Results
Shared setup
}
The ordering is part of correctness. Subscribing after .Send creates a race in which the projection
can publish before the subscription exists.
Subscribe registers the listener before returning. Disposing or cancelling an awaitable subscription
before it receives the requested number of notifications cancels its task; disposal does not report
that the projection has caught up.
using var awaiter = statement.SubscribeForFirst(cid);
var reply = await accounts(_ => true, cid, alice, command);
if (reply.Journaled?.Value != false)
await awaiter.Task.WaitAsync(cancellationToken);
// Query the statement.
Wait for the right projection
A notification means that the projection publishing it has completed its handler. It says nothing about another projection with its own progress or deployment: a statement notification does not mean a monthly report built by another projection includes the deposit. If a response depends on several read models, wait for a completion signal representing all of them.
Subscriptions are in-memory rendezvous points, not durable messages for disconnected clients. Create the subscription as part of the active request, and decide how the API reports a projection that does not catch up in time.
The timeout story differs by API. The F# sendAwaiting helper is bounded by
akka.fcqrs.command-timeout (default 30s) and raises TimeoutException. Raw Subscribe awaiters and
the C# SubscribeForFirst awaiter are not bounded by that key. Bound them with WaitAsync in C#,
as the examples on this page do, or with a cancellation token in F#.