// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "queue.h" #include #include #include namespace workerd::api { namespace { void preamble(auto callback) { TestFixture fixture; fixture.runInIoContext([&](const TestFixture::Environment& env) { callback(env.js); }); } using ReadContinuation = jsg::Promise(ReadResult&&); using CloseContinuation = jsg::Promise(ReadResult&&); using ReadErrorContinuation = jsg::Promise(jsg::Value&&); const kj::MutexGuarded unwindDetectorMutex; template struct MustCall; // Used to create a jsg::Promise continuation function that must be called // at least once during the test. If the function is not called, an error // will be thrown causing the test to fail. // TODO(cleanup): Consider adding this to jsg-test.h template struct MustNotCall; // Used to create a jsg::Promise continuation function that must not be called // during the test. If the function is called, an error will be thrown causing // the test to fail. // TODO(cleanup): Consider adding this to jsg-test.h template struct MustCall { using Func = jsg::Function; Func fn; uint expected; kj::SourceLocation location; uint called = false; MustCall(Func fn, uint expected = 1, kj::SourceLocation location = kj::SourceLocation()) : fn(kj::mv(fn)), expected(expected), location(location) {} ~MustCall() { auto unwindDetector = unwindDetectorMutex.lockExclusive(); if (!unwindDetector->isUnwinding()) { KJ_ASSERT(called == expected, kj::str("MustCall function was not called ", expected, " times. [actual: ", called, "]"), location); } } Ret operator()(jsg::Lock& js, Args&&... args) { called++; return fn(js, kj::fwd(args...)); } }; template struct MustNotCall { MustNotCall(kj::SourceLocation location = kj::SourceLocation()): location(location) {} kj::SourceLocation location; Ret operator()(jsg::Lock&, Args... args) { KJ_FAIL_REQUIRE("MustNotCall function was called!", location); } }; auto read(jsg::Lock& js, auto& consumer) { auto prp = js.newPromiseAndResolver(); consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)}); return kj::mv(prp.promise); } auto byobRead(jsg::Lock& js, auto& consumer, int size) { auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, size)), .type = ByteQueue::ReadRequest::Type::BYOB, })); return kj::mv(prp.promise); }; auto getEntry(jsg::Lock& js, auto size) { return kj::rc(js.v8Ref(v8::True(js.v8Isolate).As()), size); } #pragma region ValueQueue Tests KJ_TEST("ValueQueue basics work") { preamble([](jsg::Lock& js) { ValueQueue queue(2); // At this point, there are no consumers, data does not get enqueued. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); queue.push(js, getEntry(js, 1)); // Because there are no consumers, there is no change to backpressure. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); // Closing the queue causes the desiredSize to be zero. queue.close(js); try { queue.push(js, getEntry(js, 1)); KJ_FAIL_ASSERT("The queue push after close should have failed."); } catch (kj::Exception& ex) { KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored.")); } KJ_ASSERT(queue.desiredSize() == 0); KJ_ASSERT(queue.size() == 0); }); } KJ_TEST("ValueQueue erroring works") { preamble([](jsg::Lock& js) { ValueQueue queue(2); queue.error(js, js.v8Ref(js.v8Error("boom"_kj))); KJ_ASSERT(queue.desiredSize() == 0); try { queue.push(js, getEntry(js, 1)); KJ_FAIL_ASSERT("The queue push after close should have failed."); } catch (kj::Exception& ex) { KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored.")); } }); } KJ_TEST("ValueQueue with single consumer") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer(queue); KJ_ASSERT(queue.desiredSize() == 2); queue.push(js, getEntry(js, 2)); // The item was pushed into the consumer. KJ_ASSERT(consumer.size() == 2); // The queue size and desiredSize were updated accordingly. KJ_ASSERT(queue.size() == 2); KJ_ASSERT(queue.desiredSize() == 0); auto prp = js.newPromiseAndResolver(); consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)}); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsTrue()); KJ_ASSERT(consumer.size() == 0); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }); prp.promise.then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue with multiple consumers") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer1(queue); ValueQueue::Consumer consumer2(queue); KJ_ASSERT(queue.desiredSize() == 2); queue.push(js, getEntry(js, 2)); // The item was pushed into the consumer. KJ_ASSERT(consumer1.size() == 2); KJ_ASSERT(consumer2.size() == 2); // The queue size and desiredSize were updated accordingly. KJ_ASSERT(queue.size() == 2); KJ_ASSERT(queue.desiredSize() == 0); MustCall read1Continuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsTrue()); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 2); // Backpressure was not relieved since the other consumer has yet to read. KJ_ASSERT(queue.size() == 2); KJ_ASSERT(queue.desiredSize() == 0); return read(js, consumer2); }); MustCall read2Continuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsTrue()); KJ_ASSERT(consumer2.size() == 0); // Backpressure was relieved since both consumers have now read. KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }); MustCall close1Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(result.done); return read(js, consumer2); }); MustCall close2Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(result.done); return js.resolvedPromise(); }); read(js, consumer1).then(js, read1Continuation).then(js, read2Continuation); js.runMicrotasks(); // Closing the queue causes both consumers to be closed... queue.close(js); // After close, the consumers will still be usable, but the queue itself // has shutdown and no longer reports backpressure. KJ_ASSERT(queue.desiredSize() == 0); KJ_ASSERT(queue.size() == 0); read(js, consumer1).then(js, close1Continuation).then(js, close2Continuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue consumer with multiple-reads") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer(queue); // The first read will produce a value. MustCall read1Continuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsTrue()); return js.resolvedPromise(kj::mv(result)); }); read(js, consumer).then(js, read1Continuation); // The second and third reads will both be done = true MustCall closeContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(result.done); return js.resolvedPromise(); }, 2); read(js, consumer).then(js, closeContinuation); read(js, consumer).then(js, closeContinuation); queue.push(js, getEntry(js, 2)); // Because there is a consumer reading when the push happens, no backpressure // is applied... KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); queue.close(js); js.runMicrotasks(); }); } KJ_TEST("ValueQueue errors consumer with multiple-reads") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer(queue); MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }, 3); MustNotCall readContinuation; read(js, consumer).then(js, readContinuation, errorContinuation); read(js, consumer).then(js, readContinuation, errorContinuation); read(js, consumer).then(js, readContinuation, errorContinuation); queue.error(js, js.v8Ref(js.v8Error("boom"_kj))); js.runMicrotasks(); }); } KJ_TEST("ValueQueue with multiple consumers with pending reads") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer1(queue); ValueQueue::Consumer consumer2(queue); KJ_ASSERT(queue.desiredSize() == 2); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsTrue()); // Both reads were fulfilled immediately without buffering. KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 0); // Backpressure is not signalled since both consumer reads have been // fulfilled. KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }, 2); read(js, consumer1).then(js, readContinuation); read(js, consumer2).then(js, readContinuation); queue.push(js, getEntry(js, 2)); js.runMicrotasks(); }); } #pragma endregion ValueQueue Tests #pragma region ByteQueue Tests KJ_TEST("ByteQueue basics work") { preamble([](jsg::Lock& js) { ByteQueue queue(2); // At this point, there are no consumers, data does not get enqueued. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); auto entry = kj::rc(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4))); queue.push(js, kj::mv(entry)); // Because there are no consumers, there is no change to backpressure. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); // Closing the queue causes the desiredSize to be zero. queue.close(js); try { auto entry = kj::rc(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4))); queue.push(js, kj::mv(entry)); KJ_FAIL_ASSERT("The queue push after close should have failed."); } catch (kj::Exception& ex) { KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored.")); } KJ_ASSERT(queue.desiredSize() == 0); KJ_ASSERT(queue.size() == 0); }); } KJ_TEST("ByteQueue erroring works") { preamble([](jsg::Lock& js) { ByteQueue queue(2); queue.error(js, js.v8Ref(js.v8Error("boom"_kj))); KJ_ASSERT(queue.desiredSize() == 0); try { auto entry = kj::rc(jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4))); queue.push(js, kj::mv(entry)); KJ_FAIL_ASSERT("The queue push after close should have failed."); } catch (kj::Exception& ex) { KJ_ASSERT(ex.getDescription().endsWith("The queue is closed or errored.")); } }); } KJ_TEST("ByteQueue with single consumer") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer(queue); KJ_ASSERT(queue.desiredSize() == 2); auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr().fill('a'); auto entry = kj::rc(jsg::BufferSource(js, kj::mv(store))); queue.push(js, kj::mv(entry)); // The item was pushed into the consumer. KJ_ASSERT(consumer.size() == 4); // The queue size and desiredSize were updated accordingly. KJ_ASSERT(queue.size() == 4); KJ_ASSERT(queue.desiredSize() == -2); auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), })); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); KJ_ASSERT(source.size() == 4); KJ_ASSERT(source.asArrayPtr()[0] == 'a'); KJ_ASSERT(source.asArrayPtr()[1] == 'a'); KJ_ASSERT(source.asArrayPtr()[2] == 'a'); KJ_ASSERT(source.asArrayPtr()[3] == 'a'); KJ_ASSERT(consumer.size() == 0); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }); prp.promise.then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with single byob consumer") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer(queue); auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), .type = ByteQueue::ReadRequest::Type::BYOB, })); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); KJ_ASSERT(ptr[2] == 'b'); KJ_ASSERT(consumer.size() == 0); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }); prp.promise.then(js, readContinuation); auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); KJ_ASSERT(!pendingByob->isInvalidated()); auto& req = pendingByob->getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); ptr.first(3).fill('b'); pendingByob->respond(js, 3); KJ_ASSERT(pendingByob->isInvalidated()); // No backpressure is signaled. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(consumer.size() == 0); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with byob consumer and default consumer") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); auto prp = js.newPromiseAndResolver(); consumer1.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), .type = ByteQueue::ReadRequest::Type::BYOB, })); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); KJ_ASSERT(ptr[2] == 'b'); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 3); KJ_ASSERT(queue.size() == 3); KJ_ASSERT(queue.desiredSize() == -1); return js.resolvedPromise(kj::mv(result)); }); prp.promise.then(js, readContinuation); auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); KJ_ASSERT(!pendingByob->isInvalidated()); auto& req = pendingByob->getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); ptr.first(3).fill('b'); pendingByob->respond(js, 3); KJ_ASSERT(pendingByob->isInvalidated()); // Backpressure is signaled because the other consumer hasn't been read from. KJ_ASSERT(queue.desiredSize() == -1); KJ_ASSERT(queue.size() == 3); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 3); js.runMicrotasks(); MustCall read2Continuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); // The second consumer receives exactly the same data. KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); KJ_ASSERT(ptr[2] == 'b'); // The backpressure in the queue has been resolved. KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }); auto prp2 = js.newPromiseAndResolver(); consumer2.read(js, ByteQueue::ReadRequest(kj::mv(prp2.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), .type = ByteQueue::ReadRequest::Type::DEFAULT, })); prp2.promise.then(js, read2Continuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple byob consumers") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); KJ_ASSERT(ptr[2] == 'b'); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 0); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }, 2); // Both reads will receive the data despite there being only a single // byob read responded to. byobRead(js, consumer1, 4).then(js, readContinuation); byobRead(js, consumer2, 4).then(js, readContinuation); auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); auto nextPending = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); KJ_ASSERT(!pendingByob->isInvalidated()); auto& req = pendingByob->getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); ptr.first(3).fill('b'); pendingByob->respond(js, 3); KJ_ASSERT(pendingByob->isInvalidated()); // No backpressure is signaled because both reads were fulfilled. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 0); // The next pendingByobReadRequest was invalidated. KJ_ASSERT(nextPending->isInvalidated()); KJ_ASSERT(queue.nextPendingByobReadRequest() == kj::none); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple byob consumers") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); MustCall readContinuation([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); KJ_ASSERT(ptr[2] == 'b'); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 0); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(queue.desiredSize() == 2); return js.resolvedPromise(kj::mv(result)); }, 2); // Both reads will receive the data despite there being only a single // byob read responded to. byobRead(js, consumer1, 4).then(js, readContinuation); byobRead(js, consumer2, 4).then(js, readContinuation); auto pendingByob = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); auto nextPending = KJ_ASSERT_NONNULL(queue.nextPendingByobReadRequest()); KJ_ASSERT(!pendingByob->isInvalidated()); auto& req = pendingByob->getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); ptr.first(3).fill('b'); pendingByob->respond(js, 3); KJ_ASSERT(pendingByob->isInvalidated()); // No backpressure is signaled because both reads were fulfilled. KJ_ASSERT(queue.desiredSize() == 2); KJ_ASSERT(queue.size() == 0); KJ_ASSERT(consumer1.size() == 0); KJ_ASSERT(consumer2.size() == 0); // The next pendingByobReadRequest was invalidated. KJ_ASSERT(nextPending->isInvalidated()); KJ_ASSERT(queue.nextPendingByobReadRequest() == kj::none); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple byob consumers (multi-reads)") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); MustCall readConsumer1([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'a'); KJ_ASSERT(ptr[1] == 'a'); KJ_ASSERT(ptr[2] == 'a'); return js.resolvedPromise(kj::mv(result)); }); MustCall readConsumer2([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'a'); KJ_ASSERT(ptr[1] == 'a'); KJ_ASSERT(ptr[2] == 'a'); return byobRead(js, consumer2, 4); }); MustCall secondReadBothConsumers([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 2); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); return js.resolvedPromise(kj::mv(result)); }, 2); // All reads will be fulfilled correctly even tho there are only two byob // reads processed. byobRead(js, consumer1, 4).then(js, readConsumer1); byobRead(js, consumer1, 4).then(js, secondReadBothConsumers); byobRead(js, consumer2, 4).then(js, readConsumer2).then(js, secondReadBothConsumers); // Although there are four distinct reads happening, // there should only be two actual BYOB requests // processed by the queue, which will fulfill all four // reads. MustCall respond([&](jsg::Lock&, auto& pending) { static uint counter = 0; auto& req = pending.getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); auto num = 3 - counter; ptr.first(num).fill('a' + counter++); pending.respond(js, num); KJ_ASSERT(pending.isInvalidated()); }, 2); kj::Maybe> pendingByob; while ((pendingByob = queue.nextPendingByobReadRequest()) != kj::none) { auto& pending = KJ_ASSERT_NONNULL(pendingByob); if (pending->isInvalidated()) { continue; } respond(js, *pending); } js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple byob consumers (multi-reads, 2)") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); MustCall readConsumer1([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'a'); KJ_ASSERT(ptr[1] == 'a'); KJ_ASSERT(ptr[2] == 'a'); return js.resolvedPromise(kj::mv(result)); }); MustCall readConsumer2([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 3); KJ_ASSERT(ptr[0] == 'a'); KJ_ASSERT(ptr[1] == 'a'); KJ_ASSERT(ptr[2] == 'a'); return byobRead(js, consumer2, 4); }); MustCall secondReadBothConsumers([&](jsg::Lock& js, auto&& result) -> auto { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); KJ_ASSERT(value.getHandle(js)->IsArrayBufferView()); jsg::BufferSource source(js, value.getHandle(js)); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 2); KJ_ASSERT(ptr[0] == 'b'); KJ_ASSERT(ptr[1] == 'b'); return js.resolvedPromise(kj::mv(result)); }, 2); // All reads will be fulfilled correctly even tho there are only two BYOB reads // responded to. byobRead(js, consumer2, 4).then(js, readConsumer2).then(js, secondReadBothConsumers); byobRead(js, consumer1, 4).then(js, readConsumer1); byobRead(js, consumer1, 4).then(js, secondReadBothConsumers); // Although there are four distinct reads happening, // there should only be two actual BYOB requests // processed by the queue, which will fulfill all four // reads. MustCall respond([&](jsg::Lock&, auto& pending) { static uint counter = 0; auto& req = pending.getRequest(); auto ptr = req.pullInto.store.asArrayPtr(); auto num = 3 - counter; ptr.first(num).fill('a' + counter++); pending.respond(js, num); KJ_ASSERT(pending.isInvalidated()); }, 2); kj::Maybe> pendingByob; while ((pendingByob = queue.nextPendingByobReadRequest()) != kj::none) { auto& pending = KJ_ASSERT_NONNULL(pendingByob); if (pending->isInvalidated()) { continue; } respond(js, *pending); } js.runMicrotasks(); }); } KJ_TEST("ByteQueue with default consumer with atLeast") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer(queue); const auto read = [&](jsg::Lock& js, uint atLeast) { auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)), .atLeast = atLeast, })); return kj::mv(prp.promise); }; const auto push = [&](auto store) { try { queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); } catch (kj::Exception& ex) { KJ_DBG(ex.getDescription()); } }; MustCall readContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); auto ptr = source.asArrayPtr(); KJ_ASSERT(ptr[0] == 1); KJ_ASSERT(ptr[1] == 2); KJ_ASSERT(ptr[2] == 3); KJ_ASSERT(ptr[3] == 4); KJ_ASSERT(ptr[4] == 5); KJ_ASSERT(source.size(), 5); KJ_ASSERT(consumer.size(), 1); return read(js, 1); }); MustCall read2Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); KJ_ASSERT(source.asArrayPtr()[0], 6); KJ_ASSERT(source.size() == 1); return js.resolvedPromise(kj::mv(result)); }); read(js, 5).then(js, readContinuation).then(js, read2Continuation); auto store1 = jsg::BackingStore::alloc(js, 2); store1.asArrayPtr()[0] = 1; store1.asArrayPtr()[1] = 2; push(kj::mv(store1)); KJ_ASSERT(queue.desiredSize() == 0); auto store2 = jsg::BackingStore::alloc(js, 2); store2.asArrayPtr()[0] = 3; store2.asArrayPtr()[1] = 4; push(kj::mv(store2)); // Backpressure should be accumulating because the read has not yet fullilled. KJ_ASSERT(queue.desiredSize() == -2); auto store3 = jsg::BackingStore::alloc(js, 2); store3.asArrayPtr()[0] = 5; store3.asArrayPtr()[1] = 6; push(kj::mv(store3)); // Some backpressure should be released because pushing the final minimum // amount into the queue should have caused the read to be fulfilled. KJ_ASSERT(queue.desiredSize() == 1); // There should be one unread byte left in the queue at this point. // It will be read once the microtask queue is drained. KJ_ASSERT(queue.size() == 1); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple default consumers with atLeast (same rate)") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); const auto read = [&](jsg::Lock& js, auto& consumer, uint atLeast = 1) { auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)), .atLeast = atLeast, })); return kj::mv(prp.promise); }; const auto push = [&](auto store) { try { queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); } catch (kj::Exception& ex) { KJ_DBG(ex.getDescription()); } }; MustCall read1Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); auto ptr = source.asArrayPtr(); KJ_ASSERT(ptr[0] == 1); KJ_ASSERT(ptr[1] == 2); KJ_ASSERT(ptr[2] == 3); KJ_ASSERT(ptr[3] == 4); KJ_ASSERT(ptr[4] == 5); KJ_ASSERT(source.size(), 5); KJ_ASSERT(consumer1.size(), 1); return read(js, consumer1); }); MustCall read2Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); auto ptr = source.asArrayPtr(); KJ_ASSERT(ptr[0] == 1); KJ_ASSERT(ptr[1] == 2); KJ_ASSERT(ptr[2] == 3); KJ_ASSERT(ptr[3] == 4); KJ_ASSERT(ptr[4] == 5); KJ_ASSERT(source.size(), 5); KJ_ASSERT(consumer2.size(), 1); return read(js, consumer2); }); MustCall readFinalContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); KJ_ASSERT(source.asArrayPtr()[0], 6); KJ_ASSERT(source.size() == 1); return js.resolvedPromise(kj::mv(result)); }, 2); read(js, consumer1, 5).then(js, read1Continuation).then(js, readFinalContinuation); read(js, consumer2, 5).then(js, read2Continuation).then(js, readFinalContinuation); auto store1 = jsg::BackingStore::alloc(js, 2); store1.asArrayPtr()[0] = 1; store1.asArrayPtr()[1] = 2; push(kj::mv(store1)); KJ_ASSERT(queue.desiredSize() == 0); auto store2 = jsg::BackingStore::alloc(js, 2); store2.asArrayPtr()[0] = 3; store2.asArrayPtr()[1] = 4; push(kj::mv(store2)); // Backpressure should be accumulating because the read has not yet fullilled. KJ_ASSERT(queue.desiredSize() == -2); auto store3 = jsg::BackingStore::alloc(js, 2); store3.asArrayPtr()[0] = 5; store3.asArrayPtr()[1] = 6; push(kj::mv(store3)); // Some backpressure should be released because pushing the final minimum // amount into the queue should have caused the read to be fulfilled. KJ_ASSERT(queue.desiredSize() == 1); // There should be one unread byte left in the queue at this point. // It will be read once the microtask queue is drained. KJ_ASSERT(queue.size() == 1); js.runMicrotasks(); }); } KJ_TEST("ByteQueue with multiple default consumers with atLeast (different rate)") { preamble([](jsg::Lock& js) { ByteQueue queue(2); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); const auto read = [&](jsg::Lock& js, auto& consumer, uint atLeast = 1) { auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 5)), .atLeast = atLeast, })); return kj::mv(prp.promise); }; const auto push = [&](auto store) { try { queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); } catch (kj::Exception& ex) { KJ_DBG(ex.getDescription()); } }; MustCall read1Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); KJ_ASSERT(source.size() == 4); auto ptr = source.asArrayPtr(); // Our read was for at least 3 bytes, with a maximum of 5. // For this first read, we received 4. One the second read // we should receive 2. KJ_ASSERT(ptr[0] == 1); KJ_ASSERT(ptr[1] == 2); KJ_ASSERT(ptr[2] == 3); KJ_ASSERT(ptr[3] == 4); return js.resolvedPromise(kj::mv(result)); }); MustCall read1FinalContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); KJ_ASSERT(source.size() == 2); auto ptr = source.asArrayPtr(); KJ_ASSERT(ptr[0] == 5); KJ_ASSERT(ptr[1] == 6); return js.resolvedPromise(kj::mv(result)); }); MustCall read2Continuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); auto ptr = source.asArrayPtr(); KJ_ASSERT(source.size() == 5); KJ_ASSERT(ptr[0] == 1); KJ_ASSERT(ptr[1] == 2); KJ_ASSERT(ptr[2] == 3); KJ_ASSERT(ptr[3] == 4); KJ_ASSERT(ptr[4] == 5); KJ_ASSERT(consumer2.size() == 1); return read(js, consumer2); }); MustCall read2FinalContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); auto& value = KJ_ASSERT_NONNULL(result.value); auto view = value.getHandle(js); KJ_ASSERT(view->IsArrayBufferView()); jsg::BufferSource source(js, view); KJ_ASSERT(source.asArrayPtr()[0] == 6); KJ_ASSERT(source.size() == 1); return js.resolvedPromise(kj::mv(result)); }); // Consumer 1 will read in parallel with smaller minimum chunks... read(js, consumer1, 3).then(js, read1Continuation); read(js, consumer1).then(js, read1FinalContinuation); // Consumer 2 will read serially with a larger minimum chunk... read(js, consumer2, 5).then(js, read2Continuation).then(js, read2FinalContinuation); auto store1 = jsg::BackingStore::alloc(js, 2); store1.asArrayPtr()[0] = 1; store1.asArrayPtr()[1] = 2; push(kj::mv(store1)); KJ_ASSERT(queue.desiredSize() == 0); auto store2 = jsg::BackingStore::alloc(js, 2); store2.asArrayPtr()[0] = 3; store2.asArrayPtr()[1] = 4; push(kj::mv(store2)); // Consumer1 should not have any data buffered since its first read was for // between 3 and 5 bytes and it has received four so far. KJ_ASSERT(consumer1.size() == 0); // Consumer2 should have 4 bytes buffered since its first read was for 5 bytes // and we've only received 4 so far. KJ_ASSERT(consumer2.size() == 4); // Queue backpressure should reflect that consumer2 has data buffered. KJ_ASSERT(queue.desiredSize() == -2); auto store3 = jsg::BackingStore::alloc(js, 2); store3.asArrayPtr()[0] = 5; store3.asArrayPtr()[1] = 6; push(kj::mv(store3)); // Most of the backpressure should have been resolved since we delivered 5 bytes // to consumer2, but there's still one byte remaining. KJ_ASSERT(queue.desiredSize() = 1); KJ_ASSERT(queue.size() == 1); js.runMicrotasks(); }); } #pragma endregion ByteQueue Tests #pragma region Re-entrancy Safety Tests // Test that pushing to a closed consumer doesn't crash. // This can happen during QueueImpl::push() iteration when resolving a read // on one consumer triggers JavaScript that closes another consumer. KJ_TEST("ValueQueue push to closed consumer is safe") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer1(queue); ValueQueue::Consumer consumer2(queue); // Close consumer2 consumer2.close(js); // Now push to the queue - this pushes to all consumers // Before the fix, this would crash when trying to push to closed consumer2 queue.push(js, getEntry(js, 4)); // consumer1 should have received the data KJ_ASSERT(consumer1.size() == 4); js.runMicrotasks(); }); } // Test that pushing to a cancelled consumer doesn't crash. KJ_TEST("ValueQueue push to cancelled consumer is safe") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer1(queue); ValueQueue::Consumer consumer2(queue); // Cancel consumer2 consumer2.cancel(js, kj::none); // Now push to the queue queue.push(js, getEntry(js, 4)); // consumer1 should have received the data KJ_ASSERT(consumer1.size() == 4); js.runMicrotasks(); }); } // Test that pushing to an errored consumer doesn't crash. KJ_TEST("ValueQueue push to errored consumer is safe") { preamble([](jsg::Lock& js) { ValueQueue queue(2); ValueQueue::Consumer consumer1(queue); ValueQueue::Consumer consumer2(queue); // Error consumer2 consumer2.error(js, js.v8Ref(js.v8Error("error reason"_kj))); // Now push to the queue queue.push(js, getEntry(js, 4)); // consumer1 should have received the data KJ_ASSERT(consumer1.size() == 4); js.runMicrotasks(); }); } // Test ByteQueue version of the safety checks KJ_TEST("ByteQueue push to closed consumer is safe") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer1(queue); ByteQueue::Consumer consumer2(queue); // Close consumer2 consumer2.close(js); // Now push to the queue auto store = jsg::BackingStore::alloc(js, 4); memset(store.asArrayPtr().begin(), 'A', 4); auto entry = kj::rc(jsg::BufferSource(js, kj::mv(store))); queue.push(js, kj::mv(entry)); // consumer1 should have received the data KJ_ASSERT(consumer1.size() == 4); js.runMicrotasks(); }); } #pragma endregion Re - entrancy Safety Tests #pragma region Draining Read Tests using DrainingReadContinuation = jsg::Promise(DrainingReadResult&&); using DrainingReadErrorContinuation = jsg::Promise(jsg::Value&&); KJ_TEST("ValueQueue draining read with buffered data") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Push an ArrayBuffer auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr()[0] = 'a'; store.asArrayPtr()[1] = 'b'; store.asArrayPtr()[2] = 'c'; store.asArrayPtr()[3] = 'd'; auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab.As()), 4)); // Push a string auto str = jsg::v8Str(js.v8Isolate, "hello"); queue.push(js, kj::rc(js.v8Ref(str.As()), 5)); KJ_ASSERT(consumer.size() == 9); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 2); // First chunk is the ArrayBuffer data KJ_ASSERT(result.chunks[0].size() == 4); KJ_ASSERT(result.chunks[0][0] == 'a'); KJ_ASSERT(result.chunks[0][1] == 'b'); KJ_ASSERT(result.chunks[0][2] == 'c'); KJ_ASSERT(result.chunks[0][3] == 'd'); // Second chunk is the string converted to UTF-8 KJ_ASSERT(result.chunks[1].size() == 5); KJ_ASSERT(result.chunks[1][0] == 'h'); KJ_ASSERT(result.chunks[1][1] == 'e'); KJ_ASSERT(result.chunks[1][2] == 'l'); KJ_ASSERT(result.chunks[1][3] == 'l'); KJ_ASSERT(result.chunks[1][4] == 'o'); KJ_ASSERT(consumer.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read rejects with pending reads") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Queue a regular read auto prp = js.newPromiseAndResolver(); consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)}); KJ_ASSERT(consumer.hasReadRequests()); // Draining read should reject because there are pending reads MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue read rejects with pending draining read") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // No data in buffer, draining read will queue a pending draining read consumer.drainingRead(js); KJ_ASSERT(consumer.hasPendingDrainingRead()); // Regular read should reject because there's a pending draining read auto prp = js.newPromiseAndResolver(); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)}); prp.promise.then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read on closed stream") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); queue.close(js); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(result.done); KJ_ASSERT(result.chunks.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read on errored stream") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); queue.error(js, js.v8Ref(js.v8Error("boom"_kj))); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read with buffered data") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); // Push first chunk auto store1 = jsg::BackingStore::alloc(js, 4); store1.asArrayPtr()[0] = 'a'; store1.asArrayPtr()[1] = 'b'; store1.asArrayPtr()[2] = 'c'; store1.asArrayPtr()[3] = 'd'; queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store1)))); // Push second chunk auto store2 = jsg::BackingStore::alloc(js, 3); store2.asArrayPtr()[0] = 'e'; store2.asArrayPtr()[1] = 'f'; store2.asArrayPtr()[2] = 'g'; queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store2)))); KJ_ASSERT(consumer.size() == 7); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 2); // First chunk KJ_ASSERT(result.chunks[0].size() == 4); KJ_ASSERT(result.chunks[0][0] == 'a'); KJ_ASSERT(result.chunks[0][1] == 'b'); KJ_ASSERT(result.chunks[0][2] == 'c'); KJ_ASSERT(result.chunks[0][3] == 'd'); // Second chunk KJ_ASSERT(result.chunks[1].size() == 3); KJ_ASSERT(result.chunks[1][0] == 'e'); KJ_ASSERT(result.chunks[1][1] == 'f'); KJ_ASSERT(result.chunks[1][2] == 'g'); KJ_ASSERT(consumer.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read rejects with pending reads") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); // Queue a regular read auto prp = js.newPromiseAndResolver(); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), })); KJ_ASSERT(consumer.hasReadRequests()); // Draining read should reject because there are pending reads MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue read rejects with pending draining read") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); // No data in buffer, draining read will queue a pending draining read consumer.drainingRead(js); KJ_ASSERT(consumer.hasPendingDrainingRead()); // Regular read should reject because there's a pending draining read auto prp = js.newPromiseAndResolver(); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), })); prp.promise.then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read on closed stream") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); queue.close(js); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { KJ_ASSERT(result.done); KJ_ASSERT(result.chunks.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read on errored stream") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); queue.error(js, js.v8Ref(js.v8Error("boom"_kj))); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read with close signal") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Push some data auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr()[0] = 'a'; store.asArrayPtr()[1] = 'b'; store.asArrayPtr()[2] = 'c'; store.asArrayPtr()[3] = 'd'; auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab.As()), 4)); // Close the queue queue.close(js); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { // Should have the data and done should be true since stream is closed KJ_ASSERT(result.done); KJ_ASSERT(result.chunks.size() == 1); KJ_ASSERT(result.chunks[0].size() == 4); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read with close signal") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); // Push some data auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr()[0] = 'a'; store.asArrayPtr()[1] = 'b'; store.asArrayPtr()[2] = 'c'; store.asArrayPtr()[3] = 'd'; queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); // Close the queue queue.close(js); MustCall readContinuation([&](jsg::Lock& js, auto&& result) { // Should have the data and done should be true since stream is closed KJ_ASSERT(result.done); KJ_ASSERT(result.chunks.size() == 1); KJ_ASSERT(result.chunks[0].size() == 4); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read errors on non-byte value") { // Test that drainingRead correctly errors when encountering a value that // cannot be converted to bytes (not ArrayBuffer, ArrayBufferView, or string). preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Push a plain object - this cannot be converted to bytes auto obj = v8::Object::New(js.v8Isolate); queue.push(js, kj::rc(js.v8Ref(obj.As()), 1)); KJ_ASSERT(consumer.size() == 1); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { // Should get a TypeError about non-convertible value KJ_ASSERT(value.getHandle(js)->IsNativeError()); auto message = jsg::JsValue(value.getHandle(js)) .tryCast() .map([&](jsg::JsObject obj) { return obj.get(js, "message") .tryCast() .map([&](jsg::JsString str) { return str.toString(js); }).orDefault(kj::str()); }).orDefault(kj::str()); KJ_ASSERT(message.contains("cannot be converted to bytes")); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); // MustCall verifies the error continuation was called }); } KJ_TEST("ValueQueue draining read errors on number value") { // Another non-byte value test with a number preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Push a number - this cannot be converted to bytes auto num = v8::Number::New(js.v8Isolate, 42); queue.push(js, kj::rc(js.v8Ref(num.As()), 1)); MustNotCall readContinuation; MustCall errorContinuation([&](jsg::Lock& js, auto&& value) { KJ_ASSERT(value.getHandle(js)->IsNativeError()); return js.rejectedPromise(kj::mv(value)); }); consumer.drainingRead(js).then(js, readContinuation, errorContinuation); js.runMicrotasks(); // MustCall verifies the error continuation was called }); } #pragma endregion Draining Read Tests #pragma region Draining Read maxRead Tests // Tests for the maxRead soft limit parameter. Both the initial buffer drain and subsequent // synchronous pump attempts stop when totalRead reaches maxRead. This prevents unbounded // memory accumulation when a fast producer outpaces a slow consumer. // // Note: Testing the pump loop behavior requires the full controller infrastructure (stateListener). // These tests focus on the buffer drain behavior which can be tested without a listener. KJ_TEST("ValueQueue draining read respects maxRead during buffer drain") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Buffer 200 bytes of data (two 100-byte chunks) auto store1 = jsg::BackingStore::alloc(js, 100); store1.asArrayPtr().fill(0xAA); auto ab1 = jsg::BufferSource(js, kj::mv(store1)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab1.As()), 100)); auto store2 = jsg::BackingStore::alloc(js, 100); store2.asArrayPtr().fill(0xBB); auto ab2 = jsg::BufferSource(js, kj::mv(store2)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab2.As()), 100)); KJ_ASSERT(consumer.size() == 200); // maxRead=50: the first 100-byte chunk is drained (exceeding maxRead after the first item), // then draining stops because totalRead (100) >= maxRead (50). The second chunk stays buffered. MustCall readContinuation( [&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); // Should only have the first chunk (maxRead stopped draining before the second) KJ_ASSERT(result.chunks.size() == 1); KJ_ASSERT(result.chunks[0].size() == 100); // Second chunk should still be buffered KJ_ASSERT(consumer.size() == 100); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js, 50).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ByteQueue draining read respects maxRead during buffer drain") { preamble([](jsg::Lock& js) { ByteQueue queue(10); ByteQueue::Consumer consumer(queue); // Buffer 200 bytes of data (two 100-byte chunks) auto store1 = jsg::BackingStore::alloc(js, 100); store1.asArrayPtr().fill(0xAA); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store1)))); auto store2 = jsg::BackingStore::alloc(js, 100); store2.asArrayPtr().fill(0xBB); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store2)))); KJ_ASSERT(consumer.size() == 200); // maxRead=50: first 100-byte chunk is drained, then stops. Second chunk stays buffered. MustCall readContinuation( [&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 1); KJ_ASSERT(result.chunks[0].size() == 100); KJ_ASSERT(consumer.size() == 100); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js, 50).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read with large maxRead drains entire buffer") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Buffer 200 bytes (two 100-byte chunks) auto store1 = jsg::BackingStore::alloc(js, 100); store1.asArrayPtr().fill(0xAA); auto ab1 = jsg::BufferSource(js, kj::mv(store1)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab1.As()), 100)); auto store2 = jsg::BackingStore::alloc(js, 100); store2.asArrayPtr().fill(0xBB); auto ab2 = jsg::BufferSource(js, kj::mv(store2)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab2.As()), 100)); KJ_ASSERT(consumer.size() == 200); // maxRead=1000: both chunks should be drained since total (200) < maxRead (1000) MustCall readContinuation( [&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 2); KJ_ASSERT(result.chunks[0].size() == 100); KJ_ASSERT(result.chunks[1].size() == 100); KJ_ASSERT(consumer.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js, 1000).then(js, readContinuation); js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read with default maxRead (unlimited)") { preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Buffer some data auto store = jsg::BackingStore::alloc(js, 100); store.asArrayPtr().fill(0xAA); auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab.As()), 100)); // Default maxRead (kj::maxValue) should drain buffer normally MustCall readContinuation( [&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 1); KJ_ASSERT(result.chunks[0].size() == 100); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js).then(js, readContinuation); // No maxRead argument = default js.runMicrotasks(); }); } KJ_TEST("ValueQueue draining read maxRead bounds multiple iterations") { // Verify that successive drainingRead calls with maxRead correctly drain // a large buffer incrementally. preamble([](jsg::Lock& js) { ValueQueue queue(10); ValueQueue::Consumer consumer(queue); // Buffer 400 bytes: four 100-byte chunks for (int i = 0; i < 4; i++) { auto store = jsg::BackingStore::alloc(js, 100); store.asArrayPtr().fill(0x10 * (i + 1)); auto ab = jsg::BufferSource(js, kj::mv(store)).getHandle(js); queue.push(js, kj::rc(js.v8Ref(ab.As()), 100)); } KJ_ASSERT(consumer.size() == 400); // First read with maxRead=150: drains first chunk (100 bytes, now totalRead=100 < 150), // then drains second chunk (200 bytes total, now >= 150), stops. MustCall read1([&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 2); KJ_ASSERT(consumer.size() == 200); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js, 150).then(js, read1); js.runMicrotasks(); // Second read with maxRead=150: drains next two chunks similarly MustCall read2([&](jsg::Lock& js, DrainingReadResult&& result) { KJ_ASSERT(!result.done); KJ_ASSERT(result.chunks.size() == 2); KJ_ASSERT(consumer.size() == 0); return js.resolvedPromise(kj::mv(result)); }); consumer.drainingRead(js, 150).then(js, read2); js.runMicrotasks(); }); } #pragma endregion Draining Read maxRead Tests #pragma region Queue/Consumer Destruction Order Tests // These tests verify that destroying the queue before its consumers doesn't crash. // This can happen during isolate teardown when the destruction order isn't guaranteed // to follow the ownership hierarchy. // These will typically only catch in builds with AddressSanitizer enabled. KJ_TEST("ValueQueue destroyed before consumer doesn't crash") { preamble([](jsg::Lock& js) { // Heap-allocate the queue so we can control its destruction order auto queue = kj::heap(2); // Create a consumer attached to the queue auto consumer = kj::heap(*queue); // Push some data to make sure the consumer has state queue->push(js, getEntry(js, 4)); KJ_ASSERT(consumer->size() == 4); // Now destroy the queue FIRST - this simulates the production scenario // where wrapper cleanup destroys the controller (and its queue) before // the consumer that holds a reference to it. queue = nullptr; // When the consumer is destroyed (here, or when going out of scope), // its destructor calls queue.removeConsumer(this). // Without the fix, this is a use-after-free. // With the fix, the consumer knows the queue is gone and skips the call. consumer = nullptr; // If we get here without crashing, the test passes }); } KJ_TEST("ValueQueue destroyed before multiple consumers doesn't crash") { preamble([](jsg::Lock& js) { auto queue = kj::heap(2); auto consumer1 = kj::heap(*queue); auto consumer2 = kj::heap(*queue); queue->push(js, getEntry(js, 4)); KJ_ASSERT(consumer1->size() == 4); KJ_ASSERT(consumer2->size() == 4); // Destroy queue before consumers queue = nullptr; // Both consumers should handle destruction gracefully consumer1 = nullptr; consumer2 = nullptr; }); } KJ_TEST("ByteQueue destroyed before consumer doesn't crash") { preamble([](jsg::Lock& js) { auto queue = kj::heap(2); auto consumer = kj::heap(*queue); auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr().fill('a'); queue->push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); KJ_ASSERT(consumer->size() == 4); // Destroy queue before consumer queue = nullptr; consumer = nullptr; }); } KJ_TEST("ValueQueue destroyed with pending read requests doesn't crash") { preamble([](jsg::Lock& js) { auto queue = kj::heap(2); auto consumer = kj::heap(*queue); // Queue a read request (no data pushed, so it will be pending) auto prp = js.newPromiseAndResolver(); consumer->read(js, ValueQueue::ReadRequest{.resolver = kj::mv(prp.resolver)}); KJ_ASSERT(consumer->hasReadRequests()); // Destroy queue while there are pending reads queue = nullptr; // Consumer destruction should handle this gracefully consumer = nullptr; js.runMicrotasks(); }); } KJ_TEST("ValueQueue close then destroy before consumer doesn't crash") { preamble([](jsg::Lock& js) { auto queue = kj::heap(2); auto consumer = kj::heap(*queue); // Close the queue first queue->close(js); // Then destroy it queue = nullptr; // Consumer should still handle destruction gracefully consumer = nullptr; }); } KJ_TEST("ValueQueue error then destroy before consumer doesn't crash") { preamble([](jsg::Lock& js) { auto queue = kj::heap(2); auto consumer = kj::heap(*queue); // Error the queue first queue->error(js, js.v8Ref(js.v8Error("boom"_kj))); // Then destroy it queue = nullptr; // Consumer should still handle destruction gracefully consumer = nullptr; }); } #pragma endregion Queue / Consumer Destruction Order Tests #pragma region Consumer Destroyed During Push Tests // These tests verify that the queue handles consumer destruction gracefully. // QueueImpl::push() takes a snapshot of consumers and iterates over it. If a consumer // is destroyed (removed from allConsumers) between the snapshot and the iteration, // the code must check allConsumers.contains() before dereferencing the pointer. // That said, the tests do not actually fail without the relevant fix because the // issue is extremely timing-dependent and difficult to trigger deterministically, even // with asan enabled. These tests at least document the intended behavior and may catch // future regressions. KJ_TEST("ByteQueue push skips consumer removed from queue during iteration") { preamble([](jsg::Lock& js) { ByteQueue queue(10); // Create two consumers auto consumer1 = kj::heap(queue); auto consumer2 = kj::heap(queue); // Destroy consumer2 BEFORE pushing. This directly tests that push() // checks if consumers still exist before calling push on them. // The snapshot taken at the start of push() would have included consumer2, // but consumer2 is no longer in allConsumers when we iterate. consumer2 = nullptr; // Push data - should not crash even though consumer2 was in the queue // when it was created but is now destroyed. auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr().fill('x'); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); // consumer1 should have received the data KJ_ASSERT(consumer1->size() == 4); }); } KJ_TEST("ValueQueue push skips consumer removed from queue during iteration") { preamble([](jsg::Lock& js) { ValueQueue queue(10); auto consumer1 = kj::heap(queue); auto consumer2 = kj::heap(queue); // Destroy consumer2 before pushing consumer2 = nullptr; queue.push(js, getEntry(js, 4)); KJ_ASSERT(consumer1->size() == 4); }); } KJ_TEST("ByteQueue push handles consumer destroyed by microtask between pushes") { preamble([](jsg::Lock& js) { ByteQueue queue(10); auto consumer1 = kj::heap(queue); auto consumer2 = kj::heap(queue); // Set up a pending read on consumer1 auto prp = js.newPromiseAndResolver(); consumer1->read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, jsg::BackingStore::alloc(js, 4)), })); // The continuation destroys consumer2 MustCall readContinuation([&consumer2](jsg::Lock& js, ReadResult&& result) { consumer2 = nullptr; return js.resolvedPromise(kj::mv(result)); }); prp.promise.then(js, readContinuation); // First push - resolves consumer1's read, schedules microtask that will destroy consumer2 auto store1 = jsg::BackingStore::alloc(js, 4); store1.asArrayPtr().fill('x'); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store1)))); // Run microtasks - this destroys consumer2 js.runMicrotasks(); // Second push - consumer2 is now destroyed, should not crash auto store2 = jsg::BackingStore::alloc(js, 4); store2.asArrayPtr().fill('y'); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store2)))); // consumer1 should have the second push's data buffered KJ_ASSERT(consumer1->size() == 4); }); } KJ_TEST("ByteQueue maybeUpdateBackpressure skips destroyed consumers") { preamble([](jsg::Lock& js) { ByteQueue queue(10); auto consumer1 = kj::heap(queue); auto consumer2 = kj::heap(queue); // Push some data so consumers have size auto store = jsg::BackingStore::alloc(js, 4); store.asArrayPtr().fill('x'); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store)))); KJ_ASSERT(consumer1->size() == 4); KJ_ASSERT(consumer2->size() == 4); KJ_ASSERT(queue.size() == 4); // Destroy consumer2 consumer2 = nullptr; // Trigger backpressure recalculation by pushing more data auto store2 = jsg::BackingStore::alloc(js, 4); store2.asArrayPtr().fill('y'); queue.push(js, kj::rc(jsg::BufferSource(js, kj::mv(store2)))); // Should not crash, and size should reflect only consumer1 KJ_ASSERT(consumer1->size() == 8); KJ_ASSERT(queue.size() == 8); }); } #pragma endregion Consumer Destroyed During Push Tests } // namespace } // namespace workerd::api