Skip to content

Akka.Persistence.Embedded: SQLite journal, snapshot store and read journal (AOT-compatible) - #8733

Closed
Aaronontheweb wants to merge 22 commits into
akkadotnet:devfrom
Aaronontheweb:feature/persistence-embedded
Closed

Aaronontheweb wants to merge 22 commits into
akkadotnet:devfrom
Aaronontheweb:feature/persistence-embedded

Conversation

@Aaronontheweb

@Aaronontheweb Aaronontheweb commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Adds Akka.Persistence.Embedded: a SQLite journal, snapshot store and read journal that read and write the same tables and rows as Akka.Persistence.Sql 1.5.70 on SQLite, so a database file can move between the two plugins. It uses plain ADO.NET on Microsoft.Data.Sqlite, has no linq2db or Sql.Common dependency, and publishes with Native AOT with zero warnings from its own code.

Builds on #8730; review commits after d189acc.

The TCK race fix also ships separately in #8734; it drops out of this diff once that merges.

Usage

// HOCON: which plugins are the default, and the database. No class names, no WithFallback.
var config = ConfigurationFactory.ParseString("""
    akka.persistence.journal.plugin = "akka.persistence.journal.embedded"
    akka.persistence.snapshot-store.plugin = "akka.persistence.snapshot-store.embedded"
    akka.persistence.journal.embedded.connection-string = "Data Source=app.db"
    akka.persistence.snapshot-store.embedded.connection-string = "Data Source=app.db"
    """);

var setup = BootstrapSetup.Create().WithConfig(config)
    .And(PersistenceSetup.Create().WithEmbeddedPersistence());   // journal + snapshot store + read journal, with their defaults

