diff --git a/cmd/stellar-rpc/internal/rpcv2/observability/observability.go b/cmd/stellar-rpc/internal/rpcv2/observability/observability.go index 697972728..88ef02172 100644 --- a/cmd/stellar-rpc/internal/rpcv2/observability/observability.go +++ b/cmd/stellar-rpc/internal/rpcv2/observability/observability.go @@ -5,6 +5,7 @@ import ( "github.com/prometheus/client_golang/prometheus" + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/packfile" "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/query" "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/rocksdb" "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/stores/ledger" @@ -153,10 +154,10 @@ func NewPrometheusMetrics(registry *prometheus.Registry, namespace string) *Prom }, []string{"phase"}), } - // Three serving-invariant counters are tallied where their condition is - // detected — deep in the store and routing packages, below any metrics - // plumbing — as package-level atomics. CounterFuncs export them; each - // should flatline at zero, so one alert rule per family is "rate > 0". + // Four serving-invariant counters are tallied where their condition is + // detected — deep in the store, routing and packfile packages, below any + // metrics plumbing — as package-level atomics. CounterFuncs export them; + // each should flatline at zero, so one alert rule per family is "rate > 0". counterFunc := func(name, help string, read func() uint64) prometheus.CounterFunc { return prometheus.NewCounterFunc(prometheus.CounterOpts{ Namespace: namespace, Subsystem: subsystem, Name: name, Help: help, @@ -180,6 +181,10 @@ func NewPrometheusMetrics(registry *prometheus.Registry, namespace string) *Prom "cold ledger packs whose file was gone on first read "+ "(routing only opens packs the catalog snapshot holds; any count is an alarm)", ledger.MissingPackOpens), + counterFunc("pooled_buffer_cap_skips_total", + "pooled packfile buffers dropped for exceeding a pool's capacity cap "+ + "(the pool drains and every open allocates afresh; any count means a cap wants raising)", + packfile.PoolCapSkips), prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Namespace: namespace, Subsystem: subsystem, Name: "open_snapshots", Help: "RocksDB snapshots currently held, across all stores " + diff --git a/cmd/stellar-rpc/internal/rpcv2/packfile/index.go b/cmd/stellar-rpc/internal/rpcv2/packfile/index.go index 93b641c96..aa981faf2 100644 --- a/cmd/stellar-rpc/internal/rpcv2/packfile/index.go +++ b/cmd/stellar-rpc/internal/rpcv2/packfile/index.go @@ -79,7 +79,8 @@ func decodeIndex(buf []byte, recordCount int, indexSize int, indexBase int64) ([ // (width, min) is at its tail, so intpack.DecodeGroup naturally reads from the end of // its window. Iterating backward lets us shrink the window after each group. groupCount := (recordCount + groupSize - 1) / groupSize - deltas := make([]uint32, recordCount) + deltas := getScratch(recordCount) + defer putScratch(deltas) pos := payloadLen for g := groupCount - 1; g >= 0; g-- { @@ -102,8 +103,9 @@ func decodeIndex(buf []byte, recordCount int, indexSize int, indexBase int64) ([ return nil, fmt.Errorf("%w: index has %d unconsumed bytes after decoding all groups", ErrCorrupt, pos) } - // Forward prefix-sum to build absolute offsets from deltas. - offsets := make([]int64, recordCount+1) + // Forward prefix-sum from deltas to absolute offsets. Every entry of the + // pooled table is overwritten below, the sentinel included. + offsets := getOffsets(recordCount + 1) offset := int64(0) for i, d := range deltas { @@ -113,6 +115,7 @@ func decodeIndex(buf []byte, recordCount int, indexSize int, indexBase int64) ([ // Structural sanity check: running sum must arrive at indexBase. if offset != indexBase { + putOffsets(offsets) return nil, fmt.Errorf("%w: final offset %d != indexBase %d", ErrCorrupt, offset, indexBase) } diff --git a/cmd/stellar-rpc/internal/rpcv2/packfile/pools.go b/cmd/stellar-rpc/internal/rpcv2/packfile/pools.go new file mode 100644 index 000000000..bc1ca6b4d --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/packfile/pools.go @@ -0,0 +1,124 @@ +package packfile + +// Pools for the open path. Cold packfiles are opened per request, and every +// open decodes an offset table, with FOR-decode scratch, out of read buffers; +// these pools recycle that memory. Contents are fully rewritten before use, +// and a pooled buffer is either dead before its function returns or +// reader-private until Close returns it. +// +// Puts are capped so one pathological file cannot pin a huge array. A cap +// must sit above what a legitimate packfile needs: crossing it skips the Put, +// the pool drains, and every open allocates afresh. capSkips counts those +// skips so the drain is visible. + +import ( + "sync" + "sync/atomic" +) + +// maxPooledOffsets caps the decoded offset table (recordCount+1 int64s) a Put +// may retain. Chunk geometry, which this package cannot import, sets what it +// must exceed: the ledger pack holds one ledger per record, 10,001 entries; +// events.pack and index.pack hold 128 items per record, so 1<<20 entries +// covers about 134M events or index terms per chunk, against roughly 600K +// terms in a production chunk today. Raise it if the geometry grows. +const maxPooledOffsets = 1 << 20 // entries (8 MiB backing array) + +const ( + // maxPooledScratch tracks maxPooledOffsets: the FOR-decode scratch is + // recordCount uint32s to the offset table's recordCount+1 int64s. + maxPooledScratch = 1 << 20 // entries (4 MiB backing array) + maxPooledOpenBuf = 4 << 20 // bytes +) + +//nolint:gochecknoglobals // process-wide pools, like recordWorkspacePool +var ( + offsetsPool sync.Pool // *[]int64 + scratchPool sync.Pool // *[]uint32 + openBufPool sync.Pool // *[]byte +) + +// capSkips counts Puts dropped for exceeding a cap, across all three pools. +// A count that climbs with the open rate means the pools have stopped +// recycling. Exported through PoolCapSkips. +// +//nolint:gochecknoglobals // one tally across process-wide pools; read-only outside this file +var capSkips atomic.Uint64 + +// PoolCapSkips returns the process-wide count of pooled buffers dropped for +// exceeding a capacity cap. Zero is the only healthy value; a climbing count +// means a cap wants raising. +func PoolCapSkips() uint64 { return capSkips.Load() } + +// On a size miss the pooled buffer is returned before a larger one is +// allocated: Get already removed it, and returning it is what keeps a run of +// growing opens from draining the pool. It is under its cap, so it never +// counts as a skip. + +func getOffsets(n int) []int64 { + if p, _ := offsetsPool.Get().(*[]int64); p != nil { + if cap(*p) >= n { + return (*p)[:n] + } + putOffsets(*p) + } + return make([]int64, n) +} + +// putOffsets recycles a decoded offset table. The caller must guarantee no +// live reference remains; see Reader.Close for the in-flight handshake. +func putOffsets(s []int64) { + if cap(s) == 0 { + return + } + if cap(s) > maxPooledOffsets { + capSkips.Add(1) + return + } + s = s[:0] + offsetsPool.Put(&s) +} + +func getScratch(n int) []uint32 { + if p, _ := scratchPool.Get().(*[]uint32); p != nil { + if cap(*p) >= n { + return (*p)[:n] + } + putScratch(*p) + } + return make([]uint32, n) +} + +func putScratch(s []uint32) { + if cap(s) == 0 { + return + } + if cap(s) > maxPooledScratch { + capSkips.Add(1) + return + } + s = s[:0] + scratchPool.Put(&s) +} + +func getOpenBuf(n int) []byte { + if p, _ := openBufPool.Get().(*[]byte); p != nil { + if cap(*p) >= n { + return (*p)[:n] + } + putOpenBuf(*p) + } + return make([]byte, n) +} + +func putOpenBuf(s []byte) { + if cap(s) == 0 { + return + } + if cap(s) > maxPooledOpenBuf { + capSkips.Add(1) + return + } + s = s[:0] + openBufPool.Put(&s) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/packfile/pools_test.go b/cmd/stellar-rpc/internal/rpcv2/packfile/pools_test.go new file mode 100644 index 000000000..a2892c8e7 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/packfile/pools_test.go @@ -0,0 +1,120 @@ +package packfile + +import ( + "bytes" + "context" + "fmt" + "os" + "sync" + "testing" + + "github.com/stretchr/testify/require" +) + +// poolItems builds a deterministic item set spanning several records. +func poolItems(n int) [][]byte { + items := make([][]byte, n) + for i := range items { + items[i] = bytes.Repeat([]byte{byte(i), byte(i >> 8)}, 64) + } + return items +} + +// Several open-read-close cycles over distinct files, so a recycled buffer +// that leaked one file's decode into another's would fail the byte checks. +func TestReaderPoolReuseKeepsReadsCorrect(t *testing.T) { + for cycle := range 4 { + // Different item counts reuse recycled arrays at different lengths. + items := poolItems(300 + 40*cycle) + path := writeTestPackfile(t, items, WriterOptions{ItemsPerRecord: 16}) + r := Open(path, ReaderOptions{}) + for i, want := range items { + require.NoError(t, r.ReadItem(i, func(got []byte) error { + if !bytes.Equal(got, want) { + return fmt.Errorf("cycle %d item %d mismatch", cycle, i) + } + return nil + })) + } + require.NoError(t, r.Close()) + } +} + +// A read beginning after Close reports os.ErrClosed rather than touching +// recycled memory. +func TestReaderReadAfterCloseFails(t *testing.T) { + items := poolItems(64) + path := writeTestPackfile(t, items, WriterOptions{ItemsPerRecord: 16}) + r := Open(path, ReaderOptions{}) + require.NoError(t, r.ReadItem(0, func([]byte) error { return nil })) + require.NoError(t, r.Close()) + + err := r.ReadItem(0, func([]byte) error { return nil }) + require.ErrorIs(t, err, os.ErrClosed) + err = r.ReadItems(context.Background(), []int{0}, func(int, []byte) error { return nil }) + require.ErrorIs(t, err, os.ErrClosed) + for _, err := range r.ReadRange(0, 1) { + require.ErrorIs(t, err, os.ErrClosed) + } +} + +// Close racing an in-flight read, a caller contract violation, must leave the +// offsets to the garbage collector rather than recycle them under the reader. +func TestReaderCloseDuringReadDoesNotRecycle(t *testing.T) { + items := poolItems(256) + path := writeTestPackfile(t, items, WriterOptions{ItemsPerRecord: 16}) + r := Open(path, ReaderOptions{}) + + parked := make(chan struct{}) + unpark := make(chan struct{}) + var closeErr error + var wg sync.WaitGroup + wg.Go(func() { + <-parked + closeErr = r.Close() + close(unpark) + }) + + first := true + got := make([]byte, 0, len(items[0])) + err := r.ReadItems(context.Background(), []int{0}, func(_ int, data []byte) error { + got = append(got[:0], data...) + if first { + first = false + close(parked) + <-unpark + } + return nil + }) + wg.Wait() + require.NoError(t, closeErr) + // Either outcome is allowed; what the handshake forbids is recycled-memory + // corruption, which the byte check would catch. + if err == nil { + require.True(t, bytes.Equal(got, items[0]), "payload corrupted by Close during read") + } else { + require.ErrorIs(t, err, os.ErrClosed) + } +} + +// The counter must count exactly the Puts a cap dropped, on every pool, and +// nothing for kept buffers or ignored empty slices. It is process-wide, so +// the assertions are on the delta. +func TestPoolCapSkipsCountsDroppedPuts(t *testing.T) { + before := PoolCapSkips() + + putOffsets(make([]int64, 0, maxPooledOffsets)) + putScratch(make([]uint32, 0, maxPooledScratch)) + putOpenBuf(make([]byte, 0, maxPooledOpenBuf)) + putOffsets(nil) + putScratch(nil) + putOpenBuf(nil) + require.Equal(t, before, PoolCapSkips(), + "a Put at the cap is pooled and an empty one is ignored; neither is a cap skip") + + putOffsets(make([]int64, 0, maxPooledOffsets+1)) + putScratch(make([]uint32, 0, maxPooledScratch+1)) + putOpenBuf(make([]byte, 0, maxPooledOpenBuf+1)) + require.Equal(t, before+3, PoolCapSkips(), + "each pool must count the Put its cap dropped") +} diff --git a/cmd/stellar-rpc/internal/rpcv2/packfile/reader.go b/cmd/stellar-rpc/internal/rpcv2/packfile/reader.go index 4503a6a76..3d63edb30 100644 --- a/cmd/stellar-rpc/internal/rpcv2/packfile/reader.go +++ b/cmd/stellar-rpc/internal/rpcv2/packfile/reader.go @@ -161,6 +161,14 @@ type Reader struct { waitOpen func() error // blocks until background open completes + // closed and inflight let Close recycle the pooled offsets safely. A read + // increments inflight before checking closed; Close sets closed before + // reading inflight. So either the read backs out, or Close sees it and + // leaves the offsets to the garbage collector. A read racing Close is a + // caller bug either way; this makes its worst case a leak, not reuse. + closed atomic.Bool + inflight atomic.Int64 + closeOnce sync.Once closeErr error } @@ -235,7 +243,8 @@ func doOpen(path string) openResult { // Speculative read: last min(speculativeReadSize, fileSize) bytes. speculativeSize := min(int64(speculativeReadSize), fileSize) speculativeOff := fileSize - speculativeSize - speculativeBuf := make([]byte, speculativeSize) + speculativeBuf := getOpenBuf(int(speculativeSize)) + defer putOpenBuf(speculativeBuf) if _, err := f.ReadAt(speculativeBuf, speculativeOff); err != nil { return openResult{err: fmt.Errorf("packfile: read trailer region: %w", err)} } @@ -262,11 +271,12 @@ func doOpen(path string) openResult { var indexBuf []byte var appData []byte + // Both arms hand decodeIndex a view into a pooled buffer; only appData, + // which the Reader keeps, is copied out. if tailSize <= speculativeSize { // Index + appData are already inside the speculative read. tailStart := len(speculativeBuf) - int(tailSize) - indexBuf = make([]byte, indexSize) - copy(indexBuf, speculativeBuf[tailStart:tailStart+indexSize]) + indexBuf = speculativeBuf[tailStart : tailStart+indexSize] if appDataSize > 0 { appData = make([]byte, appDataSize) adStart := tailStart + indexSize @@ -275,7 +285,8 @@ func doOpen(path string) openResult { } else { // Single fallback read for index + appData. readSize := indexSize + appDataSize - buf := make([]byte, readSize) + buf := getOpenBuf(readSize) + defer putOpenBuf(buf) if readSize > 0 { if _, err := f.ReadAt(buf, indexBase); err != nil { return openResult{err: fmt.Errorf("packfile: read index region: %w", err)} @@ -297,11 +308,6 @@ func doOpen(path string) openResult { ErrChecksum, trailer.AppDataCRC, computed)} } - offsets, err := decodeIndex(indexBuf, recordCount, indexSize, indexBase) - if err != nil { - return openResult{err: err} - } - // Not gated on recordCount: a trailer claiming items but no itemsPerRecord // would otherwise pass Open and index past the offsets slice on first read. if itemsPerRecord <= 0 && (recordCount > 0 || totalItems > 0) { @@ -322,6 +328,13 @@ func doOpen(path string) openResult { } } + // Decoded last: the table is pooled, and every check that could still + // fail has run, so on success the Reader owns it until Close. + offsets, err := decodeIndex(indexBuf, recordCount, indexSize, indexBase) + if err != nil { + return openResult{err: err} + } + // Empty packfiles may legitimately have itemsPerRecord==0 on disk; // default the *internal* int to 1 so modulo math is well-defined. The // Trailer view keeps the on-disk value verbatim. @@ -417,6 +430,10 @@ func (r *Reader) ReadItem(position int, fn func([]byte) error) error { if err := r.waitOpen(); err != nil { return err } + if err := r.beginRead(); err != nil { + return err + } + defer r.endRead() if position < 0 || position >= r.totalItems { return ErrPositionOutOfRange } @@ -469,6 +486,11 @@ func (r *Reader) ReadRange(start, count int) iter.Seq2[[]byte, error] { yield(nil, err) return } + if err := r.beginRead(); err != nil { + yield(nil, err) + return + } + defer r.endRead() if start < 0 || count < 0 || start > r.totalItems || count > r.totalItems-start { yield(nil, fmt.Errorf("%w: ReadRange(%d, %d) out of [0, %d)", ErrPositionOutOfRange, start, count, r.totalItems)) @@ -577,6 +599,10 @@ func (r *Reader) ReadItems(ctx context.Context, positions []int, fn func(idx int if err := r.waitOpen(); err != nil { return err } + if err := r.beginRead(); err != nil { + return err + } + defer r.endRead() for i, pos := range positions { if pos < 0 || pos >= r.totalItems { @@ -768,12 +794,37 @@ func (r *Reader) Verify(ctx context.Context) error { // Readers); its lifecycle is the caller's responsibility. func (r *Reader) Close() error { r.closeOnce.Do(func() { + r.closed.Store(true) openErr := r.waitOpen() var closeErr error if r.file != nil { closeErr = r.file.Close() } + // Recycle only when no read is in flight. closed is already set, so + // no new read can begin; a read still in flight is a caller bug, and + // leaving the array to the collector keeps that a leak. + if r.inflight.Load() == 0 && r.offsets != nil { + putOffsets(r.offsets) + r.offsets = nil + } r.closeErr = errors.Join(openErr, closeErr) }) return r.closeErr } + +// errReaderClosed is returned by reads that begin after Close. It wraps +// os.ErrClosed so callers matching the closed-file error shape keep matching. +var errReaderClosed = fmt.Errorf("packfile: read after Close: %w", os.ErrClosed) + +// beginRead registers a read with the Close handshake; endRead must run when +// the read finishes. See the closed and inflight field comment. +func (r *Reader) beginRead() error { + r.inflight.Add(1) + if r.closed.Load() { + r.inflight.Add(-1) + return errReaderClosed + } + return nil +} + +func (r *Reader) endRead() { r.inflight.Add(-1) } diff --git a/cmd/stellar-rpc/internal/rpcv2/query/resolve.go b/cmd/stellar-rpc/internal/rpcv2/query/resolve.go index a74b3ed06..d738d9489 100644 --- a/cmd/stellar-rpc/internal/rpcv2/query/resolve.go +++ b/cmd/stellar-rpc/internal/rpcv2/query/resolve.go @@ -144,6 +144,16 @@ func (a *ReadView) resolveLedgers(c chunk.ID) (LedgerReader, func() error, error } } +// defaultColdEventReadConcurrency is the worker fan-out for one cold events +// read over its packfiles. A page's payload fetch is hundreds of scattered +// records with no ordering between them, and serial reads add their latencies +// together. The value depends on the storage the daemon reads through, not on +// the query, and was measured on NVMe. Its footprint is workers times +// in-flight cold pages, in goroutines and in packfile's coalesced-read +// buffers. It is compiled in; a config knob can follow if a deployment needs +// a different value. +const defaultColdEventReadConcurrency = 8 + // Events resolves chunk c's event store as the common event.Reader the // query engine consumes, uniform across tiers. A cold reader is view-owned — // Release closes it; the hot facade is registry-owned. Returns ErrUnavailable @@ -157,10 +167,8 @@ func (a *ReadView) Events(c chunk.ID) (event.Reader, error) { } switch t { case tierCold: - // TODO(events adapter / #772): thread read concurrency - // (ColdReaderOptions.Concurrency → the packfile ReadItems concurrency) here; - // decide whether it is config-driven or caller-supplied. Default for now. - cr, err := event.OpenColdReader(c, a.catalog.Layout().EventsBucketDir(c), event.ColdReaderOptions{}) + cr, err := event.OpenColdReader(c, a.catalog.Layout().EventsBucketDir(c), + event.ColdReaderOptions{Concurrency: defaultColdEventReadConcurrency}) if err != nil { return nil, err } diff --git a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go index edb3dfde4..d4bd7ec4b 100644 --- a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go +++ b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go @@ -130,6 +130,12 @@ type Store struct { ro *grocksdb.ReadOptions wo *grocksdb.WriteOptions + // roAsync is ro plus async IO, used by BatchMultiGet. A separate object + // keeps the shared ro, and every Get and iterator on it, synchronous. + // Built at open and never mutated, so it is safe to share across + // concurrent batched reads. + roAsync *grocksdb.ReadOptions + // cache is the block cache shared across every CF in this store, // created in applyTuning when BlockCacheMB is set. bbtos are the // per-CF block-based-table options (one per CF, carrying the pinned @@ -244,11 +250,14 @@ func (s *Store) GetPinned(cf string, key []byte, fn func(value []byte) error) (b // merge adjacent SST seeks across the input set. Behavior on // unsorted input is undefined per RocksDB semantics. // -// Uses async_io read options so the kernel can issue overlapping -// I/Os under the hood (notable on EBS / high random-latency -// storage). The batched call is a single CGO crossing; callers -// needing cancellation between individual key reads should not use -// this API — split into multiple calls or use Get in a loop. +// keys are not retained past the call: grocksdb copies each into C +// memory and frees the copy before returning, so a caller may carve +// the whole list out of one buffer (see event.encodeDataKeys). +// +// Reads use roAsync, the store's async_io read options, so the kernel +// can overlap the I/Os. The batched call is a single CGO crossing; +// callers needing cancellation between individual keys should split +// into multiple calls or use Get in a loop. func (s *Store) BatchMultiGet(cf string, keys [][]byte) ([][]byte, error) { if len(keys) == 0 { return nil, nil @@ -263,26 +272,33 @@ func (s *Store) BatchMultiGet(cf string, keys [][]byte) ([][]byte, error) { return nil, err } - // Fresh ReadOptions: mutating s.ro would surface async_io to - // every concurrent reader on this Store. - ro := grocksdb.NewDefaultReadOptions() - ro.SetAsyncIO(true) - defer ro.Destroy() - - pinned, err := s.db.BatchedMultiGetCF(ro, cfh, true /* sortedInput */, keys...) + pinned, err := s.db.BatchedMultiGetCF(s.roAsync, cfh, true /* sortedInput */, keys...) if err != nil { return nil, fmt.Errorf("rocksdb: batched multi get on %q: %w", cf, err) } defer pinned.Destroy() + // Copy out of the pinned cache pages, which Destroy invalidates, through + // one arena rather than a clone per value. Each pinned value is read once, + // since Data heap-allocates an out-param per call: the sizing pass keeps + // the C-backed slice, valid until Destroy, and the copy pass reads it. The + // returned slices share one backing array; retaining one retains the batch. results := make([][]byte, len(keys)) + total := 0 for i, p := range pinned { - if !p.Exists() { + if p.Exists() { + results[i] = p.Data() + total += len(results[i]) + } + } + arena := make([]byte, 0, total) + for i := range results { + if !pinned[i].Exists() { continue } - // p.Data() points into the pinned cache page; copy before - // Destroy invalidates it. - results[i] = bytes.Clone(p.Data()) + n := len(arena) + arena = append(arena, results[i]...) + results[i] = arena[n:len(arena):len(arena)] } return results, nil } @@ -542,6 +558,7 @@ func (s *Store) teardownLocked() { cfh.Destroy() } s.ro.Destroy() + s.roAsync.Destroy() s.wo.Destroy() s.db.Close() s.opts.Destroy() @@ -758,8 +775,12 @@ func (s *Store) constructAndOpen() error { s.cfOpts = cfOpts s.cfHandles = cfMap s.ro = grocksdb.NewDefaultReadOptions() + s.roAsync = grocksdb.NewDefaultReadOptions() s.wo = grocksdb.NewDefaultWriteOptions() + // Async IO is the only difference from ro; see the field. + s.roAsync.SetAsyncIO(true) + // WAL on + per-write Sync on — non-negotiable across every // rpcv2 store, so pinned here on the shared wo rather // than exposed via Tuning. The ingestion contract diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go new file mode 100644 index 000000000..7861e273e --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go @@ -0,0 +1,31 @@ +package event + +// byteArena hands out stable copies of transient byte slices, carved from +// larger chunks so a fetch copying hundreds of payloads costs a few +// allocations rather than one each. A chunk is only appended within its +// capacity, so returned copies never move. The zero value is ready. +// +// Not safe for concurrent use. +type byteArena struct { + buf []byte +} + +// The first chunk is small so a one-event page does not pay 64 KiB; each +// later chunk doubles up to arenaChunkSize. +const ( + arenaFirstChunkSize = 4 << 10 + arenaChunkSize = 64 << 10 +) + +func (a *byteArena) copy(b []byte) []byte { + if len(b) > cap(a.buf)-len(a.buf) { + next := arenaFirstChunkSize + if c := 2 * cap(a.buf); c > next { + next = min(c, arenaChunkSize) + } + a.buf = make([]byte, 0, max(next, len(b))) + } + n := len(a.buf) + a.buf = append(a.buf, b...) + return a.buf[n : n+len(b) : n+len(b)] +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/arena_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/arena_test.go new file mode 100644 index 000000000..3a537d860 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/arena_test.go @@ -0,0 +1,37 @@ +package event + +import ( + "bytes" + "testing" + + "github.com/stretchr/testify/require" +) + +// A returned copy never moves or changes, however much is copied after it. +func TestByteArenaCopiesAreStable(t *testing.T) { + var a byteArena + src := make([]byte, 300) + got := make([][]byte, 0, 3000) + want := make([][]byte, 0, 3000) + for i := range 3000 { // ~900KB total: crosses many 64KB chunks + for j := range src { + src[j] = byte(i + j) + } + c := a.copy(src) + got = append(got, c) + want = append(want, bytes.Clone(src)) + } + huge := bytes.Repeat([]byte{0xAB}, 3*arenaChunkSize) + hugeCopy := a.copy(huge) + huge[0] = 0xCD // mutate the source; the copy must not see it + require.Equal(t, byte(0xAB), hugeCopy[0]) + require.Len(t, hugeCopy, 3*arenaChunkSize) + for i := range got { + require.True(t, bytes.Equal(got[i], want[i]), "copy %d changed", i) + } + // Appending to a returned copy must not scribble into the arena. + c := a.copy([]byte{1, 2, 3}) + next := a.copy([]byte{9, 9, 9}) + _ = append(c, 7) // the append must copy, not extend in place + require.Equal(t, []byte{9, 9, 9}, next) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go index da09a6b99..071494e65 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go @@ -434,24 +434,17 @@ func buildMPHF( // /index.hash produced by an earlier buildMPHF) for // query-time lookups. // -// The file is read into memory up-front via os.ReadFile + -// streamhash.OpenBytes rather than mmapped. Rationale: a typical -// MPHF for a single Chunk is small (~hundreds of KB at production -// term counts), and on storage with expensive random IOPS (e.g. -// EBS, ~1 ms each) mmap page-faults on cold Lookups cost more than -// a single sequential read amortized across the index's lifetime. +// The file is mmapped rather than read whole. Pages fault in from the kernel +// page cache, which every reader of the same chunk shares regardless of its +// own lifetime, so a per-request open costs a map and unmap plus the pages +// its lookups touch, not a copy of a file whose size scales with the chunk's +// term count. // -// Close on the returned handle is a no-op for the OpenBytes path -// (streamhash holds no fd / mmap), but callers should still call it -// for symmetry with other open variants. +// Close unmaps; callers must call it. func openMPHF(path string) (*mphf, error) { - data, err := os.ReadFile(path) + idx, err := streamhash.Open(path) if err != nil { - return nil, fmt.Errorf("events: read %s: %w", path, err) - } - idx, err := streamhash.OpenBytes(data) - if err != nil { - return nil, fmt.Errorf("events: parse %s: %w", path, err) + return nil, fmt.Errorf("events: open %s: %w", path, err) } secret, merr := decodeEventsMeta(idx.UserMetadata()) if merr != nil { @@ -488,7 +481,7 @@ func (m *mphf) Lookup(key TermKey) (uint32, error) { return uint32(slot), nil } -// Close releases the index; a no-op for the in-memory OpenBytes path. +// Close unmaps the index file; callers must call it (see openMPHF). func (m *mphf) Close() error { return m.idx.Close() } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go index 9c0a3efa9..3edbf1933 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go @@ -450,7 +450,8 @@ func (c *ColdReader) LookupKeys(ctx context.Context, keys []TermKey) ([]*roaring // records into single ReadAt calls and optionally fans out across // the worker count set via ColdReaderOptions.Concurrency. // result[idx] writes from concurrent workers do not race — each -// idx is unique. +// idx is unique — and the payload arena they share is locked (see +// the callback). func (c *ColdReader) FetchEvents(ctx context.Context, eventIDs []uint32) ([]Payload, error) { if c.closed.Load() { return nil, stores.ErrStoreClosed @@ -474,12 +475,23 @@ func (c *ColdReader) FetchEvents(ctx context.Context, eventIDs []uint32) ([]Payl positions[i] = int(id) } results := make([]Payload, len(eventIDs)) + // One arena per call, shared by every worker, since a call's payloads + // live and die together. ReadItems calls back from up to Concurrency + // goroutines and the arena is a single appended buffer, hence the lock. + // It covers the copy alone, so the read, decode and Unmarshal overlap. + var ( + arenaMu sync.Mutex + arena byteArena + ) if err := c.events.ReadItems(ctx, positions, func(idx int, data []byte) error { // packfile.ReadItems passes a borrowed data slice valid only for // the duration of fn (see Reader.ReadItems docstring). FetchEvents - // returns the Payloads in a slice that outlives fn, so clone before + // returns the Payloads in a slice that outlives fn, so copy before // Unmarshal aliases the bytes into ContractEventBytes. - return results[idx].Unmarshal(bytes.Clone(data)) + arenaMu.Lock() + owned := arena.copy(data) + arenaMu.Unlock() + return results[idx].Unmarshal(owned) }); err != nil { // packfile.ReadItems also validates sorted positions as defense in // depth; translate its sentinel to ours so callers can errors.Is diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_fanout_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_fanout_test.go new file mode 100644 index 000000000..bbe18b157 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_fanout_test.go @@ -0,0 +1,38 @@ +package event + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/chunk" +) + +// Pins the payload arena against the reader's own fan-out: with Concurrency +// above one, ReadItems calls the callback from one goroutine per batch, so an +// unsynchronized append into the shared arena corrupts payloads or segfaults. +// Every payload is compared whole, so a torn copy fails even when -race does +// not schedule the overlap. +func TestColdReader_FetchEventsFansOutSafely(t *testing.T) { + const ( + chunkID = chunk.ID(0) + events = 4096 // 32 records at eventsPackItemsPerRecord + ) + dir, payloads := buildColdFixture(t, chunkID, events, 2) + + cr, err := OpenColdReader(chunkID, dir, ColdReaderOptions{Concurrency: 8}) + require.NoError(t, err) + t.Cleanup(func() { _ = cr.Close() }) + + ids := make([]uint32, 0, events) + for i := range uint32(events) { + ids = append(ids, i) + } + got, err := cr.FetchEvents(context.Background(), ids) + require.NoError(t, err) + require.Len(t, got, len(ids)) + for i := range ids { + require.Equal(t, payloads[i], got[i], "payload %d", i) + } +} 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/hot_store.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go index ec2e3730c..771ee0a7f 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go @@ -29,16 +29,21 @@ const ( // // - DataCF holds XDR-encoded event payloads: compressible (zstd // typically 2-3× on XDR) and read in batches via -// BatchedMultiGetCF. Larger blocks give zstd more context per -// compression unit and align with batch-fetch shapes. +// BatchedMultiGetCF. A point read decompresses one block to serve +// one ~250B event, so a 32 KiB block paid for ~128 events per cache +// miss. // - IndexCF stores 20-byte (term_hash || event_id) keys with // empty values — nothing in the values to compress, and small // blocks reduce wasted I/O per random Lookup miss (each Lookup // reads one block to find one key). // - OffsetsCF stores 8-byte (ledger_seq -> event_count) rows in // the tens-of-thousands per chunk — same shape as IndexCF. +// +// A block size applies to SSTs as they are written. Chunks already on disk +// keep theirs until compaction or rotation rewrites them; a restart changes +// nothing. const ( - dataCFBlockSize = 32 * 1024 + dataCFBlockSize = 8 * 1024 indexCFBlockSize = 4 * 1024 offsetsCFBlockSize = 4 * 1024 ) @@ -177,6 +182,12 @@ func (h *HotStore) Offsets() (*LedgerOffsets, error) { // to batch — but exposing this method satisfies the Reader // interface so callers can program against batched lookups // uniformly. +// +// Each bitmap is a point-in-time image of its term: a sparse term is +// copied out of the mirror's published id list, a dense one is +// denseState.snapshot, the immutable clone shared with every other +// reader. Neither grows under its holder, so a walk never sees an id +// written after its lookup. func (h *HotStore) LookupKeys(ctx context.Context, keys []TermKey) ([]*roaring.Bitmap, error) { if h.chunkStore.IsClosed() { return nil, stores.ErrStoreClosed @@ -232,10 +243,7 @@ func (h *HotStore) FetchEvents(ctx context.Context, eventIDs []uint32) ([]Payloa return nil, err } - keys := make([][]byte, len(eventIDs)) - for i, id := range eventIDs { - keys[i] = encodeDataKey(id) - } + keys := encodeDataKeys(eventIDs) values, err := h.chunkStore.BatchMultiGet(DataCF, keys) if err != nil { return nil, fmt.Errorf("events: batch fetch from chunk %s: %w", h.chunkID, err) @@ -675,6 +683,27 @@ func encodeDataKey(eventID uint32) []byte { return key[:] } +// encodeDataKeys encodes every id into one backing buffer and returns +// per-id sub-slices of it, in input order: two allocations for the +// batch rather than one per id, which is what encodeDataKey costs +// because its array escapes through the returned slice. +// +// The slices alias one array and must not be retained past the call +// they are passed to. BatchMultiGet qualifies: grocksdb copies each key +// into C memory and frees the copy before returning. +func encodeDataKeys(eventIDs []uint32) [][]byte { + buf := make([]byte, dataKeyLen*len(eventIDs)) + keys := make([][]byte, len(eventIDs)) + for i, id := range eventIDs { + lo := i * dataKeyLen + hi := lo + dataKeyLen + key := buf[lo:hi:hi] + binary.BigEndian.PutUint32(key, id) + keys[i] = key + } + return keys +} + func encodeIndexKey(term TermKey, eventID uint32) []byte { var key [indexKeyLen]byte copy(key[:16], term[:]) diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store_test.go index 2a2a2180d..aa3343bd6 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store_test.go @@ -346,6 +346,77 @@ func TestHotStore_FetchEventsRejectsUnsortedInput(t *testing.T) { require.ErrorIs(t, err, ErrUnsortedEventIDs, "duplicate input must error") } +// fetchEventsPerIDAllocBudget is the number of heap allocations FetchEvents +// may spend per event ID. Both are grocksdb's: one *PinnableSlice per key +// and the size out-param of the one PinnableSlice.Data call per key. Our +// own work per batch is a fixed handful. If a grocksdb bump moves this +// number, raise it after checking where the new allocation comes from; do +// not widen it to absorb a regression on our side. +const fetchEventsPerIDAllocBudget = 2 + +// TestHotStore_FetchEventsAllocationBudget pins that FetchEvents spends no +// per-ID allocation of its own; the regression it catches is a +// heap-allocated key per ID. +func TestHotStore_FetchEventsAllocationBudget(t *testing.T) { + const chunkID = chunk.ID(0) + const n = 512 + h := openHotStoreForTest(t, chunkID) + + payloads := make([]Payload, n) + for i := range n { + p, _ := makePayload(fmt.Sprintf("evt-%03d", i)) + payloads[i] = p + } + require.NoError(t, ingestLedgerEvents(h.store, 2, payloads)) + + ids := make([]uint32, n) + for i := range n { + ids[i] = uint32(i) + } + ctx := context.Background() + // Warm the block cache first: a cold read allocates in RocksDB's own + // Go-side plumbing on the way to the SSTs and would skew run one. + _, err := h.store.FetchEvents(ctx, ids) + require.NoError(t, err) + + allocs := testing.AllocsPerRun(20, func() { + if _, err := h.store.FetchEvents(ctx, ids); err != nil { + t.Error(err) + } + }) + + // Per-ID budget plus generous room for the fixed per-batch handful. + const fixedAllocSlack = 64 + budget := float64(n*fetchEventsPerIDAllocBudget + fixedAllocSlack) + assert.LessOrEqual(t, allocs, budget, + "FetchEvents of %d IDs allocated %.0f times (budget %.0f): a per-ID "+ + "allocation crept back into the fetch path", n, allocs, budget) +} + +// TestEncodeDataKeys pins that the keys equal encodeDataKey's, are distinct +// windows onto one buffer, and cost a fixed number of allocations. +func TestEncodeDataKeys(t *testing.T) { + ids := []uint32{0, 1, 7, 1 << 20, ^uint32(0)} + keys := encodeDataKeys(ids) + require.Len(t, keys, len(ids)) + for i, id := range ids { + assert.Equal(t, encodeDataKey(id), keys[i], "key %d", i) + assert.Len(t, keys[i], dataKeyLen, "key %d", i) + // Full slice expression: appending to one key must not scribble + // over the next one. + assert.Equal(t, dataKeyLen, cap(keys[i]), "key %d capacity", i) + } + + // Same allocation count two orders of magnitude apart. + small := testing.AllocsPerRun(100, func() { _ = encodeDataKeys(make([]uint32, 8)) }) + large := testing.AllocsPerRun(100, func() { _ = encodeDataKeys(make([]uint32, 4096)) }) + // One for the key buffer, one for the slice headers, one for the + // make([]uint32) the closure itself does. + const encodeDataKeysAllocs = 3 + assert.LessOrEqual(t, small, float64(encodeDataKeysAllocs), "8 IDs") + assert.LessOrEqual(t, large, float64(encodeDataKeysAllocs), "4096 IDs") +} + func TestHotStore_AllStreamsInEventIDOrder(t *testing.T) { const chunkID = chunk.ID(0) h := openHotStoreForTest(t, chunkID) 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 671bc2718..496374604 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/payload.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/payload.go index 475ac5e09..a41745648 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/payload.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/payload.go @@ -167,10 +167,11 @@ func (p *Payload) MarshalInto(dst []byte) ([]byte, error) { // slice ALIASES into data and is valid only as long as data is. The // the event store read paths apply this two ways: // -// - FetchEvents passes data that outlives the returned slice — hot from -// rocksdb.BatchMultiGet (freshly allocated, caller-owned), cold by -// cloning the borrowed packfile.ReadItems buffer — so its Payloads are -// safe to retain. +// - FetchEvents passes data that outlives the returned slice, hot from +// rocksdb.BatchMultiGet, cold copied out of the borrowed +// packfile.ReadItems buffer, so its Payloads are safe to retain. Hot +// bytes are shared by the whole batch: retaining one Payload pins +// every payload fetched with it. // - FetchRange / All pass the iterator's borrowed buffer directly // (rocksdb.IterateRange / packfile.ReadRange, valid only for the // current step), so each yielded Payload is borrowed; a consumer that 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, + }) + } + } + } + } + } +} diff --git a/go.mod b/go.mod index b1cd757dc..66a2d5bbe 100644 --- a/go.mod +++ b/go.mod @@ -4,9 +4,10 @@ go 1.26 require ( github.com/Masterminds/squirrel v1.5.4 - // Minimum v2.18.2: first upstream release with the FastOr/runContainer16 - // fix (RoaringBitmap/roaring#527) that the fork previously carried. - github.com/RoaringBitmap/roaring/v2 v2.18.2 + // v2.18.2 is the minimum (the FastOr/runContainer16 fix, #527). + // v2.26.0 adds vectorized container kernels behind x/sys/cpu and + // GODEBUG gates; the aggregation layer is unchanged since v2.18. + github.com/RoaringBitmap/roaring/v2 v2.26.0 github.com/aws/aws-sdk-go-v2 v1.45.1 github.com/aws/aws-sdk-go-v2/config v1.31.16 github.com/aws/aws-sdk-go-v2/service/s3 v1.110.0 @@ -61,7 +62,7 @@ require ( github.com/aws/aws-sdk-go-v2/service/sso v1.30.0 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.35.4 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.39.0 // indirect - github.com/bits-and-blooms/bitset v1.24.2 // indirect + github.com/bits-and-blooms/bitset v1.24.4 // indirect github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2 // indirect github.com/edsrzf/mmap-go v1.2.0 // indirect github.com/envoyproxy/go-control-plane/envoy v1.37.0 // indirect diff --git a/go.sum b/go.sum index 2828a53be..baafbb3a0 100644 --- a/go.sum +++ b/go.sum @@ -42,8 +42,8 @@ github.com/Masterminds/squirrel v1.5.4 h1:uUcX/aBc8O7Fg9kaISIUsHXdKuqehiXAMQTYX8 github.com/Masterminds/squirrel v1.5.4/go.mod h1:NNaOrjSoIDfDA40n7sr2tPNZRfjzjA400rg+riTZj10= github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= -github.com/RoaringBitmap/roaring/v2 v2.18.2 h1:oPq3Cgx//iDuJQVp6xSInAKW34J9CEwE5GmLI2z+Eic= -github.com/RoaringBitmap/roaring/v2 v2.18.2/go.mod h1:eq4wdNXxtJIS/oikeCzdX1rBzek7ANzbth041hrU8Q4= +github.com/RoaringBitmap/roaring/v2 v2.26.0 h1:K30ZxF4vZcIKvJsbmgfiep2K64f+dILJqkYGoj4xnwU= +github.com/RoaringBitmap/roaring/v2 v2.26.0/go.mod h1:BZufmFbox589n3j5eOmyTaLSGXbRLc2LmQvjKjzSEGU= github.com/ajg/form v0.0.0-20160822230020-523a5da1a92f h1:zvClvFQwU++UpIUBGC8YmDlfhUrweEy1R1Fj1gu5iIM= github.com/ajg/form v0.0.0-20160822230020-523a5da1a92f/go.mod h1:uL1WgH+h2mgNtvBq0339dVnzXdBETtL2LeUXaIv25UY= github.com/andybalholm/brotli v1.0.4 h1:V7DdXeJtZscaqfNuAdSRuRFzuiKlHSC/Zh3zl9qY3JY= @@ -94,8 +94,8 @@ github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ= github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= -github.com/bits-and-blooms/bitset v1.24.2 h1:M7/NzVbsytmtfHbumG+K2bremQPMJuqv1JD3vOaFxp0= -github.com/bits-and-blooms/bitset v1.24.2/go.mod h1:7hO7Gc7Pp1vODcmWvKMRA9BNmbv6a/7QIWpPxHddWR8= +github.com/bits-and-blooms/bitset v1.24.4 h1:95H15Og1clikBrKr/DuzMXkQzECs1M6hhoGXLwLQOZE= +github.com/bits-and-blooms/bitset v1.24.4/go.mod h1:7hO7Gc7Pp1vODcmWvKMRA9BNmbv6a/7QIWpPxHddWR8= github.com/caarlos0/env/v11 v11.4.1 h1:fYwH0sWEsBSMPG7t4e/PEfTFzrWrpjyygXyUnWiSwEw= github.com/caarlos0/env/v11 v11.4.1/go.mod h1:qupehSf/Y0TUTsxKywqRt/vJjN5nz6vauiYEUUr8P4U= github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8=