From c3db605cab98e8cae4dd7c6a9481747000aeeb11 Mon Sep 17 00:00:00 2001 From: tamirms Date: Fri, 25 Sep 2026 16:21:34 -0400 Subject: [PATCH] events: batch fetched payload copies into arenas Hot fetch: the batched read's options are built once, at open; the key list is cut from one buffer; each pinned value is read once; and the values are copied into one arena per batch instead of one allocation each. A fetch now makes two heap allocations per event id, both inside grocksdb, which TestHotStore_FetchEventsAllocationBudget pins. Cold fetch: payload copies go into arenas instead of one allocation each. A cold reader with a concurrency above one makes packfile's ReadItems call back from several goroutines, so each callback takes an arena from a per-call pool for its copy and no two share one. A single shared arena behind a lock made a fetch with 8 reads in flight 38% slower; with the pool it costs what the unpooled copies did. With warm caches, a 1,000-event page now makes 2,010 allocations instead of 4,008 hot and about 50 instead of 1,025 cold, and eight concurrent fetchers move 10% more hot pages and 17% more cold pages per second, with cold p99 21% lower. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01JCzMKKrK3NstGKvaxsHhL2 --- .../internal/rpcv2/rocksdb/rocksdb.go | 53 +++++++++----- .../internal/rpcv2/stores/event/arena.go | 31 ++++++++ .../internal/rpcv2/stores/event/arena_test.go | 37 ++++++++++ .../rpcv2/stores/event/cold_reader.go | 15 +++- .../stores/event/cold_reader_fanout_test.go | 38 ++++++++++ .../internal/rpcv2/stores/event/hot_store.go | 26 +++++-- .../rpcv2/stores/event/hot_store_test.go | 71 +++++++++++++++++++ .../internal/rpcv2/stores/event/payload.go | 9 +-- 8 files changed, 253 insertions(+), 27 deletions(-) create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/arena_test.go create mode 100644 cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_fanout_test.go 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_reader.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go index bf933a813..819648fd5 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. Nor do the payload copies: each callback copies +// into an arena of its own (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,20 @@ func (c *ColdReader) FetchEvents(ctx context.Context, eventIDs []uint32) ([]Payl positions[i] = int(id) } results := make([]Payload, len(eventIDs)) + // A call's payloads live and die together, so they are copied into + // arenas rather than one allocation each. ReadItems calls back from up + // to Concurrency goroutines and an arena is not safe for concurrent + // use, so each callback takes one from the pool for its copy. + arenas := sync.Pool{New: func() any { return new(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)) + a, _ := arenas.Get().(*byteArena) + owned := a.copy(data) + arenas.Put(a) + 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..ab759d002 --- /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 arenas against the reader's own fan-out: with Concurrency +// above one, ReadItems calls the callback from one goroutine per batch, so two +// callbacks appending into one arena corrupt payloads or segfault. +// 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, ColdDirs{Data: dir, Index: 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/hot_store.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go index ec2e3730c..c98f7faa2 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go @@ -232,10 +232,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 +672,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/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