Skip to content

The search box knows all the secrets -- try it!

Fisher is part of the Critter Stack ecosystem.

JasperFx Logo JasperFx provides formal support for Fisher and other Critter Stack libraries. Please check our Support Plans for more details.

Asynchronous Projections ​

The async daemon runs projections in the background, off the request path.

cs
builder.Services.AddFisher(opts =>
{
    opts.Connection("Data Source=app.db");
    opts.Projections.Snapshot<Order>(SnapshotLifecycle.Async);
    opts.Projections.Add(new SalesByCustomer(), ProjectionLifecycle.Async);
})
.ApplyAllDatabaseChangesOnStartup()
.AddAsyncDaemon(DaemonMode.Solo);

Or build one by hand:

cs
var daemon = await store.BuildProjectionDaemonAsync();
await daemon.StartAllAsync();

DANGER

DaemonMode.HotCold is refused, and it is a real limitation rather than an omission. Hot-cold failover means several nodes competing for a leadership lease through the database, and a Fisher store is a file SQLite does not make safe to share across nodes. Accepting the mode and running Solo would give an application the opposite of the guarantee it asked for — every node projecting at once.

WAL matters here ​

WARNING

WAL is what lets the daemon read while a session writes. It is on by default; if you turn it off, Fisher logs a warning at startup rather than refusing to run, because a non-WAL store projects correctly — it just serialises the daemon against every writer, which presents as a slow projection rather than as a misconfiguration.

The health check is how an operator finds out the warning mattered.

The high-water mark is simply max(seq_id) ​

Marten and Polecat must distinguish the highest sequence issued from the highest safe to read, because a PostgreSQL sequence or a SQL Server IDENTITY hands out numbers outside the transaction — a writer can hold 7 uncommitted while 8 commits ahead of it.

On SQLite, one writer per file plus BEGIN IMMEDIATE means a transaction's sequences fully commit before the next writer allocates any, and a rollback returns the number. Committed sequences are contiguous, so there is no separate answer to give.

WARNING

If you are extending Fisher: do not reintroduce gap-skipping. It would guard a state that cannot occur.

Liveness ​

cs
opts.Events.HighWaterLivenessInterval = TimeSpan.FromSeconds(5);   // 0 turns it off

The mark's row moves when the mark advances, which is a different question from whether the loop is running — a quiet store advances nothing and would otherwise be indistinguishable from a dead daemon. So the high-water agent re-stamps its row on an idle cycle, and that is the only liveness signal there is.

WARNING

The extended progression heartbeat column does not answer it. JasperFx returns early for the high-water shard, so nothing ever writes that column for that row — and a health check reading it looks like it has a signal it does not have.

TIP

It is throttled, where Marten's per-tenant equivalent writes on every cycle. That difference is SQLite's: a write takes the file's one write lock, so touching at the slow-polling interval would make an otherwise read-only store a permanent 1 Hz writer with a WAL to checkpoint.

The batch is atomic ​

Each batch commits the projection's document writes and the progression row in one transaction.

DANGER

Splitting them lets a crash between the two either replay events already applied or skip events never applied — permanently, and with nothing to signal it.

Sessions are collected rather than merged, and each flushes its own operations into the shared transaction, because an operation is configured against a session as its storage context and that is what carries tenancy.

Unknown event types ​

The daemon does not skip an event whose type it cannot resolve, where a stream read does. Silently skipping one would leave the projection permanently wrong.

cs
opts.Projections.Errors.SkipUnknownEvents = true;   // if you really want that

Otherwise it throws, and the exception is classified as a shard failure without the daemon needing to know Fisher's exception types.

Unreadable event bodies ​

A body the serializer cannot read is a different problem from a type it cannot resolve — a data or serializer fix rather than a deployment one — and it has its own policy:

cs
opts.Projections.Errors.SkipSerializationErrors = true;   // the default

With it on, the row is quarantined into fi_dead_letters and the shard carries on. With it off the loader throws Fisher.Exceptions.EventDeserializationFailureException — which derives from JasperFx.Events.EventDeserializationFailureException and carries the sequence, the stored event type alias and the serializer's own exception — and the shard pauses classified EventSerialization rather than Other.

WARNING

A missing IEventBinarySerializer is deliberately not covered by this policy. That is a misconfiguration of the whole store rather than one unreadable row, so it is refused by name whatever the flag says — quarantining it would turn every binary event in the store into a dead letter and bury the one thing worth saying.

Errors and dead letters ​

cs
opts.Projections.Errors.SkipApplyErrors = true;

A skipped poison event is quarantined into fi_dead_letters rather than stopping the shard.

TIP

The dead letter write goes on its own connection, outside the failing batch's transaction — that batch is about to roll back, and a dead letter written inside it would roll back with the very failure it is recording.

Read them back:

cs
foreach (var db in eventStore.AllDatabases())
{
    var letters = await db.QueryDeadLetterEventsAsync(…);
}

Rebuilds ​

cs
var daemon = await store.BuildProjectionDaemonAsync();
await daemon.RebuildProjectionAsync("SalesByCustomer", token);
await daemon.RebuildProjectionAsync<Order>(token);

Teardown clears the projection's documents and its progression rows in one transaction — doing one without the other replays a projection on top of rows it already wrote.

WARNING

Test a rebuild with a row the replay cannot recreate. A replay rewrites every row it can still produce, so a broken teardown is invisible against live data.

Two projections publishing one document type ​

Teardown clears the whole of every table the named projection publishes into. Two projections that publish the same document type share one fi_doc_* table, so rebuilding either one deletes both projections' rows — and the rebuild then replays only the projection it was asked for. The rebuild succeeds, the rebuilt projection is correct, and the other read model is left empty.

