File
Blob: src/workerd/api/tests/streams-test.js
| 1 | // Copyright (c) 2024 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 | import { strictEqual, ok, deepStrictEqual, rejects, throws } from 'node:assert'; |
| 5 | |
| 6 | const enc = new TextEncoder(); |
| 7 | |
| 8 | export const rs = { |
| 9 | async test(ctrl, env) { |
| 10 | const resp = await env.subrequest.fetch('http://example.org', { |
| 11 | method: 'POST', |
| 12 | body: new ReadableStream({ |
| 13 | expectedLength: 10, |
| 14 | start(c) { |
| 15 | c.enqueue(enc.encode('hellohello')); |
| 16 | c.close(); |
| 17 | }, |
| 18 | }), |
| 19 | }); |
| 20 | for await (const _ of resp.body) { |
| 21 | // intentionally empty |
| 22 | } |
| 23 | }, |
| 24 | }; |
| 25 | |
| 26 | export const ts = { |
| 27 | async test(ctrl, env) { |
| 28 | const { readable, writable } = new TransformStream({ |
| 29 | expectedLength: 10, |
| 30 | }); |
| 31 | const writer = writable.getWriter(); |
| 32 | writer.write(enc.encode('hellohello')); |
| 33 | writer.close(); |
| 34 | const resp = await env.subrequest.fetch('http://example.org', { |
| 35 | method: 'POST', |
| 36 | body: readable, |
| 37 | }); |
| 38 | for await (const _ of resp.body) { |
| 39 | // intentionally empty |
| 40 | } |
| 41 | }, |
| 42 | }; |
| 43 | |
| 44 | // Regression test for https://github.com/cloudflare/workerd/issues/5113 |
| 45 | export const rsRequest = { |
| 46 | async test(ctrl, env) { |
| 47 | const resp = await env.subrequest.fetch( |
| 48 | new Request('http://example.org', { |
| 49 | method: 'POST', |
| 50 | body: new ReadableStream({ |
| 51 | expectedLength: 10, |
| 52 | start(c) { |
| 53 | c.enqueue(enc.encode('hellohello')); |
| 54 | c.close(); |
| 55 | }, |
| 56 | }), |
| 57 | }) |
| 58 | ); |
| 59 | for await (const _ of resp.body) { |
| 60 | // intentionally empty |
| 61 | } |
| 62 | }, |
| 63 | }; |
| 64 | |
| 65 | // Regression test for https://github.com/cloudflare/workerd/issues/5113 |
| 66 | export const tsRequest = { |
| 67 | async test(ctrl, env) { |
| 68 | const { readable, writable } = new TransformStream({ |
| 69 | expectedLength: 10, |
| 70 | }); |
| 71 | const writer = writable.getWriter(); |
| 72 | writer.write(enc.encode('hellohello')); |
| 73 | writer.close(); |
| 74 | const resp = await env.subrequest.fetch( |
| 75 | new Request('http://example.org', { |
| 76 | method: 'POST', |
| 77 | body: readable, |
| 78 | }) |
| 79 | ); |
| 80 | for await (const _ of resp.body) { |
| 81 | // intentionally empty |
| 82 | } |
| 83 | }, |
| 84 | }; |
| 85 | |
| 86 | export const byobMin = { |
| 87 | async test() { |
| 88 | let controller; |
| 89 | const rs = new ReadableStream({ |
| 90 | type: 'bytes', |
| 91 | start(c) { |
| 92 | controller = c; |
| 93 | }, |
| 94 | }); |
| 95 | |
| 96 | async function handleRead(readable) { |
| 97 | const reader = rs.getReader({ mode: 'byob' }); |
| 98 | const result = await reader.read(new Uint8Array(10), { min: 10 }); |
| 99 | strictEqual(result.done, false); |
| 100 | strictEqual(result.value.byteLength, 10); |
| 101 | } |
| 102 | |
| 103 | async function handlePush(controller) { |
| 104 | for (let n = 0; n < 10; n++) { |
| 105 | controller.enqueue(new Uint8Array(1)); |
| 106 | await scheduler.wait(10); |
| 107 | } |
| 108 | } |
| 109 | |
| 110 | const results = await Promise.allSettled([ |
| 111 | handleRead(rs), |
| 112 | handlePush(controller), |
| 113 | ]); |
| 114 | |
| 115 | strictEqual(results[0].status, 'fulfilled'); |
| 116 | strictEqual(results[1].status, 'fulfilled'); |
| 117 | }, |
| 118 | }; |
| 119 | |
| 120 | export const cancelReadsOnReleaseLock = { |
| 121 | async test() { |
| 122 | const rs = new ReadableStream(); |
| 123 | const reader = rs.getReader(); |
| 124 | const read = reader.read(); |
| 125 | |
| 126 | const result = await Promise.allSettled([read, reader.releaseLock()]); |
| 127 | strictEqual(result[0].status, 'rejected'); |
| 128 | strictEqual( |
| 129 | result[0].reason.message, |
| 130 | 'This ReadableStream reader has been released.' |
| 131 | ); |
| 132 | strictEqual(result[1].status, 'fulfilled'); |
| 133 | |
| 134 | // Make sure we can still get another reader |
| 135 | const _reader2 = rs.getReader(); |
| 136 | }, |
| 137 | }; |
| 138 | |
| 139 | export const cancelWriteOnReleaseLock = { |
| 140 | async test() { |
| 141 | const ws = new WritableStream({ |
| 142 | write() { |
| 143 | return new Promise(() => {}); |
| 144 | }, |
| 145 | }); |
| 146 | const writer = ws.getWriter(); |
| 147 | // This first write is just to start the write queue so that the |
| 148 | // next write becomes pending in the queue. This first write will |
| 149 | // never be fulfilled since it is in-progress but the queue will |
| 150 | // be rejected. |
| 151 | writer.write('ignored'); |
| 152 | const results = await Promise.allSettled([ |
| 153 | writer.write('hello'), |
| 154 | writer.releaseLock(), |
| 155 | ]); |
| 156 | strictEqual(results[0].status, 'rejected'); |
| 157 | strictEqual( |
| 158 | results[0].reason.message, |
| 159 | 'This WritableStream writer has been released.' |
| 160 | ); |
| 161 | strictEqual(results[1].status, 'fulfilled'); |
| 162 | |
| 163 | // Make sure we can still get another writer |
| 164 | const _writer2 = ws.getWriter(); |
| 165 | }, |
| 166 | }; |
| 167 | |
| 168 | export const readAllTextRequestSmall = { |
| 169 | async test() { |
| 170 | const rs = new ReadableStream({ |
| 171 | pull(c) { |
| 172 | c.enqueue(enc.encode('hello ')); |
| 173 | c.enqueue(enc.encode('world!')); |
| 174 | c.close(); |
| 175 | }, |
| 176 | }); |
| 177 | const request = new Request('http://example.org', { |
| 178 | method: 'POST', |
| 179 | body: rs, |
| 180 | }); |
| 181 | const text = await request.text(); |
| 182 | strictEqual(text, 'hello world!'); |
| 183 | }, |
| 184 | }; |
| 185 | |
| 186 | export const readAllTextResponseSmall = { |
| 187 | async test() { |
| 188 | const rs = new ReadableStream({ |
| 189 | pull(c) { |
| 190 | c.enqueue(enc.encode('hello ')); |
| 191 | c.enqueue(enc.encode('world!')); |
| 192 | c.close(); |
| 193 | }, |
| 194 | }); |
| 195 | const response = new Response(rs); |
| 196 | const text = await response.text(); |
| 197 | strictEqual(text, 'hello world!'); |
| 198 | }, |
| 199 | }; |
| 200 | |
| 201 | export const readAllTextRequestBig = { |
| 202 | async test() { |
| 203 | const chunks = [ |
| 204 | 'a'.repeat(4097), |
| 205 | 'b'.repeat(4097 * 2), |
| 206 | 'c'.repeat(4097 * 4), |
| 207 | ]; |
| 208 | let check = ''; |
| 209 | const enc = new TextEncoder(); |
| 210 | |
| 211 | const rs = new ReadableStream({ |
| 212 | pull(c) { |
| 213 | if (chunks.length === 0) { |
| 214 | c.close(); |
| 215 | return; |
| 216 | } |
| 217 | const chunk = chunks.shift(); |
| 218 | check += chunk; |
| 219 | c.enqueue(enc.encode(chunk)); |
| 220 | }, |
| 221 | }); |
| 222 | const request = new Request('http://example.org', { |
| 223 | method: 'POST', |
| 224 | body: rs, |
| 225 | }); |
| 226 | const text = await request.text(); |
| 227 | strictEqual(text.length, check.length); |
| 228 | strictEqual(text, check); |
| 229 | }, |
| 230 | }; |
| 231 | |
| 232 | export const readAllTextResponseBig = { |
| 233 | async test() { |
| 234 | const chunks = [ |
| 235 | 'a'.repeat(4097), |
| 236 | 'b'.repeat(4097 * 2), |
| 237 | 'c'.repeat(4097 * 4), |
| 238 | ]; |
| 239 | let check = ''; |
| 240 | const enc = new TextEncoder(); |
| 241 | |
| 242 | const rs = new ReadableStream({ |
| 243 | async pull(c) { |
| 244 | await scheduler.wait(10); |
| 245 | if (chunks.length === 0) { |
| 246 | c.close(); |
| 247 | return; |
| 248 | } |
| 249 | const chunk = chunks.shift(); |
| 250 | check += chunk; |
| 251 | c.enqueue(enc.encode(chunk)); |
| 252 | }, |
| 253 | }); |
| 254 | const response = new Response(rs); |
| 255 | const promise = response.text(); |
| 256 | const text = await promise; |
| 257 | strictEqual(text.length, check.length); |
| 258 | strictEqual(text, check); |
| 259 | }, |
| 260 | }; |
| 261 | |
| 262 | export const readAllTextFailedPull = { |
| 263 | async test() { |
| 264 | const rs = new ReadableStream({ |
| 265 | async pull(c) { |
| 266 | await scheduler.wait(10); |
| 267 | throw new Error('boom'); |
| 268 | }, |
| 269 | }); |
| 270 | const response = new Response(rs); |
| 271 | await rejects(response.text(), { message: 'boom' }); |
| 272 | }, |
| 273 | }; |
| 274 | |
| 275 | export const readAllTextFailedStart = { |
| 276 | async test() { |
| 277 | const rs = new ReadableStream({ |
| 278 | async start(c) { |
| 279 | await scheduler.wait(10); |
| 280 | throw new Error('boom'); |
| 281 | }, |
| 282 | }); |
| 283 | const response = new Response(rs); |
| 284 | await rejects(response.text(), { message: 'boom' }); |
| 285 | }, |
| 286 | }; |
| 287 | |
| 288 | export const readAllTextFailed = { |
| 289 | async test() { |
| 290 | const rs = new ReadableStream({ |
| 291 | async start(c) { |
| 292 | await scheduler.wait(10); |
| 293 | c.error(new Error('boom')); |
| 294 | }, |
| 295 | }); |
| 296 | const response = new Response(rs); |
| 297 | ok(!rs.locked); |
| 298 | const promise = response.text(); |
| 299 | ok(rs.locked); |
| 300 | await rejects(promise, { message: 'boom' }); |
| 301 | }, |
| 302 | }; |
| 303 | |
| 304 | export const tsCancel = { |
| 305 | async test() { |
| 306 | // Verify that a TransformStream's cancel function is called when the |
| 307 | // readable is canceled or the writable is aborted. Verify also that |
| 308 | // errors thrown by the cancel function are propagated. |
| 309 | { |
| 310 | let cancelCalled = false; |
| 311 | const { readable } = new TransformStream({ |
| 312 | async cancel(reason) { |
| 313 | strictEqual(reason, 'boom'); |
| 314 | await scheduler.wait(10); |
| 315 | cancelCalled = true; |
| 316 | }, |
| 317 | }); |
| 318 | ok(!cancelCalled); |
| 319 | await readable.cancel('boom'); |
| 320 | ok(cancelCalled); |
| 321 | } |
| 322 | |
| 323 | { |
| 324 | let cancelCalled = false; |
| 325 | const { writable } = new TransformStream({ |
| 326 | async cancel(reason) { |
| 327 | strictEqual(reason, 'boom'); |
| 328 | await scheduler.wait(10); |
| 329 | cancelCalled = true; |
| 330 | }, |
| 331 | }); |
| 332 | ok(!cancelCalled); |
| 333 | await writable.abort('boom'); |
| 334 | ok(cancelCalled); |
| 335 | } |
| 336 | |
| 337 | { |
| 338 | const { writable } = new TransformStream({ |
| 339 | async cancel(reason) { |
| 340 | throw new Error('boomy'); |
| 341 | }, |
| 342 | }); |
| 343 | await rejects(writable.abort('boom'), { message: 'boomy' }); |
| 344 | } |
| 345 | }, |
| 346 | }; |
| 347 | |
| 348 | export const writableStreamGcTraceFinishes = { |
| 349 | test() { |
| 350 | // TODO(soon): We really need better testing for GC visitation. |
| 351 | const _ws = new WritableStream(); |
| 352 | gc(); |
| 353 | }, |
| 354 | }; |
| 355 | |
| 356 | export const readableStreamFromAsyncGenerator = { |
| 357 | async test() { |
| 358 | async function* gen() { |
| 359 | await scheduler.wait(10); |
| 360 | yield 'hello'; |
| 361 | await scheduler.wait(10); |
| 362 | yield 'world'; |
| 363 | } |
| 364 | const rs = ReadableStream.from(gen()); |
| 365 | const chunks = []; |
| 366 | for await (const chunk of rs) { |
| 367 | chunks.push(chunk); |
| 368 | } |
| 369 | deepStrictEqual(chunks, ['hello', 'world']); |
| 370 | }, |
| 371 | }; |
| 372 | |
| 373 | export const readableStreamFromSyncGenerator = { |
| 374 | async test() { |
| 375 | const rs = ReadableStream.from(['hello', 'world']); |
| 376 | const chunks = []; |
| 377 | for await (const chunk of rs) { |
| 378 | chunks.push(chunk); |
| 379 | } |
| 380 | deepStrictEqual(chunks, ['hello', 'world']); |
| 381 | }, |
| 382 | }; |
| 383 | |
| 384 | export const readableStreamFromSyncGenerator2 = { |
| 385 | async test() { |
| 386 | function* gen() { |
| 387 | yield 'hello'; |
| 388 | yield 'world'; |
| 389 | } |
| 390 | const rs = ReadableStream.from(gen()); |
| 391 | const chunks = []; |
| 392 | for await (const chunk of rs) { |
| 393 | chunks.push(chunk); |
| 394 | } |
| 395 | deepStrictEqual(chunks, ['hello', 'world']); |
| 396 | }, |
| 397 | }; |
| 398 | |
| 399 | export const readableStreamFromAsyncCanceled = { |
| 400 | async test() { |
| 401 | async function* gen() { |
| 402 | let count = 0; |
| 403 | try { |
| 404 | count++; |
| 405 | yield 'hello'; |
| 406 | count++; |
| 407 | yield 'world'; |
| 408 | } finally { |
| 409 | strictEqual(count, 1); |
| 410 | } |
| 411 | } |
| 412 | const rs = ReadableStream.from(gen()); |
| 413 | const chunks = []; |
| 414 | for await (const chunk of rs) { |
| 415 | chunks.push(chunk); |
| 416 | return; |
| 417 | } |
| 418 | deepStrictEqual(chunks, ['hello']); |
| 419 | }, |
| 420 | }; |
| 421 | |
| 422 | export const readableStreamFromThrowingAsyncGen = { |
| 423 | async test() { |
| 424 | async function* gen() { |
| 425 | yield 'hello'; |
| 426 | throw new Error('boom'); |
| 427 | } |
| 428 | const rs = ReadableStream.from(gen()); |
| 429 | const chunks = []; |
| 430 | async function consumeStream() { |
| 431 | for await (const chunk of rs) { |
| 432 | chunks.push(chunk); |
| 433 | } |
| 434 | } |
| 435 | await rejects(consumeStream, { message: 'boom' }); |
| 436 | deepStrictEqual(chunks, ['hello']); |
| 437 | }, |
| 438 | }; |
| 439 | |
| 440 | export const readableStreamFromNoopAsyncGen = { |
| 441 | async test() { |
| 442 | async function* gen() {} |
| 443 | const rs = ReadableStream.from(gen()); |
| 444 | const chunks = []; |
| 445 | for await (const chunk of rs) { |
| 446 | chunks.push(chunk); |
| 447 | } |
| 448 | deepStrictEqual(chunks, []); |
| 449 | }, |
| 450 | }; |
| 451 | |
| 452 | // Tests for ReadableStream.from() cancel behavior per WPT spec |
| 453 | export const readableStreamFromCancelRejectsWhenReturnRejects = { |
| 454 | async test() { |
| 455 | const rejectError = new Error('return error'); |
| 456 | const iterable = { |
| 457 | async next() { |
| 458 | return { value: undefined, done: true }; |
| 459 | }, |
| 460 | async return() { |
| 461 | throw rejectError; |
| 462 | }, |
| 463 | [Symbol.asyncIterator]() { |
| 464 | return this; |
| 465 | }, |
| 466 | }; |
| 467 | |
| 468 | const rs = ReadableStream.from(iterable); |
| 469 | const reader = rs.getReader(); |
| 470 | |
| 471 | await rejects(reader.cancel(), rejectError); |
| 472 | }, |
| 473 | }; |
| 474 | |
| 475 | export const readableStreamFromCancelRejectsWhenReturnThrows = { |
| 476 | async test() { |
| 477 | const throwError = new Error('return throws'); |
| 478 | const iterable = { |
| 479 | async next() { |
| 480 | return { value: undefined, done: true }; |
| 481 | }, |
| 482 | return() { |
| 483 | throw throwError; |
| 484 | }, |
| 485 | [Symbol.asyncIterator]() { |
| 486 | return this; |
| 487 | }, |
| 488 | }; |
| 489 | |
| 490 | const rs = ReadableStream.from(iterable); |
| 491 | const reader = rs.getReader(); |
| 492 | |
| 493 | await rejects(reader.cancel(), (err) => err === throwError); |
| 494 | }, |
| 495 | }; |
| 496 | |
| 497 | export const readableStreamFromCancelRejectsWhenReturnNotMethod = { |
| 498 | async test() { |
| 499 | const iterable = { |
| 500 | async next() { |
| 501 | return { value: undefined, done: true }; |
| 502 | }, |
| 503 | return: 42, // exists but not callable |
| 504 | [Symbol.asyncIterator]() { |
| 505 | return this; |
| 506 | }, |
| 507 | }; |
| 508 | |
| 509 | const rs = ReadableStream.from(iterable); |
| 510 | const reader = rs.getReader(); |
| 511 | |
| 512 | await rejects(reader.cancel(), { |
| 513 | name: 'TypeError', |
| 514 | message: /return/, |
| 515 | }); |
| 516 | }, |
| 517 | }; |
| 518 | |
| 519 | export const readableStreamFromCancelRejectsWhenReturnNonObject = { |
| 520 | async test() { |
| 521 | const iterable = { |
| 522 | async next() { |
| 523 | return { value: undefined, done: true }; |
| 524 | }, |
| 525 | async return() { |
| 526 | return 42; // fulfills with non-object |
| 527 | }, |
| 528 | [Symbol.asyncIterator]() { |
| 529 | return this; |
| 530 | }, |
| 531 | }; |
| 532 | |
| 533 | const rs = ReadableStream.from(iterable); |
| 534 | const reader = rs.getReader(); |
| 535 | |
| 536 | await rejects(reader.cancel(), { |
| 537 | name: 'TypeError', |
| 538 | }); |
| 539 | }, |
| 540 | }; |
| 541 | |
| 542 | export const readableStreamFromCancelResolvesWhenReturnMissing = { |
| 543 | async test() { |
| 544 | const iterable = { |
| 545 | async next() { |
| 546 | return { value: undefined, done: true }; |
| 547 | }, |
| 548 | // no return method |
| 549 | [Symbol.asyncIterator]() { |
| 550 | return this; |
| 551 | }, |
| 552 | }; |
| 553 | |
| 554 | const rs = ReadableStream.from(iterable); |
| 555 | const reader = rs.getReader(); |
| 556 | |
| 557 | // Should resolve without error when return() is missing |
| 558 | await Promise.all([reader.cancel(), reader.closed]); |
| 559 | }, |
| 560 | }; |
| 561 | |
| 562 | export const abortWriterAfterGc = { |
| 563 | async test() { |
| 564 | function getWriter() { |
| 565 | const { writable } = new IdentityTransformStream(); |
| 566 | return writable.getWriter(); |
| 567 | } |
| 568 | |
| 569 | const writer = getWriter(); |
| 570 | gc(); |
| 571 | await writer.abort(); |
| 572 | }, |
| 573 | }; |
| 574 | |
| 575 | export const finalReadOnInternalStreamReturnsBuffer = { |
| 576 | async test() { |
| 577 | const { readable, writable } = new IdentityTransformStream(); |
| 578 | const writer = writable.getWriter(); |
| 579 | await writer.close(); |
| 580 | |
| 581 | const reader = readable.getReader({ mode: 'byob' }); |
| 582 | let result = await reader.read(new Uint8Array(10)); |
| 583 | strictEqual(result.done, true); |
| 584 | ok(result.value instanceof Uint8Array); |
| 585 | strictEqual(result.value.byteLength, 0); |
| 586 | strictEqual(result.value.buffer.byteLength, 10); |
| 587 | |
| 588 | result = await reader.read(new Uint8Array(10)); |
| 589 | strictEqual(result.done, true); |
| 590 | ok(result.value instanceof Uint8Array); |
| 591 | strictEqual(result.value.byteLength, 0); |
| 592 | strictEqual(result.value.buffer.byteLength, 10); |
| 593 | }, |
| 594 | }; |
| 595 | |
| 596 | // Test that canceling a stream rejects body consume function |
| 597 | export const cancelStreamRejectsBodyConsume = { |
| 598 | async test() { |
| 599 | const response = new Response('foo bar'); |
| 600 | const stream = response.body; |
| 601 | |
| 602 | stream.cancel(new Error('a good reason')); |
| 603 | |
| 604 | await rejects(response.text(), TypeError); |
| 605 | }, |
| 606 | }; |
| 607 | |
| 608 | // Test that canceling a reader resolves closed promise |
| 609 | export const cancelReaderResolvesClosedPromise = { |
| 610 | async test() { |
| 611 | const response = new Response('foo bar'); |
| 612 | const stream = response.body; |
| 613 | const reader = stream.getReader(); |
| 614 | |
| 615 | reader.cancel(); |
| 616 | const closed = await reader.closed; |
| 617 | strictEqual(typeof closed, 'undefined'); |
| 618 | reader.releaseLock(); |
| 619 | |
| 620 | await rejects(response.text(), TypeError); |
| 621 | }, |
| 622 | }; |
| 623 | |
| 624 | // Test that getReader with bad mode throws |
| 625 | export const getReaderBadModeThrows = { |
| 626 | test() { |
| 627 | const response = new Response('foo bar'); |
| 628 | const stream = response.body; |
| 629 | |
| 630 | throws(() => stream.getReader({ mode: 'nope' }), TypeError); |
| 631 | }, |
| 632 | }; |
| 633 | |
| 634 | // Test that stream is locked after getReader() called |
| 635 | export const streamLockedAfterGetReader = { |
| 636 | test() { |
| 637 | const response = new Response('foo bar'); |
| 638 | const stream = response.body; |
| 639 | |
| 640 | const reader = stream.getReader(); |
| 641 | |
| 642 | ok(stream.locked); |
| 643 | |
| 644 | throws(() => stream.getReader(), TypeError); |
| 645 | |
| 646 | reader.releaseLock(); |
| 647 | ok(!stream.locked); |
| 648 | reader.releaseLock(); // Second time should be a no-op |
| 649 | }, |
| 650 | }; |
| 651 | |
| 652 | // Test BYOB reader constraints |
| 653 | export const byobReaderConstraints = { |
| 654 | async test() { |
| 655 | const response = new Response('foo bar'); |
| 656 | const stream = response.body; |
| 657 | const reader = stream.getReader({ mode: 'byob' }); |
| 658 | // Start a read - this will consume part of the stream |
| 659 | reader.read(new Uint8Array(32)).catch(() => {}); // Ignore the result |
| 660 | |
| 661 | // We use rejects() with async wrapper instead of throws() because the error |
| 662 | // is thrown synchronously without streams_enable_constructors but returned as |
| 663 | // a rejected promise when that flag is enabled. The async wrapper handles both. |
| 664 | |
| 665 | // Cannot BYOB with a zero-length buffer |
| 666 | await rejects(async () => reader.read(new Uint8Array(0)), TypeError); |
| 667 | |
| 668 | // Cannot BYOB an ArrayBuffer, only an ArrayBufferView |
| 669 | await rejects(async () => reader.read(new ArrayBuffer(32)), TypeError); |
| 670 | |
| 671 | // Cannot use BYOB reader as a non-BYOB reader |
| 672 | await rejects(async () => reader.read(), TypeError); |
| 673 | }, |
| 674 | }; |
| 675 | |
| 676 | // Test cancel error type propagation |
| 677 | export const cancelErrorTypePropagation = { |
| 678 | async test() { |
| 679 | class ExampleError extends Error { |
| 680 | constructor() { |
| 681 | super('foo bar'); |
| 682 | this.name = 'ExampleError'; |
| 683 | } |
| 684 | } |
| 685 | |
| 686 | const cancelErrorTests = [ |
| 687 | { |
| 688 | cancelWith: new Error('test'), |
| 689 | expectError: 'Error: test', |
| 690 | }, |
| 691 | { |
| 692 | cancelWith: 'test', |
| 693 | expectError: 'Error: test', |
| 694 | }, |
| 695 | { |
| 696 | cancelWith: 'jsg.Error: test', |
| 697 | expectError: 'Error: jsg.Error: test', |
| 698 | }, |
| 699 | { |
| 700 | cancelWith: new TypeError('Problems!'), |
| 701 | expectError: 'TypeError: Problems!', |
| 702 | errorType: TypeError, |
| 703 | }, |
| 704 | { |
| 705 | cancelWith: new RangeError('Problems!'), |
| 706 | expectError: 'RangeError: Problems!', |
| 707 | errorType: RangeError, |
| 708 | }, |
| 709 | { |
| 710 | cancelWith: new SyntaxError('The semicolons are bad'), |
| 711 | expectError: 'SyntaxError: The semicolons are bad', |
| 712 | errorType: SyntaxError, |
| 713 | }, |
| 714 | { |
| 715 | cancelWith: new ReferenceError("Didn't find it"), |
| 716 | expectError: "ReferenceError: Didn't find it", |
| 717 | errorType: ReferenceError, |
| 718 | }, |
| 719 | { |
| 720 | cancelWith: undefined, |
| 721 | expectError: 'Error: Stream was cancelled.', |
| 722 | }, |
| 723 | { |
| 724 | cancelWith: new ExampleError(), |
| 725 | expectError: 'ExampleError: foo bar', |
| 726 | }, |
| 727 | ]; |
| 728 | |
| 729 | for (const testCase of cancelErrorTests) { |
| 730 | const ts = new IdentityTransformStream(); |
| 731 | |
| 732 | const writer = ts.writable.getWriter(); |
| 733 | const reader = ts.readable.getReader(); |
| 734 | const writePromise = writer.write(new TextEncoder().encode('a')); |
| 735 | const writerActualClosed = writer.close(); |
| 736 | await reader.cancel(testCase.cancelWith); |
| 737 | |
| 738 | for (const promise of [writePromise, writerActualClosed]) { |
| 739 | await rejects(promise, (e) => { |
| 740 | strictEqual(String(e), testCase.expectError); |
| 741 | if (testCase.errorType) { |
| 742 | ok(e instanceof testCase.errorType); |
| 743 | } |
| 744 | return true; |
| 745 | }); |
| 746 | } |
| 747 | } |
| 748 | }, |
| 749 | }; |
| 750 | |
| 751 | // Test IdentityTransformStream write before read |
| 752 | export const identityTransformWriteBeforeRead = { |
| 753 | async test() { |
| 754 | const MAX_RW = 10; |
| 755 | const { readable, writable } = new IdentityTransformStream(); |
| 756 | const writer = writable.getWriter(); |
| 757 | const reader = readable.getReader(); |
| 758 | |
| 759 | const writePromises = []; |
| 760 | for (let i = 0; i < MAX_RW; i++) { |
| 761 | writePromises.push(writer.write(new Uint8Array([i]))); |
| 762 | } |
| 763 | |
| 764 | const chunks = []; |
| 765 | for (let i = 0; i < MAX_RW; i++) { |
| 766 | chunks.push(await reader.read()); |
| 767 | } |
| 768 | |
| 769 | await Promise.all(writePromises); |
| 770 | |
| 771 | for (let i = 0; i < chunks.length; i++) { |
| 772 | deepStrictEqual([...chunks[i].value], [i]); |
| 773 | strictEqual(chunks[i].done, false); |
| 774 | } |
| 775 | |
| 776 | const writeClosePromise = writer.close(); |
| 777 | const chunk = await reader.read(); |
| 778 | await writeClosePromise; |
| 779 | |
| 780 | strictEqual(chunk.done, true); |
| 781 | |
| 782 | await writer.closed; |
| 783 | await reader.closed; |
| 784 | }, |
| 785 | }; |
| 786 | |
| 787 | // Test IdentityTransformStream read before write |
| 788 | export const identityTransformReadBeforeWrite = { |
| 789 | async test() { |
| 790 | const MAX_RW = 10; |
| 791 | const { readable, writable } = new IdentityTransformStream(); |
| 792 | const writer = writable.getWriter(); |
| 793 | const reader = readable.getReader(); |
| 794 | |
| 795 | // IdentityTransformStream only supports one pending read at a time, |
| 796 | // so we test read-before-write by starting each read before its write |
| 797 | const chunks = []; |
| 798 | for (let i = 0; i < MAX_RW; i++) { |
| 799 | const readPromise = reader.read(); |
| 800 | await writer.write(new Uint8Array([i])); |
| 801 | chunks.push(await readPromise); |
| 802 | } |
| 803 | |
| 804 | for (let i = 0; i < chunks.length; i++) { |
| 805 | deepStrictEqual([...chunks[i].value], [i]); |
| 806 | strictEqual(chunks[i].done, false); |
| 807 | } |
| 808 | |
| 809 | const readClosePromise = reader.read(); |
| 810 | await writer.close(); |
| 811 | const chunk = await readClosePromise; |
| 812 | |
| 813 | strictEqual(chunk.done, true); |
| 814 | |
| 815 | await writer.closed; |
| 816 | await reader.closed; |
| 817 | }, |
| 818 | }; |
| 819 | |
| 820 | // Test closed promise under lock release |
| 821 | export const closedPromiseUnderLockRelease = { |
| 822 | async test() { |
| 823 | const { readable, writable } = new IdentityTransformStream(); |
| 824 | |
| 825 | const writer = writable.getWriter(); |
| 826 | const reader = readable.getReader(); |
| 827 | |
| 828 | const writerClosed = writer.closed; |
| 829 | const readerClosed = reader.closed; |
| 830 | |
| 831 | writer.releaseLock(); |
| 832 | |
| 833 | await rejects(writerClosed, TypeError); |
| 834 | |
| 835 | reader.releaseLock(); |
| 836 | |
| 837 | await rejects(readerClosed, TypeError); |
| 838 | }, |
| 839 | }; |
| 840 | |
| 841 | // Test closed promise under writer abort |
| 842 | export const closedPromiseUnderWriterAbort = { |
| 843 | async test() { |
| 844 | const { readable, writable } = new IdentityTransformStream(); |
| 845 | |
| 846 | const writer = writable.getWriter(); |
| 847 | const reader = readable.getReader(); |
| 848 | |
| 849 | const writerClosed = writer.closed; |
| 850 | const readerClosed = reader.closed; |
| 851 | |
| 852 | const readPromise = reader.read(); |
| 853 | await writer.abort(new Error('Some arbitrary, capricious reason.')); |
| 854 | |
| 855 | await rejects(writerClosed, Error); |
| 856 | await rejects(readPromise, Error); |
| 857 | await rejects(readerClosed, Error); |
| 858 | }, |
| 859 | }; |
| 860 | |
| 861 | // Test FixedLengthStream constructor preconditions |
| 862 | export const fixedLengthStreamPreconditions = { |
| 863 | test() { |
| 864 | // Can construct with negative zero |
| 865 | new FixedLengthStream(-0.0); |
| 866 | |
| 867 | // Can construct with fraction (coerced to 0) |
| 868 | new FixedLengthStream(0.00001); |
| 869 | |
| 870 | // Can construct with MAX_SAFE_INTEGER |
| 871 | new FixedLengthStream(Number.MAX_SAFE_INTEGER); |
| 872 | |
| 873 | // Cannot construct with unsafe integer |
| 874 | throws(() => new FixedLengthStream(Number.MAX_SAFE_INTEGER + 1), TypeError); |
| 875 | |
| 876 | // Cannot construct with negative integer |
| 877 | throws(() => new FixedLengthStream(-1), TypeError); |
| 878 | }, |
| 879 | }; |
| 880 | |
| 881 | // Test non-standard readAtLeast() extension with default reader (should throw) |
| 882 | export const readAtLeastDefaultReaderThrows = { |
| 883 | async test() { |
| 884 | const rs = new ReadableStream({ |
| 885 | type: 'bytes', |
| 886 | pull(c) { |
| 887 | c.enqueue(enc.encode('hello')); |
| 888 | c.close(); |
| 889 | }, |
| 890 | }); |
| 891 | |
| 892 | const reader = rs.getReader(); |
| 893 | throws(() => reader.readAtLeast(1), TypeError); |
| 894 | reader.releaseLock(); |
| 895 | |
| 896 | // Consume the stream to clean up |
| 897 | for await (const _ of rs) { |
| 898 | // intentionally empty |
| 899 | } |
| 900 | }, |
| 901 | }; |
| 902 | |
| 903 | // Test non-standard readAtLeast() extension with BYOB reader |
| 904 | // Note: The original ew-test expected value=undefined on done, which was the legacy |
| 905 | // behavior of internal streams. With `internal_stream_byob_return_view` compat flag |
| 906 | // (enabled since 2024-05-13), the spec-compliant behavior returns an empty view. |
| 907 | export const readAtLeastByobReader = { |
| 908 | async test(ctrl, env) { |
| 909 | // Use service binding to get chunked response |
| 910 | const response = await env.subrequest.fetch('http://test/chunked'); |
| 911 | const reader = response.body.getReader({ mode: 'byob' }); |
| 912 | |
| 913 | // First readAtLeast: request min 4 bytes |
| 914 | // Server sends: 'foo' (3) + 'bar' (3) = 6 bytes, first chunk 'foo' only 3 bytes |
| 915 | // so readAtLeast(4) should wait for more data |
| 916 | let result = await reader.readAtLeast(4, new Uint8Array(20)); |
| 917 | let value = new TextDecoder().decode(result.value); |
| 918 | strictEqual(result.done, false); |
| 919 | strictEqual(value.length, 6); |
| 920 | strictEqual(value, 'foobar'); |
| 921 | |
| 922 | // Regular read |
| 923 | result = await reader.read(new Uint8Array(20)); |
| 924 | value = new TextDecoder().decode(result.value); |
| 925 | strictEqual(value.length, 1); |
| 926 | strictEqual(value, 'b'); |
| 927 | strictEqual(result.done, false); |
| 928 | |
| 929 | // Second readAtLeast: request min 4 bytes, only 'az' (2 bytes) remain |
| 930 | // Server sends: 'a' (1) + 'z' (1) = 2 bytes, then closes |
| 931 | result = await reader.readAtLeast(4, new Uint8Array(20)); |
| 932 | value = new TextDecoder().decode(result.value); |
| 933 | strictEqual(value.length, 2); |
| 934 | strictEqual(value, 'az'); |
| 935 | strictEqual(result.done, false); |
| 936 | |
| 937 | // Final read should be done - spec requires empty view, not undefined |
| 938 | result = await reader.readAtLeast(4, new Uint8Array(20)); |
| 939 | strictEqual(result.done, true); |
| 940 | ok(result.value instanceof Uint8Array); |
| 941 | strictEqual(result.value.byteLength, 0); |
| 942 | }, |
| 943 | }; |
| 944 | |
| 945 | export const writeSubarray = { |
| 946 | async test() { |
| 947 | const { readable, writable } = new IdentityTransformStream(); |
| 948 | |
| 949 | const u8 = new Uint8Array([1, 2, 3, 4]); |
| 950 | |
| 951 | const writer = writable.getWriter(); |
| 952 | const reader = readable.getReader(); |
| 953 | |
| 954 | writer.write(u8.subarray(1, 3)); |
| 955 | writer.close(); |
| 956 | |
| 957 | const { value } = await reader.read(); |
| 958 | |
| 959 | strictEqual(value.length, 2); |
| 960 | strictEqual(value[0], u8[1]); |
| 961 | strictEqual(value[1], u8[2]); |
| 962 | }, |
| 963 | }; |
| 964 | |
| 965 | export const writableStreamWriterConstructor = { |
| 966 | test() { |
| 967 | const t = new IdentityTransformStream(); |
| 968 | new WritableStreamDefaultWriter(t.writable); |
| 969 | }, |
| 970 | }; |
| 971 | |
| 972 | export const readableStreamDefaultReaderConstructor = { |
| 973 | test() { |
| 974 | const t = new IdentityTransformStream(); |
| 975 | new ReadableStreamDefaultReader(t.readable); |
| 976 | }, |
| 977 | }; |
| 978 | |
| 979 | export const readableStreamByobReaderConstructor = { |
| 980 | test() { |
| 981 | const t = new IdentityTransformStream(); |
| 982 | new ReadableStreamBYOBReader(t.readable); |
| 983 | }, |
| 984 | }; |
| 985 | |
| 986 | export const byobReaderDetachesBuffer = { |
| 987 | async test() { |
| 988 | const ts = new IdentityTransformStream(); |
| 989 | const view = new Uint8Array(10); |
| 990 | const buffer = view.buffer; |
| 991 | const writer = ts.writable.getWriter(); |
| 992 | const reader = ts.readable.getReader({ mode: 'byob' }); |
| 993 | strictEqual(view.byteLength, 10); |
| 994 | strictEqual(view.buffer.byteLength, 10); |
| 995 | const res = await Promise.all([ |
| 996 | writer.write(new Uint8Array(10)), |
| 997 | reader.read(view), |
| 998 | ]); |
| 999 | |
| 1000 | strictEqual(view.byteLength, 0); |
| 1001 | strictEqual(view.buffer.byteLength, 0); |
| 1002 | |
| 1003 | ok(res[1].value.buffer instanceof ArrayBuffer); |
| 1004 | ok(res[1].value.buffer !== buffer); |
| 1005 | |
| 1006 | await rejects(async () => reader.read(view), TypeError); |
| 1007 | |
| 1008 | // Using a non-detachable ArrayBuffer must fail with a rejection |
| 1009 | const memory = new WebAssembly.Memory({ |
| 1010 | initial: 10, |
| 1011 | maximum: 10, |
| 1012 | shared: true, |
| 1013 | }); |
| 1014 | await rejects( |
| 1015 | async () => reader.read(new Uint8Array(memory.buffer)), |
| 1016 | TypeError |
| 1017 | ); |
| 1018 | |
| 1019 | await rejects( |
| 1020 | async () => reader.read(new Uint8Array(new SharedArrayBuffer(10))), |
| 1021 | TypeError |
| 1022 | ); |
| 1023 | }, |
| 1024 | }; |
| 1025 | |
| 1026 | export const captureSyncThrows = { |
| 1027 | async test() { |
| 1028 | const { readable } = new IdentityTransformStream(); |
| 1029 | const reader = readable.getReader({ mode: 'byob' }); |
| 1030 | // Without the captureThrowsAsRejections flag enabled, this would throw synchronously. |
| 1031 | // With the flag enabled, however, the synchronous throw is changed into a promise rejection. |
| 1032 | await rejects(async () => reader.read(new ArrayBuffer(10)), TypeError); |
| 1033 | }, |
| 1034 | }; |
| 1035 | |
| 1036 | export const teeFixedLengthStreamNoHang = { |
| 1037 | async test() { |
| 1038 | const ts = new FixedLengthStream(11); |
| 1039 | const writer = ts.writable.getWriter(); |
| 1040 | writer.write(new TextEncoder().encode('foo bar baz')); |
| 1041 | writer.close(); |
| 1042 | const [left, _right] = ts.readable.tee(); |
| 1043 | const response = new Response(left); |
| 1044 | strictEqual(await response.text(), 'foo bar baz'); |
| 1045 | }, |
| 1046 | }; |
| 1047 | |
| 1048 | export const transformStreamReadAllBytes = { |
| 1049 | async test() { |
| 1050 | const { readable, writable } = new IdentityTransformStream(); |
| 1051 | const response = new Response(readable); |
| 1052 | const writer = writable.getWriter(); |
| 1053 | |
| 1054 | const N = 8; |
| 1055 | const M = 5000; |
| 1056 | |
| 1057 | const writePromise = (async () => { |
| 1058 | for (let i = 0; i < N; i++) { |
| 1059 | const chunk = new Uint8Array(M); |
| 1060 | chunk.fill(i + 1); |
| 1061 | await writer.write(chunk); |
| 1062 | } |
| 1063 | await writer.close(); |
| 1064 | })(); |
| 1065 | |
| 1066 | const body = new Uint8Array(await response.arrayBuffer()); |
| 1067 | strictEqual(body.byteLength, N * M); |
| 1068 | for (let i = 0; i < body.length; i++) { |
| 1069 | strictEqual(body[i], Math.floor(i / M) + 1); |
| 1070 | } |
| 1071 | await writePromise; |
| 1072 | }, |
| 1073 | }; |
| 1074 | |
| 1075 | export const transformStreamReadAllText = { |
| 1076 | async test() { |
| 1077 | const { readable, writable } = new IdentityTransformStream(); |
| 1078 | const response = new Response(readable); |
| 1079 | const writer = writable.getWriter(); |
| 1080 | |
| 1081 | const lowerCaseA = 97; |
| 1082 | const N = 8; |
| 1083 | const M = 5000; |
| 1084 | |
| 1085 | const writePromise = (async () => { |
| 1086 | for (let i = 0; i < N; i++) { |
| 1087 | const chunk = new Uint8Array(M); |
| 1088 | chunk.fill(i + lowerCaseA); |
| 1089 | await writer.write(chunk); |
| 1090 | } |
| 1091 | await writer.close(); |
| 1092 | })(); |
| 1093 | |
| 1094 | const body = await response.text(); |
| 1095 | strictEqual(body.length, N * M); |
| 1096 | for (let i = 0; i < N; i++) { |
| 1097 | const expected = String.fromCharCode(i + lowerCaseA).repeat(M); |
| 1098 | strictEqual(body.slice(i * M, i * M + M), expected); |
| 1099 | } |
| 1100 | await writePromise; |
| 1101 | }, |
| 1102 | }; |
| 1103 | |
| 1104 | export const concurrentReadsRejected = { |
| 1105 | async test() { |
| 1106 | const { readable } = new IdentityTransformStream(); |
| 1107 | const reader = readable.getReader(); |
| 1108 | const _p0 = reader.read(); |
| 1109 | await rejects(reader.read(), TypeError); |
| 1110 | }, |
| 1111 | }; |
| 1112 | |
| 1113 | export default { |
| 1114 | async fetch(request, env) { |
| 1115 | const url = new URL(request.url); |
| 1116 | |
| 1117 | // Endpoint for chunked data for readAtLeast tests |
| 1118 | if (url.pathname === '/chunked') { |
| 1119 | const rs = new ReadableStream({ |
| 1120 | type: 'bytes', |
| 1121 | async pull(controller) { |
| 1122 | // Simulate chunked input: foo, bar, b, a, z |
| 1123 | const chunks = [ |
| 1124 | enc.encode('foo'), |
| 1125 | enc.encode('bar'), |
| 1126 | enc.encode('b'), |
| 1127 | enc.encode('a'), |
| 1128 | enc.encode('z'), |
| 1129 | ]; |
| 1130 | for (const chunk of chunks) { |
| 1131 | controller.enqueue(chunk); |
| 1132 | await scheduler.wait(1); |
| 1133 | } |
| 1134 | controller.close(); |
| 1135 | }, |
| 1136 | }); |
| 1137 | return new Response(rs); |
| 1138 | } |
| 1139 | |
| 1140 | strictEqual(request.headers.get('content-length'), '10'); |
| 1141 | return new Response(request.body); |
| 1142 | }, |
| 1143 | }; |