Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 56 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -1355,10 +1355,65 @@ reloads an inline `Snapshot<T>`.
projection statically does not help: `JasperFxSingleStreamProjectionBase`'s constructor compiles the
wrapper's accessors with FastExpressionCompiler, which throws there. That was measured, with a
`Snapshot<T, TId>` overload built and then removed because it could not achieve its purpose. The
factory refuses by name, naming #942. The rest is fisher#412.
factory refuses by name, naming #942. It is the one item of fisher#412 still open.
- **The smoke consumer references `JasperFx.Events.SourceGenerator` itself**, because a project
reference does not carry the analyzer the package bundles. Its version is kept in step by hand.

### Native AOT: the daemon, LINQ and step-through — fisher#412

`smoke/aot-consumer` now runs the async daemon over an async `Snapshot<T>` and a multi-stream
projection, walks a keyset cursor, `Include()`s by a strong-typed id, orders by full-text relevance,
compares a `DateOnly`, and runs both projection step-through methods — all natively.

- **The async daemon needed no Fisher change.** It ran natively the first time it was measured. A
run with only the cursor fix reverted fails at the cursor walk, which comes after the daemon
section, so the daemon passing is a fact rather than an inference. It is built by hand in the
smoke; `AddAsyncDaemon()` reaches the same `BuildProjectionDaemonsAsync`.
- ⚠️ **The keyset cursor is a fixed-shape `Utf8JsonWriter` and stays byte-identical to the old
`JsonSerializer` output.** A cursor goes to a client and comes back to a possibly newer deployment,
and Polecat shares the format. `cursor_encoding` pins the exact string. It was captured from the old
implementation over every key type the provider reads back: string with escapes, every integer
width, float, double, decimal, bool, null, `DBNull`, Guid, `DateTimeOffset`, `DateTime`, enum,
`DateOnly`, `TimeOnly`, `TimeSpan` and char. A key type outside that list is refused by name rather
than guessed at. Decoding is `JsonDocument`.
- **A strong-typed identity is carried as its inner value**, which fixed a bug nobody had filed. The
reflection serializer wrote the wrapper as `{"Value":…}`, which `ConvertSlot` could never bind
back, so the second page of any keyset walk over a wrapper-keyed document failed.
`a_cursor_walk_over_a_strong_typed_identity` fails against the old code.
- **`Include()` builds `Enumerable.Contains<object>` over an `object[]`**, closed statically, where it
used to close `Contains` and `Array.CreateInstance` over the member's runtime type. The predicate is
translated and never compiled, and `EnumerableContains` strips the `Convert` around the member.
⚠️ This one was a real native failure, not just a warning: over a strong-typed id it threw
"`TicketId[]` is missing native code". A Guid member happened to work, because something else
closes those instantiations.
- **The marker methods are closed through a delegate** (`new Func<…>(OrderByRelevance).Method`), which
is how `System.Linq.Queryable` does it. `OrderByRelevance` and its siblings, and `Include`'s marker,
used to look up the method by name and close it with `MakeGenericMethod`. Over reference types that
happened to work natively. The change removes the warning, and the smoke does not tell the two
versions apart.
- **`DateMember` renders through `options.GetTypeInfo(type)`**, which asks the configured resolver.
It produces the same text as before. The old call worked natively too, because the application's
context covers a document's member types, so this removes the warning only.
- ⚠️ **Step-through renders every state through the store's serializer, a visible change.**
`RunProjectionByNameAsync` used `JsonSerializer`'s default options. A console therefore saw
`Loaded` where the store persists `loaded`, and the call threw natively. It now reads the store's
`Stream` overloads, the unannotated contract every document read already uses. It also reads no
step property by reflection: a generic hop returns `ProjectionTimelineRaw` directly.
`replay_by_name_renders_state_through_the_stores_serializer` pins agreement with the typed path.
A consumer that keyed on the PascalCase names sees camelCase, so this belongs in the release notes.
- ⚠️ **The step-through's `MakeGenericMethod` lives in a non-generic helper, because ILC 10.0.1 crashed
on the old spelling.** `MakeGenericMethod(typeof(TState))` inside the generic interface method made
ILC throw `IndexOutOfRangeException` from `MakeGenericMethodSite.InstantiateDependencies` as soon as
the method was reachable. So no native publish of an application calling either step-through method
produced a binary. `InvokeReplay(string, Type, …)` hands the analyzer a plain `Type`, and it warns
instead. Do not inline it back. The remaining `MakeGenericMethod` is inherent (the interface leaves
`TState` unconstrained, or names it only as a `Type`). It carries a method-level suppression, and the
smoke runs it natively.
- **What #412 did not touch:** strong-typed aggregate ids (jasperfx#942, above), and the IL warnings
outside the five paths the issue named. `dotnet build` still lists about twenty against Fisher
(`MessagePublishing`, `AdvancedSqlResultReader`, `SecondaryStoreProxy`, `EnumMember` and others).
ILC reports some of them as reachable from the smoke, and none of them failed when it ran.

