jetstream v2 in zig stream.waow.tech
stream docs upstream-compaction-spec.md
15 kB
Markdown

upstream compaction + tombstone spec (condensed from agent study 2026-07-09) #

how jetstream physically removes deleted/updated data from sealed segments.

Everything below describes upstream, and every path in this section is an upstream Go path — none of them exist in this repository. That matters most for internal/tombstone/, which looks exactly like Stream's own src/internal/… convention and is not ours.

  • upstream sources: internal/ingest/orchestrator/compact_deletes.go, compaction_watermark.go, internal/tombstone/, segment/rewrite.go, and §3.3 of upstream's docs/README.md (upstream's, not this docs/ directory — there is no docs/README.md here).
  • Stream's counterparts, for reading the two side by side: storage/tombstone.zig (which rows may be dropped), storage/segment_rewrite.zig (the physical rewrite), compact/watermark.zig (the durable watermark), and compact/pass.zig / compact/steady.zig.

Two references in the original source list are not committed in this repository: the 2026-07-09 agent study named in the title, and the "2026-06-12 cache-invalidation design note". This condensation plus upstream source at the pin is the available authority. upstream-bootstrap-spec.md carries the same caveat.

pebble below is upstream's embedded key-value store (Go); Stream uses RocksDB in the same role.

tombstones #

BOTH record-level and account-level:

  • record tombstone: key (did, collection, rkey) → max superseding seq. produced by KindDelete and KindUpdate rows (update = delete+replace; only latest record version is kept).
  • did tombstone: key did → {seq, reason} where reason ∈ {"sync", "account"}. produced by KindSync rows (repo diverged, full resync follows) and by KindAccount rows whose payload decodes to active=false && status=="deleted".

semantics (strict inequality everywhere):

  • suppression drops only materialization rows (KindCreate, KindUpdate, KindCreateResync) with seq STRICTLY LESS than the tombstone seq (ShouldDrop: ts.Seq > ev.Seq). the marker row itself is never dropped.
  • tombstone rows are retained as event rows FOREVER: delete/account/identity/ sync rows survive compaction so mid-stream clients and audit tooling see history. compaction removes only superseded creates/updates.
  • there is NO read-time overlay: the server never suppresses at delivery; clients fold the stream themselves (creates/updates apply, deletes/ account-deletes/syncs remove).

merge semantics: Snapshot.Merge takes per-key max seq. Fold/FoldRange build a snapshot from decoded events within a (lowExclusive, highInclusive] seq window. inserted map keys are cloned so they don't pin the decoded block buffer (decoded events alias the decompressed block).

persistence: tombstones are NOT persisted as a separate structure — they ARE the event rows in segments. the in-memory tombstone.Set (mutex-guarded maps

  • approx-bytes gauge, 64-byte per-entry overhead constant) is rebuilt at startup by folding all rows above the watermark (rebuildLiveTombstones; cost scales with the watermark backlog, not the archive — segments/blocks with MaxSeq <= watermark are skipped). the live set is fed by the writer's OnAppend hook under the writer mutex, and is used ONLY for the cap trigger (Len) and size gauge — never for the pass's drop decisions.

the compaction watermark (compaction/seq in pebble) #

the highest seq physical compaction has covered. invariant: below W, superseded create/update rows are physically gone; the uncompacted tail (W, tip] may still carry rows a later marker kills. versioned uint64-LE value under key compaction/seq (compactionWatermarkV1 = 0x01).

initCompactionWatermarkFloor(store, nextSeq): called once when merge opens the destination writer. if a watermark already exists, no-op; else it writes nextSeq - 1 (or 0 when nextSeq == 0). this floors W at "everything before this instance's first real seq", so the first pass doesn't think the whole seq space below is uncompacted.

boundary-seq rules:

  • target watermark of a pass = max MaxSeq across SEALED segments with events. active segments are skipped; W never advances past them (their tombstones stay in the live set — Evict only removes entries ≤ committed W — and re-apply after seal).
  • pass no-ops when targetWatermark <= watermark.
  • eviction: after committing a chunk watermark, Tombstones.Evict(chunkEnd) removes live-set entries with seq ≤ chunkEnd.

