jetstream v2 in zig stream.waow.tech
stream docs upstream-subscribe-protocol.md
10 kB
Markdown

jetstream v1 /subscribe wire protocol (from upstream v2 rewrite) #

Superseded for v2 (2026-08-12): upstream replaced the /subscribe-v2 superset wire this document describes with network.bsky.jetstream.subscribeEvents (atproto proposal 0015: xrpc.v1.json message frames, lexicon-governed payloads, no record_cbor, RFC 3339 time). The lexicon file in upstream's lexicons/network/bsky/jetstream/subscribeEvents.json and upstream docs/README.md §5.2 are now authoritative for v2; the v1 sections below remain accurate (v1 additionally gained an optional envelope cursor).

condensed from bluesky-social/jetstream internal/subscribe/{doc,handler,filter,encoder,cursor,compress}.go and §5 of upstream's docs/README.md. Those are upstream Go paths — none of them are in this repository, and the docs/README.md is upstream's, not this docs/ directory. written 2026-07-08.

v1 and v2, since both appear throughout #

Stream serves both endpoints, /subscribe and /subscribe-v2 (serve/server.zig). They are the same event stream and the same JSON schema; they differ only in compression and one envelope field:

/subscribe (v1) /subscribe-v2
compression opt-in compress=true or Socket-Encoding: zstd; also negotiates permessage-deflate zstdDictionary=<id> only; silently declines permessage-deflate
cursor envelope field omitted optional

"v1 quirk" below therefore means deliberately preserved bug-compatible behaviour — clients depend on it, so it is held rather than fixed. The checklist at the end is the list of those.

Why the cursor magnitude split. A cursor below 1e15 is read as a sequence number and at or above it as microseconds since the epoch. That works because microsecond timestamps are already ~1.7e15 while sequence numbers start at 1, so the ranges cannot overlap until a single instance has assigned 10^15 events. The constant is 1_000_000_000_000_000 in serve/server.zig, and both endpoints use it.

The v2 dictionary ID is derived, not declared. v2_dictionary_id is computed from the embedded dictionary's own header (storage/zstd.zig), so the advertised ID and the served bytes cannot disagree. Tests pin the length (65,536), the ID (20260709), and a SHA-256 of the artifact. Note Zig writes these with digit separators — 20_260_709 — so grepping for 20260709 in the source finds nothing.

connection & query params #

endpoint: GET /subscribe → websocket upgrade. pre-upgrade errors are plain HTTP 400/503. malformed query string → 400.

wantedCollections (repeatable) #

  • absent → match-all
  • full NSID: strict validation, else 400 invalid options: invalid collection: <raw>
  • prefix pattern: any value ending .* → trim *, keep prefix ending in . (app.bsky.* → app.bsky.). NO NSID validation of head (v1 quirk — app.bsky.* with 2 segments must work). value ending * without preceding . → falls to NSID branch, rejected.
  • dedupe FIRST, then cap: unique count > 100 → 400 too many wanted collections
  • match: exact map hit, else prefix match

wantedDids (repeatable) #

  • each must be valid DID else 400 invalid options: invalid DID: <raw>; dedupe then cap 10_000 (unique)
  • applies to ALL event kinds (incl. identity/account)