### Document write SQL

`SqliteDocumentStorageDescriptorBuilder` emits four statements whose **column order and `?` order are
Expand Down
2 changes: 1 addition & 1 deletion HANDOFF.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ equivalent for and never will.
[CLAUDE.md](CLAUDE.md) has the architecture and the SQLite traps. This document is the compliance
scoreboard and the things that are true right now but not obvious from either.

**2575 tests green on net9.0 and net10.0** — 2508 in `Fisher.Tests`, 36 in
**2580 tests green on net9.0 and net10.0** — 2513 in `Fisher.Tests`, 36 in
`Fisher.AspNetCore.Tests` and 31 in `Fisher.EntityFrameworkCore.Tests`. 669 of
them are shared cross-store compliance tests — 509 event sourcing and 160 document.
On JasperFx **2.79.2** / Weasel **9.40.0**.
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ process**. There is no server to install, nothing to provision, and nothing to k
> `Fisher.AspNetCore` and `Fisher.EntityFrameworkCore` companion packages all work and are tested.
>
> Fisher passes **all 58 suites and 669 tests** it enrolls from `JasperFx.Events.ComplianceTests`,
> the shared cross-store suite Marten and Polecat also enroll in, alongside its own 2,508.
> the shared cross-store suite Marten and Polecat also enroll in, alongside its own 2,513.
>
> That suite pins **API portability, not behavioural equivalence** — code written against one store
> compiles and runs against another. It does not pin that the three behave identically, and they do
Expand Down
21 changes: 15 additions & 6 deletions docs/configuration/native-aot.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,11 @@

Fisher's document storage works in a Native AOT image (`PublishAot=true`): storing, loading, upserting
and LINQ queries, over Guid, string, `int` and `long` identities, strong-typed id wrappers and document
hierarchies. So does the core of the event store: starting and appending to streams, live aggregation,
`FetchForWriting`, and an inline `Snapshot<T>`. A smoke application is published natively and run in
Fisher's CI on every change.
hierarchies. So does the event store: starting and appending to streams, live aggregation,
`FetchForWriting`, inline snapshots, and the async daemon running async snapshots and multi-stream
projections. So do keyset paging (`ToCursorPageAsync`), `Include()`, full-text search with relevance
ordering, and projection step-through. A smoke application exercising all of these is published
natively and run in Fisher's CI on every change.

