Skip to content

Ingestion: all data types for both cold and hot stores #765

Description

@tamirms

TL;DR

Implement the data-ingestion path for RPC v2: for each ledger close meta, write the relevant data into the configured storage tiers. Three data types (ledger / events / txhash) × two tiers (hot RocksDB / cold immutable files) = six per-tier ingesters, plus a top-level service per tier that fans an LCM out across the enabled ingesters with hot-tier parallelism, configurable data-type selection, and Prometheus metrics. A working prototype of all of this exists in the rpc-hack bench harness and can serve as inspiration (see "Reference / prior work" below).

What to build

Code lands in cmd/stellar-rpc/internal/fullhistory/ingest/.

Per-ingester contract

Two interfaces, one per tier:

type HotIngester interface {
    Ingest(ctx context.Context, lcm xdr.LedgerCloseMetaView) error
}

type ColdIngester interface {
    Ingest(ctx context.Context, lcm xdr.LedgerCloseMetaView) error
    Finalize(ctx context.Context) error
    Close() error
}

Ownership model. Hot ingesters receive a long-lived store via dependency injection (the daemon owns the store's lifecycle). Cold ingesters open their own per-chunk writer in the constructor and own it for the chunk's duration — Finalize commits on success, Close drops a partial output file on failure.

Six implementations:

Hot (RocksDB) Cold (immutable files)
Ledger ledger.HotStore.AddLedgers ledger.ColdWriter.AppendLedger + Commit
Events eventstore.HotStore.IngestLedgerEvents eventstore.ColdWriter + WriteColdIndex
TxHash txhash.HotStore.AddEntries append (txHash, ledgerSeq) per ledger → sorted dump per chunk on Finalize

The global txhash MPHF build over the per-chunk sorted dumps is a separate downstream step, not in scope here.

Metrics

Metrics are emitted through a small sink interface so the same ingester code can report to Prometheus (production) or CSV (the bench-command migration, follow-up task) interchangeably. The Prometheus implementation ships with this task and registers via the existing daemon convention (internal/daemon/metrics.go). Exact sink shape is up to the implementor.

Required signals:

  • Per-ingester — hot: per-ledger ingest latency, throughput, item count, error count. Cold: per-chunk ingest latency, item count, error count.
  • Aggregate (across all enabled data types) — emitted by the top-level service for each tier. Hot: per-ledger total wall-clock time (≈ max of the parallel per-ingester times). Cold: per-chunk total wall-clock time (sequential ingests across all ledgers + all Finalizes).

Per-LCM fan-out

For each tier, a top-level service holds the slice of enabled ingesters and dispatches each incoming LCM:

  • Hot tier — runs each ingester's Ingest concurrently within the same ledger and waits for all to complete before advancing. Distinct ingesters write disjoint state (different RocksDB DBs / column families), so per-LCM concurrency is safe and recovers wall-time.
  • Cold tier — runs ingesters sequentially per ledger, then on success calls Finalize on each at end of chunk. On the failure path, deferred Close drops partial files (each store's writer removes its in-flight output when Finalize never ran).

Configurable data types

The set of enabled data types is config-driven — a deployment can run ledger-only, events-only, or all three. At startup, the top-level service for each tier is constructed with the slice of ingesters matching the enabled types.

Dependencies

Tests

  • Per-ingester unit tests with mock stores.
  • Top-level-service tests: enabled-subset configurations, hot-tier per-LCM parallel fan-out, cold-tier sequential + Finalize success/failure paths.
  • Metrics emission verified via the sink (Prometheus impl).

Acceptance

  • Six ingesters implemented per the table.
  • HotIngester / ColdIngester interfaces published.
  • Metrics sink interface published + a Prometheus implementation registered via the daemon convention.
  • Required metrics emitted: per-ingester (hot per-ledger, cold per-chunk) and aggregate (hot per-ledger total wall-clock, cold per-chunk total wall-clock).
  • Top-level services for hot + cold tiers with config-driven data-type selection.
  • Hot-tier per-LCM parallel fan-out.
  • Tests present.

Reference / prior work

A working prototype of most of this exists in the rpc-hack bench harness and can serve as inspiration:

  • interfaces — cmd/stellar-rpc/scripts/bench-fullhistory/ingester.go.
  • ledger ingesters — ingest_ledgers.go.
  • events ingesters — ingest_events.go.
  • txhash ingesters — ingest_txhash.go.

These files mix the ingestion logic with bench-specific CSV Collector instrumentation; the production version routes metrics through the sink interface (Prometheus impl now, CSV impl in the follow-up task). The bench commands themselves stay on rpc-hack for now.

Follow-up

Bench command migration: rewrite the rpc-hack bench commands to use these production ingesters with a CSV sink implementation. To be filed separately, depends on this issue. End state: no duplication of ingestion logic between production and bench.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Type

No type

Projects

Relationships

None yet

Development

No branches or pull requests

Issue actions