Skip to content

Streaming daemon: Phase 2 — hot store, lifecycle + live ingestion (closes #816, #808) - #820

Merged
chowbao merged 59 commits into
feature/full-historyfrom
streaming-phase2-lifecycle
Jul 6, 2026
Merged

chowbao merged 59 commits into
feature/full-historyfrom
streaming-phase2-lifecycle

Conversation

@chowbao

@chowbao chowbao commented Jun 24, 2026 •

Copy link
Copy Markdown
Contributor

Closes #816, #808 — the complete Phase 2 (hot tier + lifecycle + live ingestion) in one PR.

Update: this PR now contains all of Phase 2. The live-ingestion layer originally split out as #821 (layer 2) was folded in here, so #820 alone completes #816 (and #808) rather than the two-PR split. #821 is closed in favor of this PR. Rebased onto the merged Phase 1 (#819) on feature/full-history.

Phase 2 introduces the hot tier + lifecycle machinery + live ingestion on top of Phase 1's source-blind cold pipeline:

  • Hot storage — the per-chunk hotchunk DB

    • single shared multi-CF RocksDB: ledgers + events + tx-hash
    • transient/ready state machine; a read-only view serves both the freeze source and the watermark refiner
    • ingest.HotService commits each ledger as one atomic synced WriteBatch across all CFs (decision (a)) — one fsync, a single per-chunk MaxCommittedSeq, no per-store frontiers, no fan-out
    • one ExtractLedgerEvents walk per ledger feeds both the tx-hash and events CFs (event-ID assignment order unchanged)
  • Backfill integration — hot source by path (no probe seam)

    • backfillSource's hot branch opens the chunk's hot DB read-only straight from its geometry.Layout path (hotchunk.OpenReadOnly) and yields a ledgerbackend.LedgerStream
    • a ready chunk whose DB is missing/gutted fails the must-exist open — an ordinary restartable error, never auto-healed into a fresh empty DB (no watermark regression)
  • Progress — derived watermark

    • LastCommittedLedger(cat, logger) maxes the cold term (highest fully-durable chunk) vs the highest ready hot DB's MaxCommittedSeq (one read-only open, which replays any synced WAL after an ungraceful crash) vs the earliest-pin floor — all in the signed domain, never stored; restart re-derives from durable state
  • Lifecycle + live ingestion (the folded-in layer 2)

    • run() transitions from backfill to a serve+ingest steady state: start captive core (injected CoreOpener) → serve reads (injected) → run the ingestion loop and the lifecycle loop as a joined errgroup.WithContext pair (whichever returns first tears down the other; both joined before run returns — the single-lifecycle-goroutine invariant across supervisor restarts)
    • the ingestion loop owns the hot tier: it opens the resume chunk's hot DB itself, consumes one captive-core RawLedgers stream, commits each ledger, and at each chunk boundary closes the filled DB → opens the next → publishes the completed chunk (the handoff fence)
    • the lifecycle loop (lifecycle.Loop): freeze → index-aware discard → prune, driven by a latest-cell BoundarySignal — a slow lifecycle can never fall behind, since one tick over [floor, latest] subsumes every skipped boundary
    • the freeze stage reuses backfill.RunBackfill — the same path catch-up uses
    • supervise is the single clean-vs-restart decision point (a canceled ctx is a clean shutdown; anything else is a warn + backoff restart). There is no fatal-and-exit class: genuine volume loss presents as a supervised crash-loop, upheld by the must-exist hot-DB open rather than by a hard exit
  • Daemon wiring

  • Observability

    • hot metrics are batch-scoped: one HotLedgerTotal(duration, err) per ledger + per-type HotItems volume + per-phase timings (extract / ledgers / txhash / events / commit)
    • cold ColdIngest is emitted only on a terminal step (a Finalize, or an Ingest error) — never from Close — so a rolled-back or sibling-abandoned ingester leaves no phantom-success sample
  • Folded-in cleanup

    • deletes the now-dead RunHot/HotStores stream-drain orchestration and the HotProbe/HotChunk/ErrHotVolumeLost probe machinery (verified zero production callers)

Verification: go build + go vet + go test (incl. the full-lifecycle E2E: ingest → freeze → fold → discard → cold+hot lookup → restart → prune) green on ./cmd/stellar-rpc/internal/fullhistory/... (cgo RocksDB toolchain). golangci-lint runs in CI.

Follow-ups: #772 (captive-core config unification + read-serving cutover), #835 (LastCommittedLedger(cat) signature cleanup), #836 (halve cold-path extraction / thin the ColdIngester seam).

@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from 3d12c9e to 84ff8c2 Compare June 24, 2026 20:39
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch 2 times, most recently from df8ec80 to 17b5c39 Compare June 24, 2026 20:51
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from 84ff8c2 to 419f7ec Compare June 24, 2026 22:05
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch from 17b5c39 to 145c1cc Compare June 24, 2026 22:05
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from 04f9931 to 7f8e58f Compare June 25, 2026 04:09
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch from e15575a to bc56b0a Compare June 25, 2026 04:25
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from 7f8e58f to f3431cd Compare June 25, 2026 05:19
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch from bc56b0a to 440443b Compare June 25, 2026 12:24
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from fa0083f to cbc80ab Compare June 25, 2026 14:42
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch from 440443b to c4944fb Compare June 25, 2026 15:07
chowbao added a commit that referenced this pull request Jun 25, 2026
…hanges

Rebased onto the updated #820 and propagated #817's API changes into the Phase 2
live-ingestion/daemon layer:
- window -> tx-hash index rename + key prefix index: -> txhash_index:
  (TxHashIndexCoverage.Index, Catalog.txhashIndex), Catalog.Get/Has -> get/has,
  config sections regrouped (cfg.Retention/Layout/Storage/Ingestion), pins via
  PinLayout.
- daemon.go merge: kept #821's live-ingestion wiring (LifecycleConfig + Core) and
  deduped the HotProbe line (#821's Phase-2 wiring already set it, so #820's
  HotProbe fix is redundant here).
- removed the #819 cold-only catch-up E2E (TestRunDaemon_CatchUpMaterializes...)
  + its someTxBackend/oneTxLCMBytes helpers: #821's daemon now requires
  Boundaries.Core and runs a continuous live loop, so a cold-only "catch up then
  return" test can't fit — and TestE2E_DaemonLifecycle covers it end to end.

Mechanical propagation only; build/vet/test -short green (the heavy lifecycle E2E
stays -short-gated).
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch from cbc80ab to ae91d20 Compare June 25, 2026 19:14
@chowbao
chowbao force-pushed the streaming-phase2-lifecycle branch from c4944fb to aeca6a0 Compare June 25, 2026 20:02
chowbao added a commit that referenced this pull request Jun 25, 2026
Rebased the live-ingestion capstone onto the reorganized #820 and propagated:
- qualify moved symbols (geometry./catalog.) in daemon.go, startup.go, e2e_test.go
- window->tx-hash-index + RetentionGate->RetentionFloor renames; cat.layout->Layout(),
  cat.Has->public HotState shim, .IndexFilePath->.TxHashIndexFilePath
- config regroup: cfg.Streaming.CaptiveCoreConfig -> cfg.Ingestion.CaptiveCoreConfig
- restored #821's daemon_test.go (drops the cold-only catch-up test the full daemon
  supersedes; adds the supervise/backend-tip/boundaries tests) + the HotProbe/Core wiring
