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'sdocs/README.md(upstream's, not thisdocs/directory — there is nodocs/README.mdhere). - 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), andcompact/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 byKindDeleteandKindUpdaterows (update = delete+replace; only latest record version is kept). - did tombstone: key
did→{seq, reason}where reason ∈ {"sync", "account"}. produced byKindSyncrows (repo diverged, full resync follows) and byKindAccountrows whose payload decodes toactive=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 withMaxSeq <= watermarkare 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
MaxSeqacross 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) #
- cleanup: delete stale
*.jss.tmpin 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. - 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).
- 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. - chunk loop,
current = WuntiltargetWatermark: 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 tomin(NumCPU, 8)workers; each callssegment.Rewrite(path, decide)with decide = keep ifseq > chunkEnd, else drop ifsnap.ShouldDrop. a bloom prefilter narrows by candidate DIDs unless the snapshot has more thanCompactionBloomNarrowMaxDIDs(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, fireOnSegmentCompacted(idx, path)(manifest refresh + cache invalidation, below). refresh failure fails the pass; W does not advance. e. commit:saveCompactionWatermark(chunkEnd), thenEvict(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=0block, 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
rewriteMuserializes 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
compactionTriggerchannel (non-blocking, coalesced) when the in-memory set's Len reachesCompactionTombstoneCap.minCompactionTriggerSpacing = 30sfloors 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:
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.- only then
ColdReader.InvalidateSegment(idx): generation-aware — agenerationBySegmentcounter 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, plusIOFaultInjectorconsulted 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'sCheckCompactedmodel 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.