WithEmbeddedPersistence() registers all three plugins by their default ids, each with its reference configuration attached (#8730's PersistenceSetup). Pass eventAdapters: for the journal's adapters. WithEmbeddedJournal/SnapshotStore/ReadJournal(pluginId) register one plugin under another id.

Changes

  • New package Akka.Persistence.Embedded (net10.0, IsAotCompatible), plugin ids akka.persistence.{journal,snapshot-store,query.journal}.embedded.
  • Schema module builds the exact DDL text Akka.Persistence.Sql sends (byte-identical sqlite_master in all five modes), creates missing tables, verifies columns, never alters.
  • Journal: one writer thread, BEGIN IMMEDIATE batches, tag rows, tombstone deletes, delete-compatibility metadata like Akka.Persistence.Sql.
  • Snapshot store: one worker thread, update-then-insert upsert, same four-way criteria dispatch.
  • Read journal: all eight query interfaces, FromEnd, keyset-paged ids, optional gap tracking (journal-sequence-retrieval.enabled).
  • AOT canary src/aot/Akka.Persistence.Embedded.AOT.App in the AotCanary CI job: TagTable, Csv and Both databases, an event adapter, every query, and the unregistered failure. It keeps StreamsRoots.xml (Akka.Streams: island boundary subscribers use runtime generics and fail under Native AOT #8731) until Akka.Streams: build island-boundary subscribers without runtime generics (Native AOT) #8732 has merged.

Testing

  • Full TCK in TagTable, Csv, Both, delete-compat and no-writer-uuid modes, plus plugin tests for schema, row format, batching, deletes, queries, registration, forced connection reset and shutdown (410 passed, 5 skipped for documented reasons).
  • AOT canary: PublishAot on linux-x64, 0 warnings from plugin code, runs persist, recover, snapshot, delete and every query, ships libe_sqlite3.so.

Breaking changes

None. New package.

Deliberate differences from Akka.Persistence.Sql: D1 tags with ; come back intact, D2 live PersistenceIds polls on an interval, D3 gap tracking is opt-in, D4 snapshot timestamps are UTC, D5 delete pre-query runs in the delete transaction, D6 CurrentPersistenceIds is keyset-paged, D7 no blocking read journal constructor, D8 tag-read-mode = auto.

Deviations from the spec text: Buffer(1) goes before SelectMany in query streams so a finished query completes without extra demand (costs one batch of read-ahead); worker connections open with Pooling=False instead of calling ClearPool at shutdown, so a reset gets a new handle and nothing else's pooled handles are dropped; schema errors name the setting that needs each missing column.

Dependencies and licenses: Microsoft.Data.Sqlite and .Core 10.0.12 (MIT), SQLitePCLRaw 2.1.12 (Apache-2.0), SQLite (public domain). The gap tracker is ported from Akka.Persistence.Sql (Apache-2.0).

Add PersistencePluginSetup and PersistenceQuerySetup so a Native AOT or
trimmed app can start journals, snapshot stores, event adapters, bindings,
stash-overflow configurators and read journals without Type.GetType.

Every lookup site now runs Setup, then built-in, then a guard on
AkkaFeatures.IsDynamicTypeLoadingSupported, then reflection moved into
[RequiresUnreferencedCode] methods. With the switch on and nothing
registered, the same objects get built. With it off, an unregistered type
throws a ConfigurationException that names the HOCON setting.
src/aot/Akka.Persistence.AOT.App runs the in-memory journal, snapshot store
and read journal with Akka.DynamicTypeLoading off and everything registered
through PersistencePluginSetup and PersistenceQuerySetup. It persists,
recovers from a snapshot, runs CurrentEventsByPersistenceId and
CurrentEventsByTag, and checks that an unregistered journal fails at start.

The AotCanary CI job publishes it, runs it and compares its warnings with a
baseline. StreamsRoots.xml roots the Akka.Streams boundary types the
in-memory read journal's queries build by reflection.
Add a "Persistence plugins" section to the Native AOT page and a
"Registering your plugin for Native AOT" section to the custom persistence
provider guide.
The public EventAdapters.Create overload passed "event-adapters" as the
plugin path, so a failure read event-adapters.event-adapters.<name>. It now
passes an empty path, and the setting names read event-adapters.<name> and
event-adapter-bindings. Add a test that asserts the exact messages.
New package with the SQLite plugin's configuration (reference HOCON, validation,
rejected and ignored Akka.Persistence.Sql keys) and the schema module that builds
the same DDL text as Akka.Persistence.Sql, creates missing tables and verifies
columns without ever altering a table.
One writer thread owns the write connection and groups queued write requests into
a single BEGIN IMMEDIATE transaction. Recovery reads run on a small reader pool.
Rows, tag rows, tombstones and delete-compatibility metadata match what
Akka.Persistence.Sql writes. Adds the journal TCK in TagTable, Csv, Both,
delete-compat and no-writer-uuid modes plus plugin tests for row format, batching,
rejection versus failure, deletes, schema variants and shutdown.
One worker thread with one connection runs every operation in arrival order.
Saves are an update-then-insert upsert like Akka.Persistence.Sql, loads and
deletes use the same four-way criteria dispatch, and timestamps come back as UTC.
Adds the snapshot TCK specs and plugin tests for each load and delete case.
Implements all eight query interfaces over the journal tables with the same
offsets, sequence numbers, tags and Csv/TagTable matching rules as
Akka.Persistence.Sql, FromEnd offsets, keyset-paged persistence ids and optional
gap tracking. Streams use value-free unfold sources so they stay AOT-safe.
Adds the query TCK specs in every tag mode and plugin tests, and the API approval
file for the new assembly.
…ginSetup

WithEmbeddedPersistence() and WithEmbeddedReadJournal() let core build the journal,
snapshot store and read journal provider without reflection, which is how they
start when Akka.DynamicTypeLoading is off. Adds a test that boots the plugin that way.
Publishes with PublishAot on linux-x64, runs persist, snapshot, recovery, delete and
every query (including live ones) against SQLite with Akka.DynamicTypeLoading off,
and checks that the native library ships next to the executable. The warning
baseline is empty: the plugin's own code publishes with zero IL2xxx/IL3xxx warnings.
StreamsRoots.xml works around akkadotnet#8731 until Akka.Streams is fixed.
Two new headings in the Native AOT and custom persistence docs failed the
titlecase lint rule. The word "configurators" failed cspell, so the
sentence now says "stash overflow strategies".
The spec stopped an actor, waited for Terminated and recreated it under the same
name. The parent frees the name only after the watcher saw Terminated, so the
recreate failed with InvalidActorNameException in roughly one run in three.
…number from a snapshot

Custom table and column names pad like linq2db 5.4.1.9 does, and the highest
sequence number query with a lower bound (recovery from a snapshot) matches the
SQL Akka.Persistence.Sql 1.5.70 sends.
The persistence canary now registers a journal, a snapshot store and a stash
overflow configurator that Akka.Persistence does not ship and runs them with
Akka.DynamicTypeLoading off, next to the built-in inmem plugins. The comment
that called two sends a burst now says what the actor does.

Tests: add the unregistered snapshot store guard, switch-on tests for a
registered snapshot store, event adapter, binding and stash configurator, and
a run of PersistencePluginProxy with the switch off. Fold the exact-message
event adapter test into the existing one, keep Merge in its own test, and
make the read journal switch-on test assert the default config injection.

Docs: show how to add a read journal's default HOCON with the switch off, and
note in the baseline that Streams joins the gate once akkadotnet#8732 merges.
@Aaronontheweb Aaronontheweb added this to the 1.6.0 milestone Oct 2, 2026
@Aaronontheweb Aaronontheweb added akka-persistence AOT Ahead-of-Time (AOT) Compilation labels Oct 2, 2026

@Aaronontheweb Aaronontheweb left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review: what each change and test is for. Inline comments have the detail.

What this PR does

  • Adds Akka.Persistence.Embedded: a SQLite journal, snapshot store and read journal on plain Microsoft.Data.Sqlite, with no linq2db.
  • Writes the same tables, DDL text and rows as Akka.Persistence.Sql 1.5.70, so one .db file can move between the two plugins.
  • Runs one writer thread that batches writes into BEGIN IMMEDIATE transactions, with separate reader and query thread pools.
  • Registers through PersistencePluginSetup and PersistenceQuerySetup from PR #8730, so it starts with Akka.DynamicTypeLoading off.
  • Adds a Native AOT canary app and CI steps for publish, run and zero warnings from plugin code.
  • Adds tests: the TCK matrix (TagTable, Csv, Both, delete-compat, no writer uuid) plus plugin specs for schema, row format, journal, snapshots, queries, settings and shutdown.

Deliberate differences from Akka.Persistence.Sql

  • D1: TagTable tags use group_concat(tag, char(31)), so a tag containing ; stays whole.
  • D2: live PersistenceIds() polls every refresh-interval, with no busy loop.
  • D3: live and by-tag batches are bounded by MAX(ordering) read in the same transaction. Gap tracking is opt-in and off by default.
  • D4: snapshot timestamps load as DateTimeKind.Utc.
  • D5: the delete pre-query runs inside the delete transaction.
  • D6: CurrentPersistenceIds pages by persistence_id.
  • D7: each query awaits an async init handshake, not a blocking constructor.
  • D8: tag-read-mode = auto follows the write mode.
  • Not in the D list: Flatten puts Buffer(1) before SelectMany, so one batch is read ahead of demand.
  • Not in the D list: the read journal caches a ready flag, not the first handshake task that spec 9.8 asks for.

Risks to notice

  • Row and SQL parity is checked against captured APS text, not a live APS. The live cross-compat suite is PR 3.
  • Unverified claims: custom-name DDL padding (U1), ~ escaping in the Csv LIKE pattern (U2), and the highest > from SQL (U8).
  • The perf spec runs 300 events, not the 2000 in the spec, so it cannot catch a throughput regression.
  • SqliteReadJournal uses the default plugin path as a constant. Two read journals with gap tracking on and the same table name clash on the sequence actor name.
  • The canary passes only with StreamsRoots.xml (#8731). AOT users who run queries need the same file, and nothing ships or documents it.
  • ConnectionHolder.Reset may get the same pooled handle back, and no test forces that path.
  • Several tests prove writer behavior through test seams, not the full actor path.

Overlap

  • Plugin-specific specs overlap the TCK on query and journal basics. The TCK classes also overlap each other across modes by design.
  • AllocatesAllPersistenceIDsPublisher = false gates a TCK test that is already skipped, so those overrides do nothing.
  • The three docs-DDL classes run one scenario and differ only in mode.
  • This PR sits on #8730. Only the commits after f0c0f082 are reviewed here.

Not covered

  • Cross-compat with real Akka.Persistence.Sql files in both directions, mixed writers and a row diff against APS. This is PR 3.
  • Windows and macOS native packaging and AOT publish. CI publishes linux-x64 only.
  • A failed journal or snapshot init and its restart loop.
  • An unserializable snapshot on save.
  • Delete in Csv and Both mode.
  • A non-default tag-separator.
  • Snapshot row storage classes and journal_metadata row content.
  • Gap tracking under AOT, and its failure and backoff paths.
  • Event adapters, Both mode and delete-compat in the AOT canary.
  • A forced connection error that triggers a reset.
  • Two read journals in one system.
  • SQLITE_BUSY after the Default Timeout expires.

<PackageTags>$(AkkaPackageTags);persistence;eventsource;sql;sqlite;aot</PackageTags>
<GenerateDocumentationFile>true</GenerateDocumentationFile>
<Nullable>enable</Nullable>
<IsAotCompatible>true</IsAotCompatible>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We set IsAotCompatible so the trim and AOT analyzers run on this project. With repo-wide TreatWarningsAsErrors, any IL2xxx/IL3xxx fails the build, so the plugin carries no suppressions.

Comment thread Directory.Build.props
<GrpcToolsVersion>2.82.0</GrpcToolsVersion>
<BenchmarkDotNetVersion>0.15.8</BenchmarkDotNetVersion>
<MessagePackVersion>3.1.8</MessagePackVersion>
<MicrosoftDataSqliteVersion>10.0.12</MicrosoftDataSqliteVersion>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

New dependency: Microsoft.Data.Sqlite 10.0.12 (MIT). It pulls SQLitePCLRaw 2.1.x (Apache-2.0) and the native e_sqlite3. The PR body should carry the license table from spec section 2.4.

// ---- DDL text ------------------------------------------------------------------------------

/// <summary>J1, J2 and J3: the journal table plus its three indexes.</summary>
public static string JournalDdl(JournalSettings settings)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We build the DDL byte for byte like linq2db does for Akka.Persistence.Sql: bracketed names, tab-indented columns, padded types, ;\r\n between statements. Why: auto-initialize on a fresh file must give the same sqlite_master text, so both plugins can share the file. SchemaCreationSpec checks it against text captured from APS 1.5.70, not against a live APS (PR 3).

return sb.ToString();
}

private static void AppendCreateTable(StringBuilder sb, string table, List<Column> columns, int typePadding, string? constraint)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The padding rule (name width = longest bracketed name + 1, type width = longest type + typePadding) is our reading of linq2db. It is proven only for the default names. For custom names it is unverified (spec U1).

}

