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, + }) + } + } + } + } + } +}