Skip to content
File

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

69.6 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "queue.h"
6 
7#include <workerd/jsg/jsg-test.h>
8#include <workerd/jsg/jsg.h>
9#include <workerd/tests/test-fixture.h>
10 
11namespace workerd::api {
12namespace {
13 
14void preamble(auto callback) {
15 TestFixture fixture;
16 fixture.runInIoContext([&](const TestFixture::Environment& env) { callback(env.js); });
17}
18 
19using ReadContinuation = jsg::Promise<ReadResult>(ReadResult&&);
20using CloseContinuation = jsg::Promise<void>(ReadResult&&);
21using ReadErrorContinuation = jsg::Promise<ReadResult>(jsg::Value&&);
22 
23const kj::MutexGuarded<kj::UnwindDetector> unwindDetectorMutex;
24 
25template <typename Signature>
26struct MustCall;
27// Used to create a jsg::Promise continuation function that must be called
28// at least once during the test. If the function is not called, an error
29// will be thrown causing the test to fail.
30// TODO(cleanup): Consider adding this to jsg-test.h
31 
32template <typename Signature>
33struct MustNotCall;
34// Used to create a jsg::Promise continuation function that must not be called
35// during the test. If the function is called, an error will be thrown causing
36// the test to fail.
37// TODO(cleanup): Consider adding this to jsg-test.h
38 
39template <typename Ret, typename... Args>
40struct MustCall<Ret(Args...)> {
41 using Func = jsg::Function<Ret(Args...)>;
42 Func fn;
43 uint expected;
44 kj::SourceLocation location;
45 uint called = false;
46 
47 MustCall(Func fn, uint expected = 1, kj::SourceLocation location = kj::SourceLocation())
48 : fn(kj::mv(fn)),
49 expected(expected),
50 location(location) {}
51 
52 ~MustCall() {
53 auto unwindDetector = unwindDetectorMutex.lockExclusive();
54 if (!unwindDetector->isUnwinding()) {
55 KJ_ASSERT(called == expected,
56 kj::str("MustCall function was not called ", expected, " times. [actual: ", called, "]"),
57 location);
58 }
59 }
60 
61 Ret operator()(jsg::Lock& js, Args&&... args) {
62 called++;
63 return fn(js, kj::fwd<Args...>(args...));
64 }
65};
66 
67template <typename Ret, typename... Args>
68struct MustNotCall<Ret(Args...)> {
69 MustNotCall(kj::SourceLocation location = kj::SourceLocation()): location(location) {}
70 kj::SourceLocation location;
71 Ret operator()(jsg::Lock&, Args... args) {
72 KJ_FAIL_REQUIRE("MustNotCall function was called!", location);
73 }
74};
75 
76auto read(jsg::Lock& js, auto& consumer) {
77 auto prp = js.newPromiseAndResolver<ReadResult>();
78 consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)});
79 return kj::mv(prp.promise);
80}
81 
82auto byobRead(jsg::Lock& js, auto& consumer, int size) {
83 auto prp = js.newPromiseAndResolver<ReadResult>();
84 consumer.read(js,
85 ByteQueue::ReadRequest(kj::mv(prp.resolver),
86 {
87 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, size)),
88 .type = ByteQueue::ReadRequest::Type::BYOB,
89 }));
90 return kj::mv(prp.promise);
91};
92 
93auto getEntry(jsg::Lock& js, auto size) {
94 return kj::rc<ValueQueue::Entry>(js.v8Ref(v8::True(js.v8Isolate).As<v8::Value>()), size);
95}
96 
97#pragma region ValueQueue Tests
98 
99KJ_TEST("ValueQueue basics work") {
100 preamble([](jsg::Lock& js) {
101 ValueQueue queue(2);
102 
103 // At this point, there are no consumers, data does not get enqueued.
104 KJ_ASSERT(queue.desiredSize() == 2);
105 KJ_ASSERT(queue.size() == 0);
106 
107 queue.push(js, getEntry(js, 1));
108 
109 // Because there are no consumers, there is no change to backpressure.
110 KJ_ASSERT(queue.desiredSize() == 2);
111 KJ_ASSERT(queue.size() == 0);
112 
113 // Closing the queue causes the desiredSize to be zero.
114 queue.close(js);
115 
116 try {
117 queue.push(js, getEntry(js, 1));
118 KJ_FAIL_ASSERT("The queue push after close should have failed.");
119 } catch (kj::Exception& ex) {
120 KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored."));
121 }
122 
123 KJ_ASSERT(queue.desiredSize() == 0);
124 KJ_ASSERT(queue.size() == 0);
125 });
126}
127 
128KJ_TEST("ValueQueue erroring works") {
129 preamble([](jsg::Lock& js) {
130 ValueQueue queue(2);
131 
132 queue.error(js, js.v8Ref(js.v8Error("boom"_kj)));
133 
134 KJ_ASSERT(queue.desiredSize() == 0);
135 
136 try {
137 queue.push(js, getEntry(js, 1));
138 KJ_FAIL_ASSERT("The queue push after close should have failed.");
139 } catch (kj::Exception& ex) {
140 KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored."));
141 }
142 });
143}
144 
145KJ_TEST("ValueQueue with single consumer") {
146 preamble([](jsg::Lock& js) {
147 ValueQueue queue(2);
148 
149 ValueQueue::Consumer consumer(queue);
150 
151 KJ_ASSERT(queue.desiredSize() == 2);
152 
153 queue.push(js, getEntry(js, 2));
154 
155 // The item was pushed into the consumer.
156 KJ_ASSERT(consumer.size() == 2);
157 
158 // The queue size and desiredSize were updated accordingly.
159 KJ_ASSERT(queue.size() == 2);
160 KJ_ASSERT(queue.desiredSize() == 0);
161 
162 auto prp = js.newPromiseAndResolver<ReadResult>();
163 consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)});
164 
165 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
166 KJ_ASSERT(!result.done);
167 auto& value = KJ_ASSERT_NONNULL(result.value);
168 KJ_ASSERT(value.getHandle(js)->IsTrue());
169 
170 KJ_ASSERT(consumer.size() == 0);
171 KJ_ASSERT(queue.size() == 0);
172 KJ_ASSERT(queue.desiredSize() == 2);
173 
174 return js.resolvedPromise(kj::mv(result));
175 });
176 
177 prp.promise.then(js, readContinuation);
178 
179 js.runMicrotasks();
180 });
181}
182 
183KJ_TEST("ValueQueue with multiple consumers") {
184 preamble([](jsg::Lock& js) {
185 ValueQueue queue(2);
186 
187 ValueQueue::Consumer consumer1(queue);
188 ValueQueue::Consumer consumer2(queue);
189 
190 KJ_ASSERT(queue.desiredSize() == 2);
191 
192 queue.push(js, getEntry(js, 2));
193 
194 // The item was pushed into the consumer.
195 KJ_ASSERT(consumer1.size() == 2);
196 KJ_ASSERT(consumer2.size() == 2);
197 
198 // The queue size and desiredSize were updated accordingly.
199 KJ_ASSERT(queue.size() == 2);
200 KJ_ASSERT(queue.desiredSize() == 0);
201 
202 MustCall<ReadContinuation> read1Continuation([&](jsg::Lock& js, auto&& result) -> auto {
203 KJ_ASSERT(!result.done);
204 auto& value = KJ_ASSERT_NONNULL(result.value);
205 KJ_ASSERT(value.getHandle(js)->IsTrue());
206 
207 KJ_ASSERT(consumer1.size() == 0);
208 KJ_ASSERT(consumer2.size() == 2);
209 
210 // Backpressure was not relieved since the other consumer has yet to read.
211 KJ_ASSERT(queue.size() == 2);
212 KJ_ASSERT(queue.desiredSize() == 0);
213 
214 return read(js, consumer2);
215 });
216 
217 MustCall<ReadContinuation> read2Continuation([&](jsg::Lock& js, auto&& result) -> auto {
218 KJ_ASSERT(!result.done);
219 auto& value = KJ_ASSERT_NONNULL(result.value);
220 KJ_ASSERT(value.getHandle(js)->IsTrue());
221 
222 KJ_ASSERT(consumer2.size() == 0);
223 
224 // Backpressure was relieved since both consumers have now read.
225 KJ_ASSERT(queue.size() == 0);
226 KJ_ASSERT(queue.desiredSize() == 2);
227 
228 return js.resolvedPromise(kj::mv(result));
229 });
230 
231 MustCall<ReadContinuation> close1Continuation([&](jsg::Lock& js, auto&& result) {
232 KJ_ASSERT(result.done);
233 return read(js, consumer2);
234 });
235 
236 MustCall<CloseContinuation> close2Continuation([&](jsg::Lock& js, auto&& result) {
237 KJ_ASSERT(result.done);
238 return js.resolvedPromise();
239 });
240 
241 read(js, consumer1).then(js, read1Continuation).then(js, read2Continuation);
242 
243 js.runMicrotasks();
244 
245 // Closing the queue causes both consumers to be closed...
246 queue.close(js);
247 
248 // After close, the consumers will still be usable, but the queue itself
249 // has shutdown and no longer reports backpressure.
250 KJ_ASSERT(queue.desiredSize() == 0);
251 KJ_ASSERT(queue.size() == 0);
252 
253 read(js, consumer1).then(js, close1Continuation).then(js, close2Continuation);
254 
255 js.runMicrotasks();
256 });
257}
258 
259KJ_TEST("ValueQueue consumer with multiple-reads") {
260 preamble([](jsg::Lock& js) {
261 ValueQueue queue(2);
262 ValueQueue::Consumer consumer(queue);
263 
264 // The first read will produce a value.
265 MustCall<ReadContinuation> read1Continuation([&](jsg::Lock& js, auto&& result) -> auto {
266 KJ_ASSERT(!result.done);
267 auto& value = KJ_ASSERT_NONNULL(result.value);
268 KJ_ASSERT(value.getHandle(js)->IsTrue());
269 return js.resolvedPromise(kj::mv(result));
270 });
271 read(js, consumer).then(js, read1Continuation);
272 
273 // The second and third reads will both be done = true
274 MustCall<CloseContinuation> closeContinuation([&](jsg::Lock& js, auto&& result) {
275 KJ_ASSERT(result.done);
276 return js.resolvedPromise();
277 }, 2);
278 
279 read(js, consumer).then(js, closeContinuation);
280 read(js, consumer).then(js, closeContinuation);
281 
282 queue.push(js, getEntry(js, 2));
283 
284 // Because there is a consumer reading when the push happens, no backpressure
285 // is applied...
286 KJ_ASSERT(queue.desiredSize() == 2);
287 KJ_ASSERT(queue.size() == 0);
288 
289 queue.close(js);
290 
291 js.runMicrotasks();
292 });
293}
294 
295KJ_TEST("ValueQueue errors consumer with multiple-reads") {
296 preamble([](jsg::Lock& js) {
297 ValueQueue queue(2);
298 ValueQueue::Consumer consumer(queue);
299 
300 MustCall<ReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
301 KJ_ASSERT(value.getHandle(js)->IsNativeError());
302 return js.rejectedPromise<ReadResult>(kj::mv(value));
303 }, 3);
304 MustNotCall<ReadContinuation> readContinuation;
305 
306 read(js, consumer).then(js, readContinuation, errorContinuation);
307 read(js, consumer).then(js, readContinuation, errorContinuation);
308 read(js, consumer).then(js, readContinuation, errorContinuation);
309 
310 queue.error(js, js.v8Ref(js.v8Error("boom"_kj)));
311 
312 js.runMicrotasks();
313 });
314}
315 
316KJ_TEST("ValueQueue with multiple consumers with pending reads") {
317 preamble([](jsg::Lock& js) {
318 ValueQueue queue(2);
319 
320 ValueQueue::Consumer consumer1(queue);
321 ValueQueue::Consumer consumer2(queue);
322 
323 KJ_ASSERT(queue.desiredSize() == 2);
324 
325 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
326 KJ_ASSERT(!result.done);
327 auto& value = KJ_ASSERT_NONNULL(result.value);
328 KJ_ASSERT(value.getHandle(js)->IsTrue());
329 
330 // Both reads were fulfilled immediately without buffering.
331 KJ_ASSERT(consumer1.size() == 0);
332 KJ_ASSERT(consumer2.size() == 0);
333 
334 // Backpressure is not signalled since both consumer reads have been
335 // fulfilled.
336 KJ_ASSERT(queue.size() == 0);
337 KJ_ASSERT(queue.desiredSize() == 2);
338 
339 return js.resolvedPromise(kj::mv(result));
340 }, 2);
341 
342 read(js, consumer1).then(js, readContinuation);
343 read(js, consumer2).then(js, readContinuation);
344 
345 queue.push(js, getEntry(js, 2));
346 
347 js.runMicrotasks();
348 });
349}
350 
351#pragma endregion ValueQueue Tests
352 
353#pragma region ByteQueue Tests
354 
355KJ_TEST("ByteQueue basics work") {
356 preamble([](jsg::Lock& js) {
357 ByteQueue queue(2);
358 
359 // At this point, there are no consumers, data does not get enqueued.
360 KJ_ASSERT(queue.desiredSize() == 2);
361 KJ_ASSERT(queue.size() == 0);
362 
363 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)));
364 
365 queue.push(js, kj::mv(entry));
366 
367 // Because there are no consumers, there is no change to backpressure.
368 KJ_ASSERT(queue.desiredSize() == 2);
369 KJ_ASSERT(queue.size() == 0);
370 
371 // Closing the queue causes the desiredSize to be zero.
372 queue.close(js);
373 
374 try {
375 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)));
376 queue.push(js, kj::mv(entry));
377 KJ_FAIL_ASSERT("The queue push after close should have failed.");
378 } catch (kj::Exception& ex) {
379 KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored."));
380 }
381 
382 KJ_ASSERT(queue.desiredSize() == 0);
383 KJ_ASSERT(queue.size() == 0);
384 });
385}
386 
387KJ_TEST("ByteQueue erroring works") {
388 preamble([](jsg::Lock& js) {
389 ByteQueue queue(2);
390 
391 queue.error(js, js.v8Ref(js.v8Error("boom"_kj)));
392 
393 KJ_ASSERT(queue.desiredSize() == 0);
394 
395 try {
396 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)));
397 queue.push(js, kj::mv(entry));
398 KJ_FAIL_ASSERT("The queue push after close should have failed.");
399 } catch (kj::Exception& ex) {
400 KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored."));
401 }
402 });
403}
404 
405KJ_TEST("ByteQueue with single consumer") {
406 preamble([](jsg::Lock& js) {
407 ByteQueue queue(2);
408 
409 ByteQueue::Consumer consumer(queue);
410 
411 KJ_ASSERT(queue.desiredSize() == 2);
412 
413 auto store = jsg::BackingStore::alloc(js, 4);
414 store.asArrayPtr().fill('a');
415 
416 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store)));
417 queue.push(js, kj::mv(entry));
418 
419 // The item was pushed into the consumer.
420 KJ_ASSERT(consumer.size() == 4);
421 
422 // The queue size and desiredSize were updated accordingly.
423 KJ_ASSERT(queue.size() == 4);
424 KJ_ASSERT(queue.desiredSize() == -2);
425 
426 auto prp = js.newPromiseAndResolver<ReadResult>();
427 consumer.read(js,
428 ByteQueue::ReadRequest(kj::mv(prp.resolver),
429 {
430 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
431 }));
432 
433 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
434 KJ_ASSERT(!result.done);
435 auto& value = KJ_ASSERT_NONNULL(result.value);
436 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
437 jsg::BufferSource source(js, value.getHandle(js));
438 KJ_ASSERT(source.size() == 4);
439 KJ_ASSERT(source.asArrayPtr()[0] == 'a');
440 KJ_ASSERT(source.asArrayPtr()[1] == 'a');
441 KJ_ASSERT(source.asArrayPtr()[2] == 'a');
442 KJ_ASSERT(source.asArrayPtr()[3] == 'a');
443 
444 KJ_ASSERT(consumer.size() == 0);
445 KJ_ASSERT(queue.size() == 0);
446 KJ_ASSERT(queue.desiredSize() == 2);
447 
448 return js.resolvedPromise(kj::mv(result));
449 });
450 
451 prp.promise.then(js, readContinuation);
452 
453 js.runMicrotasks();
454 });
455}
456 
457KJ_TEST("ByteQueue with single byob consumer") {
458 preamble([](jsg::Lock& js) {
459 ByteQueue queue(2);
460 
461 ByteQueue::Consumer consumer(queue);
462 
463 auto prp = js.newPromiseAndResolver<ReadResult>();
464 consumer.read(js,
465 ByteQueue::ReadRequest(kj::mv(prp.resolver),
466 {
467 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
468 .type = ByteQueue::ReadRequest::Type::BYOB,
469 }));
470 
471 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
472 KJ_ASSERT(!result.done);
473 auto& value = KJ_ASSERT_NONNULL(result.value);
474 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
475 jsg::BufferSource source(js, value.getHandle(js));
476 auto ptr = source.asArrayPtr();
477 KJ_ASSERT(source.size() == 3);
478 KJ_ASSERT(ptr[0] == 'b');
479 KJ_ASSERT(ptr[1] == 'b');
480 KJ_ASSERT(ptr[2] == 'b');
481 
482 KJ_ASSERT(consumer.size() == 0);
483 KJ_ASSERT(queue.size() == 0);
484 KJ_ASSERT(queue.desiredSize() == 2);
485 
486 return js.resolvedPromise(kj::mv(result));
487 });
488 
489 prp.promise.then(js, readContinuation);
490 
491 auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
492 
493 KJ_ASSERT(!pendingByob->isInvalidated());
494 
495 auto& req = pendingByob->getRequest();
496 auto ptr = req.pullInto.store.asArrayPtr();
497 ptr.first(3).fill('b');
498 pendingByob->respond(js, 3);
499 KJ_ASSERT(pendingByob->isInvalidated());
500 
501 // No backpressure is signaled.
502 KJ_ASSERT(queue.desiredSize() == 2);
503 KJ_ASSERT(queue.size() == 0);
504 KJ_ASSERT(consumer.size() == 0);
505 
506 js.runMicrotasks();
507 });
508}
509 
510KJ_TEST("ByteQueue with byob consumer and default consumer") {
511 preamble([](jsg::Lock& js) {
512 ByteQueue queue(2);
513 
514 ByteQueue::Consumer consumer1(queue);
515 ByteQueue::Consumer consumer2(queue);
516 
517 auto prp = js.newPromiseAndResolver<ReadResult>();
518 consumer1.read(js,
519 ByteQueue::ReadRequest(kj::mv(prp.resolver),
520 {
521 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
522 .type = ByteQueue::ReadRequest::Type::BYOB,
523 }));
524 
525 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
526 KJ_ASSERT(!result.done);
527 auto& value = KJ_ASSERT_NONNULL(result.value);
528 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
529 jsg::BufferSource source(js, value.getHandle(js));
530 auto ptr = source.asArrayPtr();
531 KJ_ASSERT(source.size() == 3);
532 KJ_ASSERT(ptr[0] == 'b');
533 KJ_ASSERT(ptr[1] == 'b');
534 KJ_ASSERT(ptr[2] == 'b');
535 
536 KJ_ASSERT(consumer1.size() == 0);
537 KJ_ASSERT(consumer2.size() == 3);
538 KJ_ASSERT(queue.size() == 3);
539 KJ_ASSERT(queue.desiredSize() == -1);
540 
541 return js.resolvedPromise(kj::mv(result));
542 });
543 
544 prp.promise.then(js, readContinuation);
545 
546 auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
547 
548 KJ_ASSERT(!pendingByob->isInvalidated());
549 
550 auto& req = pendingByob->getRequest();
551 auto ptr = req.pullInto.store.asArrayPtr();
552 ptr.first(3).fill('b');
553 pendingByob->respond(js, 3);
554 KJ_ASSERT(pendingByob->isInvalidated());
555 
556 // Backpressure is signaled because the other consumer hasn't been read from.
557 KJ_ASSERT(queue.desiredSize() == -1);
558 KJ_ASSERT(queue.size() == 3);
559 KJ_ASSERT(consumer1.size() == 0);
560 KJ_ASSERT(consumer2.size() == 3);
561 
562 js.runMicrotasks();
563 
564 MustCall<ReadContinuation> read2Continuation([&](jsg::Lock& js, auto&& result) -> auto {
565 KJ_ASSERT(!result.done);
566 auto& value = KJ_ASSERT_NONNULL(result.value);
567 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
568 jsg::BufferSource source(js, value.getHandle(js));
569 auto ptr = source.asArrayPtr();
570 // The second consumer receives exactly the same data.
571 KJ_ASSERT(source.size() == 3);
572 KJ_ASSERT(ptr[0] == 'b');
573 KJ_ASSERT(ptr[1] == 'b');
574 KJ_ASSERT(ptr[2] == 'b');
575 
576 // The backpressure in the queue has been resolved.
577 KJ_ASSERT(queue.size() == 0);
578 KJ_ASSERT(queue.desiredSize() == 2);
579 
580 return js.resolvedPromise(kj::mv(result));
581 });
582 
583 auto prp2 = js.newPromiseAndResolver<ReadResult>();
584 consumer2.read(js,
585 ByteQueue::ReadRequest(kj::mv(prp2.resolver),
586 {
587 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
588 .type = ByteQueue::ReadRequest::Type::DEFAULT,
589 }));
590 prp2.promise.then(js, read2Continuation);
591 
592 js.runMicrotasks();
593 });
594}
595 
596KJ_TEST("ByteQueue with multiple byob consumers") {
597 preamble([](jsg::Lock& js) {
598 ByteQueue queue(2);
599 
600 ByteQueue::Consumer consumer1(queue);
601 ByteQueue::Consumer consumer2(queue);
602 
603 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
604 KJ_ASSERT(!result.done);
605 auto& value = KJ_ASSERT_NONNULL(result.value);
606 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
607 jsg::BufferSource source(js, value.getHandle(js));
608 auto ptr = source.asArrayPtr();
609 KJ_ASSERT(source.size() == 3);
610 KJ_ASSERT(ptr[0] == 'b');
611 KJ_ASSERT(ptr[1] == 'b');
612 KJ_ASSERT(ptr[2] == 'b');
613 
614 KJ_ASSERT(consumer1.size() == 0);
615 KJ_ASSERT(consumer2.size() == 0);
616 KJ_ASSERT(queue.size() == 0);
617 KJ_ASSERT(queue.desiredSize() == 2);
618 
619 return js.resolvedPromise(kj::mv(result));
620 }, 2);
621 
622 // Both reads will receive the data despite there being only a single
623 // byob read responded to.
624 byobRead(js, consumer1, 4).then(js, readContinuation);
625 byobRead(js, consumer2, 4).then(js, readContinuation);
626 
627 auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
628 auto nextPending = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
629 
630 KJ_ASSERT(!pendingByob->isInvalidated());
631 
632 auto& req = pendingByob->getRequest();
633 auto ptr = req.pullInto.store.asArrayPtr();
634 ptr.first(3).fill('b');
635 pendingByob->respond(js, 3);
636 KJ_ASSERT(pendingByob->isInvalidated());
637 
638 // No backpressure is signaled because both reads were fulfilled.
639 KJ_ASSERT(queue.desiredSize() == 2);
640 KJ_ASSERT(queue.size() == 0);
641 KJ_ASSERT(consumer1.size() == 0);
642 KJ_ASSERT(consumer2.size() == 0);
643 
644 // The next pendingByobReadRequest was invalidated.
645 KJ_ASSERT(nextPending->isInvalidated());
646 KJ_ASSERT(queue.nextPendingByobReadRequest() == kj::none);
647 
648 js.runMicrotasks();
649 });
650}
651 
652KJ_TEST("ByteQueue with multiple byob consumers") {
653 preamble([](jsg::Lock& js) {
654 ByteQueue queue(2);
655 
656 ByteQueue::Consumer consumer1(queue);
657 ByteQueue::Consumer consumer2(queue);
658 
659 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) -> auto {
660 KJ_ASSERT(!result.done);
661 auto& value = KJ_ASSERT_NONNULL(result.value);
662 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
663 jsg::BufferSource source(js, value.getHandle(js));
664 auto ptr = source.asArrayPtr();
665 KJ_ASSERT(source.size() == 3);
666 KJ_ASSERT(ptr[0] == 'b');
667 KJ_ASSERT(ptr[1] == 'b');
668 KJ_ASSERT(ptr[2] == 'b');
669 
670 KJ_ASSERT(consumer1.size() == 0);
671 KJ_ASSERT(consumer2.size() == 0);
672 KJ_ASSERT(queue.size() == 0);
673 KJ_ASSERT(queue.desiredSize() == 2);
674 
675 return js.resolvedPromise(kj::mv(result));
676 }, 2);
677 
678 // Both reads will receive the data despite there being only a single
679 // byob read responded to.
680 byobRead(js, consumer1, 4).then(js, readContinuation);
681 byobRead(js, consumer2, 4).then(js, readContinuation);
682 
683 auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
684 auto nextPending = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest());
685 
686 KJ_ASSERT(!pendingByob->isInvalidated());
687 
688 auto& req = pendingByob->getRequest();
689 auto ptr = req.pullInto.store.asArrayPtr();
690 ptr.first(3).fill('b');
691 pendingByob->respond(js, 3);
692 KJ_ASSERT(pendingByob->isInvalidated());
693 
694 // No backpressure is signaled because both reads were fulfilled.
695 KJ_ASSERT(queue.desiredSize() == 2);
696 KJ_ASSERT(queue.size() == 0);
697 KJ_ASSERT(consumer1.size() == 0);
698 KJ_ASSERT(consumer2.size() == 0);
699 
700 // The next pendingByobReadRequest was invalidated.
701 KJ_ASSERT(nextPending->isInvalidated());
702 KJ_ASSERT(queue.nextPendingByobReadRequest() == kj::none);
703 
704 js.runMicrotasks();
705 });
706}
707 
708KJ_TEST("ByteQueue with multiple byob consumers (multi-reads)") {
709 preamble([](jsg::Lock& js) {
710 ByteQueue queue(2);
711 
712 ByteQueue::Consumer consumer1(queue);
713 ByteQueue::Consumer consumer2(queue);
714 
715 MustCall<ReadContinuation> readConsumer1([&](jsg::Lock& js, auto&& result) -> auto {
716 KJ_ASSERT(!result.done);
717 auto& value = KJ_ASSERT_NONNULL(result.value);
718 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
719 jsg::BufferSource source(js, value.getHandle(js));
720 auto ptr = source.asArrayPtr();
721 KJ_ASSERT(source.size() == 3);
722 KJ_ASSERT(ptr[0] == 'a');
723 KJ_ASSERT(ptr[1] == 'a');
724 KJ_ASSERT(ptr[2] == 'a');
725 
726 return js.resolvedPromise(kj::mv(result));
727 });
728 
729 MustCall<ReadContinuation> readConsumer2([&](jsg::Lock& js, auto&& result) -> auto {
730 KJ_ASSERT(!result.done);
731 auto& value = KJ_ASSERT_NONNULL(result.value);
732 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
733 jsg::BufferSource source(js, value.getHandle(js));
734 auto ptr = source.asArrayPtr();
735 KJ_ASSERT(source.size() == 3);
736 KJ_ASSERT(ptr[0] == 'a');
737 KJ_ASSERT(ptr[1] == 'a');
738 KJ_ASSERT(ptr[2] == 'a');
739 
740 return byobRead(js, consumer2, 4);
741 });
742 
743 MustCall<ReadContinuation> secondReadBothConsumers([&](jsg::Lock& js, auto&& result) -> auto {
744 KJ_ASSERT(!result.done);
745 auto& value = KJ_ASSERT_NONNULL(result.value);
746 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
747 jsg::BufferSource source(js, value.getHandle(js));
748 auto ptr = source.asArrayPtr();
749 KJ_ASSERT(source.size() == 2);
750 KJ_ASSERT(ptr[0] == 'b');
751 KJ_ASSERT(ptr[1] == 'b');
752 
753 return js.resolvedPromise(kj::mv(result));
754 }, 2);
755 
756 // All reads will be fulfilled correctly even tho there are only two byob
757 // reads processed.
758 byobRead(js, consumer1, 4).then(js, readConsumer1);
759 byobRead(js, consumer1, 4).then(js, secondReadBothConsumers);
760 byobRead(js, consumer2, 4).then(js, readConsumer2).then(js, secondReadBothConsumers);
761 
762 // Although there are four distinct reads happening,
763 // there should only be two actual BYOB requests
764 // processed by the queue, which will fulfill all four
765 // reads.
766 MustCall<void(ByteQueue::ByobRequest&)> respond([&](jsg::Lock&, auto& pending) {
767 static uint counter = 0;
768 auto& req = pending.getRequest();
769 auto ptr = req.pullInto.store.asArrayPtr();
770 auto num = 3 - counter;
771 ptr.first(num).fill('a' + counter++);
772 pending.respond(js, num);
773 KJ_ASSERT(pending.isInvalidated());
774 }, 2);
775 
776 kj::Maybe<kj::Own<ByteQueue::ByobRequest>> pendingByob;
777 while ((pendingByob = queue.nextPendingByobReadRequest()) != kj::none) {
778 auto& pending = KJ_ASSERT_NONNULL(pendingByob);
779 if (pending->isInvalidated()) {
780 continue;
781 }
782 respond(js, *pending);
783 }
784 
785 js.runMicrotasks();
786 });
787}
788 
789KJ_TEST("ByteQueue with multiple byob consumers (multi-reads, 2)") {
790 preamble([](jsg::Lock& js) {
791 ByteQueue queue(2);
792 
793 ByteQueue::Consumer consumer1(queue);
794 ByteQueue::Consumer consumer2(queue);
795 
796 MustCall<ReadContinuation> readConsumer1([&](jsg::Lock& js, auto&& result) -> auto {
797 KJ_ASSERT(!result.done);
798 auto& value = KJ_ASSERT_NONNULL(result.value);
799 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
800 jsg::BufferSource source(js, value.getHandle(js));
801 auto ptr = source.asArrayPtr();
802 KJ_ASSERT(source.size() == 3);
803 KJ_ASSERT(ptr[0] == 'a');
804 KJ_ASSERT(ptr[1] == 'a');
805 KJ_ASSERT(ptr[2] == 'a');
806 return js.resolvedPromise(kj::mv(result));
807 });
808 
809 MustCall<ReadContinuation> readConsumer2([&](jsg::Lock& js, auto&& result) -> auto {
810 KJ_ASSERT(!result.done);
811 auto& value = KJ_ASSERT_NONNULL(result.value);
812 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
813 jsg::BufferSource source(js, value.getHandle(js));
814 auto ptr = source.asArrayPtr();
815 KJ_ASSERT(source.size() == 3);
816 KJ_ASSERT(ptr[0] == 'a');
817 KJ_ASSERT(ptr[1] == 'a');
818 KJ_ASSERT(ptr[2] == 'a');
819 
820 return byobRead(js, consumer2, 4);
821 });
822 
823 MustCall<ReadContinuation> secondReadBothConsumers([&](jsg::Lock& js, auto&& result) -> auto {
824 KJ_ASSERT(!result.done);
825 auto& value = KJ_ASSERT_NONNULL(result.value);
826 KJ_ASSERT(value.getHandle(js)->IsArrayBufferView());
827 jsg::BufferSource source(js, value.getHandle(js));
828 auto ptr = source.asArrayPtr();
829 KJ_ASSERT(source.size() == 2);
830 KJ_ASSERT(ptr[0] == 'b');
831 KJ_ASSERT(ptr[1] == 'b');
832 
833 return js.resolvedPromise(kj::mv(result));
834 }, 2);
835 
836 // All reads will be fulfilled correctly even tho there are only two BYOB reads
837 // responded to.
838 byobRead(js, consumer2, 4).then(js, readConsumer2).then(js, secondReadBothConsumers);
839 byobRead(js, consumer1, 4).then(js, readConsumer1);
840 byobRead(js, consumer1, 4).then(js, secondReadBothConsumers);
841 
842 // Although there are four distinct reads happening,
843 // there should only be two actual BYOB requests
844 // processed by the queue, which will fulfill all four
845 // reads.
846 MustCall<void(ByteQueue::ByobRequest&)> respond([&](jsg::Lock&, auto& pending) {
847 static uint counter = 0;
848 auto& req = pending.getRequest();
849 auto ptr = req.pullInto.store.asArrayPtr();
850 auto num = 3 - counter;
851 ptr.first(num).fill('a' + counter++);
852 pending.respond(js, num);
853 KJ_ASSERT(pending.isInvalidated());
854 }, 2);
855 
856 kj::Maybe<kj::Own<ByteQueue::ByobRequest>> pendingByob;
857 while ((pendingByob = queue.nextPendingByobReadRequest()) != kj::none) {
858 auto& pending = KJ_ASSERT_NONNULL(pendingByob);
859 if (pending->isInvalidated()) {
860 continue;
861 }
862 respond(js, *pending);
863 }
864 
865 js.runMicrotasks();
866 });
867}
868 
869KJ_TEST("ByteQueue with default consumer with atLeast") {
870 preamble([](jsg::Lock& js) {
871 ByteQueue queue(2);
872 
873 ByteQueue::Consumer consumer(queue);
874 
875 const auto read = [&](jsg::Lock& js, uint atLeast) {
876 auto prp = js.newPromiseAndResolver<ReadResult>();
877 consumer.read(js,
878 ByteQueue::ReadRequest(kj::mv(prp.resolver),
879 {
880 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)),
881 .atLeast = atLeast,
882 }));
883 return kj::mv(prp.promise);
884 };
885 
886 const auto push = [&](auto store) {
887 try {
888 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
889 } catch (kj::Exception& ex) {
890 KJ_DBG(ex.getDescription());
891 }
892 };
893 
894 MustCall<ReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
895 KJ_ASSERT(!result.done);
896 auto& value = KJ_ASSERT_NONNULL(result.value);
897 auto view = value.getHandle(js);
898 KJ_ASSERT(view->IsArrayBufferView());
899 jsg::BufferSource source(js, view);
900 auto ptr = source.asArrayPtr();
901 KJ_ASSERT(ptr[0] == 1);
902 KJ_ASSERT(ptr[1] == 2);
903 KJ_ASSERT(ptr[2] == 3);
904 KJ_ASSERT(ptr[3] == 4);
905 KJ_ASSERT(ptr[4] == 5);
906 KJ_ASSERT(source.size(), 5);
907 KJ_ASSERT(consumer.size(), 1);
908 return read(js, 1);
909 });
910 
911 MustCall<ReadContinuation> read2Continuation([&](jsg::Lock& js, auto&& result) {
912 KJ_ASSERT(!result.done);
913 auto& value = KJ_ASSERT_NONNULL(result.value);
914 auto view = value.getHandle(js);
915 KJ_ASSERT(view->IsArrayBufferView());
916 jsg::BufferSource source(js, view);
917 KJ_ASSERT(source.asArrayPtr()[0], 6);
918 KJ_ASSERT(source.size() == 1);
919 return js.resolvedPromise(kj::mv(result));
920 });
921 
922 read(js, 5).then(js, readContinuation).then(js, read2Continuation);
923 
924 auto store1 = jsg::BackingStore::alloc(js, 2);
925 store1.asArrayPtr()[0] = 1;
926 store1.asArrayPtr()[1] = 2;
927 push(kj::mv(store1));
928 
929 KJ_ASSERT(queue.desiredSize() == 0);
930 
931 auto store2 = jsg::BackingStore::alloc(js, 2);
932 store2.asArrayPtr()[0] = 3;
933 store2.asArrayPtr()[1] = 4;
934 push(kj::mv(store2));
935 
936 // Backpressure should be accumulating because the read has not yet fullilled.
937 KJ_ASSERT(queue.desiredSize() == -2);
938 
939 auto store3 = jsg::BackingStore::alloc(js, 2);
940 store3.asArrayPtr()[0] = 5;
941 store3.asArrayPtr()[1] = 6;
942 push(kj::mv(store3));
943 
944 // Some backpressure should be released because pushing the final minimum
945 // amount into the queue should have caused the read to be fulfilled.
946 KJ_ASSERT(queue.desiredSize() == 1);
947 
948 // There should be one unread byte left in the queue at this point.
949 // It will be read once the microtask queue is drained.
950 KJ_ASSERT(queue.size() == 1);
951 
952 js.runMicrotasks();
953 });
954}
955 
956KJ_TEST("ByteQueue with multiple default consumers with atLeast (same rate)") {
957 preamble([](jsg::Lock& js) {
958 ByteQueue queue(2);
959 
960 ByteQueue::Consumer consumer1(queue);
961 ByteQueue::Consumer consumer2(queue);
962 
963 const auto read = [&](jsg::Lock& js, auto& consumer, uint atLeast = 1) {
964 auto prp = js.newPromiseAndResolver<ReadResult>();
965 consumer.read(js,
966 ByteQueue::ReadRequest(kj::mv(prp.resolver),
967 {
968 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)),
969 .atLeast = atLeast,
970 }));
971 return kj::mv(prp.promise);
972 };
973 
974 const auto push = [&](auto store) {
975 try {
976 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
977 } catch (kj::Exception& ex) {
978 KJ_DBG(ex.getDescription());
979 }
980 };
981 
982 MustCall<ReadContinuation> read1Continuation([&](jsg::Lock& js, auto&& result) {
983 KJ_ASSERT(!result.done);
984 auto& value = KJ_ASSERT_NONNULL(result.value);
985 auto view = value.getHandle(js);
986 KJ_ASSERT(view->IsArrayBufferView());
987 jsg::BufferSource source(js, view);
988 auto ptr = source.asArrayPtr();
989 KJ_ASSERT(ptr[0] == 1);
990 KJ_ASSERT(ptr[1] == 2);
991 KJ_ASSERT(ptr[2] == 3);
992 KJ_ASSERT(ptr[3] == 4);
993 KJ_ASSERT(ptr[4] == 5);
994 KJ_ASSERT(source.size(), 5);
995 KJ_ASSERT(consumer1.size(), 1);
996 return read(js, consumer1);
997 });
998 
999 MustCall<ReadContinuation> read2Continuation([&](jsg::Lock& js, auto&& result) {
1000 KJ_ASSERT(!result.done);
1001 auto& value = KJ_ASSERT_NONNULL(result.value);
1002 auto view = value.getHandle(js);
1003 KJ_ASSERT(view->IsArrayBufferView());
1004 jsg::BufferSource source(js, view);
1005 auto ptr = source.asArrayPtr();
1006 KJ_ASSERT(ptr[0] == 1);
1007 KJ_ASSERT(ptr[1] == 2);
1008 KJ_ASSERT(ptr[2] == 3);
1009 KJ_ASSERT(ptr[3] == 4);
1010 KJ_ASSERT(ptr[4] == 5);
1011 KJ_ASSERT(source.size(), 5);
1012 KJ_ASSERT(consumer2.size(), 1);
1013 return read(js, consumer2);
1014 });
1015 
1016 MustCall<ReadContinuation> readFinalContinuation([&](jsg::Lock& js, auto&& result) {
1017 KJ_ASSERT(!result.done);
1018 auto& value = KJ_ASSERT_NONNULL(result.value);
1019 auto view = value.getHandle(js);
1020 KJ_ASSERT(view->IsArrayBufferView());
1021 jsg::BufferSource source(js, view);
1022 KJ_ASSERT(source.asArrayPtr()[0], 6);
1023 KJ_ASSERT(source.size() == 1);
1024 return js.resolvedPromise(kj::mv(result));
1025 }, 2);
1026 
1027 read(js, consumer1, 5).then(js, read1Continuation).then(js, readFinalContinuation);
1028 read(js, consumer2, 5).then(js, read2Continuation).then(js, readFinalContinuation);
1029 
1030 auto store1 = jsg::BackingStore::alloc(js, 2);
1031 store1.asArrayPtr()[0] = 1;
1032 store1.asArrayPtr()[1] = 2;
1033 push(kj::mv(store1));
1034 
1035 KJ_ASSERT(queue.desiredSize() == 0);
1036 
1037 auto store2 = jsg::BackingStore::alloc(js, 2);
1038 store2.asArrayPtr()[0] = 3;
1039 store2.asArrayPtr()[1] = 4;
1040 push(kj::mv(store2));
1041 
1042 // Backpressure should be accumulating because the read has not yet fullilled.
1043 KJ_ASSERT(queue.desiredSize() == -2);
1044 
1045 auto store3 = jsg::BackingStore::alloc(js, 2);
1046 store3.asArrayPtr()[0] = 5;
1047 store3.asArrayPtr()[1] = 6;
1048 push(kj::mv(store3));
1049 
1050 // Some backpressure should be released because pushing the final minimum
1051 // amount into the queue should have caused the read to be fulfilled.
1052 KJ_ASSERT(queue.desiredSize() == 1);
1053 
1054 // There should be one unread byte left in the queue at this point.
1055 // It will be read once the microtask queue is drained.
1056 KJ_ASSERT(queue.size() == 1);
1057 
1058 js.runMicrotasks();
1059 });
1060}
1061 
1062KJ_TEST("ByteQueue with multiple default consumers with atLeast (different rate)") {
1063 preamble([](jsg::Lock& js) {
1064 ByteQueue queue(2);
1065 
1066 ByteQueue::Consumer consumer1(queue);
1067 ByteQueue::Consumer consumer2(queue);
1068 
1069 const auto read = [&](jsg::Lock& js, auto& consumer, uint atLeast = 1) {
1070 auto prp = js.newPromiseAndResolver<ReadResult>();
1071 consumer.read(js,
1072 ByteQueue::ReadRequest(kj::mv(prp.resolver),
1073 {
1074 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)),
1075 .atLeast = atLeast,
1076 }));
1077 return kj::mv(prp.promise);
1078 };
1079 
1080 const auto push = [&](auto store) {
1081 try {
1082 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
1083 } catch (kj::Exception& ex) {
1084 KJ_DBG(ex.getDescription());
1085 }
1086 };
1087 
1088 MustCall<ReadContinuation> read1Continuation([&](jsg::Lock& js, auto&& result) {
1089 KJ_ASSERT(!result.done);
1090 auto& value = KJ_ASSERT_NONNULL(result.value);
1091 auto view = value.getHandle(js);
1092 KJ_ASSERT(view->IsArrayBufferView());
1093 jsg::BufferSource source(js, view);
1094 KJ_ASSERT(source.size() == 4);
1095 auto ptr = source.asArrayPtr();
1096 // Our read was for at least 3 bytes, with a maximum of 5.
1097 // For this first read, we received 4. One the second read
1098 // we should receive 2.
1099 KJ_ASSERT(ptr[0] == 1);
1100 KJ_ASSERT(ptr[1] == 2);
1101 KJ_ASSERT(ptr[2] == 3);
1102 KJ_ASSERT(ptr[3] == 4);
1103 return js.resolvedPromise(kj::mv(result));
1104 });
1105 
1106 MustCall<ReadContinuation> read1FinalContinuation([&](jsg::Lock& js, auto&& result) {
1107 KJ_ASSERT(!result.done);
1108 auto& value = KJ_ASSERT_NONNULL(result.value);
1109 auto view = value.getHandle(js);
1110 KJ_ASSERT(view->IsArrayBufferView());
1111 jsg::BufferSource source(js, view);
1112 KJ_ASSERT(source.size() == 2);
1113 auto ptr = source.asArrayPtr();
1114 KJ_ASSERT(ptr[0] == 5);
1115 KJ_ASSERT(ptr[1] == 6);
1116 return js.resolvedPromise(kj::mv(result));
1117 });
1118 
1119 MustCall<ReadContinuation> read2Continuation([&](jsg::Lock& js, auto&& result) {
1120 KJ_ASSERT(!result.done);
1121 auto& value = KJ_ASSERT_NONNULL(result.value);
1122 auto view = value.getHandle(js);
1123 KJ_ASSERT(view->IsArrayBufferView());
1124 jsg::BufferSource source(js, view);
1125 auto ptr = source.asArrayPtr();
1126 KJ_ASSERT(source.size() == 5);
1127 KJ_ASSERT(ptr[0] == 1);
1128 KJ_ASSERT(ptr[1] == 2);
1129 KJ_ASSERT(ptr[2] == 3);
1130 KJ_ASSERT(ptr[3] == 4);
1131 KJ_ASSERT(ptr[4] == 5);
1132 KJ_ASSERT(consumer2.size() == 1);
1133 return read(js, consumer2);
1134 });
1135 
1136 MustCall<ReadContinuation> read2FinalContinuation([&](jsg::Lock& js, auto&& result) {
1137 KJ_ASSERT(!result.done);
1138 auto& value = KJ_ASSERT_NONNULL(result.value);
1139 auto view = value.getHandle(js);
1140 KJ_ASSERT(view->IsArrayBufferView());
1141 jsg::BufferSource source(js, view);
1142 KJ_ASSERT(source.asArrayPtr()[0] == 6);
1143 KJ_ASSERT(source.size() == 1);
1144 return js.resolvedPromise(kj::mv(result));
1145 });
1146 
1147 // Consumer 1 will read in parallel with smaller minimum chunks...
1148 read(js, consumer1, 3).then(js, read1Continuation);
1149 read(js, consumer1).then(js, read1FinalContinuation);
1150 
1151 // Consumer 2 will read serially with a larger minimum chunk...
1152 read(js, consumer2, 5).then(js, read2Continuation).then(js, read2FinalContinuation);
1153 
1154 auto store1 = jsg::BackingStore::alloc(js, 2);
1155 store1.asArrayPtr()[0] = 1;
1156 store1.asArrayPtr()[1] = 2;
1157 push(kj::mv(store1));
1158 
1159 KJ_ASSERT(queue.desiredSize() == 0);
1160 
1161 auto store2 = jsg::BackingStore::alloc(js, 2);
1162 store2.asArrayPtr()[0] = 3;
1163 store2.asArrayPtr()[1] = 4;
1164 push(kj::mv(store2));
1165 
1166 // Consumer1 should not have any data buffered since its first read was for
1167 // between 3 and 5 bytes and it has received four so far.
1168 KJ_ASSERT(consumer1.size() == 0);
1169 
1170 // Consumer2 should have 4 bytes buffered since its first read was for 5 bytes
1171 // and we've only received 4 so far.
1172 KJ_ASSERT(consumer2.size() == 4);
1173 
1174 // Queue backpressure should reflect that consumer2 has data buffered.
1175 KJ_ASSERT(queue.desiredSize() == -2);
1176 
1177 auto store3 = jsg::BackingStore::alloc(js, 2);
1178 store3.asArrayPtr()[0] = 5;
1179 store3.asArrayPtr()[1] = 6;
1180 push(kj::mv(store3));
1181 
1182 // Most of the backpressure should have been resolved since we delivered 5 bytes
1183 // to consumer2, but there's still one byte remaining.
1184 KJ_ASSERT(queue.desiredSize() = 1);
1185 KJ_ASSERT(queue.size() == 1);
1186 
1187 js.runMicrotasks();
1188 });
1189}
1190 
1191#pragma endregion ByteQueue Tests
1192 
1193#pragma region Re-entrancy Safety Tests
1194 
1195// Test that pushing to a closed consumer doesn't crash.
1196// This can happen during QueueImpl::push() iteration when resolving a read
1197// on one consumer triggers JavaScript that closes another consumer.
1198KJ_TEST("ValueQueue push to closed consumer is safe") {
1199 preamble([](jsg::Lock& js) {
1200 ValueQueue queue(2);
1201 ValueQueue::Consumer consumer1(queue);
1202 ValueQueue::Consumer consumer2(queue);
1203 
1204 // Close consumer2
1205 consumer2.close(js);
1206 
1207 // Now push to the queue - this pushes to all consumers
1208 // Before the fix, this would crash when trying to push to closed consumer2
1209 queue.push(js, getEntry(js, 4));
1210 
1211 // consumer1 should have received the data
1212 KJ_ASSERT(consumer1.size() == 4);
1213 
1214 js.runMicrotasks();
1215 });
1216}
1217 
1218// Test that pushing to a cancelled consumer doesn't crash.
1219KJ_TEST("ValueQueue push to cancelled consumer is safe") {
1220 preamble([](jsg::Lock& js) {
1221 ValueQueue queue(2);
1222 ValueQueue::Consumer consumer1(queue);
1223 ValueQueue::Consumer consumer2(queue);
1224 
1225 // Cancel consumer2
1226 consumer2.cancel(js, kj::none);
1227 
1228 // Now push to the queue
1229 queue.push(js, getEntry(js, 4));
1230 
1231 // consumer1 should have received the data
1232 KJ_ASSERT(consumer1.size() == 4);
1233 
1234 js.runMicrotasks();
1235 });
1236}
1237 
1238// Test that pushing to an errored consumer doesn't crash.
1239KJ_TEST("ValueQueue push to errored consumer is safe") {
1240 preamble([](jsg::Lock& js) {
1241 ValueQueue queue(2);
1242 ValueQueue::Consumer consumer1(queue);
1243 ValueQueue::Consumer consumer2(queue);
1244 
1245 // Error consumer2
1246 consumer2.error(js, js.v8Ref(js.v8Error("error reason"_kj)));
1247 
1248 // Now push to the queue
1249 queue.push(js, getEntry(js, 4));
1250 
1251 // consumer1 should have received the data
1252 KJ_ASSERT(consumer1.size() == 4);
1253 
1254 js.runMicrotasks();
1255 });
1256}
1257 
1258// Test ByteQueue version of the safety checks
1259KJ_TEST("ByteQueue push to closed consumer is safe") {
1260 preamble([](jsg::Lock& js) {
1261 ByteQueue queue(10);
1262 ByteQueue::Consumer consumer1(queue);
1263 ByteQueue::Consumer consumer2(queue);
1264 
1265 // Close consumer2
1266 consumer2.close(js);
1267 
1268 // Now push to the queue
1269 auto store = jsg::BackingStore::alloc(js, 4);
1270 memset(store.asArrayPtr().begin(), 'A', 4);
1271 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store)));
1272 queue.push(js, kj::mv(entry));
1273 
1274 // consumer1 should have received the data
1275 KJ_ASSERT(consumer1.size() == 4);
1276 
1277 js.runMicrotasks();
1278 });
1279}
1280 
1281#pragma endregion Re - entrancy Safety Tests
1282 
1283#pragma region Draining Read Tests
1284 
1285using DrainingReadContinuation = jsg::Promise<DrainingReadResult>(DrainingReadResult&&);
1286using DrainingReadErrorContinuation = jsg::Promise<DrainingReadResult>(jsg::Value&&);
1287 
1288KJ_TEST("ValueQueue draining read with buffered data") {
1289 preamble([](jsg::Lock& js) {
1290 ValueQueue queue(10);
1291 ValueQueue::Consumer consumer(queue);
1292 
1293 // Push an ArrayBuffer
1294 auto store = jsg::BackingStore::alloc(js, 4);
1295 store.asArrayPtr()[0] = 'a';
1296 store.asArrayPtr()[1] = 'b';
1297 store.asArrayPtr()[2] = 'c';
1298 store.asArrayPtr()[3] = 'd';
1299 auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js);
1300 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab.As<v8::Value>()), 4));
1301 
1302 // Push a string
1303 auto str = jsg::v8Str(js.v8Isolate, "hello");
1304 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(str.As<v8::Value>()), 5));
1305 
1306 KJ_ASSERT(consumer.size() == 9);
1307 
1308 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1309 KJ_ASSERT(!result.done);
1310 KJ_ASSERT(result.chunks.size() == 2);
1311 
1312 // First chunk is the ArrayBuffer data
1313 KJ_ASSERT(result.chunks[0].size() == 4);
1314 KJ_ASSERT(result.chunks[0][0] == 'a');
1315 KJ_ASSERT(result.chunks[0][1] == 'b');
1316 KJ_ASSERT(result.chunks[0][2] == 'c');
1317 KJ_ASSERT(result.chunks[0][3] == 'd');
1318 
1319 // Second chunk is the string converted to UTF-8
1320 KJ_ASSERT(result.chunks[1].size() == 5);
1321 KJ_ASSERT(result.chunks[1][0] == 'h');
1322 KJ_ASSERT(result.chunks[1][1] == 'e');
1323 KJ_ASSERT(result.chunks[1][2] == 'l');
1324 KJ_ASSERT(result.chunks[1][3] == 'l');
1325 KJ_ASSERT(result.chunks[1][4] == 'o');
1326 
1327 KJ_ASSERT(consumer.size() == 0);
1328 return js.resolvedPromise(kj::mv(result));
1329 });
1330 
1331 consumer.drainingRead(js).then(js, readContinuation);
1332 js.runMicrotasks();
1333 });
1334}
1335 
1336KJ_TEST("ValueQueue draining read rejects with pending reads") {
1337 preamble([](jsg::Lock& js) {
1338 ValueQueue queue(10);
1339 ValueQueue::Consumer consumer(queue);
1340 
1341 // Queue a regular read
1342 auto prp = js.newPromiseAndResolver<ReadResult>();
1343 consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)});
1344 
1345 KJ_ASSERT(consumer.hasReadRequests());
1346 
1347 // Draining read should reject because there are pending reads
1348 MustNotCall<DrainingReadContinuation> readContinuation;
1349 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1350 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1351 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1352 });
1353 
1354 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1355 js.runMicrotasks();
1356 });
1357}
1358 
1359KJ_TEST("ValueQueue read rejects with pending draining read") {
1360 preamble([](jsg::Lock& js) {
1361 ValueQueue queue(10);
1362 ValueQueue::Consumer consumer(queue);
1363 
1364 // No data in buffer, draining read will queue a pending draining read
1365 consumer.drainingRead(js);
1366 
1367 KJ_ASSERT(consumer.hasPendingDrainingRead());
1368 
1369 // Regular read should reject because there's a pending draining read
1370 auto prp = js.newPromiseAndResolver<ReadResult>();
1371 
1372 MustNotCall<ReadContinuation> readContinuation;
1373 MustCall<ReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1374 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1375 return js.rejectedPromise<ReadResult>(kj::mv(value));
1376 });
1377 
1378 consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)});
1379 prp.promise.then(js, readContinuation, errorContinuation);
1380 js.runMicrotasks();
1381 });
1382}
1383 
1384KJ_TEST("ValueQueue draining read on closed stream") {
1385 preamble([](jsg::Lock& js) {
1386 ValueQueue queue(10);
1387 ValueQueue::Consumer consumer(queue);
1388 
1389 queue.close(js);
1390 
1391 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1392 KJ_ASSERT(result.done);
1393 KJ_ASSERT(result.chunks.size() == 0);
1394 return js.resolvedPromise(kj::mv(result));
1395 });
1396 
1397 consumer.drainingRead(js).then(js, readContinuation);
1398 js.runMicrotasks();
1399 });
1400}
1401 
1402KJ_TEST("ValueQueue draining read on errored stream") {
1403 preamble([](jsg::Lock& js) {
1404 ValueQueue queue(10);
1405 ValueQueue::Consumer consumer(queue);
1406 
1407 queue.error(js, js.v8Ref(js.v8Error("boom"_kj)));
1408 
1409 MustNotCall<DrainingReadContinuation> readContinuation;
1410 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1411 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1412 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1413 });
1414 
1415 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1416 js.runMicrotasks();
1417 });
1418}
1419 
1420KJ_TEST("ByteQueue draining read with buffered data") {
1421 preamble([](jsg::Lock& js) {
1422 ByteQueue queue(10);
1423 ByteQueue::Consumer consumer(queue);
1424 
1425 // Push first chunk
1426 auto store1 = jsg::BackingStore::alloc(js, 4);
1427 store1.asArrayPtr()[0] = 'a';
1428 store1.asArrayPtr()[1] = 'b';
1429 store1.asArrayPtr()[2] = 'c';
1430 store1.asArrayPtr()[3] = 'd';
1431 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store1))));
1432 
1433 // Push second chunk
1434 auto store2 = jsg::BackingStore::alloc(js, 3);
1435 store2.asArrayPtr()[0] = 'e';
1436 store2.asArrayPtr()[1] = 'f';
1437 store2.asArrayPtr()[2] = 'g';
1438 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store2))));
1439 
1440 KJ_ASSERT(consumer.size() == 7);
1441 
1442 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1443 KJ_ASSERT(!result.done);
1444 KJ_ASSERT(result.chunks.size() == 2);
1445 
1446 // First chunk
1447 KJ_ASSERT(result.chunks[0].size() == 4);
1448 KJ_ASSERT(result.chunks[0][0] == 'a');
1449 KJ_ASSERT(result.chunks[0][1] == 'b');
1450 KJ_ASSERT(result.chunks[0][2] == 'c');
1451 KJ_ASSERT(result.chunks[0][3] == 'd');
1452 
1453 // Second chunk
1454 KJ_ASSERT(result.chunks[1].size() == 3);
1455 KJ_ASSERT(result.chunks[1][0] == 'e');
1456 KJ_ASSERT(result.chunks[1][1] == 'f');
1457 KJ_ASSERT(result.chunks[1][2] == 'g');
1458 
1459 KJ_ASSERT(consumer.size() == 0);
1460 return js.resolvedPromise(kj::mv(result));
1461 });
1462 
1463 consumer.drainingRead(js).then(js, readContinuation);
1464 js.runMicrotasks();
1465 });
1466}
1467 
1468KJ_TEST("ByteQueue draining read rejects with pending reads") {
1469 preamble([](jsg::Lock& js) {
1470 ByteQueue queue(10);
1471 ByteQueue::Consumer consumer(queue);
1472 
1473 // Queue a regular read
1474 auto prp = js.newPromiseAndResolver<ReadResult>();
1475 consumer.read(js,
1476 ByteQueue::ReadRequest(kj::mv(prp.resolver),
1477 {
1478 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
1479 }));
1480 
1481 KJ_ASSERT(consumer.hasReadRequests());
1482 
1483 // Draining read should reject because there are pending reads
1484 MustNotCall<DrainingReadContinuation> readContinuation;
1485 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1486 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1487 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1488 });
1489 
1490 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1491 js.runMicrotasks();
1492 });
1493}
1494 
1495KJ_TEST("ByteQueue read rejects with pending draining read") {
1496 preamble([](jsg::Lock& js) {
1497 ByteQueue queue(10);
1498 ByteQueue::Consumer consumer(queue);
1499 
1500 // No data in buffer, draining read will queue a pending draining read
1501 consumer.drainingRead(js);
1502 
1503 KJ_ASSERT(consumer.hasPendingDrainingRead());
1504 
1505 // Regular read should reject because there's a pending draining read
1506 auto prp = js.newPromiseAndResolver<ReadResult>();
1507 
1508 MustNotCall<ReadContinuation> readContinuation;
1509 MustCall<ReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1510 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1511 return js.rejectedPromise<ReadResult>(kj::mv(value));
1512 });
1513 
1514 consumer.read(js,
1515 ByteQueue::ReadRequest(kj::mv(prp.resolver),
1516 {
1517 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
1518 }));
1519 prp.promise.then(js, readContinuation, errorContinuation);
1520 js.runMicrotasks();
1521 });
1522}
1523 
1524KJ_TEST("ByteQueue draining read on closed stream") {
1525 preamble([](jsg::Lock& js) {
1526 ByteQueue queue(10);
1527 ByteQueue::Consumer consumer(queue);
1528 
1529 queue.close(js);
1530 
1531 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1532 KJ_ASSERT(result.done);
1533 KJ_ASSERT(result.chunks.size() == 0);
1534 return js.resolvedPromise(kj::mv(result));
1535 });
1536 
1537 consumer.drainingRead(js).then(js, readContinuation);
1538 js.runMicrotasks();
1539 });
1540}
1541 
1542KJ_TEST("ByteQueue draining read on errored stream") {
1543 preamble([](jsg::Lock& js) {
1544 ByteQueue queue(10);
1545 ByteQueue::Consumer consumer(queue);
1546 
1547 queue.error(js, js.v8Ref(js.v8Error("boom"_kj)));
1548 
1549 MustNotCall<DrainingReadContinuation> readContinuation;
1550 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1551 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1552 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1553 });
1554 
1555 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1556 js.runMicrotasks();
1557 });
1558}
1559 
1560KJ_TEST("ValueQueue draining read with close signal") {
1561 preamble([](jsg::Lock& js) {
1562 ValueQueue queue(10);
1563 ValueQueue::Consumer consumer(queue);
1564 
1565 // Push some data
1566 auto store = jsg::BackingStore::alloc(js, 4);
1567 store.asArrayPtr()[0] = 'a';
1568 store.asArrayPtr()[1] = 'b';
1569 store.asArrayPtr()[2] = 'c';
1570 store.asArrayPtr()[3] = 'd';
1571 auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js);
1572 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab.As<v8::Value>()), 4));
1573 
1574 // Close the queue
1575 queue.close(js);
1576 
1577 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1578 // Should have the data and done should be true since stream is closed
1579 KJ_ASSERT(result.done);
1580 KJ_ASSERT(result.chunks.size() == 1);
1581 KJ_ASSERT(result.chunks[0].size() == 4);
1582 return js.resolvedPromise(kj::mv(result));
1583 });
1584 
1585 consumer.drainingRead(js).then(js, readContinuation);
1586 js.runMicrotasks();
1587 });
1588}
1589 
1590KJ_TEST("ByteQueue draining read with close signal") {
1591 preamble([](jsg::Lock& js) {
1592 ByteQueue queue(10);
1593 ByteQueue::Consumer consumer(queue);
1594 
1595 // Push some data
1596 auto store = jsg::BackingStore::alloc(js, 4);
1597 store.asArrayPtr()[0] = 'a';
1598 store.asArrayPtr()[1] = 'b';
1599 store.asArrayPtr()[2] = 'c';
1600 store.asArrayPtr()[3] = 'd';
1601 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
1602 
1603 // Close the queue
1604 queue.close(js);
1605 
1606 MustCall<DrainingReadContinuation> readContinuation([&](jsg::Lock& js, auto&& result) {
1607 // Should have the data and done should be true since stream is closed
1608 KJ_ASSERT(result.done);
1609 KJ_ASSERT(result.chunks.size() == 1);
1610 KJ_ASSERT(result.chunks[0].size() == 4);
1611 return js.resolvedPromise(kj::mv(result));
1612 });
1613 
1614 consumer.drainingRead(js).then(js, readContinuation);
1615 js.runMicrotasks();
1616 });
1617}
1618 
1619KJ_TEST("ValueQueue draining read errors on non-byte value") {
1620 // Test that drainingRead correctly errors when encountering a value that
1621 // cannot be converted to bytes (not ArrayBuffer, ArrayBufferView, or string).
1622 preamble([](jsg::Lock& js) {
1623 ValueQueue queue(10);
1624 ValueQueue::Consumer consumer(queue);
1625 
1626 // Push a plain object - this cannot be converted to bytes
1627 auto obj = v8::Object::New(js.v8Isolate);
1628 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(obj.As<v8::Value>()), 1));
1629 
1630 KJ_ASSERT(consumer.size() == 1);
1631 
1632 MustNotCall<DrainingReadContinuation> readContinuation;
1633 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1634 // Should get a TypeError about non-convertible value
1635 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1636 auto message = jsg::JsValue(value.getHandle(js))
1637 .tryCast<jsg::JsObject>()
1638 .map([&](jsg::JsObject obj) {
1639 return obj.get(js, "message")
1640 .tryCast<jsg::JsString>()
1641 .map([&](jsg::JsString str) {
1642 return str.toString(js);
1643 }).orDefault(kj::str());
1644 }).orDefault(kj::str());
1645 KJ_ASSERT(message.contains("cannot be converted to bytes"));
1646 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1647 });
1648 
1649 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1650 js.runMicrotasks();
1651 // MustCall verifies the error continuation was called
1652 });
1653}
1654 
1655KJ_TEST("ValueQueue draining read errors on number value") {
1656 // Another non-byte value test with a number
1657 preamble([](jsg::Lock& js) {
1658 ValueQueue queue(10);
1659 ValueQueue::Consumer consumer(queue);
1660 
1661 // Push a number - this cannot be converted to bytes
1662 auto num = v8::Number::New(js.v8Isolate, 42);
1663 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(num.As<v8::Value>()), 1));
1664 
1665 MustNotCall<DrainingReadContinuation> readContinuation;
1666 MustCall<DrainingReadErrorContinuation> errorContinuation([&](jsg::Lock& js, auto&& value) {
1667 KJ_ASSERT(value.getHandle(js)->IsNativeError());
1668 return js.rejectedPromise<DrainingReadResult>(kj::mv(value));
1669 });
1670 
1671 consumer.drainingRead(js).then(js, readContinuation, errorContinuation);
1672 js.runMicrotasks();
1673 // MustCall verifies the error continuation was called
1674 });
1675}
1676 
1677#pragma endregion Draining Read Tests
1678 
1679#pragma region Draining Read maxRead Tests
1680 
1681// Tests for the maxRead soft limit parameter. Both the initial buffer drain and subsequent
1682// synchronous pump attempts stop when totalRead reaches maxRead. This prevents unbounded
1683// memory accumulation when a fast producer outpaces a slow consumer.
1684//
1685// Note: Testing the pump loop behavior requires the full controller infrastructure (stateListener).
1686// These tests focus on the buffer drain behavior which can be tested without a listener.
1687 
1688KJ_TEST("ValueQueue draining read respects maxRead during buffer drain") {
1689 preamble([](jsg::Lock& js) {
1690 ValueQueue queue(10);
1691 ValueQueue::Consumer consumer(queue);
1692 
1693 // Buffer 200 bytes of data (two 100-byte chunks)
1694 auto store1 = jsg::BackingStore::alloc(js, 100);
1695 store1.asArrayPtr().fill(0xAA);
1696 auto ab1 = jsg::BufferSource(js, kj::mv(store1)).getHandle(js);
1697 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab1.As<v8::Value>()), 100));
1698 
1699 auto store2 = jsg::BackingStore::alloc(js, 100);
1700 store2.asArrayPtr().fill(0xBB);
1701 auto ab2 = jsg::BufferSource(js, kj::mv(store2)).getHandle(js);
1702 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab2.As<v8::Value>()), 100));
1703 
1704 KJ_ASSERT(consumer.size() == 200);
1705 
1706 // maxRead=50: the first 100-byte chunk is drained (exceeding maxRead after the first item),
1707 // then draining stops because totalRead (100) >= maxRead (50). The second chunk stays buffered.
1708 MustCall<DrainingReadContinuation> readContinuation(
1709 [&](jsg::Lock& js, DrainingReadResult&& result) {
1710 KJ_ASSERT(!result.done);
1711 // Should only have the first chunk (maxRead stopped draining before the second)
1712 KJ_ASSERT(result.chunks.size() == 1);
1713 KJ_ASSERT(result.chunks[0].size() == 100);
1714 // Second chunk should still be buffered
1715 KJ_ASSERT(consumer.size() == 100);
1716 return js.resolvedPromise(kj::mv(result));
1717 });
1718 
1719 consumer.drainingRead(js, 50).then(js, readContinuation);
1720 js.runMicrotasks();
1721 });
1722}
1723 
1724KJ_TEST("ByteQueue draining read respects maxRead during buffer drain") {
1725 preamble([](jsg::Lock& js) {
1726 ByteQueue queue(10);
1727 ByteQueue::Consumer consumer(queue);
1728 
1729 // Buffer 200 bytes of data (two 100-byte chunks)
1730 auto store1 = jsg::BackingStore::alloc(js, 100);
1731 store1.asArrayPtr().fill(0xAA);
1732 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store1))));
1733 
1734 auto store2 = jsg::BackingStore::alloc(js, 100);
1735 store2.asArrayPtr().fill(0xBB);
1736 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store2))));
1737 
1738 KJ_ASSERT(consumer.size() == 200);
1739 
1740 // maxRead=50: first 100-byte chunk is drained, then stops. Second chunk stays buffered.
1741 MustCall<DrainingReadContinuation> readContinuation(
1742 [&](jsg::Lock& js, DrainingReadResult&& result) {
1743 KJ_ASSERT(!result.done);
1744 KJ_ASSERT(result.chunks.size() == 1);
1745 KJ_ASSERT(result.chunks[0].size() == 100);
1746 KJ_ASSERT(consumer.size() == 100);
1747 return js.resolvedPromise(kj::mv(result));
1748 });
1749 
1750 consumer.drainingRead(js, 50).then(js, readContinuation);
1751 js.runMicrotasks();
1752 });
1753}
1754 
1755KJ_TEST("ValueQueue draining read with large maxRead drains entire buffer") {
1756 preamble([](jsg::Lock& js) {
1757 ValueQueue queue(10);
1758 ValueQueue::Consumer consumer(queue);
1759 
1760 // Buffer 200 bytes (two 100-byte chunks)
1761 auto store1 = jsg::BackingStore::alloc(js, 100);
1762 store1.asArrayPtr().fill(0xAA);
1763 auto ab1 = jsg::BufferSource(js, kj::mv(store1)).getHandle(js);
1764 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab1.As<v8::Value>()), 100));
1765 
1766 auto store2 = jsg::BackingStore::alloc(js, 100);
1767 store2.asArrayPtr().fill(0xBB);
1768 auto ab2 = jsg::BufferSource(js, kj::mv(store2)).getHandle(js);
1769 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab2.As<v8::Value>()), 100));
1770 
1771 KJ_ASSERT(consumer.size() == 200);
1772 
1773 // maxRead=1000: both chunks should be drained since total (200) < maxRead (1000)
1774 MustCall<DrainingReadContinuation> readContinuation(
1775 [&](jsg::Lock& js, DrainingReadResult&& result) {
1776 KJ_ASSERT(!result.done);
1777 KJ_ASSERT(result.chunks.size() == 2);
1778 KJ_ASSERT(result.chunks[0].size() == 100);
1779 KJ_ASSERT(result.chunks[1].size() == 100);
1780 KJ_ASSERT(consumer.size() == 0);
1781 return js.resolvedPromise(kj::mv(result));
1782 });
1783 
1784 consumer.drainingRead(js, 1000).then(js, readContinuation);
1785 js.runMicrotasks();
1786 });
1787}
1788 
1789KJ_TEST("ValueQueue draining read with default maxRead (unlimited)") {
1790 preamble([](jsg::Lock& js) {
1791 ValueQueue queue(10);
1792 ValueQueue::Consumer consumer(queue);
1793 
1794 // Buffer some data
1795 auto store = jsg::BackingStore::alloc(js, 100);
1796 store.asArrayPtr().fill(0xAA);
1797 auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js);
1798 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab.As<v8::Value>()), 100));
1799 
1800 // Default maxRead (kj::maxValue) should drain buffer normally
1801 MustCall<DrainingReadContinuation> readContinuation(
1802 [&](jsg::Lock& js, DrainingReadResult&& result) {
1803 KJ_ASSERT(!result.done);
1804 KJ_ASSERT(result.chunks.size() == 1);
1805 KJ_ASSERT(result.chunks[0].size() == 100);
1806 return js.resolvedPromise(kj::mv(result));
1807 });
1808 
1809 consumer.drainingRead(js).then(js, readContinuation); // No maxRead argument = default
1810 js.runMicrotasks();
1811 });
1812}
1813 
1814KJ_TEST("ValueQueue draining read maxRead bounds multiple iterations") {
1815 // Verify that successive drainingRead calls with maxRead correctly drain
1816 // a large buffer incrementally.
1817 preamble([](jsg::Lock& js) {
1818 ValueQueue queue(10);
1819 ValueQueue::Consumer consumer(queue);
1820 
1821 // Buffer 400 bytes: four 100-byte chunks
1822 for (int i = 0; i < 4; i++) {
1823 auto store = jsg::BackingStore::alloc(js, 100);
1824 store.asArrayPtr().fill(0x10 * (i + 1));
1825 auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js);
1826 queue.push(js, kj::rc<ValueQueue::Entry>(js.v8Ref(ab.As<v8::Value>()), 100));
1827 }
1828 KJ_ASSERT(consumer.size() == 400);
1829 
1830 // First read with maxRead=150: drains first chunk (100 bytes, now totalRead=100 < 150),
1831 // then drains second chunk (200 bytes total, now >= 150), stops.
1832 MustCall<DrainingReadContinuation> read1([&](jsg::Lock& js, DrainingReadResult&& result) {
1833 KJ_ASSERT(!result.done);
1834 KJ_ASSERT(result.chunks.size() == 2);
1835 KJ_ASSERT(consumer.size() == 200);
1836 return js.resolvedPromise(kj::mv(result));
1837 });
1838 consumer.drainingRead(js, 150).then(js, read1);
1839 js.runMicrotasks();
1840 
1841 // Second read with maxRead=150: drains next two chunks similarly
1842 MustCall<DrainingReadContinuation> read2([&](jsg::Lock& js, DrainingReadResult&& result) {
1843 KJ_ASSERT(!result.done);
1844 KJ_ASSERT(result.chunks.size() == 2);
1845 KJ_ASSERT(consumer.size() == 0);
1846 return js.resolvedPromise(kj::mv(result));
1847 });
1848 consumer.drainingRead(js, 150).then(js, read2);
1849 js.runMicrotasks();
1850 });
1851}
1852 
1853#pragma endregion Draining Read maxRead Tests
1854 
1855#pragma region Queue/Consumer Destruction Order Tests
1856 
1857// These tests verify that destroying the queue before its consumers doesn't crash.
1858// This can happen during isolate teardown when the destruction order isn't guaranteed
1859// to follow the ownership hierarchy.
1860// These will typically only catch in builds with AddressSanitizer enabled.
1861 
1862KJ_TEST("ValueQueue destroyed before consumer doesn't crash") {
1863 preamble([](jsg::Lock& js) {
1864 // Heap-allocate the queue so we can control its destruction order
1865 auto queue = kj::heap<ValueQueue>(2);
1866 
1867 // Create a consumer attached to the queue
1868 auto consumer = kj::heap<ValueQueue::Consumer>(*queue);
1869 
1870 // Push some data to make sure the consumer has state
1871 queue->push(js, getEntry(js, 4));
1872 KJ_ASSERT(consumer->size() == 4);
1873 
1874 // Now destroy the queue FIRST - this simulates the production scenario
1875 // where wrapper cleanup destroys the controller (and its queue) before
1876 // the consumer that holds a reference to it.
1877 queue = nullptr;
1878 
1879 // When the consumer is destroyed (here, or when going out of scope),
1880 // its destructor calls queue.removeConsumer(this).
1881 // Without the fix, this is a use-after-free.
1882 // With the fix, the consumer knows the queue is gone and skips the call.
1883 consumer = nullptr;
1884 
1885 // If we get here without crashing, the test passes
1886 });
1887}
1888 
1889KJ_TEST("ValueQueue destroyed before multiple consumers doesn't crash") {
1890 preamble([](jsg::Lock& js) {
1891 auto queue = kj::heap<ValueQueue>(2);
1892 
1893 auto consumer1 = kj::heap<ValueQueue::Consumer>(*queue);
1894 auto consumer2 = kj::heap<ValueQueue::Consumer>(*queue);
1895 
1896 queue->push(js, getEntry(js, 4));
1897 KJ_ASSERT(consumer1->size() == 4);
1898 KJ_ASSERT(consumer2->size() == 4);
1899 
1900 // Destroy queue before consumers
1901 queue = nullptr;
1902 
1903 // Both consumers should handle destruction gracefully
1904 consumer1 = nullptr;
1905 consumer2 = nullptr;
1906 });
1907}
1908 
1909KJ_TEST("ByteQueue destroyed before consumer doesn't crash") {
1910 preamble([](jsg::Lock& js) {
1911 auto queue = kj::heap<ByteQueue>(2);
1912 auto consumer = kj::heap<ByteQueue::Consumer>(*queue);
1913 
1914 auto store = jsg::BackingStore::alloc(js, 4);
1915 store.asArrayPtr().fill('a');
1916 queue->push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
1917 KJ_ASSERT(consumer->size() == 4);
1918 
1919 // Destroy queue before consumer
1920 queue = nullptr;
1921 consumer = nullptr;
1922 });
1923}
1924 
1925KJ_TEST("ValueQueue destroyed with pending read requests doesn't crash") {
1926 preamble([](jsg::Lock& js) {
1927 auto queue = kj::heap<ValueQueue>(2);
1928 auto consumer = kj::heap<ValueQueue::Consumer>(*queue);
1929 
1930 // Queue a read request (no data pushed, so it will be pending)
1931 auto prp = js.newPromiseAndResolver<ReadResult>();
1932 consumer->read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)});
1933 
1934 KJ_ASSERT(consumer->hasReadRequests());
1935 
1936 // Destroy queue while there are pending reads
1937 queue = nullptr;
1938 
1939 // Consumer destruction should handle this gracefully
1940 consumer = nullptr;
1941 
1942 js.runMicrotasks();
1943 });
1944}
1945 
1946KJ_TEST("ValueQueue close then destroy before consumer doesn't crash") {
1947 preamble([](jsg::Lock& js) {
1948 auto queue = kj::heap<ValueQueue>(2);
1949 auto consumer = kj::heap<ValueQueue::Consumer>(*queue);
1950 
1951 // Close the queue first
1952 queue->close(js);
1953 
1954 // Then destroy it
1955 queue = nullptr;
1956 
1957 // Consumer should still handle destruction gracefully
1958 consumer = nullptr;
1959 });
1960}
1961 
1962KJ_TEST("ValueQueue error then destroy before consumer doesn't crash") {
1963 preamble([](jsg::Lock& js) {
1964 auto queue = kj::heap<ValueQueue>(2);
1965 auto consumer = kj::heap<ValueQueue::Consumer>(*queue);
1966 
1967 // Error the queue first
1968 queue->error(js, js.v8Ref(js.v8Error("boom"_kj)));
1969 
1970 // Then destroy it
1971 queue = nullptr;
1972 
1973 // Consumer should still handle destruction gracefully
1974 consumer = nullptr;
1975 });
1976}
1977 
1978#pragma endregion Queue / Consumer Destruction Order Tests
1979 
1980#pragma region Consumer Destroyed During Push Tests
1981// These tests verify that the queue handles consumer destruction gracefully.
1982// QueueImpl::push() takes a snapshot of consumers and iterates over it. If a consumer
1983// is destroyed (removed from allConsumers) between the snapshot and the iteration,
1984// the code must check allConsumers.contains() before dereferencing the pointer.
1985// That said, the tests do not actually fail without the relevant fix because the
1986// issue is extremely timing-dependent and difficult to trigger deterministically, even
1987// with asan enabled. These tests at least document the intended behavior and may catch
1988// future regressions.
1989 
1990KJ_TEST("ByteQueue push skips consumer removed from queue during iteration") {
1991 preamble([](jsg::Lock& js) {
1992 ByteQueue queue(10);
1993 
1994 // Create two consumers
1995 auto consumer1 = kj::heap<ByteQueue::Consumer>(queue);
1996 auto consumer2 = kj::heap<ByteQueue::Consumer>(queue);
1997 
1998 // Destroy consumer2 BEFORE pushing. This directly tests that push()
1999 // checks if consumers still exist before calling push on them.
2000 // The snapshot taken at the start of push() would have included consumer2,
2001 // but consumer2 is no longer in allConsumers when we iterate.
2002 consumer2 = nullptr;
2003 
2004 // Push data - should not crash even though consumer2 was in the queue
2005 // when it was created but is now destroyed.
2006 auto store = jsg::BackingStore::alloc(js, 4);
2007 store.asArrayPtr().fill('x');
2008 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
2009 
2010 // consumer1 should have received the data
2011 KJ_ASSERT(consumer1->size() == 4);
2012 });
2013}
2014 
2015KJ_TEST("ValueQueue push skips consumer removed from queue during iteration") {
2016 preamble([](jsg::Lock& js) {
2017 ValueQueue queue(10);
2018 
2019 auto consumer1 = kj::heap<ValueQueue::Consumer>(queue);
2020 auto consumer2 = kj::heap<ValueQueue::Consumer>(queue);
2021 
2022 // Destroy consumer2 before pushing
2023 consumer2 = nullptr;
2024 
2025 queue.push(js, getEntry(js, 4));
2026 
2027 KJ_ASSERT(consumer1->size() == 4);
2028 });
2029}
2030 
2031KJ_TEST("ByteQueue push handles consumer destroyed by microtask between pushes") {
2032 preamble([](jsg::Lock& js) {
2033 ByteQueue queue(10);
2034 
2035 auto consumer1 = kj::heap<ByteQueue::Consumer>(queue);
2036 auto consumer2 = kj::heap<ByteQueue::Consumer>(queue);
2037 
2038 // Set up a pending read on consumer1
2039 auto prp = js.newPromiseAndResolver<ReadResult>();
2040 consumer1->read(js,
2041 ByteQueue::ReadRequest(kj::mv(prp.resolver),
2042 {
2043 .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)),
2044 }));
2045 
2046 // The continuation destroys consumer2
2047 MustCall<ReadContinuation> readContinuation([&consumer2](jsg::Lock& js, ReadResult&& result) {
2048 consumer2 = nullptr;
2049 return js.resolvedPromise(kj::mv(result));
2050 });
2051 prp.promise.then(js, readContinuation);
2052 
2053 // First push - resolves consumer1's read, schedules microtask that will destroy consumer2
2054 auto store1 = jsg::BackingStore::alloc(js, 4);
2055 store1.asArrayPtr().fill('x');
2056 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store1))));
2057 
2058 // Run microtasks - this destroys consumer2
2059 js.runMicrotasks();
2060 
2061 // Second push - consumer2 is now destroyed, should not crash
2062 auto store2 = jsg::BackingStore::alloc(js, 4);
2063 store2.asArrayPtr().fill('y');
2064 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store2))));
2065 
2066 // consumer1 should have the second push's data buffered
2067 KJ_ASSERT(consumer1->size() == 4);
2068 });
2069}
2070 
2071KJ_TEST("ByteQueue maybeUpdateBackpressure skips destroyed consumers") {
2072 preamble([](jsg::Lock& js) {
2073 ByteQueue queue(10);
2074 
2075 auto consumer1 = kj::heap<ByteQueue::Consumer>(queue);
2076 auto consumer2 = kj::heap<ByteQueue::Consumer>(queue);
2077 
2078 // Push some data so consumers have size
2079 auto store = jsg::BackingStore::alloc(js, 4);
2080 store.asArrayPtr().fill('x');
2081 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store))));
2082 
2083 KJ_ASSERT(consumer1->size() == 4);
2084 KJ_ASSERT(consumer2->size() == 4);
2085 KJ_ASSERT(queue.size() == 4);
2086 
2087 // Destroy consumer2
2088 consumer2 = nullptr;
2089 
2090 // Trigger backpressure recalculation by pushing more data
2091 auto store2 = jsg::BackingStore::alloc(js, 4);
2092 store2.asArrayPtr().fill('y');
2093 queue.push(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, kj::mv(store2))));
2094 
2095 // Should not crash, and size should reflect only consumer1
2096 KJ_ASSERT(consumer1->size() == 8);
2097 KJ_ASSERT(queue.size() == 8);
2098 });
2099}
2100#pragma endregion Consumer Destroyed During Push Tests
2101 
2102} // namespace
2103} // namespace workerd::api