jetstream client
atproto jetstream client
jetstream docs client.md
5.7 kB
Markdown

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 upstream decodeLiveRecord (whole numbers to integers, $bytes/$link sentinels, finite f64 floats); archive records alias the stored payload verbatim (block-buffer lifetime).
  • record — the generic std.json.Value form, 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.