Repository navigation
Akka.Persistence.Embedded: SQLite journal, snapshot store and read journal (AOT-compatible) - #8739
Conversation
551eae7 to
e2c6943
Compare
Aaronontheweb
left a comment
There was a problem hiding this comment.
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 ADO.NET (Microsoft.Data.Sqlite). It aims to write the same tables and rows as Akka.Persistence.Sql 1.5.70. - Adds
Akka.Persistence.Embedded.Hosting, the only way to register the plugin. It uses #8738'sCreatePluginActorFactoryandWithReadJournal, so nothing needsAkka.DynamicTypeLoading. SqliteSchemabuilds the DDL text APS sends, creates missing tables and checks columns. It never alters a table.RowCodecandJournalSqldefine the row format: storage classes,manifestnever NULL, onecreatedvalue per write call, one writer uuid per journal instance.JournalWriteris one thread, one connection,BEGIN IMMEDIATEbatches.SqliteWorkerPoolgives reads and the snapshot store their own threads. SQLite's one-writer rule lives in one place.- The read journal has all eight query interfaces,
FromEnd, and an opt-in gap tracker (JournalSequenceActor). - Tests: the TCK in five modes (thin wrappers), plugin specs for schema, row format, journal, snapshot, queries, lifecycle and settings, nine Hosting specs, and an AOT canary app in the
AotCanaryjob.
Public API (Hosting)
WithEmbeddedPersistence(connectionString, ...)adds journal, snapshot store and read journal as the defaults.mode,tagWriteMode,pluginIdentifierandisDefaultPluginwork likeWithSqlPersistence.- A second overload takes
EmbeddedJournalOptions?,EmbeddedSnapshotOptions?andEmbeddedReadJournalOptions?. WithEmbeddedJournal,WithEmbeddedSnapshotStoreandWithEmbeddedReadJournaladd one part each.EmbeddedJournalOptionsandEmbeddedSnapshotOptionsderive fromJournalOptionsandSnapshotOptionsand overrideCreatePluginActorFactory.EmbeddedReadJournalOptionsstands alone and exposes three settings:RefreshInterval,MaxBufferSize,JournalSequenceRetrievalEnabled.- Parameter order and names differ from
WithSqlPersistenceand from the spec's string overload (callbacks second,configureSnapshotlast,readJournalOptionsthird in the options overload). See the inline comments. - Public API beyond the spec's Hosting section:
EmbeddedReadJournalOptions, the three single-part methods,ReadThreads, and the constructorSqliteReadJournalProvider(system, config, pluginPath). The API approval files lock them. - No breaking change, so
BREAKING_CHANGES_V1.6.mdis untouched.
Deliberate differences from Akka.Persistence.Sql
- D1: tags with
;come back whole.QuerySqlaggregates withgroup_concat(tag, char(31)). Test:Should_return_tags_containing_semicolon_intact_When_TagTable. - D2: live
PersistenceIdspolls everyrefresh-interval.LiveIdsStep. Only the TCKPersistenceIdsSpeccovers it, so nothing checks that it does not spin. - D3: batches are bounded by
MAX(ordering)read in the same transaction (ReadMaxAndBatch). The tracker is opt-in. Tests: the stall spec and the gap spec. - D4: snapshot
Timestampis UTC.Should_return_utc_timestamp_When_loading. - D5: the delete pre-query runs inside the delete transaction (
RunDelete). No test can show the race it removes. - D6:
CurrentPersistenceIdsis keyset-paged (ReadPersistenceIdsPage).Should_page_persistence_ids_When_more_than_max_buffer_size. - D7: no blocking constructor. Queries await
EnsureInitializedAsync.Should_wait_for_journal_initialization_When_query_starts_first. - D8:
tag-read-mode = auto.Should_resolve_tag_read_mode_from_write_mode_When_set_to_auto, plus the Both-mode query specs. - Also in the PR body:
Buffer(1)beforeSelectMany,Pooling=Falseinstead ofClearPool, schema errors that name the setting.
Risks to notice
- One failed insert fails every request in its batch, including other persistence ids. APS does the same; a test pins it.
Pooling=Falseoverrides aPooling=Truethe user set, silently.- The Csv
LIKEescape of~is unverified against linq2db. Our test only shows self-consistency. - DDL padding for custom names is our reading of linq2db, checked against hand-captured text.
JournalWriterhas a static test seam and a hook in the production loop.FindPluginPathmatches config text, so the HOCONclasspath can fall back to the default id with no warning.- The write path can commit after core's 10 s circuit breaker reports failure, because
Default Timeoutis 30 s. The spec notes it; the code does not warn. PostStopblocks a dispatcher thread for up to 5 s per thread group while it joins.- Windows and macOS Native AOT are not checked; the canary runs on linux-x64 only.
Overlap
- Builds on #8738 (
PluginActorFactory,WithReadJournal, persistence setups). It replaces #8733, whose plugin review still applies. - Same schema and settings names as Akka.Persistence.Sql, so a config copied across mostly works.
table-mappingand Sql.Common keys fail on purpose. - The TCK wrappers cover behavior the plugin specs also touch (
FromEnd, deletes, tags). The plugin specs add APS-specific numbers. OptionsSpecandShould_write_the_options_into_config_...check some of the same HOCON.- Cross-compat with real APS files, docs and migration pages are PR 3.
Not covered
- Any real Akka.Persistence.Sql database file, in either direction. Expected values in these tests come from capture notes, not from running APS.
- Reads and writes with custom table and column names. Only DDL text and HOCON are tested.
- Multi-page replay: no test sets a small
replay-batch-size. orderingthat is not an INTEGER primary key, andauto-initialize = trueover docs-DDL tables.WithEmbeddedJournal,WithEmbeddedSnapshotStore,PersistenceMode.Journal,tagWriteMode: CsvandautoInitialize: falsethrough Hosting.- That a later
WithEmbeddedReadJournalcall replaces the oneWithEmbeddedPersistenceadded. - NULL
serializer_idon a snapshot row, WAL mode, two processes on one file, and a crash mid-transaction. - Hosting specs run on the JIT with the switch off. Only the canary shows ILC is clean.
| // ---- DDL text ------------------------------------------------------------------------------ | ||
|
|
||
| /// <summary>J1, J2 and J3: the journal table plus its three indexes.</summary> | ||
| public static string JournalDdl(JournalSettings settings) |
There was a problem hiding this comment.
We build the CREATE text by hand, byte for byte, because a file that moves between this plugin and Akka.Persistence.Sql must keep the same sqlite_master text. ExpectedSchemas pins it for default names in five modes. The ;\r\n between statements never reaches sqlite_master, so no test can see it.
| return sb.ToString(); | ||
| } | ||
|
|
||
| private static void AppendCreateTable(StringBuilder sb, string table, List<Column> columns, int typePadding, string? constraint) |
There was a problem hiding this comment.
We pad columns with one rule (name width = longest + 1, type width = longest + 1, or + 2 for the journal). It matches linq2db for default names. For custom names it is our reading of linq2db, checked only against text captured by hand. A miss changes the sqlite_master text, not behaviour.
| } | ||
|
|
||
| var (bytes, manifest, identifier) = SerializePayload(payload, _settings.DefaultSerializer); | ||
| var tagArray = tags.ToArray(); |
There was a problem hiding this comment.
We write tags in the set's enumeration order, as APS does. String hashes are random per process, so ;red;blue; and ;blue;red; can both occur, and tag rows insert in that order. A row-level comparison with APS has to compare tags as sets.
| var tagArray = tags.ToArray(); | ||
| return new JournalRow | ||
| { | ||
| Created = representation.Timestamp == 0 ? batchTicks : representation.Timestamp, |
There was a problem hiding this comment.
We use the call's single tick value for created unless the representation already has a timestamp. Akka clears that timestamp before the journal sees an actor's write, so the second branch only runs for direct callers.
| /// A plain Microsoft host with Akka.Hosting, no TestKit: the TestKit itself loads types by name, so it cannot run | ||
| /// with <c>Akka.DynamicTypeLoading</c> off. Flips the switch around the host's life. | ||
| /// </summary> | ||
| internal sealed class HostedSystem : IAsyncDisposable |
There was a problem hiding this comment.
We use a plain Microsoft host and flip Akka.DynamicTypeLoading around its life, because the Akka TestKit loads types by name. The switch is process-wide, so these classes run in a non-parallel collection. This proves the registration path needs no type loading on the JIT. Only the AOT canary proves ILC.
| .WithHealthCheck(), | ||
| configureSnapshot: snapshot => snapshot.WithHealthCheck()); | ||
|
|
||
| [Fact(DisplayName = "Should_persist_recover_snapshot_and_query_by_tag_When_dynamic_type_loading_is_off")] |
There was a problem hiding this comment.
We persist, snapshot, recover and query by tag through a Hosting-built system with the switch off. The tag comes from an adapter registered with a factory, so it proves AddWriteEventAdapter(factory) and WithEmbeddedPersistence need no type names.
d81340c to
d1734ae
Compare
There was a problem hiding this comment.
... Would hit 'approve' but the comment I left outside of this review means I selected what I selected ;)
Edit: My comment outside of this review was incorrect after digging deeper. This looks good to me. Below is one question, one low hanging fruit that is easy to solve before merge, and possible future improvements.
| /// would wait forever. Cost: one batch is read ahead of demand, and a stream that is cancelled may already have | ||
| /// issued one more query. Spec 9.9 does not list this operator; it is a deliberate deviation. | ||
| /// </remarks> | ||
| public static Source<TElem, NotUsed> Flatten<TElem>(this Source<IReadOnlyList<TElem>, NotUsed> batches) |
There was a problem hiding this comment.
Curious as to whether this exists in Persistence.Sql and I didn't notice, or if not why it's not needed there.
| transaction.Dispose(); | ||
| } | ||
|
|
||
| foreach (var request in batch) |
There was a problem hiding this comment.
We could re-run the other requests alone after a constraint error.
FWIW, The general justification in Persistence.Sql (Inherited from pekko justification) is that 'if you are failing a write on one actor, you probably want everything to fail' since it could just be adding more questionable (may be related, may be unrelated) state. (The whole journal winds up shutting down when that happens, right?)
| foreach (var tag in row.Tags) | ||
| { | ||
| tagParameters["@ordering_id"].Value = ordering; | ||
| tagParameters["@tag"].Value = tag; | ||
| tagParameters["@sequence_nr"].Value = row.SequenceNr; | ||
| tagParameters["@persistence_id"].Value = row.PersistenceId; | ||
| _insertTag.ExecuteNonQuery(); | ||
| } |
There was a problem hiding this comment.
Fair. It's SQLite and p/invoke is probably improved enough we don't need to really batch for tags here.
| foreach (var request in batch) | ||
| { | ||
| foreach (var row in request.Rows) | ||
| InsertRow(row); | ||
| } |
There was a problem hiding this comment.
It's probably fine, I still have to wonder whether for cases of non-tag-table if a lot of records going in would benefit from more proper batching or not due to p/invoke and marshalling of params or not.
(If nothing else, something for someone to explore in future. I get this is MVP-ish)
Edit: I am impressed by p/invoke improvements in .NET, don't think we need to worry here.
6eb0b58 to
1b3ccb1
Compare
1b3ccb1 to
0a24fc9
Compare
Aaronontheweb
left a comment
There was a problem hiding this comment.
I'm trimming this plugin down - we over-ported a ton of the compat bits from Akka.Persistence.Sql and IMHO that's not really necessary.
- Stick with tag tables only
- WriterUuid always enforced
- Delete compatibility mode and metadata table
0a8f43f to
e3dc090
Compare
…urnal (AOT-compatible) A new SQLite-only plugin that writes the same tables and rows as Akka.Persistence.Sql 1.5.70 on SQLite, so either plugin can open the other's database. Plain ADO.NET on Microsoft.Data.Sqlite, no linq2db. Full persistence TCK in every tag mode, plus an AOT canary.
…osting New package with EmbeddedJournalOptions, EmbeddedSnapshotOptions and EmbeddedReadJournalOptions, and AkkaConfigurationBuilder.WithEmbeddedPersistence. The options supply the plugin factories through Akka.Persistence.Hosting's CreatePluginActorFactory, and the read journal registers with WithReadJournal, so nothing needs a HOCON class and the plugin starts with Akka.DynamicTypeLoading off. The plugin no longer exposes PersistenceSetup extensions. The AOT canary builds its systems with Akka.Hosting and WithEmbeddedPersistence, drops StreamsRoots.xml and gates Akka.Streams and the Hosting package at zero warnings.
…om names and paging, harden pooling and provider path Hosting: WithEmbeddedPersistence now has the same overloads and parameter order as WithSqlPersistence. The read journal's settings are properties of EmbeddedJournalOptions, so customizing it needs no second registration. Removes WithEmbeddedJournal, WithEmbeddedSnapshotStore, WithEmbeddedReadJournal and EmbeddedReadJournalOptions. Plugin: keep an explicit Pooling setting and log once when forcing it off; the read journal provider refuses to guess its plugin path; a CrashForTests message replaces the WriteFinished(null) trick. Tests: custom table and column names end to end (Both with delete compatibility and Csv), multi-page replay, Hosting surface, pooling, provider path, ordering primary key error, tighter gap timing.
…nnection strings The static CreatedForTests seam matched the writer's connection string against Db.ConnectionString. The writer gets the HOCON spelling (forward slashes), Db.ConnectionString has the OS spelling, so the match missed on Windows and eight journal behavior tests plus the pooling test failed there. The journal actor now answers GetWriterForTests with its writer, so no string is compared.
… struct QuerySql builds every statement variant in its constructor (by tag or not, tag table or Csv, first or next page), so a query only binds parameters. NumericRangeEntry is a readonly struct: it is only held in an ImmutableList, so it needs no object per entry.
EmbeddedJournalOptions derives from JournalOptions<SqliteWriteJournal, SqliteReadJournalProvider> and EmbeddedSnapshotOptions from SnapshotOptions<SqliteSnapshotStore>, so WithJournal and WithSnapshot register the plugin types. The PluginActorFactory overrides and the explicit WithReadJournal call are gone, and the typed AddWriteEventAdapter<T> replaces the factory overload. The read journal id follows the journal's identifier, as the base class defines it. A default plugin under another identifier no longer also answers to akka.persistence.query.journal.embedded, because the base registers one read journal.
… SQLite one The plugin now writes and reads only the layout Akka.Persistence.Sql creates on SQLite by default. The other options came from the obsolete Sql.Common plugins. Removed: - TagWriteMode (public), Csv and Both tag modes, tag-separator, tag-read-mode. Tags always go to the tag table; writer_uuid is always written. - delete-compatibility-mode, the journal_metadata table and its SQL. - use-writer-uuid-column, custom column names, schema-name, table-mapping, table-compatibility-mode, and the code that accepted, ignored or rejected Akka.Persistence.Sql keys. Unknown keys are not read. - warn-on-auto-init-fail, the serializer fallback, the read journal's own connection-string, journal-sequence-retrieval and JournalSequenceActor. - max-concurrent-queries, query-throttle-timeout and write-plugin-init-timeout are internal constants with the old defaults. - The matching Hosting options and the tagStorageMode, deleteCompatibilityMode, useWriterUuidColumn and maxConcurrentQueries parameters. Changed: - Flat HOCON keys: table-name and tag-table-name on the journal, table-name on the snapshot store. - The delete tombstone now gets an empty message. The highest sequence number still comes from MAX(sequence_number), and queries skip deleted rows. - Tests: the Csv, DeleteCompat and NoWriterUuid suites and the column-name and schema-variant tests are gone. The TagTable TCK suites lost their prefix. New tests cover the tombstone and tag queries.
Flatten now buffers one element after SelectMany instead of one batch before it. The next query runs once the current batch is drained, and a finished query still completes without extra demand. This matches Akka.Persistence.Sql, which gets the same read-ahead from SelectAsync(1, ...) and ConcatMany. The doc comment now says the buffer is load-bearing and what it costs.
Akka.Persistence.Sql looks up the writer with FindSerializerForType(type, serializer): a serialization binding wins, otherwise the named serializer replaces the System.Object fallback. Do the same for journal events and snapshots. Reads still use the stored serializer id and manifest. The reference HOCON gets `serializer = null` back, and the Hosting README no longer says the option is ignored.
TestData/aps-compat holds a SQLite file that Akka.Persistence.Sql 1.5.70 wrote with its default settings: three persistence ids, string and record events, tags in the tag table, one snapshot, and two tombstones. The generator source and a README (versions, regeneration, reverse check) sit next to it as .txt files so the test project does not compile them. ApsCompatSpec copies the file, starts Embedded on it, and checks that Embedded leaves the existing tables untouched, replays every id with the right highest sequence number, answers the tag and all-events queries with the same offsets, and loads the snapshot. Then it writes, deletes and restarts, and reads old and new data back.
…_uuid non-null on write The Embedded canary now links src/aot/Shared/LogWatchdogFilter.cs like the other canaries. JournalRow.WriterUuid is required and non-null because every write sets it. RawJournalRow.WriterUuid stays nullable: rows written by Akka.Persistence.Sql may lack it.
e3dc090 to
dda6bdb
Compare
Stacked on #8738 (Hosting-based plugin registration). Replaces #8733.
Adds
Akka.Persistence.Embedded: a SQLite journal, snapshot store and read journal that read and write the tables Akka.Persistence.Sql 1.5.70 creates on SQLite by default, so a database file can move between the two plugins. It supports that one layout only. 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.Akka.Persistence.Embedded.Hostingis the way to use it.Usage
That adds the journal, snapshot store and read journal as the default plugins, creates the tables, and needs no HOCON, no
classsetting and noWithFallback. It works withAkka.DynamicTypeLoadingoff (Native AOT).WithEmbeddedPersistencehas the same overloads and parameter names asWithSqlPersistence(connection string, options objects, delegates), minus the linq2db ones. Event adapters and health checks go injournalBuilderandsnapshotBuilder; call it again with anotherpluginIdentifierfor a second database. The read journal comes with the journal: its settings areQuery*properties ofEmbeddedJournalOptions, so nothing registers twice. The options classes derive from #8738'sJournalOptions<SqliteWriteJournal, SqliteReadJournalProvider>andSnapshotOptions<SqliteSnapshotStore>, so Hosting registers the plugin types. The read journal id isakka.persistence.query.journal.{identifier}.Layout and settings
tagstable,writer_uuidis always written, there is nojournal_metadatatable and no CSV tags column. Column names are fixed; table names are settings.deleted = 1, emptymessage) and deletes lower rows and their tag rows. The highest sequence number isMAX(sequence_number)over all rows, deleted ones included. Queries and replay skip deleted rows, tag queries too.connection-string,auto-initialize,table-name/JournalTableName,tag-table-name/TagTableName,buffer-size,batch-size,replay-batch-size,read-threads. Snapshot store:connection-string,auto-initialize,table-name/TableName. Both also readserializer/Serializer: a serialization binding wins, otherwise the named serializer replaces the System.Object fallback (as in Akka.Persistence.Sql). Read journal:write-plugin,max-buffer-size,refresh-interval,query-threads. Other keys are not read.tagscolumn) opens.Changes
Akka.Persistence.Embedded(net10.0,IsAotCompatible), plugin idsakka.persistence.{journal,snapshot-store,query.journal}.embedded.Akka.Persistence.Embedded.Hosting:EmbeddedJournalOptions(journal and read journal settings),EmbeddedSnapshotOptionsand the threeWithEmbeddedPersistence(...)overloads, health checks included.sqlite_master), creates missing tables, verifies columns, never alters.BEGIN IMMEDIATEbatches, tag rows, tombstone deletes.FromEnd, keyset-paged ids. Each batch is bounded byMAX(ordering)read in the same transaction; there is no gap-tracking actor (SQLite has one writer).src/aot/Akka.Persistence.Embedded.AOT.Appin theAotCanaryCI job, built with Akka.Hosting: persist, snapshot, recover, delete, every query, an event adapter, health checks, custom table names, and the unregistered failure. Its warning gate covers the plugin, the Hosting package and Akka.Streams.Testing
Akka.DynamicTypeLoadingoff persists, snapshots, recovers and queries by tag through an event adapter; the same with it on; every overload, journal-only, snapshot-only,autoInitialize: false, custom table names, query settings, health checks, two journals via successive calls, two read journals.PublishAoton linux-x64, 0 warnings from the plugin, its Hosting package and Akka.Streams, runs persist, recover, snapshot, delete and every query, shipslibe_sqlite3.so.Breaking changes
None. New packages.
Deliberate differences from Akka.Persistence.Sql: D1 tags with
;come back intact, D2 livePersistenceIdspolls on an interval, D3 no gap-tracking actor, D4 snapshot timestamps are UTC, D5 delete pre-query runs in the delete transaction, D6CurrentPersistenceIdsis keyset-paged, D7 no blocking read journal constructor, D8 the tombstone row'smessageis emptied.Deviations from the spec text: the read journal reads ahead one element (
Buffer(1)afterSelectMany), like Akka.Persistence.Sql does throughSelectAsync(1, ...)andConcatMany, so a finished query completes without extra demand; worker connections open withPooling=Falseunless the connection string setsPoolingitself (logged once at Debug) instead of callingClearPoolat shutdown, so a reset gets a new handle and nothing else's pooled handles are dropped.Dependencies and licenses: Microsoft.Data.Sqlite and .Core 10.0.12 (MIT), SQLitePCLRaw 2.1.12 (Apache-2.0), SQLite (public domain).