Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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
53 changes: 37 additions & 16 deletions cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
}
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand Down
31 changes: 31 additions & 0 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go
Original file line number Diff line number Diff line change
@@ -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)]
}
37 changes: 37 additions & 0 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/arena_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
15 changes: 12 additions & 3 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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)
}
}
26 changes: 22 additions & 4 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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[:])
Expand Down
71 changes: 71 additions & 0 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/hot_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
9 changes: 5 additions & 4 deletions cmd/stellar-rpc/internal/rpcv2/stores/event/payload.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading