Skip to content
File

Blob: src/workerd/api/streams/internal-test.c++

32.8 KB
1// Copyright (c) 2017-2022 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#include "internal.h"
6#include "readable.h"
7#include "standard.h"
8#include "writable.h"
9 
10#include <workerd/jsg/jsg-test.h>
11#include <workerd/jsg/jsg.h>
12#include <workerd/tests/test-fixture.h>
13 
14#include <openssl/rand.h>
15 
16namespace workerd::api {
17namespace {
18 
19// ======================================================================================
20// Shared test helpers
21 
22// Simple source that returns EOF immediately
23class EofSource final: public ReadableStreamSource {
24 public:
25 kj::Promise<size_t> tryRead(void*, size_t, size_t) override {
26 return static_cast<size_t>(0);
27 }
28};
29 
30// Simple sink that accepts all writes
31class NoopSink final: public WritableStreamSink {
32 public:
33 kj::Promise<void> write(kj::ArrayPtr<const byte>) override {
34 return kj::READY_NOW;
35 }
36 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>>) override {
37 return kj::READY_NOW;
38 }
39 kj::Promise<void> end() override {
40 return kj::READY_NOW;
41 }
42 void abort(kj::Exception) override {}
43};
44 
45// Creates a TestFixture with common flags for stream tests
46TestFixture makeStreamTestFixture() {
47 capnp::MallocMessageBuilder message;
48 auto flags = message.initRoot<CompatibilityFlags>();
49 flags.setStreamsJavaScriptControllers(true);
50 return TestFixture({.featureFlags = flags.asReader()});
51}
52 
53// Creates a TestFixture with the abortClearsQueue flag for testing abort behavior
54TestFixture makeAbortClearsQueueTestFixture() {
55 capnp::MallocMessageBuilder message;
56 auto flags = message.initRoot<CompatibilityFlags>();
57 flags.setStreamsJavaScriptControllers(true);
58 flags.setInternalWritableStreamAbortClearsQueue(true);
59 return TestFixture({.featureFlags = flags.asReader()});
60}
61 
62// Creates a BYOB-capable ReadableStream
63jsg::Ref<ReadableStream> makeByteStream(jsg::Lock& js) {
64 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
65 rs->getController().setup(
66 js, UnderlyingSource{.type = kj::str("bytes")}, StreamQueuingStrategy{});
67 return rs;
68}
69 
70// ======================================================================================
71// ReadableStreamSource test implementations
72 
73template <int size>
74class FooStream: public ReadableStreamSource {
75 public:
76 FooStream(): ptr(&data[0]), remaining_(size) {
77 KJ_ASSERT(RAND_bytes(data, size) == 1);
78 }
79 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
80 maxMaxBytesSeen_ = kj::max(maxMaxBytesSeen_, maxBytes);
81 numreads_++;
82 if (remaining_ == 0) return static_cast<size_t>(0);
83 KJ_ASSERT(minBytes == maxBytes);
84 auto amount = kj::min(remaining_, maxBytes);
85 memcpy(buffer, ptr, amount);
86 ptr += amount;
87 remaining_ -= amount;
88 return amount;
89 }
90 
91 kj::ArrayPtr<uint8_t> buf() {
92 return data;
93 }
94 
95 size_t remaining() {
96 return remaining_;
97 }
98 
99 size_t numreads() {
100 return numreads_;
101 }
102 
103 size_t maxMaxBytesSeen() {
104 return maxMaxBytesSeen_;
105 }
106 
107 private:
108 uint8_t data[size];
109 uint8_t* ptr;
110 size_t remaining_;
111 size_t numreads_ = 0;
112 size_t maxMaxBytesSeen_ = 0;
113};
114 
115template <int size>
116class BarStream: public FooStream<size> {
117 public:
118 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override {
119 return size;
120 }
121};
122 
123KJ_TEST("test") {
124 kj::EventLoop loop;
125 kj::WaitScope waitScope(loop);
126 
127 // In this first case, the stream does not report a length. The read size
128 // is min(limit, DEFAULT_BUFFER_CHUNK) = min(10001, 131072) = 10001, so the
129 // entire stream is consumed in a single read that returns a short read (10000 < 10001).
130 FooStream<10000> stream;
131 
132 stream.readAllBytes(10001)
133 .then([&](auto bytes) {
134 KJ_ASSERT(bytes.size() == 10000);
135 KJ_ASSERT(bytes == stream.buf().first(10000));
136 }).wait(waitScope);
137 
138 KJ_ASSERT(stream.numreads() == 1);
139 KJ_ASSERT(stream.maxMaxBytesSeen() == 10001);
140}
141 
142KJ_TEST("test (text)") {
143 kj::EventLoop loop;
144 kj::WaitScope waitScope(loop);
145 
146 // In this first case, the stream does not report a length. The read size
147 // is min(limit, DEFAULT_BUFFER_CHUNK) = min(10001, 131072) = 10001, so the
148 // entire stream is consumed in a single read that returns a short read (10000 < 10001).
149 FooStream<10000> stream;
150 
151 stream.readAllText(10001)
152 .then([&](auto bytes) {
153 KJ_ASSERT(bytes.size() == 10000);
154 KJ_ASSERT(bytes.asBytes() == stream.buf().first(10000));
155 }).wait(waitScope);
156 
157 KJ_ASSERT(stream.numreads() == 1);
158 KJ_ASSERT(stream.maxMaxBytesSeen() == 10001);
159}
160 
161KJ_TEST("test2") {
162 kj::EventLoop loop;
163 kj::WaitScope waitScope(loop);
164 
165 // In this second case, the stream does report a size, so we should see
166 // only one read.
167 BarStream<10000> stream;
168 
169 stream.readAllBytes(10001)
170 .then([&](auto bytes) {
171 KJ_ASSERT(bytes.size() == 10000);
172 KJ_ASSERT(bytes == stream.buf().first(10000));
173 }).wait(waitScope);
174 
175 KJ_ASSERT(stream.numreads() == 2);
176 KJ_ASSERT(stream.maxMaxBytesSeen() == 10000);
177}
178 
179KJ_TEST("test2 (text)") {
180 kj::EventLoop loop;
181 kj::WaitScope waitScope(loop);
182 
183 // In this second case, the stream does report a size, so we should see
184 // only one read.
185 BarStream<10000> stream;
186 
187 stream.readAllText(10001)
188 .then([&](auto bytes) {
189 KJ_ASSERT(bytes.size() == 10000);
190 KJ_ASSERT(bytes.asBytes() == stream.buf().first(10000));
191 }).wait(waitScope);
192 
193 KJ_ASSERT(stream.numreads() == 2);
194 KJ_ASSERT(stream.maxMaxBytesSeen() == 10000);
195}
196 
197KJ_TEST("zero-length stream") {
198 kj::EventLoop loop;
199 kj::WaitScope waitScope(loop);
200 
201 class Zero: public ReadableStreamSource {
202 public:
203 kj::Promise<size_t> tryRead(void*, size_t, size_t) override {
204 return static_cast<size_t>(0);
205 }
206 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override {
207 return static_cast<size_t>(0);
208 }
209 };
210 
211 Zero zero;
212 zero.readAllBytes(10).then([&](kj::Array<kj::byte> bytes) {
213 KJ_ASSERT(bytes.size() == 0);
214 }).wait(waitScope);
215}
216 
217KJ_TEST("lying stream") {
218 kj::EventLoop loop;
219 kj::WaitScope waitScope(loop);
220 
221 class Dishonest: public FooStream<10000> {
222 public:
223 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override {
224 return static_cast<size_t>(10);
225 }
226 };
227 
228 Dishonest stream;
229 stream.readAllBytes(10001)
230 .then([&](kj::Array<kj::byte> bytes) {
231 // The stream lies! it says there are only 10 bytes but there are more.
232 // oh well, we at least make sure we get the right result.
233 KJ_ASSERT(bytes.size() == 10000);
234 }).wait(waitScope);
235 
236 KJ_ASSERT(stream.numreads() == 1001);
237 KJ_ASSERT(stream.maxMaxBytesSeen() == 10);
238}
239 
240KJ_TEST("honest small stream") {
241 kj::EventLoop loop;
242 kj::WaitScope waitScope(loop);
243 
244 class HonestSmall: public FooStream<100> {
245 public:
246 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override {
247 return static_cast<size_t>(100);
248 }
249 };
250 
251 HonestSmall stream;
252 stream.readAllBytes(1001).then([&](kj::Array<kj::byte> bytes) {
253 KJ_ASSERT(bytes.size() == 100);
254 }).wait(waitScope);
255 
256 KJ_ASSERT(stream.numreads() == 2);
257 KJ_ASSERT(stream.maxMaxBytesSeen(), 100);
258}
259 
260KJ_TEST("WritableStreamInternalController queue size assertion") {
261 auto fixture = makeStreamTestFixture();
262 fixture.runInIoContext([&](const TestFixture::Environment& env) {
263 // Make sure that while an internal sink is being piped into, no other writes are
264 // allowed to be queued.
265 
266 jsg::Ref<ReadableStream> source = ReadableStream::constructor(env.js, kj::none, kj::none);
267 jsg::Ref<WritableStream> sink =
268 env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none);
269 
270 auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{.preventClose = true});
271 
272 KJ_ASSERT(sink->isLocked());
273 try {
274 sink->getWriter(env.js);
275 KJ_FAIL_ASSERT("Expected getWriter to throw");
276 } catch (...) {
277 auto ex = kj::getCaughtExceptionAsKj();
278 KJ_ASSERT(ex.getDescription() ==
279 "expected !stream->isLocked(); jsg.TypeError: This WritableStream "
280 "is currently locked to a writer.");
281 }
282 
283 auto buffersource = env.js.bytes(kj::heapArray<kj::byte>(10));
284 
285 bool writeFailed = false;
286 
287 auto write = sink->getController()
288 .write(env.js, buffersource.getHandle(env.js))
289 .catch_(env.js, [&](jsg::Lock& js, jsg::Value value) {
290 writeFailed = true;
291 auto ex = js.exceptionToKj(kj::mv(value));
292 KJ_ASSERT(
293 ex.getDescription() == "jsg.TypeError: This WritableStream is currently being piped to.");
294 });
295 
296 source->getController().cancel(env.js, kj::none);
297 
298 env.js.runMicrotasks();
299 
300 KJ_ASSERT(!sink->isLocked());
301 KJ_ASSERT(!sink->getController().isClosedOrClosing());
302 KJ_ASSERT(!sink->getController().isErrored());
303 KJ_ASSERT(sink->getController().isErroring(env.js) == kj::none);
304 
305 // Getting a writer at this point does not throw...
306 sink->getWriter(env.js);
307 });
308}
309 
310KJ_TEST("WritableStreamInternalController operations reject when piped to") {
311 // Tests that close/flush/tryPipeFrom reject with "currently being piped to"
312 // during an active pipe operation.
313 auto fixture = makeStreamTestFixture();
314 
315 fixture.runInIoContext([&](const TestFixture::Environment& env) {
316 auto source = ReadableStream::constructor(env.js, kj::none, kj::none);
317 auto source2 = ReadableStream::constructor(env.js, kj::none, kj::none);
318 auto sink = env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none);
319 
320 auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{.preventClose = true});
321 KJ_ASSERT(sink->isLocked());
322 
323 constexpr auto expectedError =
324 "jsg.TypeError: This WritableStream is currently being piped to."_kj;
325 
326 auto expectReject = [&](jsg::Promise<void> promise, bool& flag) {
327 promise.catch_(env.js, [&](jsg::Lock& js, jsg::Value value) {
328 flag = true;
329 KJ_ASSERT(js.exceptionToKj(kj::mv(value)).getDescription() == expectedError);
330 });
331 };
332 
333 bool closeFailed = false, flushFailed = false, pipeFailed = false;
334 
335 expectReject(sink->getController().close(env.js), closeFailed);
336 expectReject(sink->getController().flush(env.js), flushFailed);
337 
338 KJ_IF_SOME(secondPipe,
339 sink->getController().tryPipeFrom(
340 env.js, source2.addRef(), PipeToOptions{.preventClose = true})) {
341 expectReject(kj::mv(secondPipe), pipeFailed);
342 } else {
343 KJ_FAIL_ASSERT("Expected tryPipeFrom to return a promise");
344 }
345 
346 source->getController().cancel(env.js, kj::none);
347 env.js.runMicrotasks();
348 
349 KJ_ASSERT(closeFailed);
350 KJ_ASSERT(flushFailed);
351 KJ_ASSERT(pipeFailed);
352 });
353}
354 
355KJ_TEST("WritableStreamInternalController observability") {
356 auto fixture = makeStreamTestFixture();
357 
358 class MyObserver final: public ByteStreamObserver {
359 public:
360 void onChunkEnqueued(size_t bytes) override {
361 ++queueSize;
362 queueSizeBytes += bytes;
363 };
364 void onChunkDequeued(size_t bytes) override {
365 queueSizeBytes -= bytes;
366 --queueSize;
367 };
368 uint64_t queueSize = 0;
369 uint64_t queueSizeBytes = 0;
370 };
371 
372 auto myObserver = kj::heap<MyObserver>();
373 auto& observer = *myObserver;
374 kj::Maybe<jsg::Ref<WritableStream>> stream;
375 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
376 stream = env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::mv(myObserver));
377 
378 auto write = [&](size_t size) {
379 auto buffersource = env.js.bytes(kj::heapArray<kj::byte>(size));
380 return env.context.awaitJs(env.js,
381 KJ_ASSERT_NONNULL(stream)->getController().write(env.js, buffersource.getHandle(env.js)));
382 };
383 
384 KJ_ASSERT(observer.queueSize == 0);
385 KJ_ASSERT(observer.queueSizeBytes == 0);
386 
387 auto builder = kj::heapArrayBuilder<kj::Promise<void>>(2);
388 builder.add(write(1));
389 
390 KJ_ASSERT(observer.queueSize == 1);
391 KJ_ASSERT(observer.queueSizeBytes == 1);
392 
393 builder.add(write(10));
394 
395 KJ_ASSERT(observer.queueSize == 2);
396 KJ_ASSERT(observer.queueSizeBytes == 11);
397 
398 return kj::joinPromises(builder.finish());
399 });
400 
401 KJ_ASSERT(observer.queueSize == 0);
402 KJ_ASSERT(observer.queueSizeBytes == 0);
403}
404 
405// Test for use-after-free fix in pipeLoop when abort is called during pending read.
406// The fix ensures the Pipe::State is ref-counted and survives until all callbacks complete.
407KJ_TEST("WritableStreamInternalController pipeLoop abort during pending read") {
408 auto fixture = makeAbortClearsQueueTestFixture();
409 fixture.runInIoContext([&](const TestFixture::Environment& env) {
410 // Create a JavaScript-backed ReadableStream.
411 // The pull function will be called when the pipe tries to read.
412 // We use a JS-backed stream so that pipeLoop is used (not the kj pipe path).
413 //
414 // We need to simulate:
415 // 1. First read succeeds with some data
416 // 2. Second read is pending (the promise from pull is not resolved)
417 // 3. While pending, we abort the writable stream
418 //
419 // Using an UnderlyingSource with a pull callback that enqueues data once,
420 // then on the second call returns without enqueuing (leaving the read pending).
421 
422 int pullCount = 0;
423 jsg::Ref<ReadableStream> source = ReadableStream::constructor(env.js,
424 UnderlyingSource{.pull =
425 [&pullCount](jsg::Lock& js, UnderlyingSource::Controller controller) {
426 pullCount++;
427 auto& c = KJ_ASSERT_NONNULL(controller.tryGet<jsg::Ref<ReadableStreamDefaultController>>());
428 if (pullCount == 1) {
429 // First pull: enqueue some data so the pipe loop can make progress
430 auto data = js.bytes(kj::heapArray<kj::byte>({1, 2, 3, 4}));
431 c->enqueue(js, data.getHandle(js));
432 }
433 // Second pull onwards: don't enqueue anything, leaving the read pending.
434 // This simulates an async data source that hasn't received data yet.
435 // The promise returned by read() will be pending.
436 return js.resolvedPromise();
437 }},
438 kj::none);
439 
440 jsg::Ref<WritableStream> sink =
441 env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none);
442 
443 auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{});
444 pipeTo.markAsHandled(env.js);
445 env.js.runMicrotasks();
446 
447 // Abort while pipeLoop is waiting for a pending read
448 auto abortPromise = sink->getController().abort(env.js, env.js.v8TypeError("Test abort"_kj));
449 abortPromise.markAsHandled(env.js);
450 env.js.runMicrotasks();
451 
452 // If we get here without crashing, the test passes
453 KJ_ASSERT(pullCount >= 1);
454 });
455}
456 
457// ======================================================================================
458// DrainingReader tests for internal streams
459//
460// The internal stream implementation's drainingRead() behaves like a normal read() -
461// it returns at most one chunk at a time rather than draining all buffered data.
462// This is because internal streams are backed by kj I/O which is inherently async
463// and doesn't have internal JS-side buffering.
464 
465KJ_TEST("DrainingReader basic creation and locking (internal stream)") {
466 auto fixture = makeStreamTestFixture();
467 fixture.runInIoContext([&](const TestFixture::Environment& env) {
468 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>());
469 KJ_ASSERT(!rs->isLocked());
470 
471 KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) {
472 KJ_ASSERT(rs->isLocked());
473 KJ_ASSERT(reader->isAttached());
474 
475 reader->releaseLock(env.js);
476 KJ_ASSERT(!rs->isLocked());
477 KJ_ASSERT(!reader->isAttached());
478 } else {
479 KJ_FAIL_ASSERT("Failed to create DrainingReader");
480 }
481 });
482}
483 
484KJ_TEST("DrainingReader cannot be created on locked internal stream") {
485 auto fixture = makeStreamTestFixture();
486 fixture.runInIoContext([&](const TestFixture::Environment& env) {
487 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>());
488 
489 KJ_IF_SOME(reader1, DrainingReader::create(env.js, *rs)) {
490 KJ_ASSERT(rs->isLocked());
491 KJ_ASSERT(DrainingReader::create(env.js, *rs) == kj::none);
492 reader1->releaseLock(env.js);
493 } else {
494 KJ_FAIL_ASSERT("Failed to create first DrainingReader");
495 }
496 });
497}
498 
499KJ_TEST("DrainingReader read after releaseLock rejects (internal stream)") {
500 auto fixture = makeStreamTestFixture();
501 fixture.runInIoContext([&](const TestFixture::Environment& env) {
502 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>());
503 
504 KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) {
505 reader->releaseLock(env.js);
506 
507 bool rejected = false;
508 reader->read(env.js).catch_(env.js, [&](jsg::Lock&, jsg::Value) -> DrainingReadResult {
509 rejected = true;
510 return {.done = true};
511 });
512 env.js.runMicrotasks();
513 KJ_ASSERT(rejected);
514 } else {
515 KJ_FAIL_ASSERT("Failed to create DrainingReader");
516 }
517 });
518}
519 
520KJ_TEST("DrainingReader with maxRead parameter (internal stream)") {
521 // Test that the maxRead parameter is respected for internal streams
522 capnp::MallocMessageBuilder message;
523 auto flags = message.initRoot<CompatibilityFlags>();
524 flags.setStreamsJavaScriptControllers(true);
525 
526 TestFixture fixture({.featureFlags = flags.asReader()});
527 
528 bool testCompleted = false;
529 size_t lastMaxBytes = 0;
530 
531 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
532 class TestSource final: public ReadableStreamSource {
533 public:
534 explicit TestSource(size_t& maxBytesOut): lastMaxBytesOut(maxBytesOut) {}
535 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
536 readCount++;
537 // Note: maxBytes should be limited by the maxRead parameter
538 lastMaxBytesOut = maxBytes;
539 if (readCount == 1) {
540 // Return less than maxBytes
541 auto toWrite = kj::min(maxBytes, static_cast<size_t>(100));
542 memset(buffer, 'x', toWrite);
543 return toWrite;
544 }
545 return static_cast<size_t>(0); // EOF
546 }
547 uint readCount = 0;
548 size_t& lastMaxBytesOut;
549 };
550 
551 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>(lastMaxBytes));
552 
553 auto maybeReader = DrainingReader::create(env.js, *rs);
554 KJ_ASSERT(maybeReader != kj::none, "Failed to create DrainingReader");
555 auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader));
556 
557 // Pass a small maxRead value
558 auto readPromise = reader->read(env.js, 50);
559 
560 return env.context.awaitJs(env.js,
561 kj::mv(readPromise)
562 .then(env.js,
563 [&testCompleted, reader = kj::mv(reader)](
564 jsg::Lock& js, DrainingReadResult&& result) mutable {
565 KJ_ASSERT(result.chunks.size() == 1);
566 // The internal implementation uses maxRead to allocate the buffer
567 KJ_ASSERT(result.chunks[0].size() <= 50);
568 KJ_ASSERT(!result.done);
569 reader->releaseLock(js);
570 testCompleted = true;
571 }));
572 });
573 
574 KJ_ASSERT(testCompleted);
575 // Verify maxBytes was limited
576 KJ_ASSERT(lastMaxBytes == 50);
577}
578 
579KJ_TEST("DrainingReader with maxRead = 0 (internal stream)") {
580 // Test that the maxRead = 0 parameter is respected for internal streams
581 capnp::MallocMessageBuilder message;
582 auto flags = message.initRoot<CompatibilityFlags>();
583 flags.setStreamsJavaScriptControllers(true);
584 
585 TestFixture fixture({.featureFlags = flags.asReader()});
586 
587 bool testCompleted = false;
588 
589 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
590 class TestSource final: public ReadableStreamSource {
591 public:
592 explicit TestSource() = default;
593 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
594 KJ_FAIL_ASSERT("tryRead should not be called when maxRead = 0");
595 }
596 };
597 
598 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>());
599 
600 auto maybeReader = DrainingReader::create(env.js, *rs);
601 KJ_ASSERT(maybeReader != kj::none, "Failed to create DrainingReader");
602 auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader));
603 
604 auto readPromise = reader->read(env.js, 0);
605 
606 return env.context.awaitJs(env.js,
607 kj::mv(readPromise)
608 .then(env.js,
609 [&testCompleted, reader = kj::mv(reader)](
610 jsg::Lock& js, DrainingReadResult&& result) mutable {
611 KJ_ASSERT(result.chunks.size() == 0);
612 KJ_ASSERT(!result.done);
613 reader->releaseLock(js);
614 testCompleted = true;
615 }));
616 });
617 
618 KJ_ASSERT(testCompleted);
619}
620 
621KJ_TEST("DrainingReader on stream with pending closure (internal stream)") {
622 auto fixture = makeStreamTestFixture();
623 fixture.runInIoContext([&](const TestFixture::Environment& env) {
624 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>());
625 rs->getController().setPendingClosure();
626 
627 KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) {
628 bool rejected = false;
629 reader->read(env.js).catch_(
630 env.js, [&](jsg::Lock& js, jsg::Value reason) -> DrainingReadResult {
631 rejected = true;
632 auto msg = kj::str(reason.getHandle(js));
633 KJ_ASSERT(msg.contains("closing"), msg);
634 return {.done = true};
635 });
636 env.js.runMicrotasks();
637 KJ_ASSERT(rejected);
638 } else {
639 KJ_FAIL_ASSERT("Failed to create DrainingReader");
640 }
641 });
642}
643 
644KJ_TEST("DrainingReader on closed stream (internal stream)") {
645 auto fixture = makeStreamTestFixture();
646 
647 fixture.runInIoContext([&](const TestFixture::Environment& env) {
648 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>());
649 rs->getController().cancel(env.js, kj::none);
650 
651 KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) {
652 bool done = false;
653 reader->read(env.js).then(env.js, [&](jsg::Lock&, DrainingReadResult&& result) {
654 done = true;
655 KJ_ASSERT(result.done);
656 KJ_ASSERT(result.chunks.size() == 0);
657 });
658 env.js.runMicrotasks();
659 KJ_ASSERT(done);
660 } else {
661 KJ_FAIL_ASSERT("Failed to create DrainingReader");
662 }
663 });
664}
665 
666KJ_TEST("DrainingReader error propagation (internal stream)") {
667 // Tests that I/O errors from tryRead are propagated and put the stream in errored state.
668 auto fixture = makeStreamTestFixture();
669 
670 bool testCompleted = false;
671 KJ_EXPECT_LOG(ERROR, "Simulated I/O error");
672 
673 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
674 class TestSource final: public ReadableStreamSource {
675 public:
676 kj::Promise<size_t> tryRead(void*, size_t, size_t) override {
677 return KJ_EXCEPTION(FAILED, "Simulated I/O error");
678 }
679 };
680 
681 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>());
682 auto maybeReader = DrainingReader::create(env.js, *rs);
683 auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader));
684 
685 return env.context.awaitJs(env.js,
686 reader->read(env.js)
687 .catch_(env.js,
688 [&, reader = kj::mv(reader)](
689 jsg::Lock& js, jsg::Value) mutable -> DrainingReadResult {
690 reader->releaseLock(js);
691 testCompleted = true;
692 return {.done = true};
693 }).then(env.js, [](jsg::Lock&, DrainingReadResult&&) {}));
694 });
695 
696 KJ_ASSERT(testCompleted);
697}
698 
699KJ_TEST("DrainingReader concurrent read rejection (internal stream)") {
700 auto fixture = makeStreamTestFixture();
701 bool testCompleted = false;
702 
703 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
704 auto paf = kj::newPromiseAndFulfiller<size_t>();
705 
706 class TestSource final: public ReadableStreamSource {
707 public:
708 explicit TestSource(kj::Promise<size_t> p): promise(kj::mv(p)) {}
709 kj::Promise<size_t> tryRead(void*, size_t, size_t) override {
710 return kj::mv(promise);
711 }
712 kj::Promise<size_t> promise;
713 };
714 
715 auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>(kj::mv(paf.promise)));
716 auto maybeReader = DrainingReader::create(env.js, *rs);
717 auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader));
718 
719 auto firstRead = reader->read(env.js);
720 
721 // Second read while first is pending should reject synchronously
722 bool rejected = false;
723 reader->read(env.js).catch_(
724 env.js, [&](jsg::Lock& js, jsg::Value reason) -> DrainingReadResult {
725 rejected = true;
726 auto msg = kj::str(reason.getHandle(js));
727 KJ_ASSERT(msg.contains("single pending read"), msg);
728 return {.done = true};
729 });
730 env.js.runMicrotasks();
731 KJ_ASSERT(rejected);
732 
733 paf.fulfiller->fulfill(0); // Complete first read with EOF
734 
735 return env.context.awaitJs(env.js,
736 kj::mv(firstRead).then(
737 env.js, [&, reader = kj::mv(reader)](jsg::Lock& js, DrainingReadResult&&) mutable {
738 reader->releaseLock(js);
739 testCompleted = true;
740 }));
741 });
742 
743 KJ_ASSERT(testCompleted);
744}
745 
746// ======================================================================================
747// ReadableStreamBYOBReader validation tests
748 
749KJ_TEST("ReadableStreamBYOBReader rejects read with zero-sized buffer") {
750 KJ_EXPECT_LOG(ERROR, "read() on a BYOB reader requires a positive-sized TypedArray");
751 auto fixture = makeStreamTestFixture();
752 fixture.runInIoContext([&](const TestFixture::Environment& env) {
753 auto rs = makeByteStream(env.js);
754 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
755 
756 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 0);
757 auto view = v8::Uint8Array::New(buffer, 0, 0);
758 
759 bool rejected = false;
760 reader->read(env.js, view, kj::none)
761 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
762 rejected = true;
763 auto ex = js.exceptionToKj(kj::mv(reason));
764 KJ_ASSERT(ex.getDescription().contains(
765 "read() on a BYOB reader requires a positive-sized TypedArray"),
766 ex);
767 return {.done = true};
768 });
769 env.js.runMicrotasks();
770 KJ_ASSERT(rejected, "Expected read() to reject with zero-sized buffer");
771 });
772}
773 
774KJ_TEST("ReadableStreamBYOBReader rejects read with atLeast=0") {
775 auto fixture = makeStreamTestFixture();
776 fixture.runInIoContext([&](const TestFixture::Environment& env) {
777 auto rs = makeByteStream(env.js);
778 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
779 
780 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10);
781 auto view = v8::Uint8Array::New(buffer, 0, 10);
782 
783 bool rejected = false;
784 reader->readAtLeast(env.js, 0, view)
785 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
786 rejected = true;
787 auto ex = js.exceptionToKj(kj::mv(reason));
788 KJ_ASSERT(
789 ex.getDescription().contains("Requested invalid minimum number of bytes to read (0)"),
790 ex);
791 return {.done = true};
792 });
793 env.js.runMicrotasks();
794 KJ_ASSERT(rejected, "Expected readAtLeast() to reject with atLeast=0");
795 });
796}
797 
798KJ_TEST("ReadableStreamBYOBReader rejects read when atLeast exceeds buffer size") {
799 auto fixture = makeStreamTestFixture();
800 fixture.runInIoContext([&](const TestFixture::Environment& env) {
801 auto rs = makeByteStream(env.js);
802 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
803 
804 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10);
805 auto view = v8::Uint8Array::New(buffer, 0, 10);
806 
807 bool rejected = false;
808 reader->readAtLeast(env.js, 20, view)
809 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
810 rejected = true;
811 auto ex = js.exceptionToKj(kj::mv(reason));
812 KJ_ASSERT(
813 ex.getDescription().contains("Minimum bytes to read (20) exceeds size of buffer (10)"),
814 ex);
815 return {.done = true};
816 });
817 env.js.runMicrotasks();
818 KJ_ASSERT(rejected, "Expected readAtLeast() to reject when atLeast exceeds buffer size");
819 });
820}
821 
822KJ_TEST("ReadableStreamBYOBReader readAtLeast with element count within capacity succeeds") {
823 // readAtLeast() treats its first argument as an element count.
824 // readAtLeast(10, Uint32Array(10)) → 10 elements * 4 bytes = 40 bytes, buffer is 40 → OK.
825 auto fixture = makeStreamTestFixture();
826 fixture.runInIoContext([&](const TestFixture::Environment& env) {
827 auto rs = makeByteStream(env.js);
828 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
829 
830 // Uint32Array: element size 4, byteLength 40, length 10
831 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40);
832 auto view = v8::Uint32Array::New(buffer, 0, 10);
833 
834 bool rejected = false;
835 reader->readAtLeast(env.js, 10, view)
836 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
837 rejected = true;
838 auto ex = js.exceptionToKj(kj::mv(reason));
839 KJ_FAIL_ASSERT("readAtLeast(10) on 10-element Uint32Array should not reject", ex);
840 return {.done = true};
841 });
842 env.js.runMicrotasks();
843 KJ_ASSERT(!rejected);
844 });
845}
846 
847KJ_TEST("ReadableStreamBYOBReader readAtLeast rejects when element count exceeds capacity") {
848 // Regression test for the bug: before the fix, validation compared
849 // the raw element count against byteLength (mixed units), so large element
850 // counts slipped through and produced an impossible (minBytes > maxBytes) pair
851 // that caused stack overflow in decoders like brotli.
852 // readAtLeast(11, Uint32Array(10)) → 11 * 4 = 44 bytes > 40 byte buffer → reject.
853 auto fixture = makeStreamTestFixture();
854 fixture.runInIoContext([&](const TestFixture::Environment& env) {
855 auto rs = makeByteStream(env.js);
856 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
857 
858 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40);
859 auto view = v8::Uint32Array::New(buffer, 0, 10);
860 
861 bool rejected = false;
862 reader->readAtLeast(env.js, 11, view)
863 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
864 rejected = true;
865 auto ex = js.exceptionToKj(kj::mv(reason));
866 KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex);
867 return {.done = true};
868 });
869 env.js.runMicrotasks();
870 KJ_ASSERT(rejected, "readAtLeast(11) on 10-element Uint32Array should reject");
871 });
872}
873 
874KJ_TEST("ReadableStreamBYOBReader readAtLeast rejects byteLength as element count") {
875 // readAtLeast(view.byteLength, view) with a Uint32Array.
876 // byteLength=4096, element count interpretation → 4096 * 4 = 16384 > 4096 → must reject.
877 auto fixture = makeStreamTestFixture();
878 fixture.runInIoContext([&](const TestFixture::Environment& env) {
879 auto rs = makeByteStream(env.js);
880 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
881 
882 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 4096);
883 auto view = v8::Uint32Array::New(buffer, 0, 1024);
884 
885 bool rejected = false;
886 reader->readAtLeast(env.js, 4096, view)
887 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
888 rejected = true;
889 auto ex = js.exceptionToKj(kj::mv(reason));
890 KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex);
891 return {.done = true};
892 });
893 env.js.runMicrotasks();
894 KJ_ASSERT(rejected, "readAtLeast(4096) on 1024-element Uint32Array must reject");
895 });
896}
897 
898KJ_TEST("ReadableStreamBYOBReader read() with min exceeding element capacity rejects") {
899 // read(view, {min: N}) where N is in elements. Before the fix, validation
900 // compared element count against byte length (mixed units), so large values
901 // slipped through.
902 // min=11 elements on Uint32Array(10) → 44 bytes > 40 → reject.
903 auto fixture = makeStreamTestFixture();
904 fixture.runInIoContext([&](const TestFixture::Environment& env) {
905 auto rs = makeByteStream(env.js);
906 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
907 
908 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40);
909 auto view = v8::Uint32Array::New(buffer, 0, 10);
910 
911 ReadableStreamBYOBReader::ReadableStreamBYOBReaderReadOptions opts;
912 opts.min = 11;
913 bool rejected = false;
914 reader->read(env.js, view, kj::mv(opts))
915 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
916 rejected = true;
917 auto ex = js.exceptionToKj(kj::mv(reason));
918 KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex);
919 return {.done = true};
920 });
921 env.js.runMicrotasks();
922 KJ_ASSERT(rejected, "read() with min=11 on 10-element Uint32Array should reject");
923 });
924}
925 
926KJ_TEST("ReadableStreamBYOBReader rejects read after releaseLock") {
927 auto fixture = makeStreamTestFixture();
928 fixture.runInIoContext([&](const TestFixture::Environment& env) {
929 auto rs = makeByteStream(env.js);
930 auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef());
931 reader->releaseLock(env.js);
932 
933 auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10);
934 auto view = v8::Uint8Array::New(buffer, 0, 10);
935 
936 bool rejected = false;
937 reader->read(env.js, view, kj::none)
938 .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult {
939 rejected = true;
940 auto ex = js.exceptionToKj(kj::mv(reason));
941 KJ_ASSERT(ex.getDescription().contains("This ReadableStream reader has been released"), ex);
942 return {.done = true};
943 });
944 env.js.runMicrotasks();
945 KJ_ASSERT(rejected, "Expected read() to reject after releaseLock");
946 });
947}
948 
949} // namespace
950} // namespace workerd::api