the pass algorithm (runDeleteCompaction, one pass, chunked) #

  1. cleanup: delete stale *.jss.tmp in segments/ (also done at process boot). must hold the rewrite lock — a live timestamp-import rewrite may have a tmp open; under the lock any tmp seen is genuinely stale.
  2. steady mode only: force-rotate the live writer's active segment, so rows deleted while their segment was active get compacted this pass. race-free because tombstone Observe runs as the OnAppend hook under the writer mutex: by ForceRotate's return every sealed event's tombstone is in the live set, and later events land above the sealed MaxSeq (above this pass's target). merge-tail passes hand liveWriter=nil (tree already sealed).
  3. sweep segments/: open each file's header (SkipChecksum: true), skip active ones, compute targetWatermark. steady mode then reconciles the manifest (see below). exit if nothing to do.
  4. chunk loop, current = W until targetWatermark: a. build the chunk's tombstone snapshot by FOLDING SEALED ROWS FROM DISK in (current, targetWatermark] — deliberately NOT the in-memory set, which collapses each key to the GLOBAL max seq; a key updated again above the watermark would vanish from an upper-bounded readout and an older superseded row would survive forever below the committed W (issue #100 dual). folding the on-disk window can only see seqs ≤ target. block-index seq bounds prune blocks (bounds survive rewrites as historical supersets, so pruning is safe). b. cap: CompactionTombstoneCap (default 32_000_000 entries) bounds snapshot size; when hit, chunkEnd = that segment's MaxSeq and the loop does multiple chunks. c. rewrite chunk under the rewrite lock (scoped per chunk, not per pass, so a timestamp import may interleave between chunks — safe because both rewrites are per-segment atomic + idempotent): fan segments to min(NumCPU, 8) workers; each calls segment.Rewrite(path, decide) with decide = keep if seq > chunkEnd, else drop if snap.ShouldDrop. a bloom prefilter narrows by candidate DIDs unless the snapshot has more than CompactionBloomNarrowMaxDIDs (default 100_000) distinct DIDs (merge-tail passes run millions — probing costs more than it saves), in which case exact-scan every segment. IMPORTANT: after g.Wait() the ctx error is folded in — a cancel can let workers return nil for a chunk cut short, and committing that watermark would evict tombstones for rewrites that never ran. d. steady mode: for each actually-rewritten segment, fire OnSegmentCompacted(idx, path) (manifest refresh + cache invalidation, below). refresh failure fails the pass; W does not advance. e. commit: saveCompactionWatermark(chunkEnd), then Evict(chunkEnd). crashpoints straddle both sides (AfterCompactionRewriteBeforeWatermark, AfterCompactionChunkWatermark).

