From 76598f842031920ff04be04c4a925bcc361e924b Mon Sep 17 00:00:00 2001 From: tamirms Date: Fri, 25 Sep 2026 14:56:52 -0400 Subject: [PATCH] events: match queries one slab at a time A getEvents query was answered by intersecting each filter's term bitmaps across the whole chunk, unioning the results across filters, and only then restricting to the window. Its cost tracked the terms' size across the chunk, not the window the client asked for or the page it would receive. The engine now walks the window one slab at a time: 65,536 consecutive event ids, roaring's container size, in the requested direction. Each filter becomes plans that are plain ANDs of terms. On each slab, every plan is one FastAnd over the slab's ids in the window and the plan's terms, and the plans' results are ORed and handed to the fetch. Before opening a slab, the walk asks each term for its next id and jumps to the first slab that can hold a candidate, so a rare term opens only the slabs it appears in. Work scales with the window and the page, memory per query with one container, and both directions share one path. The first batch is the caller's page size, so a full page arrives in one fetch. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01QHnF5BhsuoxpWxGmatQzmt --- .../rpcv2/stores/event/concurrent_bitmaps.go | 14 +- .../internal/rpcv2/stores/event/match.go | 353 +++----- .../internal/rpcv2/stores/event/match_test.go | 92 +- .../stores/event/matches_differential_test.go | 271 ++++++ .../stores/event/roaring_contract_test.go | 341 +++++++ .../internal/rpcv2/stores/event/slab_match.go | 391 ++++++++ .../rpcv2/stores/event/slab_match_test.go | 850 ++++++++++++++++++ 7 files changed, 2011 insertions(+), 301 deletions(-) create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/matches_differential_test.go create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/roaring_contract_test.go create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match.go create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match_test.go diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps.go index ad2d2cefd..102e2f841 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps.go @@ -93,15 +93,19 @@ func NewConcurrentBitmapsFromBitmaps(b Bitmaps) *ConcurrentBitmaps { // - RunOptimize, AddRange, RemoveRange, FlipInt // - Add, AddMany, Remove, CheckedAdd, CheckedRemove, AddInt // - SetCopyOnWrite -// - single-input roaring.FastAnd / roaring.FastOr (roaring takes a -// Clone-the-input shortcut when there is only one input) +// - roaring.FastAnd / roaring.FastOr with a single input, which Clone it // // Safe: any non-mutating read (Contains, GetCardinality, Iterator, -// ToArray, IsEmpty, Minimum, Maximum) plus roaring.And / FastAnd / -// FastOr with 2+ inputs. +// NextValue, PreviousValue, ToArray, IsEmpty, Minimum, Maximum) and +// passing it to FastAnd beside at least one other input; +// roaring_contract_test.go pins these against the pinned roaring version. // // A Get that starts after an AddTo returns sees that AddTo's IDs, and -// the pointer stays valid for as long as the caller holds it. +// the pointer stays valid for as long as the caller holds it. The +// result is a point-in-time image: a sparse term is copied out of the +// published id list, and a dense term's snapshot is never mutated once +// published. Neither grows as ingest continues, so a query can hold it +// for a whole walk. func (s *ConcurrentBitmaps) Get(key TermKey) (*roaring.Bitmap, error) { s.rwmu.RLock() p := s.terms[key] diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/match.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/match.go index d263667dd..4a471dd95 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/match.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/match.go @@ -7,20 +7,20 @@ package event // Matches. // // Optimization shape: terms are deduped across filters and issued as -// a single Reader.LookupKeys call at iteration start; payload fetches -// then stream in internal batches. On the cold path this is one -// MPHF+index.pack round trip per Matches call, not per batch. +// a single batched Reader.LookupKeys at iteration start, whose bitmaps +// the walk holds for the whole query; payload fetches stream in +// internal batches. On the cold path the lookup is one MPHF+index.pack +// round trip per Matches call. The candidate set comes from the slab +// engine in slab_match.go, which serves both directions from one walk +// over the window. import ( "bytes" - "cmp" "context" "fmt" "iter" "slices" - "github.com/RoaringBitmap/roaring/v2" - protocol "github.com/stellar/go-stellar-sdk/protocols/rpc" "github.com/stellar/go-stellar-sdk/xdr" ) @@ -85,9 +85,9 @@ func (f TopicCountFilter) termKeys() []TermKey { // valueTermKeys returns one term per constrained value field // (contract ID, event type, topics): the single enumeration -// termGroups and CountDistinctTerms share, so the two cannot drift +// termPlans and CountDistinctTerms share, so the two cannot drift // over which values a filter names. The topic-count buckets are not -// value terms; termGroups adds them separately and the budget does +// value terms; termPlans adds them separately and the budget does // not count them. func (f *Filter) valueTermKeys() []TermKey { var keys []TermKey @@ -106,25 +106,30 @@ func (f *Filter) valueTermKeys() []TermKey { return keys } -// termGroups returns the indexed terms this filter constrains, grouped -// by field: the bitmaps within a group are OR-ed and the groups are -// AND-ed. Only the topic-count group ever holds more than one term. -func (f *Filter) termGroups() [][]TermKey { - var groups [][]TermKey - for _, key := range f.valueTermKeys() { - groups = append(groups, []TermKey{key}) - } - // A constrained topic position already implies an "at least" count at - // or below it, since a topic term is only indexed for events carrying - // that position. Skipping the group there keeps the common - // ["a", "**"] shape from OR-ing chunk-sized bucket bitmaps; the - // post-filter enforces the count either way. +// termPlans returns the index terms a candidate must carry, one +// conjunction per plan. A filter yields one plan, its value terms, unless +// it constrains the topic count, when it yields one plan per bucket, a +// single one for an exact count: A and (b or c) is (A and b) or (A and +// c), and the union across plans keeps the or. A constrained topic +// position already implies an "at least" count at or below it, since a +// topic term is only indexed for events carrying that position. Skipping +// the buckets there keeps the common ["a", "**"] shape from fanning out +// over chunk-sized bucket bitmaps; the post-filter enforces the count +// either way. +func (f *Filter) termPlans() [][]TermKey { + values := f.valueTermKeys() + var buckets []TermKey if !f.impliesTopicCount() { - if keys := f.TopicCount.termKeys(); len(keys) > 0 { - groups = append(groups, keys) - } + buckets = f.TopicCount.termKeys() + } + if len(buckets) == 0 { + return [][]TermKey{values} + } + plans := make([][]TermKey, 0, len(buckets)) + for _, bucket := range buckets { + plans = append(plans, append(slices.Clone(values), bucket)) } - return groups + return plans } // impliesTopicCount reports whether f's constrained topic positions @@ -214,18 +219,18 @@ type Match struct { Ordinal uint32 } -// termPlan is a filter's termGroups resolved to slots in the batched -// LookupKeys result. -type termPlan [][]int +// termPlan is one of a filter's termPlans, as slots in the batched term +// lookup's result. +type termPlan []int -// batchSizes resolves the first and following internal batch sizes -// from the caller's hint. The hint applies only when it is positive -// and below the default. Both results are clamped positive, so a zero -// test seam cannot stall a stream (a zero step never advances). +// batchSizes resolves the first and following internal batch sizes from the +// caller's hint. The hint is a page size the handler has already validated, +// so it is honored in full and a page arrives in one fetch. Both sizes are +// clamped positive: a zero step would stall the stream. func batchSizes(hint int) (int, int) { rest := max(1, matchBatchSize) first := rest - if hint > 0 && hint < rest { + if hint > 0 { first = hint } return first, rest @@ -251,11 +256,9 @@ func batchSizes(hint int) (int, int) { // drops are invisible: the iterator advances past them internally, so // consumers never see or reason about resume state. // -// firstBatch sizes the first internal fetch batch. A consumer that -// will stop after N matches passes N, so the first round trip fetches -// no more than it needs; 0 or any out-of-range value uses the -// default. Later batches use the default size. The hint changes I/O -// counts only, never what the stream yields. +// firstBatch sizes the first internal fetch: a consumer that will stop after +// N matches passes N. Zero and negative hints use the default. The hint +// changes I/O counts only, never what the stream yields. func Matches( ctx context.Context, r Reader, filters []Filter, window IDRange, descending bool, firstBatch int, @@ -268,11 +271,7 @@ func Matches( if window.isEmpty() { return } - union, matchAll, err := unionForFilters(ctx, r, filters, window) - if err != nil { - yield(Match{}, err) - return - } + plans, uniqueKeys, matchAll := planIndexTerms(filters) // Match-all path: empty filter slice or any filter that asks the // index for no terms. Serves without touching the index: the // window is dense, so it streams Reader.FetchRange directly. @@ -280,10 +279,18 @@ func Matches( streamRange(ctx, r, window, descending, firstBatch, yield) return } - if union.IsEmpty() { + sources, err := r.LookupKeys(ctx, uniqueKeys) + if err != nil { + yield(Match{}, fmt.Errorf("events: query lookup: %w", err)) + return + } + st := newSlabStepper(plans, sources, window, descending) + // No plan survived term resolution, so nothing can match and no slab + // is worth evaluating. + if len(st.plans) == 0 { return } - streamUnion(ctx, r, filters, union, descending, firstBatch, yield) + streamSlabs(ctx, r, filters, st, descending, firstBatch, yield) } } @@ -312,197 +319,73 @@ func validateMatchCall(ctx context.Context, r Reader, filters []Filter, window I return nil } -// unionForFilters runs the index side once per Matches call (the -// numbered steps below). matchAll reports that some filter (or the -// empty slice) constrains nothing, detected before any index I/O; the -// caller then streams the window directly. Otherwise the result is -// empty when no candidate falls in the window, and is never a -// borrowed mirror snapshot (the window AND allocates on the borrowing -// path), so downstream iteration is safe. -func unionForFilters( - ctx context.Context, r Reader, filters []Filter, window IDRange, -) (*roaring.Bitmap, bool, error) { - // ───── 1. Dedupe terms across filters ───── - // - // filterPlans[i] holds the slots filter i needs out of the batched - // lookup: the bitmaps within a group are OR-ed, the groups AND-ed. - // - // A filter that asks the index for no terms constrains nothing, and - // so does an empty filter slice: both take the match-all path. - // Reading the condition off the term groups themselves is what keeps - // an unconstrained filter from intersecting nothing and coming back - // empty instead. +// planIndexTerms maps every filter's plans to slots in the single batched +// lookup that follows; it runs before any index I/O. A plan that repeats +// an earlier one is dropped: plans only pick candidates, and the +// post-filter still runs every filter. +// +// matchAll reports that some filter, or the empty slice, constrains +// nothing, so the caller streams the window directly rather than +// intersecting nothing and returning empty. +func planIndexTerms(filters []Filter) ([]termPlan, []TermKey, bool) { if len(filters) == 0 { - return nil, true, nil + return nil, nil, true } - filterPlans := make([]termPlan, len(filters)) var uniqueKeys []TermKey - + plans := make([]termPlan, 0, len(filters)) for i := range filters { - groups := filters[i].termGroups() - if len(groups) == 0 { - return nil, true, nil - } - plan := make(termPlan, len(groups)) - for g, keys := range groups { - slots := make([]int, len(keys)) + for _, keys := range filters[i].termPlans() { + if len(keys) == 0 { + return nil, nil, true + } + plan := make(termPlan, len(keys)) for j, key := range keys { - slots[j] = indexOfOrAddTerm(&uniqueKeys, key) + plan[j] = indexOfOrAddTerm(&uniqueKeys, key) } - plan[g] = slots - } - filterPlans[i] = plan - } - - // ───── 2. Single batched lookup for all unique terms ───── - bitmaps, err := r.LookupKeys(ctx, uniqueKeys) - if err != nil { - return nil, false, fmt.Errorf("events: query lookup: %w", err) - } - - // ───── 3. Per-filter intersect ───── - // - // If a whole group is absent from the index (every bitmap in it is - // nil), that filter's intersection is empty — skip it without - // contributing to the union. - // - // Bitmap ownership in perFilter is mixed: - // - Single-constraint filter: we borrow bitmaps[s] directly (a - // mirror snapshot from LookupKeys), skipping FastAnd's Clone. - // - Multi-constraint filter: FastAnd allocates a fresh result. - // Either way the downstream union (FastOr) and the window AND never - // mutate their inputs, so a borrowed entry stays valid through - // the rest of the function. FastAnd never mutates its inputs - // either, so the same bitmap may appear across multiple filters - // safely. - perFilter := make([]*roaring.Bitmap, 0, len(filterPlans)) - for _, plan := range filterPlans { - inputs := make([]*roaring.Bitmap, 0, len(plan)) - missed := false - for _, slots := range plan { - group := unionSlots(bitmaps, slots) - if group == nil { - missed = true - break + // Slots follow field order, so equal plans are equal slices. + dup := slices.ContainsFunc(plans, func(p termPlan) bool { + return slices.Equal(p, plan) + }) + if dup { + continue } - inputs = append(inputs, group) - } - if missed { - continue + plans = append(plans, plan) } - if len(inputs) == 1 { - perFilter = append(perFilter, inputs[0]) - continue - } - // FastAnd intersects left-to-right — putting the smallest - // bitmap first shrinks the accumulator fastest. roaring's own - // docs call this out as the recommended caller-side prep. - slices.SortFunc(inputs, func(a, b *roaring.Bitmap) int { - return cmp.Compare(a.GetCardinality(), b.GetCardinality()) - }) - perFilter = append(perFilter, roaring.FastAnd(inputs...)) - } - - if len(perFilter) == 0 { - return roaring.New(), false, nil - } - - // ───── 4. Union across filters ───── - // Single-filter case: FastOr would Clone — skip it and use the - // already-computed bitmap directly. That bitmap may be borrowed - // (from LookupKeys), so step 5's window And uses the fresh-result - // variant on that path to avoid mutating shared state. - var union *roaring.Bitmap - singleFilter := len(perFilter) == 1 - if singleFilter { - union = perFilter[0] - } else { - union = roaring.FastOr(perFilter...) - } - - // ───── 5. Apply the event-ID window ───── - // - // The window AND enforces the caller's pinned range. It also clips - // phantom IDs from a concurrent hot-store ingest: the mirror - // publishes entries before offsets, so LookupKeys can briefly - // surface IDs past EventCount. The AND keeps the stream strictly - // within the snapshot the caller pinned at request entry. - // - // This covers the multi-term group too. Its bitmaps are separate - // mirror snapshots taken at different instants, but an event never - // moves between the terms of one group once ingested, so a torn read - // across them can only surface IDs past the pinned End. - rangeBM := roaring.New() - rangeBM.AddRange(uint64(window.Start), uint64(window.End)) - if singleFilter { - union = roaring.And(union, rangeBM) // fresh result; union may be borrowed - } else { - union.And(rangeBM) // FastOr output is owned; in-place is fine - } - return union, false, nil + } + return plans, uniqueKeys, false } -// streamUnion walks the union bitmap in internal batches: collect -// candidate ordinals up to the batch size, fetch, post-filter, yield -// the survivors. Drops advance the walk with no yield. The first -// batch is sized to firstBatch (see Matches); later batches use the -// default. -// -// FetchEvents requires ascending ids, so a descending batch is -// collected highest-first and flipped before the fetch, then the -// fetched matches are flipped back. Stepping one id at a time is fine -// here: the fetch I/O dominates a 512-step loop. -func streamUnion( - ctx context.Context, r Reader, filters []Filter, union *roaring.Bitmap, - descending bool, firstBatch int, yield func(Match, error) bool, -) { - var it interface { - HasNext() bool - Next() uint32 +// emitBatch fetches one batch of candidate ordinals, drops the bitmap-side +// false positives and yields the survivors, reporting whether the stream +// should continue. FetchEvents requires ascending ids, so a descending batch +// is flipped in place before the fetch and flipped back before yielding. +func emitBatch( + ctx context.Context, r Reader, filters []Filter, ids []uint32, + descending bool, yield func(Match, error) bool, +) bool { + if descending { + slices.Reverse(ids) + } + payloads, err := r.FetchEvents(ctx, ids) + if err != nil { + yield(Match{}, err) + return false + } + // Drop bitmap-side false positives (see postFilter for the rationale). + matched, err := postFilter(payloads, ids, filters) + if err != nil { + yield(Match{}, err) + return false } if descending { - it = union.ReverseIterator() - } else { - it = union.Iterator() + slices.Reverse(matched) } - batch, rest := batchSizes(firstBatch) - ids := make([]uint32, 0, batch) - for { - if err := ctx.Err(); err != nil { - yield(Match{}, err) - return - } - ids = ids[:0] - for it.HasNext() && len(ids) < batch { - ids = append(ids, it.Next()) - } - batch = rest - if len(ids) == 0 { - return - } - if descending { - slices.Reverse(ids) - } - payloads, err := r.FetchEvents(ctx, ids) - if err != nil { - yield(Match{}, err) - return - } - // Drop bitmap-side false positives (see postFilter for the rationale). - matched, err := postFilter(payloads, ids, filters) - if err != nil { - yield(Match{}, err) - return - } - if descending { - slices.Reverse(matched) - } - for i := range matched { - if !yield(matched[i], nil) { - return - } + for i := range matched { + if !yield(matched[i], nil) { + return false } } + return true } // ValidateFilters rejects filters that would silently never match @@ -561,9 +444,9 @@ func ValidateFilters(filters []Filter) error { // filters name, deduped by field and value together: one contract ID // in five filters counts once, the same bytes in two topic positions // count twice. Topic-count buckets are excluded: they are an -// implementation detail of the engine's grouping, not a value the +// implementation detail of the engine's plans, not a value the // client named. Exported for the v2 handler's term-budget check. It -// lives here, beside termGroups, so the budget and the engine's +// lives here, beside termPlans, so the budget and the engine's // lookups agree on what a value term is: TermKey over the store's // canonical bytes. func CountDistinctTerms(filters []Filter) int { @@ -576,30 +459,6 @@ func CountDistinctTerms(filters []Filter) int { return len(unique) } -// unionSlots ORs the bitmaps at slots, and returns nil when every one -// of them is absent from the index. A lone present bitmap is borrowed -// rather than cloned, like the single-constraint path in -// unionForFilters. -func unionSlots(bitmaps []*roaring.Bitmap, slots []int) *roaring.Bitmap { - if len(slots) == 1 { - return bitmaps[slots[0]] - } - present := make([]*roaring.Bitmap, 0, len(slots)) - for _, s := range slots { - if bitmaps[s] != nil { - present = append(present, bitmaps[s]) - } - } - switch len(present) { - case 0: - return nil - case 1: - return present[0] - default: - return roaring.FastOr(present...) - } -} - // indexOfOrAddTerm returns the index of key inside *keys, appending // it first if absent. func indexOfOrAddTerm(keys *[]TermKey, key TermKey) int { diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/match_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/match_test.go index 4d26c5353..69eada19b 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/match_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/match_test.go @@ -620,9 +620,9 @@ func TestQuery_ChunkWithLedgersButZeroEvents(t *testing.T) { // TestQuery_DescendingWithRangeAndMaxEvents covers the // three-way combination — order × range × cap — that no other test -// hits together. Forces the descending branch of streamUnion -// (ReverseIterator with a per-batch flip) over a range-narrowed -// union, then the shim's MaxEvents truncation. +// hits together. Forces the descending slab walk (each slab's result +// read backwards, with a per-batch flip for the fetch) over a +// range-narrowed window, then the shim's MaxEvents truncation. func TestQuery_DescendingWithRangeAndMaxEvents(t *testing.T) { fx := newMultiLedgerQueryFixture(t) first := chunk.ID(0).FirstLedger() @@ -1126,26 +1126,6 @@ func TestQuery_InvalidFilterRejected(t *testing.T) { } } -// TestUnionSlots covers the OR-within-a-group step directly, including the -// all-absent case a fixture cannot reach: the topic-count buckets are the only -// multi-term group, and the overflow bucket is populated in any chunk holding -// an event with topics. -func TestUnionSlots(t *testing.T) { - first := roaring.BitmapOf(1, 2) - second := roaring.BitmapOf(3) - bitmaps := []*roaring.Bitmap{first, nil, second, nil} - - assert.Same(t, first, unionSlots(bitmaps, []int{0}), - "a lone bitmap is borrowed, not cloned") - assert.Nil(t, unionSlots(bitmaps, []int{1})) - assert.Nil(t, unionSlots(bitmaps, []int{1, 3}), - "a group absent from the index empties the filter") - assert.Same(t, second, unionSlots(bitmaps, []int{1, 2}), - "the one present bitmap in a group is borrowed too") - assert.Equal(t, []uint32{1, 2, 3}, unionSlots(bitmaps, []int{0, 2}).ToArray()) - assert.Equal(t, []uint32{1, 2}, first.ToArray(), "inputs must not be mutated") -} - // ─── Cold-reader parity coverage ──────────────────────────────────────── // // The hot tests above prove Query works against *HotStore. The whole @@ -1162,12 +1142,13 @@ func TestUnionSlots(t *testing.T) { // // - match-all asc → streamRange + cold FetchRange // - match-all desc + cap → streamRange top-down, slices.Backward -// - single-filter (contractID) → LookupKeys + streamUnion asc -// - multi-term filter (AND) → FastAnd over multiple cold bitmaps -// - cross-filter (OR) → FastOr across filters -// - ledger range + filter → roaring.And with the range bitmap -// - descending + range + cap → ReverseIterator on cold-derived -// union, single-filter And path +// - single-filter (contractID) → one LookupKeys term, then the +// ascending slab walk +// - multi-term filter (AND) → one FastAnd per plan over cold bitmaps +// - cross-filter (OR) → the in-place Or across plans +// - ledger range + filter → the window clipped to each slab's +// id range +// - descending + range + cap → the slab walk run high to low // // What we don't replay against cold: // - The mirror-poisoning collision test (mutating an mmap'd cold @@ -1625,21 +1606,11 @@ func TestMatches_EmptyStreams(t *testing.T) { wholeChunk(t, fx.store), false)) } -// TestMatches_WindowANDLeavesBorrowedBitmapUntouched pins the -// singleFilter branch of the window AND: a single-constraint filter -// borrows the hot mirror's bitmap directly from LookupKeys, and the -// narrowing AND must allocate a fresh result rather than shrink the -// mirror's live state in place. -// -// The borrow only exists for DENSE terms (the mirror's sparse mode -// materializes a fresh bitmap per Get, which no mutation can corrupt), -// so the term is first promoted past the mirror's promotion threshold -// with injected ids outside the query window (never fetched). -// TestQuery_DoesNotMutateMirrorBitmaps cannot catch the mutation: -// its filters carry two constraints (FastAnd-owned inputs) and it -// compares only cardinality over whole-chunk ranges, where the AND is -// a no-op. -func TestMatches_WindowANDLeavesBorrowedBitmapUntouched(t *testing.T) { +// TestMatches_LeavesSharedSnapshotUntouched pins that a narrowing window over +// a single-term filter leaves the hot mirror's shared bitmap as it was. Only +// dense terms are shared, so the term is first promoted with injected ids +// above the query window. +func TestMatches_LeavesSharedSnapshotUntouched(t *testing.T) { fx := newQueryFixture(t) key := ComputeTermKey(fx.contractA[:], FieldContractID) // Promote contract A's term (real matches: ids 0, 1, 4) to dense @@ -1653,18 +1624,16 @@ func TestMatches_WindowANDLeavesBorrowedBitmapUntouched(t *testing.T) { "fixture sanity: the term must be dense so LookupKeys borrows") snapshot := before.Clone() - // Single filter, single constraint, narrowing range: the borrowed - // path with an AND that actually removes ids. + // A narrowing window over a single-term filter. got := collectMatches(t, fx.store, []Filter{{ContractID: fx.contractA[:]}}, IDRange{Start: 0, End: 2}, false) assert.Equal(t, []uint32{0, 1}, matchOrdinals(got)) after := lookupOne(t, fx.store, key) assert.True(t, snapshot.Equals(after), - "the window AND must not mutate the mirror's term bitmap in place") + "Matches must not mutate the mirror's shared term bitmap") - // End to end: the same filter over the whole chunk still sees the - // ids an in-place AND would have destroyed. + // The same filter over the whole chunk still sees every id. full := collectMatches(t, fx.store, []Filter{{ContractID: fx.contractA[:]}}, wholeChunk(t, fx.store), false) assert.Equal(t, []uint32{0, 1, 4}, matchOrdinals(full)) @@ -1779,3 +1748,28 @@ func TestCountDistinctTerms(t *testing.T) { {ContractID: cid, TopicCount: TopicCountFilter{Count: 1}}, }), "topic-count buckets are not value terms and are not counted") } + +// The first-batch hint contract: a positive hint sizes the first fetch in +// full, since every caller passes a validated page size, and later batches +// use the default. +func TestBatchSizes(t *testing.T) { + first, rest := batchSizes(0) + require.Equal(t, matchBatchSize, first) + require.Equal(t, matchBatchSize, rest) + + first, rest = batchSizes(-3) + require.Equal(t, matchBatchSize, first) + require.Equal(t, matchBatchSize, rest) + + first, rest = batchSizes(7) + require.Equal(t, 7, first) + require.Equal(t, matchBatchSize, rest) + + first, rest = batchSizes(1000) + require.Equal(t, 1000, first, "a page-sized hint is the first fetch size") + require.Equal(t, matchBatchSize, rest) + + first, rest = batchSizes(10_000) + require.Equal(t, 10_000, first, "a v1 page-sized hint is honored in full") + require.Equal(t, matchBatchSize, rest) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/matches_differential_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/matches_differential_test.go new file mode 100644 index 000000000..677f5b24c --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/matches_differential_test.go @@ -0,0 +1,271 @@ +package event + +// The in-memory chunk the match-path tests run against, served through +// LookupKeys, and the borrow-safety gate over it. + +import ( + "context" + "errors" + "iter" + "math/rand" + "testing" + + "github.com/RoaringBitmap/roaring/v2" + "github.com/stretchr/testify/require" + + protocol "github.com/stellar/go-stellar-sdk/protocols/rpc" + "github.com/stellar/go-stellar-sdk/xdr" + + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/chunk" +) + +// diffCorpus is an in-memory chunk with one distinct event per id. +type diffCorpus struct { + raw [][]byte + mirror *ConcurrentBitmaps +} + +// diffReader serves the corpus through LookupKeys, the seam Matches reads the +// index by. +type diffReader struct{ c *diffCorpus } + +func (r diffReader) ChunkID() chunk.ID { return chunk.ID(0) } +func (r diffReader) EventCount() (uint32, error) { return uint32(len(r.c.raw)), nil } + +func (r diffReader) Offsets() (*LedgerOffsets, error) { + return nil, errors.New("diffReader: Offsets is not part of the match path") +} + +func (r diffReader) LookupKeys(_ context.Context, keys []TermKey) ([]*roaring.Bitmap, error) { + out := make([]*roaring.Bitmap, len(keys)) + for i, k := range keys { + bm, err := r.c.mirror.Get(k) + if err != nil { + return nil, err + } + out[i] = bm + } + return out, nil +} + +func (r diffReader) FetchEvents(_ context.Context, ids []uint32) ([]Payload, error) { + // A dedup bug in the union surfaces here, not as a doubled result. + if err := validateSortedEventIDs(ids); err != nil { + return nil, err + } + out := make([]Payload, len(ids)) + for i, id := range ids { + out[i] = Payload{ContractEventBytes: r.c.raw[id]} + } + return out, nil +} + +func (r diffReader) FetchRange(_ context.Context, start, count uint32) iter.Seq2[Payload, error] { + return func(yield func(Payload, error) bool) { + total, _ := r.EventCount() + if err := validateFetchRange(start, count, total, r.ChunkID()); err != nil { + yield(Payload{}, err) + return + } + for id := start; id < start+count; id++ { + if !yield(Payload{ContractEventBytes: r.c.raw[id]}, nil) { + return + } + } + } +} + +func (r diffReader) All(ctx context.Context) iter.Seq2[Payload, error] { + total, _ := r.EventCount() + return r.FetchRange(ctx, 0, total) +} + +var _ Reader = diffReader{} + +// diffVocab is the closed vocabulary the corpus and the random filters share. +type diffVocab struct { + contracts [][]byte + topics []xdr.ScVal + topicRaw [][]byte + types []xdr.ContractEventType +} + +func newDiffVocab(tb testing.TB) *diffVocab { + tb.Helper() + v := &diffVocab{types: []xdr.ContractEventType{ + xdr.ContractEventTypeSystem, + xdr.ContractEventTypeContract, + xdr.ContractEventTypeDiagnostic, + }} + for i := range 4 { + cid := xdr.ContractId{0: byte(0xC0 + i)} + v.contracts = append(v.contracts, cid[:]) + } + for _, name := range []string{"alpha", "beta", "gamma", "delta", "epsilon"} { + sym := xdr.ScSymbol(name) + val := xdr.ScVal{Type: xdr.ScValTypeScvSymbol, Sym: &sym} + raw, err := val.MarshalBinary() + require.NoError(tb, err) + v.topics, v.topicRaw = append(v.topics, val), append(v.topicRaw, raw) + } + return v +} + +func newDiffCorpus(t *testing.T, rng *rand.Rand, v *diffVocab, n int) *diffCorpus { + t.Helper() + c := &diffCorpus{mirror: NewConcurrentBitmapsFromBitmaps(NewBitmaps())} + for id := range n { + var cid xdr.ContractId + copy(cid[:], v.contracts[rng.Intn(len(v.contracts))]) + nTopics := rng.Intn(protocol.MaxTopicCount + 2) + topics := make([]xdr.ScVal, 0, nTopics) + for range nTopics { + topics = append(topics, v.topics[rng.Intn(len(v.topics))]) + } + sym := xdr.ScSymbol("data") + ev := xdr.ContractEvent{ + ContractId: &cid, + Type: v.types[rng.Intn(len(v.types))], + Body: xdr.ContractEventBody{ + V: 0, + V0: &xdr.ContractEventV0{ + Topics: topics, + Data: xdr.ScVal{Type: xdr.ScValTypeScvSymbol, Sym: &sym}, + }, + }, + } + raw, err := ev.MarshalBinary() + require.NoError(t, err) + c.raw = append(c.raw, raw) + keys, err := TermsForBytes(raw) + require.NoError(t, err) + for _, k := range keys { + c.mirror.AddTo(k, uint32(id)) + } + } + return c +} + +// randomFilters builds a filter list over the shared vocabulary, including the +// unconstrained shape that routes to the match-all path. +func randomFilters(rng *rand.Rand, v *diffVocab) []Filter { + filters := make([]Filter, 0, 3) + for range 1 + rng.Intn(3) { + var f Filter + if rng.Intn(3) > 0 { + f.ContractID = v.contracts[rng.Intn(len(v.contracts))] + } + if rng.Intn(3) == 0 { + et := xdr.ContractEventTypeContract + if rng.Intn(2) == 0 { + et = xdr.ContractEventTypeSystem + } + f.EventType = &et + } + for pos := range min(3, protocol.MaxTopicCount) { + if rng.Intn(4) == 0 { + f.Topics[pos] = v.topicRaw[rng.Intn(len(v.topicRaw))] + } + } + if rng.Intn(3) == 0 { + f.TopicCount = TopicCountFilter{ + Count: rng.Intn(protocol.MaxTopicCount + 1), + Exact: rng.Intn(2) == 0, + } + } + filters = append(filters, f) + } + if rng.Intn(20) == 0 { + return nil // the empty-slice match-all shape + } + return filters +} + +func collectOrdinals(t *testing.T, r Reader, filters []Filter, w IDRange, desc bool) []uint32 { + t.Helper() + var out []uint32 + for m, err := range Matches(context.Background(), r, filters, w, desc, 0) { + require.NoError(t, err) + out = append(out, m.Ordinal) + } + return out +} + +// The match path holds mirror snapshots across a whole walk while AddTo +// publishes new termStates on the same keys, sparse-to-dense promotion +// included. Under -race any write reaching a held snapshot fails the run; +// without it, the identity check pins that a pinned window ignores later ingest. +func TestMatches_ConcurrentIngestBorrowSafety(t *testing.T) { + rng := rand.New(rand.NewSource(20260830)) + v := newDiffVocab(t) + const corpusSize = 400 + const pinned = corpusSize / 2 + + corpus := &diffCorpus{mirror: NewConcurrentBitmapsFromBitmaps(NewBitmaps())} + keysByID := make([][]TermKey, corpusSize) + for id := range corpusSize { + var cid xdr.ContractId + copy(cid[:], v.contracts[rng.Intn(len(v.contracts))]) + topics := make([]xdr.ScVal, 0, 3) + for range 1 + rng.Intn(3) { + topics = append(topics, v.topics[rng.Intn(len(v.topics))]) + } + sym := xdr.ScSymbol("data") + ev := xdr.ContractEvent{ + ContractId: &cid, + Type: v.types[rng.Intn(len(v.types))], + Body: xdr.ContractEventBody{ + V: 0, + V0: &xdr.ContractEventV0{ + Topics: topics, + Data: xdr.ScVal{Type: xdr.ScValTypeScvSymbol, Sym: &sym}, + }, + }, + } + raw, err := ev.MarshalBinary() + require.NoError(t, err) + corpus.raw = append(corpus.raw, raw) + keys, err := TermsForBytes(raw) + require.NoError(t, err) + keysByID[id] = keys + } + // Only the pinned window is indexed up front; the writer feeds the rest + // live, and most keys cross the promotion threshold mid-run. + for id := range pinned { + for _, k := range keysByID[id] { + corpus.mirror.AddTo(k, uint32(id)) + } + } + + r := diffReader{corpus} + et := xdr.ContractEventTypeContract + filters := []Filter{ + {ContractID: v.contracts[0]}, + {Topics: [protocol.MaxTopicCount][]byte{0: v.topicRaw[1]}, EventType: &et}, + {TopicCount: TopicCountFilter{Count: 2}}, + } + window := IDRange{Start: 0, End: pinned} + want := collectOrdinals(t, r, filters, window, false) + require.NotEmpty(t, want, "fixture sanity: the pinned window must match something") + + done := make(chan struct{}) + go func() { + defer close(done) + for id := pinned; id < corpusSize; id++ { + for _, k := range keysByID[id] { + corpus.mirror.AddTo(k, uint32(id)) + } + } + }() + for { + select { + case <-done: + require.Equal(t, want, collectOrdinals(t, r, filters, window, false), + "pinned window changed after ingest completed") + return + default: + require.Equal(t, want, collectOrdinals(t, r, filters, window, false), + "pinned window changed mid-ingest") + } + } +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/roaring_contract_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/roaring_contract_test.go new file mode 100644 index 000000000..25d637d37 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/roaring_contract_test.go @@ -0,0 +1,341 @@ +package event + +// roaring_contract_test.go pins the roaring properties the event index rests +// on. ConcurrentBitmaps.Get publishes one bitmap to every concurrent reader +// and never mutates it, and the match path hands those bitmaps to FastAnd, +// NextValue and PreviousValue from many goroutines at once. That is sound +// only while FastAnd reads its arguments and returns containers that share +// no storage with them, and the searches are read-only, inclusive of the +// target, and return -1 for none. The tests are written against roaring +// alone, so they are the first thing to run on a version bump. +// +// Pinned version: github.com/RoaringBitmap/roaring/v2 v2.26.0. + +import ( + "math/rand" + "slices" + "sync" + "testing" + + "github.com/RoaringBitmap/roaring/v2" + "github.com/stretchr/testify/require" +) + +// sharedBitmaps returns bitmaps in the shape the index publishes: the +// copy-on-write Clone that denseState.snapshot hands to readers, sharing its +// containers with the writer's bitmap. The set spans all three container +// kinds and several high keys, since the roaring paths branch on both. +func sharedBitmaps(t *testing.T) []*roaring.Bitmap { + t.Helper() + rng := rand.New(rand.NewSource(20260909)) + + // Array container: a few hundred scattered ids inside one high key. + sparse := roaring.New() + for range 400 { + sparse.Add(uint32(rng.Intn(1 << 16))) + } + // Bitmap container: dense enough that roaring stores it as a bitset. + dense := roaring.New() + for range 40_000 { + dense.Add(uint32(rng.Intn(1 << 16))) + } + // Run container: contiguous stretches, RunOptimize'd so the container + // really is a run and not an array or a bitmap. + runs := roaring.New() + for i := range uint32(50) { + runs.AddRange(uint64(i*1000), uint64(i*1000+600)) + } + runs.RunOptimize() + // Containers under several high keys, so the aggregation's key-advance + // paths run too. + wide := roaring.New() + for key := range uint32(4) { + for i := range uint32(5000) { + wide.Add(key<<16 + i*7) + } + } + + out := make([]*roaring.Bitmap, 0, 4) + for _, wbm := range []*roaring.Bitmap{sparse, dense, runs, wide} { + wbm.SetCopyOnWrite(true) + snap := wbm.Clone() + require.True(t, snap.GetCopyOnWrite(), + "fixture: the published snapshot must carry the copy-on-write mark") + out = append(out, snap) + } + return out +} + +func bitmapBytes(t *testing.T, bm *roaring.Bitmap) []byte { + t.Helper() + b, err := bm.ToBytes() + require.NoError(t, err) + return b +} + +func bitmapImages(t *testing.T, bms []*roaring.Bitmap) [][]byte { + t.Helper() + out := make([][]byte, len(bms)) + for i, bm := range bms { + out[i] = bitmapBytes(t, bm) + } + return out +} + +func requireUnchanged(t *testing.T, bms []*roaring.Bitmap, before [][]byte, op string) { + t.Helper() + for i, bm := range bms { + require.Equal(t, before[i], bitmapBytes(t, bm), + "%s mutated shared argument %d: roaring can no longer be handed "+ + "denseState.snapshot bitmaps at this version", op, i) + } +} + +// freshRange is the slab window the match path hands to FastAnd: one +// contiguous id range, in a bitmap the test may write. +func freshRange(lo, hi uint64) *roaring.Bitmap { + bm := roaring.New() + bm.AddRange(lo, hi) + return bm +} + +// scratchBitmap is a bitmap the test owns outright and may write: one array +// container of scattered ids. +func scratchBitmap() *roaring.Bitmap { + scratch := roaring.New() + for i := range uint32(3000) { + scratch.Add(i * 3) + } + return scratch +} + +// windowRanges are slab windows that take different paths inside roaring: a +// whole high key, a partial one, one across a key boundary, a span over +// several keys, one the inputs do not reach, and the four-id span that is the +// smallest range roaring keeps as a run. +var windowRanges = []struct { + name string + lo, hi uint64 +}{ + {"whole key", 0, 1 << 16}, + {"partial key", 1000, 40_000}, + {"key boundary", 65_000, 66_000}, + {"several keys", 0, 4 << 16}, + {"disjoint key", 40 << 16, 41 << 16}, + {"four ids", 8000, 8004}, +} + +// TestRoaringContract_AggregationDoesNotMutateInputs pins that FastAnd +// leaves shared arguments byte-identical and equals the pairwise And chain, +// for every non-empty subset of the shared bitmaps beside the window: each +// container kind meets the window alone, in pairs and all together, which +// is every kernel a plan of one to seven terms can reach. +func TestRoaringContract_AggregationDoesNotMutateInputs(t *testing.T) { + t.Parallel() + + for _, w := range windowRanges { + t.Run(w.name, func(t *testing.T) { + t.Parallel() + shared := sharedBitmaps(t) + for pick := 1; pick < 1< 0 { + wantPrev = int64(ids[at-1]) + } + + require.Equal(t, wantNext, shared.NextValue(target), + "bitmap %d: NextValue(%d) must be the first id at or above the target", + i, target) + require.Equal(t, wantPrev, shared.PreviousValue(target), + "bitmap %d: PreviousValue(%d) must be the last id at or below the target", + i, target) + } + + require.Equal(t, before, bitmapBytes(t, shared), + "bitmap %d: a value search mutated the bitmap it searched: the slab "+ + "walk cannot run them on denseState.snapshot bitmaps at this version", i) + } + + empty := roaring.New() + require.Equal(t, int64(-1), empty.NextValue(0), + "an empty bitmap must report no next value") + require.Equal(t, int64(-1), empty.PreviousValue(1<<20), + "an empty bitmap must report no previous value") +} + +// searchTargets returns each id, its neighbors, the container boundaries the +// ids span, and the ends of the uint32 range. +func searchTargets(ids []uint32) []uint32 { + out := []uint32{0, 1<<32 - 1} + for _, id := range ids { + out = append(out, id) + if id > 0 { + out = append(out, id-1) + } + if id < 1<<32-1 { + out = append(out, id+1) + } + } + for key := range uint32(6) { + out = append(out, key<<16) + if key > 0 { + out = append(out, key<<16-1) + } + } + slices.Sort(out) + return slices.Compact(out) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match.go new file mode 100644 index 000000000..51aec3a04 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match.go @@ -0,0 +1,391 @@ +package event + +// slab_match.go produces the candidate ids behind Matches. The window is +// walked one slab at a time, 65536 ids, the span of one roaring container, +// and the whole filter algebra is evaluated inside each slab. Direction is +// only the walk order: ascending walks slabs low to high and reads each +// result forward, descending walks high to low and reads backward. +// +// Work is lazy per slab. A consumer that stops after one page has paid for +// the slabs that page spans, one container per term per plan each. Slabs +// that can hold no candidate are skipped: before evaluating a slab the walk +// asks the term bitmaps where the next candidate can be and jumps there (see +// the bound helpers below). +// +// The term bitmaps come from the single Reader.LookupKeys call at query start +// and are held for the whole walk. They are read-only and may be snapshots +// shared with other readers. FastAnd reads its arguments and returns fresh +// containers, which roaring_contract_test.go pins against the pinned roaring +// version; the only bitmaps this file mutates are the ones it builds for a +// slab and the results FastAnd hands back. Because the lookup is a +// point-in-time image, ids ingested during the walk are invisible to it, as +// IDRange's snapshot-isolation contract already requires. + +import ( + "cmp" + "context" + "slices" + + "github.com/RoaringBitmap/roaring/v2" +) + +// slabShift is the slab width as a power of two; 1<<16 is one roaring +// container. A var so tests can shrink it. It never changes what a stream +// yields. +// +//nolint:gochecknoglobals // test seam; production never writes it +var slabShift uint = 16 + +// slabPlan is one plan's term bitmaps, held for the whole query and ordered +// rarest first. +type slabPlan []*roaring.Bitmap + +// resolveSlabPlans resolves every plan's terms, drops the plans that name +// a term absent from the index, since such a conjunction matches nothing, +// and orders each survivor's terms rarest first, the order the pinned +// roaring's FastAnd intersects them in. A present but empty term keeps its +// plan. +func resolveSlabPlans(plans []termPlan, sources []*roaring.Bitmap) []slabPlan { + // A term's cardinality is counted once, however many plans name it: on + // run containers it walks every run. + cards := make([]uint64, len(sources)) + for i, bm := range sources { + if bm != nil { + cards[i] = bm.GetCardinality() + } + } + absent := func(slot int) bool { return sources[slot] == nil } + out := make([]slabPlan, 0, len(plans)) + for _, plan := range plans { + if slices.ContainsFunc(plan, absent) { + continue + } + slots := slices.Clone(plan) + slices.SortStableFunc(slots, func(a, b int) int { + return cmp.Compare(cards[a], cards[b]) + }) + p := make(slabPlan, 0, len(slots)) + for _, slot := range slots { + p = append(p, sources[slot]) + } + out = append(out, p) + } + return out +} + +// eval returns p's matches inside slab, the bitmap of one slab's ids, as a +// fresh bitmap the caller owns, or nil when there are none. The slab and +// the plan's terms go to roaring in one FastAnd call, the shape in which a +// count-first FastAnd can reject an empty intersection without allocating +// for it; the pinned v2.26.0 still intersects pairwise and allocates the +// first intermediate, so that saving arrives with the roaring bump. +func (p slabPlan) eval(slab *roaring.Bitmap) *roaring.Bitmap { + ops := make([]*roaring.Bitmap, 0, len(p)+1) + ops = append(ops, slab) + ops = append(ops, p...) + res := roaring.FastAnd(ops...) + if res.IsEmpty() { + return nil + } + return res +} + +// Bounds. A slab is skipped only on a bound proved from the term bitmaps +// themselves. Ascending, at position pos: +// +// - a term's bound is NextValue(pos), the first id at or above pos it +// holds; a term holding none proves its plan matches nothing from pos on; +// - a plan's bound is the largest of its terms' bounds, since a candidate +// satisfies every term; +// - the query's bound is the smallest of the live plans' bounds. +// +// Descending mirrors this with PreviousValue and the min and max swapped. +// The post-filter only drops candidates, so a bound proved on the index +// bounds the stream. NextValue and PreviousValue are inclusive of the target, +// return -1 for none, and do not write the bitmap they search; +// roaring_contract_test.go pins all three properties. + +// boundRetired is the bound of a plan whose terms ran out ahead of the +// cursor: the cursor never comes back, so the plan is dropped for the rest +// of the walk. It is roaring's own "none" and sorts below every id, so the +// "outside this slab" test covers it. +const boundRetired = int64(-1) + +// nextBound is the plan's bound: the largest of its terms' first ids at or +// above pos, or boundRetired as soon as one term holds none, which proves +// the plan is done. +func (p slabPlan) nextBound(pos uint32) int64 { + bound := int64(pos) + for _, bm := range p { + v := bm.NextValue(pos) + if v < 0 { + return boundRetired + } + bound = max(bound, v) + } + return bound +} + +// prevBound is nextBound descending: the smallest of the terms' last ids at +// or below pos. +func (p slabPlan) prevBound(pos uint32) int64 { + bound := int64(pos) + for _, bm := range p { + v := bm.PreviousValue(pos) + if v < 0 { + return boundRetired + } + bound = min(bound, v) + } + return bound +} + +// slabStepper walks one query's slabs in emission order, evaluating a slab +// only when the consumer has drained the previous one. +type slabStepper struct { + plans []slabPlan + window IDRange + desc bool + + // bounds holds, per plan, the id its next candidate is proved to lie at + // or past (at or before, descending), or boundRetired once it has none + // left. A bound is monotone in the walk direction, so one proved earlier + // still holds at every position up to it: it is re-proved only once the + // cursor reaches it, and a plan whose bound lies past the current slab + // is not evaluated there. A held bound can be weaker than a fresh one, + // which costs an evaluated slab, never a match. Bounds start at the + // cursor so the first step proves them all. + bounds []int64 + + // cursor is the next unevaluated boundary: the inclusive low bound + // ascending, the exclusive high bound descending. + cursor uint32 + done bool + + // cur is the current slab's result, held only for its iterator. + cur *roaring.Bitmap + asc roaring.ManyIntIterable + rev roaring.IntIterable +} + +func newSlabStepper( + plans []termPlan, sources []*roaring.Bitmap, window IDRange, descending bool, +) *slabStepper { + s := &slabStepper{ + plans: resolveSlabPlans(plans, sources), + window: window, + desc: descending, + } + if descending { + s.cursor = window.End + } else { + s.cursor = window.Start + } + s.bounds = make([]int64, len(s.plans)) + for i := range s.bounds { + s.bounds[i] = int64(s.cursor) + } + return s +} + +// seekAsc returns the lowest position at or above pos where some plan can +// still match, or false when none can inside the window. A plan is asked +// again only once pos has reached its held bound; a plan with no bound left +// is retired rather than ending the walk. +func (s *slabStepper) seekAsc(pos uint32) (uint32, bool) { + var best int64 + found := false + for i := range s.plans { + if s.bounds[i] == boundRetired { + continue + } + if s.bounds[i] <= int64(pos) { + s.bounds[i] = s.plans[i].nextBound(pos) + if s.bounds[i] == boundRetired { + continue + } + } + if !found || s.bounds[i] < best { + best, found = s.bounds[i], true + } + } + if !found || best >= int64(s.window.End) { + return 0, false + } + return uint32(best), true +} + +// seekDesc is seekAsc mirrored: the largest of the plans' bounds at or below +// hi-1, returned as the exclusive high bound the walk resumes at. +func (s *slabStepper) seekDesc(hi uint32) (uint32, bool) { + pos := hi - 1 + var best int64 + found := false + for i := range s.plans { + if s.bounds[i] == boundRetired { + continue + } + if s.bounds[i] >= int64(pos) { + s.bounds[i] = s.plans[i].prevBound(pos) + if s.bounds[i] == boundRetired { + continue + } + } + if !found || s.bounds[i] > best { + best, found = s.bounds[i], true + } + } + if !found || best < int64(s.window.Start) { + return 0, false + } + return uint32(best) + 1, true +} + +// nextBounds returns the next slab's [lo, hi), clipped to the window, in the +// walk's direction. The cursor first moves to the proved bound, so the slab +// returned holds the next possible candidate and is entered at that candidate +// rather than at its base. +func (s *slabStepper) nextBounds() (uint32, uint32, bool) { + if s.done { + return 0, 0, false + } + if s.desc { + // hi > window.Start here, because an empty window never reaches the + // stepper and the walk stops at window.Start, so seekDesc's hi-1 + // cannot underflow. + hi, ok := s.seekDesc(s.cursor) + if !ok { + s.done = true + return 0, 0, false + } + lo := s.window.Start + // The base of the slab holding hi-1. + if base := ((uint64(hi) - 1) >> slabShift) << slabShift; base > uint64(lo) { + lo = uint32(base) //nolint:gosec // base < hi <= MaxUint32 + } else { + s.done = true + } + s.cursor = lo + return lo, hi, true + } + lo, ok := s.seekAsc(s.cursor) + if !ok { + s.done = true + return 0, 0, false + } + hi := s.window.End + // The base of the slab above the one holding lo. + if next := (uint64(lo)>>slabShift + 1) << slabShift; next < uint64(hi) { + hi = uint32(next) //nolint:gosec // next < hi <= MaxUint32 + } else { + s.done = true + } + s.cursor = hi + return lo, hi, true +} + +// evalSlab unions the per-plan results for [lo, hi), or returns nil when +// the slab holds nothing. A plan whose bound lies outside the slab is +// skipped, since the bound already proves it matches nothing here. The +// results are this call's own bitmaps, so the union runs in place into the +// first of them. +func (s *slabStepper) evalSlab(lo, hi uint32) *roaring.Bitmap { + var acc *roaring.Bitmap + slab := roaring.New() + slab.AddRange(uint64(lo), uint64(hi)) + for i := range s.plans { + if b := s.bounds[i]; b < int64(lo) || b >= int64(hi) { + continue + } + bm := s.plans[i].eval(slab) + switch { + case bm == nil: + case acc == nil: + acc = bm + default: + acc.Or(bm) + } + } + return acc +} + +// ensureSlab advances to the next slab that holds a match, reporting false +// once the window is exhausted. +func (s *slabStepper) ensureSlab() bool { + for s.cur == nil { + lo, hi, ok := s.nextBounds() + if !ok { + return false + } + bm := s.evalSlab(lo, hi) + if bm == nil { + continue + } + s.cur = bm + if s.desc { + s.rev = bm.ReverseIterator() + } else { + s.asc = bm.ManyIterator() + } + } + return true +} + +func (s *slabStepper) dropSlab() { + s.cur, s.asc, s.rev = nil, nil, nil +} + +// appendUpTo appends at most n ids in emission order to dst, evaluating slabs +// as it fills, and returns the extended slice. A result shorter than n means +// the query is exhausted. +func (s *slabStepper) appendUpTo(dst []uint32, n int) []uint32 { + for len(dst) < n { + if !s.ensureSlab() { + return dst + } + if s.desc { + for len(dst) < n && s.rev.HasNext() { + dst = append(dst, s.rev.Next()) + } + if !s.rev.HasNext() { + s.dropSlab() + } + continue + } + // NextMany fills the tail in bulk and returns short only at the end + // of the bitmap, so a short fill is this slab's last id. + want := n - len(dst) + base := len(dst) + dst = slices.Grow(dst, want)[:base+want] + got := s.asc.NextMany(dst[base:]) + dst = dst[:base+got] + if got < want { + s.dropSlab() + } + } + return dst +} + +// streamSlabs is the streaming loop for both directions: fill one batch of +// candidate ordinals from the stepper, fetch, post-filter, yield. +func streamSlabs( + ctx context.Context, r Reader, filters []Filter, st *slabStepper, + descending bool, firstBatch int, yield func(Match, error) bool, +) { + batch, rest := batchSizes(firstBatch) + ids := make([]uint32, 0, batch) + for { + if err := ctx.Err(); err != nil { + yield(Match{}, err) + return + } + ids = st.appendUpTo(ids[:0], batch) + batch = rest + if len(ids) == 0 { + return + } + if !emitBatch(ctx, r, filters, ids, descending, yield) { + return + } + } +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match_test.go new file mode 100644 index 000000000..f54040b2a --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/slab_match_test.go @@ -0,0 +1,850 @@ +package event + +// slab_match_test.go covers the slab engine over a corpus built to hold each +// named query shape by construction, and again at slab widths narrow enough +// to put many empty slabs between consecutive ids. Every case is checked +// against an answer computed without the index: postFilter over every +// ordinal in the corpus, clipped to the window, the direction and the page. + +import ( + "cmp" + "context" + "iter" + "math/rand" + "slices" + "testing" + + "github.com/RoaringBitmap/roaring/v2" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + protocol "github.com/stellar/go-stellar-sdk/protocols/rpc" + "github.com/stellar/go-stellar-sdk/xdr" +) + +// drainMatches drains seq into a slice, stopping after limit items when limit +// is positive. +func drainMatches(tb testing.TB, seq iter.Seq2[Match, error], limit int) []Match { + tb.Helper() + out := []Match{} + for m, err := range seq { + require.NoError(tb, err) + out = append(out, m) + if limit > 0 && len(out) == limit { + break + } + } + return out +} + +// ───────────────────────── the answer without the index ───────────────────── + +// matchingEvents is the corpus's whole answer to filters, computed by running +// the post-filter over every ordinal rather than by asking the index. An +// empty filter slice selects everything. +func matchingEvents(tb testing.TB, c *diffCorpus, filters []Filter) []Match { + tb.Helper() + ids := make([]uint32, len(c.raw)) + payloads := make([]Payload, len(c.raw)) + for id := range c.raw { + ids[id] = uint32(id) + payloads[id] = Payload{ContractEventBytes: c.raw[id]} + } + if len(filters) == 0 { + out := make([]Match, len(ids)) + for i := range ids { + out[i] = Match{Payload: payloads[i], Ordinal: ids[i]} + } + return out + } + out, err := postFilter(payloads, ids, filters) + require.NoError(tb, err) + return out +} + +// expectedStream clips the whole-corpus answer to one query: the window, the +// direction, and the page the consumer stops on. all is ascending by ordinal, +// so the window is a slice of it. +func expectedStream(all []Match, w IDRange, desc bool, limit int) []Match { + byOrdinal := func(m Match, id uint32) int { return cmp.Compare(m.Ordinal, id) } + lo, _ := slices.BinarySearchFunc(all, w.Start, byOrdinal) + hi, _ := slices.BinarySearchFunc(all, w.End, byOrdinal) + in := all[lo:hi] + + n := len(in) + if limit > 0 && n > limit { + n = limit + } + out := make([]Match, 0, n) + if !desc { + return append(out, in[:n]...) + } + for i := len(in) - 1; len(out) < n; i-- { + out = append(out, in[i]) + } + return out +} + +// queryCase is one (filters, window, direction, limit) query. +type queryCase struct { + name string + filters []Filter + window IDRange + desc bool + limit int +} + +// requireStream drives one case through Matches and requires the stream the +// corpus says it must be. +func requireStream(tb testing.TB, r Reader, all []Match, c queryCase) []Match { + tb.Helper() + got := drainMatches(tb, + Matches(context.Background(), r, c.filters, c.window, c.desc, c.limit), c.limit) + require.Equal(tb, expectedStream(all, c.window, c.desc, c.limit), got, + "case %q window %v desc=%v limit=%d", c.name, c.window, c.desc, c.limit) + return got +} + +// ───────────────────────── the shaped corpus ───────────────────────── + +// slabShape is one distinct event in the shaped corpus; ids map to shapes by +// rule, so the corpus costs a few dozen XDR marshals. +type slabShape struct { + contract int + topic0 int + topic1 int // -1 for an event carrying one topic + evType int +} + +// shapedFixture is the shaped corpus plus the vocabulary its filters are +// written against. +type shapedFixture struct { + corpus *diffCorpus + vocab *diffVocab + // thin holds the ids where the two fat terms of the overlap shape meet. + thin []uint32 + // rareContract and rareTopic hold the ids carrying the sparse terms. + rareContract []uint32 + rareTopic []uint32 +} + +// shapeFor is the corpus's id to event rule: +// +// - contracts 0 and 1 split the corpus in half: two dense terms; +// - contract 2 is carried by a handful of ids: a sparse term; +// - topic 0 tracks the id's parity except on the thin-overlap ids, so +// contract 0 AND the odd topic is two chunk-sized terms meeting on a few; +// - a second topic appears on a handful of ids, so the topic-count buckets +// are one dense and one sparse; +// - the system event type stays sparse. +func (f *shapedFixture) shapeFor(id uint32) slabShape { + sh := slabShape{topic1: -1} + switch { + case slices.Contains(f.rareContract, id): + sh.contract = 2 + case id%2 == 0: + sh.contract = 0 + default: + sh.contract = 1 + } + if id%2 == 1 || slices.Contains(f.thin, id) { + sh.topic0 = 1 + } + if slices.Contains(f.rareTopic, id) { + sh.topic1 = 2 + } + if id%2000 == 0 { + sh.evType = 0 // system + } else { + sh.evType = 1 // contract + } + return sh +} + +// shapedCorpusSize spans slab 0 whole and part of slab 1, so every window +// bound below is a real slab-relative position rather than a synthetic one. +const shapedCorpusSize uint32 = 70_000 + +// newShapedFixture builds the shaped corpus. +func newShapedFixture(tb testing.TB) *shapedFixture { + tb.Helper() + const n = shapedCorpusSize + v := newDiffVocab(tb) + f := &shapedFixture{vocab: v} + // Thin-overlap ids: even, spread across the window, and deliberately + // sitting on both sides of a slab boundary. + for _, id := range []uint32{4, 30_000, 65_534, 65_536, 65_538, n - 2} { + if id < n { + f.thin = append(f.thin, id) + } + } + for _, id := range []uint32{1, 12_345, 65_535, 66_000, n - 1} { + if id < n { + f.rareContract = append(f.rareContract, id) + } + } + for id := uint32(0); id < n; id += n/37 + 1 { + f.rareTopic = append(f.rareTopic, id) + } + slices.Sort(f.thin) + slices.Sort(f.rareContract) + slices.Sort(f.rareTopic) + + raws := make(map[slabShape][]byte) + keys := make(map[slabShape][]TermKey) + c := &diffCorpus{ + raw: make([][]byte, n), + mirror: NewConcurrentBitmapsFromBitmaps(NewBitmaps()), + } + idsByKey := make(map[TermKey][]uint32) + for id := range n { + sh := f.shapeFor(id) + raw, ok := raws[sh] + if !ok { + raw = f.marshalShape(tb, sh) + ks, err := TermsForBytes(raw) + require.NoError(tb, err) + raws[sh], keys[sh] = raw, ks + } + c.raw[id] = raw + for _, k := range keys[sh] { + idsByKey[k] = append(idsByKey[k], id) + } + } + for k, ids := range idsByKey { + c.mirror.AddTo(k, ids...) + } + f.corpus = c + return f +} + +func (f *shapedFixture) marshalShape(tb testing.TB, sh slabShape) []byte { + tb.Helper() + var cid xdr.ContractId + copy(cid[:], f.vocab.contracts[sh.contract]) + topics := []xdr.ScVal{f.vocab.topics[sh.topic0]} + if sh.topic1 >= 0 { + topics = append(topics, f.vocab.topics[sh.topic1]) + } + sym := xdr.ScSymbol("data") + ev := xdr.ContractEvent{ + ContractId: &cid, + Type: f.vocab.types[sh.evType], + Body: xdr.ContractEventBody{ + V: 0, + V0: &xdr.ContractEventV0{ + Topics: topics, + Data: xdr.ScVal{Type: xdr.ScValTypeScvSymbol, Sym: &sym}, + }, + }, + } + raw, err := ev.MarshalBinary() + require.NoError(tb, err) + return raw +} + +// named query shapes over the shaped corpus. +func (f *shapedFixture) filterDenseOnly() []Filter { + return []Filter{{ContractID: f.vocab.contracts[0]}} +} + +func (f *shapedFixture) filterSparseOnly() []Filter { + return []Filter{{ContractID: f.vocab.contracts[2]}} +} + +func (f *shapedFixture) filterMixedGroups() []Filter { + // A sparse group AND a dense group inside one filter. + return []Filter{{ + ContractID: f.vocab.contracts[2], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[1]}, + }} +} + +func (f *shapedFixture) filterMixedOrGroup() []Filter { + // One group that ORs a dense bucket with a sparse one: the topic-count + // family, where "one topic" is the whole corpus and "two topics" is the + // handful of ids carrying the second topic. + return []Filter{{TopicCount: TopicCountFilter{Count: 1}}} +} + +func (f *shapedFixture) filterThinOverlap() []Filter { + // Two chunk-sized terms meeting on |thin| ids. + return []Filter{{ + ContractID: f.vocab.contracts[0], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[1]}, + }} +} + +func (f *shapedFixture) filterAbsentTerm() []Filter { + // contracts[3] is in the vocabulary but never used by the corpus. + return []Filter{{ContractID: f.vocab.contracts[3]}} +} + +func (f *shapedFixture) filterAbsentPlusPresent() []Filter { + return []Filter{ + {ContractID: f.vocab.contracts[3]}, + {ContractID: f.vocab.contracts[2]}, + } +} + +// Both positions of an absent group are named: an engine that stopped +// resolving at the first miss but kept the filter would only get the trailing +// shape wrong. +func (f *shapedFixture) filterAbsentGroupTrailing() []Filter { + // topicRaw[4] is in the vocabulary but no corpus event carries it. + return []Filter{{ + ContractID: f.vocab.contracts[0], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[4]}, + }} +} + +func (f *shapedFixture) filterAbsentGroupLeading() []Filter { + return []Filter{{ + ContractID: f.vocab.contracts[3], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[0]}, + }} +} + +// The corpus carries one and two topics, so "at least two" ORs one sparse +// bucket with three absent ones and "at least three" is absent in full. +func (f *shapedFixture) filterPartlyAbsentGroup() []Filter { + return []Filter{{TopicCount: TopicCountFilter{Count: 2}}} +} + +func (f *shapedFixture) filterWhollyAbsentGroup() []Filter { + return []Filter{{TopicCount: TopicCountFilter{Count: 3}}} +} + +func (f *shapedFixture) filterUnion() []Filter { + sysType := xdr.ContractEventTypeSystem + return []Filter{ + {ContractID: f.vocab.contracts[2]}, + {EventType: &sysType}, + { + ContractID: f.vocab.contracts[0], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[1]}, + }, + } +} + +// filterDeepAND names every constrainable field, so the AND runs five groups +// deep and meets on the handful of ids carrying a second topic. Exact keeps +// the count group in the plan. +func (f *shapedFixture) filterDeepAND() []Filter { + evType := xdr.ContractEventTypeContract + return []Filter{{ + ContractID: f.vocab.contracts[0], + EventType: &evType, + Topics: [protocol.MaxTopicCount][]byte{ + 0: f.vocab.topicRaw[0], + 1: f.vocab.topicRaw[2], + }, + TopicCount: TopicCountFilter{Count: 2, Exact: true}, + }} +} + +// filterWideUnion unions eight filters, mixing dense, sparse, absent and +// multi-group ones, so the union drops a filter and dedups shared ids. +func (f *shapedFixture) filterWideUnion() []Filter { + sysType := xdr.ContractEventTypeSystem + return []Filter{ + { + ContractID: f.vocab.contracts[0], + Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[1]}, + }, + {ContractID: f.vocab.contracts[1]}, + {ContractID: f.vocab.contracts[2]}, + {ContractID: f.vocab.contracts[3]}, + {EventType: &sysType}, + {Topics: [protocol.MaxTopicCount][]byte{0: f.vocab.topicRaw[0]}}, + {Topics: [protocol.MaxTopicCount][]byte{1: f.vocab.topicRaw[2]}}, + {TopicCount: TopicCountFilter{Count: 2, Exact: true}}, + } +} + +// namedShape is one query shape with the name its failures report under. +type namedShape struct { + name string + filters []Filter +} + +func (f *shapedFixture) namedShapes() []namedShape { + sysType := xdr.ContractEventTypeSystem + return []namedShape{ + {"dense only", f.filterDenseOnly()}, + {"sparse only", f.filterSparseOnly()}, + {"mixed groups", f.filterMixedGroups()}, + {"mixed or group", f.filterMixedOrGroup()}, + {"thin overlap", f.filterThinOverlap()}, + {"absent term", f.filterAbsentTerm()}, + {"absent plus present", f.filterAbsentPlusPresent()}, + {"absent group trailing", f.filterAbsentGroupTrailing()}, + {"absent group leading", f.filterAbsentGroupLeading()}, + {"partly absent group", f.filterPartlyAbsentGroup()}, + {"wholly absent group", f.filterWhollyAbsentGroup()}, + {"union of filters", f.filterUnion()}, + {"deep and", f.filterDeepAND()}, + {"wide union", f.filterWideUnion()}, + {"match all empty slice", nil}, + {"match all wildcard filter", []Filter{{}}}, + {"match all beside constrained", []Filter{{EventType: &sysType}, {}}}, + {"exact topic count", []Filter{{TopicCount: TopicCountFilter{Count: 2, Exact: true}}}}, + {"term and count range", []Filter{{ContractID: f.vocab.contracts[0], TopicCount: TopicCountFilter{Count: 1}}}}, + } +} + +// ───────────────────────── the shaped matrix ───────────────────────── + +func TestMatches_ShapedMatrix(t *testing.T) { + f := newShapedFixture(t) + r := diffReader{f.corpus} + + const slab = 1 << 16 + windows := []struct { + name string + w IDRange + }{ + {"whole corpus", IDRange{0, shapedCorpusSize}}, + {"empty at zero", IDRange{0, 0}}, + {"empty mid slab", IDRange{12_345, 12_345}}, + {"empty at boundary", IDRange{slab, slab}}, + {"first slab exactly", IDRange{0, slab}}, + {"second slab only", IDRange{slab, shapedCorpusSize}}, + {"cursor mid slab", IDRange{33_333, shapedCorpusSize}}, + {"cursor one below boundary", IDRange{slab - 1, shapedCorpusSize}}, + {"cursor on boundary", IDRange{slab, shapedCorpusSize}}, + {"cursor one above boundary", IDRange{slab + 1, shapedCorpusSize}}, + {"end one below boundary", IDRange{0, slab - 1}}, + {"end on boundary", IDRange{0, slab}}, + {"end one above boundary", IDRange{0, slab + 1}}, + {"single id at boundary", IDRange{slab, slab + 1}}, + {"straddles boundary", IDRange{slab - 3, slab + 3}}, + {"tail", IDRange{shapedCorpusSize - 5, shapedCorpusSize}}, + } + + // One page and a single item: both a page ending mid-slab and one ending + // on a slab's last id, at every window bound. + limits := []int{1, 1000} + + for _, sh := range f.namedShapes() { + all := matchingEvents(t, f.corpus, sh.filters) + for _, w := range windows { + for _, desc := range []bool{false, true} { + for _, limit := range limits { + requireStream(t, r, all, queryCase{ + name: sh.name + "/" + w.name, + filters: sh.filters, + window: w.w, + desc: desc, + limit: limit, + }) + } + } + } + } +} + +// The bounded matrix never reaches the end of a fat stream; this runs whole +// streams over the windows where the end differs. +func TestMatches_WholeStreams(t *testing.T) { + f := newShapedFixture(t) + r := diffReader{f.corpus} + const slab = 1 << 16 + + for _, sh := range f.namedShapes() { + all := matchingEvents(t, f.corpus, sh.filters) + for _, w := range []IDRange{ + {0, shapedCorpusSize}, + {slab, shapedCorpusSize}, + {slab - 3, slab + 3}, + } { + for _, desc := range []bool{false, true} { + requireStream(t, r, all, queryCase{ + name: sh.name, filters: sh.filters, window: w, desc: desc, + }) + } + } + } +} + +// The corpus must hold the shapes its filters are named for; drift in its +// rules would weaken the matrix without failing it. +func TestMatches_ShapedFixtureIsWhatItClaims(t *testing.T) { + f := newShapedFixture(t) + r := diffReader{f.corpus} + ctx := context.Background() + window := IDRange{0, shapedCorpusSize} + + card := func(filters []Filter) int { + return len(drainMatches(t, Matches(ctx, r, filters, window, false, 0), 0)) + } + + require.Greater(t, card(f.filterDenseOnly()), 30_000, "dense term must be chunk-sized") + require.Len(t, drainMatches(t, + Matches(ctx, r, f.filterSparseOnly(), window, false, 0), 0), len(f.rareContract), + "sparse term must hold exactly the rare ids") + require.Less(t, len(f.rareContract), promotionThreshold, + "the sparse term must stay under the promotion threshold") + require.Len(t, drainMatches(t, + Matches(ctx, r, f.filterThinOverlap(), window, false, 0), 0), len(f.thin), + "the thin overlap must be exactly the constructed ids") + require.Greater(t, len(f.thin), 1) + require.Less(t, len(f.thin), 10, "the overlap must be thin") + require.Zero(t, card(f.filterAbsentTerm()), "the absent term must select nothing") + require.Zero(t, card(f.filterAbsentGroupTrailing()), + "a filter whose second group is absent must select nothing, even though "+ + "its first group is chunk-sized") + require.Zero(t, card(f.filterAbsentGroupLeading()), + "a filter whose first group is absent must select nothing") + require.Zero(t, card(f.filterWhollyAbsentGroup()), + "a group whose every bucket is absent must select nothing") + require.Len(t, drainMatches(t, + Matches(ctx, r, f.filterPartlyAbsentGroup(), window, false, 0), 0), len(f.rareTopic), + "the partly-absent group must select exactly its one present bucket") + + // The high-arity shapes must select something, or the matrix passes + // vacuously. + require.Greater(t, card(f.filterDeepAND()), 10, + "the five-group AND must still select something") + require.Greater(t, card(f.filterWideUnion()), 30_000, + "the wide union must span the corpus") + + // Both sides of the overlap must be chunk-sized. + fat, err := r.LookupKeys(ctx, []TermKey{ + ComputeTermKey(f.vocab.contracts[0], FieldContractID), + ComputeTermKey(f.vocab.topicRaw[1], topicField(0)), + }) + require.NoError(t, err) + for i, bm := range fat { + require.NotNil(t, bm, "thin-overlap term %d must be indexed", i) + require.Greater(t, bm.GetCardinality(), uint64(shapedCorpusSize/3), + "thin-overlap term %d must be chunk-sized", i) + } +} + +// ───────────────────────── the randomized matrix ───────────────────────── + +// TestMatches_RandomizedAgainstPostFilter drives random filters, windows and +// page sizes over a random corpus at several slab widths. +func TestMatches_RandomizedAgainstPostFilter(t *testing.T) { + v := newDiffVocab(t) + const corpusSize = 300 + corpus := newDiffCorpus(t, rand.New(rand.NewSource(20260829)), v, corpusSize) + + // A small batch exercises batch seams on a small corpus. + defer func(n int) { matchBatchSize = n }(matchBatchSize) + matchBatchSize = 7 + defer func(s uint) { slabShift = s }(slabShift) + + r := diffReader{corpus} + // Every width must yield the same stream: 1, 2 and 4 put many slab seams + // inside the corpus, 8 leaves one, and 16 is the production width. + for _, shift := range []uint{1, 2, 4, 8, 16} { + slabShift = shift + rng := rand.New(rand.NewSource(int64(20260909 + shift))) + matched := 0 + for range 400 { + matched += randomizedTrial(t, r, corpus, v, rng, corpusSize) + } + require.Greater(t, matched, 2000, + "fixture sanity: randomized queries selected too little") + } +} + +// randomizedTrial runs one random query in both directions and returns how +// many matches the ascending run selected. +func randomizedTrial( + t *testing.T, r Reader, corpus *diffCorpus, v *diffVocab, rng *rand.Rand, + corpusSize int, +) int { + t.Helper() + filters := randomFilters(rng, v) + start := uint32(rng.Intn(corpusSize + 1)) + end := start + uint32(rng.Intn(corpusSize+1-int(start))) + limit := []int{0, 0, 1, 3, 17, 200}[rng.Intn(6)] + all := matchingEvents(t, corpus, filters) + + matched := 0 + for _, desc := range []bool{false, true} { + got := requireStream(t, r, all, queryCase{ + name: "randomized", + filters: filters, + window: IDRange{Start: start, End: end}, + desc: desc, + limit: limit, + }) + if !desc { + matched = len(got) + } + } + return matched +} + +// ───────────────────────── the candidate-set pin ───────────────────────── + +// fetchTracer records the ordinals of every FetchEvents call, in call order. +// Output equality cannot see an over-wide candidate set, since postFilter +// drops the extra fetches; recording them can. +type fetchTracer struct { + diffReader + + batches *[][]uint32 +} + +func (r fetchTracer) FetchEvents(ctx context.Context, ids []uint32) ([]Payload, error) { + *r.batches = append(*r.batches, slices.Clone(ids)) + return r.diffReader.FetchEvents(ctx, ids) +} + +var _ Reader = fetchTracer{} + +// TestMatches_FetchesOnlyTrueCandidates pins that the ordinals the engine +// fetches are exactly the query's matches in emission order, and that a +// consumer stopping after a page fetched only the batches it spans. Match-all +// shapes stream FetchRange and never reach the index. +func TestMatches_FetchesOnlyTrueCandidates(t *testing.T) { + f := newShapedFixture(t) + + const slab = 1 << 16 + shapes := [][]Filter{ + f.filterDenseOnly(), + f.filterSparseOnly(), + f.filterMixedGroups(), + f.filterMixedOrGroup(), + f.filterThinOverlap(), + f.filterAbsentTerm(), + f.filterAbsentGroupTrailing(), + f.filterAbsentGroupLeading(), + f.filterPartlyAbsentGroup(), + f.filterWhollyAbsentGroup(), + f.filterUnion(), + } + windows := []IDRange{ + {0, shapedCorpusSize}, + {slab, shapedCorpusSize}, + {slab - 3, slab + 3}, + {33_333, shapedCorpusSize}, + } + + for si, filters := range shapes { + all := matchingEvents(t, f.corpus, filters) + for _, w := range windows { + for _, desc := range []bool{false, true} { + for _, limit := range []int{0, 1, 1000} { + want := expectedStream(all, w, desc, 0) + fetched := traceFetches(t, f, filters, w, desc, limit) + + wantIDs := make([]uint32, 0, len(want)) + for _, m := range want { + wantIDs = append(wantIDs, m.Ordinal) + } + require.LessOrEqual(t, len(fetched), len(wantIDs), + "shape %d window %v desc=%v limit=%d: fetched an ordinal "+ + "outside the answer", si, w, desc, limit) + require.Equal(t, wantIDs[:len(fetched)], fetched, + "shape %d window %v desc=%v limit=%d: the fetched "+ + "candidates are not the answer's leading run", + si, w, desc, limit) + require.GreaterOrEqual(t, len(fetched), pageFloor(limit, len(wantIDs)), + "shape %d window %v desc=%v limit=%d: the page was served "+ + "without fetching enough candidates", si, w, desc, limit) + } + } + } + } +} + +// traceFetches drives one query and returns the ordinals it fetched in +// emission order; a descending batch is fetched flipped, so it is flipped +// back. +func traceFetches( + t *testing.T, f *shapedFixture, filters []Filter, w IDRange, desc bool, limit int, +) []uint32 { + t.Helper() + batches := [][]uint32{} + r := fetchTracer{diffReader{f.corpus}, &batches} + drainMatches(t, Matches(context.Background(), r, filters, w, desc, limit), limit) + + out := []uint32{} + for _, b := range batches { + if desc { + slices.Reverse(b) + } + out = append(out, b...) + } + return out +} + +// pageFloor is how many candidates a query must have fetched to have served +// its page: the page itself, or the whole answer when it is shorter. +func pageFloor(limit, answer int) int { + if limit <= 0 { + return answer + } + return min(limit, answer) +} + +// ───────────────────────── the planning step ───────────────────────── + +// termBitmap is a materialized term, the shape LookupKeys hands the planner. +func termBitmap(ids ...uint32) *roaring.Bitmap { + bm := roaring.New() + bm.AddMany(ids) + return bm +} + +// resolveSlabPlans keeps the plans whose every term is present, orders each +// one's terms rarest first, and holds the lookup's bitmaps rather than +// copying them. +func TestResolveSlabPlans(t *testing.T) { + sources := []*roaring.Bitmap{ + termBitmap(1, 2), + nil, // absent + termBitmap(2, 3, 4), + termBitmap(), // present, holding nothing + } + got := resolveSlabPlans([]termPlan{{0}, {2, 0}, {1, 2}, {3}, {1}}, sources) + require.Len(t, got, 3, "a plan naming an absent term matches nothing and is dropped") + require.Len(t, got[0], 1) + assert.Same(t, sources[0], got[0][0], "the lookup's bitmap is held, not copied") + require.Len(t, got[1], 2) + assert.Same(t, sources[0], got[1][0], "terms are ordered rarest first") + assert.Same(t, sources[2], got[1][1]) + require.Len(t, got[2], 1) + assert.Same(t, sources[3], got[2][0], "a present but empty term keeps its plan") +} + +// A topic-count range fans out to one plan per bucket, and a repeated filter +// adds no plan. +func TestPlanIndexTermsSplitsRanges(t *testing.T) { + cid := []byte{0xAB} + ranged := Filter{ContractID: cid, TopicCount: TopicCountFilter{Count: 2}} + plans, keys, matchAll := planIndexTerms([]Filter{ranged, ranged, {ContractID: cid}}) + require.False(t, matchAll) + buckets := len(TopicCountTermKeysAtLeast(2)) + require.Len(t, plans, buckets+1) + assert.Len(t, keys, buckets+1) + for _, p := range plans[:buckets] { + assert.Equal(t, 0, p[0], "the contract slot leads every plan of the range") + assert.Len(t, p, 2) + } + assert.Equal(t, termPlan{0}, plans[buckets]) +} + +// ───────────────────────── the skip ───────────────────────── + +// The skip is invisible in a stream, so this drives nextBounds directly and +// requires the exact slabs the walk opens. +func TestSlabStepperSkipsCandidateFreeSlabs(t *testing.T) { + defer func(s uint) { slabShift = s }(slabShift) + slabShift = 16 + const slab = 1 << 16 + const slabs = 10 + window := IDRange{0, slabs * slab} + + // A rare term in slabs 0, 3 and 9; another in slabs 1 and 6; a term + // holding every id in the window; and a term present but empty. + fat := roaring.New() + fat.AddRange(uint64(window.Start), uint64(window.End)) + sources := []*roaring.Bitmap{ + termBitmap(5, 3*slab+7, 9*slab+1), + termBitmap(slab+1, 6*slab+3), + fat, + termBitmap(), + } + + walk := func(plans []termPlan, desc bool) [][2]uint32 { + st := newSlabStepper(plans, sources, window, desc) + out := [][2]uint32{} + for { + lo, hi, ok := st.nextBounds() + if !ok { + return out + } + out = append(out, [2]uint32{lo, hi}) + } + } + + // The rare term alone opens its three slabs and no others, in either + // direction, each entered at the candidate rather than at the slab base. + rareAsc := [][2]uint32{{5, slab}, {3*slab + 7, 4 * slab}, {9*slab + 1, 10 * slab}} + rareDesc := [][2]uint32{ + {9 * slab, 9*slab + 2}, {3 * slab, 3*slab + 8}, {0, 6}, + } + assert.Equal(t, rareAsc, walk([]termPlan{{0}}, false)) + assert.Equal(t, rareDesc, walk([]termPlan{{0}}, true)) + + // ANDing it with a chunk-sized term changes nothing: the plan's bound is + // the strongest of its terms'. + assert.Equal(t, rareAsc, walk([]termPlan{{0, 2}}, false), + "a chunk-sized term must not weaken the rare term's bound") + assert.Equal(t, rareDesc, walk([]termPlan{{0, 2}}, true)) + + // The chunk-sized term alone opens every slab. + full := make([][2]uint32, 0, slabs) + for i := range uint32(slabs) { + full = append(full, [2]uint32{i * slab, (i + 1) * slab}) + } + assert.Len(t, full, slabs) + assert.Equal(t, full, walk([]termPlan{{2}}, false), + "a term holding every id must not skip a slab") + + // OR-ed plans open the union of their slabs: 0, 1, 3, 6, 9. + assert.Equal(t, [][2]uint32{ + {5, slab}, + {slab + 1, 2 * slab}, + {3*slab + 7, 4 * slab}, + {6*slab + 3, 7 * slab}, + {9*slab + 1, 10 * slab}, + }, walk([]termPlan{{0}, {1}}, false)) + + // A present but empty term ends the walk before the first slab. + assert.Empty(t, walk([]termPlan{{3}}, false)) + assert.Empty(t, walk([]termPlan{{3}}, true)) +} + +// TestMatches_RareTermsSpanSlabs is the oracle gate on the skip: the shapes +// a rare term dominates run at slab widths that put hundreds of empty slabs +// between consecutive ids. The windows start and end between rare ids, so +// the first and last slab are clipped by the window rather than a candidate. +func TestMatches_RareTermsSpanSlabs(t *testing.T) { + f := newShapedFixture(t) + r := diffReader{f.corpus} + defer func(s uint) { slabShift = s }(slabShift) + + require.LessOrEqual(t, len(f.rareContract), 5, + "fixture: the rare term must stay rare for the skip to matter") + + shapes := []namedShape{ + {"rare alone", f.filterSparseOnly()}, + {"rare and chunk-sized", f.filterMixedGroups()}, + {"rare or chunk-sized", f.filterUnion()}, + } + windows := []IDRange{ + {0, shapedCorpusSize}, + {20_000, shapedCorpusSize}, + {0, 60_000}, + {20_000, 60_000}, + } + + for _, sh := range shapes { + all := matchingEvents(t, f.corpus, sh.filters) + // 64-, 1024- and 8192-wide slabs. + for _, shift := range []uint{6, 10, 13} { + slabShift = shift + for _, w := range windows { + for _, desc := range []bool{false, true} { + for _, limit := range []int{0, 1, 3} { + requireStream(t, r, all, queryCase{ + name: sh.name, + filters: sh.filters, + window: w, + desc: desc, + limit: limit, + }) + } + } + } + } + } +}