Skip to content

createShapeGroup

createShapeGroup(options): object

Defined in: packages/client/src/circuits/shape-group.ts:57

Run K shape streams as one group.

This layer is ours permanently, not a bridge to something (ADR-0055 decision 10): the stream reader (readShapeStream) reads one stream and has no multi-stream coordinator, so the K-streams side of decision 4 has to exist either way. What it deliberately does NOT do is imitate Electric’s MultiShapeStream — it hands back envelopes and per-stream offsets, which is what the native path actually has, rather than reshaping them into a wire format we no longer speak.

Delivery is serialized. K streams arrive concurrently, but the subscriber is the apply path, whose inbox and commit queue assume one batch at a time. Rather than push that assumption onto every caller, the group chains deliveries: batches are handed over in arrival order, one at a time, and a slow subscriber exerts backpressure on every stream behind it. Per-shape ordering is preserved because each stream’s own reader is sequential.

ShapeGroupOptions

unsubscribeAll: () => void = closeEveryStream

void

get isUpToDate(): boolean

true once every shape in the group has reported up-to-date at least once — the group catch-up alignment of ADR-0031, which the native path inherits unchanged.

Latching, not instantaneous: a live stream that later has new data to deliver has not stopped being caught up, and the boot gate that reads this must not reopen once it has closed.

boolean

get pending(): string[]

Which shapes have not yet caught up — the boot rail’s answer to “what are we waiting on?”.

string[]

start(): Promise<void>

Open every stream. Resolves once each has produced its first response.

A start that FAILS closes whatever it managed to open. K streams open concurrently, so one refusing leaves the others reading a group the caller is about to discard — and since the caller’s answer to a failed open is another attempt (startCircuitsSync’s restart ladder), every attempt would strand another live reader. The group is spent either way: closed latches, and a stream still resolving its first response finds it set and closes itself below.

Promise<void>

subscribe(onBatch): void

(batch) => void | Promise<void>

void