Three things are different from a JIT application. Each has to be configured, because the reflection
that works them out under the JIT isn't available in a native image.
Expand Down Expand Up @@ -82,8 +84,15 @@ compiles the wrapper's accessors with FastExpressionCompiler, which throws there
([jasperfx#942](https://github.com/JasperFx/jasperfx/issues/942)). Fisher refuses such an aggregate by
name rather than failing inside JasperFx. Key it on a Guid, string, `int` or `long` until that ships.

## Async projections and the daemon

Async projections need nothing extra. Name the projected document types in the JSON context, as you
do for any document. The CI smoke builds the daemon with `BuildProjectionDaemonAsync()`, starts it,
and waits for it to catch up, all in the native image.

::: warning What has not been measured
The CI smoke covers document storage and the event-store paths above, through `AddFisher`. Async
projections, the async daemon, and the LINQ operators that still serialize through reflection (cursor
paging, `Include`, full-text extracts) have not been run in a native image.
The CI smoke covers the paths above, through `AddFisher`. It does not cover the hosted daemon
(`AddAsyncDaemon()`), although that starts the same daemon. Nor does it cover subscriptions, event
publishing from projections, raw SQL (`AdvancedSql`), or a second store registered with
`AddFisherStore<T>`. Fisher's build still reports trimming and AOT warnings on some of those paths.
:::
5 changes: 4 additions & 1 deletion smoke/aot-consumer/AotConsumer.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,10 @@
fisher#384 / #386. A document write in a Native AOT image: insert, load, upsert and a LINQ query,
over all four canonical identity types, strong-typed id wrappers and a document hierarchy — and since
fisher#398, the event store: starting and appending to a stream, live aggregation,
FetchForWriting, and an inline Snapshot<T>. Every one of those works under CoreCLR, so no test project
FetchForWriting, and an inline Snapshot<T> — and since fisher#412, the async daemon projecting an
async snapshot and a multi-stream projection, keyset paging, Include() over a strong-typed id,
full-text relevance ordering, a DateOnly comparison and projection step-through. Every one of
those works under CoreCLR, so no test project
here can see an AOT failure; this one is published with PublishAot and RUN, which is the only
way to find out. The two failures #384 reported both "published clean": ILC printed only an
assembly-level IL2104/IL3053 rollup and the first sign was a crash on the first write.
Expand Down
187 changes: 187 additions & 0 deletions smoke/aot-consumer/Program.cs
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
using System.Text.Json;
using System.Text.Json.Serialization;
using Fisher;
using Fisher.Linq;
using Fisher.Linq.Includes;
using JasperFx;
using JasperFx.Descriptors;
using JasperFx.Events;
using JasperFx.Events.Projections;
using Microsoft.Extensions.DependencyInjection;

Expand Down Expand Up @@ -31,6 +35,14 @@
// fisher#398. An inline snapshot is a projection closed over (aggregate, id) — reflective unless
// the registration names both while they are still generic arguments.
options.Projections.Snapshot<Voyage>(SnapshotLifecycle.Inline);

// fisher#412. The async daemon: an async snapshot and a multi-stream projection, both projected by
// a daemon running in the native image.
options.Projections.Snapshot<Ledger>(SnapshotLifecycle.Async);
options.Projections.Add(new PortVisitsProjection(), ProjectionLifecycle.Async);

// fisher#412. The LINQ paths that used to serialize or close generics by reflection.
options.Schema.For<Note>().FullTextIndex(x => x.Body);
});

await using var provider = services.BuildServiceProvider();
Expand Down Expand Up @@ -130,6 +142,115 @@
Expect(snapshot is { Port: "Tromso", Legs: 3 }, "the inline snapshot was written and reloads");
}

// ---- the async daemon (fisher#412) ----

var ledger = Guid.NewGuid();

await using (var session = store.LightweightSession())
{
session.Events.StartStream<Ledger>(ledger, new Deposited(100), new Withdrawn(30));
await session.SaveChangesAsync();
}

var daemon = await store.BuildProjectionDaemonAsync();
try
{
await daemon.StartAllAsync();
await daemon.WaitForNonStaleData(TimeSpan.FromSeconds(60));
}
finally
{
await daemon.StopAllAsync();
daemon.Dispose();
}

await using (var query = store.QuerySession())
{
Expect((await query.LoadAsync<Ledger>(ledger))?.Balance == 70, "the async snapshot was projected by the daemon");

var visits = await query.LoadAsync<PortVisits>("Bergen");
Expect(visits?.Count == 1, "the multi-stream projection was projected by the daemon");
Expect((await query.LoadAsync<PortVisits>("Tromso"))?.Count == 1, "the multi-stream projection grouped by port");
}

// ---- LINQ and the step-through (fisher#412) ----

var boat = Guid.NewGuid();

// The full-text index's FTS5 table and triggers are created by the migration and not by the
// on-demand path that creates a document table at first write, so the schema is applied here.
await store.ApplyAllConfiguredChangesToDatabaseAsync();

await using (var session = store.LightweightSession())
{
session.Store(new Boat { Id = boat, Name = "Belle" });

for (var i = 1; i <= 5; i++)
{
session.Store(new Catch
{
Id = Guid.NewGuid(), Weight = i, BoatId = boat, Landed = new DateOnly(2026, 8, i)
});
}

session.Store(new Escalation { Id = Guid.NewGuid(), TicketId = ticket.Id });
session.Store(new Note { Id = Guid.NewGuid(), Body = "corrosion on the hull" });
session.Store(new Note { Id = Guid.NewGuid(), Body = "corrosion corrosion corrosion everywhere" });
await session.SaveChangesAsync();
}

