the client surface #
mirrors the official Go client's shape: one subscribe covers the live tail, filtered archive replay, and a gapless replay→live cutover. pinned upstream: bluesky-social/jetstream 289b032 (v0.2.0).
subscribe — the unified client #
the port of upstream jetstream.Subscribe. after_seq sweeps the sealed
archive and cuts over to the live tail at the sealed tip; the client's
per-event seq dedup makes the seam at-least-once with no gap, and a
CursorTooOld refusal re-enters backfill from the last delivered seq
(bounded by max_rebackfill_stalls, monotonic non-decreasing).
snapshot_only stops at the sealed range. without after_seq it is a
pure live tail. live_cursor accepts a host-local sequence or a Unix
microsecond timestamp (values >= 1e15); null starts at the tip. Sequence
resumes deduplicate the inclusive boundary. Timestamp resumes establish a
new sequence watermark from received events.
The archive sweep follows every plan page, pinning the first page's sealed
tip as the upper bound. A plan that does not advance fails with PlanStalled.
Collection filters accept exact names and namespace wildcards such as
app.bsky.feed.*; archive and live delivery both enforce filters locally.
DID-level markers bypass collection filters and remain subject to DID and
kind filters.
multi-host failover covers the live tail. hosts lists the primary first;
after repeated failures, the client uses a rewound witnessed timestamp to
enter the next host's live stream without archive requests. See
failover.md for rotation, retention, and delivery limits.
LiveClient — the live tail standalone #
the single-host subscribeEvents WebSocket client: lexicon frames,
sequence or timestamp cursors, kinds/collections/dids filters,
dictionary-zstd with rotation recovery, and reconnect backoff. Reconnects
use the highest received sequence after the initial cursor resolves.
dedup_floor is always a sequence in that host's namespace, never a timestamp.
TCP keepalive is enabled; event-based stall detection is opt-in through
stall_timeout_ns.
A rejected sequence returns CursorTooOld. A timestamp older than server
retention can be clamped: OutdatedCursor reaches onInfo, or is logged
when no callback is supplied. The live-only path does not start archive
replay automatically. InvalidRequest is terminal; an
UnknownZstdDictionary rejection refreshes the dictionary or falls back to
uncompressed delivery.
ArchiveBackfill — the archive half standalone #
planSnapshot (collections/dids, seq bounds, api_key for token-gated
instances) → getBlock/getSegment → jss v1 columnar decode. two
handler contracts: the classic onEvent(zat.JetstreamEvent) (commits
only, kept for existing consumers) and the unified client's
onRow(jetstream.Event) (every row with its seq, including
identity/account/sync markers). fetchSeqBounds maps a witnessed-time
window to plan seq bounds.
archive errors #
An optional onError(anyerror) bool callback chooses whether to continue
after a recoverable archive error. Returning true accepts the error and
continues on the same instance; false stops cleanly. Omitting the callback
returns the error to the caller, so data is never skipped without a decision.
A failed block download emits the entry's good prefix, reports the error, and skips its remaining blocks before proceeding to the next entry. A decompression failure skips that block; a malformed record preserves other valid rows from the block and reports the decode error after them. Whole segment download/header failures can likewise proceed to the next entry. Plan failures, cancellation, and allocation failure remain terminal. Accepting an error means accepting an incomplete portion of the archive; it does not repair that portion or trigger failover.
records #
commit events carry both representations, matching upstream Event:
record_cbor— the record's canonical DAG-CBOR. live records are canonicalized from the wire JSON exactly like upstreamdecodeLiveRecord(whole numbers to integers,$bytes/$linksentinels, finite f64 floats); archive records alias the stored payload verbatim (block-buffer lifetime).record— the genericstd.json.Valueform, derived from the canonical bytes (not the wire JSON), so both representations agree.
deletes carry neither. a create/update record that cannot canonicalize
is a MalformedRecord — surfaced but non-fatal, like every malformed
data frame (one bad frame never drops the tail).
the v1 /subscribe wire client (time_us cursors, host round-robin)
remains zat.JetstreamClient in zat.
Authentication: the unified client’s api_key authenticates archive requests
only. Live WebSocket connections are public and never send that credential.
Standalone live options have no api_key field.
Archive GET retries make up to three attempts. The fallback starts at 500ms
and doubles, capped at 30s. Block/list requests add 0–20% jitter and honor a
positive server delay in place of the fallback; whole-segment requests use
no jitter and the longer of the fallback and server delay. Absolute
RateLimit-Reset takes precedence over Retry-After (seconds or IMF-fixdate).
Past hints use the fallback, and all server-directed waits are capped at 30s.
Live dictionary-zstd compression is enabled by default. Set
zstd_compression = false to disable dictionary fetching and negotiation.
A rejected dictionary is refetched from the same live host; a failed, invalid,
or unchanged refresh falls back to uncompressed delivery.