From d4b920c4f5a4a6bc99902950ae33c8ae450f7fb4 Mon Sep 17 00:00:00 2001 From: "Kamat, Trivikram" <16024985+trivikr@users.noreply.github.com> Date: Sun, 30 Aug 2026 15:08:33 -0700 Subject: [PATCH] Rename broadcast result property to channel --- API.md | 12 +++--- SLIDES.md | 6 +-- benchmarks/10-advanced-features.ts | 20 ++++----- benchmarks/20-memory-allocations.ts | 6 +-- benchmarks/21-memory-sustained.ts | 6 +-- benchmarks/22-memory-backpressure.ts | 6 +-- benchmarks/html/04-branching.html | 16 +++---- benchmarks/html/benchmarks.js | 6 +-- benchmarks/profile-systematic.ts | 4 +- docs/COMPLETENESS-ANALYSIS.md | 8 ++-- docs/DESIGN.md | 10 ++--- docs/MIGRATION-NODEJS.md | 6 +-- docs/MIGRATION.md | 14 +++---- docs/TRANSFER-INTEGRATION.md | 6 +-- index.bs | 6 +-- samples/03-branching.ts | 22 +++++----- samples/06-buffer-configuration.ts | 4 +- samples/09-resource-management.ts | 22 +++++----- samples/11-real-world-patterns.ts | 6 +-- samples/12-nodejs-interop.ts | 6 +-- samples/13-web-streams-interop.ts | 6 +-- samples/html/06-branching.html | 26 ++++++------ src/broadcast.js | 4 +- src/broadcast.test.ts | 62 ++++++++++++++-------------- src/broadcast.ts | 4 +- src/types.ts | 4 +- 26 files changed, 149 insertions(+), 149 deletions(-) diff --git a/API.md b/API.md index 7ad9366..7a26cdd 100644 --- a/API.md +++ b/API.md @@ -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; } ``` @@ -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. diff --git a/SLIDES.md b/SLIDES.md index 7876833..719d98e 100644 --- a/SLIDES.md +++ b/SLIDES.md @@ -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 () => { diff --git a/benchmarks/10-advanced-features.ts b/benchmarks/10-advanced-features.ts index f39cbf9..486cfd2 100644 --- a/benchmarks/10-advanced-features.ts +++ b/benchmarks/10-advanced-features.ts @@ -66,11 +66,11 @@ async function runBenchmarks(): Promise { 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 () => { @@ -140,10 +140,10 @@ async function runBenchmarks(): Promise { 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) { @@ -597,13 +597,13 @@ async function runBenchmarks(): Promise { 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 () => { @@ -735,7 +735,7 @@ async function runBenchmarks(): Promise { const result = await benchmark( `broadcast ${numConsumers} consumers`, async () => { - const { writer, broadcast } = Stream.broadcast({ + const { writer, channel } = Stream.broadcast({ highWaterMark: 50, backpressure: 'block', }); @@ -743,7 +743,7 @@ async function runBenchmarks(): Promise { // Create N consumers const consumers: AsyncIterable[] = []; for (let i = 0; i < numConsumers; i++) { - consumers.push(broadcast.push()); + consumers.push(channel.push()); } // Producer writes all chunks then ends diff --git a/benchmarks/20-memory-allocations.ts b/benchmarks/20-memory-allocations.ts index 6a7aced..05f600e 100644 --- a/benchmarks/20-memory-allocations.ts +++ b/benchmarks/20-memory-allocations.ts @@ -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(); diff --git a/benchmarks/21-memory-sustained.ts b/benchmarks/21-memory-sustained.ts index 2d74c8c..35cb2a4 100644 --- a/benchmarks/21-memory-sustained.ts +++ b/benchmarks/21-memory-sustained.ts @@ -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(); diff --git a/benchmarks/22-memory-backpressure.ts b/benchmarks/22-memory-backpressure.ts index 67c34ad..f396571 100644 --- a/benchmarks/22-memory-backpressure.ts +++ b/benchmarks/22-memory-backpressure.ts @@ -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') { diff --git a/benchmarks/html/04-branching.html b/benchmarks/html/04-branching.html index f26f7ad..f13404b 100644 --- a/benchmarks/html/04-branching.html +++ b/benchmarks/html/04-branching.html @@ -90,11 +90,11 @@

Raw Output

// ======================================== 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 () => { @@ -114,14 +114,14 @@

Raw Output

} 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 diff --git a/benchmarks/html/benchmarks.js b/benchmarks/html/benchmarks.js index acf05ed..40f1d6d 100644 --- a/benchmarks/html/benchmarks.js +++ b/benchmarks/html/benchmarks.js @@ -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(); diff --git a/benchmarks/profile-systematic.ts b/benchmarks/profile-systematic.ts index 864f3fe..442c7db 100644 --- a/benchmarks/profile-systematic.ts +++ b/benchmarks/profile-systematic.ts @@ -162,10 +162,10 @@ async function main() { }); await measure('broadcast() single consumer', iterations, async () => { - const { writer, broadcast } = Stream.broadcast(); + const { writer, channel } = Stream.broadcast(); 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; } })(); diff --git a/docs/COMPLETENESS-ANALYSIS.md b/docs/COMPLETENESS-ANALYSIS.md index d39181b..7d8f6bb 100644 --- a/docs/COMPLETENESS-ANALYSIS.md +++ b/docs/COMPLETENESS-ANALYSIS.md @@ -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([ diff --git a/docs/DESIGN.md b/docs/DESIGN.md index 46a86ef..8e8c1e4 100644 --- a/docs/DESIGN.md +++ b/docs/DESIGN.md @@ -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 { @@ -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) { @@ -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 | diff --git a/docs/MIGRATION-NODEJS.md b/docs/MIGRATION-NODEJS.md index c921508..33e75cb 100644 --- a/docs/MIGRATION-NODEJS.md +++ b/docs/MIGRATION-NODEJS.md @@ -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'); diff --git a/docs/MIGRATION.md b/docs/MIGRATION.md index 720519b..61b5919 100644 --- a/docs/MIGRATION.md +++ b/docs/MIGRATION.md @@ -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'); @@ -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) }); @@ -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... } ``` diff --git a/docs/TRANSFER-INTEGRATION.md b/docs/TRANSFER-INTEGRATION.md index 11b58e9..1aaadea 100644 --- a/docs/TRANSFER-INTEGRATION.md +++ b/docs/TRANSFER-INTEGRATION.md @@ -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 | @@ -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` representing one consumer's view of the shared/broadcast data. +Each call to `share.pull()` or `channel.push()` returns an `AsyncIterable` representing one consumer's view of the shared/broadcast data. **Transfer behavior:** @@ -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 }); diff --git a/index.bs b/index.bs index fdb075c..0948305 100644 --- a/index.bs +++ b/index.bs @@ -326,7 +326,7 @@ dictionary PushStreamResult { dictionary BroadcastResult { required Writer writer; - required Broadcast broadcast; + required Broadcast channel; }; @@ -999,8 +999,8 @@ The broadcast(options) method creates a [=broadca
  • Let |backpressure| be |options|["{{BroadcastOptions/backpressure}}"].
  • Create a shared circular buffer with [=byte budget=] |budget|.
  • Let |writer| be a new {{Writer}} backed by the shared buffer with |backpressure| policy. -
  • Let |broadcast| be a new {{Broadcast}} object backed by the shared buffer. -
  • Return «[ "writer" → |writer|, "broadcast" → |broadcast| ]». +
  • Let |channel| be a new {{Broadcast}} object backed by the shared buffer. +
  • Return «[ "writer" → |writer|, "channel" → |channel| ]». diff --git a/samples/03-branching.ts b/samples/03-branching.ts index a769dea..554c72f 100644 --- a/samples/03-branching.ts +++ b/samples/03-branching.ts @@ -19,11 +19,11 @@ async function main() { // Create a broadcast - writer pushes to all consumers { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Create multiple consumers - const consumer1 = broadcast.push(); - const consumer2 = broadcast.push(); + const consumer1 = channel.push(); + const consumer2 = channel.push(); // Write data - both consumers will see it (async () => { @@ -43,13 +43,13 @@ async function main() { // Multiple consumers with different transforms { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Consumer without transform - const raw = broadcast.push(); + const raw = channel.push(); // Consumer with transform (transforms are applied lazily when consumer pulls) - const transformed = broadcast.push(uppercaseTransform()); + const transformed = channel.push(uppercaseTransform()); (async () => { await writer.write('Hello World'); @@ -155,10 +155,10 @@ async function main() { yield 'broadcast data 2'; } - const { broadcast } = Broadcast.from(Stream.from(dataSource())); + const { channel } = Broadcast.from(Stream.from(dataSource())); - const consumer1 = broadcast.push(); - const consumer2 = broadcast.push(); + const consumer1 = channel.push(); + const consumer2 = channel.push(); await new Promise(resolve => setTimeout(resolve, 50)); const results = await Promise.all([ @@ -240,12 +240,12 @@ async function main() { // Broadcast with 'strict' - uses writeSync to check buffer space { - const { writer, broadcast } = Stream.broadcast({ + const { writer, channel } = Stream.broadcast({ highWaterMark: 3, backpressure: 'strict' }); - const consumer = broadcast.push(); + const consumer = channel.push(); console.log("'strict' policy with writeSync:"); diff --git a/samples/06-buffer-configuration.ts b/samples/06-buffer-configuration.ts index f968c74..79bab90 100644 --- a/samples/06-buffer-configuration.ts +++ b/samples/06-buffer-configuration.ts @@ -215,12 +215,12 @@ async function main() { section('Broadcast with backpressure'); { - const { writer, broadcast } = Stream.broadcast({ + const { writer, channel } = Stream.broadcast({ highWaterMark: 3, backpressure: 'strict' }); - const consumer = broadcast.push(); + const consumer = channel.push(); console.log('Broadcast highWaterMark: 3'); diff --git a/samples/09-resource-management.ts b/samples/09-resource-management.ts index ebf230a..bc49c57 100644 --- a/samples/09-resource-management.ts +++ b/samples/09-resource-management.ts @@ -106,16 +106,16 @@ async function main() { section('Broadcast with Symbol.dispose'); { - const { writer, broadcast } = Stream.broadcast(); + const { writer, channel } = Stream.broadcast(); // Check for Symbol.dispose - console.log('Broadcast has Symbol.dispose:', Symbol.dispose in broadcast); + console.log('Broadcast has Symbol.dispose:', Symbol.dispose in channel); // Create consumers - const consumer1 = broadcast.push(); - const consumer2 = broadcast.push(); + const consumer1 = channel.push(); + const consumer2 = channel.push(); - console.log('Consumer count:', broadcast.consumerCount); + console.log('Consumer count:', channel.consumerCount); // Start consuming in background const promise1 = Stream.text(consumer1).catch((e) => `Consumer1 error: ${(e as Error).message}`); @@ -125,9 +125,9 @@ async function main() { await writer.write('broadcast data'); // Cancel all consumers using Symbol.dispose (same as .cancel()) - broadcast[Symbol.dispose](); + channel[Symbol.dispose](); - // Or equivalently: broadcast.cancel(); + // Or equivalently: channel.cancel(); const [result1, result2] = await Promise.all([promise1, promise2]); console.log('Results after dispose:', { result1: result1.slice(0, 30), result2: result2.slice(0, 30) }); @@ -140,10 +140,10 @@ async function main() { let consumerCleanedUp = false; { - const { writer, broadcast } = Stream.broadcast(); - using _ = broadcast; // Will call Symbol.dispose on scope exit + const { writer, channel } = Stream.broadcast(); + using _ = channel; // Will call Symbol.dispose on scope exit - const consumer = broadcast.push(); + const consumer = channel.push(); // Start consuming const textPromise = Stream.text(consumer).catch(() => { @@ -154,7 +154,7 @@ async function main() { await writer.write('data in using block'); console.log(' Leaving using block...'); - // broadcast[Symbol.dispose]() called automatically here + // channel[Symbol.dispose]() called automatically here } await new Promise((r) => setTimeout(r, 50)); diff --git a/samples/11-real-world-patterns.ts b/samples/11-real-world-patterns.ts index 562a627..3a27aef 100644 --- a/samples/11-real-world-patterns.ts +++ b/samples/11-real-world-patterns.ts @@ -186,11 +186,11 @@ Charlie,35,Chicago`; { // Use broadcast for multi-consumer pattern - const { writer, broadcast } = Stream.broadcast(); + const { writer, channel } = Stream.broadcast(); // Create branches - const logBranch = broadcast.push(); - const processBranch = broadcast.push(); + const logBranch = channel.push(); + const processBranch = channel.push(); // Log in background const logPromise = (async () => { diff --git a/samples/12-nodejs-interop.ts b/samples/12-nodejs-interop.ts index 900a3eb..218c040 100644 --- a/samples/12-nodejs-interop.ts +++ b/samples/12-nodejs-interop.ts @@ -444,11 +444,11 @@ async function main() { section('Example 8: Tee to File While Processing'); { - const { writer, broadcast } = Stream.broadcast(); + const { writer, channel } = Stream.broadcast(); // Create branches - const fileBranch = broadcast.push(); - const processBranch = broadcast.push(); + const fileBranch = channel.push(); + const processBranch = channel.push(); // Save to file in background const savePromise = (async () => { diff --git a/samples/13-web-streams-interop.ts b/samples/13-web-streams-interop.ts index 299726a..0d22d5e 100644 --- a/samples/13-web-streams-interop.ts +++ b/samples/13-web-streams-interop.ts @@ -541,11 +541,11 @@ async function main() { section('Example 11: Tee to Web and New Stream'); { - const { writer, broadcast } = Stream.broadcast(); + const { writer, channel } = Stream.broadcast(); // Create branches - const webBranch = broadcast.push(); - const mainBranch = broadcast.push(); + const webBranch = channel.push(); + const mainBranch = channel.push(); // Convert branch to Web ReadableStream for web consumers const webReadable = toReadableStream(webBranch); diff --git a/samples/html/06-branching.html b/samples/html/06-branching.html index cc36662..7075543 100644 --- a/samples/html/06-branching.html +++ b/samples/html/06-branching.html @@ -124,11 +124,11 @@

    06 - Branching

    // Basic broadcast - multiple consumers see same data { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Create multiple consumers - const consumer1 = broadcast.push(); - const consumer2 = broadcast.push(); + const consumer1 = channel.push(); + const consumer2 = channel.push(); // Write data - both consumers will see it (async () => { @@ -149,13 +149,13 @@

    06 - Branching

    // Multiple consumers with different transforms { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Consumer without transform - const raw = broadcast.push(); + const raw = channel.push(); // Consumer with transform - const transformed = broadcast.push(uppercaseTransform()); + const transformed = channel.push(uppercaseTransform()); (async () => { await writer.write('Hello World'); @@ -228,7 +228,7 @@

    06 - Branching

    section('Dynamic consumers'); { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Start writing const writePromise = (async () => { @@ -241,11 +241,11 @@

    06 - Branching

    })(); // Early consumer - sees everything - const early = broadcast.push(); + const early = channel.push(); // Wait a bit, then add late consumer await new Promise(r => setTimeout(r, 15)); - const late = broadcast.push(); + const late = channel.push(); await writePromise; @@ -266,7 +266,7 @@

    06 - Branching

    section('Fan-out pattern'); { - const { writer, broadcast } = Stream.broadcast({ highWaterMark: 100 }); + const { writer, channel } = Stream.broadcast({ highWaterMark: 100 }); // Create a transform factory that adds a prefix to each chunk function addPrefix(prefix) { @@ -278,9 +278,9 @@

    06 - Branching

    }; } - const logProcessor = broadcast.push(addPrefix('LOG')); - const alertProcessor = broadcast.push(addPrefix('ALERT')); - const archiveProcessor = broadcast.push(addPrefix('ARCHIVE')); + const logProcessor = channel.push(addPrefix('LOG')); + const alertProcessor = channel.push(addPrefix('ALERT')); + const archiveProcessor = channel.push(addPrefix('ARCHIVE')); (async () => { await writer.write('system event'); diff --git a/src/broadcast.js b/src/broadcast.js index deac4d6..e065bd4 100644 --- a/src/broadcast.js +++ b/src/broadcast.js @@ -772,7 +772,7 @@ function broadcast(options) { }, { once: true }); } } - return { writer: writer, broadcast: broadcastImpl }; + return { writer: writer, channel: broadcastImpl }; } /** * Check if value implements Broadcastable protocol. @@ -800,7 +800,7 @@ exports.Broadcast = { var bc = input[types_js_1.broadcastProtocol](options); // The protocol returns Broadcast, we need to create a writer // This is a simplification - in practice the protocol would return the full result - return { writer: {}, broadcast: bc }; + return { writer: {}, channel: bc }; } // Create broadcast and pump from source var result = broadcast(options); diff --git a/src/broadcast.test.ts b/src/broadcast.test.ts index fd2d2a6..0b47a05 100644 --- a/src/broadcast.test.ts +++ b/src/broadcast.test.ts @@ -27,8 +27,8 @@ async function collect(source: AsyncIterable): Promise { describe('basic usage', () => { - it('should create writer and broadcast pair [BCAST-001]', () => { - const { writer, broadcast: bc } = broadcast(); + it('should create writer and channel pair [BCAST-001]', () => { + const { writer, channel: bc } = broadcast(); assert.ok(writer); assert.ok(bc); assert.strictEqual(bc.consumerCount, 0); @@ -36,7 +36,7 @@ describe('broadcast()', () => { }); it('should allow single consumer to receive data [BCAST-002]', async () => { - const { writer, broadcast: bc } = broadcast(); + const { writer, channel: bc } = broadcast(); const consumer = bc.push(); // Write in background @@ -55,7 +55,7 @@ describe('broadcast()', () => { }); it('should allow multiple consumers to receive same data [BCAST-003]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer1 = bc.push(); const consumer2 = bc.push(); @@ -80,7 +80,7 @@ describe('broadcast()', () => { }); it('should track consumer count [BCAST-004]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); assert.strictEqual(bc.consumerCount, 0); @@ -103,7 +103,7 @@ describe('broadcast()', () => { describe('buffer management', () => { it('should respect buffer limit [BCAST-010]', () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 2 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 2 }); assert.strictEqual(writer.writeSync('chunk1'), true); assert.strictEqual(writer.writeSync('chunk2'), true); @@ -114,7 +114,7 @@ describe('broadcast()', () => { }); it('should trim buffer as consumers advance [BCAST-011]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer = bc.push(); const iter = consumer[Symbol.asyncIterator](); @@ -135,7 +135,7 @@ describe('broadcast()', () => { }); it('should use drop-oldest policy [BCAST-012]', () => { - const { writer, broadcast: bc } = broadcast({ + const { writer, channel: bc } = broadcast({ highWaterMark: 2, backpressure: 'drop-oldest', }); @@ -148,7 +148,7 @@ describe('broadcast()', () => { }); it('should use drop-newest policy [BCAST-013]', () => { - const { writer, broadcast: bc } = broadcast({ + const { writer, channel: bc } = broadcast({ highWaterMark: 2, backpressure: 'drop-newest', }); @@ -174,7 +174,7 @@ describe('broadcast()', () => { }); it('should support writev [BCAST-021]', async () => { - const { writer, broadcast: bc } = broadcast(); + const { writer, channel: bc } = broadcast(); const consumer = bc.push(); await writer.writev(['part1', 'part2', 'part3']); @@ -186,7 +186,7 @@ describe('broadcast()', () => { }); it('should propagate errors via fail [BCAST-022]', async () => { - const { writer, broadcast: bc } = broadcast(); + const { writer, channel: bc } = broadcast(); const consumer = bc.push(); const error = new Error('Test error'); @@ -200,7 +200,7 @@ describe('broadcast()', () => { describe('cancel()', () => { it('should cancel all consumers without error [BCAST-030]', async () => { - const { writer, broadcast: bc } = broadcast(); + const { writer, channel: bc } = broadcast(); const consumer = bc.push(); bc.cancel(); @@ -210,7 +210,7 @@ describe('broadcast()', () => { }); it('should cancel all consumers with error [BCAST-031]', async () => { - const { writer, broadcast: bc } = broadcast(); + const { writer, channel: bc } = broadcast(); const consumer = bc.push(); // Start iteration before cancelling @@ -225,7 +225,7 @@ describe('broadcast()', () => { }); it('should be idempotent [BCAST-032]', () => { - const { broadcast: bc } = broadcast(); + const { channel: bc } = broadcast(); bc.cancel(); bc.cancel(); // Should not throw bc.cancel(new Error('test')); // Should not throw @@ -234,7 +234,7 @@ describe('broadcast()', () => { describe('Symbol.dispose', () => { it('should cancel on dispose [BCAST-040]', async () => { - const { broadcast: bc } = broadcast(); + const { channel: bc } = broadcast(); const consumer = bc.push(); bc[Symbol.dispose](); @@ -249,7 +249,7 @@ describe('broadcast()', () => { const controller = new AbortController(); controller.abort(); - const { broadcast: bc } = broadcast({ signal: controller.signal }); + const { channel: bc } = broadcast({ signal: controller.signal }); const consumer = bc.push(); const chunks = await collect(consumer); @@ -258,7 +258,7 @@ describe('broadcast()', () => { it('should cancel on signal abort [BCAST-051]', async () => { const controller = new AbortController(); - const { writer, broadcast: bc } = broadcast({ signal: controller.signal }); + const { writer, channel: bc } = broadcast({ signal: controller.signal }); const consumer = bc.push(); writer.writeSync('chunk1'); @@ -274,7 +274,7 @@ describe('broadcast()', () => { describe('late subscribers', () => { it('should allow late subscribers to receive new data [BCAST-005]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); // Write before any subscribers writer.writeSync('chunk1'); @@ -297,7 +297,7 @@ describe('broadcast()', () => { describe('Broadcast.from()', () => { it('should create broadcast from async iterable [BCAST-060]', async () => { const source = from(['chunk1', 'chunk2']); - const { broadcast: bc } = Broadcast.from(source); + const { channel: bc } = Broadcast.from(source); const consumer = bc.push(); @@ -310,7 +310,7 @@ describe('Broadcast.from()', () => { it('should create broadcast from sync iterable [BCAST-061]', async () => { const source = fromSync(['chunk1', 'chunk2']); - const { broadcast: bc } = Broadcast.from(source); + const { channel: bc } = Broadcast.from(source); const consumer = bc.push(); @@ -342,7 +342,7 @@ describe('broadcast() with transforms', () => { }; it('should apply single transform to consumer [BCAST-070]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer = bc.push(uppercase); @@ -359,7 +359,7 @@ describe('broadcast() with transforms', () => { }); it('should apply multiple transforms in order [BCAST-071]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer = bc.push(uppercase, prefix); @@ -373,7 +373,7 @@ describe('broadcast() with transforms', () => { }); it('should allow different transforms per consumer [BCAST-072]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer1 = bc.push(uppercase); const consumer2 = bc.push(prefix); @@ -440,7 +440,7 @@ describe('broadcast writer drainable protocol', () => { // BCAST-084: ondrain returns pending Promise when desiredSize === 0 it('should return pending Promise when desiredSize === 0 [BCAST-084]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 1 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 1 }); // Create a consumer to read data const consumer = bc.push(); @@ -512,7 +512,7 @@ describe('broadcast writer drainable protocol', () => { // BCAST-087: Multiple drain waiters all resolve together it('should resolve multiple drain waiters together [BCAST-087]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 1 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 1 }); // Create a consumer to read data const consumer = bc.push(); @@ -547,7 +547,7 @@ describe('broadcast writer drainable protocol', () => { describe('broadcast write signal cancellation', () => { it('should reject blocked write when signal fires [BCAST-090]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 1 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 1 }); // Create a consumer so writes go through const consumer = bc.push(); @@ -586,7 +586,7 @@ describe('broadcast write signal cancellation', () => { }); it('should clean up signal listener on normal write completion [BCAST-092]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 1 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 1 }); const consumer = bc.push(); const iter = consumer[Symbol.asyncIterator](); @@ -621,7 +621,7 @@ describe('Broadcast.from() cancel bug fix', () => { })()); const controller = new AbortController(); - const { broadcast: bc } = Broadcast.from(source, { + const { channel: bc } = Broadcast.from(source, { highWaterMark: 1, signal: controller.signal, }); @@ -691,7 +691,7 @@ describe('Broadcast.from() cancel bug fix', () => { })()); const controller = new AbortController(); - const { broadcast: bc } = Broadcast.from(source, { + const { channel: bc } = Broadcast.from(source, { highWaterMark: 100, // Large buffer so writes don't block signal: controller.signal, }); @@ -713,7 +713,7 @@ describe('Broadcast.from() cancel bug fix', () => { describe('broadcast _abort consumer cleanup', () => { it('should detach consumers and clear the set on fail [BCAST-110]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); const consumer1 = bc.push(); const consumer2 = bc.push(); @@ -731,7 +731,7 @@ describe('broadcast _abort consumer cleanup', () => { }); it('should detach consumers even without pending reads [BCAST-111]', async () => { - const { writer, broadcast: bc } = broadcast({ highWaterMark: 100 }); + const { writer, channel: bc } = broadcast({ highWaterMark: 100 }); // Create consumers but don't start reading — no pending reads bc.push(); diff --git a/src/broadcast.ts b/src/broadcast.ts index 92c47ff..994d2a4 100644 --- a/src/broadcast.ts +++ b/src/broadcast.ts @@ -788,7 +788,7 @@ export function broadcast(options?: BroadcastOptions): BroadcastResult { } } - return { writer, broadcast: broadcastImpl }; + return { writer, channel: broadcastImpl }; } /** @@ -822,7 +822,7 @@ export const Broadcast = { const bc = input[broadcastProtocol](options); // The protocol returns Broadcast, we need to create a writer // This is a simplification - in practice the protocol would return the full result - return { writer: {} as Writer, broadcast: bc }; + return { writer: {} as Writer, channel: bc }; } // Create broadcast and pump from source diff --git a/src/types.ts b/src/types.ts index f6449aa..46d7b6c 100644 --- a/src/types.ts +++ b/src/types.ts @@ -634,11 +634,11 @@ export interface Broadcast { } /** - * Result of Stream.broadcast() - writer + broadcast pair. + * Result of Stream.broadcast() - writer + channel pair. */ export interface BroadcastResult { writer: Writer; - broadcast: Broadcast; + channel: Broadcast; } // =============================================================================