cursor #

  • absent/empty → live tail
  • not base-10 int64 or negative → HTTP 400 pre-upgrade
  • both endpoints split by magnitude: < 1e15 = seq, >= 1e15 = time_us (µs since epoch)
  • future cursor → silently live mode, no error. cursor 0 → floored to first event.
  • too-old cursor on v1 → silently clamped to oldest retained (never rejected)
  • seq replay: first event with local seq >= cursor; timestamp replay: first event with witnessed_at >= cursor (inclusive-overlapping; clients dedupe)
  • timestamp translation happens before WebSocket upgrade using sealed segment headers, the selected block index, and the first matching row. A segment read/decode/index fault is a server fault: HTTP 503 with the generic body service not ready: cursor resolution failed, metric mode resolve_failed, and internal detail only in server logs. It is not a permanent client-input 400.
  • hot-window timestamp cursors resolve against the hot tail, not the archive (Stream translation of upstream's active-segment replay walk). Upstream's replay engine walks sealed segments and then the active segment before handing to live; Stream's port replaced the active-segment walk with the hot-tail deque, so a time cursor at or after the oldest retained hot frame must start on the hot tail directly. Resolving it through resolveTimeToSeq instead lands at the last sealed-segment boundary and strands the subscriber on the cold reader, which delivers read-batch bursts at block-seal cadence — the jerky-resume defect observed live 2026-08-08 (coral and bufo-bot resumed with near-now cursors after a proxy restart and trailed seals indefinitely). Known-open follow-up: a genuinely cold resume (cursor older than hot retention) still trails seal cadence after catching up — the cold→hot seam handoff in coldBatch does not fire at the tip (reproduced with a 10s-rewind probe that stayed on 1024-row bursts 7+ minutes after connect).

compress #

  • an ordinary /subscribe connection negotiates RFC 7692 permessage-deflate when offered, with 32 KiB context takeover and a 128-byte minimum message size
  • compress=true or header Socket-Encoding: zstd → per-message zstd, binary frames
  • zstd dictionary: dict ID 1612007021, 112640 bytes, upstream internal/subscribe/zstd_dictionary (verbatim v1 pkg/models/zstd_dictionary). out-of-band, both sides ship identical file. 128KiB window, each message a standalone zstd frame.
  • offering zstd AND permessage-deflate simultaneously → 400

v2 compression #

  • /subscribe-v2 rejects the v1 compress=true and Socket-Encoding: zstd opt-ins. Its only compression scheme is zstdDictionary=<id>.
  • /subscribe-v2 silently declines permessage-deflate; an offering client proceeds without that extension unless it selected the v2 zstd dictionary
  • the current structured dictionary is the exact pinned upstream 65,536-byte artifact with ID 20260709. Clients fetch it from network.bsky.jetstream.getZstdDictionary, then reconnect with that ID.
  • a malformed ID is 400. An unknown/retired ID is 400 with the stable unknown zstd dictionary id marker and current ID. V1 ignores this v2-only parameter.
  • the dictionary XRPC serves exact bytes with ETag "zstd-dict-20260709", X-Zstd-Dictionary-Id: 20260709, and one-year immutable public caching; exact-ID lookup and conditional 304 are supported.

requireHello #

  • true iff value is exactly literal "true"
  • when true: upgrade accepted, no events until first valid options_update; no pre-hello buffering; invalid update during wait disconnects

maxMessageSizeBytes #

  • empty/malformed/negative/overflow → silently 0 = no cap, never an error
  • enforced post-encode per event on uncompressed JSON length; oversize event silently skipped, cursor advances

keepalive/limits #

  • server ping every 30s; 5s per-frame write timeout
  • client→server read limit 10_000_000 bytes (decimal)
  • close-frame reason truncated to 123 bytes rune-aligned + "..."

JSON event schema #

{"did":"did:plc:...","time_us":1725911162329308,"kind":"commit|identity|account",
 "commit":{"rev":"...","operation":"create|update|delete","collection":"...","rkey":"...",
           "record":{...},"cid":"bafyrei..."}}
  • one websocket text frame per event (binary if zstd)
  • delete: only rev/operation/collection/rkey (no record/cid). rev/operation/collection/rkey always present.
  • record = DAG-CBOR block decoded to JSON; cid = CID over raw DAG-CBOR (string form)
  • identity: "identity":{"did","handle"(opt),"seq","time"} — seq is UPSTREAM relay seq, time RFC3339
  • account: "account":{"active","did","seq","status"(opt),"time"}
  • error frames: top-level {"error","message"}
  • upstream v2 adds optional "cursor" (omitempty) envelope field; pure v1 omits

sourcing from firehose #

  • #commit → one wire event PER OP. drops: non-TID rev → whole event; invalid collection NSID or rkey → that op; create/update op with record block missing from CAR → that op. all ops share witnessed_at.
  • #identity/#account → one event each
  • #sync → never emitted on v1 wire (skipped, cursor advances)
  • #info → nothing
  • identity/account events bypass wantedCollections (respect wantedDids only)
  • commit with EMPTY collection bypasses collection filter

client→server messages #

text frame: {"type":"options_update","payload":{"wantedCollections":[],"wantedDids":[],"maxMessageSizeBytes":0}}

  • envelope parse fail → close 1007 "bad SubscriberSourcedMessage envelope"
  • payload parse fail → close 1007; payload validation fail → close 1008 with error text
  • valid update REPLACES whole filter atomically; signals hello
  • unknown type → logged and ignored (never fatal); non-text frames ignored

v1 quirks checklist #

  1. maxMessageSizeBytes garbage → 0, never error
  2. identity/account bypass wantedCollections
  3. #sync not emitted
  4. unknown SubscriberSourcedMessage type ignored
  5. prefix patterns: no NSID validation of head
  6. caps fire post-dedupe
  7. empty-collection commits bypass collection filter
  8. requireHello exact-literal "true" only
  9. below-floor cursor clamped, never rejected
  10. read limit 10_000_000 decimal

stream-only observability additions #

  • stream_subscribe_worker_exits_total{reason} — why a subscriber's delivery worker ended (tail_closed, read_batch_error, write_error, write_closed, ping_write_error, too_slow, cold_failed, cold_stall). Upstream has no equivalent: its handler exits are as silent as ours were. Added 2026-08-13 after eval connections died during compaction passes with no log line and no counter anywhere. Every reason renders even at zero, so an absent series means "not scraped", never "never happened". The wedged-writer kill itself (frame_write_timeout_ms = 5s) IS upstream behavior (frameWriteTimeout, internal/subscribe/handler.go) — only the accounting is ours.