#include #include #include #include #include namespace workerd::api { namespace { class FakeStreamSource final: public ReadableStreamSource { public: FakeStreamSource(size_t length): length(length) {} kj::Promise tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { return kj::evalNow([this, maxBytes, buffer] { auto amount = kj::min(maxBytes, length); memset(buffer, 0, amount); length -= amount; return amount; }); } private: size_t length; }; KJ_TEST("Streams tee stack overflow regression") { // Verify that deeply nested tee() chains don't cause a stack overflow. // This is a regression test for a fix that removed deep recursion from tee(). static constexpr size_t teeDepth = 200 * 1024 / sizeof(void*); TestFixture testFixture; testFixture.runInIoContext([](const TestFixture::Environment& env) { auto& js = jsg::Lock::from(env.isolate); ReadableStream s(env.context, kj::heap(10 * 1024 * 1024)); auto readableStreams = s.tee(js); for (size_t i = 0; i < teeDepth; i++) { readableStreams = readableStreams[0]->tee(js); } }); } KJ_TEST("Reading from default reader") { static constexpr size_t streamLength = 10 * 1024; TestFixture testFixture; testFixture.runInIoContext([](const TestFixture::Environment& env) -> kj::Promise { auto& js = jsg::Lock::from(env.isolate); auto stream = js.alloc(env.context, kj::heap(streamLength)); auto reader = stream->getReader(js, {}); KJ_REQUIRE(reader.is>()); auto& defaultReader = reader.get>(); return env.context.awaitJs(js, defaultReader->read(js).then(js, JSG_VISITABLE_LAMBDA((reader = defaultReader.addRef(), stream = stream.addRef()), (reader, stream), (jsg::Lock& js, ReadResult readResult) { KJ_ASSERT(!readResult.done); auto& value = KJ_REQUIRE_NONNULL(readResult.value); auto handle = value.getHandle(js); KJ_ASSERT(handle->IsUint8Array()); if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) { // With 16KB buffer, the entire 10KB stream fits in one read. KJ_ASSERT(streamLength == handle.As()->ByteLength()); } else { KJ_ASSERT(4 * 1024 == handle.As()->ByteLength()); } }))); }); } KJ_TEST("Reading from byob reader") { TestFixture testFixture; struct TestData { size_t streamLength = 10 * 1024; size_t bufferSize; bool expectDone; }; TestData tests[] = { {.streamLength = 10 * 1024, .bufferSize = 100}, {.streamLength = 10 * 1024, .bufferSize = 100 * 1024}, {.streamLength = 10, .bufferSize = 100}, {.streamLength = 1024, .bufferSize = 1024}, }; for (auto test: tests) { testFixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise { auto& js = jsg::Lock::from(env.isolate); auto stream = js.alloc(env.context, kj::heap(test.streamLength)); ReadableStream::GetReaderOptions getReaderOptons = {.mode = kj::str("byob")}; auto reader = stream->getReader(js, kj::mv(getReaderOptons)); KJ_REQUIRE(reader.is>()); auto& byobReader = reader.get>(); auto buffer = v8::Uint8Array::New( v8::ArrayBuffer::New(js.v8Isolate, test.bufferSize), 0, test.bufferSize); return env.context.awaitJs(js, byobReader->read(js, buffer, {}).then(js, JSG_VISITABLE_LAMBDA( (test, reader = byobReader.addRef(), stream = stream.addRef()), (reader, stream), (jsg::Lock& js, ReadResult readResult) { KJ_ASSERT(!readResult.done); auto& value = KJ_REQUIRE_NONNULL(readResult.value); auto handle = value.getHandle(js); KJ_ASSERT(handle->IsUint8Array()); auto view = handle.As(); KJ_ASSERT(kj::min(test.streamLength, test.bufferSize) == view->ByteLength()); KJ_ASSERT(test.bufferSize == view->Buffer()->ByteLength()); }))); return kj::READY_NOW; }); } } KJ_TEST("PumpToReader regression") { // If the promise holding the PumpToReader is dropped while the inner // write to the sink is pending, the PumpToReader can free the sink. // In some cases, this means that the sink can error because shutdownWrite // is called while there is still a pending write promise. This test verifies // that PumpToReader cancels any pending write promise when it is destroyed. struct TestSink final: public WritableStreamSink { kj::TwoWayPipe pipe; kj::PromiseFulfillerPair paf; kj::Vector& events; TestSink(kj::Vector& events) : pipe(kj::newTwoWayPipe()), paf(kj::newPromiseAndFulfiller()), events(events) {} ~TestSink() { events.add(kj::str("sink was destroyed")); pipe.ends[0]->shutdownWrite(); } kj::Promise write(kj::ArrayPtr buffer) override { events.add(kj::str("got the write")); paf.fulfiller->fulfill(); return pipe.ends[0]->write(buffer).attach( kj::defer([this] { events.add(kj::str("write promise was dropped")); })); } kj::Promise write(kj::ArrayPtr> pieces) override { events.add(kj::str("got the write")); paf.fulfiller->fulfill(); // Concatenate pieces into a single buffer for the pipe write. kj::Vector data; for (auto& piece: pieces) { data.addAll(piece); } auto arr = data.releaseAsArray(); return pipe.ends[0]->write(arr).attach( kj::mv(arr), kj::defer([this] { events.add(kj::str("write promise was dropped")); })); } kj::Promise end() override { return kj::READY_NOW; } void abort(kj::Exception reason) override {} }; kj::Vector events; capnp::MallocMessageBuilder flagsBuilder; auto featureFlags = flagsBuilder.initRoot(); featureFlags.setStreamsJavaScriptControllers(true); TestFixture testFixture({.featureFlags = featureFlags.asReader()}); testFixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise { auto& js = jsg::Lock::from(env.isolate); auto stream = ReadableStream::constructor(js, UnderlyingSource{.start = [](jsg::Lock& js, auto controller) { auto& c = KJ_REQUIRE_NONNULL( controller.template tryGet>()); c->enqueue(js, v8::ArrayBuffer::New(js.v8Isolate, 10)); c->close(js); return js.resolvedPromise(); }}, kj::none); auto sink = kj::heap(events); auto writePromise = kj::mv(sink->paf.promise); auto promise = stream->pumpTo(js, kj::mv(sink), true); return writePromise.attach(kj::mv(promise)); }); KJ_ASSERT(events.size() == 3); KJ_ASSERT(events[0] == "got the write"); KJ_ASSERT(events[1] == "write promise was dropped"); KJ_ASSERT(events[2] == "sink was destroyed"); } } // namespace } // namespace workerd::api