Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
41 commits
Select commit Hold shift + click to select a range
5edf0d1
events: 8 KiB DataCF blocks — the block is getEvents' decompression unit
tamirms Aug 29, 2026
a1de1b6
events: batch arenas for fetched payload copies; page-sized first mat…
tamirms Aug 29, 2026
1e6aea8
events: ascending match pulls candidates through a cursor tree — no m…
tamirms Aug 29, 2026
4d35a28
events: bound the intersect gallop — rarest group leads, overrun spil…
tamirms Aug 29, 2026
ee9d2cd
deps: roaring v2.18.2 -> v2.26.0 — bulk aggregation without the pre-O…
tamirms Aug 29, 2026
e1b86ca
events: lock the cold fetch arena against packfile fan-out
tamirms Aug 29, 2026
cad54c7
query: cold events reads fan out over the packfiles — Concurrency=8
tamirms Aug 29, 2026
73259dd
events: mmap the MPHF — the page cache is the cross-request index cache
tamirms Aug 29, 2026
3d4c91e
packfile: pool the open path — recycle offsets, decode scratch, and o…
tamirms Aug 29, 2026
8f470ef
review: the cold fan-out seam is a real Registry field; two test gaps
tamirms Aug 30, 2026
14e335e
review: Matches doc — the hint cap, not the default, handles oversize…
tamirms Aug 30, 2026
5d7d9d2
lint: CI findings — WaitGroup.Go, funcorder on the read handshake, ar…
tamirms Aug 30, 2026
0c8bfc9
lint: prealloc the arena test's expectation slices
tamirms Aug 30, 2026
484379f
review: the fan-out knob is deployment-reachable; the arena's first c…
tamirms Aug 30, 2026
78931d0
events: trim this branch's comments to contracts
tamirms Aug 31, 2026
e742b01
query: the fan-out doc names its override
tamirms Aug 31, 2026
958a0c5
events: pin the postings accessors' freshness against the lazy dense …
tamirms Sep 9, 2026
8424908
query: drop the cold fan-out override — the constant is the contract
tamirms Sep 9, 2026
72bb3c0
packfile: a pool size miss returns the buffer; the offsets cap states…
tamirms Sep 9, 2026
6409e4f
events, deps: comments that describe what the code actually does
tamirms Sep 9, 2026
b75de99
events: pin roaring's read-only argument contract for shared index bi…
tamirms Sep 9, 2026
32eee60
events: a slab-stepped match engine beside the cursor tree, proved by…
tamirms Sep 9, 2026
d503ce4
events: the slab walk is the match engine — drop the cursor tree and …
tamirms Sep 9, 2026
f011769
events: one oracle for the match path — drop the benchmarks and the t…
tamirms Sep 9, 2026
96aa40d
events: LookupKeys is the only index path — drop the postings seam
tamirms Sep 9, 2026
709514a
events: the slab walk skips slabs it can prove hold no candidate
tamirms Sep 9, 2026
ae923c4
events: union the per-filter slab results in place
tamirms Sep 9, 2026
3c5a904
packfile: count the pooled buffers a capacity cap drops
tamirms Sep 9, 2026
3586a7c
events: the slab walk holds the bounds it proves
tamirms Sep 9, 2026
27fdd80
observability: export the packfile pool cap-skip counter
tamirms Sep 9, 2026
9f1db00
events: one buffer for the hot fetch's key list
tamirms Sep 9, 2026
a0691ad
rocksdb: build the batched read's options once, at open
tamirms Sep 9, 2026
6789545
rocksdb: read each pinned value once in the batched get
tamirms Sep 9, 2026
1ba88a1
events: the fetch hint is a validated page size, honored in full
tamirms Sep 9, 2026
1494493
packfile: return the pooled offsets on every failed open
tamirms Sep 9, 2026
a84845b
events: compare whole payloads in the fan-out test; fix two stale tes…
tamirms Sep 9, 2026
32f544a
rpcv2: trim the comments this branch added
tamirms Sep 9, 2026
17f89ba
event: trim the test suite and its comments
tamirms Sep 9, 2026
9f03272
event: preallocate the value-search results in the roaring race gate
tamirms Sep 9, 2026
39a0c94
event: evaluate each slab with one FastAnd per plan
tamirms Sep 11, 2026
c9b4053
event: say what the pinned FastAnd allocates in eval's doc
tamirms Sep 12, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 9 additions & 4 deletions cmd/stellar-rpc/internal/rpcv2/observability/observability.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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,
Expand All @@ -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 " +
Expand Down
9 changes: 6 additions & 3 deletions cmd/stellar-rpc/internal/rpcv2/packfile/index.go
Original file line number Diff line number Diff line change
Expand Up @@ -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-- {
Expand All @@ -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 {
Expand All @@ -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)
}

Expand Down
124 changes: 124 additions & 0 deletions cmd/stellar-rpc/internal/rpcv2/packfile/pools.go
Original file line number Diff line number Diff line change
@@ -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)
}
120 changes: 120 additions & 0 deletions cmd/stellar-rpc/internal/rpcv2/packfile/pools_test.go
Original file line number Diff line number Diff line change
@@ -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
Comment on lines +84 to +85
}
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")
}
Loading
Loading