Skip to content
File

Blob: src/workerd/api/streams/readable-source-adapter-test.c++

69.0 KB
1#include "readable-source-adapter.h"
2#include "standard.h"
3#include "writable-sink.h"
4 
5#include <workerd/api/system-streams.h>
6#include <workerd/jsg/jsg-test.h>
7#include <workerd/jsg/jsg.h>
8#include <workerd/tests/test-fixture.h>
9#include <workerd/util/own-util.h>
10#include <workerd/util/stream-utils.h>
11 
12namespace workerd::api::streams {
13 
14namespace {
15 
16struct RecordingSource final: public kj::AsyncInputStream {
17 size_t readCalled = 0;
18 
19 kj::Promise<size_t> tryRead(void*, size_t minBytes, size_t maxBytes) override {
20 readCalled++;
21 co_return 0;
22 }
23 
24 kj::Maybe<uint64_t> tryGetLength() override {
25 static const uint64_t length = 42;
26 return length;
27 }
28};
29 
30struct NeverDoneSource final: public kj::AsyncInputStream {
31 size_t readCalled = 0;
32 
33 kj::Promise<size_t> tryRead(void* ptr, size_t minBytes, size_t maxBytes) override {
34 readCalled++;
35 kj::ArrayPtr<kj::byte> buffer(static_cast<kj::byte*>(ptr), maxBytes);
36 buffer.fill('a');
37 return maxBytes;
38 }
39 
40 kj::Maybe<uint64_t> tryGetLength() override {
41 return kj::none;
42 }
43};
44 
45struct MinimalReadSource final: public kj::AsyncInputStream {
46 size_t readCalled = 0;
47 
48 kj::Promise<size_t> tryRead(void* ptr, size_t minBytes, size_t maxBytes) override {
49 readCalled++;
50 kj::ArrayPtr<kj::byte> buffer(static_cast<kj::byte*>(ptr), minBytes);
51 buffer.fill('a');
52 return minBytes;
53 }
54 
55 kj::Maybe<uint64_t> tryGetLength() override {
56 return kj::none;
57 }
58};
59 
60struct FiniteReadSource final: public kj::AsyncInputStream {
61 size_t readCalled = 0;
62 size_t maxReads;
63 
64 FiniteReadSource(size_t maxReads): maxReads(maxReads) {}
65 
66 kj::Promise<size_t> tryRead(void* ptr, size_t minBytes, size_t maxBytes) override {
67 if (readCalled >= maxReads) {
68 co_return 0;
69 }
70 readCalled++;
71 kj::ArrayPtr<kj::byte> buffer(static_cast<kj::byte*>(ptr), minBytes);
72 buffer.fill('a');
73 co_return minBytes;
74 }
75 
76 kj::Maybe<uint64_t> tryGetLength() override {
77 return kj::none;
78 }
79};
80 
81} // namespace
82 
83KJ_TEST("Test successful construction with valid ReadableStreamSource") {
84 TestFixture fixture;
85 RecordingSource source;
86 
87 fixture.runInIoContext([&](const TestFixture::Environment& env) {
88 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
89 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
90 env.js, env.context, newReadableSource(kj::mv(fake)));
91 
92 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
93 KJ_ASSERT(
94 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
95 
96 return kj::READY_NOW;
97 });
98}
99 
100KJ_TEST("Adapter shutdown with no reads") {
101 TestFixture fixture;
102 RecordingSource source;
103 
104 fixture.runInIoContext([&](const TestFixture::Environment& env) {
105 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
106 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
107 env.js, env.context, newReadableSource(kj::mv(fake)));
108 
109 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
110 KJ_ASSERT(
111 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
112 
113 adapter->shutdown(env.js);
114 adapter->shutdown(env.js); // second call is no-op
115 
116 // Read after shutdown should be resolved immediate
117 auto read = adapter->read(env.js,
118 ReadableStreamSourceJsAdapter::ReadOptions{
119 .buffer = jsg::BufferSource(env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, 10)),
120 });
121 KJ_ASSERT(read.getState(env.js) ==
122 jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult>::State::FULFILLED,
123 "Read after shutdown should be resolved immediately");
124 
125 KJ_ASSERT(adapter->isClosed(), "Adapter shoud be closed after shutdown()");
126 KJ_ASSERT(adapter->isCanceled() == kj::none, "Adapter should not be canceled after shutdown()");
127 
128 return kj::READY_NOW;
129 });
130}
131 
132KJ_TEST("Adapter cancel with no reads") {
133 TestFixture fixture;
134 RecordingSource source;
135 
136 fixture.runInIoContext([&](const TestFixture::Environment& env) {
137 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
138 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
139 env.js, env.context, newReadableSource(kj::mv(fake)));
140 
141 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
142 KJ_ASSERT(
143 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
144 
145 adapter->cancel(env.js, env.js.error("boom"));
146 
147 auto read = adapter->read(env.js,
148 ReadableStreamSourceJsAdapter::ReadOptions{
149 .buffer = jsg::BufferSource(env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, 10)),
150 });
151 KJ_ASSERT(read.getState(env.js) ==
152 jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult>::State::REJECTED,
153 "Read after shutdown should be rejected immediately");
154 
155 adapter->shutdown(env.js); // shutdown after cancel is no-op
156 
157 KJ_ASSERT(!adapter->isClosed(), "Adapter shoud be canceled, not closed");
158 auto& ex = KJ_ASSERT_NONNULL(adapter->isCanceled());
159 KJ_ASSERT(ex.getDescription().contains("boom"),
160 "Adapter should be in canceled state with provided exception");
161 
162 return kj::READY_NOW;
163 });
164}
165 
166KJ_TEST("Adapter cancel (kj::Exception) with no reads") {
167 TestFixture fixture;
168 RecordingSource source;
169 
170 fixture.runInIoContext([&](const TestFixture::Environment& env) {
171 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
172 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
173 env.js, env.context, newReadableSource(kj::mv(fake)));
174 
175 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
176 KJ_ASSERT(
177 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
178 
179 adapter->cancel(KJ_EXCEPTION(FAILED, "boom"));
180 
181 KJ_ASSERT(!adapter->isClosed(), "Adapter shoud be canceled, not closed");
182 auto& ex = KJ_ASSERT_NONNULL(adapter->isCanceled());
183 KJ_ASSERT(ex.getDescription().contains("boom"),
184 "Adapter should be in canceled state with provided exception");
185 
186 return kj::READY_NOW;
187 });
188}
189 
190KJ_TEST("Adapter with single read (ArrayBuffer)") {
191 TestFixture fixture;
192 NeverDoneSource source;
193 
194 fixture.runInIoContext([&](const TestFixture::Environment& env) {
195 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
196 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
197 env.js, env.context, newReadableSource(kj::mv(fake)));
198 
199 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
200 KJ_ASSERT(
201 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
202 
203 const size_t bufferSize = 10;
204 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, bufferSize);
205 
206 return env.context
207 .awaitJs(env.js,
208 adapter
209 ->read(env.js,
210 ReadableStreamSourceJsAdapter::ReadOptions{
211 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
212 .minBytes = 5,
213 })
214 .then(env.js, [](jsg::Lock& js, auto result) {
215 KJ_ASSERT(!result.done, "Stream should not be done yet");
216 KJ_ASSERT(result.buffer.asArrayPtr().size() == 10, "Read buffer should be full size");
217 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaa"_kjb);
218 
219 // BufferSource should be an ArrayBuffer
220 auto handle = result.buffer.getHandle(js);
221 KJ_ASSERT(handle->IsArrayBuffer());
222 })).attach(kj::mv(adapter));
223 });
224}
225 
226KJ_TEST("Adapter with single read (Uint8Array)") {
227 TestFixture fixture;
228 NeverDoneSource source;
229 
230 fixture.runInIoContext([&](const TestFixture::Environment& env) {
231 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
232 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
233 env.js, env.context, newReadableSource(kj::mv(fake)));
234 
235 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
236 KJ_ASSERT(
237 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
238 
239 const size_t bufferSize = 10;
240 auto backing = jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize);
241 
242 return env.context
243 .awaitJs(env.js,
244 adapter
245 ->read(env.js,
246 ReadableStreamSourceJsAdapter::ReadOptions{
247 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
248 .minBytes = 5,
249 })
250 .then(env.js, [](jsg::Lock& js, auto result) {
251 KJ_ASSERT(!result.done, "Stream should not be done yet");
252 KJ_ASSERT(result.buffer.asArrayPtr().size() == 10, "Read buffer should be full size");
253 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaa"_kjb);
254 
255 // BufferSource should be an ArrayBuffer
256 auto handle = result.buffer.getHandle(js);
257 KJ_ASSERT(handle->IsUint8Array());
258 })).attach(kj::mv(adapter));
259 });
260}
261 
262KJ_TEST("Adapter with single read (Int32Array)") {
263 TestFixture fixture;
264 NeverDoneSource source;
265 
266 fixture.runInIoContext([&](const TestFixture::Environment& env) {
267 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
268 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
269 env.js, env.context, newReadableSource(kj::mv(fake)));
270 
271 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
272 KJ_ASSERT(
273 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
274 
275 const size_t bufferSize = 16;
276 auto backing = jsg::BackingStore::alloc<v8::Int32Array>(env.js, bufferSize);
277 
278 return env.context
279 .awaitJs(env.js,
280 adapter
281 ->read(env.js,
282 ReadableStreamSourceJsAdapter::ReadOptions{
283 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
284 .minBytes = 5,
285 })
286 .then(env.js, [](jsg::Lock& js, auto result) {
287 KJ_ASSERT(!result.done, "Stream should not be done yet");
288 KJ_ASSERT(result.buffer.asArrayPtr().size() == 16, "Read buffer should be full size");
289 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaaaaaaaa"_kjb);
290 
291 // BufferSource should be an ArrayBuffer
292 auto handle = result.buffer.getHandle(js);
293 KJ_ASSERT(handle->IsInt32Array());
294 })).attach(kj::mv(adapter));
295 });
296}
297 
298KJ_TEST("Adapter with single large read (ArrayBuffer)") {
299 TestFixture fixture;
300 NeverDoneSource source;
301 
302 fixture.runInIoContext([&](const TestFixture::Environment& env) {
303 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
304 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
305 env.js, env.context, newReadableSource(kj::mv(fake)));
306 
307 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
308 KJ_ASSERT(
309 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
310 
311 const size_t bufferSize = 16 * 1024;
312 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, bufferSize);
313 
314 return env.context
315 .awaitJs(env.js,
316 adapter
317 ->read(env.js,
318 ReadableStreamSourceJsAdapter::ReadOptions{
319 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
320 .minBytes = 5,
321 })
322 .then(env.js, [](jsg::Lock& js, auto result) {
323 KJ_ASSERT(!result.done, "Stream should not be done yet");
324 KJ_ASSERT(result.buffer.asArrayPtr().size() == 16 * 1024, "Read buffer should be full size");
325 
326 // BufferSource should be an ArrayBuffer
327 auto handle = result.buffer.getHandle(js);
328 KJ_ASSERT(handle->IsArrayBuffer());
329 })).attach(kj::mv(adapter));
330 });
331}
332 
333KJ_TEST("Adapter with single small read (ArrayBuffer)") {
334 TestFixture fixture;
335 NeverDoneSource source;
336 
337 fixture.runInIoContext([&](const TestFixture::Environment& env) {
338 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
339 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
340 env.js, env.context, newReadableSource(kj::mv(fake)));
341 
342 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
343 KJ_ASSERT(
344 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
345 
346 const size_t bufferSize = 1;
347 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, bufferSize);
348 
349 return env.context
350 .awaitJs(env.js,
351 adapter
352 ->read(env.js,
353 ReadableStreamSourceJsAdapter::ReadOptions{
354 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
355 .minBytes = 5,
356 })
357 .then(env.js, [](jsg::Lock& js, auto result) {
358 KJ_ASSERT(!result.done, "Stream should not be done yet");
359 KJ_ASSERT(result.buffer.asArrayPtr().size() == 1, "Read buffer should be full size");
360 
361 // BufferSource should be an ArrayBuffer
362 auto handle = result.buffer.getHandle(js);
363 KJ_ASSERT(handle->IsArrayBuffer());
364 })).attach(kj::mv(adapter));
365 });
366}
367 
368KJ_TEST("Adapter with minimal reads (Uint8Array)") {
369 TestFixture fixture;
370 MinimalReadSource source;
371 
372 fixture.runInIoContext([&](const TestFixture::Environment& env) {
373 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
374 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
375 env.js, env.context, newReadableSource(kj::mv(fake)));
376 
377 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
378 KJ_ASSERT(
379 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
380 
381 const size_t bufferSize = 10;
382 auto backing = jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize);
383 
384 auto promise = adapter
385 ->read(env.js,
386 ReadableStreamSourceJsAdapter::ReadOptions{
387 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
388 .minBytes = 3,
389 })
390 .then(env.js, [](jsg::Lock& js, auto result) {
391 KJ_ASSERT(!result.done, "Stream should not be done yet");
392 KJ_ASSERT(result.buffer.asArrayPtr().size() == 3, "Read buffer should be three bytes");
393 KJ_ASSERT(result.buffer.asArrayPtr() == "aaa"_kjb);
394 
395 // BufferSource should be an ArrayBuffer
396 auto handle = result.buffer.getHandle(js);
397 KJ_ASSERT(handle->IsUint8Array());
398 });
399 
400 return env.context.awaitJs(env.js, kj::mv(promise)).attach(kj::mv(adapter));
401 });
402}
403 
404KJ_TEST("Adapter with minimal reads (Uint32Array)") {
405 TestFixture fixture;
406 MinimalReadSource source;
407 
408 fixture.runInIoContext([&](const TestFixture::Environment& env) {
409 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
410 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
411 env.js, env.context, newReadableSource(kj::mv(fake)));
412 
413 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
414 KJ_ASSERT(
415 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
416 
417 const size_t bufferSize = 16;
418 auto backing = jsg::BackingStore::alloc<v8::Uint32Array>(env.js, bufferSize);
419 
420 auto promise = adapter
421 ->read(env.js,
422 ReadableStreamSourceJsAdapter::ReadOptions{
423 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
424 .minBytes = 3, // Impl with round up to 4
425 })
426 .then(env.js, [](jsg::Lock& js, auto result) {
427 KJ_ASSERT(!result.done, "Stream should not be done yet");
428 KJ_ASSERT(result.buffer.asArrayPtr().size() == 4, "Read buffer should be four bytes");
429 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaa"_kjb);
430 
431 // BufferSource should be an ArrayBuffer
432 auto handle = result.buffer.getHandle(js);
433 KJ_ASSERT(handle->IsUint32Array());
434 });
435 
436 return env.context.awaitJs(env.js, kj::mv(promise)).attach(kj::mv(adapter));
437 });
438}
439 
440KJ_TEST("Adapter with over large min reads (Uint32Array)") {
441 TestFixture fixture;
442 MinimalReadSource source;
443 
444 fixture.runInIoContext([&](const TestFixture::Environment& env) {
445 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
446 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
447 env.js, env.context, newReadableSource(kj::mv(fake)));
448 
449 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
450 KJ_ASSERT(
451 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
452 
453 const size_t bufferSize = 16;
454 auto backing = jsg::BackingStore::alloc<v8::Uint32Array>(env.js, bufferSize);
455 
456 auto promise = adapter
457 ->read(env.js,
458 ReadableStreamSourceJsAdapter::ReadOptions{
459 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
460 .minBytes = 24, // Impl with round up to 4
461 })
462 .then(env.js, [](jsg::Lock& js, auto result) {
463 KJ_ASSERT(!result.done, "Stream should not be done yet");
464 KJ_ASSERT(result.buffer.asArrayPtr().size() == 16, "Read buffer should be four bytes");
465 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaaaaaaaa"_kjb);
466 
467 // BufferSource should be an ArrayBuffer
468 auto handle = result.buffer.getHandle(js);
469 KJ_ASSERT(handle->IsUint32Array());
470 });
471 
472 return env.context.awaitJs(env.js, kj::mv(promise)).attach(kj::mv(adapter));
473 });
474}
475 
476KJ_TEST("Adapter with over large min reads (Uint32Array)") {
477 TestFixture fixture;
478 
479 fixture.runInIoContext([&](const TestFixture::Environment& env) {
480 auto source = newReadableSource(newNullInputStream());
481 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(env.js, env.context, kj::mv(source));
482 
483 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
484 KJ_ASSERT(
485 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
486 
487 const size_t bufferSize = 1;
488 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, bufferSize);
489 
490 auto promise = adapter
491 ->read(env.js,
492 ReadableStreamSourceJsAdapter::ReadOptions{
493 .buffer = jsg::BufferSource(env.js, kj::mv(backing)),
494 })
495 .then(env.js, [](jsg::Lock& js, auto result) {
496 KJ_ASSERT(result.done, "Stream should be done");
497 KJ_ASSERT(result.buffer.asArrayPtr().size() == 0, "Read buffer should be 0 bytes");
498 auto handle = result.buffer.getHandle(js);
499 KJ_ASSERT(handle->IsArrayBuffer());
500 });
501 
502 return env.context.awaitJs(env.js, kj::mv(promise)).attach(kj::mv(adapter));
503 });
504}
505 
506KJ_TEST("Adapter with multiple reads (Uint8Array)") {
507 TestFixture fixture;
508 NeverDoneSource source;
509 
510 fixture.runInIoContext([&](const TestFixture::Environment& env) {
511 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
512 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
513 env.js, env.context, newReadableSource(kj::mv(fake)));
514 
515 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
516 KJ_ASSERT(
517 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
518 
519 const size_t bufferSize = 10;
520 
521 auto read1 = adapter->read(env.js,
522 ReadableStreamSourceJsAdapter::ReadOptions{
523 .buffer = jsg::BufferSource(
524 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
525 });
526 auto read2 = adapter->read(env.js,
527 ReadableStreamSourceJsAdapter::ReadOptions{
528 .buffer = jsg::BufferSource(
529 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
530 });
531 auto read3 = adapter->read(env.js,
532 ReadableStreamSourceJsAdapter::ReadOptions{
533 .buffer = jsg::BufferSource(
534 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
535 });
536 
537 return env.context
538 .awaitJs(env.js,
539 read1
540 .then(env.js,
541 [read2 = kj::mv(read2)](jsg::Lock& js, auto result) mutable {
542 KJ_ASSERT(!result.done, "Stream should not be done yet");
543 KJ_ASSERT(result.buffer.asArrayPtr().size() == 10, "Read buffer should be full size");
544 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaa"_kjb);
545 return kj::mv(read2);
546 })
547 .then(env.js, [read3 = kj::mv(read3)](jsg::Lock& js, auto result) mutable {
548 KJ_ASSERT(!result.done, "Stream should not be done yet");
549 KJ_ASSERT(result.buffer.asArrayPtr().size() == 10, "Read buffer should be full size");
550 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaa"_kjb);
551 return kj::mv(read3);
552 }).then(env.js, [](jsg::Lock& js, auto result) mutable {
553 KJ_ASSERT(!result.done, "Stream should not be done yet");
554 KJ_ASSERT(result.buffer.asArrayPtr().size() == 10, "Read buffer should be full size");
555 KJ_ASSERT(result.buffer.asArrayPtr() == "aaaaaaaaaa"_kjb);
556 return js.resolvedPromise();
557 })).attach(kj::mv(adapter));
558 });
559}
560 
561KJ_TEST("Adapter with multiple reads shutdown") {
562 TestFixture fixture;
563 NeverDoneSource source;
564 
565 fixture.runInIoContext([&](const TestFixture::Environment& env) {
566 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
567 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
568 env.js, env.context, newReadableSource(kj::mv(fake)));
569 
570 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
571 KJ_ASSERT(
572 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
573 
574 const size_t bufferSize = 10;
575 
576 auto read1 = adapter->read(env.js,
577 ReadableStreamSourceJsAdapter::ReadOptions{
578 .buffer = jsg::BufferSource(
579 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
580 });
581 auto read2 = adapter->read(env.js,
582 ReadableStreamSourceJsAdapter::ReadOptions{
583 .buffer = jsg::BufferSource(
584 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
585 });
586 auto read3 = adapter->read(env.js,
587 ReadableStreamSourceJsAdapter::ReadOptions{
588 .buffer = jsg::BufferSource(
589 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
590 });
591 
592 adapter->shutdown(env.js);
593 
594 return env.context
595 .awaitJs(env.js,
596 read1
597 .then(env.js,
598 [](jsg::Lock& js, auto result) {
599 return js.rejectedPromise<ReadableStreamSourceJsAdapter::ReadResult>(
600 js.error("Should not have completed read after shutdown"));
601 },
602 [read2 = kj::mv(read2)](jsg::Lock& js, jsg::Value exception) mutable
603 -> jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> {
604 return kj::mv(read2);
605 })
606 .then(env.js,
607 [](jsg::Lock& js, auto result) {
608 return js.rejectedPromise<ReadableStreamSourceJsAdapter::ReadResult>(
609 js.error("Should not have completed read after shutdown"));
610 },
611 [read3 = kj::mv(read3)](jsg::Lock& js, jsg::Value exception) mutable
612 -> jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> {
613 return kj::mv(read3);
614 }).then(env.js, [](jsg::Lock& js, auto result) {
615 return js.rejectedPromise<void>(js.error("Should not have completed read after shutdown"));
616 }, [](jsg::Lock& js, jsg::Value exception) mutable -> jsg::Promise<void> {
617 return js.resolvedPromise();
618 })).attach(kj::mv(adapter));
619 });
620}
621 
622KJ_TEST("Adapter with multiple reads cancel") {
623 TestFixture fixture;
624 NeverDoneSource source;
625 
626 fixture.runInIoContext([&](const TestFixture::Environment& env) {
627 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
628 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
629 env.js, env.context, newReadableSource(kj::mv(fake)));
630 
631 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
632 KJ_ASSERT(
633 adapter->isCanceled() == kj::none, "Adapter should not be canceled upon construction");
634 
635 const size_t bufferSize = 10;
636 
637 auto read1 = adapter->read(env.js,
638 ReadableStreamSourceJsAdapter::ReadOptions{
639 .buffer = jsg::BufferSource(
640 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
641 });
642 auto read2 = adapter->read(env.js,
643 ReadableStreamSourceJsAdapter::ReadOptions{
644 .buffer = jsg::BufferSource(
645 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
646 });
647 auto read3 = adapter->read(env.js,
648 ReadableStreamSourceJsAdapter::ReadOptions{
649 .buffer = jsg::BufferSource(
650 env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, bufferSize)),
651 });
652 
653 adapter->cancel(env.js, env.js.error("boom"));
654 adapter->cancel(env.js, env.js.error("bang"));
655 
656 return env.context
657 .awaitJs(env.js,
658 read1
659 .then(env.js,
660 [](jsg::Lock& js, auto result) {
661 return js.rejectedPromise<ReadableStreamSourceJsAdapter::ReadResult>(
662 js.error("Should not have completed read after shutdown"));
663 },
664 [read2 = kj::mv(read2)](jsg::Lock& js, jsg::Value exception) mutable
665 -> jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> {
666 auto handle = exception.getHandle(js);
667 KJ_ASSERT(kj::str(handle).contains("boom"),
668 "Read should have been rejected with cancelation error");
669 return kj::mv(read2);
670 })
671 .then(env.js,
672 [](jsg::Lock& js, auto result) {
673 return js.rejectedPromise<ReadableStreamSourceJsAdapter::ReadResult>(
674 js.error("Should not have completed read after shutdown"));
675 },
676 [read3 = kj::mv(read3)](jsg::Lock& js, jsg::Value exception) mutable
677 -> jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> {
678 auto handle = exception.getHandle(js);
679 KJ_ASSERT(kj::str(handle).contains("boom"),
680 "Read should have been rejected with cancelation error");
681 return kj::mv(read3);
682 }).then(env.js, [](jsg::Lock& js, auto result) {
683 return js.rejectedPromise<void>(js.error("Should not have completed read after shutdown"));
684 }, [](jsg::Lock& js, jsg::Value exception) mutable -> jsg::Promise<void> {
685 auto handle = exception.getHandle(js);
686 KJ_ASSERT(kj::str(handle).contains("boom"),
687 "Read should have been rejected with cancelation error");
688 return js.resolvedPromise();
689 })).attach(kj::mv(adapter));
690 });
691}
692 
693KJ_TEST("Adapter close after read") {
694 TestFixture fixture;
695 NeverDoneSource source;
696 
697 fixture.runInIoContext([&](const TestFixture::Environment& env) {
698 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
699 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
700 env.js, env.context, newReadableSource(kj::mv(fake)));
701 
702 auto read = adapter->read(env.js,
703 ReadableStreamSourceJsAdapter::ReadOptions{
704 .buffer = jsg::BufferSource(env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, 10)),
705 });
706 
707 auto closePromise = adapter->close(env.js);
708 
709 return env.context
710 .awaitJs(env.js,
711 closePromise.then(
712 env.js, [&adapter = *adapter, read = kj::mv(read)](jsg::Lock& js) mutable {
713 KJ_ASSERT(adapter.isClosed(), "Adapter should be closed after close()");
714 KJ_ASSERT(adapter.isCanceled() == kj::none, "Adapter should not be canceled after close()");
715 
716 KJ_ASSERT(read.getState(js) ==
717 jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult>::State::FULFILLED,
718 "Read should have completed successfully before close()");
719 })).attach(kj::mv(adapter));
720 });
721}
722 
723KJ_TEST("Adapter close") {
724 TestFixture fixture;
725 NeverDoneSource source;
726 
727 fixture.runInIoContext([&](const TestFixture::Environment& env) {
728 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
729 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
730 env.js, env.context, newReadableSource(kj::mv(fake)));
731 auto closePromise = adapter->close(env.js);
732 
733 // reads after close should be resoved immediately.
734 auto read = adapter->read(env.js,
735 ReadableStreamSourceJsAdapter::ReadOptions{
736 .buffer = jsg::BufferSource(env.js, jsg::BackingStore::alloc<v8::Uint8Array>(env.js, 10)),
737 });
738 KJ_ASSERT(read.getState(env.js) ==
739 jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult>::State::FULFILLED,
740 "Read after close should be fullfilled immediately");
741 
742 return env.context
743 .awaitJs(env.js, closePromise.then(env.js, [&adapter = *adapter](jsg::Lock& js) {
744 KJ_ASSERT(adapter.isClosed(), "Adapter should be closed after close()");
745 KJ_ASSERT(adapter.isCanceled() == kj::none, "Adapter should not be canceled after close()");
746 })).attach(kj::mv(adapter));
747 });
748}
749 
750KJ_TEST("Adapter close superseded by cancel") {
751 TestFixture fixture;
752 NeverDoneSource source;
753 
754 fixture.runInIoContext([&](const TestFixture::Environment& env) {
755 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
756 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
757 env.js, env.context, newReadableSource(kj::mv(fake)));
758 
759 auto closePromise = adapter->close(env.js);
760 
761 adapter->cancel(env.js, env.js.error("boom"));
762 
763 return env.context
764 .awaitJs(env.js, closePromise.then(env.js, [](jsg::Lock& js) {
765 return js.rejectedPromise<void>(js.error("Should not have completed close after cancel"));
766 }, [](jsg::Lock& js, jsg::Value exception) {
767 auto handle = exception.getHandle(js);
768 KJ_ASSERT(kj::str(handle).contains("boom"),
769 "Close should have been rejected with cancelation error");
770 return js.resolvedPromise();
771 })).attach(kj::mv(adapter));
772 });
773}
774 
775KJ_TEST("After read BackingStore maintains identity") {
776 TestFixture fixture;
777 NeverDoneSource source;
778 
779 fixture.runInIoContext([&](const TestFixture::Environment& env) {
780 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
781 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
782 env.js, env.context, newReadableSource(kj::mv(fake)));
783 
784 std::unique_ptr<v8::BackingStore> backing =
785 v8::ArrayBuffer::NewBackingStore(env.js.v8Isolate, 10);
786 auto* backingPtr = backing.get();
787 v8::Local<v8::ArrayBuffer> originalArrayBuffer =
788 v8::ArrayBuffer::New(env.js.v8Isolate, kj::mv(backing));
789 jsg::BufferSource source(env.js, originalArrayBuffer);
790 
791 return env.context
792 .awaitJs(env.js,
793 adapter
794 ->read(env.js,
795 ReadableStreamSourceJsAdapter::ReadOptions{
796 .buffer = jsg::BufferSource(env.js, originalArrayBuffer),
797 .minBytes = 5,
798 })
799 .then(env.js, [backingPtr](jsg::Lock& js, auto result) {
800 auto handle = result.buffer.getHandle(js);
801 KJ_ASSERT(handle->IsArrayBuffer());
802 auto backing = handle.template As<v8::ArrayBuffer>()->GetBackingStore();
803 KJ_ASSERT(backing.get() == backingPtr);
804 return js.resolvedPromise();
805 })).attach(kj::mv(adapter));
806 });
807}
808 
809KJ_TEST("Read all text") {
810 TestFixture fixture;
811 FiniteReadSource source(4);
812 
813 fixture.runInIoContext([&](const TestFixture::Environment& env) {
814 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
815 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
816 env.js, env.context, newReadableSource(kj::mv(fake)));
817 
818 return env.context
819 .awaitJs(env.js,
820 adapter->readAllText(env.js).then(
821 env.js, [&adapter = *adapter](jsg::Lock& js, jsg::JsRef<jsg::JsString> result) {
822 auto str = result.getHandle(js).toString(js);
823 // With exponential growth strategy: 1024 + 2048 + 4096 + 8192 = 15360
824 KJ_ASSERT(str.size() == 15360);
825 KJ_ASSERT(adapter.isClosed(), "Adapter should be closed after readAllText()");
826 })).attach(kj::mv(adapter));
827 });
828}
829 
830KJ_TEST("Read all bytes") {
831 TestFixture fixture;
832 FiniteReadSource source(4);
833 
834 fixture.runInIoContext([&](const TestFixture::Environment& env) {
835 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
836 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
837 env.js, env.context, newReadableSource(kj::mv(fake)));
838 
839 return env.context
840 .awaitJs(env.js,
841 adapter->readAllBytes(env.js).then(
842 env.js, [&adapter = *adapter](jsg::Lock& js, jsg::BufferSource result) {
843 // With exponential growth strategy: 1024 + 2048 + 4096 + 8192 = 15360
844 KJ_ASSERT(result.size() == 15360);
845 KJ_ASSERT(adapter.isClosed(), "Adapter should be closed after readAllText()");
846 })).attach(kj::mv(adapter));
847 });
848}
849 
850KJ_TEST("Read all text (limit)") {
851 TestFixture fixture;
852 FiniteReadSource source(2);
853 
854 fixture.runInIoContext([&](const TestFixture::Environment& env) {
855 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
856 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
857 env.js, env.context, newReadableSource(kj::mv(fake)));
858 
859 return env.context
860 .awaitJs(env.js,
861 adapter->readAllText(env.js, 100)
862 .then(env.js,
863 [](jsg::Lock& js, jsg::JsRef<jsg::JsString> result) -> jsg::Promise<void> {
864 KJ_FAIL_ASSERT("Should not have completed readAllText within limit");
865 }, [&adapter = *adapter](jsg::Lock& js, jsg::Value exception) {
866 return js.resolvedPromise();
867 })).attach(kj::mv(adapter));
868 });
869}
870 
871KJ_TEST("Read all bytes (limit)") {
872 TestFixture fixture;
873 FiniteReadSource source(2);
874 
875 fixture.runInIoContext([&](const TestFixture::Environment& env) {
876 kj::Own<kj::AsyncInputStream> fake(&source, kj::NullDisposer::instance);
877 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(
878 env.js, env.context, newReadableSource(kj::mv(fake)));
879 
880 return env.context
881 .awaitJs(env.js,
882 adapter->readAllBytes(env.js, 100)
883 .then(env.js, [](jsg::Lock&, auto) -> jsg::Promise<void> {
884 KJ_FAIL_ASSERT("Should not have completed readAllBytes within limit");
885 }, [&adapter = *adapter](jsg::Lock& js, jsg::Value exception) {
886 return js.resolvedPromise();
887 })).attach(kj::mv(adapter));
888 });
889}
890 
891KJ_TEST("tryGetLength") {
892 TestFixture fixture;
893 
894 fixture.runInIoContext([&](const TestFixture::Environment& env) {
895 auto source = newReadableSource(newNullInputStream());
896 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(env.js, env.context, kj::mv(source));
897 auto length = KJ_ASSERT_NONNULL(adapter->tryGetLength(StreamEncoding::IDENTITY));
898 KJ_ASSERT(length == 0, "Length of empty stream should be 0");
899 
900 adapter->shutdown(env.js);
901 
902 KJ_ASSERT(adapter->tryGetLength(StreamEncoding::IDENTITY) == kj::none,
903 "Length after shutdown should be none");
904 
905 return kj::READY_NOW;
906 });
907}
908 
909KJ_TEST("tee successful") {
910 TestFixture fixture;
911 
912 fixture.runInIoContext([&](const TestFixture::Environment& env) {
913 auto dataSource = newMemoryInputStream("hello world"_kjb);
914 auto source = newReadableSource(kj::mv(dataSource));
915 auto adapter = kj::heap<ReadableStreamSourceJsAdapter>(env.js, env.context, kj::mv(source));
916 
917 auto [branch1, branch2] = KJ_ASSERT_NONNULL(adapter->tryTee(env.js));
918 
919 KJ_ASSERT(adapter->isClosed(), "Original adapter should be closed after tee");
920 KJ_ASSERT(
921 adapter->isCanceled() == kj::none, "Original adapter should not be canceled after tee");
922 
923 KJ_ASSERT(!branch1->isClosed(), "Branch1 should not be closed after tee");
924 KJ_ASSERT(branch1->isCanceled() == kj::none, "Branch1 should not be canceled after tee");
925 
926 KJ_ASSERT(!branch2->isClosed(), "Branch2 should not be closed after tee");
927 KJ_ASSERT(branch2->isCanceled() == kj::none, "Branch2 should not be canceled after tee");
928 
929 auto backing1 = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 11);
930 auto buffer1 = jsg::BufferSource(env.js, kj::mv(backing1));
931 auto read1 = branch1->read(env.js,
932 ReadableStreamSourceJsAdapter::ReadOptions{
933 .buffer = kj::mv(buffer1),
934 });
935 auto backing2 = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 11);
936 auto buffer2 = jsg::BufferSource(env.js, kj::mv(backing2));
937 auto read2 = branch2->read(env.js,
938 ReadableStreamSourceJsAdapter::ReadOptions{
939 .buffer = kj::mv(buffer2),
940 });
941 
942 return env.context
943 .awaitJs(env.js,
944 kj::mv(read1)
945 .then(env.js, [read2 = kj::mv(read2)](jsg::Lock& js, auto result1) mutable {
946 KJ_ASSERT(!result1.done, "Stream should not be done yet");
947 KJ_ASSERT(result1.buffer.asArrayPtr().size() == 11);
948 KJ_ASSERT(result1.buffer.asArrayPtr() == "hello world"_kjb);
949 return kj::mv(read2);
950 }).then(env.js, [](jsg::Lock& js, auto result2) {
951 KJ_ASSERT(!result2.done, "Stream should not be done yet");
952 KJ_ASSERT(result2.buffer.asArrayPtr().size() == 11);
953 KJ_ASSERT(result2.buffer.asArrayPtr() == "hello world"_kjb);
954 return js.resolvedPromise();
955 })).attach(kj::mv(branch1), kj::mv(branch2));
956 });
957}
958 
959// ===========================================================================================
960 
961namespace {
962static size_t countStatic = 0;
963jsg::Ref<ReadableStream> createFiniteBytesReadableStream(
964 jsg::Lock& js, size_t chunkSize = 1024, size_t* count = nullptr) {
965 if (count == nullptr) {
966 countStatic = 0;
967 count = &countStatic;
968 }
969 return ReadableStream::constructor(js,
970 UnderlyingSource{
971 .pull =
972 [chunkSize, count](jsg::Lock& js, auto controller) {
973 auto c = kj::mv(
974 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>()));
975 auto& counter = *count;
976 if (counter++ < 10) {
977 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
978 jsg::BufferSource buffer(js, kj::mv(backing));
979 buffer.asArrayPtr().fill(96 + counter); // fill with 'a'...'j'
980 c->enqueue(js, buffer.getHandle(js));
981 }
982 if (counter == 10) {
983 c->close(js);
984 }
985 return js.resolvedPromise();
986 },
987 .expectedLength = 10 * chunkSize,
988 },
989 StreamQueuingStrategy{
990 .highWaterMark = 0,
991 });
992}
993 
994jsg::Ref<ReadableStream> createFiniteByobReadableStream(jsg::Lock& js, size_t chunkSize = 1024) {
995 return ReadableStream::constructor(js,
996 UnderlyingSource{
997 .type = kj::str("bytes"),
998 .pull =
999 [chunkSize](jsg::Lock& js, auto controller) {
1000 auto c = kj::mv(
1001 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableByteStreamController>>()));
1002 static int count = 0;
1003 if (count++ < 10) {
1004 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
1005 jsg::BufferSource buffer(js, kj::mv(backing));
1006 c->enqueue(js, kj::mv(buffer));
1007 }
1008 if (count == 10) {
1009 c->close(js);
1010 }
1011 return js.resolvedPromise();
1012 },
1013 .expectedLength = 10 * chunkSize,
1014 },
1015 kj::none);
1016}
1017 
1018jsg::Ref<ReadableStream> createErroredStream(jsg::Lock& js) {
1019 return ReadableStream::constructor(js,
1020 UnderlyingSource{.start =
1021 [](jsg::Lock& js, auto controller) {
1022 auto c = kj::mv(
1023 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>()));
1024 c->error(js, js.error("boom"));
1025 return js.resolvedPromise();
1026 }},
1027 kj::none);
1028}
1029 
1030jsg::Ref<ReadableStream> createClosedStream(jsg::Lock& js) {
1031 return ReadableStream::constructor(js,
1032 UnderlyingSource{.start =
1033 [](jsg::Lock& js, auto controller) {
1034 auto c = kj::mv(
1035 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>()));
1036 c->close(js);
1037 return js.resolvedPromise();
1038 }},
1039 kj::none);
1040}
1041 
1042struct RecordingSink final: public kj::AsyncOutputStream {
1043 kj::Vector<kj::byte> data;
1044 
1045 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
1046 data.addAll(buffer.begin(), buffer.end());
1047 co_return;
1048 }
1049 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
1050 for (auto piece: pieces) {
1051 data.addAll(piece.begin(), piece.end());
1052 }
1053 co_return;
1054 }
1055 
1056 kj::Promise<void> whenWriteDisconnected() override {
1057 return kj::NEVER_DONE;
1058 }
1059};
1060 
1061struct ErrorSink final: public kj::AsyncOutputStream {
1062 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
1063 KJ_FAIL_REQUIRE("worker_do_not_log; Write failed");
1064 co_return;
1065 }
1066 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
1067 KJ_FAIL_REQUIRE("worker_do_not_log; Write failed");
1068 co_return;
1069 }
1070 
1071 kj::Promise<void> whenWriteDisconnected() override {
1072 return kj::NEVER_DONE;
1073 }
1074};
1075} // namespace
1076 
1077KJ_TEST("KjAdapter constructor with valid normal ReadableStream") {
1078 capnp::MallocMessageBuilder message;
1079 auto flags = message.initRoot<CompatibilityFlags>();
1080 flags.setStreamsJavaScriptControllers(true);
1081 TestFixture fixture({.featureFlags = flags.asReader()});
1082 
1083 // Constructs and drops without failures
1084 
1085 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1086 auto stream = createFiniteBytesReadableStream(env.js, 16 * 1024);
1087 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1088 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1089 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1090 
1091 // The size is known because we provided expectedLength in the source.
1092 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->tryGetLength(StreamEncoding::IDENTITY)), 16 * 1024);
1093 
1094 // The encoding is always IDENTITY
1095 KJ_ASSERT(adapter->getEncoding() == StreamEncoding::IDENTITY);
1096 
1097 // Teeing is unsupported so always throws
1098 try {
1099 adapter->tee(1);
1100 } catch (...) {
1101 auto ex = kj::getCaughtExceptionAsKj();
1102 KJ_ASSERT(ex.getDescription().contains("not supported"));
1103 }
1104 
1105 return kj::READY_NOW;
1106 });
1107}
1108 
1109KJ_TEST("KjAdapter constructor with valid byob ReadableStream") {
1110 capnp::MallocMessageBuilder message;
1111 auto flags = message.initRoot<CompatibilityFlags>();
1112 flags.setStreamsJavaScriptControllers(true);
1113 TestFixture fixture({.featureFlags = flags.asReader()});
1114 
1115 // Constructs and drops without failures
1116 
1117 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1118 auto stream = createFiniteByobReadableStream(env.js, 16 * 1024);
1119 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1120 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1121 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1122 
1123 // The size is known because we provided expectedLength in the source.
1124 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->tryGetLength(StreamEncoding::IDENTITY)), 16 * 1024);
1125 
1126 // The encoding is always IDENTITY
1127 KJ_ASSERT(adapter->getEncoding() == StreamEncoding::IDENTITY);
1128 
1129 return kj::READY_NOW;
1130 });
1131}
1132 
1133KJ_TEST("KjAdapter constructor with valid ReadableStream manual cancel") {
1134 capnp::MallocMessageBuilder message;
1135 auto flags = message.initRoot<CompatibilityFlags>();
1136 flags.setStreamsJavaScriptControllers(true);
1137 TestFixture fixture({.featureFlags = flags.asReader()});
1138 
1139 // Constructs and drops without failures
1140 
1141 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1142 auto stream = createFiniteBytesReadableStream(env.js, 16 * 1024);
1143 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1144 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1145 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1146 
1147 adapter->cancel(KJ_EXCEPTION(FAILED, "Manual cancel"));
1148 
1149 KJ_ASSERT(stream->isLocked(), "Stream should remain locked after adapter cancel");
1150 
1151 KJ_ASSERT(adapter->tryGetLength(StreamEncoding::IDENTITY) == kj::none,
1152 "Length after cancel should be none");
1153 
1154 return kj::READY_NOW;
1155 });
1156}
1157 
1158KJ_TEST("KjAdapter constructor with locked/disturbed stream fails") {
1159 capnp::MallocMessageBuilder message;
1160 auto flags = message.initRoot<CompatibilityFlags>();
1161 flags.setStreamsJavaScriptControllers(true);
1162 TestFixture fixture({.featureFlags = flags.asReader()});
1163 
1164 // Constructs and drops without failures
1165 
1166 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1167 auto stream = createFiniteBytesReadableStream(env.js, 16 * 1024);
1168 auto reader = stream->getReader(env.js, kj::none);
1169 try {
1170 kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1171 KJ_FAIL_ASSERT("Should not be able to get adapter");
1172 } catch (...) {
1173 // Expected.
1174 auto ex = kj::getCaughtExceptionAsKj();
1175 KJ_ASSERT(ex.getDescription().contains("ReadableStream is locked"));
1176 }
1177 
1178 auto& r = KJ_ASSERT_NONNULL(reader.tryGet<jsg::Ref<ReadableStreamDefaultReader>>());
1179 r->read(env.js);
1180 r->releaseLock(env.js);
1181 KJ_ASSERT(stream->isDisturbed());
1182 
1183 // Disturbed streams are also fatal, even if not locked.
1184 KJ_ASSERT(stream->isDisturbed());
1185 
1186 try {
1187 kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1188 KJ_FAIL_ASSERT("Should not be able to get adapter");
1189 } catch (...) {
1190 // Expected.
1191 auto ex = kj::getCaughtExceptionAsKj();
1192 KJ_ASSERT(ex.getDescription().contains("ReadableStream is disturbed"));
1193 }
1194 
1195 return kj::READY_NOW;
1196 });
1197}
1198 
1199KJ_TEST("KjAdapter read with valid buffer and byte ranges") {
1200 capnp::MallocMessageBuilder message;
1201 auto flags = message.initRoot<CompatibilityFlags>();
1202 flags.setStreamsJavaScriptControllers(true);
1203 TestFixture fixture({.featureFlags = flags.asReader()});
1204 size_t counter = 0;
1205 
1206 // Constructs and drops without failures
1207 
1208 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1209 auto stream = createFiniteBytesReadableStream(env.js, 1024, &counter);
1210 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1211 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1212 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1213 
1214 auto buffer = kj::heapArray<kj::byte>(2049);
1215 
1216 return adapter->read(buffer, 512)
1217 .then([buffer = kj::mv(buffer), &adapter = *adapter](size_t bytesRead) mutable {
1218 KJ_ASSERT(bytesRead >= 512 && bytesRead <= buffer.size());
1219 KJ_ASSERT(bytesRead == 2048);
1220 
1221 kj::FixedArray<kj::byte, 2048> expected;
1222 expected.asPtr().first(1024).fill(97); // 'a'
1223 expected.asPtr().slice(1024).fill(98); // 'b'
1224 KJ_ASSERT(buffer.asPtr().first(bytesRead) == expected.asPtr());
1225 
1226 // Perform another read...
1227 return adapter.read(buffer, 1).then([buffer = kj::mv(buffer)](size_t bytesRead) {
1228 KJ_ASSERT(bytesRead >= 1 && bytesRead <= buffer.size());
1229 KJ_ASSERT(bytesRead == 2048);
1230 
1231 kj::FixedArray<kj::byte, 2048> expected;
1232 expected.asPtr().first(1024).fill(99); // 'c'
1233 expected.asPtr().slice(1024).fill(100); // 'd'
1234 KJ_ASSERT(buffer.asPtr().first(bytesRead) == expected.asPtr());
1235 
1236 return kj::READY_NOW;
1237 });
1238 }).attach(kj::mv(adapter));
1239 });
1240}
1241 
1242KJ_TEST("KjAdapter read with left over (source provides more than requested)") {
1243 capnp::MallocMessageBuilder message;
1244 auto flags = message.initRoot<CompatibilityFlags>();
1245 flags.setStreamsJavaScriptControllers(true);
1246 TestFixture fixture({.featureFlags = flags.asReader()});
1247 size_t counter = 0;
1248 
1249 // Constructs and drops without failures
1250 
1251 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1252 auto stream = createFiniteBytesReadableStream(env.js, 1024, &counter);
1253 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1254 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1255 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1256 
1257 auto buffer = kj::heapArray<kj::byte>(1000);
1258 
1259 return adapter->read(buffer, 1000)
1260 .then([buffer = kj::mv(buffer), &adapter = *adapter](size_t bytesRead) mutable {
1261 KJ_ASSERT(bytesRead >= 512 && bytesRead <= buffer.size());
1262 KJ_ASSERT(bytesRead == 1000);
1263 
1264 kj::FixedArray<kj::byte, 1000> expected;
1265 expected.asPtr().fill(97); // 'a'
1266 KJ_ASSERT(buffer.asPtr().first(bytesRead) == expected.asPtr());
1267 
1268 // Perform another read...
1269 return adapter.read(buffer, 1).then([buffer = kj::mv(buffer)](size_t bytesRead) {
1270 // The next read should be only for the 24 remaining bytes leftover from the first chunk.
1271 KJ_ASSERT(bytesRead >= 1 && bytesRead <= buffer.size());
1272 KJ_ASSERT(bytesRead == 24);
1273 
1274 kj::FixedArray<kj::byte, 24> expected;
1275 expected.asPtr().fill(97); // 'a'
1276 KJ_ASSERT(buffer.asPtr().first(bytesRead) == expected.asPtr());
1277 
1278 return kj::READY_NOW;
1279 });
1280 }).attach(kj::mv(adapter));
1281 });
1282}
1283 
1284KJ_TEST("KjAdapter read with clamped minBytes (minBytes=0)") {
1285 capnp::MallocMessageBuilder message;
1286 auto flags = message.initRoot<CompatibilityFlags>();
1287 flags.setStreamsJavaScriptControllers(true);
1288 TestFixture fixture({.featureFlags = flags.asReader()});
1289 size_t counter = 0;
1290 
1291 // Constructs and drops without failures
1292 
1293 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1294 auto stream = createFiniteBytesReadableStream(env.js, 5, &counter);
1295 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1296 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1297 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1298 
1299 auto buffer = kj::heapArray<kj::byte>(3);
1300 
1301 return adapter->read(buffer, 0)
1302 .then([buffer = kj::mv(buffer), &adapter = *adapter](size_t bytesRead) mutable {
1303 // Should return exactly 1 byte, since minBytes is clamped to 1.
1304 KJ_ASSERT(bytesRead >= 1);
1305 }).attach(kj::mv(adapter));
1306 });
1307}
1308 
1309KJ_TEST("KjAdapter read with clamped minBytes (minBytes > maxBytes)") {
1310 capnp::MallocMessageBuilder message;
1311 auto flags = message.initRoot<CompatibilityFlags>();
1312 flags.setStreamsJavaScriptControllers(true);
1313 TestFixture fixture({.featureFlags = flags.asReader()});
1314 size_t counter = 0;
1315 
1316 // Constructs and drops without failures
1317 
1318 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1319 auto stream = createFiniteBytesReadableStream(env.js, 5, &counter);
1320 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1321 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1322 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1323 
1324 auto buffer = kj::heapArray<kj::byte>(3);
1325 
1326 return adapter->read(buffer, 4)
1327 .then([buffer = kj::mv(buffer), &adapter = *adapter](size_t bytesRead) mutable {
1328 // Should return exactly 3 byte, since minBytes is clamped to 3.
1329 KJ_ASSERT(bytesRead == 3);
1330 }).attach(kj::mv(adapter));
1331 });
1332}
1333 
1334KJ_TEST("KjAdapter read with zero length buffer") {
1335 capnp::MallocMessageBuilder message;
1336 auto flags = message.initRoot<CompatibilityFlags>();
1337 flags.setStreamsJavaScriptControllers(true);
1338 TestFixture fixture({.featureFlags = flags.asReader()});
1339 size_t counter = 0;
1340 
1341 // Constructs and drops without failures
1342 
1343 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1344 auto stream = createFiniteBytesReadableStream(env.js, 5, &counter);
1345 KJ_ASSERT(!stream->isLocked(), "Stream should not be locked before adapter construction");
1346 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1347 KJ_ASSERT(stream->isLocked(), "Stream should be locked after adapter construction");
1348 
1349 auto buffer = kj::heapArray<kj::byte>(0);
1350 
1351 return adapter->read(buffer, 1)
1352 .then([buffer = kj::mv(buffer), &adapter = *adapter](size_t bytesRead) mutable {
1353 // Should return exactly 0 byte
1354 KJ_ASSERT(bytesRead == 0);
1355 }).attach(kj::mv(adapter));
1356 });
1357}
1358 
1359KJ_TEST("KjAdapter forbid concurrent reads") {
1360 capnp::MallocMessageBuilder message;
1361 auto flags = message.initRoot<CompatibilityFlags>();
1362 flags.setStreamsJavaScriptControllers(true);
1363 TestFixture fixture({.featureFlags = flags.asReader()});
1364 size_t counter = 0;
1365 
1366 // Constructs and drops without failures
1367 
1368 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1369 auto stream = createFiniteBytesReadableStream(env.js, 5, &counter);
1370 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1371 
1372 auto buffer = kj::heapArray<kj::byte>(2);
1373 
1374 // Concurrent reads are not allowed.
1375 auto read1 = adapter->read(buffer, 1);
1376 
1377 try {
1378 auto read2 KJ_UNUSED = adapter->read(buffer, 1);
1379 } catch (...) {
1380 auto ex = kj::getCaughtExceptionAsKj();
1381 KJ_ASSERT(ex.getDescription().contains("Cannot have multiple concurrent reads"));
1382 }
1383 
1384 return kj::READY_NOW;
1385 });
1386}
1387 
1388KJ_TEST("KjAdapter cancel in-flight reads") {
1389 capnp::MallocMessageBuilder message;
1390 auto flags = message.initRoot<CompatibilityFlags>();
1391 flags.setStreamsJavaScriptControllers(true);
1392 TestFixture fixture({.featureFlags = flags.asReader()});
1393 size_t counter = 0;
1394 
1395 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1396 auto stream = createFiniteBytesReadableStream(env.js, 5, &counter);
1397 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1398 
1399 auto buffer = kj::heapArray<kj::byte>(2);
1400 
1401 // Concurrent reads are not allowed.
1402 auto read1 = adapter->read(buffer, 1);
1403 
1404 adapter->cancel(KJ_EXCEPTION(FAILED, "worker_do_not_log; Manual cancel"));
1405 
1406 return read1
1407 .then([](size_t) { KJ_FAIL_ASSERT("Should not have completed read after cancel"); },
1408 [](kj::Exception exception) {
1409 KJ_ASSERT(exception.getDescription().contains("Manual cancel"));
1410 }).attach(kj::mv(adapter));
1411 });
1412}
1413 
1414KJ_TEST("KjAdapter read errored stream") {
1415 capnp::MallocMessageBuilder message;
1416 auto flags = message.initRoot<CompatibilityFlags>();
1417 flags.setStreamsJavaScriptControllers(true);
1418 TestFixture fixture({.featureFlags = flags.asReader()});
1419 
1420 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1421 auto stream = createErroredStream(env.js);
1422 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1423 
1424 auto buffer = kj::heapArray<kj::byte>(2);
1425 
1426 // Concurrent reads are not allowed.
1427 auto read1 = adapter->read(buffer, 1);
1428 
1429 return read1
1430 .then([](size_t) { KJ_FAIL_ASSERT("Should not have completed read after cancel"); },
1431 [&adapter = *adapter](kj::Exception exception) {
1432 KJ_ASSERT(exception.getDescription().contains("boom"));
1433 })
1434 .then([&adapter = *adapter]() {
1435 // The adapter should be in the errored state now.
1436 kj::FixedArray<kj::byte, 1> buf;
1437 return adapter.read(buf, 1).then([](auto) {
1438 KJ_FAIL_ASSERT("Should not have completed read on errored adapter");
1439 }, [](kj::Exception exception) { KJ_ASSERT(exception.getDescription().contains("boom")); });
1440 }).attach(kj::mv(adapter));
1441 });
1442}
1443 
1444KJ_TEST("KjAdapter read closed stream") {
1445 capnp::MallocMessageBuilder message;
1446 auto flags = message.initRoot<CompatibilityFlags>();
1447 flags.setStreamsJavaScriptControllers(true);
1448 TestFixture fixture({.featureFlags = flags.asReader()});
1449 
1450 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1451 auto stream = createClosedStream(env.js);
1452 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1453 
1454 auto buffer = kj::heapArray<kj::byte>(2);
1455 
1456 auto read1 = adapter->read(buffer, 1);
1457 
1458 return read1.then([](size_t size) { KJ_ASSERT(size == 0); }).attach(kj::mv(adapter));
1459 });
1460}
1461 
1462KJ_TEST("KjAdapter pumpTo") {
1463 capnp::MallocMessageBuilder message;
1464 auto flags = message.initRoot<CompatibilityFlags>();
1465 flags.setStreamsJavaScriptControllers(true);
1466 TestFixture fixture({.featureFlags = flags.asReader()});
1467 RecordingSink sink;
1468 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
1469 auto writableSink = newWritableSink(kj::mv(fakeOwn));
1470 
1471 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1472 // The assertion in pumpToImpl requires totalRead <= buffer.size() after reading,
1473 // which means the stream size must be <= bufferSize / 2. With bufferSize = 1024
1474 // (for streams < 4096 bytes), we need total size <= 512 bytes.
1475 auto stream = createFiniteBytesReadableStream(env.js, 50);
1476 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1477 
1478 return adapter->pumpTo(*writableSink, EndAfterPump::YES).attach(kj::mv(adapter));
1479 });
1480 
1481 kj::FixedArray<kj::byte, 10 * 50> expected;
1482 expected.asPtr().first(50).fill(97); // 'a'
1483 expected.asPtr().slice(50, 100).fill(98); // 'b'
1484 expected.asPtr().slice(100, 150).fill(99); // 'c'
1485 expected.asPtr().slice(150, 200).fill(100); // 'd'
1486 expected.asPtr().slice(200, 250).fill(101); // 'e'
1487 expected.asPtr().slice(250, 300).fill(102); // 'f'
1488 expected.asPtr().slice(300, 350).fill(103); // 'g'
1489 expected.asPtr().slice(350, 400).fill(104); // 'h'
1490 expected.asPtr().slice(400, 450).fill(105); // 'i'
1491 expected.asPtr().slice(450, 500).fill(106); // 'j'
1492 
1493 KJ_ASSERT(sink.data.size() == 10 * 50);
1494 KJ_ASSERT(sink.data.asPtr() == expected.asPtr());
1495}
1496 
1497KJ_TEST("KjAdapter pumpTo (no end)") {
1498 capnp::MallocMessageBuilder message;
1499 auto flags = message.initRoot<CompatibilityFlags>();
1500 flags.setStreamsJavaScriptControllers(true);
1501 TestFixture fixture({.featureFlags = flags.asReader()});
1502 RecordingSink sink;
1503 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
1504 auto writableSink = newWritableSink(kj::mv(fakeOwn));
1505 
1506 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1507 // Use the same size constraint as the previous test.
1508 auto stream = createFiniteBytesReadableStream(env.js, 50);
1509 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1510 
1511 return adapter->pumpTo(*writableSink, EndAfterPump::NO).attach(kj::mv(adapter));
1512 });
1513 
1514 kj::FixedArray<kj::byte, 10 * 50> expected;
1515 expected.asPtr().first(50).fill(97); // 'a'
1516 expected.asPtr().slice(50, 100).fill(98); // 'b'
1517 expected.asPtr().slice(100, 150).fill(99); // 'c'
1518 expected.asPtr().slice(150, 200).fill(100); // 'd'
1519 expected.asPtr().slice(200, 250).fill(101); // 'e'
1520 expected.asPtr().slice(250, 300).fill(102); // 'f'
1521 expected.asPtr().slice(300, 350).fill(103); // 'g'
1522 expected.asPtr().slice(350, 400).fill(104); // 'h'
1523 expected.asPtr().slice(400, 450).fill(105); // 'i'
1524 expected.asPtr().slice(450, 500).fill(106); // 'j'
1525 
1526 KJ_ASSERT(sink.data.size() == 10 * 50);
1527 KJ_ASSERT(sink.data.asPtr() == expected.asPtr());
1528}
1529 
1530KJ_TEST("KjAdapter pumpTo (errored)") {
1531 capnp::MallocMessageBuilder message;
1532 auto flags = message.initRoot<CompatibilityFlags>();
1533 flags.setStreamsJavaScriptControllers(true);
1534 TestFixture fixture({.featureFlags = flags.asReader()});
1535 RecordingSink sink;
1536 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
1537 auto writableSink = newWritableSink(kj::mv(fakeOwn));
1538 
1539 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1540 auto stream = createErroredStream(env.js);
1541 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1542 
1543 return env.context.waitForDeferredProxy(adapter->pumpTo(*writableSink, EndAfterPump::NO))
1544 .then([]() -> kj::Promise<void> {
1545 KJ_FAIL_ASSERT("Should not have completed pumpTo on errored stream");
1546 }, [](kj::Exception exception) {
1547 }).attach(kj::mv(adapter));
1548 });
1549}
1550 
1551KJ_TEST("KjAdapter pumpTo (error sink)") {
1552 capnp::MallocMessageBuilder message;
1553 auto flags = message.initRoot<CompatibilityFlags>();
1554 flags.setStreamsJavaScriptControllers(true);
1555 TestFixture fixture({.featureFlags = flags.asReader()});
1556 ErrorSink sink;
1557 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
1558 auto writableSink = newWritableSink(kj::mv(fakeOwn));
1559 
1560 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1561 auto stream = createFiniteBytesReadableStream(env.js, 1000);
1562 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1563 
1564 return env.context.waitForDeferredProxy(adapter->pumpTo(*writableSink, EndAfterPump::NO))
1565 .then([]() -> kj::Promise<void> {
1566 KJ_FAIL_ASSERT("Should not have completed pumpTo on errored stream");
1567 }, [](kj::Exception exception) {
1568 KJ_ASSERT(exception.getDescription().contains("Write failed"));
1569 }).attach(kj::mv(adapter));
1570 });
1571}
1572 
1573KJ_TEST("KjAdapter MinReadPolicy IMMEDIATE behavior") {
1574 capnp::MallocMessageBuilder message;
1575 auto flags = message.initRoot<CompatibilityFlags>();
1576 flags.setStreamsJavaScriptControllers(true);
1577 TestFixture fixture({.featureFlags = flags.asReader()});
1578 size_t counter = 0;
1579 
1580 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1581 // Create a stream that returns data in small chunks to test the policy difference
1582 auto stream = ReadableStream::constructor(env.js,
1583 UnderlyingSource{
1584 .pull =
1585 [&counter](jsg::Lock& js, auto controller) {
1586 auto& c = KJ_ASSERT_NONNULL(
1587 controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
1588 if (counter < 8) {
1589 // Return 256 bytes per chunk, 8 chunks total (2048 bytes)
1590 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 256);
1591 jsg::BufferSource buffer(js, kj::mv(backing));
1592 buffer.asArrayPtr().fill(97 + counter); // 'a', 'b', 'c', etc.
1593 c->enqueue(js, buffer.getHandle(js));
1594 counter++;
1595 } else {
1596 c->close(js);
1597 }
1598 return js.resolvedPromise();
1599 },
1600 .expectedLength = 2048,
1601 },
1602 StreamQueuingStrategy{.highWaterMark = 0});
1603 
1604 // Test IMMEDIATE policy - should return as soon as minBytes is satisfied
1605 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef(),
1606 ReadableSourceKjAdapter::Options{
1607 .minReadPolicy = ReadableSourceKjAdapter::MinReadPolicy::IMMEDIATE});
1608 
1609 auto buffer = kj::heapArray<kj::byte>(2048);
1610 
1611 return adapter->read(buffer, 512)
1612 .then([buffer = kj::mv(buffer)](size_t bytesRead) {
1613 // With IMMEDIATE policy, should return as soon as minBytes (512) is satisfied
1614 KJ_ASSERT(bytesRead == 512, "Should have read exactly minBytes");
1615 
1616 // Verify the data content matches expected pattern
1617 for (size_t i = 0; i < bytesRead; i++) {
1618 size_t chunkIndex = i / 256;
1619 KJ_ASSERT(buffer[i] == static_cast<kj::byte>(97 + chunkIndex),
1620 "Data should match expected pattern");
1621 }
1622 
1623 return kj::READY_NOW;
1624 }).attach(kj::mv(adapter));
1625 });
1626}
1627 
1628KJ_TEST("KjAdapter MinReadPolicy OPPORTUNISTIC behavior") {
1629 capnp::MallocMessageBuilder message;
1630 auto flags = message.initRoot<CompatibilityFlags>();
1631 flags.setStreamsJavaScriptControllers(true);
1632 TestFixture fixture({.featureFlags = flags.asReader()});
1633 size_t counter = 0;
1634 
1635 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1636 // Create a stream that returns data in small chunks to test the policy difference
1637 auto stream = ReadableStream::constructor(env.js,
1638 UnderlyingSource{
1639 .pull =
1640 [&counter](jsg::Lock& js, auto controller) {
1641 auto c = kj::mv(KJ_ASSERT_NONNULL(
1642 controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>()));
1643 
1644 if (counter < 8) {
1645 // Return 256 bytes per chunk, 8 chunks total (2048 bytes)
1646 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 256);
1647 jsg::BufferSource buffer(js, kj::mv(backing));
1648 buffer.asArrayPtr().fill(97 + counter); // 'a', 'b', 'c', etc.
1649 c->enqueue(js, buffer.getHandle(js));
1650 counter++;
1651 } else {
1652 c->close(js);
1653 }
1654 return js.resolvedPromise();
1655 },
1656 .expectedLength = 2048,
1657 },
1658 StreamQueuingStrategy{.highWaterMark = 0});
1659 
1660 // Test OPPORTUNISTIC policy - should try to fill buffer more completely
1661 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef(),
1662 ReadableSourceKjAdapter::Options{
1663 .minReadPolicy = ReadableSourceKjAdapter::MinReadPolicy::OPPORTUNISTIC});
1664 
1665 auto buffer = kj::heapArray<kj::byte>(2048);
1666 
1667 return adapter->read(buffer, 512)
1668 .then([buffer = kj::mv(buffer)](size_t bytesRead) {
1669 // With OPPORTUNISTIC policy, should try to fill buffer more completely
1670 // when data is readily available
1671 KJ_ASSERT(bytesRead == 1792, "Should have read as much as possible up to maxBytes");
1672 
1673 // Verify the data content matches expected pattern
1674 for (size_t i = 0; i < bytesRead; i++) {
1675 size_t chunkIndex = i / 256;
1676 KJ_ASSERT(buffer[i] == static_cast<kj::byte>(97 + chunkIndex),
1677 "Data should match expected pattern");
1678 }
1679 
1680 return kj::READY_NOW;
1681 }).attach(kj::mv(adapter));
1682 });
1683}
1684 
1685KJ_TEST("KjAdapter readAllBytes") {
1686 capnp::MallocMessageBuilder message;
1687 auto flags = message.initRoot<CompatibilityFlags>();
1688 flags.setStreamsJavaScriptControllers(true);
1689 TestFixture fixture({.featureFlags = flags.asReader()});
1690 
1691 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
1692 auto stream = createFiniteBytesReadableStream(env.js, 1024);
1693 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1694 auto bytes = co_await adapter->readAllBytes(kj::maxValue).attach(kj::mv(adapter));
1695 kj::FixedArray<kj::byte, 10 * 1024> expected;
1696 expected.asPtr().first(1024).fill(97); // 'a'
1697 expected.asPtr().slice(1024, 2048).fill(98); // 'b'
1698 expected.asPtr().slice(2048, 3072).fill(99); // 'c'
1699 expected.asPtr().slice(3072, 4096).fill(100); // 'd'
1700 expected.asPtr().slice(4096, 5120).fill(101); // 'e'
1701 expected.asPtr().slice(5120, 6144).fill(102); // 'f'
1702 expected.asPtr().slice(6144, 7168).fill(103); // 'g'
1703 expected.asPtr().slice(7168, 8192).fill(104); // 'h'
1704 expected.asPtr().slice(8192, 9216).fill(105); // 'i'
1705 expected.asPtr().slice(9216, 10240).fill(106); // 'j'
1706 
1707 KJ_ASSERT(bytes.size() == 10 * 1024);
1708 KJ_ASSERT(bytes == expected);
1709 });
1710}
1711 
1712KJ_TEST("KjAdapter readAllBytes (limit exceeded)") {
1713 capnp::MallocMessageBuilder message;
1714 auto flags = message.initRoot<CompatibilityFlags>();
1715 flags.setStreamsJavaScriptControllers(true);
1716 TestFixture fixture({.featureFlags = flags.asReader()});
1717 
1718 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
1719 auto stream = createFiniteBytesReadableStream(env.js, 1024);
1720 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1721 try {
1722 co_await adapter->readAllBytes(100).attach(kj::mv(adapter));
1723 KJ_FAIL_ASSERT("should have failed");
1724 } catch (...) {
1725 auto ex = kj::getCaughtExceptionAsKj();
1726 KJ_ASSERT(ex.getDescription().contains("would be exceeded"));
1727 }
1728 });
1729}
1730 
1731KJ_TEST("KjAdapter readAllText") {
1732 capnp::MallocMessageBuilder message;
1733 auto flags = message.initRoot<CompatibilityFlags>();
1734 flags.setStreamsJavaScriptControllers(true);
1735 TestFixture fixture({.featureFlags = flags.asReader()});
1736 
1737 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
1738 auto stream = createFiniteBytesReadableStream(env.js, 2048);
1739 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1740 
1741 auto text = co_await adapter->readAllText(kj::maxValue).attach(kj::mv(adapter));
1742 kj::FixedArray<char, 10 * 2048> expected;
1743 expected.asPtr().first(2048).fill(97); // 'a'
1744 expected.asPtr().slice(2048, 4096).fill(98); // 'b'
1745 expected.asPtr().slice(4096, 6144).fill(99); // 'c'
1746 expected.asPtr().slice(6144, 8192).fill(100); // 'd'
1747 expected.asPtr().slice(8192, 10240).fill(101); // 'e'
1748 expected.asPtr().slice(10240, 12288).fill(102); // 'f'
1749 expected.asPtr().slice(12288, 14336).fill(103); // 'g'
1750 expected.asPtr().slice(14336, 16384).fill(104); // 'h'
1751 expected.asPtr().slice(16384, 18432).fill(105); // 'i'
1752 expected.asPtr().slice(18432, 20480).fill(106); // 'j'
1753 
1754 KJ_ASSERT(text.size() == 10 * 2048);
1755 KJ_ASSERT(text == expected.asPtr());
1756 });
1757}
1758 
1759KJ_TEST("KjAdapter readAllText (limit exceeded)") {
1760 capnp::MallocMessageBuilder message;
1761 auto flags = message.initRoot<CompatibilityFlags>();
1762 flags.setStreamsJavaScriptControllers(true);
1763 TestFixture fixture({.featureFlags = flags.asReader()});
1764 
1765 fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> {
1766 auto stream = createFiniteBytesReadableStream(env.js, 1024);
1767 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
1768 try {
1769 co_await adapter->readAllText(100).attach(kj::mv(adapter));
1770 KJ_FAIL_ASSERT("should have failed");
1771 } catch (...) {
1772 auto ex = kj::getCaughtExceptionAsKj();
1773 KJ_ASSERT(ex.getDescription().contains("would be exceeded"));
1774 }
1775 });
1776}
1777 
1778} // namespace workerd::api::streams