Skip to content
File

Blob: src/workerd/api/tests/streams-r2-patterns-test.js

javascript502 lines
1// Copyright (c) 2025 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 
5import { strictEqual, ok } from 'node:assert';
6 
7// Test BYOB readAtLeast with automatic atLeast handling
8export const byobReadAtLeastAutomatic = {
9 async test() {
10 const enc = new TextEncoder();
11 const dec = new TextDecoder();
12 const chunks = ['hello', 'there'];
13 const rs = new ReadableStream({
14 type: 'bytes',
15 pull(c) {
16 // When using enqueue, the stream impl will take care of properly handling the
17 // at least requirement...
18 c.enqueue(enc.encode(chunks.shift()));
19 if (chunks.length === 0) c.close();
20 },
21 });
22 
23 const reader = rs.getReader({ mode: 'byob' });
24 
25 const res = await reader.readAtLeast(100, new Uint8Array(100));
26 
27 strictEqual(dec.decode(res.value), 'hellothere');
28 },
29};
30 
31// Test BYOB readAtLeast with manual atLeast handling
32export const byobReadAtLeastManual = {
33 async test() {
34 const enc = new TextEncoder();
35 const dec = new TextDecoder();
36 const chunks = ['hello', 'there'];
37 const expectedAtLeasts = [100, 95];
38 const rs = new ReadableStream({
39 type: 'bytes',
40 pull(c) {
41 if (chunks.length === 0) {
42 c.close();
43 c.byobRequest.respond(0);
44 } else {
45 // The respond() can partially fulfill the minRead requirement over
46 // multiple calls to pull.
47 strictEqual(c.byobRequest.atLeast, expectedAtLeasts.shift());
48 
49 enc.encodeInto(chunks.shift(), c.byobRequest.view);
50 c.byobRequest.respond(5);
51 }
52 },
53 });
54 
55 const reader = rs.getReader({ mode: 'byob' });
56 
57 const res = await reader.readAtLeast(100, new Uint8Array(100));
58 
59 strictEqual(dec.decode(res.value), 'hellothere');
60 },
61};
62 
63// Test IdentityTransformStream with readAtLeast incremental writes
64export const identityTransformReadAtLeast = {
65 async test() {
66 const { readable, writable } = new IdentityTransformStream();
67 
68 const reader = readable.getReader({ mode: 'byob' });
69 const writer = writable.getWriter();
70 
71 // There's been a latent bug in IdentityTransformStream ever since
72 // readAtLeast was introduced that caused it to mishandle the atLeast
73 // calculation when individual writes were < atLeast.
74 
75 for (let n = 0; n < 8; n++) {
76 writer.write(new Uint8Array(1));
77 }
78 writer.write(new Uint8Array([0x1]));
79 writer.write(new Uint8Array([0x2]));
80 writer.write(new Uint8Array([0x3]));
81 writer.write(new Uint8Array([0x4]));
82 
83 const res = await reader.readAtLeast(8, new Uint8Array(8));
84 strictEqual(res.value.byteLength, 8);
85 
86 const res2 = await reader.readAtLeast(2, new Uint8Array(4));
87 const res3 = await reader.readAtLeast(2, new Uint8Array(4));
88 
89 strictEqual(res2.value.byteLength, 2);
90 strictEqual(res2.value[0], 0x1);
91 strictEqual(res2.value[1], 0x2);
92 
93 strictEqual(res3.value.byteLength, 2);
94 strictEqual(res3.value[0], 0x3);
95 strictEqual(res3.value[1], 0x4);
96 },
97};
98 
99// Test FixedLengthStream with readAtLeast incremental writes
100export const fixedLengthStreamReadAtLeast = {
101 async test() {
102 const { readable, writable } = new FixedLengthStream(12);
103 
104 const reader = readable.getReader({ mode: 'byob' });
105 const writer = writable.getWriter();
106 
107 // There's been a latent bug in IdentityTransformStream ever since
108 // readAtLeast was introduced that caused it to mishandle the atLeast
109 // calculation when individual writes were < atLeast.
110 
111 for (let n = 0; n < 8; n++) {
112 writer.write(new Uint8Array(1));
113 }
114 writer.write(new Uint8Array([0x1]));
115 writer.write(new Uint8Array([0x2]));
116 writer.write(new Uint8Array([0x3]));
117 writer.write(new Uint8Array([0x4]));
118 
119 const res = await reader.readAtLeast(8, new Uint8Array(8));
120 strictEqual(res.value.byteLength, 8);
121 
122 const res2 = await reader.readAtLeast(2, new Uint8Array(4));
123 const res3 = await reader.readAtLeast(2, new Uint8Array(4));
124 
125 strictEqual(res2.value.byteLength, 2);
126 strictEqual(res2.value[0], 0x1);
127 strictEqual(res2.value[1], 0x2);
128 
129 strictEqual(res3.value.byteLength, 2);
130 strictEqual(res3.value[0], 0x3);
131 strictEqual(res3.value[1], 0x4);
132 },
133};
134 
135// Test BYOB stream tee closed on start with waitUntil
136// Tests that a teed BYOB stream that closes immediately after enqueuing
137// still works correctly when one branch is consumed via waitUntil
138export const closedByobTeeOnStart = {
139 async test(ctrl, env, ctx) {
140 const enc = new TextEncoder();
141 const dec = new TextDecoder();
142 
143 async function consume(rs) {
144 const reader = rs.getReader({ mode: 'byob' });
145 let result = '';
146 for (;;) {
147 const res = await reader.readAtLeast(10, new Uint8Array(10));
148 if (res.done) break;
149 result += dec.decode(res.value, { stream: true });
150 }
151 result += dec.decode();
152 if (result !== 'hello') throw new Error('Incorrect result in branch');
153 return result;
154 }
155 
156 const rs = new ReadableStream({
157 type: 'bytes',
158 start(c) {
159 c.enqueue(enc.encode('hello'));
160 c.close();
161 },
162 });
163 
164 const [b1, b2] = rs.tee();
165 
166 const branch2Promise = consume(b2);
167 ctx.waitUntil(branch2Promise);
168 
169 const result1 = await consume(b1);
170 strictEqual(result1, 'hello');
171 
172 const result2 = await branch2Promise;
173 strictEqual(result2, 'hello');
174 },
175};
176 
177// Test IdentityTransformStream properly handles readAtLeast
178export const identityTransformStreamReadAtLeast = {
179 async test() {
180 const { readable, writable } = new IdentityTransformStream();
181 
182 const reader = readable.getReader({ mode: 'byob' });
183 const writer = writable.getWriter();
184 
185 const expectedReads = [100, 100, 1, 0];
186 
187 async function consume(reader) {
188 const res = await reader.readAtLeast(100, new Uint8Array(100));
189 if (!res.done) {
190 strictEqual(res.value.byteLength, expectedReads.shift());
191 return consume(reader);
192 }
193 }
194 
195 await Promise.all([
196 consume(reader),
197 writer.write(new Uint8Array(100)),
198 writer.write(new Uint8Array(1)),
199 writer.write(new Uint8Array(100)),
200 writer.close(),
201 ]);
202 },
203};
204 
205// Test BYOB readAtLeast partially filled
206export const partiallyFilledByobAtLeast = {
207 async test() {
208 const { readable, writable } = new IdentityTransformStream();
209 const reader = readable.getReader({ mode: 'byob' });
210 const rs = new ReadableStream({
211 type: 'bytes',
212 async pull(controller) {
213 const chunk = await reader.readAtLeast(100, new Uint8Array(100));
214 if (!chunk.done) {
215 controller.enqueue(chunk.value);
216 } else {
217 controller.close();
218 }
219 },
220 });
221 
222 async function consume(readable) {
223 let ab = new ArrayBuffer(102);
224 const dec = new TextDecoder();
225 let ret = '';
226 const reader = readable.getReader({ mode: 'byob' });
227 for (;;) {
228 const read = await reader.readAtLeast(102, new Uint8Array(ab));
229 if (!read.done) {
230 ret += dec.decode(read.value);
231 ab = read.value.buffer;
232 continue;
233 } else {
234 break;
235 }
236 }
237 strictEqual(ret, 'hello'.repeat(1000));
238 return ret.length;
239 }
240 
241 const p = consume(rs);
242 
243 const enc = new TextEncoder();
244 const writer = writable.getWriter();
245 writer.write(enc.encode('hello'.repeat(1000)));
246 writer.close();
247 
248 strictEqual(await p, 5000);
249 },
250};
251 
252// Test BYOB readAtLeast with tee
253export const byobReadAtLeastTee = {
254 async test() {
255 const enc = new TextEncoder();
256 const dec = new TextDecoder();
257 const chunks = ['hello', 'there'];
258 const expectedAtLeasts = [100, 95];
259 const rs = new ReadableStream({
260 type: 'bytes',
261 pull(c) {
262 if (chunks.length === 0) {
263 c.close();
264 c.byobRequest.respond(0);
265 } else {
266 strictEqual(c.byobRequest.atLeast, expectedAtLeasts.shift());
267 enc.encodeInto(chunks.shift(), c.byobRequest.view);
268 c.byobRequest.respond(5);
269 }
270 },
271 });
272 
273 const [branch1, branchB] = rs.tee();
274 const [branch2, branch3] = branchB.tee();
275 
276 const reader = branch1.getReader({ mode: 'byob' });
277 const reader2 = branch2.getReader({ mode: 'byob' });
278 const reader3 = branch3.getReader({ mode: 'byob' });
279 
280 const p1 = reader.readAtLeast(100, new Uint8Array(100));
281 const p2 = reader2.readAtLeast(5, new Uint8Array(100));
282 const p3 = reader3.readAtLeast(3, new Uint8Array(3));
283 
284 const res = await Promise.all([p1, p2, p3]);
285 
286 strictEqual(dec.decode(res[0].value), 'hellothere');
287 strictEqual(dec.decode(res[1].value), 'hello');
288 strictEqual(dec.decode(res[2].value), 'hel');
289 
290 const res2 = await reader2.readAtLeast(5, new Uint8Array(100));
291 strictEqual(dec.decode(res2.value), 'there');
292 
293 const res3 = await reader3.readAtLeast(4, new Uint8Array(4));
294 strictEqual(dec.decode(res3.value), 'loth');
295 
296 const res4 = await reader.readAtLeast(100, new Uint8Array(100));
297 strictEqual(res4.done, true);
298 
299 const res5 = await reader2.readAtLeast(5, new Uint8Array(100));
300 strictEqual(res5.done, true);
301 
302 const res6 = await reader3.readAtLeast(4, new Uint8Array(4));
303 strictEqual(dec.decode(res6.value), 'ere');
304 
305 const res7 = await reader2.readAtLeast(5, new Uint8Array(100));
306 strictEqual(res7.done, true);
307 },
308};
309 
310// Test BYOB readAtLeast with tee complex variant 1
311export const byobReadAtLeastTeeComplex1 = {
312 async test() {
313 const enc = new TextEncoder();
314 const dec = new TextDecoder();
315 const chunks = ['helloth', 'ere'];
316 let previousByobRequest;
317 const rs = new ReadableStream({
318 type: 'bytes',
319 pull(c) {
320 const req = c.byobRequest;
321 if (chunks.length === 0) {
322 c.close();
323 req.respond(0);
324 } else {
325 ok(!(req === previousByobRequest));
326 const chunk = chunks.shift();
327 enc.encodeInto(chunk, req.view);
328 req.respond(chunk.length);
329 }
330 },
331 });
332 
333 const [branch1, branchB] = rs.tee();
334 const [branch2, branch3] = branchB.tee();
335 
336 const reader1 = branch1.getReader({ mode: 'byob' });
337 const reader2 = branch2.getReader({ mode: 'byob' });
338 const reader3 = branch3.getReader({ mode: 'byob' });
339 
340 const res1 = await reader1.readAtLeast(5, new Uint8Array(10));
341 strictEqual(dec.decode(res1.value), 'helloth');
342 const res2 = await reader2.readAtLeast(10, new Uint8Array(10));
343 strictEqual(dec.decode(res2.value), 'hellothere');
344 
345 const res3 = await reader1.readAtLeast(5, new Uint8Array(10));
346 strictEqual(dec.decode(res3.value), 'ere');
347 
348 const res4 = await reader3.readAtLeast(2, new Uint8Array(12));
349 strictEqual(dec.decode(res4.value), 'hellothere');
350 },
351};
352 
353// Test BYOB readAtLeast with tee complex variant 2
354export const byobReadAtLeastTeeComplex2 = {
355 async test() {
356 const enc = new TextEncoder();
357 const dec = new TextDecoder();
358 const chunks = ['helloth', 'ere'];
359 let previousByobRequest;
360 const rs = new ReadableStream({
361 type: 'bytes',
362 pull(c) {
363 if (chunks.length === 0) {
364 c.close();
365 c.byobRequest.respond(0);
366 } else {
367 ok(!(c.byobRequest === previousByobRequest));
368 const chunk = chunks.shift();
369 enc.encodeInto(chunk, c.byobRequest.view);
370 c.byobRequest.respond(chunk.length);
371 }
372 },
373 });
374 
375 const [branch1, branchB] = rs.tee();
376 const [branch2, branch3] = branchB.tee();
377 
378 const reader1 = branch1.getReader({ mode: 'byob' });
379 const reader2 = branch2.getReader({ mode: 'byob' });
380 const reader3 = branch3.getReader({ mode: 'byob' });
381 
382 const res1 = await reader1.readAtLeast(5, new Uint8Array(10));
383 strictEqual(dec.decode(res1.value), 'helloth');
384 const res2 = await reader2.readAtLeast(10, new Uint8Array(10));
385 strictEqual(dec.decode(res2.value), 'hellothere');
386 
387 const res3 = await reader1.readAtLeast(5, new Uint8Array(10));
388 strictEqual(dec.decode(res3.value), 'ere');
389 
390 const res4 = await reader3.readAtLeast(2, new Uint8Array(12));
391 strictEqual(dec.decode(res4.value), 'hellothere');
392 },
393};
394 
395// Test BYOB readAtLeast with tee complex variant 3 (typed arrays)
396export const byobReadAtLeastTeeComplex3 = {
397 async test() {
398 const chunks = [
399 new Uint8Array([0x01]),
400 new Uint8Array([0x02]),
401 new Uint8Array([0x03]),
402 new Uint8Array([0x04]),
403 new Uint8Array([0x05, 0x06]),
404 ];
405 
406 const rs = new ReadableStream({
407 type: 'bytes',
408 pull(c) {
409 if (chunks.length === 0) {
410 c.close();
411 c.byobRequest.respond(0);
412 } else {
413 const view = c.byobRequest.view;
414 const chunk = chunks.shift();
415 for (let n = 0; n < chunk.length; n++) {
416 view[n] = chunk[n];
417 }
418 c.byobRequest.respond(chunk.length);
419 }
420 },
421 });
422 
423 const [branch1, branch2] = rs.tee();
424 
425 const reader1 = branch1.getReader({ mode: 'byob' });
426 const reader2 = branch2.getReader({ mode: 'byob' });
427 
428 const [res1, res2, res3, res4] = await Promise.all([
429 reader1.readAtLeast(2, new Uint16Array(2)),
430 reader1.readAtLeast(2, new Uint8Array(2)),
431 reader2.readAtLeast(2, new Uint8Array(2)),
432 reader2.readAtLeast(1, new Uint32Array(1)),
433 ]);
434 
435 strictEqual(res1.value instanceof Uint16Array, true);
436 strictEqual(res2.value instanceof Uint8Array, true);
437 strictEqual(res1.value[0], 0x0201);
438 strictEqual(res1.value[1], 0x0403);
439 strictEqual(res2.value[0], 0x05);
440 strictEqual(res2.value[1], 0x06);
441 
442 strictEqual(res3.value instanceof Uint8Array, true);
443 strictEqual(res4.value instanceof Uint32Array, true);
444 strictEqual(res3.value[0], 0x1);
445 strictEqual(res3.value[1], 0x2);
446 strictEqual(res4.value[0], 0x06050403);
447 },
448};
449 
450export const requestCloneByob = {
451 async test() {
452 const enc = new TextEncoder();
453 const dec = new TextDecoder();
454 const chunks = [
455 enc.encode('hello'),
456 enc.encode('there'),
457 enc.encode('!!!!!'),
458 ];
459 const rs = new ReadableStream({
460 type: 'bytes',
461 pull(c) {
462 c.enqueue(chunks.shift());
463 if (chunks.length === 0) c.close();
464 },
465 });
466 
467 const newRequest = new Request('http://example.org', {
468 method: 'POST',
469 body: rs,
470 });
471 const reader = newRequest.clone().body.getReader({ mode: 'byob' });
472 
473 strictEqual(
474 dec.decode((await reader.readAtLeast(10, new Uint8Array(10))).value),
475 'hellothere'
476 );
477 },
478};
479 
480export const textDecoderStreamRequest = {
481 async test() {
482 const enc = new TextEncoder();
483 const rs = new ReadableStream({
484 type: 'bytes',
485 start(c) {
486 c.enqueue(enc.encode('hello'));
487 c.close();
488 },
489 });
490 
491 const request = new Request('http://example.org', {
492 method: 'POST',
493 body: rs,
494 });
495 
496 const reader = request.body
497 .pipeThrough(new TextDecoderStream('utf-8'))
498 .getReader();
499 strictEqual(typeof (await reader.read()).value, 'string');
500 },
501};