Skip to content
File

Blob: src/workerd/api/tests/streams-tee-edge-cases-test.js

javascript386 lines
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 
13import { 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)
17export 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)
64export 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)
118export 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)
171export 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)
219export 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)
256export 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
302export 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)
344export 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};