File
Blob: src/workerd/api/tests/streams-tee-edge-cases-test.js
| 1 | // Copyright (c) 2026 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | // Tests for tee() edge cases with asymmetric consumption patterns. |
| 6 | // These tests focus on scenarios where tee branches are consumed at |
| 7 | // different rates or only partially consumed. |
| 8 | // |
| 9 | // Test inspirations: |
| 10 | // - Bun: test/js/web/streams/streams.test.js (tee for default and direct streams) |
| 11 | // - Deno: tests/unit/streams_test.ts (tee tests) |
| 12 | |
| 13 | import { strictEqual, ok, deepStrictEqual } from 'node:assert'; |
| 14 | |
| 15 | // Test consuming only one branch of a tee completely |
| 16 | // Inspired by: Bun test/js/web/streams/streams.test.js (tee tests) |
| 17 | export const teeConsumeOneBranchFully = { |
| 18 | async test() { |
| 19 | let pullCount = 0; |
| 20 | const rs = new ReadableStream({ |
| 21 | pull(controller) { |
| 22 | pullCount++; |
| 23 | if (pullCount <= 5) { |
| 24 | controller.enqueue(pullCount); |
| 25 | } else { |
| 26 | controller.close(); |
| 27 | } |
| 28 | }, |
| 29 | }); |
| 30 | |
| 31 | const [branch1, branch2] = rs.tee(); |
| 32 | |
| 33 | // Only consume branch1 fully |
| 34 | const reader1 = branch1.getReader(); |
| 35 | const values = []; |
| 36 | |
| 37 | while (true) { |
| 38 | const { value, done } = await reader1.read(); |
| 39 | if (done) break; |
| 40 | values.push(value); |
| 41 | } |
| 42 | |
| 43 | deepStrictEqual(values, [1, 2, 3, 4, 5]); |
| 44 | |
| 45 | // branch2 should still be readable (though branch1 consumed the data) |
| 46 | ok(!branch2.locked); |
| 47 | |
| 48 | // Now consume branch2 |
| 49 | const reader2 = branch2.getReader(); |
| 50 | const values2 = []; |
| 51 | |
| 52 | while (true) { |
| 53 | const { value, done } = await reader2.read(); |
| 54 | if (done) break; |
| 55 | values2.push(value); |
| 56 | } |
| 57 | |
| 58 | deepStrictEqual(values2, [1, 2, 3, 4, 5]); |
| 59 | }, |
| 60 | }; |
| 61 | |
| 62 | // Test tee with different read rates on branches |
| 63 | // Inspired by: Deno tests/unit/streams_test.ts (async stream tests) |
| 64 | export const teeDifferentReadRates = { |
| 65 | async test() { |
| 66 | let counter = 0; |
| 67 | const rs = new ReadableStream({ |
| 68 | pull(controller) { |
| 69 | counter++; |
| 70 | if (counter <= 10) { |
| 71 | controller.enqueue(counter); |
| 72 | } else { |
| 73 | controller.close(); |
| 74 | } |
| 75 | }, |
| 76 | }); |
| 77 | |
| 78 | const [branch1, branch2] = rs.tee(); |
| 79 | |
| 80 | const reader1 = branch1.getReader(); |
| 81 | const reader2 = branch2.getReader(); |
| 82 | |
| 83 | const results1 = []; |
| 84 | const results2 = []; |
| 85 | |
| 86 | // Read from branch1 fast, branch2 slow |
| 87 | for (let i = 0; i < 10; i++) { |
| 88 | // Read 2 from branch1 |
| 89 | results1.push((await reader1.read()).value); |
| 90 | if (i % 2 === 1) { |
| 91 | // Read 1 from branch2 every other iteration |
| 92 | results2.push((await reader2.read()).value); |
| 93 | } |
| 94 | } |
| 95 | |
| 96 | // Finish reading branch2 |
| 97 | while (true) { |
| 98 | const { value, done } = await reader2.read(); |
| 99 | if (done) break; |
| 100 | results2.push(value); |
| 101 | } |
| 102 | |
| 103 | // Finish reading branch1 |
| 104 | while (true) { |
| 105 | const { value, done } = await reader1.read(); |
| 106 | if (done) break; |
| 107 | results1.push(value); |
| 108 | } |
| 109 | |
| 110 | // Both branches should have received all values |
| 111 | deepStrictEqual(results1, [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]); |
| 112 | deepStrictEqual(results2, [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]); |
| 113 | }, |
| 114 | }; |
| 115 | |
| 116 | // Test canceling the slower branch mid-stream |
| 117 | // Inspired by: Bun test/js/web/streams/streams.test.js (cancel tests) |
| 118 | export const teeCancelSlowBranch = { |
| 119 | async test() { |
| 120 | let counter = 0; |
| 121 | let sourceCancelled = false; |
| 122 | |
| 123 | const rs = new ReadableStream({ |
| 124 | pull(controller) { |
| 125 | counter++; |
| 126 | if (counter <= 20) { |
| 127 | controller.enqueue(counter); |
| 128 | } else { |
| 129 | controller.close(); |
| 130 | } |
| 131 | }, |
| 132 | cancel() { |
| 133 | sourceCancelled = true; |
| 134 | }, |
| 135 | }); |
| 136 | |
| 137 | const [branch1, branch2] = rs.tee(); |
| 138 | |
| 139 | const reader1 = branch1.getReader(); |
| 140 | const reader2 = branch2.getReader(); |
| 141 | |
| 142 | // Read 5 from both branches |
| 143 | for (let i = 0; i < 5; i++) { |
| 144 | await reader1.read(); |
| 145 | await reader2.read(); |
| 146 | } |
| 147 | |
| 148 | // Cancel branch2 (the "slow" one) |
| 149 | await reader2.cancel('No longer needed'); |
| 150 | |
| 151 | // Source should NOT be cancelled yet (branch1 still active) |
| 152 | ok(!sourceCancelled); |
| 153 | |
| 154 | // Continue reading branch1 to completion |
| 155 | const remaining = []; |
| 156 | while (true) { |
| 157 | const { value, done } = await reader1.read(); |
| 158 | if (done) break; |
| 159 | remaining.push(value); |
| 160 | } |
| 161 | |
| 162 | // Branch1 should have received remaining values |
| 163 | strictEqual(remaining.length, 15); |
| 164 | strictEqual(remaining[0], 6); |
| 165 | strictEqual(remaining[14], 20); |
| 166 | }, |
| 167 | }; |
| 168 | |
| 169 | // Test tee with byte stream using default readers |
| 170 | // Inspired by: workerd streams-js-test.js (byte stream tee tests) |
| 171 | export const teeByteStreamDefaultReaders = { |
| 172 | async test() { |
| 173 | const data = new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8]); |
| 174 | let offset = 0; |
| 175 | |
| 176 | const rs = new ReadableStream({ |
| 177 | type: 'bytes', |
| 178 | pull(controller) { |
| 179 | if (offset < data.length) { |
| 180 | const chunk = data.slice(offset, offset + 2); |
| 181 | offset += 2; |
| 182 | controller.enqueue(chunk); |
| 183 | } else { |
| 184 | controller.close(); |
| 185 | } |
| 186 | }, |
| 187 | }); |
| 188 | |
| 189 | const [branch1, branch2] = rs.tee(); |
| 190 | |
| 191 | // Use default readers (not BYOB) |
| 192 | const reader1 = branch1.getReader(); |
| 193 | const reader2 = branch2.getReader(); |
| 194 | |
| 195 | const bytes1 = []; |
| 196 | const bytes2 = []; |
| 197 | |
| 198 | // Read all from both branches, collecting individual bytes |
| 199 | while (true) { |
| 200 | const { value, done } = await reader1.read(); |
| 201 | if (done) break; |
| 202 | for (const b of value) bytes1.push(b); |
| 203 | } |
| 204 | |
| 205 | while (true) { |
| 206 | const { value, done } = await reader2.read(); |
| 207 | if (done) break; |
| 208 | for (const b of value) bytes2.push(b); |
| 209 | } |
| 210 | |
| 211 | // Both branches should have received all 8 bytes with same values |
| 212 | deepStrictEqual(bytes1, [1, 2, 3, 4, 5, 6, 7, 8]); |
| 213 | deepStrictEqual(bytes2, [1, 2, 3, 4, 5, 6, 7, 8]); |
| 214 | }, |
| 215 | }; |
| 216 | |
| 217 | // Test tee with byte stream using mixed reader types |
| 218 | // Inspired by: workerd streams-js-test.js (BYOB tee tests) |
| 219 | export const teeByteStreamMixedReaders = { |
| 220 | async test() { |
| 221 | const enc = new TextEncoder(); |
| 222 | const dec = new TextDecoder(); |
| 223 | |
| 224 | let controller; |
| 225 | const rs = new ReadableStream({ |
| 226 | type: 'bytes', |
| 227 | start(c) { |
| 228 | controller = c; |
| 229 | }, |
| 230 | }); |
| 231 | |
| 232 | const [branch1, branch2] = rs.tee(); |
| 233 | |
| 234 | // Use BYOB reader on branch1, default reader on branch2 |
| 235 | const reader1 = branch1.getReader({ mode: 'byob' }); |
| 236 | const reader2 = branch2.getReader(); |
| 237 | |
| 238 | // Start reads |
| 239 | const read1Promise = reader1.read(new Uint8Array(5)); |
| 240 | const read2Promise = reader2.read(); |
| 241 | |
| 242 | // Enqueue data |
| 243 | controller.enqueue(enc.encode('hello')); |
| 244 | controller.close(); |
| 245 | |
| 246 | const [result1, result2] = await Promise.all([read1Promise, read2Promise]); |
| 247 | |
| 248 | // Both should receive the data |
| 249 | strictEqual(dec.decode(result1.value), 'hello'); |
| 250 | strictEqual(dec.decode(result2.value), 'hello'); |
| 251 | }, |
| 252 | }; |
| 253 | |
| 254 | // Test tee with large number of chunks |
| 255 | // Inspired by: Deno tests/unit/streams_test.ts (large stream tests) |
| 256 | export const teeLargeChunkCount = { |
| 257 | async test() { |
| 258 | const CHUNK_COUNT = 1000; |
| 259 | let counter = 0; |
| 260 | |
| 261 | const rs = new ReadableStream({ |
| 262 | pull(controller) { |
| 263 | if (counter < CHUNK_COUNT) { |
| 264 | controller.enqueue(counter++); |
| 265 | } else { |
| 266 | controller.close(); |
| 267 | } |
| 268 | }, |
| 269 | }); |
| 270 | |
| 271 | const [branch1, branch2] = rs.tee(); |
| 272 | |
| 273 | // Read both branches in parallel |
| 274 | async function consumeBranch(branch) { |
| 275 | const reader = branch.getReader(); |
| 276 | let count = 0; |
| 277 | let sum = 0; |
| 278 | while (true) { |
| 279 | const { value, done } = await reader.read(); |
| 280 | if (done) break; |
| 281 | count++; |
| 282 | sum += value; |
| 283 | } |
| 284 | return { count, sum }; |
| 285 | } |
| 286 | |
| 287 | const [result1, result2] = await Promise.all([ |
| 288 | consumeBranch(branch1), |
| 289 | consumeBranch(branch2), |
| 290 | ]); |
| 291 | |
| 292 | strictEqual(result1.count, CHUNK_COUNT); |
| 293 | strictEqual(result2.count, CHUNK_COUNT); |
| 294 | strictEqual(result1.sum, result2.sum); |
| 295 | // Sum of 0 to 999 = 999 * 1000 / 2 = 499500 |
| 296 | strictEqual(result1.sum, 499500); |
| 297 | }, |
| 298 | }; |
| 299 | |
| 300 | // Test tee after partial read from original stream |
| 301 | // Inspired by: Bun test/js/web/streams/streams.test.js |
| 302 | export const teeAfterPartialRead = { |
| 303 | async test() { |
| 304 | let counter = 0; |
| 305 | const rs = new ReadableStream({ |
| 306 | pull(controller) { |
| 307 | counter++; |
| 308 | if (counter <= 10) { |
| 309 | controller.enqueue(counter); |
| 310 | } else { |
| 311 | controller.close(); |
| 312 | } |
| 313 | }, |
| 314 | }); |
| 315 | |
| 316 | // Read some values before tee |
| 317 | const originalReader = rs.getReader(); |
| 318 | const firstValue = await originalReader.read(); |
| 319 | strictEqual(firstValue.value, 1); |
| 320 | |
| 321 | // Release lock and tee |
| 322 | originalReader.releaseLock(); |
| 323 | |
| 324 | const [branch1, branch2] = rs.tee(); |
| 325 | |
| 326 | // Both branches should start from where original left off |
| 327 | const reader1 = branch1.getReader(); |
| 328 | const reader2 = branch2.getReader(); |
| 329 | |
| 330 | const value1 = await reader1.read(); |
| 331 | const value2 = await reader2.read(); |
| 332 | |
| 333 | // Both should get value 2 (first value after the pre-tee read) |
| 334 | strictEqual(value1.value, 2); |
| 335 | strictEqual(value2.value, 2); |
| 336 | |
| 337 | reader1.releaseLock(); |
| 338 | reader2.releaseLock(); |
| 339 | }, |
| 340 | }; |
| 341 | |
| 342 | // Test that cancel reason is passed through tee |
| 343 | // Inspired by: Deno tests/unit/streams_test.ts (cancel reason tests) |
| 344 | export const teeCancelReason = { |
| 345 | async test() { |
| 346 | let receivedReason = null; |
| 347 | let cancelCalled = false; |
| 348 | |
| 349 | const rs = new ReadableStream({ |
| 350 | pull(controller) { |
| 351 | controller.enqueue('data'); |
| 352 | }, |
| 353 | cancel(reason) { |
| 354 | cancelCalled = true; |
| 355 | receivedReason = reason; |
| 356 | }, |
| 357 | }); |
| 358 | |
| 359 | const [branch1, branch2] = rs.tee(); |
| 360 | |
| 361 | const reader1 = branch1.getReader(); |
| 362 | const reader2 = branch2.getReader(); |
| 363 | |
| 364 | // Read one value from each |
| 365 | await reader1.read(); |
| 366 | await reader2.read(); |
| 367 | |
| 368 | // Cancel both with specific reasons |
| 369 | await reader1.cancel('Reason from branch 1'); |
| 370 | await reader2.cancel('Reason from branch 2'); |
| 371 | |
| 372 | // The source cancel should be called when both branches are cancelled |
| 373 | ok(cancelCalled, 'Source cancel should be called'); |
| 374 | // The reason format may vary by implementation - it could be: |
| 375 | // - An array of reasons |
| 376 | // - The first reason |
| 377 | // - The second reason |
| 378 | // - A composite reason |
| 379 | // Just verify cancel was called with some reason |
| 380 | ok( |
| 381 | receivedReason !== null && receivedReason !== undefined, |
| 382 | 'Cancel reason should be provided' |
| 383 | ); |
| 384 | }, |
| 385 | }; |