Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 6 additions & 6 deletions API.md
Original file line number Diff line number Diff line change
Expand Up @@ -1017,15 +1017,15 @@ Two patterns for sharing a single source among multiple consumers:
### `Stream.broadcast(options?)`

Create a push-model multi-consumer channel. Data written to the [writer](#writer-interface) is
delivered to all consumers that have subscribed via `broadcast.push()`.
delivered to all consumers that have subscribed via `channel.push()`.
Data remains in the buffer and is available to consumers that attach before
it is overwritten. Late-joining consumers begin reading from the oldest entry
still in the buffer at the time they call `broadcast.push()`.
still in the buffer at the time they call `channel.push()`.

```typescript
function broadcast(options?: BroadcastOptions): {
writer: Writer;
broadcast: Broadcast;
channel: Broadcast;
}
```

Expand All @@ -1051,11 +1051,11 @@ interface Broadcast {

**Example:**
```typescript
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 100 });

// Create consumers before writing
const consumer1 = broadcast.push();
const consumer2 = broadcast.push(decompress);
const consumer1 = channel.push();
const consumer2 = channel.push(decompress);

// Producer and consumers must run concurrently. Awaited writes
// block when the buffer fills until consumers read.
Expand Down
6 changes: 3 additions & 3 deletions SLIDES.md
Original file line number Diff line number Diff line change
Expand Up @@ -350,11 +350,11 @@ const bytesWritten = await Stream.pipeTo(
### Broadcast (Push)

```typescript
const { writer, broadcast } =
const { writer, channel } =
Stream.broadcast({ highWaterMark: 100 });

const c1 = broadcast.push();
const c2 = broadcast.push(decompress);
const c1 = channel.push();
const c2 = channel.push(decompress);

// Producer and consumers run concurrently
(async () => {
Expand Down
20 changes: 10 additions & 10 deletions benchmarks/10-advanced-features.ts
Original file line number Diff line number Diff line change
Expand Up @@ -66,11 +66,11 @@ async function runBenchmarks(): Promise<void> {
const broadcastResult = await benchmark(
'broadcast()',
async () => {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 1000 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 1000 });

// Create two consumers
const consumer1 = broadcast.push();
const consumer2 = broadcast.push();
const consumer1 = channel.push();
const consumer2 = channel.push();

// Producer task
const producerTask = (async () => {
Expand Down Expand Up @@ -140,10 +140,10 @@ async function runBenchmarks(): Promise<void> {
const broadcastResult = await benchmark(
'broadcast()+transform',
async () => {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 1000 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 1000 });

const consumer1 = broadcast.push();
const consumer2 = broadcast.push(xorTransform);
const consumer1 = channel.push();
const consumer2 = channel.push(xorTransform);

const producerTask = (async () => {
for (const chunk of chunks) {
Expand Down Expand Up @@ -597,13 +597,13 @@ async function runBenchmarks(): Promise<void> {
const result = await benchmark(
`${policy}`,
async () => {
const { writer, broadcast } = Stream.broadcast({
const { writer, channel } = Stream.broadcast({
highWaterMark: 50,
backpressure: policy,
});

// Create a slow consumer (will cause backpressure)
const consumer = broadcast.push();
const consumer = channel.push();

// Fast producer
const producerTask = (async () => {
Expand Down Expand Up @@ -735,15 +735,15 @@ async function runBenchmarks(): Promise<void> {
const result = await benchmark(
`broadcast ${numConsumers} consumers`,
async () => {
const { writer, broadcast } = Stream.broadcast({
const { writer, channel } = Stream.broadcast({
highWaterMark: 50,
backpressure: 'block',
});

// Create N consumers
const consumers: AsyncIterable<Uint8Array[]>[] = [];
for (let i = 0; i < numConsumers; i++) {
consumers.push(broadcast.push());
consumers.push(channel.push());
}

// Producer writes all chunks then ends
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/20-memory-allocations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -233,9 +233,9 @@ boxplot(() => {
summary(() => {
bench('broadcast 2 consumers (new)', function* () {
yield async () => {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 });
const c1 = broadcast.push();
const c2 = broadcast.push();
const { writer, channel } = Stream.broadcast({ highWaterMark: 100 });
const c1 = channel.push();
const c2 = channel.push();
const producing = (async () => {
for (const chunk of bcastChunks) await writer.write(chunk);
await writer.end();
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/21-memory-sustained.ts
Original file line number Diff line number Diff line change
Expand Up @@ -339,9 +339,9 @@ async function main() {

console.log('Running: Broadcast / tee (2 consumers)...');
const bcastNew = await measureSustained('broadcast 2x (new)', async () => {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 });
const c1 = broadcast.push();
const c2 = broadcast.push();
const { writer, channel } = Stream.broadcast({ highWaterMark: 100 });
const c1 = channel.push();
const c2 = channel.push();
const producing = (async () => {
for (const chunk of bcastChunks) await writer.write(chunk);
await writer.end();
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/22-memory-backpressure.ts
Original file line number Diff line number Diff line change
Expand Up @@ -92,12 +92,12 @@ for (const policy of policies) {
for (const hwm of hwmValues) {
bench(`bcast ${policy} hwm=${hwm}`, function* () {
yield async () => {
const { writer, broadcast } = Stream.broadcast({
const { writer, channel } = Stream.broadcast({
highWaterMark: hwm,
backpressure: policy,
});
const c1 = broadcast.push();
const c2 = broadcast.push();
const c1 = channel.push();
const c2 = channel.push();

const producing = (async () => {
if (policy === 'drop-oldest' || policy === 'drop-newest') {
Expand Down
16 changes: 8 additions & 8 deletions benchmarks/html/04-branching.html
Original file line number Diff line number Diff line change
Expand Up @@ -90,11 +90,11 @@ <h2>Raw Output</h2>
// ========================================

async function newStreamBroadcast2(chunks) {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 1000 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 1000 });

// Create 2 consumers
const consumer1 = broadcast.push();
const consumer2 = broadcast.push();
const consumer1 = channel.push();
const consumer2 = channel.push();

// Write all chunks
(async () => {
Expand All @@ -114,14 +114,14 @@ <h2>Raw Output</h2>
}

async function newStreamBroadcast4(chunks) {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 1000 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 1000 });

// Create 4 consumers
const consumers = [
broadcast.push(),
broadcast.push(),
broadcast.push(),
broadcast.push()
channel.push(),
channel.push(),
channel.push(),
channel.push()
];

// Write all chunks
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/html/benchmarks.js
Original file line number Diff line number Diff line change
Expand Up @@ -208,9 +208,9 @@ export async function webStreamText(chunks) {
// =============================================================================

export async function newStreamBroadcast2(Stream, chunks) {
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 1000 });
const c1 = broadcast.push();
const c2 = broadcast.push();
const { writer, channel } = Stream.broadcast({ highWaterMark: 1000 });
const c1 = channel.push();
const c2 = channel.push();
(async () => {
for (const chunk of chunks) await writer.write(chunk);
await writer.end();
Expand Down
4 changes: 2 additions & 2 deletions benchmarks/profile-systematic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -162,10 +162,10 @@ async function main() {
});

await measure('broadcast() single consumer', iterations, async () => {
const { writer, broadcast } = Stream.broadcast<Uint8Array>();
const { writer, channel } = Stream.broadcast<Uint8Array>();
let total = 0;
const readPromise = (async () => {
for await (const batch of broadcast.consume()) {
for await (const batch of channel.consume()) {
for (const chunk of batch) total += chunk.length;
}
})();
Expand Down
8 changes: 4 additions & 4 deletions docs/COMPLETENESS-ANALYSIS.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,12 +112,12 @@ handled at the source level, not the stream level.

```typescript
// Video transcoding pipeline
const { writer, broadcast } = Stream.broadcast();
const { writer, channel } = Stream.broadcast();

// Multiple output qualities
const hd = broadcast.push(transcodeToHD);
const sd = broadcast.push(transcodeToSD);
const thumbnail = broadcast.push(extractThumbnails);
const hd = channel.push(transcodeToHD);
const sd = channel.push(transcodeToSD);
const thumbnail = channel.push(extractThumbnails);

// Process all in parallel
await Promise.all([
Expand Down
10 changes: 5 additions & 5 deletions docs/DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -416,7 +416,7 @@ Creates a multi-consumer broadcast channel where a writer pushes to all consumer
```typescript
function broadcast(options?: BroadcastOptions): {
writer: Writer;
broadcast: Broadcast;
channel: Broadcast;
}

interface Broadcast {
Expand All @@ -441,11 +441,11 @@ interface BroadcastOptions {

**Example:**
```typescript
const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 });
const { writer, channel } = Stream.broadcast({ highWaterMark: 100 });

// Create consumers with different transforms
const consumer1 = broadcast.push();
const consumer2 = broadcast.push(decompress);
const consumer1 = channel.push();
const consumer2 = channel.push(decompress);

// Producer
for await (const chunk of source) {
Expand Down Expand Up @@ -507,7 +507,7 @@ const [raw, decompressed, parsed] = await Promise.all([
|--------|---------------|-----------|
| Model | Push (writer -> consumers) | Pull (source -> consumers) |
| Data source | Writer pushes explicitly | Source pulled on demand |
| Create consumer | `broadcast.push()` | `share.pull()` |
| Create consumer | `channel.push()` | `share.pull()` |
| Sync version | No | Yes (`shareSync`) |
| Use case | Event sources, WebSocket | File, response body, iterables |

Expand Down
6 changes: 3 additions & 3 deletions docs/MIGRATION-NODEJS.md
Original file line number Diff line number Diff line change
Expand Up @@ -894,10 +894,10 @@ const [result1, result2] = await Promise.all([
**New Stream API - Push Model (broadcast):**
```javascript
// broadcast() - producer pushes to all consumers
const { writer, broadcast } = Stream.broadcast();
const { writer, channel } = Stream.broadcast();

const consumer1 = broadcast.push();
const consumer2 = broadcast.push();
const consumer1 = channel.push();
const consumer2 = channel.push();

// Push data
await writer.write('shared data');
Expand Down
14 changes: 7 additions & 7 deletions docs/MIGRATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -599,10 +599,10 @@ const [result1, result2] = await Promise.all([
**New Stream API - Push Model (broadcast):**
```javascript
// broadcast() - push-based, producer controls data flow
const { writer, broadcast } = Stream.broadcast();
const { writer, channel } = Stream.broadcast();

const consumer1 = broadcast.push();
const consumer2 = broadcast.push();
const consumer1 = channel.push();
const consumer2 = channel.push();

// Producer pushes to all consumers
await writer.write('shared data');
Expand Down Expand Up @@ -642,7 +642,7 @@ const shared = Stream.share(source, {
backpressure: 'drop-oldest' // or 'strict', 'block', 'drop-newest'
});

const { writer, broadcast } = Stream.broadcast({
const { writer, channel } = Stream.broadcast({
highWaterMark: 100,
backpressure: 'block' // Wait for space (use 'strict' to reject)
});
Expand Down Expand Up @@ -772,9 +772,9 @@ try {

// Async cleanup with 'await using'
{
const { writer, broadcast } = Stream.broadcast();
await using _ = broadcast; // Will cancel on scope exit
// Use broadcast...
const { writer, channel } = Stream.broadcast();
await using _ = channel; // Will cancel on scope exit
// Use channel...
}
```

Expand Down
6 changes: 3 additions & 3 deletions docs/TRANSFER-INTEGRATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ This section identifies which types in the new streams API would implement `[Sym
| Writer (from `Stream.push()`) | Yes | Single-owner write endpoint; transfer moves write authority |
| DuplexChannel | Yes | Bundles writer + readable for one endpoint; single unit of ownership |
| Share consumer (from `share.pull()`) | Yes | Each consumer iterable is single-consumer |
| Broadcast consumer (from `broadcast.push()`) | Yes | Each consumer iterable is single-consumer |
| Broadcast consumer (from `channel.push()`) | Yes | Each consumer iterable is single-consumer |
| Share instance | Possible | Multi-consumer wrapper owns the source; transfer moves management authority |
| Broadcast instance | Possible | Less clear value; the writer side is the primary ownership concern |
| `WriterIterablePair` (from `Stream.push()`) | Possible | Atomic transfer of both writer and readable together |
Expand Down Expand Up @@ -193,7 +193,7 @@ A `DuplexChannel` bundles a writer (sends to the peer) and a readable (receives

### 3.4 Share Consumer / Broadcast Consumer

Each call to `share.pull()` or `broadcast.push()` returns an `AsyncIterable<Uint8Array[]>` representing one consumer's view of the shared/broadcast data.
Each call to `share.pull()` or `channel.push()` returns an `AsyncIterable<Uint8Array[]>` representing one consumer's view of the shared/broadcast data.

**Transfer behavior:**

Expand Down Expand Up @@ -512,7 +512,7 @@ writer.desiredSize; // null

### 5.5 Transfer and Multi-Consumer Patterns

Transferring a consumer from `share.pull()` or `broadcast.push()` affects only that consumer. Other consumers are unaffected.
Transferring a consumer from `share.pull()` or `channel.push()` affects only that consumer. Other consumers are unaffected.

```js
const shared = Stream.share(source, { highWaterMark: 100 });
Expand Down
6 changes: 3 additions & 3 deletions index.bs
Original file line number Diff line number Diff line change
Expand Up @@ -326,7 +326,7 @@ dictionary PushStreamResult {

dictionary BroadcastResult {
required Writer writer;
required Broadcast broadcast;
required Broadcast channel;
};
</pre>

Expand Down Expand Up @@ -999,8 +999,8 @@ The <dfn method for="Stream">broadcast(options)</dfn> method creates a [=broadca
<li>Let |backpressure| be |options|["{{BroadcastOptions/backpressure}}"].
<li>Create a shared circular buffer with [=byte budget=] |budget|.
<li>Let |writer| be a new {{Writer}} backed by the shared buffer with |backpressure| policy.
<li>Let |broadcast| be a new {{Broadcast}} object backed by the shared buffer.
<li>Return «[ "writer" → |writer|, "broadcast" → |broadcast| ]».
<li>Let |channel| be a new {{Broadcast}} object backed by the shared buffer.
<li>Return «[ "writer" → |writer|, "channel" → |channel| ]».
</ol>
</div>

Expand Down
Loading
Loading