Skip to content
File

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

99.3 KB
1#include "readable.h"
2#include "standard.h"
3#include "writable.h"
4 
5#include <workerd/jsg/jsg-test.h>
6#include <workerd/jsg/jsg.h>
7#include <workerd/jsg/observer.h>
8#include <workerd/tests/test-fixture.h>
9 
10namespace workerd::api {
11namespace {
12 
13void preamble(auto callback) {
14 TestFixture fixture;
15 fixture.runInIoContext([&](const TestFixture::Environment& env) { callback(env.js); });
16}
17 
18v8::Local<v8::Value> toBytes(jsg::Lock& js, kj::String str) {
19 return jsg::BackingStore::from(js, str.asBytes().attach(kj::mv(str))).createHandle(js);
20}
21 
22jsg::BufferSource toBufferSource(jsg::Lock& js, kj::String str) {
23 auto backing = jsg::BackingStore::from(js, str.asBytes().attach(kj::mv(str))).createHandle(js);
24 return jsg::BufferSource(js, kj::mv(backing));
25}
26 
27jsg::BufferSource toBufferSource(jsg::Lock& js, kj::Array<kj::byte> bytes) {
28 auto backing = jsg::BackingStore::from(js, kj::mv(bytes)).createHandle(js);
29 return jsg::BufferSource(js, kj::mv(backing));
30}
31 
32// ======================================================================================
33// Happy Cases
34 
35KJ_TEST("ReadableStream read all text (value readable)") {
36 preamble([](jsg::Lock& js) {
37 uint checked = 0;
38 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
39 // clang-format off
40 rs->getController().setup(js, UnderlyingSource{
41 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
42 // Because we're using a value-based stream, two enqueue operations will
43 // require at least three reads to complete: one for the first chunk, 'hello, ',
44 // one for the second chunk, 'world!', and one to signal close.
45 KJ_SWITCH_ONEOF(controller) {
46 // Because we're using a value-based stream, two enqueue operations will
47 // require at least three reads to complete: one for the first chunk, 'hello, ',
48 // one for the second chunk, 'world!', and one to signal close.
49 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
50 checked++;
51 c->enqueue(js, toBytes(js, kj::str("Hello, ")));
52 c->enqueue(js, toBytes(js, kj::str("world!")));
53 c->close(js);
54 return js.resolvedPromise();
55 }
56 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
57 }
58 KJ_UNREACHABLE;
59 }
60 // Setting a highWaterMark of 0 means the pull function above will not be called
61 // immediately on creation of the stream, but only when the first read in the
62 // readall call below happens.
63 }, StreamQueuingStrategy{.highWaterMark = 0});
64 // clang-format on
65 
66 // Starts a read loop of javascript promises.
67 auto promise =
68 rs->getController().readAllText(js, 20).then(js, [&](jsg::Lock& js, kj::String&& text) {
69 KJ_ASSERT(text == "Hello, world!"_kjc);
70 checked++;
71 });
72 
73 // Reading left the stream locked and disturbed
74 KJ_ASSERT(rs->isLocked());
75 KJ_ASSERT(rs->isDisturbed());
76 
77 // Run the microtasks to completion. This should resolve the promise and
78 // run it to completion. The test is buggy if it fails to do so.
79 js.runMicrotasks();
80 KJ_ASSERT(checked == 2);
81 
82 // Reading everything successfully should cause the stream to close.
83 KJ_ASSERT(rs->getController().isClosed());
84 
85 // Add we should still be locked and disturbed.
86 KJ_ASSERT(rs->isLocked());
87 KJ_ASSERT(rs->isDisturbed());
88 });
89}
90 
91KJ_TEST("ReadableStream read all text, rs ref held (value readable)") {
92 preamble([](jsg::Lock& js) {
93 uint checked = 0;
94 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
95 // clang-format off
96 rs->getController().setup(js, UnderlyingSource{
97 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
98 // Because we're using a value-based stream, two enqueue operations will
99 // require at least three reads to complete: one for the first chunk, 'hello, ',
100 // one for the second chunk, 'world!', and one to signal close.
101 KJ_SWITCH_ONEOF(controller) {
102 // Because we're using a value-based stream, two enqueue operations will
103 // require at least three reads to complete: one for the first chunk, 'hello, ',
104 // one for the second chunk, 'world!', and one to signal close.
105 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
106 checked++;
107 c->enqueue(js, toBytes(js, kj::str("Hello, ")));
108 c->enqueue(js, toBytes(js, kj::str("world!")));
109 c->close(js);
110 return js.resolvedPromise();
111 }
112 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
113 }
114 KJ_UNREACHABLE;
115 }
116 // Setting a highWaterMark of 0 means the pull function above will not be called
117 // immediately on creation of the stream, but only when the first read in the
118 // readall call below happens.
119 }, StreamQueuingStrategy{.highWaterMark = 0});
120 // clang-format on
121 
122 // Starts a read loop of javascript promises.
123 auto promise =
124 rs->getController().readAllText(js, 20).then(js, [&](jsg::Lock& js, kj::String&& text) {
125 KJ_ASSERT(text == "Hello, world!"_kjc);
126 checked++;
127 });
128 
129 // Reading left the stream locked and disturbed
130 KJ_ASSERT(rs->isLocked());
131 KJ_ASSERT(rs->isDisturbed());
132 
133 // Let's drop our ref to rs, things should still work as expected.
134 { auto drop = kj::mv(rs); }
135 
136 // Run the microtasks to completion. This should resolve the promise and
137 // run it to completion. The test is buggy if it fails to do so.
138 js.runMicrotasks();
139 KJ_ASSERT(checked == 2);
140 });
141}
142 
143KJ_TEST("ReadableStream read all text (byte readable)") {
144 preamble([](jsg::Lock& js) {
145 uint checked = 0;
146 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
147 // clang-format off
148 rs->getController().setup(js, UnderlyingSource{
149 .type = kj::str("bytes"),
150 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
151 // Because we're using a value-based stream, two enqueue operations will
152 // require at least three reads to complete: one for the first chunk, 'hello, ',
153 // one for the second chunk, 'world!', and one to signal close.
154 KJ_SWITCH_ONEOF(controller) {
155 // Because we're using a value-based stream, two enqueue operations will
156 // require at least three reads to complete: one for the first chunk, 'hello, ',
157 // one for the second chunk, 'world!', and one to signal close.
158 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
159 checked++;
160 c->enqueue(js, toBufferSource(js, kj::str("Hello, ")));
161 c->enqueue(js, toBufferSource(js, kj::str("world!")));
162 c->close(js);
163 return js.resolvedPromise();
164 }
165 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
166 }
167 KJ_UNREACHABLE;
168 }
169 // Setting a highWaterMark of 0 means the pull function above will not be called
170 // immediately on creation of the stream, but only when the first read in the
171 // readall call below happens.
172 }, StreamQueuingStrategy{.highWaterMark = 0});
173 // clang-format on
174 
175 // Starts a read loop of javascript promises.
176 auto promise =
177 rs->getController().readAllText(js, 20).then(js, [&](jsg::Lock& js, kj::String&& text) {
178 KJ_ASSERT(text == "Hello, world!"_kjc);
179 checked++;
180 });
181 
182 // Reading left the stream locked and disturbed
183 KJ_ASSERT(rs->isLocked());
184 KJ_ASSERT(rs->isDisturbed());
185 
186 // Run the microtasks to completion. This should resolve the promise and
187 // run it to completion. The test is buggy if it fails to do so.
188 js.runMicrotasks();
189 KJ_ASSERT(checked == 2);
190 
191 // Reading everything successfully should cause the stream to close.
192 KJ_ASSERT(rs->getController().isClosed());
193 
194 // Add we should still be locked and disturbed.
195 KJ_ASSERT(rs->isLocked());
196 KJ_ASSERT(rs->isDisturbed());
197 });
198}
199 
200KJ_TEST("ReadableStream read all bytes (value readable)") {
201 preamble([](jsg::Lock& js) {
202 uint checked = 0;
203 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
204 // clang-format off
205 rs->getController().setup(js, UnderlyingSource{
206 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
207 // Because we're using a value-based stream, two enqueue operations will
208 // require at least three reads to complete: one for the first chunk, 'hello, ',
209 // one for the second chunk, 'world!', and one to signal close.
210 KJ_SWITCH_ONEOF(controller) {
211 // Because we're using a value-based stream, two enqueue operations will
212 // require at least three reads to complete: one for the first chunk, 'hello, ',
213 // one for the second chunk, 'world!', and one to signal close.
214 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
215 checked++;
216 c->enqueue(js, toBytes(js, kj::str("Hello, ")));
217 c->enqueue(js, toBytes(js, kj::str("world!")));
218 c->close(js);
219 return js.resolvedPromise();
220 }
221 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
222 }
223 KJ_UNREACHABLE;
224 }
225 // Setting a highWaterMark of 0 means the pull function above will not be called
226 // immediately on creation of the stream, but only when the first read in the
227 // readall call below happens.
228 }, StreamQueuingStrategy{.highWaterMark = 0});
229 // clang-format on
230 
231 // Starts a read loop of javascript promises.
232 auto promise = rs->getController().readAllBytes(js, 20).then(
233 js, [&](jsg::Lock& js, jsg::BufferSource&& text) {
234 KJ_ASSERT(text.asArrayPtr() == "Hello, world!"_kjb);
235 checked++;
236 });
237 
238 // Reading left the stream locked and disturbed
239 KJ_ASSERT(rs->isLocked());
240 KJ_ASSERT(rs->isDisturbed());
241 
242 // Run the microtasks to completion. This should resolve the promise and
243 // run it to completion. The test is buggy if it fails to do so.
244 js.runMicrotasks();
245 KJ_ASSERT(checked == 2);
246 
247 // Reading everything successfully should cause the stream to close.
248 KJ_ASSERT(rs->getController().isClosed());
249 
250 // Add we should still be locked and disturbed.
251 KJ_ASSERT(rs->isLocked());
252 KJ_ASSERT(rs->isDisturbed());
253 });
254}
255 
256KJ_TEST("ReadableStream read all bytes (byte readable)") {
257 preamble([](jsg::Lock& js) {
258 uint checked = 0;
259 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
260 // clang-format off
261 rs->getController().setup(js, UnderlyingSource{
262 .type = kj::str("bytes"),
263 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
264 // Because we're using a value-based stream, two enqueue operations will
265 // require at least three reads to complete: one for the first chunk, 'hello, ',
266 // one for the second chunk, 'world!', and one to signal close.
267 KJ_SWITCH_ONEOF(controller) {
268 // Because we're using a value-based stream, two enqueue operations will
269 // require at least three reads to complete: one for the first chunk, 'hello, ',
270 // one for the second chunk, 'world!', and one to signal close.
271 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
272 checked++;
273 c->enqueue(js, toBufferSource(js, kj::str("Hello, ")));
274 c->enqueue(js, toBufferSource(js, kj::str("world!")));
275 c->close(js);
276 return js.resolvedPromise();
277 }
278 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
279 }
280 KJ_UNREACHABLE;
281 }
282 // Setting a highWaterMark of 0 means the pull function above will not be called
283 // immediately on creation of the stream, but only when the first read in the
284 // readall call below happens.
285 }, StreamQueuingStrategy{.highWaterMark = 0});
286 // clang-format on
287 
288 // Starts a read loop of javascript promises.
289 auto promise = rs->getController().readAllBytes(js, 20).then(
290 js, [&](jsg::Lock& js, jsg::BufferSource&& text) {
291 KJ_ASSERT(text.asArrayPtr() == "Hello, world!"_kjb);
292 checked++;
293 });
294 
295 // Reading left the stream locked and disturbed
296 KJ_ASSERT(rs->isLocked());
297 KJ_ASSERT(rs->isDisturbed());
298 
299 // Run the microtasks to completion. This should resolve the promise and
300 // run it to completion. The test is buggy if it fails to do so.
301 js.runMicrotasks();
302 KJ_ASSERT(checked == 2);
303 
304 // Reading everything successfully should cause the stream to close.
305 KJ_ASSERT(rs->getController().isClosed());
306 
307 // Add we should still be locked and disturbed.
308 KJ_ASSERT(rs->isLocked());
309 KJ_ASSERT(rs->isDisturbed());
310 });
311}
312 
313KJ_TEST("ReadableStream read all bytes (value readable, more reads)") {
314 preamble([](jsg::Lock& js) {
315 uint checked = 0;
316 uint counter = 0;
317 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
318 auto chunks = kj::arr<kj::String>(kj::str("H"), kj::str("e"), kj::str("l"), kj::str("l"),
319 kj::str("o"), kj::str(","), kj::str(" "), kj::str("w"), kj::str("o"), kj::str("r"),
320 kj::str("l"), kj::str("d"), kj::str("!"));
321 // clang-format off
322 rs->getController().setup(js, UnderlyingSource{
323 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
324 // Because we're using a value-based stream, two enqueue operations will
325 // require at least three reads to complete: one for the first chunk, 'hello, ',
326 // one for the second chunk, 'world!', and one to signal close.
327 KJ_SWITCH_ONEOF(controller) {
328 // Because we're using a value-based stream, two enqueue operations will
329 // require at least three reads to complete: one for the first chunk, 'hello, ',
330 // one for the second chunk, 'world!', and one to signal close.
331 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
332 checked++;
333 c->enqueue(js, toBytes(js, kj::mv(chunks[counter++])));
334 if (counter == chunks.size()) {
335 c->close(js);
336 }
337 
338 return js.resolvedPromise();
339 }
340 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
341 }
342 KJ_UNREACHABLE;
343 }
344 // Setting a highWaterMark of 0 means the pull function above will not be called
345 // immediately on creation of the stream, but only when the first read in the
346 // readall call below happens.
347 }, StreamQueuingStrategy{.highWaterMark = 0});
348 // clang-format on
349 
350 // Starts a read loop of javascript promises.
351 auto promise = rs->getController().readAllBytes(js, 20).then(
352 js, [&](jsg::Lock& js, jsg::BufferSource&& text) {
353 KJ_ASSERT(text.asArrayPtr() == "Hello, world!"_kjb);
354 checked++;
355 });
356 
357 // Reading left the stream locked and disturbed
358 KJ_ASSERT(rs->isLocked());
359 KJ_ASSERT(rs->isDisturbed());
360 
361 // Run the microtasks to completion. This should resolve the promise and
362 // run it to completion. The test is buggy if it fails to do so.
363 js.runMicrotasks();
364 KJ_ASSERT(checked == 14);
365 
366 // Reading everything successfully should cause the stream to close.
367 KJ_ASSERT(rs->getController().isClosed());
368 
369 // Add we should still be locked and disturbed.
370 KJ_ASSERT(rs->isLocked());
371 KJ_ASSERT(rs->isDisturbed());
372 });
373}
374 
375KJ_TEST("ReadableStream read all bytes (byte readable, more reads)") {
376 preamble([](jsg::Lock& js) {
377 uint checked = 0;
378 uint counter = 0;
379 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
380 auto chunks = kj::arr<kj::String>(kj::str("H"), kj::str("e"), kj::str("l"), kj::str("l"),
381 kj::str("o"), kj::str(","), kj::str(" "), kj::str("w"), kj::str("o"), kj::str("r"),
382 kj::str("l"), kj::str("d"), kj::str("!"));
383 // clang-format off
384 rs->getController().setup(js, UnderlyingSource{
385 .type = kj::str("bytes"),
386 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
387 // Because we're using a value-based stream, two enqueue operations will
388 // require at least three reads to complete: one for the first chunk, 'hello, ',
389 // one for the second chunk, 'world!', and one to signal close.
390 KJ_SWITCH_ONEOF(controller) {
391 // Because we're using a value-based stream, two enqueue operations will
392 // require at least three reads to complete: one for the first chunk, 'hello, ',
393 // one for the second chunk, 'world!', and one to signal close.
394 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
395 checked++;
396 c->enqueue(js, toBufferSource(js, kj::mv(chunks[counter++])));
397 if (counter == chunks.size()) {
398 c->close(js);
399 }
400 
401 return js.resolvedPromise();
402 }
403 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
404 }
405 KJ_UNREACHABLE;
406 }
407 // Setting a highWaterMark of 0 means the pull function above will not be called
408 // immediately on creation of the stream, but only when the first read in the
409 // readall call below happens.
410 }, StreamQueuingStrategy{.highWaterMark = 0});
411 // clang-format on
412 
413 // Starts a read loop of javascript promises.
414 auto promise = rs->getController().readAllBytes(js, 20).then(
415 js, [&](jsg::Lock& js, jsg::BufferSource&& text) {
416 KJ_ASSERT(text.asArrayPtr() == "Hello, world!"_kjb);
417 checked++;
418 });
419 
420 // Reading left the stream locked and disturbed
421 KJ_ASSERT(rs->isLocked());
422 KJ_ASSERT(rs->isDisturbed());
423 
424 // Run the microtasks to completion. This should resolve the promise and
425 // run it to completion. The test is buggy if it fails to do so.
426 js.runMicrotasks();
427 KJ_ASSERT(checked == 14);
428 
429 // Reading everything successfully should cause the stream to close.
430 KJ_ASSERT(rs->getController().isClosed());
431 
432 // Add we should still be locked and disturbed.
433 KJ_ASSERT(rs->isLocked());
434 KJ_ASSERT(rs->isDisturbed());
435 });
436}
437 
438KJ_TEST("ReadableStream read all bytes (byte readable, large data)") {
439 preamble([](jsg::Lock& js) {
440 uint checked = 0;
441 uint counter = 0;
442 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
443 static constexpr uint BASE = 4097;
444 auto chunks = kj::arr<kj::Array<kj::byte>>(kj::heapArray<kj::byte>(BASE),
445 kj::heapArray<kj::byte>(BASE * 2), kj::heapArray<kj::byte>(BASE * 4));
446 chunks[0].asPtr().fill('A');
447 chunks[1].asPtr().fill('B');
448 chunks[2].asPtr().fill('C');
449 // clang-format off
450 rs->getController().setup(js, UnderlyingSource{
451 .type = kj::str("bytes"),
452 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
453 // Because we're using a value-based stream, two enqueue operations will
454 // require at least three reads to complete: one for the first chunk, 'hello, ',
455 // one for the second chunk, 'world!', and one to signal close.
456 KJ_SWITCH_ONEOF(controller) {
457 // Because we're using a value-based stream, two enqueue operations will
458 // require at least three reads to complete: one for the first chunk, 'hello, ',
459 // one for the second chunk, 'world!', and one to signal close.
460 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
461 checked++;
462 c->enqueue(js, toBufferSource(js, kj::mv(chunks[counter++])));
463 if (counter == chunks.size()) {
464 c->close(js);
465 }
466 
467 return js.resolvedPromise();
468 }
469 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
470 }
471 KJ_UNREACHABLE;
472 }
473 // Setting a highWaterMark of 0 means the pull function above will not be called
474 // immediately on creation of the stream, but only when the first read in the
475 // readall call below happens.
476 }, StreamQueuingStrategy{.highWaterMark = 0});
477 // clang-format on
478 
479 // Starts a read loop of javascript promises.
480 auto promise = rs->getController()
481 .readAllBytes(js, (BASE * 7) + 1)
482 .then(js, [&](jsg::Lock& js, jsg::BufferSource&& text) {
483 kj::byte check[BASE * 7]{};
484 kj::arrayPtr(check).first(BASE).fill('A');
485 kj::arrayPtr(check).slice(BASE).first(BASE * 2).fill('B');
486 kj::arrayPtr(check).slice(BASE * 3).fill('C');
487 KJ_ASSERT(text.size() == BASE * 7);
488 KJ_ASSERT(text.asArrayPtr() == check);
489 checked++;
490 });
491 
492 // Reading left the stream locked and disturbed
493 KJ_ASSERT(rs->isLocked());
494 KJ_ASSERT(rs->isDisturbed());
495 
496 // Run the microtasks to completion. This should resolve the promise and
497 // run it to completion. The test is buggy if it fails to do so.
498 js.runMicrotasks();
499 KJ_ASSERT(checked == 4);
500 
501 // Reading everything successfully should cause the stream to close.
502 KJ_ASSERT(rs->getController().isClosed());
503 
504 // Add we should still be locked and disturbed.
505 KJ_ASSERT(rs->isLocked());
506 KJ_ASSERT(rs->isDisturbed());
507 });
508}
509 
510// ======================================================================================
511// Fail cases
512 
513KJ_TEST("ReadableStream read all bytes (value readable, wrong type)") {
514 preamble([](jsg::Lock& js) {
515 uint checked = 0;
516 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
517 // clang-format off
518 rs->getController().setup(js, UnderlyingSource{
519 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
520 // Because we're using a value-based stream, two enqueue operations will
521 // require at least three reads to complete: one for the first chunk, 'hello, ',
522 // one for the second chunk, 'world!', and one to signal close.
523 KJ_SWITCH_ONEOF(controller) {
524 // Because we're using a value-based stream, two enqueue operations will
525 // require at least three reads to complete: one for the first chunk, 'hello, ',
526 // one for the second chunk, 'world!', and one to signal close.
527 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
528 c->enqueue(js, js.str("wrong type"_kjc));
529 checked++;
530 return js.resolvedPromise();
531 }
532 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
533 }
534 KJ_UNREACHABLE;
535 },
536 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
537 KJ_ASSERT(kj::str(reason) == "TypeError: This ReadableStream did not return bytes.");
538 checked++;
539 return js.resolvedPromise();
540 }
541 // Setting a highWaterMark of 0 means the pull function above will not be called
542 // immediately on creation of the stream, but only when the first read in the
543 // readall call below happens.
544 }, StreamQueuingStrategy{.highWaterMark = 0});
545 // clang-format on
546 
547 // Starts a read loop of javascript promises.
548 auto promise = rs->getController().readAllBytes(js, 20).then(js,
549 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
550 [&](jsg::Lock& js, jsg::Value&& exception) {
551 KJ_ASSERT(kj::str(exception.getHandle(js)) ==
552 "TypeError: This ReadableStream did not return bytes.");
553 checked++;
554 });
555 
556 // Reading left the stream locked and disturbed
557 KJ_ASSERT(rs->isLocked());
558 KJ_ASSERT(rs->isDisturbed());
559 
560 // Run the microtasks to completion. This should resolve the promise and
561 // run it to completion. The test is buggy if it fails to do so.
562 js.runMicrotasks();
563 KJ_ASSERT(checked == 3);
564 
565 KJ_ASSERT(rs->getController().isClosedOrErrored());
566 
567 // Add we should still be locked and disturbed.
568 KJ_ASSERT(rs->isLocked());
569 KJ_ASSERT(rs->isDisturbed());
570 });
571}
572 
573KJ_TEST("ReadableStream read all bytes (value readable, to many bytes)") {
574 preamble([](jsg::Lock& js) {
575 uint checked = 0;
576 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
577 // clang-format off
578 rs->getController().setup(js, UnderlyingSource{
579 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
580 // Because we're using a value-based stream, two enqueue operations will
581 // require at least three reads to complete: one for the first chunk, 'hello, ',
582 // one for the second chunk, 'world!', and one to signal close.
583 KJ_SWITCH_ONEOF(controller) {
584 // Because we're using a value-based stream, two enqueue operations will
585 // require at least three reads to complete: one for the first chunk, 'hello, ',
586 // one for the second chunk, 'world!', and one to signal close.
587 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
588 c->enqueue(js, toBytes(js, kj::str("123456789012345678901")));
589 checked++;
590 return js.resolvedPromise();
591 }
592 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
593 }
594 KJ_UNREACHABLE;
595 }
596 // Setting a highWaterMark of 0 means the pull function above will not be called
597 // immediately on creation of the stream, but only when the first read in the
598 // readall call below happens.
599 }, StreamQueuingStrategy{.highWaterMark = 0});
600 // clang-format on
601 
602 // Starts a read loop of javascript promises.
603 auto promise = rs->getController().readAllBytes(js, 20).then(js,
604 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
605 [&](jsg::Lock& js, jsg::Value&& exception) {
606 KJ_ASSERT(kj::str(exception.getHandle(js)) == "TypeError: Memory limit exceeded before EOF.");
607 checked++;
608 });
609 
610 // Reading left the stream locked and disturbed
611 KJ_ASSERT(rs->isLocked());
612 KJ_ASSERT(rs->isDisturbed());
613 
614 // Run the microtasks to completion. This should resolve the promise and
615 // run it to completion. The test is buggy if it fails to do so.
616 js.runMicrotasks();
617 KJ_ASSERT(checked == 2);
618 
619 KJ_ASSERT(rs->getController().isClosedOrErrored());
620 
621 // Add we should still be locked and disturbed.
622 KJ_ASSERT(rs->isLocked());
623 KJ_ASSERT(rs->isDisturbed());
624 });
625}
626 
627KJ_TEST("ReadableStream read all bytes (byte readable, to many bytes)") {
628 preamble([](jsg::Lock& js) {
629 uint checked = 0;
630 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
631 // clang-format off
632 rs->getController().setup(js, UnderlyingSource{
633 .type = kj::str("bytes"),
634 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
635 // Because we're using a value-based stream, two enqueue operations will
636 // require at least three reads to complete: one for the first chunk, 'hello, ',
637 // one for the second chunk, 'world!', and one to signal close.
638 KJ_SWITCH_ONEOF(controller) {
639 // Because we're using a value-based stream, two enqueue operations will
640 // require at least three reads to complete: one for the first chunk, 'hello, ',
641 // one for the second chunk, 'world!', and one to signal close.
642 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
643 c->enqueue(js, toBufferSource(js, kj::str("123456789012345678901")));
644 checked++;
645 return js.resolvedPromise();
646 }
647 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
648 }
649 KJ_UNREACHABLE;
650 }
651 // Setting a highWaterMark of 0 means the pull function above will not be called
652 // immediately on creation of the stream, but only when the first read in the
653 // readall call below happens.
654 }, StreamQueuingStrategy{.highWaterMark = 0});
655 // clang-format on
656 
657 // Starts a read loop of javascript promises.
658 auto promise = rs->getController().readAllBytes(js, 20).then(js,
659 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
660 [&](jsg::Lock& js, jsg::Value&& exception) {
661 KJ_ASSERT(kj::str(exception.getHandle(js)) == "TypeError: Memory limit exceeded before EOF.");
662 checked++;
663 });
664 
665 // Reading left the stream locked and disturbed
666 KJ_ASSERT(rs->isLocked());
667 KJ_ASSERT(rs->isDisturbed());
668 
669 // Run the microtasks to completion. This should resolve the promise and
670 // run it to completion. The test is buggy if it fails to do so.
671 js.runMicrotasks();
672 KJ_ASSERT(checked == 2);
673 
674 KJ_ASSERT(rs->getController().isClosedOrErrored());
675 
676 // Add we should still be locked and disturbed.
677 KJ_ASSERT(rs->isLocked());
678 KJ_ASSERT(rs->isDisturbed());
679 });
680}
681 
682KJ_TEST("ReadableStream read all bytes (byte readable, failed read)") {
683 preamble([](jsg::Lock& js) {
684 uint checked = 0;
685 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
686 // clang-format off
687 rs->getController().setup(js, UnderlyingSource{
688 .type = kj::str("bytes"),
689 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
690 checked++;
691 return js.rejectedPromise<void>(js.error("boom"));
692 }
693 // Setting a highWaterMark of 0 means the pull function above will not be called
694 // immediately on creation of the stream, but only when the first read in the
695 // readall call below happens.
696 }, StreamQueuingStrategy{.highWaterMark = 0});
697 // clang-format on
698 
699 // Starts a read loop of javascript promises.
700 auto promise = rs->getController().readAllBytes(js, 20).then(js,
701 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
702 [&](jsg::Lock& js, jsg::Value&& exception) {
703 KJ_ASSERT(kj::str(exception.getHandle(js)) == "Error: boom");
704 checked++;
705 });
706 
707 // Reading left the stream locked and disturbed
708 KJ_ASSERT(rs->isLocked());
709 KJ_ASSERT(rs->isDisturbed());
710 
711 // Run the microtasks to completion. This should resolve the promise and
712 // run it to completion. The test is buggy if it fails to do so.
713 js.runMicrotasks();
714 KJ_ASSERT(checked == 2);
715 
716 KJ_ASSERT(rs->getController().isClosedOrErrored());
717 
718 // Add we should still be locked and disturbed.
719 KJ_ASSERT(rs->isLocked());
720 KJ_ASSERT(rs->isDisturbed());
721 });
722}
723 
724KJ_TEST("ReadableStream read all bytes (value readable, failed read)") {
725 preamble([](jsg::Lock& js) {
726 uint checked = 0;
727 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
728 // clang-format off
729 rs->getController().setup(js, UnderlyingSource{
730 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
731 checked++;
732 return js.rejectedPromise<void>(js.error("boom"));
733 }
734 // Setting a highWaterMark of 0 means the pull function above will not be called
735 // immediately on creation of the stream, but only when the first read in the
736 // readall call below happens.
737 }, StreamQueuingStrategy{.highWaterMark = 0});
738 // clang-format on
739 
740 // Starts a read loop of javascript promises.
741 auto promise = rs->getController().readAllBytes(js, 20).then(js,
742 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
743 [&](jsg::Lock& js, jsg::Value&& exception) {
744 KJ_ASSERT(kj::str(exception.getHandle(js)) == "Error: boom");
745 checked++;
746 });
747 
748 // Reading left the stream locked and disturbed
749 KJ_ASSERT(rs->isLocked());
750 KJ_ASSERT(rs->isDisturbed());
751 
752 // Run the microtasks to completion. This should resolve the promise and
753 // run it to completion. The test is buggy if it fails to do so.
754 js.runMicrotasks();
755 KJ_ASSERT(checked == 2);
756 
757 KJ_ASSERT(rs->getController().isClosedOrErrored());
758 
759 // Add we should still be locked and disturbed.
760 KJ_ASSERT(rs->isLocked());
761 KJ_ASSERT(rs->isDisturbed());
762 });
763}
764 
765KJ_TEST("ReadableStream read all bytes (byte readable, failed start)") {
766 preamble([](jsg::Lock& js) {
767 uint checked = 0;
768 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
769 // clang-format off
770 rs->getController().setup(js, UnderlyingSource{
771 .type = kj::str("bytes"),
772 .start = [&](jsg::Lock& js, UnderlyingSource::Controller controller) -> jsg::Promise<void> {
773 checked++;
774 return js.rejectedPromise<void>(js.error("boom"));
775 }
776 // Setting a highWaterMark of 0 means the pull function above will not be called
777 // immediately on creation of the stream, but only when the first read in the
778 // readall call below happens.
779 }, StreamQueuingStrategy{.highWaterMark = 0});
780 // clang-format on
781 
782 // Starts a read loop of javascript promises.
783 auto promise = rs->getController().readAllBytes(js, 20).then(js,
784 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
785 [&](jsg::Lock& js, jsg::Value&& exception) {
786 KJ_ASSERT(kj::str(exception.getHandle(js)) == "Error: boom");
787 checked++;
788 });
789 
790 // Reading left the stream locked and disturbed
791 KJ_ASSERT(rs->isLocked());
792 KJ_ASSERT(rs->isDisturbed());
793 
794 // Run the microtasks to completion. This should resolve the promise and
795 // run it to completion. The test is buggy if it fails to do so.
796 js.runMicrotasks();
797 KJ_ASSERT(checked == 2);
798 
799 KJ_ASSERT(rs->getController().isClosedOrErrored());
800 
801 // Add we should still be locked and disturbed.
802 KJ_ASSERT(rs->isLocked());
803 KJ_ASSERT(rs->isDisturbed());
804 });
805}
806 
807KJ_TEST("ReadableStream read all bytes (byte readable, failed start 2)") {
808 preamble([](jsg::Lock& js) {
809 uint checked = 0;
810 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
811 // clang-format off
812 rs->getController().setup(js, UnderlyingSource{
813 .type = kj::str("bytes"),
814 .start = [&](jsg::Lock& js, UnderlyingSource::Controller controller) -> jsg::Promise<void> {
815 checked++;
816 JSG_FAIL_REQUIRE(Error, "boom");
817 }
818 // Setting a highWaterMark of 0 means the pull function above will not be called
819 // immediately on creation of the stream, but only when the first read in the
820 // readall call below happens.
821 }, StreamQueuingStrategy{.highWaterMark = 0});
822 // clang-format on
823 
824 // Starts a read loop of javascript promises.
825 auto promise = rs->getController().readAllBytes(js, 20).then(js,
826 [](jsg::Lock& js, jsg::BufferSource&& text) { KJ_UNREACHABLE; },
827 [&](jsg::Lock& js, jsg::Value&& exception) {
828 KJ_ASSERT(kj::str(exception.getHandle(js)) == "Error: boom");
829 checked++;
830 });
831 
832 // Reading left the stream locked and disturbed
833 KJ_ASSERT(rs->isLocked());
834 KJ_ASSERT(rs->isDisturbed());
835 
836 // Run the microtasks to completion. This should resolve the promise and
837 // run it to completion. The test is buggy if it fails to do so.
838 js.runMicrotasks();
839 KJ_ASSERT(checked == 2);
840 
841 KJ_ASSERT(rs->getController().isClosedOrErrored());
842 
843 // Add we should still be locked and disturbed.
844 KJ_ASSERT(rs->isLocked());
845 KJ_ASSERT(rs->isDisturbed());
846 });
847}
848 
849// ======================================================================================
850// DrainingReader tests
851 
852KJ_TEST("DrainingReader basic creation and locking (value stream)") {
853 preamble([](jsg::Lock& js) {
854 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
855 rs->getController().setup(js, UnderlyingSource{}, StreamQueuingStrategy{.highWaterMark = 0});
856 
857 // Stream should not be locked initially
858 KJ_ASSERT(!rs->isLocked());
859 
860 // Create DrainingReader
861 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
862 // Stream should now be locked
863 KJ_ASSERT(rs->isLocked());
864 KJ_ASSERT(reader->isAttached());
865 
866 // Release the lock
867 reader->releaseLock(js);
868 KJ_ASSERT(!rs->isLocked());
869 KJ_ASSERT(!reader->isAttached());
870 } else {
871 KJ_FAIL_ASSERT("Failed to create DrainingReader");
872 }
873 });
874}
875 
876KJ_TEST("DrainingReader cannot be created on locked stream") {
877 preamble([](jsg::Lock& js) {
878 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
879 rs->getController().setup(js, UnderlyingSource{}, StreamQueuingStrategy{.highWaterMark = 0});
880 
881 // Create first reader to lock the stream
882 KJ_IF_SOME(reader1, DrainingReader::create(js, *rs)) {
883 KJ_ASSERT(rs->isLocked());
884 
885 // Try to create another reader - should fail
886 auto maybeReader2 = DrainingReader::create(js, *rs);
887 KJ_ASSERT(maybeReader2 == kj::none);
888 
889 reader1->releaseLock(js);
890 } else {
891 KJ_FAIL_ASSERT("Failed to create first DrainingReader");
892 }
893 });
894}
895 
896KJ_TEST("DrainingReader read drains buffered data (value stream)") {
897 preamble([](jsg::Lock& js) {
898 uint pullCount = 0;
899 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
900 // clang-format off
901 rs->getController().setup(js, UnderlyingSource{
902 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
903 KJ_SWITCH_ONEOF(controller) {
904 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
905 pullCount++;
906 if (pullCount == 1) {
907 // First pull - enqueue multiple chunks
908 c->enqueue(js, toBytes(js, kj::str("Hello, ")));
909 c->enqueue(js, toBytes(js, kj::str("world!")));
910 } else {
911 // Second pull - close the stream
912 c->close(js);
913 }
914 return js.resolvedPromise();
915 }
916 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
917 }
918 KJ_UNREACHABLE;
919 }
920 }, StreamQueuingStrategy{.highWaterMark = 0});
921 // clang-format on
922 
923 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
924 bool readCompleted = false;
925 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
926 // Should have drained both chunks
927 KJ_ASSERT(result.chunks.size() == 2);
928 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "Hello, ");
929 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "world!");
930 KJ_ASSERT(!result.done); // Stream not closed yet
931 readCompleted = true;
932 });
933 
934 js.runMicrotasks();
935 KJ_ASSERT(readCompleted);
936 KJ_ASSERT(pullCount == 1); // Only one pull needed
937 
938 reader->releaseLock(js);
939 } else {
940 KJ_FAIL_ASSERT("Failed to create DrainingReader");
941 }
942 });
943}
944 
945KJ_TEST("DrainingReader read drains buffered data (byte stream)") {
946 preamble([](jsg::Lock& js) {
947 uint pullCount = 0;
948 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
949 // clang-format off
950 rs->getController().setup(js, UnderlyingSource{
951 .type = kj::str("bytes"),
952 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
953 KJ_SWITCH_ONEOF(controller) {
954 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
955 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
956 pullCount++;
957 if (pullCount == 1) {
958 c->enqueue(js, toBufferSource(js, kj::str("Hello, ")));
959 c->enqueue(js, toBufferSource(js, kj::str("world!")));
960 } else {
961 c->close(js);
962 }
963 return js.resolvedPromise();
964 }
965 }
966 KJ_UNREACHABLE;
967 }
968 }, StreamQueuingStrategy{.highWaterMark = 0});
969 // clang-format on
970 
971 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
972 bool readCompleted = false;
973 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
974 // Should have drained both chunks
975 KJ_ASSERT(result.chunks.size() == 2);
976 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "Hello, ");
977 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "world!");
978 KJ_ASSERT(!result.done);
979 readCompleted = true;
980 });
981 
982 js.runMicrotasks();
983 KJ_ASSERT(readCompleted);
984 KJ_ASSERT(pullCount == 1);
985 
986 reader->releaseLock(js);
987 } else {
988 KJ_FAIL_ASSERT("Failed to create DrainingReader");
989 }
990 });
991}
992 
993KJ_TEST("DrainingReader read on closed stream returns done") {
994 preamble([](jsg::Lock& js) {
995 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
996 // clang-format off
997 rs->getController().setup(js, UnderlyingSource{
998 .start = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
999 KJ_SWITCH_ONEOF(controller) {
1000 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1001 c->close(js);
1002 return js.resolvedPromise();
1003 }
1004 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1005 }
1006 KJ_UNREACHABLE;
1007 }
1008 }, StreamQueuingStrategy{.highWaterMark = 0});
1009 // clang-format on
1010 
1011 js.runMicrotasks();
1012 
1013 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1014 bool readCompleted = false;
1015 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1016 KJ_ASSERT(result.chunks.size() == 0);
1017 KJ_ASSERT(result.done);
1018 readCompleted = true;
1019 });
1020 
1021 js.runMicrotasks();
1022 KJ_ASSERT(readCompleted);
1023 
1024 reader->releaseLock(js);
1025 } else {
1026 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1027 }
1028 });
1029}
1030 
1031KJ_TEST("DrainingReader read after releaseLock rejects") {
1032 preamble([](jsg::Lock& js) {
1033 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1034 rs->getController().setup(js, UnderlyingSource{}, StreamQueuingStrategy{.highWaterMark = 0});
1035 
1036 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1037 reader->releaseLock(js);
1038 
1039 bool readRejected = false;
1040 auto promise =
1041 reader->read(js).catch_(js, [&](jsg::Lock& js, jsg::Value reason) -> DrainingReadResult {
1042 readRejected = true;
1043 return DrainingReadResult{
1044 .chunks = kj::Array<kj::Array<kj::byte>>(),
1045 .done = true,
1046 };
1047 });
1048 
1049 js.runMicrotasks();
1050 KJ_ASSERT(readRejected);
1051 } else {
1052 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1053 }
1054 });
1055}
1056 
1057KJ_TEST("DrainingReader sync data then async pull waits") {
1058 // Test case: pull enqueues some data synchronously, then returns a pending promise.
1059 // The first draining read should get the sync data immediately.
1060 // A second draining read should wait for the async pull to complete.
1061 preamble([](jsg::Lock& js) {
1062 uint pullCount = 0;
1063 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1064 kj::Maybe<jsg::Ref<ReadableStreamDefaultController>> savedController;
1065 
1066 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1067 // clang-format off
1068 rs->getController().setup(js, UnderlyingSource{
1069 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1070 KJ_SWITCH_ONEOF(controller) {
1071 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1072 pullCount++;
1073 if (pullCount == 1) {
1074 // First pull: enqueue data synchronously, but return async promise
1075 c->enqueue(js, toBytes(js, kj::str("sync-chunk")));
1076 // Return a promise that resolves later
1077 auto prp = js.newPromiseAndResolver<void>();
1078 asyncResolver = kj::mv(prp.resolver);
1079 savedController = c.addRef();
1080 return kj::mv(prp.promise);
1081 } else if (pullCount == 2) {
1082 // Second pull after async resolution: enqueue more data
1083 c->enqueue(js, toBytes(js, kj::str("async-chunk")));
1084 return js.resolvedPromise();
1085 }
1086 return js.resolvedPromise();
1087 }
1088 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1089 }
1090 KJ_UNREACHABLE;
1091 }
1092 }, StreamQueuingStrategy{.highWaterMark = 0});
1093 // clang-format on
1094 
1095 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1096 // First read - should get sync data immediately
1097 bool firstReadCompleted = false;
1098 auto promise1 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1099 KJ_ASSERT(result.chunks.size() == 1);
1100 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "sync-chunk");
1101 KJ_ASSERT(!result.done);
1102 firstReadCompleted = true;
1103 });
1104 
1105 js.runMicrotasks();
1106 KJ_ASSERT(firstReadCompleted);
1107 KJ_ASSERT(pullCount == 1); // Only first pull happened
1108 
1109 // Second read - should wait for async data
1110 bool secondReadCompleted = false;
1111 auto promise2 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1112 // Should get the async chunk
1113 KJ_ASSERT(result.chunks.size() >= 1);
1114 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "async-chunk");
1115 KJ_ASSERT(!result.done);
1116 secondReadCompleted = true;
1117 });
1118 
1119 js.runMicrotasks();
1120 // Second read should NOT complete yet - still waiting for async pull
1121 KJ_ASSERT(!secondReadCompleted);
1122 
1123 // Now resolve the async pull
1124 KJ_ASSERT_NONNULL(asyncResolver).resolve(js);
1125 js.runMicrotasks();
1126 
1127 // Now second read should complete
1128 KJ_ASSERT(secondReadCompleted);
1129 KJ_ASSERT(pullCount == 2);
1130 
1131 reader->releaseLock(js);
1132 } else {
1133 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1134 }
1135 });
1136}
1137 
1138KJ_TEST("DrainingReader with fully async pull") {
1139 // Test case: pull returns a promise without enqueueing anything synchronously.
1140 // The draining read should wait for the pull to complete and then get the data.
1141 preamble([](jsg::Lock& js) {
1142 uint pullCount = 0;
1143 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1144 kj::Maybe<jsg::Ref<ReadableStreamDefaultController>> savedController;
1145 
1146 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1147 // clang-format off
1148 rs->getController().setup(js, UnderlyingSource{
1149 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1150 KJ_SWITCH_ONEOF(controller) {
1151 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1152 pullCount++;
1153 // Return a promise and save the controller - will enqueue data when resolved
1154 auto prp = js.newPromiseAndResolver<void>();
1155 asyncResolver = kj::mv(prp.resolver);
1156 savedController = c.addRef();
1157 return kj::mv(prp.promise);
1158 }
1159 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1160 }
1161 KJ_UNREACHABLE;
1162 }
1163 }, StreamQueuingStrategy{.highWaterMark = 0});
1164 // clang-format on
1165 
1166 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1167 bool readCompleted = false;
1168 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1169 KJ_ASSERT(result.chunks.size() == 1);
1170 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "async-data");
1171 KJ_ASSERT(!result.done);
1172 readCompleted = true;
1173 });
1174 
1175 js.runMicrotasks();
1176 // Read should NOT complete yet - waiting for async pull
1177 KJ_ASSERT(!readCompleted);
1178 KJ_ASSERT(pullCount == 1);
1179 
1180 // Enqueue data and resolve the pull
1181 KJ_ASSERT_NONNULL(savedController)->enqueue(js, toBytes(js, kj::str("async-data")));
1182 KJ_ASSERT_NONNULL(asyncResolver).resolve(js);
1183 js.runMicrotasks();
1184 
1185 // Now read should complete
1186 KJ_ASSERT(readCompleted);
1187 
1188 reader->releaseLock(js);
1189 } else {
1190 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1191 }
1192 });
1193}
1194 
1195KJ_TEST("DrainingReader byte stream with async pull") {
1196 // Test async behavior with byte streams
1197 preamble([](jsg::Lock& js) {
1198 uint pullCount = 0;
1199 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1200 kj::Maybe<jsg::Ref<ReadableByteStreamController>> savedController;
1201 
1202 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1203 // clang-format off
1204 rs->getController().setup(js, UnderlyingSource{
1205 .type = kj::str("bytes"),
1206 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1207 KJ_SWITCH_ONEOF(controller) {
1208 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
1209 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
1210 pullCount++;
1211 if (pullCount == 1) {
1212 // Enqueue sync data but return async
1213 c->enqueue(js, toBufferSource(js, kj::str("sync-bytes")));
1214 auto prp = js.newPromiseAndResolver<void>();
1215 asyncResolver = kj::mv(prp.resolver);
1216 savedController = c.addRef();
1217 return kj::mv(prp.promise);
1218 }
1219 return js.resolvedPromise();
1220 }
1221 }
1222 KJ_UNREACHABLE;
1223 }
1224 }, StreamQueuingStrategy{.highWaterMark = 0});
1225 // clang-format on
1226 
1227 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1228 // First read gets sync data
1229 bool firstReadCompleted = false;
1230 auto promise1 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1231 KJ_ASSERT(result.chunks.size() == 1);
1232 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "sync-bytes");
1233 KJ_ASSERT(!result.done);
1234 firstReadCompleted = true;
1235 });
1236 
1237 js.runMicrotasks();
1238 KJ_ASSERT(firstReadCompleted);
1239 
1240 // Resolve async pull to allow future pulls
1241 KJ_ASSERT_NONNULL(asyncResolver).resolve(js);
1242 js.runMicrotasks();
1243 
1244 reader->releaseLock(js);
1245 } else {
1246 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1247 }
1248 });
1249}
1250 
1251KJ_TEST("DrainingReader multiple sync chunks then close") {
1252 // Test: Multiple sync chunks followed by close in the same pull
1253 preamble([](jsg::Lock& js) {
1254 uint pullCount = 0;
1255 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1256 // clang-format off
1257 rs->getController().setup(js, UnderlyingSource{
1258 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1259 KJ_SWITCH_ONEOF(controller) {
1260 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1261 pullCount++;
1262 // Enqueue multiple chunks then close
1263 c->enqueue(js, toBytes(js, kj::str("chunk1")));
1264 c->enqueue(js, toBytes(js, kj::str("chunk2")));
1265 c->enqueue(js, toBytes(js, kj::str("chunk3")));
1266 c->close(js);
1267 return js.resolvedPromise();
1268 }
1269 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1270 }
1271 KJ_UNREACHABLE;
1272 }
1273 }, StreamQueuingStrategy{.highWaterMark = 0});
1274 // clang-format on
1275 
1276 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1277 bool readCompleted = false;
1278 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1279 // Should get all 3 chunks and done=true
1280 KJ_ASSERT(result.chunks.size() == 3);
1281 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "chunk1");
1282 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "chunk2");
1283 KJ_ASSERT(kj::str(result.chunks[2].asChars()) == "chunk3");
1284 KJ_ASSERT(result.done);
1285 readCompleted = true;
1286 });
1287 
1288 js.runMicrotasks();
1289 KJ_ASSERT(readCompleted);
1290 KJ_ASSERT(pullCount == 1);
1291 
1292 reader->releaseLock(js);
1293 } else {
1294 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1295 }
1296 });
1297}
1298 
1299KJ_TEST("DrainingReader read from teed branches") {
1300 // Test: DrainingReader works correctly on both branches of a teed stream
1301 preamble([](jsg::Lock& js) {
1302 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1303 // clang-format off
1304 rs->getController().setup(js, UnderlyingSource{
1305 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
1306 KJ_SWITCH_ONEOF(controller) {
1307 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1308 c->enqueue(js, toBytes(js, kj::str("chunk1")));
1309 c->enqueue(js, toBytes(js, kj::str("chunk2")));
1310 c->close(js);
1311 return js.resolvedPromise();
1312 }
1313 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1314 }
1315 KJ_UNREACHABLE;
1316 }
1317 }, StreamQueuingStrategy{.highWaterMark = 0});
1318 // clang-format on
1319 
1320 // Tee the stream into two branches
1321 auto branches = rs->tee(js);
1322 KJ_ASSERT(branches.size() == 2);
1323 auto& branch1 = branches[0];
1324 auto& branch2 = branches[1];
1325 
1326 // Create DrainingReader on branch1
1327 KJ_IF_SOME(reader1, DrainingReader::create(js, *branch1)) {
1328 bool read1Completed = false;
1329 auto promise1 = reader1->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1330 KJ_ASSERT(result.chunks.size() == 2);
1331 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "chunk1");
1332 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "chunk2");
1333 KJ_ASSERT(result.done);
1334 read1Completed = true;
1335 });
1336 
1337 js.runMicrotasks();
1338 KJ_ASSERT(read1Completed, "Branch1 read should complete");
1339 
1340 reader1->releaseLock(js);
1341 } else {
1342 KJ_FAIL_ASSERT("Failed to create DrainingReader for branch1");
1343 }
1344 
1345 // Create DrainingReader on branch2 - should get the same data
1346 KJ_IF_SOME(reader2, DrainingReader::create(js, *branch2)) {
1347 bool read2Completed = false;
1348 auto promise2 = reader2->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1349 KJ_ASSERT(result.chunks.size() == 2);
1350 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "chunk1");
1351 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "chunk2");
1352 KJ_ASSERT(result.done);
1353 read2Completed = true;
1354 });
1355 
1356 js.runMicrotasks();
1357 KJ_ASSERT(read2Completed, "Branch2 read should complete");
1358 
1359 reader2->releaseLock(js);
1360 } else {
1361 KJ_FAIL_ASSERT("Failed to create DrainingReader for branch2");
1362 }
1363 });
1364}
1365 
1366KJ_TEST("DrainingReader read from byte stream with BYOB support") {
1367 // Test: DrainingReader works correctly with byte streams (which support BYOB reads)
1368 // even though DrainingReader itself uses a default reader. This tests that closing
1369 // the controller synchronously during draining works correctly - the controller's
1370 // close triggers doClose() which must be deferred while onConsumerWantsData is active.
1371 preamble([](jsg::Lock& js) {
1372 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1373 // clang-format off
1374 rs->getController().setup(js, UnderlyingSource{
1375 .type = kj::str("bytes"),
1376 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1377 KJ_SWITCH_ONEOF(controller) {
1378 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
1379 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
1380 // Enqueue multiple byte chunks - verifies DrainingReader handles
1381 // byte stream chunks correctly and preserves order
1382 c->enqueue(js, toBufferSource(js, kj::str("byob-chunk1")));
1383 c->enqueue(js, toBufferSource(js, kj::str("byob-chunk2")));
1384 c->enqueue(js, toBufferSource(js, kj::str("byob-chunk3")));
1385 // Close synchronously - this tests that the fix for use-after-free works.
1386 // Without the fix, this would cause ByteReadable to be destroyed while
1387 // onConsumerWantsData is still on the stack.
1388 c->close(js);
1389 return js.resolvedPromise();
1390 }
1391 }
1392 KJ_UNREACHABLE;
1393 }
1394 }, StreamQueuingStrategy{.highWaterMark = 0});
1395 // clang-format on
1396 
1397 // Use DrainingReader (which uses default reader) to drain the BYOB-capable byte stream
1398 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1399 bool readCompleted = false;
1400 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1401 // Should successfully drain all byte chunks in order
1402 KJ_ASSERT(result.chunks.size() == 3, "Should get 3 chunks");
1403 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "byob-chunk1");
1404 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "byob-chunk2");
1405 KJ_ASSERT(kj::str(result.chunks[2].asChars()) == "byob-chunk3");
1406 KJ_ASSERT(result.done); // Stream closed, so done should be true
1407 readCompleted = true;
1408 });
1409 
1410 js.runMicrotasks();
1411 KJ_ASSERT(readCompleted);
1412 
1413 reader->releaseLock(js);
1414 } else {
1415 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1416 }
1417 });
1418}
1419 
1420KJ_TEST("DrainingReader error during pull in value stream") {
1421 // Test: Calling controller.error() synchronously inside pull during a draining
1422 // read must not cause a use-after-free. The pull callback transitions the
1423 // ConsumerImpl from Ready→Errored which destroys the Ready struct (and its
1424 // RingBuffer). The drainingRead loop must detect this and stop.
1425 preamble([](jsg::Lock& js) {
1426 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1427 // clang-format off
1428 rs->getController().setup(js, UnderlyingSource{
1429 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
1430 KJ_SWITCH_ONEOF(controller) {
1431 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1432 c->enqueue(js, toBytes(js, kj::str("before-error")));
1433 c->error(js, js.error("deliberate error"));
1434 return js.resolvedPromise();
1435 }
1436 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1437 }
1438 KJ_UNREACHABLE;
1439 }
1440 }, StreamQueuingStrategy{.highWaterMark = 0});
1441 // clang-format on
1442 
1443 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1444 bool readCompleted = false;
1445 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1446 KJ_FAIL_ASSERT("Should have rejected, not resolved");
1447 }, [&](jsg::Lock& js, jsg::Value&& err) {
1448 // The draining read should reject with the error from controller.error().
1449 readCompleted = true;
1450 });
1451 
1452 js.runMicrotasks();
1453 KJ_ASSERT(readCompleted, "Read should have completed with rejection");
1454 
1455 reader->releaseLock(js);
1456 } else {
1457 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1458 }
1459 });
1460}
1461 
1462KJ_TEST("DrainingReader error during pull in byte stream") {
1463 // Test: Same as above but for byte streams (ByteQueue path).
1464 preamble([](jsg::Lock& js) {
1465 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1466 // clang-format off
1467 rs->getController().setup(js, UnderlyingSource{
1468 .type = kj::str("bytes"),
1469 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
1470 KJ_SWITCH_ONEOF(controller) {
1471 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
1472 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
1473 c->enqueue(js, toBufferSource(js, kj::str("before-error")));
1474 c->error(js, js.error("deliberate error"));
1475 return js.resolvedPromise();
1476 }
1477 }
1478 KJ_UNREACHABLE;
1479 }
1480 }, StreamQueuingStrategy{.highWaterMark = 0});
1481 // clang-format on
1482 
1483 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1484 bool readCompleted = false;
1485 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1486 KJ_FAIL_ASSERT("Should have rejected, not resolved");
1487 }, [&](jsg::Lock& js, jsg::Value&& err) { readCompleted = true; });
1488 
1489 js.runMicrotasks();
1490 KJ_ASSERT(readCompleted, "Read should have completed with rejection");
1491 
1492 reader->releaseLock(js);
1493 } else {
1494 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1495 }
1496 });
1497}
1498 
1499KJ_TEST("DrainingReader read from stream with transform-like pattern") {
1500 // Test: DrainingReader works correctly with a stream that simulates the
1501 // TransformStream pattern where data is written to writable and read from readable
1502 preamble([](jsg::Lock& js) {
1503 // Create a stream where the controller is stored and chunks are enqueued asynchronously
1504 // (simulating how TransformStream's readable side receives transformed data)
1505 kj::Maybe<jsg::Ref<ReadableStreamDefaultController>> savedController;
1506 bool startResolved = false;
1507 
1508 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1509 // clang-format off
1510 rs->getController().setup(js, UnderlyingSource{
1511 .start = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1512 KJ_SWITCH_ONEOF(controller) {
1513 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1514 savedController = c.addRef();
1515 startResolved = true;
1516 return js.resolvedPromise();
1517 }
1518 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1519 }
1520 KJ_UNREACHABLE;
1521 },
1522 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
1523 // No-op pull - data comes from external enqueue calls (like transform writes)
1524 return js.resolvedPromise();
1525 }
1526 }, StreamQueuingStrategy{.highWaterMark = 0});
1527 // clang-format on
1528 
1529 js.runMicrotasks();
1530 KJ_ASSERT(startResolved, "Stream should have started");
1531 
1532 auto& controller = KJ_ASSERT_NONNULL(savedController, "Controller should be saved");
1533 
1534 // Simulate TransformStream write->transform->enqueue pattern
1535 // Enqueue transformed chunks (like what TransformStream's transform callback would do)
1536 controller->enqueue(js, toBytes(js, kj::str("transformed-a")));
1537 controller->enqueue(js, toBytes(js, kj::str("transformed-b")));
1538 
1539 // Create DrainingReader to drain all buffered transformed data
1540 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1541 bool readCompleted = false;
1542 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1543 // Should drain all enqueued chunks
1544 KJ_ASSERT(result.chunks.size() == 2);
1545 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "transformed-a");
1546 KJ_ASSERT(kj::str(result.chunks[1].asChars()) == "transformed-b");
1547 KJ_ASSERT(!result.done); // Stream not closed yet
1548 readCompleted = true;
1549 });
1550 
1551 js.runMicrotasks();
1552 KJ_ASSERT(readCompleted);
1553 
1554 // Simulate more data being written/transformed
1555 controller->enqueue(js, toBytes(js, kj::str("transformed-c")));
1556 controller->close(js);
1557 
1558 bool finalReadCompleted = false;
1559 auto finalPromise =
1560 reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1561 KJ_ASSERT(result.chunks.size() == 1);
1562 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "transformed-c");
1563 KJ_ASSERT(result.done); // Stream now closed
1564 finalReadCompleted = true;
1565 });
1566 
1567 js.runMicrotasks();
1568 KJ_ASSERT(finalReadCompleted);
1569 
1570 reader->releaseLock(js);
1571 } else {
1572 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1573 }
1574 });
1575}
1576 
1577KJ_TEST("DrainingReader cancel while read is pending (value stream)") {
1578 // Test: Calling cancel() on the reader while a read() is pending should
1579 // cause the pending read to reject and the stream to be canceled.
1580 preamble([](jsg::Lock& js) {
1581 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1582 bool cancelCalled = false;
1583 
1584 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1585 // clang-format off
1586 rs->getController().setup(js, UnderlyingSource{
1587 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1588 // Return a pending promise to keep the read waiting
1589 auto prp = js.newPromiseAndResolver<void>();
1590 asyncResolver = kj::mv(prp.resolver);
1591 return kj::mv(prp.promise);
1592 },
1593 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
1594 cancelCalled = true;
1595 KJ_ASSERT(kj::str(reason) == "canceled by reader");
1596 return js.resolvedPromise();
1597 }
1598 }, StreamQueuingStrategy{.highWaterMark = 0});
1599 // clang-format on
1600 
1601 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1602 // Start a read that will be pending (waiting for async pull)
1603 bool readRejected = false;
1604 bool readResolved = false;
1605 auto readPromise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1606 readResolved = true;
1607 }, [&](jsg::Lock& js, jsg::Value&& reason) { readRejected = true; });
1608 
1609 js.runMicrotasks();
1610 // Read should still be pending - waiting for async pull
1611 KJ_ASSERT(!readResolved);
1612 KJ_ASSERT(!readRejected);
1613 
1614 // Now cancel while read is pending
1615 bool cancelResolved = false;
1616 auto cancelPromise =
1617 reader->cancel(js, js.str("canceled by reader"_kjc)).then(js, [&](jsg::Lock& js) {
1618 cancelResolved = true;
1619 });
1620 
1621 js.runMicrotasks();
1622 
1623 // Cancel should have completed
1624 KJ_ASSERT(cancelResolved, "cancel() should resolve");
1625 KJ_ASSERT(cancelCalled, "underlying source cancel should be called");
1626 
1627 // The pending read should have been rejected or resolved with done
1628 KJ_ASSERT(readResolved || readRejected, "pending read should complete after cancel");
1629 
1630 // Stream should be in closed/errored state
1631 KJ_ASSERT(rs->getController().isClosedOrErrored());
1632 } else {
1633 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1634 }
1635 });
1636}
1637 
1638KJ_TEST("DrainingReader cancel while read is pending (byte stream)") {
1639 // Test: Same as above but with byte stream
1640 preamble([](jsg::Lock& js) {
1641 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1642 bool cancelCalled = false;
1643 
1644 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1645 // clang-format off
1646 rs->getController().setup(js, UnderlyingSource{
1647 .type = kj::str("bytes"),
1648 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1649 // Return a pending promise to keep the read waiting
1650 auto prp = js.newPromiseAndResolver<void>();
1651 asyncResolver = kj::mv(prp.resolver);
1652 return kj::mv(prp.promise);
1653 },
1654 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
1655 cancelCalled = true;
1656 KJ_ASSERT(kj::str(reason) == "canceled by reader");
1657 return js.resolvedPromise();
1658 }
1659 }, StreamQueuingStrategy{.highWaterMark = 0});
1660 // clang-format on
1661 
1662 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1663 // Start a read that will be pending
1664 bool readRejected = false;
1665 bool readResolved = false;
1666 auto readPromise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1667 readResolved = true;
1668 }, [&](jsg::Lock& js, jsg::Value&& reason) { readRejected = true; });
1669 
1670 js.runMicrotasks();
1671 KJ_ASSERT(!readResolved);
1672 KJ_ASSERT(!readRejected);
1673 
1674 // Cancel while read is pending
1675 bool cancelResolved = false;
1676 auto cancelPromise =
1677 reader->cancel(js, js.str("canceled by reader"_kjc)).then(js, [&](jsg::Lock& js) {
1678 cancelResolved = true;
1679 });
1680 
1681 js.runMicrotasks();
1682 
1683 KJ_ASSERT(cancelResolved, "cancel() should resolve");
1684 KJ_ASSERT(cancelCalled, "underlying source cancel should be called");
1685 KJ_ASSERT(readResolved || readRejected, "pending read should complete after cancel");
1686 KJ_ASSERT(rs->getController().isClosedOrErrored());
1687 } else {
1688 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1689 }
1690 });
1691}
1692 
1693KJ_TEST("DrainingReader cancel while read is pending with buffered data") {
1694 // Test: Cancel while read is pending, but there's already some buffered data.
1695 // The buffered data should be discarded and the stream canceled.
1696 preamble([](jsg::Lock& js) {
1697 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1698 kj::Maybe<jsg::Ref<ReadableStreamDefaultController>> savedController;
1699 bool cancelCalled = false;
1700 
1701 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1702 // clang-format off
1703 rs->getController().setup(js, UnderlyingSource{
1704 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1705 KJ_SWITCH_ONEOF(controller) {
1706 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
1707 // Enqueue some data synchronously
1708 c->enqueue(js, toBytes(js, kj::str("buffered-data")));
1709 savedController = c.addRef();
1710 // But return a pending promise (more data coming)
1711 auto prp = js.newPromiseAndResolver<void>();
1712 asyncResolver = kj::mv(prp.resolver);
1713 return kj::mv(prp.promise);
1714 }
1715 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
1716 }
1717 KJ_UNREACHABLE;
1718 },
1719 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
1720 cancelCalled = true;
1721 return js.resolvedPromise();
1722 }
1723 }, StreamQueuingStrategy{.highWaterMark = 0});
1724 // clang-format on
1725 
1726 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1727 // First read gets the buffered data
1728 bool firstReadCompleted = false;
1729 auto readPromise1 =
1730 reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1731 KJ_ASSERT(result.chunks.size() == 1);
1732 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "buffered-data");
1733 KJ_ASSERT(!result.done);
1734 firstReadCompleted = true;
1735 });
1736 
1737 js.runMicrotasks();
1738 KJ_ASSERT(firstReadCompleted);
1739 
1740 // Second read will be pending (waiting for async pull resolution)
1741 bool secondReadRejected = false;
1742 bool secondReadResolved = false;
1743 auto readPromise2 = reader->read(js).then(js,
1744 [&](jsg::Lock& js, DrainingReadResult&& result) { secondReadResolved = true; },
1745 [&](jsg::Lock& js, jsg::Value&& reason) { secondReadRejected = true; });
1746 
1747 js.runMicrotasks();
1748 KJ_ASSERT(!secondReadResolved);
1749 KJ_ASSERT(!secondReadRejected);
1750 
1751 // Cancel while second read is pending
1752 bool cancelResolved = false;
1753 auto cancelPromise =
1754 reader->cancel(js, js.str("cancel reason"_kjc)).then(js, [&](jsg::Lock& js) {
1755 cancelResolved = true;
1756 });
1757 
1758 js.runMicrotasks();
1759 
1760 KJ_ASSERT(cancelResolved);
1761 KJ_ASSERT(cancelCalled);
1762 KJ_ASSERT(secondReadResolved || secondReadRejected);
1763 KJ_ASSERT(rs->getController().isClosedOrErrored());
1764 } else {
1765 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1766 }
1767 });
1768}
1769 
1770KJ_TEST("DrainingReader cancel while read pending - UAF safety (value stream)") {
1771 // This test specifically exercises the potential UAF scenario where:
1772 // 1. A draining read creates a promise with lambdas capturing `this` (the Consumer)
1773 // 2. Cancel is called, which rejects the pending read (scheduling the error lambda)
1774 // 3. doClose() destroys the Consumer
1775 // 4. The error lambda runs and must NOT access the destroyed Consumer
1776 //
1777 // The lambdas in ValueQueue::Consumer::drainingRead capture `this` to clear
1778 // hasPendingDrainingRead. If the Consumer is destroyed before the lambda runs,
1779 // this would be a use-after-free.
1780 preamble([](jsg::Lock& js) {
1781 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1782 bool cancelCalled = false;
1783 bool pullCalled = false;
1784 
1785 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1786 // clang-format off
1787 rs->getController().setup(js, UnderlyingSource{
1788 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1789 pullCalled = true;
1790 // Return a pending promise - this keeps the read waiting
1791 auto prp = js.newPromiseAndResolver<void>();
1792 asyncResolver = kj::mv(prp.resolver);
1793 return kj::mv(prp.promise);
1794 },
1795 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
1796 cancelCalled = true;
1797 return js.resolvedPromise();
1798 }
1799 }, StreamQueuingStrategy{.highWaterMark = 0});
1800 // clang-format on
1801 
1802 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1803 // Start a draining read - this will:
1804 // 1. Call pull (which returns pending promise)
1805 // 2. Queue a ReadRequest
1806 // 3. Return a promise with lambdas capturing `this` (the Consumer)
1807 bool readRejected = false;
1808 bool readResolved = false;
1809 auto readPromise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1810 readResolved = true;
1811 }, [&](jsg::Lock& js, jsg::Value&& reason) {
1812 // This error handler runs after cancel rejects the pending read.
1813 // The lambda in drainingRead also runs to clear hasPendingDrainingRead.
1814 // If that lambda accesses a destroyed Consumer, we have UAF.
1815 readRejected = true;
1816 });
1817 
1818 js.runMicrotasks();
1819 KJ_ASSERT(pullCalled, "pull should have been called");
1820 KJ_ASSERT(!readResolved);
1821 KJ_ASSERT(!readRejected);
1822 
1823 // Now cancel. This will:
1824 // 1. Call cancelPendingReads() which rejects the ReadRequest
1825 // 2. The rejection schedules the error lambda as a microtask
1826 // 3. doClose() runs (via KJ_DEFER) and may destroy the Consumer
1827 // 4. Microtasks run - the error lambda in drainingRead accesses `this`
1828 //
1829 // If `this` is destroyed before the lambda runs, we have UAF.
1830 bool cancelResolved = false;
1831 auto cancelPromise =
1832 reader->cancel(js, js.str("cancel for UAF test"_kjc)).then(js, [&](jsg::Lock& js) {
1833 cancelResolved = true;
1834 });
1835 
1836 // Run microtasks - this is where UAF would occur if the bug exists
1837 js.runMicrotasks();
1838 
1839 KJ_ASSERT(cancelResolved, "cancel should resolve");
1840 KJ_ASSERT(cancelCalled, "underlying source cancel should be called");
1841 KJ_ASSERT(readResolved || readRejected, "read should complete after cancel");
1842 KJ_ASSERT(rs->getController().isClosedOrErrored(), "stream should be closed/errored");
1843 } else {
1844 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1845 }
1846 });
1847}
1848 
1849KJ_TEST("DrainingReader cancel while read pending - UAF safety (byte stream)") {
1850 // Same test as above but for byte streams (ByteQueue::Consumer)
1851 preamble([](jsg::Lock& js) {
1852 kj::Maybe<jsg::Promise<void>::Resolver> asyncResolver;
1853 bool cancelCalled = false;
1854 bool pullCalled = false;
1855 
1856 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1857 // clang-format off
1858 rs->getController().setup(js, UnderlyingSource{
1859 .type = kj::str("bytes"),
1860 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1861 pullCalled = true;
1862 auto prp = js.newPromiseAndResolver<void>();
1863 asyncResolver = kj::mv(prp.resolver);
1864 return kj::mv(prp.promise);
1865 },
1866 .cancel = [&](jsg::Lock& js, auto reason) -> jsg::Promise<void> {
1867 cancelCalled = true;
1868 return js.resolvedPromise();
1869 }
1870 }, StreamQueuingStrategy{.highWaterMark = 0});
1871 // clang-format on
1872 
1873 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
1874 bool readRejected = false;
1875 bool readResolved = false;
1876 auto readPromise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1877 readResolved = true;
1878 }, [&](jsg::Lock& js, jsg::Value&& reason) { readRejected = true; });
1879 
1880 js.runMicrotasks();
1881 KJ_ASSERT(pullCalled);
1882 KJ_ASSERT(!readResolved);
1883 KJ_ASSERT(!readRejected);
1884 
1885 bool cancelResolved = false;
1886 auto cancelPromise =
1887 reader->cancel(js, js.str("cancel for UAF test"_kjc)).then(js, [&](jsg::Lock& js) {
1888 cancelResolved = true;
1889 });
1890 
1891 js.runMicrotasks();
1892 
1893 KJ_ASSERT(cancelResolved);
1894 KJ_ASSERT(cancelCalled);
1895 KJ_ASSERT(readResolved || readRejected);
1896 KJ_ASSERT(rs->getController().isClosedOrErrored());
1897 } else {
1898 KJ_FAIL_ASSERT("Failed to create DrainingReader");
1899 }
1900 });
1901}
1902 
1903// ======================================================================================
1904// Execution Termination Tests
1905 
1906// Test that execution termination during a read doesn't cause an assertion failure.
1907// This tests the fix for the production error:
1908// "expected !js.v8Isolate->IsExecutionTerminating()"
1909//
1910// Note: This test verifies that if IsExecutionTerminating() is true after readCallback()
1911// returns, we handle it gracefully instead of asserting. However, actually triggering
1912// the termination check in the exact right spot in a unit test is difficult because
1913// V8 only checks the termination flag during JS execution, not during C++ callbacks.
1914// The production scenario occurs when the termination happens during actual JS execution
1915// inside readCallback() (e.g., user's pull function runs too long).
1916KJ_TEST("ReadableStream handles execution termination during read") {
1917 preamble([](jsg::Lock& js) {
1918 // We need to handle termination at this level because the context scope
1919 // will propagate the termination exception.
1920 v8::TryCatch outerTryCatch(js.v8Isolate);
1921 bool testCompleted = false;
1922 
1923 try {
1924 v8::TryCatch tryCatch(js.v8Isolate);
1925 
1926 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
1927 
1928 // Set up a stream where the pull callback terminates execution
1929 rs->getController().setup(js,
1930 UnderlyingSource{
1931 .pull =
1932 [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
1933 // Terminate execution - this simulates CPU time limit exceeded
1934 js.terminateNextExecution();
1935 
1936 // Force V8 to check the termination flag by doing JS work.
1937 // This is similar to what terminateExecutionNow() does.
1938 jsg::check(v8::JSON::Stringify(js.v8Context(), js.str("test"_kj)));
1939 
1940 // We shouldn't get here - termination should have been triggered
1941 return js.resolvedPromise();
1942 },
1943 },
1944 StreamQueuingStrategy{.highWaterMark = 0});
1945 
1946 // Start a read - this will call pull which terminates execution
1947 auto promise = rs->getController().readAllText(js, 100);
1948 
1949 // Run microtasks - this will propagate the termination
1950 js.runMicrotasks();
1951 
1952 // If we get here without the assertion failing, the fix works
1953 testCompleted = true;
1954 } catch (jsg::JsExceptionThrown&) {
1955 // Expected - execution was terminated
1956 // The important thing is we didn't hit the assertion failure
1957 testCompleted = true;
1958 }
1959 
1960 // Cancel termination so cleanup can proceed
1961 if (js.v8Isolate->IsExecutionTerminating()) {
1962 js.v8Isolate->CancelTerminateExecution();
1963 }
1964 
1965 // The test passes if we got here without an assertion failure
1966 KJ_ASSERT(testCompleted, "Test did not complete as expected");
1967 });
1968}
1969 
1970// ======================================================================================
1971// WritableStream re-entrancy tests
1972 
1973KJ_TEST("WritableStream close during abort algorithm returns rejected promise") {
1974 // Regression test for a bug where calling writer.close() re-entrantly from
1975 // the underlying sink's abort algorithm would hit a KJ_ASSERT in WritableImpl::close.
1976 //
1977 // The issue: finishErroring transitions WritableImpl to Errored state, then runs
1978 // the abort algorithm (user JS code). During that window, WritableStreamJsController
1979 // is still in Controller state (doError hasn't propagated yet). If the abort
1980 // algorithm calls writer.close(), the JsController delegates to WritableImpl::close
1981 // which previously asserted isWritable()||isErroring() -- but state was Errored.
1982 //
1983 // The fix adds defensive checks in WritableImpl::close for Closed/Errored states
1984 // that return rejected promises instead of asserting.
1985 preamble([](jsg::Lock& js) {
1986 bool abortCalled = false;
1987 bool closeRejected = false;
1988 bool abortResolved = false;
1989 
1990 auto ws = js.alloc<WritableStream>(newWritableStreamJsController());
1991 
1992 // We capture a raw pointer to the writer so the abort callback can call close().
1993 WritableStreamDefaultWriter* writerPtr = nullptr;
1994 
1995 // clang-format off
1996 ws->getController().setup(js, UnderlyingSink{
1997 .abort = [&](jsg::Lock& js, v8::Local<v8::Value> reason) -> jsg::Promise<void> {
1998 abortCalled = true;
1999 // Re-entrantly call close() on the writer during the abort algorithm.
2000 // At this point, WritableImpl has already transitioned to Errored state
2001 // but WritableStreamJsController hasn't been updated yet.
2002 // This should return a rejected promise instead of crashing.
2003 writerPtr->close(js).then(js,
2004 [](jsg::Lock& js) {},
2005 [&](jsg::Lock& js, jsg::Value reason) {
2006 closeRejected = true;
2007 });
2008 return js.resolvedPromise();
2009 }
2010 }, StreamQueuingStrategy{});
2011 // clang-format on
2012 
2013 auto writer = js.alloc<WritableStreamDefaultWriter>();
2014 writer->lockToStream(js, *ws);
2015 writerPtr = &*writer;
2016 
2017 // Abort the writer. This triggers:
2018 // 1. WritableImpl::abort() sets pending abort, KJ_DEFER fires startErroring
2019 // 2. startErroring -> finishErroring (no in-flight ops, stream is started)
2020 // 3. finishErroring transitions WritableImpl to Errored
2021 // 4. finishErroring calls the abort algorithm (our callback above)
2022 // 5. The abort callback calls writer.close() while WritableImpl is Errored
2023 // but WritableStreamJsController still shows Controller state
2024 writer->abort(js, kj::none).then(js, [&](jsg::Lock& js) {
2025 abortResolved = true;
2026 }, [&](jsg::Lock& js, jsg::Value reason) {});
2027 
2028 js.runMicrotasks();
2029 
2030 KJ_ASSERT(abortCalled, "abort algorithm was not called");
2031 KJ_ASSERT(closeRejected, "re-entrant close() should have been rejected");
2032 KJ_ASSERT(abortResolved, "abort should have resolved successfully");
2033 });
2034}
2035 
2036// Regression test: DrainingReader with pull that synchronously closes the stream.
2037// This exercises the case where the pull callback (triggered by onConsumerWantsData inside
2038// consumer->drainingRead()) synchronously closes the controller, causing a deferred state
2039// transition. Previously, ValueReadable::drainingRead() had its own beginOperation/endOperation
2040// pair that would fire the deferred transition BEFORE the returned promise's .then() callbacks
2041// (which capture `this` pointing at the Consumer) had a chance to run, causing use-after-free.
2042// The fix moves beginOperation() to before consumer->drainingRead() at the controller level,
2043// and endOperation() into the .then() callbacks, so the deferred transition only fires after
2044// the Consumer's callbacks have run.
2045KJ_TEST("DrainingReader: pull that synchronously closes does not UAF (value stream)") {
2046 preamble([](jsg::Lock& js) {
2047 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2048 // clang-format off
2049 rs->getController().setup(js, UnderlyingSource{
2050 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2051 KJ_SWITCH_ONEOF(controller) {
2052 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2053 // Synchronously close without enqueuing any data.
2054 // This triggers close -> onConsumerClose -> doClose -> deferTransitionTo<Closed>
2055 // while the drainingRead's pending read is still in flight.
2056 c->close(js);
2057 return js.resolvedPromise();
2058 }
2059 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2060 }
2061 KJ_UNREACHABLE;
2062 }
2063 }, StreamQueuingStrategy{.highWaterMark = 0});
2064 // clang-format on
2065 
2066 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2067 bool readCompleted = false;
2068 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2069 // Stream closed immediately, so we should get done=true.
2070 KJ_ASSERT(result.done);
2071 readCompleted = true;
2072 });
2073 
2074 js.runMicrotasks();
2075 KJ_ASSERT(readCompleted);
2076 } else {
2077 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2078 }
2079 });
2080}
2081 
2082KJ_TEST("DrainingReader: pull that synchronously closes does not UAF (byte stream)") {
2083 preamble([](jsg::Lock& js) {
2084 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2085 // clang-format off
2086 rs->getController().setup(js, UnderlyingSource{
2087 .type = kj::str("bytes"),
2088 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2089 KJ_SWITCH_ONEOF(controller) {
2090 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
2091 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
2092 c->close(js);
2093 return js.resolvedPromise();
2094 }
2095 }
2096 KJ_UNREACHABLE;
2097 }
2098 }, StreamQueuingStrategy{.highWaterMark = 0});
2099 // clang-format on
2100 
2101 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2102 bool readCompleted = false;
2103 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2104 KJ_ASSERT(result.done);
2105 readCompleted = true;
2106 });
2107 
2108 js.runMicrotasks();
2109 KJ_ASSERT(readCompleted);
2110 } else {
2111 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2112 }
2113 });
2114}
2115 
2116// Test that a pull which synchronously errors the stream doesn't UAF either.
2117KJ_TEST("DrainingReader: pull that synchronously errors does not UAF (value stream)") {
2118 preamble([](jsg::Lock& js) {
2119 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2120 // clang-format off
2121 rs->getController().setup(js, UnderlyingSource{
2122 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2123 KJ_SWITCH_ONEOF(controller) {
2124 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2125 c->error(js, js.v8TypeError("test error"_kj));
2126 return js.resolvedPromise();
2127 }
2128 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2129 }
2130 KJ_UNREACHABLE;
2131 }
2132 }, StreamQueuingStrategy{.highWaterMark = 0});
2133 // clang-format on
2134 
2135 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2136 bool readRejected = false;
2137 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2138 KJ_FAIL_ASSERT("Should have been rejected");
2139 }, [&](jsg::Lock& js, jsg::Value reason) { readRejected = true; });
2140 
2141 js.runMicrotasks();
2142 KJ_ASSERT(readRejected);
2143 } else {
2144 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2145 }
2146 });
2147}
2148 
2149// Test drainingRead when pull enqueues data then closes on the next pull.
2150KJ_TEST("DrainingReader: pull enqueues then closes on next pull (value stream)") {
2151 preamble([](jsg::Lock& js) {
2152 uint pullCount = 0;
2153 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2154 // clang-format off
2155 rs->getController().setup(js, UnderlyingSource{
2156 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2157 KJ_SWITCH_ONEOF(controller) {
2158 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2159 pullCount++;
2160 if (pullCount == 1) {
2161 c->enqueue(js, toBytes(js, kj::str("data")));
2162 } else {
2163 // Second pull: close synchronously without enqueuing.
2164 c->close(js);
2165 }
2166 return js.resolvedPromise();
2167 }
2168 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2169 }
2170 KJ_UNREACHABLE;
2171 }
2172 }, StreamQueuingStrategy{.highWaterMark = 0});
2173 // clang-format on
2174 
2175 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2176 // First read should get the data.
2177 bool firstReadDone = false;
2178 auto p1 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2179 KJ_ASSERT(result.chunks.size() == 1);
2180 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "data");
2181 KJ_ASSERT(!result.done);
2182 firstReadDone = true;
2183 });
2184 
2185 js.runMicrotasks();
2186 KJ_ASSERT(firstReadDone);
2187 
2188 // Second read triggers the close.
2189 bool secondReadDone = false;
2190 auto p2 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2191 KJ_ASSERT(result.done);
2192 secondReadDone = true;
2193 });
2194 
2195 js.runMicrotasks();
2196 KJ_ASSERT(secondReadDone);
2197 
2198 reader->releaseLock(js);
2199 } else {
2200 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2201 }
2202 });
2203}
2204 
2205// Regression test: DrainingReader with pull that synchronously cancels the controller.
2206// When cancel is called on the controller from within the pull callback, the controller's
2207// doCancel() transitions its impl state to Closed and destroys the Queue. The QueueImpl
2208// destructor detaches consumers (sets queue = kj::none) but does NOT cancel them — the
2209// consumer state remains Active. Without the fix, drainingRead falls through to the
2210// pending-read path and queues a ReadRequest that will never be fulfilled (promise hangs).
2211// The fix detects the detached queue and returns done immediately.
2212KJ_TEST("DrainingReader: pull that synchronously cancels does not hang (value stream)") {
2213 preamble([](jsg::Lock& js) {
2214 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2215 // clang-format off
2216 rs->getController().setup(js, UnderlyingSource{
2217 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2218 KJ_SWITCH_ONEOF(controller) {
2219 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2220 // Synchronously cancel without enqueuing any data.
2221 auto promise KJ_UNUSED = c->cancel(js, kj::none);
2222 return js.resolvedPromise();
2223 }
2224 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2225 }
2226 KJ_UNREACHABLE;
2227 }
2228 }, StreamQueuingStrategy{.highWaterMark = 0});
2229 // clang-format on
2230 
2231 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2232 bool readCompleted = false;
2233 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2234 // Stream canceled immediately with no data — should resolve with done=true.
2235 KJ_ASSERT(result.done);
2236 KJ_ASSERT(result.chunks.size() == 0);
2237 readCompleted = true;
2238 });
2239 
2240 js.runMicrotasks();
2241 KJ_ASSERT(readCompleted);
2242 } else {
2243 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2244 }
2245 });
2246}
2247 
2248KJ_TEST("DrainingReader: pull that synchronously cancels does not hang (byte stream)") {
2249 preamble([](jsg::Lock& js) {
2250 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2251 // clang-format off
2252 rs->getController().setup(js, UnderlyingSource{
2253 .type = kj::str("bytes"),
2254 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2255 KJ_SWITCH_ONEOF(controller) {
2256 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
2257 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
2258 auto promise KJ_UNUSED = c->cancel(js, kj::none);
2259 return js.resolvedPromise();
2260 }
2261 }
2262 KJ_UNREACHABLE;
2263 }
2264 }, StreamQueuingStrategy{.highWaterMark = 0});
2265 // clang-format on
2266 
2267 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2268 bool readCompleted = false;
2269 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2270 KJ_ASSERT(result.done);
2271 KJ_ASSERT(result.chunks.size() == 0);
2272 readCompleted = true;
2273 });
2274 
2275 js.runMicrotasks();
2276 KJ_ASSERT(readCompleted);
2277 } else {
2278 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2279 }
2280 });
2281}
2282 
2283// Test drainingRead when pull enqueues data then cancels on the next pull.
2284// The first read should deliver the enqueued data, and the second read
2285// (which triggers the cancel) should resolve with done=true and no data.
2286KJ_TEST("DrainingReader: pull enqueues then cancels on next pull (value stream)") {
2287 preamble([](jsg::Lock& js) {
2288 uint pullCount = 0;
2289 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2290 // clang-format off
2291 rs->getController().setup(js, UnderlyingSource{
2292 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2293 KJ_SWITCH_ONEOF(controller) {
2294 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2295 pullCount++;
2296 if (pullCount == 1) {
2297 c->enqueue(js, toBytes(js, kj::str("data")));
2298 } else {
2299 // Second pull: cancel synchronously without enqueuing.
2300 auto promise KJ_UNUSED = c->cancel(js, kj::none);
2301 }
2302 return js.resolvedPromise();
2303 }
2304 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2305 }
2306 KJ_UNREACHABLE;
2307 }
2308 }, StreamQueuingStrategy{.highWaterMark = 0});
2309 // clang-format on
2310 
2311 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2312 // First read should get the data.
2313 bool firstReadDone = false;
2314 auto p1 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2315 KJ_ASSERT(result.chunks.size() == 1);
2316 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "data");
2317 KJ_ASSERT(!result.done);
2318 firstReadDone = true;
2319 });
2320 
2321 js.runMicrotasks();
2322 KJ_ASSERT(firstReadDone);
2323 
2324 // Second read triggers the cancel — should resolve with done=true.
2325 bool secondReadDone = false;
2326 auto p2 = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2327 KJ_ASSERT(result.done);
2328 secondReadDone = true;
2329 });
2330 
2331 js.runMicrotasks();
2332 KJ_ASSERT(secondReadDone);
2333 
2334 reader->releaseLock(js);
2335 } else {
2336 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2337 }
2338 });
2339}
2340 
2341// Test that a pending error applied during wrapDrainingRead's endOperation() propagates
2342// as a rejection rather than silently returning data. This exercises the defense-in-depth
2343// fix in wrapDrainingRead where endOperation() applies a deferred Errored state.
2344//
2345// Scenario: Pull enqueues data synchronously but returns a rejected promise. The rejection
2346// handler (which errors the stream) runs as a microtask between drainingRead's resolved
2347// promise and wrapDrainingRead's .then() callback. The data collected by drainingRead should
2348// be discarded because the error was applied during the operation.
2349KJ_TEST("DrainingReader: pending error in endOperation rejects read (value stream)") {
2350 preamble([](jsg::Lock& js) {
2351 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2352 // clang-format off
2353 rs->getController().setup(js, UnderlyingSource{
2354 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
2355 KJ_SWITCH_ONEOF(controller) {
2356 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2357 // Enqueue data synchronously — drainingRead will collect it.
2358 c->enqueue(js, toBytes(js, kj::str("should-be-discarded")));
2359 // Return rejected promise — the pull failure handler runs as a microtask
2360 // and calls doError(), which defers the error because beginOperation() is
2361 // active. When wrapDrainingRead's endOperation() fires, it applies the
2362 // pending error and should throw rather than returning the data.
2363 return js.rejectedPromise<void>(js.v8TypeError("pull failed"_kj));
2364 }
2365 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2366 }
2367 KJ_UNREACHABLE;
2368 }
2369 }, StreamQueuingStrategy{.highWaterMark = 0});
2370 // clang-format on
2371 
2372 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2373 bool readRejected = false;
2374 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2375 KJ_FAIL_ASSERT("Should have rejected, not resolved with data");
2376 }, [&](jsg::Lock& js, jsg::Value&& err) { readRejected = true; });
2377 
2378 js.runMicrotasks();
2379 KJ_ASSERT(readRejected, "Read should have rejected due to pending error");
2380 
2381 reader->releaseLock(js);
2382 } else {
2383 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2384 }
2385 });
2386}
2387 
2388KJ_TEST("DrainingReader: pending error in endOperation rejects read (byte stream)") {
2389 preamble([](jsg::Lock& js) {
2390 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2391 // clang-format off
2392 rs->getController().setup(js, UnderlyingSource{
2393 .type = kj::str("bytes"),
2394 .pull = [](jsg::Lock& js, UnderlyingSource::Controller controller) {
2395 KJ_SWITCH_ONEOF(controller) {
2396 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
2397 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
2398 c->enqueue(js, toBufferSource(js, kj::str("should-be-discarded")));
2399 return js.rejectedPromise<void>(js.v8TypeError("pull failed"_kj));
2400 }
2401 }
2402 KJ_UNREACHABLE;
2403 }
2404 }, StreamQueuingStrategy{.highWaterMark = 0});
2405 // clang-format on
2406 
2407 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2408 bool readRejected = false;
2409 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2410 KJ_FAIL_ASSERT("Should have rejected, not resolved with data");
2411 }, [&](jsg::Lock& js, jsg::Value&& err) { readRejected = true; });
2412 
2413 js.runMicrotasks();
2414 KJ_ASSERT(readRejected, "Read should have rejected due to pending error");
2415 
2416 reader->releaseLock(js);
2417 } else {
2418 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2419 }
2420 });
2421}
2422 
2423// Regression test: After drainingRead returns done=true because the stream was closed,
2424// the controller-level close should also be finalized (not deferred until GC).
2425//
2426// When pull enqueues data AND closes in the same call, the sequence is:
2427// 1. c->enqueue() pushes data to consumer buffer
2428// 2. c->close() → ConsumerImpl::close() pushes Close marker, calls maybeDrainAndSetState()
2429// 3. maybeDrainAndSetState() finds queueTotalSize > 0 (data still in buffer), can't finalize
2430// 4. Back in drainingRead, drainBuffer() drains the data and sees the Close marker
2431// 5. drainingRead returns done=true, but the consumer is still Active
2432//
2433// Without the fix, the consumer Close marker is never popped and maybeDrainAndSetState()
2434// is never called again, so onConsumerClose never fires and the controller stays open.
2435// The fix calls impl.maybeDrainAndSetState(js) after draining when isClosing is true.
2436KJ_TEST("DrainingReader: controller closes promptly after drainingRead done (value stream)") {
2437 preamble([](jsg::Lock& js) {
2438 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2439 // clang-format off
2440 rs->getController().setup(js, UnderlyingSource{
2441 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2442 KJ_SWITCH_ONEOF(controller) {
2443 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {
2444 // Enqueue data and close in the same pull. This causes
2445 // ConsumerImpl::close() → maybeDrainAndSetState() to find non-empty
2446 // buffer, preventing immediate finalization.
2447 c->enqueue(js, toBytes(js, kj::str("hello")));
2448 c->close(js);
2449 return js.resolvedPromise();
2450 }
2451 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {}
2452 }
2453 KJ_UNREACHABLE;
2454 }
2455 }, StreamQueuingStrategy{.highWaterMark = 0});
2456 // clang-format on
2457 
2458 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2459 bool readDone = false;
2460 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2461 // Should get data with done=true (data + close in same batch).
2462 KJ_ASSERT(result.done);
2463 KJ_ASSERT(result.chunks.size() == 1);
2464 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "hello");
2465 readDone = true;
2466 });
2467 js.runMicrotasks();
2468 KJ_ASSERT(readDone);
2469 
2470 // The controller should now be in the Closed state — not deferred.
2471 KJ_ASSERT(rs->getController().isClosed(),
2472 "controller should be closed after drainingRead returns done");
2473 
2474 reader->releaseLock(js);
2475 } else {
2476 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2477 }
2478 });
2479}
2480 
2481KJ_TEST("DrainingReader: controller closes promptly after drainingRead done (byte stream)") {
2482 preamble([](jsg::Lock& js) {
2483 auto rs = js.alloc<ReadableStream>(newReadableStreamJsController());
2484 // clang-format off
2485 rs->getController().setup(js, UnderlyingSource{
2486 .type = kj::str("bytes"),
2487 .pull = [&](jsg::Lock& js, UnderlyingSource::Controller controller) {
2488 KJ_SWITCH_ONEOF(controller) {
2489 KJ_CASE_ONEOF(c, jsg::Ref<ReadableStreamDefaultController>) {}
2490 KJ_CASE_ONEOF(c, jsg::Ref<ReadableByteStreamController>) {
2491 // Enqueue data and close in the same pull.
2492 c->enqueue(js, toBufferSource(js, kj::str("world")));
2493 c->close(js);
2494 return js.resolvedPromise();
2495 }
2496 }
2497 KJ_UNREACHABLE;
2498 }
2499 }, StreamQueuingStrategy{.highWaterMark = 0});
2500 // clang-format on
2501 
2502 KJ_IF_SOME(reader, DrainingReader::create(js, *rs)) {
2503 bool readDone = false;
2504 auto promise = reader->read(js).then(js, [&](jsg::Lock& js, DrainingReadResult&& result) {
2505 KJ_ASSERT(result.done);
2506 KJ_ASSERT(result.chunks.size() == 1);
2507 KJ_ASSERT(kj::str(result.chunks[0].asChars()) == "world");
2508 readDone = true;
2509 });
2510 js.runMicrotasks();
2511 KJ_ASSERT(readDone);
2512 
2513 // The controller should now be in the Closed state — not deferred.
2514 KJ_ASSERT(rs->getController().isClosed(),
2515 "controller should be closed after drainingRead returns done");
2516 
2517 reader->releaseLock(js);
2518 } else {
2519 KJ_FAIL_ASSERT("Failed to create DrainingReader");
2520 }
2521 });
2522}
2523 
2524} // namespace
2525} // namespace workerd::api