- avoided the txhash_txhash_index find-replace corruption (was only in the dropped restack)

build + vet + go test -short green EXCEPT the lifecycle E2E, whose generated TOML
still uses the pre-regroup [streaming]/[backfill] schema (follow-up; per maintainer
the stack will be re-rebased).
@chowbao
chowbao force-pushed the streaming-phase1-daemon branch 6 times, most recently from 14aa4c8 to aafbe0d Compare June 26, 2026 14:02
chowbao added a commit that referenced this pull request Jun 26, 2026
Relocate the one-write protocol ordering helper (mark -> create -> barrier ->
flip) from backfill/process.go to catalog_protocol.go, where the protocol's
states and mark/flip steps already live, and export it as catalog.OneWrite.
processChunk and buildTxhashIndex now call it across the package boundary.

It is a zero-dependency pure function and catalog never imports backfill, so
there is no import cycle; #820's hot-tier openHotTierForChunk adopts it as the
third caller by import alone, with no later relocation. Addresses the #818
review thread that asked to establish the shared helper here rather than
deferring the move to #820.
// sliding floor (the fixed earliest-ledger floor alone applies).
RetentionChunks uint32

// OpRetryAttempts / OpRetryBackoff bound the per-op retry the discard/prune

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

OpRetryAttempts / OpRetryBackoff have no production wiring: run() builds lifecycle.Config from ExecConfig + RetentionChunks only, and no TOML field reaches these, so production unconditionally runs the WithLifecycleDefaults constants — only tests can set them. Either plumb a config knob (the design's config schema carried the retry attempts/backoff) or demote them to constants; as they stand they advertise a configurability that doesn't exist.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sweeping for more of this class turned up five siblings:

  • backfill.ExecConfig.RetryBackoff — same shape exactly: MaxRetries is TOML-wired right next to it, but no production code sets RetryBackoff; only execute_test.go does. (Its default is also applied twice — WithDefaults and retryBackOff() — same constant today, nothing linking them.)
  • ExecConfig.Build.BuildOpts — nobody sets it, not even tests; cfg.BuildOpts... at txindex.go:101 always expands empty. Fully dead knob.
  • ColdWriterOptions{} — both cold ingesters pass the zero value, whose own doc says batch workloads should set non-zero — and the batch freeze/backfill path is the only production consumer, so a full-history backfill always runs serial zstd with no background writeback. The inline "driver-level tuning is a follow-up via Config" note has no issue home — worth folding into fullhistory: halve cold per-ledger extraction (one ExtractLedgerEvents walk) + thin the ColdIngester seam #836 or filing.
  • [backfill].workers defaults to GOMAXPROCS independently in two packages (config.go:168, execute.go:54) with nothing linking them; the second site is unreachable today but diverges silently if either changes.
  • waitForCoverage's documented zero-value fallback (backend.go:146-150) is unreachable — its only caller passes the same constants the callee falls back to.

Adjacent nit: logging.format accepts any string and silently means text unless it's exactly "json".

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All landed in 8217374, direction: demote rather than plumb — none of these have an operator asking for them yet, and TOML can be added when one does. (a) OpRetryAttempts/OpRetryBackoff → unexported opRetryAttempts/opRetryBackoff, documented as a test seam. (b) ExecConfig.RetryBackoff → unexported, and the double default is gone — retryBackOff() is now the single applier of defaultRetryBackoff (removed from WithDefaults). (c) BuildOpts deleted, field and expansion. (d) ColdWriterOptions{} tuning note now points at #836 explicitly (commented there too). (e) workers: both sites now call a shared backfill.DefaultWorkers() — one source, no silent divergence. (f) waitForCoverage's unreachable fallback deleted; params stay (tests pass explicit values). (g) logging.format now rejects unknown values at validation instead of silently meaning text.

}

// Batch is durable — now and only now apply the events mirror/offsets update.
applyEvents()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

