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.