Marten behaves the same way. Fisher warns about it at rebuild time, naming the table and the other projections that write into it:

warn: Rebuilding projection LandedTally will clear 1 table(s) that another registered projection
      also publishes into: fi_doc_tally is also published by ReleasedTally. ...

DANGER

Rebuilding the other projection afterwards does not fix it. There is no operation that rebuilds two projections together — every RebuildProjectionAsync overload names one — so a second rebuild clears the shared table again and discards what the first one wrote. Only the projection rebuilt last keeps its rows.

Rewind instead. RewindSubscriptionAsync replays a projection onto the rows that are already there rather than clearing first:

cs
await daemon.RewindSubscriptionAsync("ReleasedTally", token, sequenceFloor: 0);

Sharing a published type between two projections is legal and costs nothing until somebody rebuilds, so Fisher warns rather than refusing.

Reaching the running daemon ​

AddAsyncDaemon() registers a JasperFx IProjectionCoordinator, which is how application code gets at the daemons the host is running:

cs
var coordinator = services.GetRequiredService<IProjectionCoordinator>();

var daemon = coordinator.DaemonForMainDatabase();
await daemon.WaitForNonStaleData(TimeSpan.FromSeconds(30));

// Under database-per-tenant, one daemon per file.
var forTenant = await coordinator.DaemonForDatabase("tenant-a");
var all = await coordinator.AllDaemonsAsync();

PauseAsync() stops every agent without disposing anything, and ResumeAsync() restarts them — which is what a maintenance window or a test fixture wants:

cs
await coordinator.PauseAsync();
// ... do the thing the daemon must not be running for ...
await coordinator.ResumeAsync();

TIP

The coordinator is the hosted service, registered once and resolved from both places. Two registrations would be two daemons over one file, which on SQLite is two writers contending for the single write lock.

TIP

Pausing stops the tenant poller too. It builds and starts daemons of its own, so leaving it running would let a tenant appearing mid-pause quietly begin projecting.

WARNING

The coordinator's daemons exist only once the host has started — it is a hosted service. Resolving it while the container is still being built and calling DaemonForMainDatabase() throws saying so. DaemonMode.Disabled and DaemonMode.ExternallyManaged register no coordinator at all.

Resetting data under a running daemon ​

Advanced.ResetAllDataAsync() pauses a daemon this process is hosting, wipes, and resumes it.

Without that the wipe strands the daemon: the delete takes fi_event_progression out from under agents holding their positions in memory, so they carry on from where they were, record nothing against an event store that now starts at zero, and every later WaitForNonStaleData times out saying shards have recorded no progress. Silent until something waits — which for the spec suite that reported it (#138) was the next scenario.

WARNING

Only a daemon this process is hosting is paused, because that is the only one the store knows about. DaemonMode.ExternallyManaged, a store built with DocumentStore.For(...), or a daemon in another process all keep the hazard — pause them yourself around the wipe.

TIP

This is a deliberate divergence from Marten, which leaves its daemon alone here. The reason to take it is that the alternative was unreachable rather than merely manual: before #138 there was no way to get at the running daemon at all, so the caller this method overwhelmingly has — a test fixture resetting between scenarios — could not have paused it by hand.

Catching up and waiting ​

cs
await daemon.WaitForNonStaleData(TimeSpan.FromSeconds(5));

Or from a query:

cs
await session.Query<CustomerSales>().QueryForNonStaleData(TimeSpan.FromSeconds(5)).ToListAsync();

Non-stale means every registered async shard has reached the current high-water mark. A shard that has not started yet counts as behind — it has no progression row, and a row is evidence about a shard rather than the definition of one. So the wait blocks until it reports, and a store with no async projections returns immediately.

WARNING

Non-stale does not imply a post-commit listener has run. The progression row is written inside the batch's transaction, so non-stale is true the moment that commits — strictly before any listener. Wait on the listener's own signal instead.

WARNING

Waiting without a running daemon throws TimeoutException. Nothing will advance a registered shard, so there is no honest early return; the message names the shards that have recorded nothing so "never started" reads differently from "still catching up". A progression row for a projection that is no longer registered is ignored rather than waited on — that is what DeleteProjectionProgressByShardNameAsync is for.

Event-emitting projections ​

A projection can raise events, and they are appended inside the batch's own transaction.

Two SQLite-shaped decisions in that:

  • The version comes from a read under the write lock. A slice pre-assigns versions client-side from its own event count, which is only the stream's real version when the projection has seen every event on it. Fisher re-reads inside the batch's BEGIN IMMEDIATE and the optimistic guard runs there, so a projection raising events onto a stream another writer has moved on fails the batch instead of writing a wrong version.
  • Raised events are taggable. They are routed through Fisher's own append operation, which is the only thing that supplies the sequence a tag row is keyed by. Queuing bare per-event operations would silently make raised events untaggable.

WARNING

Polecat no-ops the equivalent members rather than throwing, so an event-raising projection there drops its events with no signal. Fisher does not.

Multi-tenancy ​

Under database-per-tenant the daemon runs one instance per tenant database:

cs
var all = await store.BuildProjectionDaemonsAsync();
var one = await store.BuildProjectionDaemonAsync("acme");

AddAsyncDaemon() hosts them all. N daemons over N files do not contend, which is the same property that makes that tenancy a performance feature.

WARNING

The no-argument BuildProjectionDaemonAsync() projects the default file and says nothing about the others.

What Fisher supplies, and what it does not ​

The daemon itself is JasperFx's — coordinator, subscription agents, shard tracker, throttled and resilient loaders, roughly 10,500 lines. What Fisher supplies is the storage seam: progress reads and writes, the high-water detector, the event loader, the projection batch, and the session and shard plumbing. That is why a projection ports between the three stores unchanged.

Released under the MIT License.