ingest: full-history ingestion for cold + hot stores (#765) - #779
Conversation
…tmaps Rework the in-memory events index and the full-history eventstore as one unit — they are tightly coupled, since the eventstore consumes the events index types and the membitmaps removal forces the eventstore migration in lockstep. The query engine (Query + postfilter) is deliberately left out and follows in a separate stacked PR. events: - Add ConcurrentBitmaps and ConcurrentLedgerOffsets for lock-free concurrent reads during ingest; remove the old membitmaps implementation. - Add the ingest_view path: build the index directly from xdr LedgerCloseMetaView / TransactionMetaView zero-copy views. - payload / index / ledgeroffsets / bitmaps reworked accordingly. eventstore: - Migrate the cold store (format / index / reader / writer) and the hot store onto the new events index API. - reader.go interface updates. deps: roaring v2.18.0 -> v2.18.2 (upstream FastOr/runContainer16 fix), go-stellar-sdk bump (XDR View types used by ingest_view), and tamirms/streamhash promoted to a direct dependency. The eventstore.Query concurrency test (TestHotStore_QueryUnderConcurrentIngest) moves to the query-engine PR alongside query.go. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
The view-based ingest extractor (ingest_view.go / LCMToPayloadsFromRaw) and its Payload term-precompute plumbing move to the separate #764 work, so remove them here along with the now-dead TermKeys() skip-branch in IngestLedgerEvents and the V3-SorobanMeta test fixture only the extractor's test used. Payload now carries the event purely as raw XDR (ContractEventBytes) — the decoded xdr.ContractEvent field is gone. Ingest (LCMToPayloads) marshals each event into the bytes, and terms are derived straight from the raw XDR via xdr.ContractEventView (events.TermsForBytes), no full UnmarshalBinary. Remove the per-Reader useXDRViews toggle from HotStore and ColdReader; the read path always decodes via views. Payload.Unmarshal is the sole consumer decoder (struct decoder removed; former UnmarshalView renamed to Unmarshal). FetchEvents returns owned Payloads; FetchRange/All yield borrowed Payloads (ContractEventBytes aliases the iterator's step buffer — clone to retain). IngestLedgerEvents marshals each payload into one reused scratch buffer (BatchWriter.Put copies the value synchronously), and is idempotent on retry: re-ingesting an already-committed ledger is a no-op (a gap or out-of-range ledger still errors). Warmup now cross-checks the per-chunk CFs on open (verifyChunkConsistency): the index may not reference an event beyond the committed count, and the data tail must align with it (event total-1 present, nothing at id >= total) — a corrupt or tampered chunk fails to open loudly instead of serving an inconsistent cache. ConcurrentLedgerOffsets.Append is now a single positional primitive (no ledger arg, no error); the sequence, capacity, and cumulative-overflow checks live at the warmup trust boundary in warmupOffsets, where on-disk rows are untrusted. Deps: - go-stellar-sdk -> latest main (v0.5.1-0.20260604220920-ff1e140adca5) - streamhash -> github.com/stellar/streamhash (was tamirms/streamhash) - roaring/v2 unchanged at v2.18.2 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
dc95b0b to
b8b1a59
Compare
|
Updated to align with #780 (the authoritative txhash cold store, #728): this PR no longer bundles a txhash cold store — it now leaves The integration seam is the per-chunk |
… ingest-core deps
Four zero-copy XDR view extractors over xdr.LedgerCloseMetaView (events, tx-hashes, tx-details-by-hash, tx-pages); outputs alias the view buffer. Transactions paired to TxSet envelopes by hash (mirroring ingest.LedgerTransactionReader); V3 contract events gated on IsSorobanTx. Differential-tested vs the parsed / db.ParseTransaction path across LCM V0/V1/V2, meta V1-V4, V0Components + ParallelTxs, order mismatch, diagnostic events, empty, sponsorship, large-tx, and protocol-transition fixtures; aliasing + negative-path coverage; per-extractor benches. Leaf package (no internal/db dep). Closes #764.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 686f3b8703
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
…, error-chain helper - ExtractEvents/ExtractTransactions/ExtractTxDetailsByHash now treat legacy TransactionMeta V0 (pre-Soroban, Operations only) as event-free instead of erroring, so full-history backfill from genesis can read those ledgers. The SDK reference path (GetTransactionEvents / LCMToPayloads / db.ParseTransaction) rejects V0, so this is deliberately more permissive and documented as such; the two tests that asserted the V0 error now assert V0 success. - envPartFromView reads the envelope type from the decode it already performs for hashing, dropping a redundant view .Type() traversal. - Add a generic short-circuiting step()/viewChain helper and use it in readLedgerHeader and readTxHash to collapse the per-accessor error ladder.
The decode is not 'purely for the hash' — it also feeds the envelope type and the soroban flag (which must inspect Tx.Ext), so hashing piggybacks on a decode that is required regardless. Fix the now-stale comments that still described it as transient/hash-only.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 9c579585b6
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
|
Requirement: events must be written to the event store in ascending getEvents cursor order — the order the SQLite path serves ( The current emission (
For two txs each with a before-fee event, one op event, and a refund: From protocol 23 onward every transaction emits fee events at the Before/After stages, so essentially every modern ledger gets stored out of cursor order — and frozen that way into immutable cold chunks. The visible effects: a term query returns a different sequence than the same query against the SQL path, and cursor pagination over the stream skips or duplicates stage events at page boundaries inside a ledger. The fix is in the emitters: per ledger, emit all BeforeAllTxs events (the traversal already yields these in apply order, so a stable partition preserves it), then per tx in apply order its op events followed by its AfterTx events, then all AfterAllTxs events — in both Once emission is cursor-ordered, payload.go's reconstruction CAVEAT should be updated too: every On tests, two layers:
|
|
Proposal: dissolve the
Tests follow the code:
A side benefit: the |
|
Going one step further: |
|
Simplification, superseding my earlier "no committed artifacts on failure" thread: the defensive machinery around cold artifacts protects against a consumer that doesn't exist, and should come out. The design's model is that the orchestrator's completion record is the only source of truth for a chunk: completion is recorded after all artifacts land, a crash or failure before that point means the chunk is re-run from scratch, and nothing consumes cold artifacts by scanning directories — the index build takes explicit paths composed from the record. Under that model, disk under What follows:
One addition: write the model down in this package's |
|
Every ingester re-derives the ledger sequence from the view ( |
|
Follow-ups on
The overall shape — accumulate everything, sort once at Finalize — is right; these are just the leaks inside it. |
|
On the metrics: the required signals from #765 are all here and correct (per-ingester hot/cold, both aggregates with the exactly-once emit, per-tier buckets). But the sink interface as shaped can't support the follow-up that is its stated reason for existing. #765's metrics section defines the sink as letting "the same ingester code report to Prometheus (production) or CSV (the bench-command migration, follow-up task) interchangeably," and the follow-up's end state is the rpc-hack bench commands rewritten on these production ingesters with no duplication of ingestion logic. The rpc-hack collectors that migration must reproduce are per-stage — Concretely: add one method, e.g. |
…tructive opens Address the remaining PR #779 review threads: - All hot stores are chunk-bound (each accumulates one chunk before being frozen into cold artifacts), so make the binding explicit on the ledger and txhash hot stores too: their constructors now take a chunk.ID and expose ChunkID(), and RunHot validates every injected store's binding up front instead of only the events store's. - Validate the probed first ledger IS the chunk's first before the destructive cold constructors run, so a corrupted/misrouted pack or a wrong-range ChunkSource cannot truncate a previously finalized chunk on its way to drain's rejection. - Re-check ctx cancellation after the first-ledger probe: a sibling chunk worker's failure cancels gctx while this worker is blocked in the probe's I/O, and the replay wrapper would otherwise hand drain the cached ledger only after the constructors already truncated the existing artifacts. - Extract the pre-build validation into probeChunkSource (keeps runOneChunkCold under the funlen limit and gives the probe a single home). - Bump go-stellar-sdk to the post-merge commit of stellar/go-stellar-sdk#5949.
…tifact model, seq-through interfaces, stage metrics Address the second review wave on PR #779 (issue comments): - views.ExtractEvents now emits each ledger's payloads in ascending getEvents cursor order (BeforeAllTxs across the ledger, then per tx its op events followed by AfterTx, then AfterAllTxs) — write order is the cursor contract, and pre-23 the old emission stored essentially every modern ledger out of cursor order. payload.go's reconstruction caveat is replaced: every (txIdx, opIdx) group is now contiguous, so the per-event index is positional within its group. - events.LCMToPayloads is deleted (no production callers, and as the views differential oracle it shared the emitter's traversal — the ordering bug above is exactly the shared failure). The differential now runs against the real oracle: db InsertEvents + GetEvents over SQLite on stage-event fixtures, plus an emission-order invariant test. StageSentinels is now the single sentinel definition; db's InsertEvents imports it instead of carrying the mapping inline. - The defensive cold-artifact machinery is removed in favor of the documented completion-record model (doc.go): the orchestrator's completion record is the only source of truth, nothing consumes artifacts by directory scan, and a chunk attempt overwrites its paths freely. Gone: the first-ledger probe + peekedStream, the Finalize unpublish rollback, the txhash .bin tmp+rename (writes in place now; fsync and the reader's header-vs-size check stay), the stale-bin constructor removal, the orphan-pack removal, and the bucket-dir pre-validation. - HotIngester/ColdIngester.Ingest now take the driver-validated seq; ledgerSeqOf and every per-ingester malformed-header branch are gone. - The txhash ingesters consume SDK ExtractTxHashes directly (the views wrapper is deleted); the cold path appends truncated ColdEntry keys straight into the accumulator, and Finalize sorts with slices.SortFunc. - MetricSink gains IngestStage(dataType, tier, stage, d, items) emitted around extract / term-index / store-write / finalize, so the CSV bench sink migration can reproduce the rpc-hack per-stage collectors; PrometheusSink maps it to per-tier stage histograms with pre-resolved children.
|
All six addressed in d3c833a: Cursor ordering (#issuecomment-4694645399): Delete LCMToPayloads (#issuecomment-4694784966): gone — Artifact-model simplification (#issuecomment-4695072420): all of it came out — first-ledger probe +
txhashCold follow-ups (#issuecomment-4695233108): Per-stage sink granularity (#issuecomment-4695488161): added |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d3c833a2c4
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
…e-histogram help strings
|
Re-reviewed 1. Eliminate the
That empties the package, so it goes away — which also stops 2. golangci-lint is red — mechanical, blocking.
3. Stage-metric coverage gaps. The stage histograms partition the per-ingester
So a CSV built from the stages reconciles for 4. Minor — Not actionable: the build-matrix red is an infra flake (canceled stellar-core download; sibling matrix entries passed and it builds clean locally), and dependency-sanity-checker is just the SDK pin being a pseudo-version — #5949 merged and go.mod is already bumped to the post-merge commit, so it clears once the SDK cuts a release tag. |
Addresses PR #779 re-review (items 1-4). - Move ExtractEvents -> events.LCMViewToPayloads, next to StageSentinels, and delete the now-empty views package (txdetails.go + ExtractTxDetailsByHash / ExtractTransactions and their tests had no production callers; the future serving PR can call the SDK ingest read path directly). Preallocate the payload slice from the per-tx event counts; drop the unused //nolint:gosec. - txhashCold.Ingest: emit the extract stage AFTER the truncate-and-append loop so the cold stages (extract + finalize) partition the per-chunk ColdIngest total with no unexplained remainder. - eventsCold.ingestSeq: fold offsets.Append into the write stage and emit term_index/write for every ledger (incl. empty/V0), so each of the three cold stage histograms carries exactly one sample per ledger. - golangci-lint: shorten TestColdService_Success and TestExtractEvents_MatchesSQLite under funlen; wrap long lines in ingest_test.go and metrics.go.
|
Thanks — all four landed in 992363a. 1. 2. golangci-lint green. The unused 3. Stage-metric coverage closed.
4. On the not-actionable notes: agreed — build-matrix red is the canceled stellar-core download (builds clean locally), and dependency-sanity is just the SDK pseudo-version pending the post-#5949 release tag. |
@chowbao What's the reconstruction strategy for the case where there are multiple events emitted per operation? Today each event in an operation gets a distinct id, so the trailing component of the eventid (as returned by Without stored If the plan is to assign events a ledger-wide incremental id in getEvents response, that would differ from the id (and cursor) currently returned by getEvents. That seems like a breaking change unless there's some mechanism to preserve or reconstruct the existing event Ids. |
feature/full-history received the official squashed #756 (concurrent index; migrate off membitmaps), while this PR's base (fh-ingest-base) already carried an equivalent variant of the same work. Both branches resolve to byte-identical trees (c388061), so the apparent conflicts are purely topological. Recording the merge with -s ours keeps the PR tree unchanged while making feature/full-history an ancestor so the PR is mergeable.
There isn't a set plan right now for I think we can defer defining the event index for now because it's pretty easy to add back into the payload but we don't know exactly what format the index should take right now and we don't know if preserving or reconstructing the rpc v1 event ids is needed or not |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 486f518abd
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { | ||
| return nil, fmt.Errorf("mkdir %s: %w", filepath.Dir(path), err) | ||
| } | ||
| w, err := ledger.NewColdWriter(path, chunkID.FirstLedger(), ledger.ColdWriterOptions{}) |
There was a problem hiding this comment.
Reject overlapping pack sources before truncating ledgers
When RunCold is used to derive stores from an existing ledger-pack root such as NewPackSource(filepath.Join(coldDir, "ledgers")) and cfg.Ledgers is left enabled, this opens the ledger cold writer at the same packPath that packStream.RawLedgers will later read. ledger.NewColdWriter truncates that file immediately, so the subsequent source read fails after the original ledger pack has already been destroyed. Either reject overlapping pack-source/output paths or open and hold the source reader before constructing the ledger writer.
Useful? React with 👍 / 👎.
Agree with deferring the format decision but I'd suggest keeping EventIdx in the payload for now even if we end up not using it in the future. Adding it back later would require a payload version bump and reingestion. We can always ignore the field if a different cursor format is ultimately chosen. Removing EventIdx from the payload isn't fully backward compatible. It either requires reconstructing EventIdx when generating event IDs (behaviorally same but algorithmically different). Until we make a deliberate decision to change the cursor/ID format, I think we should preserve backward compatibility and keep the current ID format (-) since it's part of the public API today. |
Yeah that's fair. I'll add it back in |
Per review on #779: keep the per-event index stored in the 0x01 events payload rather than dropping it for read-time positional reconstruction. Removing it isn't fully backward compatible with the public getEvents <TOID>-<eventIdx> ID format, and re-adding it later would force a payload version bump plus reingestion of frozen cold chunks. The format decision is deferred — the field can be ignored if a different cursor scheme is chosen — but it costs only 4 bytes to preserve compatibility now. - payload.go: re-add the eventIdx slot (offset 53) to the wire layout, the Payload.EventIdx field, and its marshal/unmarshal; keep the exact-length check (now a generic wrong-width guard). - extract.go: LCMViewToPayloads populates EventIdx with the same per-group counter semantics as the SQLite path (db/event.go) — ledger-wide for BeforeAllTxs/AfterAllTxs, per-tx for AfterTx, per-op index for op events. - tests: round-trip covers EventIdx; the SQLite differential now asserts the stored EventIdx equals the SQL cursor's Event index directly; reworked the loud-failure test to reject the eventIdx-less layout.
|
@urvisavla added back event index here |
Implements #765 — the full-history data-ingestion path (cold + hot stores) for each
LedgerCloseMeta. Code atcmd/stellar-rpc/internal/fullhistory/ingest/.SDK dependency (stellar/go-stellar-sdk#5949)
The view extractors now live in the SDK as the zero-copy twins of the parsed path. Per the #5949 review, the SDK exports only the complete extractors (the navigation scaffolding is unexported):
ingest.ExtractTxHashes— per-ledger transaction hashes in apply order.ingest.ExtractLedgerEvents— per-transaction contract events + hash from oneTxProcessingwalk.ingest.LedgerTransactionViewByHash/LedgerTransactionViewRange(+ingest.LedgerTransactionView) — view parallel ofLedgerTransaction/LedgerTransactionReader.network.TransactionViewHasher— view twin ofnetwork.HashTransactionInEnvelope(returns just the hash).go.mod pins the #5949 branch pseudo-version (
f92b870f); bump to the merged version once #5949 lands.RPC-side view adapters (
internal/fullhistory/views/)The local
viewspackage collapses to thin adapters that add only RPC-specific shapes/policy:ExtractEventscomposesingest.ExtractLedgerEvents(hash + events from one walk) with theevents.Payloadshape and theStage→(TxIdx, OpIdx)cursor sentinels (events.StageSentinels).ExtractTxHasheswrapsingest.ExtractTxHashesintotxhash.Entry.ExtractTxDetailsByHash/ExtractTransactionsdelegate to the SDK read path;views.Transactionaliasesingest.LedgerTransactionView.dispatch.go+envelopes.godeleted (moved to the SDK).EventIdx removed from the events payload
The per-event index is positional and reconstructed at read time, so the
eventIdxslot is dropped from the0x01payload layout.unmarshalHeadernow requires the declared ContractEvent length to consume every remaining byte — a pre-removal record fails loudly (ErrPayloadLengthMismatch) rather than silently misparsing.LCMToPayloadskeeps only theStage→(TxIdx, OpIdx)sentinels (shared with the SQL path viaevents.StageSentinels); the ledger-wide before/after counters are gone.Ingestion path (#765)
HotIngester{ Ingest(ctx, xdr.LedgerCloseMetaView) },ColdIngester{ Ingest, Finalize, Close }.HotService(per-LCM parallel fan-out, waits-all) andColdService(sequential ingest +Finalize).MetricSink(NopSinkdefault) +PrometheusSinkunder the daemon namespace.ChunkSource— pack / GCS / S3 / DataStore.Tests
Per-ingester readback (real temp-dir stores) incl. V0-as-empty events;
HotServicefan-out + failure/sibling-cancel;ColdServicesuccess + failure-path-no-artifact; the SDK differential tests prove the view extractors wire-identical to the parsed path across LCM V0/V1/V2 and meta V0–V4.-raceclean.Test plan
go test ./cmd/stellar-rpc/internal/fullhistory/... ./cmd/stellar-rpc/internal/events/(incl.-race),go vet,gofmt— green against the #5949 SDK branch.