[Fact(DisplayName = "Should_not_stall_When_live_query_crosses_physically_deleted_rows")]
public async Task Should_not_stall_When_live_query_crosses_physically_deleted_rows()

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We delete rows so orderings 1 and 2 are gone, start a live AllEvents, write one event and expect it at ordering 4. It proves no stall at physical holes with tracking off (D3). The wait is 10 s, the same as APS's stall, so a regression may pass or flake instead of failing clearly.

}

[Fact(DisplayName = "Should_treat_percent_and_underscore_in_tag_literally_When_Csv")]
public async Task Should_treat_percent_and_underscore_in_tag_literally_When_Csv()

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Csv tags with %, _ and ~ must match literally, with case-insensitive match kept. The ~ assertions test our escape choice, not APS's behavior (spec U2).

}

[Fact(DisplayName = "Should_hold_back_at_gap_When_gap_detection_is_on")]
public async Task Should_hold_back_at_gap_When_gap_detection_is_on()

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With tracking on, we insert a row at ordering 5 with 3 and 4 missing. The live query must hold back for under query-delay * max-tries, then deliver it. It uses a 700 ms negative wait, so it depends on timing. It covers one path of the tracker.

The read journal uses its real plugin id for its settings section, logs and gap
tracker actor name, so two read journals with tracking on no longer clash. It caches
the first handshake task as the spec asks. Worker connections open with Pooling=False
so a reset really gets a new handle. Schema errors name the setting that requires the
missing column. A snapshot that cannot be serialized faults the save task.

