Skip to content

readShapeStream

readShapeStream(options, onBatch, onEnd?): Promise<ShapeStreamSubscription>

Defined in: packages/client/src/circuits/stream-source.ts:157

Read one Circuits shape stream, yielding envelopes in stream order.

pgxsinkit’s own long-poll reader of the Durable Streams protocol (ADR-0065 decision 6; the wire is long-poll.ts). It catches up with plain reads from offset (so a catch-up response stays cacheable), and once a response says it is up to date, a live read long-polls the tail, echoing the server’s cursor. live: false stops at the first up-to-date response instead.

Resolves once the first response has arrived, so an immediate failure — a refused token, a stream that does not exist — rejects this call rather than dying on a loop nobody is watching. The first batch is delivered after that, never during the call: a caller that closes its subscription from inside onBatch already holds it.

Every successful response is one batch, an empty one included (a long poll that timed out delivers the empty up-to-date batch at the same offset). Backpressure, one response ahead: the request for the next response is sent when a batch is handed to onBatch, and the request after that not until onBatch has settled. So the network wait and the apply overlap instead of adding up, which is worth up to half the time of a catch-up that spans many responses, and a slow apply still throttles the read: at most one response is ever held beyond the one being applied. One request is in flight at a time, which is one connection per stream while it long-polls.

Failures, per request:

  • 401/403: one immediate retry with a token from onTokenRejected (createTokenRecovery); a second rejection of the same request, or no fresh token, is terminal. Every request of the session gets this, not only the opening one, so a read that lives for hours can re-mint many times, but a refused token never loops.
  • 429 and every 5xx, and network failures (the request, or its body, cut off): retried with backoff until they succeed or the caller stops the read. A Retry-After is a floor under the backoff.
  • Every other status is terminal, 404 and 410 (the stream was never created, or was retired) included, as is a response that breaks the protocol (readResponse).

The end of a read. onEnd is not optional decoration: without it a read that dies mid-session is silent — the rows stay, and nobody re-subscribes. It fires at most once:

  • onEnd(null) after onBatch has been handed the batch that ends the read — the up-to-date batch of a live: false read, the Stream-Closed batch of any read — and has settled on it. A caller may close on onEnd without losing that batch.
  • onEnd(error) for a terminal failure after the first response, and for an onBatch that throws. This is the mid-session reset (ADR-0056 decision 7) a caller answers with a re-subscribe.
  • onEnd(null) when the caller’s signal aborts.
  • Nothing when the caller called ShapeStreamSubscription.close: a group tearing its streams down would otherwise hear K “the stream ended” reports and try to recover from a stop it ordered.
  • Nothing for a failure before the first response: that rejects this call instead.

The transport only. Everything above it — the fold, apply modes, the boot gate — is the engine’s, which is ADR-0009’s precedent applied to a new substrate.

StreamSourceOptions

(batch) => void | Promise<void>

(error) => void

Promise<ShapeStreamSubscription>