Conversation
|
Review the following changes in direct dependencies. Learn more about Socket for GitHub.
|
There was a problem hiding this comment.
Pull request overview
Optimizes getEvents across hot and cold storage while fixing concurrent cold-fetch safety.
Changes:
- Adds lazy ascending event matching and broader fetch batching.
- Reduces allocations through arenas, pooled packfile buffers, and mmap-backed indexes.
- Enables concurrent cold reads and upgrades RoaringBitmap.
Reviewed changes
Copilot reviewed 18 out of 19 changed files in this pull request and generated 3 comments.
Show a summary per file
| File | Description |
|---|---|
go.mod |
Upgrades bitmap dependencies. |
go.sum |
Updates dependency checksums. |
event/arena.go |
Adds payload byte arena. |
event/arena_test.go |
Tests arena stability. |
event/cold_format.go |
Opens MPHF indexes via mmap. |
event/cold_reader.go |
Makes concurrent arena use safe. |
event/cold_reader_fanout_test.go |
Tests concurrent cold fetching. |
event/concurrent_bitmaps.go |
Exposes immutable postings views. |
event/hot_store.go |
Uses postings fast path and smaller blocks. |
event/match.go |
Splits ascending and descending matching. |
event/match_iter.go |
Implements lazy iterator-tree matching. |
event/match_iter_test.go |
Tests and benchmarks iterator matching. |
event/matches_differential_test.go |
Adds randomized differential coverage. |
packfile/index.go |
Pools index decoding allocations. |
packfile/pools.go |
Adds packfile allocation pools. |
packfile/pools_test.go |
Tests pooling and close/read interaction. |
packfile/reader.go |
Pools open state and guards offset recycling. |
query/resolve.go |
Enables eight-way cold-read fan-out. |
rocksdb/rocksdb.go |
Consolidates batch value copies. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 19 out of 20 changed files in this pull request and generated no new comments.
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
cmd/stellar-rpc/internal/rpcv2/query/resolve.go:183
- The production override described here is not reachable:
coldEventReadConcurrencyis an unexported field, neitherOpenRegistrynorNewRegistryaccepts it, and startup constructs the registry throughquery.OpenRegistry. Consequently every deployment is forced to fan out eight workers per cold page, despite the comment above noting that this multiplies per-request goroutines and 1 MiB buffers. Please thread an operator setting through startup/registry construction (with validation), or remove the claimed override and choose a globally safe bound.
conc := a.coldEventReadConcurrency
if conc == 0 {
conc = defaultColdEventReadConcurrency
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 19 out of 20 changed files in this pull request and generated no new comments.
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
cmd/stellar-rpc/internal/rpcv2/query/registry.go:45
- This override is not reachable by a deployment:
Registryis constructed throughOpenRegistry, the field is unexported, and there is no constructor option or setter for it. Consequently production is still fixed at 8 despite the comments saying operators can raise it. Thread the setting through the public registry construction/configuration path, or describe this as a test-only seam instead.
// coldEventReadConcurrency overrides the packfile read fan-out one cold
// events read gets. Zero means defaultColdEventReadConcurrency, resolved
// where the reader is opened (ReadView.Events). Held here, like
// maxScanLedgers, so a deployment with I/O headroom to spend — or a test
// pinning the serialized path — sets it on its Registry rather than
// rebuilding.
coldEventReadConcurrency int
cmd/stellar-rpc/internal/rpcv2/stores/event/arena.go:24
- The arena’s first non-empty copy always allocates at least 64 KiB.
getEventsaccepts a limit of 1, and the pager passesroom+1toMatches, so a small cold page can fetch only two ~250-byte payloads yet now allocates 64 KiB instead of roughly the bytes returned. This regresses allocation volume for low-limit requests; size the initial chunk from the requested batch (or grow from a smaller first chunk) and reserve 64 KiB chunks for batches that actually reach that size.
func (a *byteArena) copy(b []byte) []byte {
if len(b) > cap(a.buf)-len(a.buf) {
a.buf = make([]byte, 0, max(arenaChunkSize, len(b)))
}
806d538 to
fe01e7c
Compare
Drop the tests that duplicate what the whole-stream comparisons already
prove or pin an internal heuristic rather than behavior:
- requireStrictOrder, redundant with the whole-stream equality every
randomized trial already asserts
- TestResolveSlabFiltersOrdersRarestFirst, which pinned the rarest-first
ordering heuristic; the absent-group drop it also covered is exercised
by TestResolveSlabTerms and the shaped matrix
- the FastAnd and FastOr subtests of the roaring contract, since the
match path only calls AndAny, And, NextValue and PreviousValue on
shared bitmaps
- TestRoaringContract_ConcurrentValueSearch, folded into the
aggregation race gate as TestRoaringContract_ConcurrentReaders
Rename TestMatches_WindowANDLeavesBorrowedBitmapUntouched, whose doc
described a branch that no longer exists, and shorten the remaining test
comments.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01M23ZvEkyrFjyUUbm7zLwob
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 23 out of 24 changed files in this pull request and generated 1 comment.
Suppressed comments (1)
cmd/stellar-rpc/internal/rpcv2/stores/event/match.go:230
- Every positive hint becomes the initial slice capacity in both filtered and descending match-all streams. v1 only bounds the request by the operator-configured maximum, which may be raised as high as
MaxInt32; a selective request under such a setting can therefore allocate billions of entries before discovering that only a few events match. Preserve the 10,000-item fast path, but impose a hard internal cap or chunk oversized first fetches—the hint changes I/O shape, not result semantics.
if hint > 0 {
first = hint
}
| close(parked) | ||
| <-unpark |
Fixes the two prealloc findings by moving the NextValue/PreviousValue loop into one helper that sizes its result up front. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01M23ZvEkyrFjyUUbm7zLwob
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 23 out of 24 changed files in this pull request and generated no new comments.
Suppressed comments (2)
cmd/stellar-rpc/internal/rpcv2/packfile/pools_test.go:85
- This does not actually verify that
Closewithheld the offsets from the pool: the sole item has already been located and decoded before this callback runs, and no subsequent open overwrites a pooled array, so an implementation that always recycledr.offsetswould still pass the byte check. Assert here, afterClosehas returned while the callback still keepsinflightnonzero, thatr.offsetsremains non-nil.
<-unpark
cmd/stellar-rpc/internal/rpcv2/query/resolve.go:155
- This hard-codes eight 1 MiB packfile worker buffers per in-flight cold read for every deployment, even though the comment notes that the optimum depends on storage and scales with concurrent pages. Deployments on constrained disks or memory limits cannot tune the new fan-out without rebuilding. Please carry a registry/configured concurrency value into each
ReadView(with zero selecting this default), as is done for other per-view serving bounds.
const defaultColdEventReadConcurrency = 8
eval narrowed a per-filter accumulator through the filter's groups one AndAny at a time, so every filter paid an 8 KiB accumulator on every slab even when its groups had nothing in common there, and a topic-count range paid its union on top. A filter now yields plans that are conjunctions of single terms, one per bucket of a range, since A and (b or c) is (A and b) or (A and c) and the union across plans keeps the or; the planner drops a plan that repeats an earlier one. The slab engine then holds one bitmap per term, runs one FastAnd over the slab's id range and the terms per plan, keeps its bounds as the searches return them, and counts each term's cardinality once. The slab bitmap is built once per slab and shared by every plan. With the pinned roaring, FastAnd intersects pairwise and the cost is unchanged or lower. With roaring's count-first intersections (RoaringBitmap/roaring#566 and the FastAnd follow-up built on it) an empty result allocates nothing, which is what bounds the slab walk on filter lists that match nothing; the version bump will add the contract test for that. roaring_contract_test.go now pins that FastAnd leaves its inputs untouched, returns fresh containers on the two- and many-input paths over whole and partial slab ranges, and agrees with the pairwise And chain for every subset of the container kinds; that a slab range built with AddRange is stored as run containers; and runs FastAnd in the concurrent-readers race gate. The AndAny pins go, since the engine no longer calls it. The shaped differential matrix gains a filter with a term and a count range, and a planner test covers the fan-out and the dedup. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01M23ZvEkyrFjyUUbm7zLwob
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 23 out of 24 changed files in this pull request and generated 1 comment.
Suppressed comments (1)
cmd/stellar-rpc/internal/rpcv2/stores/event/match.go:234
- This makes the initial allocation and fetch unbounded. v1 has no upper validation for
max_items_per_responseand passes configured limits through up tomath.MaxInt32, so an operator-raised cap can makestreamSlabsallocate a multi-gigabyte ID slice (and descending match-all allocate an equally unboundedMatchblock) before serving anything. Bound the first batch independently of the page limit; subsequent batches can still complete larger pages.
if hint > 0 {
first = hint
The sentence read as if an empty plan cost no container today; at roaring v2.26.0 FastAnd intersects pairwise and allocates the first intermediate, and the count-first behaviour arrives with the roaring bump. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01M23ZvEkyrFjyUUbm7zLwob
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 23 out of 24 changed files in this pull request and generated no new comments.
Suppressed comments (1)
cmd/stellar-rpc/internal/rpcv2/stores/event/match.go:234
- A positive hint is now used as an allocation capacity with no internal ceiling (
streamSlabsallocatesidsat that capacity, and descendingstreamRangedoes the same forblock). This is not bounded by validation for v1:max_items_per_responsemay be raised arbitrarily, so an accepted request can force a multi-gigabyte allocation/OOM before the window is even clipped. Keep the first fetch optimization, but cap the internal batch to a bounded multiple ofmatchBatchSizeand let subsequent batches serve the rest of the page.
if hint > 0 {
first = hint
|
Split into smaller PRs. Everything in this one now lives in one of them:
#1040–#1043 are stacked on #1034. #1045–#1047 are independent. |
getEvents read-path performance on both tiers: the hot RocksDB store, the cold packfile store, and the match engine they share. It also fixes a latent race in the cold fetch arena.
Match engine
What the engine does. A getEvents query is a list of filters. Each filter names constraints: a contract ID, a topic at a position, an event type, a topic count. The index maps every term to a bitmap of the event IDs in the chunk that carry it. The candidates for a query are the IDs that satisfy at least one filter, and a filter is satisfied when all of its constraints hold: an OR across filters of an AND within each filter, restricted to the window: the range of event IDs inside the chunk that the client's ledger range covers. Candidates are then fetched and post-filtered, and the client receives at most a page of them, its limit.
How the base branch does it. All at once, over whole-chunk bitmaps: intersect each filter's term bitmaps across the entire chunk, union the results across filters, and intersect with the window last. Everything is computed before the first event is fetched. For any query with more than one term or filter, the cost tracks the size of the terms across the whole chunk rather than the window the client asked for or the page it will receive. A page of 1,000 events from a popular contract with a topic constraint pays to intersect that contract's millions of IDs first.
How this PR does it. Roaring bitmaps store IDs in containers of 65,536 consecutive IDs, and every bitmap operation works one container at a time. The engine follows that grain. A filter is first turned into plans, each a plain AND of single terms: a topic-count range, the one constraint that is an OR of terms, becomes one plan per bucket, since A and (b or c) is (A and b) or (A and c), and a plan that repeats an earlier one is dropped. The walk then visits the window one container-sized slab at a time, in the requested direction, and evaluates every plan on that slab alone, one
FastAndover the slab's IDs inside the window and the plan's terms, ORs the results across plans, and hands the union to the fetch. Every intermediate value is a single container, so the work per slab is bounded, and a page that fills within a couple of slabs never touches the rest of the chunk. Ascending and descending share one code path; only the walk order differs.Skipping slabs. A rare term over a wide window would still visit every slab, so before evaluating one the walk asks each term for its next ID at or after the cursor,
NextValueon a roaring bitmap, which is cheap. Within a plan every term must hold, so nothing can match before the term whose next ID is furthest away. Across plans any may match, so the earliest wins. The walk jumps straight to the slab holding that ID. Each plan's bound is kept until the cursor passes it, and a plan with no IDs left is retired. Descending usesPreviousValueand the mirror-image rules. A query for a rare term opens only the slabs where that term has events.Term lookups. Both tiers serve the walk through
Reader.LookupKeys, which returns each term's bitmap once at query start. The walk holds those bitmaps for its whole duration and never mutates them, so a query sees one consistent image of the index while ingest continues, the same guarantee the pinned window already gives. On the hot tier they are snapshots shared between readers; on the cold tier they are fresh copies.Why it is better. Work scales with the window and the page instead of the chunk: the walk opens only the slabs the window holds and stops once the page is full, so a limit-1 request that finds a match in its first slab pays for one slab whatever the size of the chunk, the ten-thousand-ledger unit of storage. Memory per query is bounded by a container instead of by chunk-sized intermediates. One engine serves both directions, replacing two paths.
roaring moves from v2.18.2 to v2.26.0. At that version
FastAndintersects pairwise, so a plan that matches nothing on a slab still allocates one intermediate there. RoaringBitmap/roaring#566 and theFastAndchange built on it count a slab's result before allocating it, which makes an empty plan cost no container at all; that bump is a follow-up, and its allocation pin forroaring_contract_test.gois written.Hot tier
dataCFBlockSize), the unit getEvents decompresses per fetched event. Existing databases read fine, since SSTs self-describe their block size, and converge through compaction. Ingest throughput is unchanged; disk grows 0.8 %.TestHotStore_FetchEventsAllocationBudget.Cold tier
defaultColdEventReadConcurrency). The fetch arena is now locked against that fan-out. Before this PR, any cold read concurrency above 1 produced SIGSEGVs or silently short results; the wired value of 1 kept the race latent.pooled_buffer_cap_skips_total; a nonzero count means a cap wants raising.Tests
TestMatches_FetchesOnlyTrueCandidatestraces every fetch and requires the fetched ordinals to be exactly the query's matches, in emission order. Output equality alone cannot see an over-wide candidate set, because the post-filter hides it as wasted I/O.slab_match_test.gocompare against an oracle computed without the index, at several slab widths.TestSlabStepperSkipsCandidateFreeSlabspins the exact slab list the walk opens per shape;TestMatches_RareTermsSpanSlabsruns the skip-dominated shapes at widths that put hundreds of empty slabs between candidates.roaring_contract_test.gopins, against roaring alone, the properties the shared index bitmaps rest on:FastAndleaves copy-on-write arguments byte-identical, returns containers that share no storage with them, and equals the pairwiseAndchain, for every subset of the container kinds over whole and partial slab ranges; a range built withAddRangeis stored as run containers; andNextValue/PreviousValueare inclusive and read-only, with eight concurrent readers under-race. It is the first thing to run on a roaring bump.TestMatches_ConcurrentIngestBorrowSafetyholds index snapshots across a walk while ingest publishes new term states on the same keys.packfile/pools_test.go; the fan-out test compares whole payloads.Measured impact
Service layer, 50 rps, K-corpus at limit=1000; sac = the phase-3 stress chunk, pubnet = chunk 6410 (9.1 M events). Baseline is the branch base (32 KiB blocks, cold concurrency 1); after is this PR.
Hot
Cold
Item counts are identical on both sides of every A/B cell (2,800,552 / 1,694,968).
Filter lists. Engine alone on one 8.4 M-event chunk (128 slabs), Xeon 8375C, roaring v2.26.0, limit 1 unless noted; before is this branch before the per-plan engine, after is this head. Wall time and allocation per request.
Matching rows are unchanged. With the count-first roaring the three never-matching rows fall to 5 to 16 ms and 1 to 4 MB.
🤖 Generated with Claude Code