Aller au contenu
Français

InMemoryQuery

Ce contenu n’est pas encore disponible dans votre langue.

Defined in: src/persistence/query/InMemoryQuery.ts:56

Reference query implementation that walks any Journal via its public read API. No backend-specific tag index — scans every persistence id on each poll and filters in-process. Correct for any journal, but only fast for the in-memory one (where the scan is just a Map walk).

Backends that ship a “real” tag index (SQLite via the tags column, Cassandra via secondary table) provide their own PersistenceQuery implementation that overrides the tag paths — see SqliteQuery and CassandraQuery.

Push delivery (#42). When the journal exposes a JournalEventBus (journal.events), the live queries — eventsByX and allPersistenceIds — subscribe to it for sub-poll-interval delivery. The polling loop stays as a fallback for cross-process journals (e.g. Cassandra) where in-process notifications can’t reach every subscriber.

Id pagination (#156). Unlike the tag path, the paginated id walk is not overridden per backend: every page goes through InMemoryQuery.readPersistenceIdPage, which forwards to Journal.persistenceIdsPaginated when the backend has a sorted key over ids and cuts the page out of the full list when it does not. Putting the seam on the journal rather than the query keeps the four query classes identical here, and — since they all inherit from this one — is also what makes the feature arrive on every backend at once.

new InMemoryQuery(journal): InMemoryQuery

Defined in: src/persistence/query/InMemoryQuery.ts:57

Journal

InMemoryQuery

allPersistenceIds(options?): AsyncIterable<string>

Defined in: src/persistence/query/InMemoryQuery.ts:149

Live stream of every persistence id the journal has ever seen, plus each new one as it first appears. The stream never completes on its own — break out of the loop or call return() on the iterator.

This is the fan-out primitive: start a per-entity projection as its entity shows up, instead of polling currentPersistenceIds and diffing the result yourself.

Once per stream, not once per journal. A fresh subscription re-emits every id, so a consumer that must not act twice across a restart needs its own checkpoint — the same at-least-once posture the event queries have.

Memory. “Once per stream” is enforced with an in-process set of the ids emitted so far, so a stream over a journal with a million ids holds a million strings for as long as it runs. A lexicographic watermark would be cheaper and wrong: ids are not created in sorted order, so an id that sorts below the mark would be skipped forever, and a fan-out projection that silently never starts is the worst failure this API could have. The catch-up sweep is paged, so the peak is the set plus one page — and a one-shot enumeration that needs no set at all is currentPersistenceIdsPaginated.

LiveQueryOptions = {}

AsyncIterable<string>

PersistenceQuery.allPersistenceIds


currentEventsByPersistenceId<E>(persistenceId, fromSeq, toSeq?): Promise<PersistentEvent<E>[]>

Defined in: src/persistence/query/InMemoryQuery.ts:61

One-shot read of every event for persistenceId whose sequenceNr >= fromSeq (and <= toSeq if given). Resolves once with the events known at call time.

E

string

number

number

Promise<PersistentEvent<E>[]>

PersistenceQuery.currentEventsByPersistenceId


currentEventsByTag<E>(filter, fromOffset): Promise<TaggedEvent<E>[]>

Defined in: src/persistence/query/InMemoryQuery.ts:85

One-shot read of every event matching filter whose offset is >= fromOffset. See TagFilter for the operator semantics.

E

TagFilter

Offset

Promise<TaggedEvent<E>[]>

PersistenceQuery.currentEventsByTag


currentPersistenceIds(): Promise<string[]>

Defined in: src/persistence/query/InMemoryQuery.ts:127

Snapshot of every persistence id known to the journal, as one array. Resolves once.

Kept, and not deprecated: for a journal with a handful of ids this is simply the convenient shape, and “small journal” is a permanent case rather than a legacy one. It is, however, the only method here with no bound on what it materialises — at a million ids that is a million strings in one allocation. Past the point where the array itself is a cost, use currentPersistenceIdsPaginated, which walks the same data one page at a time.

Promise<string[]>

PersistenceQuery.currentPersistenceIds


currentPersistenceIdsPaginated(options?): AsyncIterable<string>

Defined in: src/persistence/query/InMemoryQuery.ts:134

Cursor-paginated snapshot of the persistence ids currently in the journal, ascending, yielded one at a time. Fetches pageSize ids per round-trip, so peak memory is one page rather than the whole set. Completes when the journal is exhausted — unlike allPersistenceIds it does not wait for new ids.

The cursor is a persistence id, not an opaque token. Resume a partial walk by passing the last id you handled as afterPersistenceId. Encoding it would have bought per-backend freedom that no backend here needs — every one of them enumerates ids through an index keyed on the id itself — at the price of a checkpoint no operator can read and no other process can construct.

PaginationOptions = {}

AsyncIterable<string>

PersistenceQuery.currentPersistenceIdsPaginated


eventsByPersistenceId<E>(persistenceId, fromSeq, options?): AsyncIterable<PersistentEvent<E>>

Defined in: src/persistence/query/InMemoryQuery.ts:67

Live stream of every event for persistenceId whose sequenceNr >= fromSeq. Past events are emitted first (chronological by sequenceNr), then new events as they are appended. The stream never completes on its own — break out of the loop or call return() on the iterator to stop polling.

E

string

number

LiveQueryOptions = {}

AsyncIterable<PersistentEvent<E>>

PersistenceQuery.eventsByPersistenceId


eventsByTag<E>(filter, fromOffset, options?): AsyncIterable<TaggedEvent<E>>

Defined in: src/persistence/query/InMemoryQuery.ts:104

Live stream of every event matching filter whose offset is >= fromOffset. Yields events ordered by (timestamp, persistenceId, sequenceNr). See Offset for offset semantics — the stream emits the offset alongside the event so the consumer can persist progress.

filter accepts either a single tag string (back-compat shortcut for { all: [tag] }) or a TagFilter object that combines all (intersect), any (union), and not (exclusion) operators.

E

TagFilter

Offset

LiveQueryOptions = {}

AsyncIterable<TaggedEvent<E>>

PersistenceQuery.eventsByTag