failure policy: a failed steady pass logs and retries at the next scheduled pass — it must NEVER tear down the daemon (would turn a transient IO error into an ingestion outage). watermark_lag_seconds (tip witnessed-at minus watermark witnessed-at, floored to the span of uncompacted data so an empty watermark doesn't read as a 50-year spike) is the operator paging signal.

segment rewrite crash safety (segment.Rewrite) #

  • rewrite is read-old → write sibling path + ".tmp" (.jss.tmp) → fsync tmp → close → rename over path → fsync parent dir. failure at or before the rename leaves the original untouched; a crashed tmp is reclaimed at boot and pass start. any write/fsync/rename error is crash-loud (disk-full gets an actionable operator message).
  • block topology is PRESERVED: a fully-dropped block remains as an event_count=0 block, and historical seq/witnessed-at envelopes are kept, so block numbers, cursor translation, and cold replay stay stable across generations. rewritten files get a new checksum.
  • if the segment-level bloom rules out all candidate DIDs, or no row drops, the rewrite is a no-op (Rewritten: false — "clean" segment).
  • the orchestrator-level rewriteMu serializes delete-compaction and timestamp-import rewrites; two concurrent tmp+rename on one file would silently drop the loser's changes.

merge-tail vs steady-state #

merge-tail (compactionMergeTail), once, at the tail of the merge phase: after the dst writer is sealed (SealActiveAndClose) and before phase=steady_state is written — so the FIRST served archive view is delete/update compliant. differences from steady: liveWriter=nil (no force-rotate), no per-segment manifest refresh (the pass is manifest-oblivious — nothing is serving yet), typically millions of tombstones so the bloom narrowing is skipped. immediately after it, reconcileCompactionManifestFromDisk runs one-shot: re-list sealed segments and re-fire the manifest refresh for every entry whose resident checksum mismatches the on-disk header — serving must not ungate until the manifest matches disk. reconcile failure aborts the transition (crash-loud).

steady (compactionSteady), the runSteadyCompactor loop:

  • timer: every CompactionInterval (default 4h; 0 disables compaction entirely — passes, live-set feeding, and rebuild all skip).
  • early trigger: the live consumer sends on a 1-buffered compactionTrigger channel (non-blocking, coalesced) when the in-memory set's Len reaches CompactionTombstoneCap. minCompactionTriggerSpacing = 30s floors back-to-back cap-triggered passes (defense against a failed pass making the consumer re-fire on every event). an early pass resets the timer.
  • each steady pass also runs the cheap manifest reconcile (checksum compare against resident entries) to heal the rewrite-succeeded/refresh-failed crash window.

cache invalidation (2026-06-12 design) #

why: rename changes durable bytes instantly, but serving keeps two in-memory structures — the manifest (headers, block indexes, blooms, collection indexes per segment) and the subscribe cold-path decoded-block LRU. either serving pre-compaction data after the rename is a correctness bug (clients would observe compacted-away rows below W).

mechanism, per rewritten segment, ordering deliberate:

  1. Manifest.OnSegmentCompacted(idx, path): reopen the file sealed, verify header/footer checksum, rebuild all resident metadata, swap the entry under the manifest lock. failure ⇒ pass fails, W does not advance.
  2. only then ColdReader.InvalidateSegment(idx): generation-aware — a generationBySegment counter is part of the block cache key; bump + purge resident entries. an in-flight decode may satisfy its own read but cannot re-insert under a stale generation (closes the purge-undone-by- in-flight-decode race). checksum stays in the key too, covering missed invalidations across crash windows.

manifest is NOT in pebble — it's a directory scan + self-describing file headers, so it can't drift from disk across restarts; the block cache is process-local and empty after restart. crash matrix: rewrite ok/refresh fails → next pass reconciles by checksum; crash before refresh → startup rereads disk; crash after refresh before watermark commit → next pass observes lag and continues.

crashpoints / oracle seams #

  • orchestrator: AfterCompactionRewriteBeforeWatermark (chunk rewritten, watermark not committed — recovery re-runs the chunk, idempotent), AfterCompactionChunkWatermark (committed, eviction may or may not have happened — rebuild-from-disk makes that immaterial).
  • segment: CrashPointRewriteTempWritten / TempSynced / Renamed / DirSynced, plus IOFaultInjector consulted before every tmp write, fsync, and rename.
  • hooks: OnBeforeCompactionPass(targetWatermark) fires after force-rotate with no rewrites yet (full pre-compaction snapshot point); OnCompactionPass{Watermark, Err} after every pass. the oracle's CheckCompacted model matches the window-fold (FoldRange) semantics exactly.
  • metrics worth porting: watermark, watermark-lag seconds, segments examined/rewritten/clean, rows dropped by reason (record/sync/account), tombstones collected by reason, manifest reconciles, early passes, tombstone set bytes/len.

Stream implementation receipt #

Stream implements this contract directly over its JSS archive and shared RocksDB metadata store:

  • account CBOR is decoded during observation/fold and malformed markers fail the append or pass; tombstone keys are cloned, strict drop reasons survive through physical rewrite, and the live set publishes upstream-equivalent entry and approximate-byte gauges;
  • every pass folds its bounded window from sealed JSS blocks, narrows real segment reads with the Gloom DID bloom up to the 100,000-DID cutoff, and fans all sealed segments across min(CPU, 8) workers;
  • delete and timestamp rewrites share the archive rewrite mutex. Cleanup is locked independently and mutation is locked per chunk, allowing safe import interleaving only between chunks;
  • rewritten files preserve block topology and historical envelopes, then use tmp write/fsync/rename/directory-fsync. Manifest refresh completes before the RocksDB watermark commits; steady passes heal checksum mismatches and merge-tail passes perform the required one-shot reconcile before serving;
  • the archive hook is fallible and runs after sequence assignment but before any flush or seal. Force rotation therefore cannot expose a sealed event whose tombstone observation has not succeeded;
  • a zero compaction interval disables both merge-tail and steady compaction; the merge destination is still sealed because serving requires sealed JSS;
  • all canonical compaction counters, gauges, reason labels, and the exact 16-bucket pass histogram are populated at these boundaries. Offline tests exercise segment files, concurrent workers, bloom avoidance, malformed input, manifest refresh, and durable watermark advancement.