Tests: forced connection reset, two read journals, non-default tag separator, delete
in Csv and Both mode, snapshot and metadata row storage classes, unserializable
snapshot, core's json serializer fallback, schema error text. Drops the no-op
AllocatesAllPersistenceIDsPublisher overrides.
Replace PersistencePluginSetup and PersistenceQuerySetup with one
PersistenceSetup that holds records keyed by plugin id, the way
SerializationSetup holds serializers. A registered plugin needs no HOCON
class, and its default config sits under the plugin's section, so switch-off
apps no longer add a read journal's DefaultConfiguration by hand.

- PersistenceSetup: Create(), Create(factory), WithJournal, WithSnapshotStore,
  WithPlugin, WithPlugins, WithStashOverflowStrategy, Merge.
- JournalDetails, SnapshotStoreDetails and EventAdapterDetails (an adapter
  with its bound event types); ReadJournalDetails and a WithReadJournal
  extension in Akka.Persistence.Query.
- Akka.Persistence.Hosting: WithPersistenceSetup builds on the setup earlier
  calls left.
- One internal registry per actor system serves persistence and the query
  extension. Type-name matching is gone.
- The canary, tests and docs use the new API.
PersistenceSetup.Create().WithEmbeddedPersistence() registers the journal, snapshot
store and read journal by plugin id, each with its reference configuration, so a user
writes no class names and no WithFallback. WithEmbeddedJournal, WithEmbeddedSnapshotStore
and WithEmbeddedReadJournal take a plugin id; the journal takes event adapters. The
read journal provider gets its plugin id from the registration.

The AOT canary uses the new API and now also runs an event adapter, a Both-mode
database with a delete, and keeps StreamsRoots.xml until akkadotnet#8732 has merged.
No empty catch blocks remain. The writer's CompleteAdding cannot throw, so its
try/catch is gone. Waiting for a pending write uses Task.WhenAny, which completes
without throwing the write's error. Closing a broken connection and a stopping
worker thread catch the specific exception and log it at Debug. The rollback failure
log now carries the exception. The fd scan in the shutdown test uses a TryRead helper,
and the canary reports a temp file it could not delete.
@Aaronontheweb

Copy link
Copy Markdown
Member Author

Superseded by #8739: the branch moved to akkadotnet so the two PRs can be a GitHub stack (#8738 → #8739), with each diff holding only its own changes. Review threads here are covered by the fixes already in that branch.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

akka-persistence AOT Ahead-of-Time (AOT) Compilation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant