Repository navigation
events: snapshot dense index terms lazily so apply latency does not grow with chunk fill - #973
Conversation
There was a problem hiding this comment.
Pull request overview
Lazily snapshots dense event-index bitmaps to prevent ingest latency from growing with chunk size.
Changes:
- Adds writer-private dense bitmaps with lazy immutable snapshots.
- Adds concurrency, freshness, immutability, and warmup tests.
- Updates existing tests for the new dense-state model.
Reviewed changes
Copilot reviewed 3 out of 3 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
concurrent_bitmaps.go |
Implements lazy dense snapshots. |
concurrent_bitmaps_writer_test.go |
Adds writer/reader concurrency tests. |
concurrent_bitmaps_test.go |
Adapts existing dense-state tests. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
4e53849 to
20566b3
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 3 out of 3 changed files in this pull request and generated 1 comment.
Suppressed comments (1)
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps.go:24
- The “allocates nothing per AddTo” claim is incorrect: after a snapshot,
AddManyallocates when it copy-on-writes a touched container, and normal container growth can also allocate, as the later AddTo documentation acknowledges. Please narrow this claim so the performance contract is accurate.
// denseState is the whole state of one dense term. Exactly one is
// allocated per promoted term and it lives for the term's lifetime,
// so a dense term allocates nothing per AddTo and nothing per
// published snapshot beyond the snapshot itself.
…g with chunk fill Hot-chunk ingest gets slower as a chunk fills. On EC2, across one 10,000-ledger chunk on the sac profile, the apply phase p99 goes 28 -> 344 ms and ingest p99 163 -> 508 ms, monotonically with chunk fill rather than with load. The growth term is ConcurrentBitmaps.AddTo on dense terms: every call did Clone() + AddMany + publish. With copy-on-write enabled, roaring's Clone still allocates and copies the keys/containers/needCopyOnWrite slices (sized by the container count, which grows with the chunk's event-ID span) and marks every container on both bitmaps as needing copy, so the following AddMany deep-copies the touched container again on every ledger. A chunk with many dense terms (SAC from/to address topics) pays that once per (term, ledger). Change, one file plus tests: - A dense term owns a writer-private bitmap. AddTo mutates it in place under a per-term mutex and sets a dirty flag. No clone, no allocation, no publish. - The immutable snapshot readers receive lives in an atomic pointer next to it. Get returns it lock-free when the term is clean; otherwise the first reader clones the writer's bitmap once, publishes it, and clears the flag. - Get reads the dirty flag before the snapshot pointer, and a republish stores the pointer before clearing the flag, so a Get that starts after an AddTo returns observes that AddTo's IDs. - Snapshots are lazy at promotion and in NewConcurrentBitmapsFromBitmaps: a dense term that is written but never read never clones. - termState stays 32 bytes; sparse terms are unchanged. Local harness, 10,000 ledgers, first to last 1,000-ledger window: bounded account pool apply p50 2.3 -> 8.7 ms becomes 1.9 -> 2.1, p99 7.8 -> 55.8 becomes 3.1 -> 4.2; 4,000-contract pool p99 11.5 -> 61.1 becomes 6.7 -> 8.3. Under 2,000 concurrent reads/s readers pay +0.9 us p50 / +7.5 us p99 and apply stays flat. Semantics preserved: immutable snapshots, term-by-term visibility, same pointer when nothing changed, single-writer idempotent AddTo. Changed: a reader of a dense term written since the last read blocks for one AddMany and pays one shallow clone. Complementary to #902, which also bounds memory via a tail-delta index and a new row format; redundant if #902 lands first, no conflict in approach.
20566b3 to
bf42c94
Compare
|
@claude review |
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated no new comments.
Suppressed comments (3)
Previously missed (1) — in code that hasn't changed since the last review.
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:112
- This race test can pass without exercising any concurrent read: the writer is launched before the readers, and a valid schedule lets it set
donebefore they enter their loops. Gate both sides on a start barrier (and signal reader readiness before releasing the writer) so the test reliably coversAddManyoverlappingGet.
This issue also appears in the following locations of the same file:
- line 169
- line 405
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:173
- The writer may complete and set
donebefore any reader starts, in which case this test only checks the final bitmap and never exercises the advertised promotion-under-readers path. Add a readiness/start barrier so at least one reader is active before promotion begins.
wg.Go(func() {
defer done.Store(true)
for i := range uint32(total) {
s.AddTo(key, i)
writes.Store(i + 1)
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:407
- These liveness assertions are scheduler-dependent. Because writes are intentionally coalesced while
pubis nil, a correct run may publish far fewer than 1,000 snapshots (or the writer may finish before any reader runs, yielding zero reads). Synchronize reader startup with the writer and assert freshness from deterministic handshakes rather than requiring an arbitrary number of observed pointer changes.
require.Positive(t, reads.Load(), "the stress loop must have done real reads")
require.GreaterOrEqual(t, republishes.Load(), uint64(numBatches/10),
"readers must observe snapshots being republished, not one frozen bitmap")
a82cc31 to
64bd0cf
Compare
64bd0cf to
6508a7a
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated 3 comments.
Suppressed comments (2)
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:169
- This concurrency test can be vacuous: the writer is started first and may set
donebefore the readers are scheduled, after which all reader loops skip and only the final-state assertion runs. Add a startup/interleaving barrier and verify that reads occurred while promotion/writes were active.
wg.Go(func() {
cmd/stellar-rpc/internal/rpcv2/stores/event/reader.go:108
- This updated contract leaves the preceding “no per-key Clone” claim stale. A lookup of each dirty dense term now calls
snapshot()and clones that term, so describe this as an in-memory lookup with a possible lazy clone rather than no clone at all.
// Bitmap ownership: callers MUST treat returned bitmaps as
// read-only. The hot path returns immutable snapshots of the
// live mirror — ConcurrentBitmaps mutates a writer-owned bitmap
// under a per-term mutex and publishes snapshots lazily, so a
// returned pointer will never be mutated by anyone; see
// ConcurrentBitmaps.Get for the full contract. The cold path
6508a7a to
8e8c2ed
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated 1 comment.
Suppressed comments (2)
cmd/stellar-rpc/internal/rpcv2/stores/event/reader.go:103
- This describes all hot snapshots as shared, but sparse terms take the
st.idspath and build a fresh bitmap on every lookup. Documenting the two ownership modes separately keeps the Reader contract aligned withConcurrentBitmaps.Get.
// snapshots of the live mirror shared by all readers of a term;
// a dense term written since its last lookup is cloned once, by
// the first reader to look it up, and that clone is then shared.
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:407
- This is scheduler-dependent and
republishesis not a distinct-snapshot count: each reader has its ownlast, so the same pointer is counted once per reader. The writer may also finish before readers run, leavingreads == 0. Use a start barrier and a deterministic handshake around known publications instead of asserting an aggregate timing-dependent count.
require.Positive(t, reads.Load(), "the stress loop must have done real reads")
require.GreaterOrEqual(t, republishes.Load(), uint64(numBatches/10),
"readers must observe snapshots being republished, not one frozen bitmap")
8e8c2ed to
5713751
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated no new comments.
Suppressed comments (2)
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:196
- This loop can likewise run entirely after all 512 writes:
ready.Doneonly says the reader goroutine started, and the writer may complete immediately afterready.Waitunblocks. In that schedule the test never observes the sparse-to-dense transition it is intended to cover. Synchronize at least one reader observation before promotion and one after the writer begins crossing the threshold, rather than relying on scheduler interleaving.
for n := 0; !done.Load() || n < minReads; {
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:133
- The readiness barrier does not guarantee any read overlaps the writer: after the last
ready.Done, the scheduler can run the writer to completion before a reader enters this loop, andminReadsthen counts only post-write reads. That lets this race regression pass without exercising concurrentAddMany/snapshot access. Add a deterministic during-write handoff or count reads performed whiledoneis false and require at least one before allowing the writer to finish.
for i, n := 0, 0; !done.Load() || n < minReads; i++ {
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7opQUW9UZPKzE6B3tz9b5
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated no new comments.
Suppressed comments (3)
Previously missed (1) — in code that hasn't changed since the last review.
cmd/stellar-rpc/internal/rpcv2/stores/event/reader.go:108
- This ownership description is inaccurate for sparse hot terms:
ConcurrentBitmaps.Getbuilds and returns a fresh bitmap on every lookup for those terms (concurrent_bitmaps.go:116-118), so only dense hot snapshots are shared. Please distinguish the two representations so this public reader contract accurately describes bitmap ownership.
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:137
- The new barrier and
minReadscondition still do not guarantee reader/writer overlap: once all readers callready.Done, the writer may run to completion, after which each reader can satisfy allminReadsiterations against the final snapshot. The test can therefore pass under-racewithout exercising concurrentGetandAddTo; add a deterministic handoff that requires a read between writer batches, as the freshness stress test below does.
for i, n := 0, 0; !done.Load() || n < minReads; i++ {
key := keys[(i+r)%numTerms]
bm, err := s.Get(key)
if err != nil || bm == nil {
continue
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:200
- This loop has the same scheduling gap:
ready.Doneonly proves that reader goroutines started, whilen < minReadsallows every checked read to happen afterdonebecomes true. Thus the test can pass without observing the sparse-to-dense transition concurrently; synchronize writer progress with reader acknowledgements so at least the promotion and subsequent writes are read while the writer remains active.
for n := 0; !done.Load() || n < minReads; {
bm, err := s.Get(key)
if err != nil || bm == nil {
continue
}
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7opQUW9UZPKzE6B3tz9b5
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 5 out of 5 changed files in this pull request and generated no new comments.
Suppressed comments (2)
Previously missed (1) — in code that hasn't changed since the last review.
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:196
- This barrier has the same scheduling gap at the promotion boundary: the writer may execute all 512 additions immediately after
ready.Wait()unblocks, after which readers perform their required reads only against the final dense state. In that execution the test never observes the sparse-to-dense transition. Pause the writer aroundpromotionThresholdand require reader progress on both sides of the promotion before allowing it to finish.
cmd/stellar-rpc/internal/rpcv2/stores/event/concurrent_bitmaps_writer_test.go:133
- The readiness barrier does not guarantee that any
Getoverlaps the writer: after the lastready.Done(), the scheduler may run the writer to completion before a reader enters this loop, andminReadsthen only counts post-write reads. The test can therefore pass without exercising the reader/writer race it is intended to cover. Add a two-way progress handshake that keeps the writer active until readers have performed work during the write phase.
ready.Done()
for i, n := 0, 0; !done.Load() || n < minReads; i++ {
…snapshot #973 made a dense term's reader-visible bitmap a lazily published snapshot: the writer mutates wbm in place under the term mutex and drops pub, and a reader clones once when it finds pub nil. This branch's no-materialize seam — lookupPostings, postings.bitmap, postings.estimate — is now built on that, and its failure mode is invisible to every test the branch already has: sources are resolved once up front and only read, so an accessor handing back a stale or nil pub agrees with the materialized twin on all of them. So the accessors get their own write-then-read-through-the-accessor tests, the only shape that separates the two. They cover a sparse term, the AddTo that promotes one, a dense term whose snapshot the writer has just dropped, the estimate path (fresh count, and still no clone), and the same race the index's own freshness stress runs, but driven through lookupPostings. Each fails against a lookupPostings that returns pub directly, an AddTo that forgets to drop pub, and — for the estimate case — an estimate that reaches for snapshot(). HotStore.lookupPostings' doc said "borrowed snapshot", true of the representation #968 was written against; it now names snapshot's shared clone, matching what #973 put in Reader.LookupKeys. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7opQUW9UZPKzE6B3tz9b5
…snapshot #973 made a dense term's reader-visible bitmap a lazily published snapshot: the writer mutates wbm in place under the term mutex and drops pub, and a reader clones once when it finds pub nil. This branch's no-materialize seam — lookupPostings, postings.bitmap, postings.estimate — is now built on that, and its failure mode is invisible to every test the branch already has: sources are resolved once up front and only read, so an accessor handing back a stale or nil pub agrees with the materialized twin on all of them. So the accessors get their own write-then-read-through-the-accessor tests, the only shape that separates the two. They cover a sparse term, the AddTo that promotes one, a dense term whose snapshot the writer has just dropped, the estimate path (fresh count, and still no clone), and the same race the index's own freshness stress runs, but driven through lookupPostings. Each fails against a lookupPostings that returns pub directly, an AddTo that forgets to drop pub, and — for the estimate case — an estimate that reaches for snapshot(). HotStore.lookupPostings' doc said "borrowed snapshot", true of the representation #968 was written against; it now names snapshot's shared clone, matching what #973 put in Reader.LookupKeys. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7opQUW9UZPKzE6B3tz9b5
Three cuts at the hot ingest tail, in causal order. The zstd Encode destination is pooled — a fresh ~compressBound dst per ledger was ~43% of hot ingestion's allocations — and the pool is promoted to the canonical internal/rpcv2/zstd.Pool (CGo context + retained dst buffer; the consumer contract is synchronous copy, which BatchWriter.Put honors). The events hot index writes ONE packed row per ledger instead of a row per (term, event) — tens of thousands of memtable keys per ledger gone from the shared commit batch. The ~21ms ledger compression forks off the critical path at IngestLedger entry and joins as the batch's last queue step (AddLedgerToBatch remains as the composition of the pair, one write path). A fourth cut in the original branch made the dense ConcurrentBitmaps terms publish through an immutable tail, so a small AddTo never cloned the base bitmap. It is dropped here: #973 landed on the base branch first and removes the same per-write clone by snapshotting a dense term lazily under its own mutex, so the clone population that cut existed to delete is already gone, and carrying a second mechanism for it would mean two designs for one problem. sac-6000 @600ms, 1k: ingest_total p99 303 -> ~83ms across the four cuts as originally measured, with the dense-overlay share of that now coming from #973 instead (final figures per constituent benches; GOGC-free). The encode state is store-OWNED, not sync.Pool'd: the write side is single-flight by contract (one compression in flight, hotchunk's single-writer loop; loud CAS latch), reads stay fully concurrent with the writer. A sync.Pool here was measured losing its lone expensive state (~15MB dst + CGo context) about 1-in-5 ledgers to GC pool-emptying and silently re-allocating it; the GC-survival regression test pins the guarantee. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01C3xaaoD9LhzyKPdnZCJa2N
…lized
index.pack stored every term's postings as a roaring bitmap. Measured on a
real pubnet chunk (8.65M events, 204,752 terms), per-term min(roaring,
delta) is 22.23MB against 37.74MB for roaring alone, a 41% saving, and 71%
of it sits below 1024 postings where roaring costs 1.7x to 3.8x a delta
list. Above that, run and bitmap containers pull ahead, reaching 12.5x on
the nine terms holding more than a million postings. So an index.pack item
is now `fingerprint ‖ codec ‖ postings`, delta-varint at or below 1024 and
roaring above.
Storing them that way alone would have made queries slower: decoding a
1024-posting delta term into a bitmap costs 20.6us against 5.7us to
unmarshal an equivalent roaring body. The saving only materializes if the
planner stops materializing, so the second half of this change is
events.Postings, which carries a term in whichever form the store already
had it, and the set operations that work on either.
- Reader.LookupKeys and HotIndex.Get return Postings instead of bitmaps.
ConcurrentBitmaps.Get does too, which stops materializing its sparse
entries: that state already holds an id list, and building a bitmap
out of it on every lookup was pure loss. A dense term still hands
back the snapshot #973 publishes for it, now wrapped as Postings.
- events.Intersect drives from the smallest side whatever its form and
probes the rest, since probing beats materialize-then-FastAnd by 4.3x
at cardinality 4, 2.8x at 64 and 1.6x at 1024.
- events.Union merges ascending lists rather than building a bitmap per
input, and keeps the result a list.
- Postings.ClipRange and Postings.SelectIDs replace the range bitmap and
the bitmap drain.
query.go no longer imports roaring. Both builders (in-memory WriteColdIndex
and streaming WriteColdIndexFromRuns) share one encoder, so the
byte-identity gate between them stays meaningful.
One unbounded allocation goes away on the list path: selectEventIDs used to
do make([]uint32, cardinality) whenever MaxEvents was 0, which is 34MB per
query on a real chunk.
Two defects found in review and fixed here rather than shipped: ClipRange
and SelectIDs handed out subslices of the store's live posting list without
clamping capacity, so a downstream append would have overwritten a posting
other readers could still see, invisible to -race because it writes past a
published length; and roaring accepts a run container holding no intervals,
which reads back as a bitmap that reports itself non-empty with zero
postings, so the cold decode now validates and rejects it.
FORMAT: this changes the on-disk index.pack record layout. Split out of the
"events index follow-ons" commit so it can be evaluated, cherry-picked or
reverted independently of the bloom/fence and hardening work it shipped
alongside — cold artifacts are permanent, so this is the piece that must be
right in the first v2 release.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Note
This is a refactor of one of the optimizations originally introduced by @tamirms in #902.
👉 concurrent_bitmaps.go#902
Problem
Hot ingest latency grows as the chunk fills. For one 10,000-ledger chunk:
applyp99 goes from 28 ms (first 1,000 ledgers) to 344 ms (all 10,000). Ingest p99 goes from 163 ms to 508 ms. The other five phases do not change.Cause:
ConcurrentBitmaps.AddToclones the term bitmap on every write to a dense term. With copy-on-write enabled, roaringClonestill allocates and copies thekeys,containersandneedCopyOnWriteslices and marks every container on both bitmaps for copy. The nextAddManythen deep-copies the touched container. The cost is O(containers) per (term, ledger), and the container count grows with the chunk.Change
Changes how dense terms in the live index are updated: the writer mutates a private bitmap in place under a per-term mutex instead of cloning on every write, and readers get a snapshot that the first reader after a write clones once and shares.
Result
Full-length phase 2 run on c6id.8xlarge, hot leg 10,000 ledgers.
Before: phase2-pr961-full-c110e623-20260831T141126Z. After: phase2-lazy-snapshot-optimization-full-6cd08c13-20260903T020032Z.
closes #974