Pular para o conteúdo
Português (BR)

CassandraQuery

Este conteúdo não está disponível em sua língua ainda.

Defined in: src/persistence/query/CassandraQuery.ts:47

Cassandra/Scylla query. By default, inherits the journal-walking currentEventsByTag from InMemoryQuery — correct for any volume but only fast for small-to-medium event corpora because the default Cassandra schema has no secondary index on the tags column.

When the journal was constructed with useTagIndex: true, append dual-writes every (event, tag) pair into an events_by_tag side table partitioned by (tag) and ordered by (timestamp, persistence_id, sequence_nr). This class overrides currentEventsByTag to walk that side table — a single tag-partition scan instead of a full client-side journal sweep (#44).

Multi-tag filters. The TagFilter operators (all / any / not, see TagFilter) translate into the side-table query the same way SqliteQuery does:

  • all non-empty → walk the side-table partition for all[0], JS-refine the rest of the filter against the per-row tags set carried in the side table. Bounded scan even when the final result is the intersection of several tags.
  • any non-empty (no all) → walk one partition per any tag and merge by (persistence_id, sequence_nr). N partition scans, sequential, JS-refines not.
  • Only not (or fully empty) → fall back to the inherited journal scan; only-not queries don’t have a selective tag to seed the index walk anyway.

new CassandraQuery(cassandra): CassandraQuery

Defined in: src/persistence/query/CassandraQuery.ts:51

CassandraJournal

CassandraQuery

InMemoryQuery.constructor

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>

InMemoryQuery.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>[]>

InMemoryQuery.currentEventsByPersistenceId


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

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

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

E

TagFilter

Offset

Promise<TaggedEvent<E>[]>

InMemoryQuery.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[]>

InMemoryQuery.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>

InMemoryQuery.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>>

InMemoryQuery.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>>

InMemoryQuery.eventsByTag