Add a "from the end" (last-N) query Offset to Akka.Persistence.Query - #8245
Merged
Aaronontheweb merged 7 commits intoJun 23, 2026
Merged
Aaronontheweb merged 7 commits into
Aaronontheweb merged 7 commits into
Conversation
…Query Introduce Offset.FromEnd(int count), a query-input-only offset that begins a query at the Nth event from the end of history rather than the beginning. Read journals resolve it at materialization into a concrete Sequence start offset and then reuse their existing forward-streaming pipeline, so the change is fully additive: no interface signature changes and no wire-format changes. - Core: FromEnd : Offset type + Offset.FromEnd(int) factory - InMemory: resolve FromEnd in the by-tag and all-events publishers via a new MemoryJournal.SelectEventCount/EventCount round-trip - TCK: opt-in FromEndOffsetSpec (last-N, N > total, live continuation) and InMemoryFromEndOffsetSpec - Update API approval files for the additive surface SQL and MongoDB plugin implementations to follow. See akkadotnet#8244
- Fail the FromEnd resolution stream instead of hanging it: MemoryJournal now pipes a failure (EventCountFailure) for SelectEventCount, and the resolving state surfaces it via OnErrorThenStop, matching the sibling replay handlers. - Extract the duplicated FromEnd resolution handshake from both InMemory query publishers into a shared FromEndResolvingPublisher base. - Deduplicate the event-count predicate behind MemoryJournal.CountEvents and reuse it from the replay paths. - FromEnd no longer implements IComparable: it is a relative, input-only offset with no stream position, so CompareTo now always throws. - Document that the from-the-end window is resolved at materialization and is best-effort under concurrent writes. See akkadotnet#8244
Harden the FromEnd ("last N events") contract before it backs SQL/Mongo:
- FromEndOffsetSpec: grow from 4 to 9 tests covering the full
{by-tag, all-events} x {current, live} matrix against an interleaved,
multi-persistence-id fixture mixing tagged and untagged events. Adds
N>total and empty/zero-match cases for all-events, live AllEvents
continuation, and a discriminating test pinning that by-tag FromEnd
resolves against the per-tag count while all-events resolves against the
total count (a single-persistence-id fixture cannot tell these apart, so
a backend resolving both against a global ordinal would pass today).
Converted to async TestKit style (ExpectMsgAsync/ExpectNextAsync/
AwaitConditionAsync) to match the sibling AllEventsSpec.
- OffsetSpec/OffsetCompareSpecs: add FromEnd unit tests for the factory,
the count<=0 guard, ToString, Equals/GetHashCode, and CompareTo throwing.
See akkadotnet#8244
Fixes from a high-effort review of the FromEnd spec expansion: - FromEndOffsetSpec: the cross-query discriminator test now also waits for the green tag index to settle (not just all-events) before resolving its by-tag from-end window. Previously it could resolve against a partial tag set and flake on backends whose tag projection lags the global ordering (hidden on InMemory, where writes are immediately visible). - FromEndOffsetSpec: the live tests now require the current-query counterpart (RequireQuery) instead of silently skipping stabilization via an "is ICurrent...Query" check. The from-end window is resolved once at materialization, so a live-only backend with read-side lag would otherwise resolve it against an unsettled journal. Fail loud rather than flake. - FromEndOffsetSpec: derive AllEvents/GreenEvents from a single Fixture table (tagging by event text) instead of hand-maintaining two parallel arrays that could drift from the writes. - FromEndOffsetSpec: collapse the two near-identical WaitFor* helpers into one WaitForVisibleAsync(query, count); project the discriminator's result tuples once instead of recomputing them per assertion. - FromEndOffsetSpec: add a deep-history by-tag test (25 events, FromEnd(4)) to exercise the "start = count - N" arithmetic at depth and a long forward replay, restoring coverage the single-PID fixture had lost. - OffsetSpec/OffsetCompareSpecs: add the missing #nullable enable directive (CLAUDE.md: enable nullable in new/modified files). All green: 10 InMemory FromEnd TCK tests + 15 Offset unit tests; -warnaserror clean. See akkadotnet#8244
Aaronontheweb
marked this pull request as ready for review
June 23, 2026 00:55
Supporting a "from the end" start is an internal detail of how the InMemory read journal interprets an Offset, not something every publisher constructor should know about. The previous design threaded two parallel ints — fromOffset and fromEndCount — through Props and three constructor layers, where fromEndCount > 0 was an implicit "resolve from the end, ignore fromOffset" sentinel (primitive-obsession + boolean-blindness). Replace that pair with a single internal readonly struct, ReplayStart, with named factories (At(offset) / LastN(count)) and an IsFromEnd discriminator. The read journal's Offset resolvers now return a ReplayStart; publishers take one argument instead of two and no longer re-derive the mode from a >0 check. Also drops the now-dead FromOffset property from both publishers (only ever assigned; CurrentOffset drives the replay). Purely internal: every affected type is internal, so there is no public API or wire change and no API-approval delta. Behavior is unchanged — full InMemory query suite (56 tests incl. the 10 FromEnd TCK tests) green, -warnaserror clean. See akkadotnet#8244
Aaronontheweb
commented
Jun 23, 2026
Aaronontheweb
left a comment
Member
Author
There was a problem hiding this comment.
LGTM - good to go for a prototype I think
| /// detail of how the read journal interprets an <see cref="Query.Offset"/>, instead of an extra argument threaded | ||
| /// through every publisher constructor. | ||
| /// </summary> | ||
| internal readonly struct ReplayStart |
| var probe = queries.CurrentEventsByTag("green", FromEnd(2)) | ||
| .RunWith(this.SinkProbe<EventEnvelope>(), Materializer); | ||
| probe.Request(10); | ||
| // the last two green events, across persistence ids, in ascending order |
| } | ||
|
|
||
| [Fact] | ||
| public virtual async Task ReadJournal_query_CurrentEventsByTag_with_FromEnd_larger_than_total_should_return_all_events() |
6 tasks
This was referenced Jul 2, 2026
Aaronontheweb
added a commit
that referenced
this pull request
Jul 3, 2026
…8245) (#8308) * feat: add FromEnd ("last N events") query offset to Akka.Persistence.Query Introduce Offset.FromEnd(int count), a query-input-only offset that begins a query at the Nth event from the end of history rather than the beginning. Read journals resolve it at materialization into a concrete Sequence start offset and then reuse their existing forward-streaming pipeline, so the change is fully additive: no interface signature changes and no wire-format changes. - Core: FromEnd : Offset type + Offset.FromEnd(int) factory - InMemory: resolve FromEnd in the by-tag and all-events publishers via a new MemoryJournal.SelectEventCount/EventCount round-trip - TCK: opt-in FromEndOffsetSpec (last-N, N > total, live continuation) and InMemoryFromEndOffsetSpec - Update API approval files for the additive surface SQL and MongoDB plugin implementations to follow. See #8244 * refactor: address code review for FromEnd query offset - Fail the FromEnd resolution stream instead of hanging it: MemoryJournal now pipes a failure (EventCountFailure) for SelectEventCount, and the resolving state surfaces it via OnErrorThenStop, matching the sibling replay handlers. - Extract the duplicated FromEnd resolution handshake from both InMemory query publishers into a shared FromEndResolvingPublisher base. - Deduplicate the event-count predicate behind MemoryJournal.CountEvents and reuse it from the replay paths. - FromEnd no longer implements IComparable: it is a relative, input-only offset with no stream position, so CompareTo now always throws. - Document that the from-the-end window is resolved at materialization and is best-effort under concurrent writes. See #8244 * test: expand FromEnd TCK to full query matrix + add Offset unit tests Harden the FromEnd ("last N events") contract before it backs SQL/Mongo: - FromEndOffsetSpec: grow from 4 to 9 tests covering the full {by-tag, all-events} x {current, live} matrix against an interleaved, multi-persistence-id fixture mixing tagged and untagged events. Adds N>total and empty/zero-match cases for all-events, live AllEvents continuation, and a discriminating test pinning that by-tag FromEnd resolves against the per-tag count while all-events resolves against the total count (a single-persistence-id fixture cannot tell these apart, so a backend resolving both against a global ordinal would pass today). Converted to async TestKit style (ExpectMsgAsync/ExpectNextAsync/ AwaitConditionAsync) to match the sibling AllEventsSpec. - OffsetSpec/OffsetCompareSpecs: add FromEnd unit tests for the factory, the count<=0 guard, ToString, Equals/GetHashCode, and CompareTo throwing. See #8244 * test: address code-review findings on FromEnd TCK + Offset unit tests Fixes from a high-effort review of the FromEnd spec expansion: - FromEndOffsetSpec: the cross-query discriminator test now also waits for the green tag index to settle (not just all-events) before resolving its by-tag from-end window. Previously it could resolve against a partial tag set and flake on backends whose tag projection lags the global ordering (hidden on InMemory, where writes are immediately visible). - FromEndOffsetSpec: the live tests now require the current-query counterpart (RequireQuery) instead of silently skipping stabilization via an "is ICurrent...Query" check. The from-end window is resolved once at materialization, so a live-only backend with read-side lag would otherwise resolve it against an unsettled journal. Fail loud rather than flake. - FromEndOffsetSpec: derive AllEvents/GreenEvents from a single Fixture table (tagging by event text) instead of hand-maintaining two parallel arrays that could drift from the writes. - FromEndOffsetSpec: collapse the two near-identical WaitFor* helpers into one WaitForVisibleAsync(query, count); project the discriminator's result tuples once instead of recomputing them per assertion. - FromEndOffsetSpec: add a deep-history by-tag test (25 events, FromEnd(4)) to exercise the "start = count - N" arithmetic at depth and a long forward replay, restoring coverage the single-PID fixture had lost. - OffsetSpec/OffsetCompareSpecs: add the missing #nullable enable directive (CLAUDE.md: enable nullable in new/modified files). All green: 10 InMemory FromEnd TCK tests + 15 Offset unit tests; -warnaserror clean. See #8244 * refactor: collapse FromEnd plumbing into an internal ReplayStart value Supporting a "from the end" start is an internal detail of how the InMemory read journal interprets an Offset, not something every publisher constructor should know about. The previous design threaded two parallel ints — fromOffset and fromEndCount — through Props and three constructor layers, where fromEndCount > 0 was an implicit "resolve from the end, ignore fromOffset" sentinel (primitive-obsession + boolean-blindness). Replace that pair with a single internal readonly struct, ReplayStart, with named factories (At(offset) / LastN(count)) and an IsFromEnd discriminator. The read journal's Offset resolvers now return a ReplayStart; publishers take one argument instead of two and no longer re-derive the mode from a >0 check. Also drops the now-dead FromOffset property from both publishers (only ever assigned; CurrentOffset drives the replay). Purely internal: every affected type is internal, so there is no public API or wire change and no API-approval delta. Behavior is unchanged — full InMemory query suite (56 tests incl. the 10 FromEnd TCK tests) green, -warnaserror clean. See #8244 (cherry picked from commit 938a629)
This was referenced Jul 3, 2026
Open
Bump Akka.TestKit.Xunit2 from 1.5.38 to 1.5.70
cuteboy0323/Aaron.Akka.Streams.BackpressureMonitor#12
Open
Open
This was referenced Jul 27, 2026
This was referenced Aug 18, 2026
Closed
This was referenced Aug 27, 2026
Closed
Closed
This was referenced Sep 21, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #8244 (phase 1).
Summary
Adds a new
Offset.FromEnd(int count)query offset toAkka.Persistence.Queryso callers can retrieve the last N events (per-tag or across all events) without replaying the whole journal or maintaining a custom projection.FromEndis a query input only — it is never emitted in anEventEnvelope. A read journal resolves it at materialization into a concreteSequencestart offset (the Nth-from-last matchingOrdering) and then reuses its existing forward-streaming pipeline. The change is fully additive: no interface signature changes and no wire-format changes. Backends that don't recognize it throw on the unknown offset exactly as they do today.CurrentEventsByTag/CurrentAllEvents): resolve start → stream forward → complete = exactly the last N.EventsByTag/AllEvents): resolve start once, then continue forward = last N, then everything new.Release plan
This ships in an upcoming 1.5 beta, and the beta is the validation vehicle for the offset contract across the wider Akka.NET ecosystem. The public
Offset.FromEndsurface and its API approvals are included in this PR so that:Akka.Persistence.Sql,Akka.Persistence.MongoDb, community backends) can implement and validate against the offset + the TCK during the beta, andv1.6.0.Only InMemory implements
FromEndin this PR; every other backend throws on the unknown offset until it opts in (tracked in #8244). That is expected for the beta — the TCK below is the spec those backends build against.What's in this PR (phase 1)
FromEnd : Offset+Offset.FromEnd(int)factory (input-only;CompareTothrows because the offset is relative, not a stream position). Constructor rejectscount <= 0.FromEnd(N)via a newMemoryJournal.SelectEventCount/EventCountround-trip (sharedFromEndResolvingPublisherbase), reusing the existing replay machinery; a faulted count query fails the stream (EventCountFailure) rather than hanging it.FromEndOffsetSpec— the cross-backend contract. Opt-in spec covering the full{by-tag, all-events} × {current, live}matrix against an interleaved, multi-persistence-id fixture that mixes tagged and untagged events:FromEndresolves against the per-tag count while all-events resolves against the total count — a distinction a single-persistence-id fixture cannot express, so a backend resolving both against a global ordinal would pass a naive spec but fail this onestart = count - Narithmetic at depthOffsetSpec/OffsetCompareSpecscover theFromEndfactory, thecount <= 0guard,ToString,Equals/GetHashCode, andCompareTothrowing.Review status
Opened as a design-review draft; has since had a high-effort review pass with all findings addressed. InMemory is green against the hardened TCK (10 TCK tests + 15 Offset unit tests),
-warnaserrorclean, CI passing.Design notes
FromEndparticipates inOffset'sIComparable<Offset>butCompareTothrows — generic comparison-based code paths (e.g. sorting a mixed offset set) will throw on aFromEnd. Intentional: it is a relative, input-only offset with no stream position, and is never emitted in an envelope. Calling it out so beta consumers know not to putFromEndin a comparison/sort path.Scope / follow-ups
Akka.Persistence.Sql(Linq2Db) andAkka.Persistence.MongoDbimplementations + their TCK opt-ins. Azure is awkward (Table Storage has no server-side DESC/LIMIT); Redis is out of scope until its tag/all-events queries exist.HighestSequenceNr - N + 1helper (listed in Add a "from the end" (last-N) query Offset to Akka.Persistence.Query #8244) is not in this PR; defer or fold into a follow-up.Caveat
The from-the-end position is resolved at materialization by reading the current matching count; for live queries it is best-effort — events written between resolving the count and reading the first batch may widen the initial window beyond N (documented on the type).