applyEvents() runs outside all five phases, so the claim both doc sites make — "the phases sum to the per-ledger total" (the Phase doc and the hot_phase_duration_seconds help) — isn't quite true: the ledger's cost is extract + queues + commit + apply, and the apply share lands in no metric. It's usually small, but the mirror update clones each touched term's bitmap copy-on-write, and popular terms' bitmaps grow all chunk long, so the cost rises exactly where a stall would be hardest to diagnose — the phase histograms would all look fast while live cadence lags. The pre-unification LedgerPhases doc carried a "(minus the tiny post-commit mirror apply)" qualifier that the rework dropped while keeping the claim. Cheapest honest fix: a sixth PhaseApply stamped around this call — it only runs on success, so the emission contract stays clean. Restoring the qualifier in both doc sites is the fallback, at the cost of keeping the blind spot.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 8217374 with the sixth phase: PhaseApply stamped around the mirror apply, so the phases genuinely sum to the per-ledger total and a growing copy-on-write bitmap apply now shows up exactly where you'd look for a late-chunk stall. It's post-commit and success-only, so the failure-emission contract is untouched (Failed can never be PhaseApply). Enum, sink array, help text, and the emission tests all updated; TestPrometheusSink_Smoke and TestHotService_EmitsEveryPhaseOnSuccess now cover six phases.