await using (var query = store.QuerySession())
{
// Keyset paging: the cursor payload used to be JsonSerializer over an object?[].
var weights = new List<int>();
string? cursor = null;
do
{
var page = await query.Query<Catch>().OrderBy(x => x.Weight).ThenBy(x => x.Id)
.ToCursorPageAsync(2, cursor);
weights.AddRange(page.Items.Select(x => x.Weight));
cursor = page.NextCursor;
} while (cursor is not null);

Expect(weights.SequenceEqual([1, 2, 3, 4, 5]), "a keyset cursor walk covers every row once");

// A DateOnly comparison value is rendered through the store's serializer.
var late = await query.Query<Catch>().Where(x => x.Landed >= new DateOnly(2026, 8, 4)).ToListAsync();
Expect(late.Count == 2, "a DateOnly comparison finds the documents");

// Include() closed Enumerable.Contains over the member's runtime type.
var boats = new List<Boat>();
var catches = await query.Query<Catch>().Include(x => x.BoatId, boats).ToListAsync();
Expect(catches.Count == 5 && boats.Count == 1 && boats[0].Name == "Belle", "Include fetches the related document");

// ...over a member whose type is a value type nothing else closes Contains over.
var related = new List<Ticket>();
await query.Query<Escalation>().Include(x => x.TicketId, related).ToListAsync();
Expect(related.Count == 1 && related[0].Subject == "printer", "Include fetches by a strong-typed identity");

// OrderByRelevance closed its own marker method over T by name.
var ranked = await query.Query<Note>().Where(x => x.Search("corrosion")).OrderByRelevance().ToListAsync();
Expect(ranked.Count == 2 && ranked[0].Body.StartsWith("corrosion corrosion"), "a full-text query ranks by relevance");
}

// Projection step-through rendered state with JsonSerializer's default options.
var records = new[] { (object)new Departed("Hull"), new Arrived("Bergen") }
.Select((body, index) => new EventRecord(Guid.NewGuid(), index + 1, index + 1, voyage.ToString(),
store.Options.EventGraph.EventMappingFor(body.GetType()).EventTypeName,
JsonDocument.Parse(store.Options.Serializer.ToJson(body)).RootElement, null,
DateTimeOffset.UtcNow, null, null))
.ToList();

var timeline = await ((IEventStore)store).RunProjectionByNameAsync(nameof(Voyage), voyage, records, null,
CancellationToken.None);
Expect(timeline.Steps.Count == 2 && timeline.FinalState?.GetProperty("legs").GetInt32() == 2,
"projection step-through renders each state through the store's serializer");

var typedTimeline = await ((IEventStore)store).RunProjectionAsync<Voyage>(nameof(Voyage), voyage, records,
null, CancellationToken.None);
Expect(typedTimeline.Steps.Select(x => x.After?.Legs).SequenceEqual([1, 2]),
"typed projection step-through copies the state at every step");

Console.WriteLine("OK: documents and events written and read in a native image.");
return 0;
}
Expand Down Expand Up @@ -227,6 +348,72 @@ public void Apply(Arrived arrived)
}
}

public record Deposited(decimal Amount);

public record Withdrawn(decimal Amount);

public class Ledger
{
public Guid Id { get; set; }
public decimal Balance { get; set; }

public static Ledger Create(Deposited deposited) => new() { Balance = deposited.Amount };

public void Apply(Deposited deposited) => Balance += deposited.Amount;

public void Apply(Withdrawn withdrawn) => Balance -= withdrawn.Amount;
}

public class PortVisits
{
public string Id { get; set; } = "";
public int Count { get; set; }
}

public partial class PortVisitsProjection : Fisher.Projections.MultiStreamProjection<PortVisits, string>
{
public PortVisitsProjection()
{
Identity<Arrived>(x => x.Port);
}

public void Apply(Arrived _, PortVisits visits) => visits.Count++;
}

public class Boat
{
public Guid Id { get; set; }
public string Name { get; set; } = "";
}

public class Catch
{
public Guid Id { get; set; }
public int Weight { get; set; }
public Guid BoatId { get; set; }
public DateOnly Landed { get; set; }
}

public class Escalation
{
public Guid Id { get; set; }
public TicketId TicketId { get; set; }
}

public class Note
{
public Guid Id { get; set; }
public string Body { get; set; } = "";
}

[JsonSerializable(typeof(Deposited))]
[JsonSerializable(typeof(Withdrawn))]
[JsonSerializable(typeof(Ledger))]
[JsonSerializable(typeof(PortVisits))]
[JsonSerializable(typeof(Boat))]
[JsonSerializable(typeof(Catch))]
[JsonSerializable(typeof(Note))]
[JsonSerializable(typeof(Escalation))]
[JsonSerializable(typeof(Departed))]
[JsonSerializable(typeof(Arrived))]
[JsonSerializable(typeof(Voyage))]
Expand Down
Loading
Loading