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.
Parameters
Section titled “Parameters”options
Section titled “options”Returns
Section titled “Returns”unsubscribeAll
Section titled “unsubscribeAll”unsubscribeAll: () =>
void=closeEveryStream
Returns
Section titled “Returns”void
isUpToDate
Section titled “isUpToDate”Get Signature
Section titled “Get Signature”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.
Returns
Section titled “Returns”boolean
pending
Section titled “pending”Get Signature
Section titled “Get Signature”get pending():
string[]
Which shapes have not yet caught up — the boot rail’s answer to “what are we waiting on?”.
Returns
Section titled “Returns”string[]
start()
Section titled “start()”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.
Returns
Section titled “Returns”Promise<void>
subscribe()
Section titled “subscribe()”subscribe(
onBatch):void
Parameters
Section titled “Parameters”onBatch
Section titled “onBatch”(batch) => void | Promise<void>
Returns
Section titled “Returns”void