// NewRetentionFloor pins the floor for one (through, retentionChunks, earliest)
// snapshot. A shortened retentionChunks raises the floor at once.
func NewRetentionFloor(through, retentionChunks, earliest uint32) RetentionFloor {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

NewRetentionFloor has zero production callers: the tick builds the gate via EffectiveRetentionFloor + RetentionFloorAt (compute once, share with both scans), and startup uses EffectiveRetentionFloor directly. Only tests use this constructor, and as a second construction path it can drift from the shape that superseded it. Tests can compose the two calls; suggest deleting it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Deleted in 8217374. Tests now compose RetentionFloorAt(EffectiveRetentionFloor(...)) at all call sites — the same shape the production tick builds, so there's no second construction path to drift.

next := closed + 1
// Handoff fence: close the write handle BEFORE the next chunk's key is
// created (that key is what makes THIS chunk complete to a tick, which may
// then freeze and discard its hot DB — no writer may hold it then).

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

rocksdb.Store.Close returns nil on every path — a flush failure is logged and swallowed by design (data stays durable in the WAL) — and hotchunk.DB.Close just forwards it. So this boundary close-error branch, and the deferred close's error propagation below, are unreachable. Related doc casualty: run()'s header says it "returns nil only on a clean shutdown", but run() never returns nil at all (ingestion's nil is converted to an error, every other path errors), and supervise classifies clean shutdown by ctx.Err(), so its err == nil arm is equally dead. Worth either making Close's contract honest or trimming the dead branches and fixing the two docs — as written, a reader budgets for error paths that cannot happen.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Trimmed in 8217374 (took the trim-the-branches option — Close's log-and-swallow is deliberate, the WAL keeps the data durable, so the contract to make honest was the callers'). The boundary close-error branch and the deferred close's error propagation are gone (_ = hotDB.Close() with a one-line note); the loop dropped its named error return. run()'s header no longer claims a nil return — it states that a clean shutdown surfaces as a ctx-canceled error classified by supervise, whose dead err == nil arm is also gone. The load-bearing pair (the stream-ended-unexpectedly error and the nil-to-error guard) stays, now pinned by tests from the coverage thread.


// CompleteThrough maps a signed chunk index to its "complete through" last ledger:
// c < 0 ⇒ PreGenesisLedger; c >= 0 ⇒ chunk.ID(c).LastLedger().
func CompleteThrough(c int64) uint32 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CompleteThrough is the design pseudocode's retired name — the doc rewrite split that concept into lastCompleteChunk (which LastCompleteChunkAt implements) plus a plain chunk→ledger conversion, chunkLastLedger. This function is the second half, so name it for what it does: ChunkLastLedger(c int64), the exact companion of ChunkFirstLedger below. "Complete through" describes one caller's reading of the result, not the conversion itself — and the call sites get clearer for it (lastCommitted != geometry.ChunkLastLedger(c)).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Renamed in 8217374: geometry.ChunkLastLedger(c), doc rewritten as the plain chunk→ledger conversion companion to ChunkFirstLedger; all call sites updated (lastCommitted != geometry.ChunkLastLedger(c) reads as intended now).

// run is the daemon's startup: backfill to the tip, then serve reads (injected).
// Returns nil only on clean shutdown; any other return is restartable
// (ErrFirstStartNoTip on a first start with no reachable backend).
// run is the daemon's startup, in two steps: (1) BACKFILL to the tip, then

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A stale-comment sweep of the final diff — after this many rewrite rounds, these now contradict the code they sit on:

  • startup.go:21-25 — the run() header still describes the pre-reorder choreography ("…the live ingestion loop (which opens the resume chunk's hot DB itself)"); the body now opens it before ServeReads and hands it in.
  • ingest/metrics.go:202-204 — the NOTE says "there is no full-history ingest daemon startup path yet"; this PR built it (buildSinks in daemon.go wires NewPrometheusSink into both tiers).
  • hotloop.go:203-204 — "both surfaced as the cursor's error element": the cursor is deleted; it's the stream's error element.
  • "watermark" (renamed project-wide to last committed ledger) survives at hotloop.go:35, catalog/catalog.go:153, lifecycle/progress.go:147, observability/observability.go:17, and pkg/stores/hotchunk/hotchunk.go:90-93 (twice).
  • backfill/backend.go:142 promises "a fatal 'backend tip query' error … (a broken backend is not retried)" and backfill/execute.go:236 says "an unproducible chunk fatals" — nothing fatals: the task retries under withRetries, then the plan cancels and supervise restarts.
  • geometry/paths.go:117 cites RunCold, which no longer exists (the cold entry point is WriteColdChunk); geometry/paths.go:14 and backfill/process.go:4 cite design-docs/full-history-streaming-workflow.md, which isn't in the tree.
  • "fold" at daemon.go:65 ("fold+prune") and lifecycle/lifecycle.go:122 ("freeze + index fold") — the index path rebuilds from scratch (Rebuild metric, buildThenSweep); fold is retired design language.
  • e2e_test.go:330,464 — "doorbell" is the deleted mechanism; it's BoundarySignal now.
  • Review-history narration that stops meaning anything once this merges: ingest/driver.go:17-19 ("Close no longer emits…"), ingest/events.go:175 ("matching the old…"), config.go:27-28, catalog/catalog.go:185-186, geometry/txhash_index.go:17-19 ("It was once…") — state the current contract, drop the history.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Full sweep landed in 8217374 — every listed site: run()'s header now describes the open-before-serve choreography; the metrics NOTE acknowledges buildSinks; hotloop's cursor prose is stream prose; all five "watermark" stragglers (plus the test identifiers flagged in the residue thread) now say last-committed; the false "fatal / not retried" claims at backend.go/execute.go now describe withRetries + pass-fail + supervised restart; paths.go cites WriteColdChunk and the citations to the nonexistent workflow doc are dropped (nothing under design-docs/ matches — didn't want to re-point at a guess); fold→rebuild at both sites; doorbell→BoundarySignal; and the review-history narration in driver.go/events.go/config.go/catalog.go/txhash_index.go states the current contract with the history dropped.

}

func (s *testSink) HotIngest(dataType string, _ time.Duration, items int, err error) {
func (s *testSink) HotPhase(phase hotchunk.Phase, _ time.Duration, items int, err error) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Six invariants from this review's fixes have no test pin — each of these regressions would pass CI today:

  1. Failed-phase narrowing: the only failure ever driven through HotService is a closed DB, where Failed keeps its PhaseCommit default — no test triggers an extract or queue failure and asserts the attribution, and Failed's zero value IS PhaseExtract, so a forgotten assignment on an extract path is structurally invisible.
  2. Partial durations on failure: this testSink.HotPhase discards the duration argument and nothing anywhere reads rep.Phases[p].Dur, so a straight revert of the partial-duration stamping passes.
  3. The last-committed gauge: no ingestion-loop test wires a metrics recorder (the per-ledger emission is unobserved), and no lifecycle test would notice the tick re-emitting a chunk-aligned value — the exact regression the gauge split fixed can return silently.
  4. The handoff fence: recordingBoundary records chunk ids only; nothing observes close-before-next-key or publish-after-open. A publisher fake that, inside Publish(closed), attempts hotchunk.OpenExisting(closed) (fails on the RocksDB LOCK if the writer still held it) and reads HotState(next) would pin both edges.
  5. runOps ctx-abort mid-backoff: the cancel test cancels before the first op, so the backoff.WithContext wiring is never exercised — dropping it would block shutdown for (attempts−1)×backoff per failing op and every test stays green.
  6. The nil-return pair: no test covers a stream that ends cleanly (hotloop's "ingestion stream ended unexpectedly") or run()'s nil-guard, so both could be deleted together and a graceful stream end would hang g.Wait with nothing red.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All six pinned in 8217374 (one is sharpened by the AddEntriesToBatch thread — txhash attribution is now unrepresentable by construction, so the reachable narrowing set is extract/events/commit): (1) TestHotService_ExtractFailureLandsOnExtractPhase + TestHotService_EventsQueueFailureLandsOnEventsPhase — the latter asserts a non-zero Failed, the discriminator you called out for the zero-value trap. (2) TestHotService_FailedPhaseCarriesPartialDuration — testSink now records the duration arg; failed and completed phases both assert Dur>0. (3) TestRunIngestionLoop_LastCommittedGaugeAdvancesPerLedger asserts the exact per-ledger sequence, and TestRunLifecycleTick_DoesNotReEmitLastCommitted pins the tick to RetentionFloor only. (4) TestRunIngestionLoop_HandoffFenceClosesBeforeNextKey — your suggested shape: the publisher fake re-opens the closed chunk read-write inside Publish (LOCK-fenced) and asserts HotState(next)==ready. (5) TestRunOps_CtxCancelDuringBackoffReturnsPromptly — 30s backoff, cancel mid-sleep, asserts prompt ctx return; a dropped backoff.WithContext hangs it visibly. (6) TestRunIngestionLoop_CleanStreamEndIsError + TestRun_IngestionCleanEndSurfacesErrorNotHang. One honest note on (6): the run()-level test pins the property (graceful end → error, no g.Wait hang) but can't reach the literal nil-guard line — it's unreachable through the real loop, which is exactly what (6a) enforces.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

applySharedTableOptions (rocksdb.go:674-690) installs one shared bloom filter across every CF's table options — but grocksdb's SetFilterPolicy is a documented move ("this op is move, fp is no longer usable"): it nils the filter's C pointer after installing it. So only the FIRST CF in the loop gets the filter; every later iteration passes NULL and silently installs no filter. The if s.filter != nil guard doesn't catch it because the Go wrapper stays non-nil after its C pointer is gone. With hotchunk's CF order, the 12-bits/key bloom lands only on ledgers — sequential 4-byte keys, where a bloom is useless — and never on txhash, the random-point-lookup CF that is the filter's entire justification (its doc: every false positive at no-compaction SST counts costs a disk seek). Silent and perf-only, so no test can catch it.

This also corrects #838's premise: the tuning isn't over-applied to all CFs — the bloom is mis-applied to exactly one, the wrong one. #838's per-CF options would fix this incidentally (each CF constructing its own filter), but the PR shouldn't ship a filterless txhash CF in the meantime: the one-line interim fix is constructing a fresh NewBloomFilter per CF in the loop.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 8217374 with your interim shape: a fresh NewBloomFilter per CF constructed inside the loop, so txhash actually gets the 12-bits/key filter. Verified the move in grocksdb v1.10.7 source (SetFilterPolicy does opts.cFp = fp.c; fp.c = nil) — which is also why the old s.filter != nil guard couldn't see it. The shared filter field is gone; ownership now rides with each CF's table options (see the teardown reply for who frees what). #838's per-CF options can subsume this cleanly.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Two more instances of the same default-CF/API-semantics class as the flush finding:

  • logOpenState (rocksdb.go:696-699) reads rocksdb.cur-size-active-mem-table, num-files-at-level0, and total-sst-files-size via DB-level GetProperty, which resolves against the default CF — always empty here. So the [ROCKSDB:OPEN] line's memtable/L0/SST figures are permanently ~zero and its stated diagnostic purpose (big WAL → replay, high L0 → pending compaction) can never fire; only the WAL figure is real. Fix: GetPropertyCF summed over s.cfHandles.
  • Teardown: the per-CF BlockBasedTableOptions created in applySharedTableOptions are never destroyed (grocksdb's Options.Destroy doesn't free them), and s.filter.Destroy() in Close is a no-op once the move has nil'd the pointer. A small C allocation leaks on every open — recurring, since the daemon opens a DB at every chunk boundary and every freeze.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Both fixed in 8217374. logOpenState now sums memtable/L0/SST figures per-CF via GetPropertyCF (WAL stays DB-level), so the diagnostic line can actually fire. Teardown: verified in grocksdb source that SetBlockBasedTableFactory copies the BBTO rep (caller keeps ownership) and Options.Destroy never frees it — the Store now retains the per-CF BBTOs and destroys them in Close and on the open-failure path; the no-op filter.Destroy() is deleted (post-move it was rocksdb_filterpolicy_destroy(nil)). The moved-in filter policy is freed via its owning BBTO. One residual noted so nobody re-hunts it: grocksdb's BBTO keeps a cFp field it never frees — a few bytes per CF per open, unreachable through the wrapper's API.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The applyEvents finding generalizes — holding every metric's help text to its literal claim against the actual timed window:

  • phase_duration_seconds{phase="freeze"}: the timer starts after resolve() (execute.go:243 vs 247), but the claim is "plan-and-execute". Resolve is range-proportional catalog I/O — two metastore reads per chunk over the whole range plus per-window coverage scans — so the steady-state tick, whose plan is empty, scans [floor, lastChunk] every tick while Freeze observes ~0. The "plan" half is invisible.
  • phase="backfill_pass": the pass's own definition includes sampling the network tip (startup.go:170-173), but the timer wraps only runBackfill — on a flaky bulk backend the retried tip call (up to attempts×interval of sleep) escapes. One-line fix now; full-history: frontfill-only deployment cannot bootstrap (no tip source on first start) #833's tipSampler will restructure this path anyway.
  • cold_chunk_duration_seconds: the type doc says "times from the first Ingest" but the timer starts at construction; the window includes drain's source-stream time (bulk download can dominate) which the help attributes to "cold ingesters' ingests plus their Finalizes"; and Close — part of the claimed lifetime — runs after the emit. The wide window looks intended (the bucket comment says so) — the doc text is what needs fixing.
  • pruned_artifacts_total: the help's "(below the retention floor)" is false for three of the four counted categories (transient index keys from any window, in-retention .bin demotions, redundant txhash keys in finalized windows) — noting this amends the earlier rename we settled in the pruned_ops_total thread; the name is right, the parenthetical isn't.
  • phase="rebuild": help says "one index rebuild's wall-clock" but the window wraps withRetries — up to MaxRetries+1 full attempts plus exponential sleeps in one sample.
  • Discard/Prune on failure: both emit only after runOps fully succeeds (lifecycle.go:142-145, 163-166), so a mid-sweep failure permanently loses counts for ops that already retired DBs or swept artifacts — they don't reappear in the next scan. The family's own convention is the opposite ("reported even on failure" at both freeze and rebuild).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All six in 8217374. Freeze: timer moved before resolve(), so "plan-and-execute" is now literally what's timed and the steady-state tick's range-proportional scan shows up. backfill_pass: timer now starts before the tip sample, so retried tip calls are in the window. cold_chunk: kept the wide window (it's intended) and fixed the doc to say construction-to-emit including source-stream drain. pruned_artifacts_total: parenthetical replaced with the four real categories. rebuild: help now states the sample spans all retry attempts plus backoff sleeps. Discard/Prune on failure: real fix, not a doc fix — runOps returns the completed-op count and the prune scan returns per-op artifact weights, so a mid-sweep failure meters the ops that actually retired DBs/swept artifacts before the error surfaces, matching the family's reported-even-on-failure convention.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AddEntriesToBatch is a loop of error-less Puts ending in return nil — it structurally cannot fail (CF errors are latched inside BatchWriter and surfaced by Store.Batch itself). That makes hotchunk's queue-tx-hashes failure branch dead and Failed == PhaseTxhash unrepresentable — unlike its two siblings, which have real error paths (zstd encode; the events facade's four returns). Either drop the error from this signature and the dead branch with it, or leave a one-word note that the error exists for signature symmetry. Related: this sharpens the test-gap thread's first item — one of the three narrowing branches isn't just unpinned, it can't fire.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dropped the error in 8217374 — signature is now func(...) with a doc stating it cannot fail (Put latches, Store.Batch surfaces), and hotchunk's dead queue-tx-hashes branch went with it, making Failed==PhaseTxhash unrepresentable by construction rather than merely unexercised. PhaseTxhash survives as a duration phase. The narrowing tests from the coverage thread pin the two attributions that remain reachable (extract, events).

}

// ChunkID returns the chunk this DB is bound to.
func (d *DB) ChunkID() chunk.ID { return d.chunkID }

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sweep residue in two small classes:

  • Test-only exports missing the seam doc their siblings carry: DB.Ledgers() and DB.ChunkID() (hotchunk.go:127-130) have zero production callers — Txhash()/Events() carry the explicit "Wire the v2 stores into the API handlers and support both v1 and v2 #772 read seam" doc, these two don't; same for eventstore.HotStore.ChunkID(). To be clear about history: the earlier accessor thread deleted the ledger/txhash facade accessors and kept eventstore's chunkID field as load-bearing — these are the accessors that ruling didn't cover. Also eventstore.HotStore.All's doc claims "used by the freeze loop" — the freeze path this PR built never calls it (it's test-only now), and catalog.AllArtifacts+NewArtifactSet are production-unreachable (the resolver builds sets per-kind; pre-existing, but worth a doc fix or a fullhistory/streaming: split the streaming package into purpose-named packages #824 line).
  • Retired vocabulary survives only in test identifiers: mustDeriveWatermark/wmBeforeRestart (e2e_test.go), seedWatermark/TestRunIngestionLoop_RestartResumesFromWatermark (hotloop_test.go), and two ...Folds... test names — production identifiers are clean apart from the already-flagged CompleteThrough. Bycatch: rocksdb/encoding.go's EncodeUint64/DecodeUint64 are dead tree-wide (pre-PR file).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Landed in 8217374. Seam docs: DB.Ledgers()/DB.ChunkID() and eventstore.HotStore.ChunkID() now carry the same #772 read-seam doc as their siblings; HotStore.All's doc no longer claims the freeze loop uses it (states it's the Reader full-scan, test-only until #772); catalog.AllArtifacts/NewArtifactSet are doc-marked as test-only seams. Retired vocabulary: mustDeriveWatermark→mustDeriveLastCommitted, wmBeforeRestart→lastCommittedBeforeRestart, seedWatermark→seedLastCommitted, TestRunIngestionLoop_RestartResumesFromWatermark→…FromLastCommitted, and both "Folds" test names renamed to rebuild/covers vocabulary. Bycatch: EncodeUint64/DecodeUint64 confirmed dead tree-wide and deleted.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One residue the identifier renames didn't reach: ~28 "watermark" mentions survive in test comments and assertion strings (e2e_test, startup_test, hotloop_test, and five more files) — including the e2e narration of the resume semantics, which is exactly where retired vocabulary re-seeds itself. Mechanical find-replace to "last committed ledger" / lastCommitted.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in fab9b92 — swept all residual watermark mentions from test comments and assertion strings, and renamed the deriveWatermark test helper (and its TestDeriveWatermark* cases) to deriveLastCommitted. Prose now reads "last committed ledger" / "last committed seq"; the "single watermark" phrasings became "committed frontier" (the term already used in progress.go) to avoid tautologies like "the last committed seq is the last committed seq". grep -ni watermark over the fullhistory tree is now 0.

chowbao added 2 commits July 6, 2026 03:09
- rocksdb: flush all CFs on close (FlushCFs), per-CF bloom filter
  (SetFilterPolicy is a move), per-CF property sums in the open log,
  destroy per-CF table options on close; drop dead EncodeUint64/DecodeUint64
- hot ingest: sixth PhaseApply phase around the mirror apply;
  AddEntriesToBatch cannot fail - drop its error and the dead txhash
  failure branch; remove the redundant closed-check in ledger iterate
- knobs: unexport test-seam retry knobs (lifecycle op retry, backfill
  RetryBackoff), delete dead BuildOpts, share one DefaultWorkers source,
  drop waitForCoverage's unreachable fallback, validate logging.format,
  delete NewRetentionFloor
- lifecycle: meter Discard/Prune counts even when a sweep fails mid-way
- metrics: align each help text with the actual timed window (freeze
  plan+execute, backfill_pass incl. tip sample, cold chunk lifetime,
  pruned categories, rebuild retries)
- rename geometry.CompleteThrough -> ChunkLastLedger; trim unreachable
  hot-DB close-error branches; stale-comment sweep (watermark -> last
  committed, fold -> rebuild, doorbell -> BoundarySignal, dead citations)
- tests: pin failed-phase attribution, partial phase durations,
  last-committed gauge ownership, boundary handoff fence, runOps
  ctx-abort mid-backoff, clean-stream-end error
- geometry/paths.go: name the cold entry point (backfill.WriteColdChunk)
  in the TxHashRawRoot rationale
- backfill/backend.go: replace the last 'fatal / don't retry' wording with
  the actual poll-abort -> pass failure -> supervised restart flow
@tamirms

tamirms commented Jul 6, 2026

Copy link
Copy Markdown
Contributor

Final round from the review — six small items worth landing in this PR. The first three are in pre-existing files, same rationale as the flush/bloom fixes: this PR's daemon is what makes them bite.

  1. pkg/stores/txhash/cold_merge.go:193-200 — openDirect is the tree's only raw syscall.Open and omits O_CLOEXEC. A captive-core child (re)started during a multi-minute merge inherits every .bin fd; the post-build sweep then unlinks files the long-lived child keeps pinned — tens of GB of deleted data held until core exits. Fix shape: route it through os.OpenFile(path, os.O_RDONLY|directOpenFlag(), 0) — the stdlib sets O_CLOEXEC unconditionally and returns the *os.File directly, so the manual fd-wrapping goes away and forgetting the flag becomes impossible (nonstandard bits pass through to openat; the EINVAL fallback keeps its shape).
  2. pkg/stores/txhash/cold_index.go:168-169 — scanBinHeader multiplies an untrusted header count with no bound check. Since 20·2⁶² ≡ 0 mod 2⁶⁴, one flipped high bit passes the size check and feeds ~2⁶² into streamhash.NewSortedBuilder — an OOM, and since resolve re-emits the same build every restart, a crash loop until the file is repaired. Fix shape: invert the arithmetic — validate (size-header) % entrySize == 0 and require the header count to equal (size-header)/entrySize; division of the trusted size can't overflow. ReadColdBin does the same validation with different code, so extract one shared header-check helper both call. The adjacent size < 0 branch is structurally dead and can go with it.
  3. pkg/stores/txhash/cold_index.go:30-37 — the .bin entry layout is defined twice in the package. The writer derives ColdKeySize = streamhash.MinKeySize; the index/merge side hardcodes binKeySize = 16 — they agree only by coincidence, and a streamhash bump would wedge every subsequent index build. Fix shape: delete the bin* constant family and the duplicated format doc; cold_index/cold_merge use the writer's ColdKeySize/coldBin* constants directly (same package), restoring the single owner cold_bin.go already claims to be.
  4. geometry/paths.go:114 — "backfill.WriteColdChunk" → it's ingest.WriteColdChunk (backend.go got it right in the same commit).
  5. hotloop.go:201-203 — the nil-return rationale still says a nil "supervise would read as a clean shutdown and silently stop ingesting"; since this round's supervise change, classification is ctx-only and run()'s guard converts a nil anyway. Simplest fix: drop the speculation — the comment only needs "a source that stops without an error while ctx is live is abnormal; surface a restartable error." The real what-if analysis already lives on run()'s guard.
  6. pkg/stores/ledger/hot_store_test.go:234-241 — the rewrite dropped the only assertion that a closed store surfaces stores.ErrStoreClosed through the iterate path (the new NoError is right for the inverted-range case it now exercises). Two lines: iterate a valid range after Close, require.ErrorIs(err, stores.ErrStoreClosed).


// The one snapshot every stage shares. earliest and the retention gate are read
// and computed ONCE here (not re-derived per scan), then passed to both scans.
through := lastChunk.LastLedger()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

through is retired design vocabulary that outlived its source — it was named from CompleteThrough(...), which this PR just renamed to ChunkLastLedger, leaving a bare preposition. Two fixes, neither a rename-in-place: in LastCommittedLedger, through is the value the function returns, so call it lastCommitted; in the tick, it exists only to compare chunk completeness in the ledger domain — c.LastLedger() <= through is c <= lastChunk (LastLedger is monotonic), so pass lastChunk down, compare in the chunk domain, and convert at the single site that needs a ledger (EffectiveRetentionFloor). The log field renames with it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in fab9b92. In LastCommittedLedger, through (the returned value) → lastCommitted. In the tick, through is gone entirely: eligibleDiscardOps now takes lastChunk and compares in the chunk domain (c <= lastChunk, via LastLedger monotonicity), and the only ledger-domain conversion is EffectiveRetentionFloor(lastChunk.LastLedger(), …). The debug log field renamed to last_chunk.

Cold txhash pipeline:
- openDirect goes through os.OpenFile so .bin fds get O_CLOEXEC and a
  captive-core child spawned mid-merge can't pin unlinked files.
- scanBinHeader and ReadColdBin share one overflow-safe header check
  (coldBinCount) that divides the trusted file size instead of
  multiplying the untrusted header count; drops the dead size<0 branch.
- Delete the duplicated bin* constant/format-doc family; cold_index and
  cold_merge use cold_bin.go's ColdKeySize/coldBin* directly.

Comments/vocab:
- geometry/paths.go: backfill.WriteColdChunk -> ingest.WriteColdChunk.
- hotloop.go: drop stale supervise speculation on the fall-through error.
- lifecycle: rename the retired "through" vocabulary — lastCommitted in
  LastCommittedLedger, and the tick/discard scan compares in the chunk
  domain against lastChunk (converting to a ledger only at
  EffectiveRetentionFloor).
- Sweep residual "watermark" mentions from test prose and the
  deriveWatermark test helper.

Tests:
- ledger hot_store post-close ErrStoreClosed iterate assertion already
  present (no change needed).
- Add TestBuildColdIndex_HeaderOverflowRejected for the count overflow.
@chowbao

chowbao commented Jul 6, 2026

Copy link
Copy Markdown
Contributor Author

All six addressed in fab9b92:

  1. openDirect O_CLOEXEC — now routes through os.OpenFile(path, os.O_RDONLY|directOpenFlag(), 0) (EINVAL fallback kept), so the fd gets O_CLOEXEC unconditionally and the manual os.NewFile fd-wrapping is gone.
  2. scanBinHeader overflow — scanBinHeader and ReadColdBin now share one coldBinCount helper that divides the trusted size (body % entrySize == 0 and count == body/entrySize) instead of multiplying the untrusted count. The dead size < 0 branch is gone, and TestBuildColdIndex_HeaderOverflowRejected covers the MaxUint64 count.
  3. Duplicated .bin layout — deleted the bin* constant family and the duplicated format doc; cold_index/cold_merge use cold_bin.go's ColdKeySize/coldBin* directly (added a named coldBinSeqSize so the entry width is no longer a magic +4).
  4. geometry/paths.go:114 — backfill.WriteColdChunk → ingest.WriteColdChunk.
  5. hotloop.go — dropped the stale supervise speculation; the comment now just states that a source stopping without an error while ctx is live is abnormal and surfaces a restartable error, pointing at run()'s guard for the clean-vs-restart classification.
  6. hot_store_test.go iterate assertion — this one needed no change: the valid-range-after-Close require.ErrorIs(..., stores.ErrStoreClosed) iterate assertion is still present at lines 229–233 (added in 6e773e2, retained through the round-5 rewrite, which only added the inverted-range NoError case at 237–243). So the "only assertion" was not dropped — both cases are asserted.

@chowbao

chowbao commented Jul 6, 2026 •

Copy link
Copy Markdown
Contributor Author

PR #820 — review summary

Phase 2 of the streaming full-history daemon: live captive-core ingestion + the hot→cold freeze/discard/prune lifecycle (closes #816, #808). It went through ~5 rounds of review (primarily @tamirms). Concise recap of the back-and-forth and every change that landed.

Design alignment (first pass — design comments #14–#39)

Round 3 — polish: dead code, metric-name fix, stale docs; added DiscardHotChunk crash-resume + absent-key-noop tests.

Round 4 — lifecycle tick cleanups + gauge correctness (don't regress last-committed from the chunk-aligned value); deleted dead seams; pinned missing tests; golangci fixes.

Round 5 — unified hot metrics into one phase-keyed family; deleted SeqValidatedCursor (enforce in-order at the source); partial-duration metric on a failed hot phase + restored lifecycle op-retry; open the resume hot DB before serving reads.

Final round

  • openDirect via os.OpenFile so .bin fds get O_CLOEXEC (a captive-core child can't pin unlinked files mid-merge).
  • Overflow-safe shared .bin header check (coldBinCount divides the trusted size, not the untrusted count) + dedup of the duplicated bin* constants/format doc.
  • Retired the through/watermark vocabulary: chunk-domain comparison against lastChunk, lastCommitted naming.
  • Comment fixes (paths.go, hotloop.go); golangci follow-ups (gofumpt / unused-nolint / testifylint).

Deferred (tracked, no behavior change): LastCommittedLedger signature cleanup → #835; earlier rounds deferred smaller cleanups to their own follow-up issues.

Status: vet + package tests + golangci-lint (--new-from-rev, v2.11.3/go1.26) pass locally after f17378a1; CI re-running.

@tamirms

tamirms commented Jul 6, 2026

Copy link
Copy Markdown
Contributor

Two last things before this wraps:

  1. golangci-lint broke on fab9b92c (green at d67bf79e): cold_merge.go:384 gofumpt formatting; cold_bin.go:130 unused //nolint:gosec directive; lifecycle_helpers_test.go:189 testifylint require-error. The summary's tree check ran gofmt — CI's linter is gofumpt, which is stricter.

  2. One review thread is still open: the LastCommittedLedger cleanup (Streaming daemon: Phase 2 — hot store, lifecycle + live ingestion (closes #816, #808) #820 (comment)) — drop the logger parameter via a logger on Catalog, and move progress.go next to its only caller in the fullhistory root (nothing in lifecycle uses it since the tick stopped calling it). Fine to land here or to explicitly keep it deferred on fullhistory: simplify LastCommittedLedger to the design signature LastCommittedLedger(cat) #835 — it just needs one or the other stated.

- cold_bin.go: drop the now-unused //nolint:gosec on coldBinCount's
  size division (gosec doesn't flag it).
- lifecycle_helpers_test.go: assert.NoError -> require.NoError in
  assertQuiescent (testifylint require-error).
- cold_merge.go: gofumpt blank lines between the streamReader accessors
  that gofmt wrapped when the constant rename lengthened k0().
@chowbao

chowbao commented Jul 6, 2026

Copy link
Copy Markdown
Contributor Author

Both handled:

  1. Lint fixed in f17378a — dropped the unused //nolint:gosec on coldBinCount's size division (cold_bin.go), assert.NoError → require.NoError in assertQuiescent (lifecycle_helpers_test.go, testifylint require-error), and the gofumpt blank lines between the streamReader accessors that gofmt wrapped when the constant rename lengthened k0() (cold_merge.go). Verified 0 issues locally with golangci-lint v2.11.3 built on go1.26 via --new-from-rev=d67bf79e (my earlier tree check used gofmt, not the stricter gofumpt — thanks for the catch).
  2. LastCommittedLedger cleanup → deferred to fullhistory: simplify LastCommittedLedger to the design signature LastCommittedLedger(cat) #835 (stated on the thread). fullhistory: simplify LastCommittedLedger to the design signature LastCommittedLedger(cat) #835 already captures it exactly — drop the logger param for the design signature LastCommittedLedger(cat), convert the nil-logger positional tests to real-DB, and relocate progress.go. It's pure cleanup blocked on the coupled test rework, so it stays out of this PR.

@tamirms tamirms left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎉

@chowbao
chowbao merged commit e45425d into feature/full-history Jul 6, 2026
8 of 15 checks passed
@chowbao
chowbao deleted the streaming-phase2-lifecycle branch July 6, 2026 18:27
@github-project-automation github-project-automation Bot moved this from Needs Review to Done in Platform Scrum Jul 6, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

Status: Done

Development

Successfully merging this pull request may close these issues.

3 participants