jetstream v2 in zig stream.waow.tech
stream docs timestamp-import.md
7.0 kB
Markdown

Timestamp import #

Stream follows Jetstream V2's operator timestamp-import design. An import changes the timestamp presented to subscribers without changing the immutable witness timestamp used for archive ranges and timestamp cursors.

What this is for, and the one distinction the whole design turns on #

Records carry two timestamps, and an import moves only one of them:

changed by an import? used for
display timestamp (time_us on the wire) yes what subscribers see
witness timestamp (witnessed_at) never archive ranges, timestamp cursors, segment header min/max

That separation is what makes an import safe. witnessed_at is when this instance observed the event, so it defines segment ordering and is what a ?cursor=<microseconds> replay resolves against. Moving it would invalidate every sealed segment's header range and every timestamp cursor, and break the "segment files sort in creation order, and that order is time order" invariant. Imports are a presentation layer over immutable history.

The operator use case is backfilled or migrated content whose true creation time is older than when this instance first saw it — without an import, a repository imported today would present all its history as today's events.

Verified 2026-07-30: --timestamp-import-dir and --timestamp-import-token both exist in runtime/cli.zig, and /xrpc/network.bsky.jetstream.importTimestamps and /xrpc/network.bsky.jetstream.getImportStatus are both routed in serve/server.zig. (Those are matched as full literal paths, so grepping the source for the bare NSID finds nothing.)

Source format #

The source is a plain, seekable RFC 4180 CSV. uri and timestamp are required header columns; scope and cid are optional. Header order and case do not matter, but unknown and duplicate columns fail the import because they make every row ambiguous.

uri,timestamp,scope,cid
at://did:plc:example/app.bsky.feed.post/3abc,2026-07-15T10:00:00Z,all_versions,

An empty scope means all_versions. specific_version requires a valid CIDv1 using dag-cbor or raw content with SHA-256. URI authorities must be DIDs and must include both collection and record key.

Parsing is a single forward-only pass with a fixed 64 KiB row window. Every accepted row carries its source byte offset so the apply pass can seek back to the same instruction. Ordinary malformed rows are counted and skipped; structurally invalid headers and quote errors that consume the remainder of a CSV fail the job. Rejection samples are bounded independently of full per-reason counts.

Durable rule map #

Validated rules are bulk-loaded into <data-dir>/import-rules, a dedicated RocksDB database with full Bloom filters. Imports build sorted external SSTs in bounded chunks and ingest them in CSV order. Duplicate keys are therefore last-write-wins both within and across chunks without routing a potentially huge cold keyspace through the metadata database or its WAL.

The append lookup runs under the archive writer lock before sequence assignment or buffering. It first checks an in-memory set of imported collections, then checks a specific-version rule using the DAG-CBOR payload CID, and finally falls back to an all-versions path rule. Specific-version matches take precedence. The stamped row is returned through the archive transaction so immediate v1/v2 delivery and durable JSS replay encode the same display timestamp. Resync replacements, bootstrap capture, and merge all use the same central hook. The collection set is reconstructed from durable marker keys on every open; RocksDB read failures poison the archive transaction rather than silently changing the timestamp selected for a record.

Implementation state #

The streaming validation boundary, durable external-SST rule database, archive-transaction stamping, rule precedence, hot/cold subscriber display-time agreement, and atomic topology-preserving JSS patch primitive are in place and covered by offline production-path tests. The patcher accepts an upstream-produced sealed segment, changes only indexed_at, copies clean zstd frames and the bloom/collection footer tail byte-for-byte, rejects forbidden cross-row mutation, and makes an already-applied patch a no-op. The sealed-segment catalog, bloom-routed bounded-LRU bucketer, durable packed offset files, positioned row revalidation, last-write-wins patch planner, and specific-CID precedence are also implemented. An offline disk-to-disk test forces descriptor/cache eviction while importing one CSV across two real JSS files and verifies the final rows plus idempotent reapplication. The unified metadata RocksDB owns the synchronous current-job pointer, complete job record and per-segment completion set. Their reopen test proves records and checkpoints survive process restart, segment keys retain their numeric identity, and terminal cleanup is durable. The production runner binds that state to rule activation, forced sealing, parse/bucket, per-segment apply checkpoints, cancellation and resume. Its offline receipt stops after one of two real JSS patches, closes the archive and both RocksDB stores, reopens everything, skips the durable done marker, finishes the second patch, and verifies terminal cleanup and both on-disk timestamps. The manager canonicalizes the staging root and target through symlinks, admits regular in-root files only, enforces the durable single-job slot, runs imports in the background, joins cancellation on shutdown, and adopts a nonterminal current job at startup.

Operator surface #

POST /xrpc/network.bsky.jetstream.importTimestamps accepts {"path":"..."} and immediately returns {"job":"..."}. GET /xrpc/network.bsky.jetstream.getImportStatus?job=... returns the canonical lifecycle, phase, timestamps, parse totals, apply progress, and mutation totals; omitting job selects the current or most recent import. Both endpoints require the exact bearer secret configured by --timestamp-import-token. An empty configured token disables access, and disabled/missing/wrong credentials have the same response. The comparison hashes both tokens to fixed-width SHA-256 and compares the digests in constant time. TLS terminates at the operator's reverse proxy.

The default staging root is <data-dir>/imports; override it with --timestamp-import-dir. The manager reports ImportNotReady until bootstrap has reached steady state, but authenticated status remains available during bootstrap. The runner exports the canonical job-result/duration, phase, parse/rejection, routing, match, corruption, mutation, segment, and rewritten byte metrics at the same lifecycle boundaries as upstream. Delete compaction and timestamp patching share a dedicated sealed-segment rewrite lock which does not block live append; a concurrent production-path test verifies the compactor waits and advances its watermark only after the lock is released.