File
Blob: src/workerd/api/streams-test.c++
| 1 | #include <workerd/api/streams/readable.h> |
| 2 | #include <workerd/api/streams/standard.h> |
| 3 | #include <workerd/tests/test-fixture.h> |
| 4 | #include <workerd/util/autogate.h> |
| 5 | |
| 6 | #include <kj/test.h> |
| 7 | |
| 8 | namespace workerd::api { |
| 9 | |
| 10 | namespace { |
| 11 | |
| 12 | class FakeStreamSource final: public ReadableStreamSource { |
| 13 | public: |
| 14 | FakeStreamSource(size_t length): length(length) {} |
| 15 | |
| 16 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 17 | return kj::evalNow([this, maxBytes, buffer] { |
| 18 | auto amount = kj::min(maxBytes, length); |
| 19 | memset(buffer, 0, amount); |
| 20 | length -= amount; |
| 21 | return amount; |
| 22 | }); |
| 23 | } |
| 24 | |
| 25 | private: |
| 26 | size_t length; |
| 27 | }; |
| 28 | |
| 29 | KJ_TEST("Streams tee stack overflow regression") { |
| 30 | // Verify that deeply nested tee() chains don't cause a stack overflow. |
| 31 | // This is a regression test for a fix that removed deep recursion from tee(). |
| 32 | static constexpr size_t teeDepth = 200 * 1024 / sizeof(void*); |
| 33 | TestFixture testFixture; |
| 34 | testFixture.runInIoContext([](const TestFixture::Environment& env) { |
| 35 | auto& js = jsg::Lock::from(env.isolate); |
| 36 | ReadableStream s(env.context, kj::heap<FakeStreamSource>(10 * 1024 * 1024)); |
| 37 | auto readableStreams = s.tee(js); |
| 38 | for (size_t i = 0; i < teeDepth; i++) { |
| 39 | readableStreams = readableStreams[0]->tee(js); |
| 40 | } |
| 41 | }); |
| 42 | } |
| 43 | |
| 44 | KJ_TEST("Reading from default reader") { |
| 45 | static constexpr size_t streamLength = 10 * 1024; |
| 46 | TestFixture testFixture; |
| 47 | |
| 48 | testFixture.runInIoContext([](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 49 | auto& js = jsg::Lock::from(env.isolate); |
| 50 | auto stream = js.alloc<ReadableStream>(env.context, kj::heap<FakeStreamSource>(streamLength)); |
| 51 | auto reader = stream->getReader(js, {}); |
| 52 | KJ_REQUIRE(reader.is<jsg::Ref<ReadableStreamDefaultReader>>()); |
| 53 | auto& defaultReader = reader.get<jsg::Ref<ReadableStreamDefaultReader>>(); |
| 54 | |
| 55 | return env.context.awaitJs(js, defaultReader->read(js).then(js, |
| 56 | JSG_VISITABLE_LAMBDA((reader = defaultReader.addRef(), stream = stream.addRef()), |
| 57 | (reader, stream), (jsg::Lock& js, ReadResult readResult) { |
| 58 | KJ_ASSERT(!readResult.done); |
| 59 | auto& value = KJ_REQUIRE_NONNULL(readResult.value); |
| 60 | auto handle = value.getHandle(js); |
| 61 | KJ_ASSERT(handle->IsUint8Array()); |
| 62 | if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) { |
| 63 | // With 16KB buffer, the entire 10KB stream fits in one read. |
| 64 | KJ_ASSERT(streamLength == handle.As<v8::Uint8Array>()->ByteLength()); |
| 65 | } else { |
| 66 | KJ_ASSERT(4 * 1024 == handle.As<v8::Uint8Array>()->ByteLength()); |
| 67 | } |
| 68 | }))); |
| 69 | }); |
| 70 | } |
| 71 | |
| 72 | KJ_TEST("Reading from byob reader") { |
| 73 | TestFixture testFixture; |
| 74 | |
| 75 | struct TestData { |
| 76 | size_t streamLength = 10 * 1024; |
| 77 | size_t bufferSize; |
| 78 | bool expectDone; |
| 79 | }; |
| 80 | |
| 81 | TestData tests[] = { |
| 82 | {.streamLength = 10 * 1024, .bufferSize = 100}, |
| 83 | {.streamLength = 10 * 1024, .bufferSize = 100 * 1024}, |
| 84 | {.streamLength = 10, .bufferSize = 100}, |
| 85 | {.streamLength = 1024, .bufferSize = 1024}, |
| 86 | }; |
| 87 | |
| 88 | for (auto test: tests) { |
| 89 | testFixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 90 | auto& js = jsg::Lock::from(env.isolate); |
| 91 | auto stream = |
| 92 | js.alloc<ReadableStream>(env.context, kj::heap<FakeStreamSource>(test.streamLength)); |
| 93 | ReadableStream::GetReaderOptions getReaderOptons = {.mode = kj::str("byob")}; |
| 94 | auto reader = stream->getReader(js, kj::mv(getReaderOptons)); |
| 95 | KJ_REQUIRE(reader.is<jsg::Ref<ReadableStreamBYOBReader>>()); |
| 96 | auto& byobReader = reader.get<jsg::Ref<ReadableStreamBYOBReader>>(); |
| 97 | |
| 98 | auto buffer = v8::Uint8Array::New( |
| 99 | v8::ArrayBuffer::New(js.v8Isolate, test.bufferSize), 0, test.bufferSize); |
| 100 | |
| 101 | return env.context.awaitJs(js, byobReader->read(js, buffer, {}).then(js, |
| 102 | JSG_VISITABLE_LAMBDA( |
| 103 | (test, reader = byobReader.addRef(), stream = stream.addRef()), |
| 104 | (reader, stream), (jsg::Lock& js, ReadResult readResult) { |
| 105 | KJ_ASSERT(!readResult.done); |
| 106 | |
| 107 | auto& value = KJ_REQUIRE_NONNULL(readResult.value); |
| 108 | auto handle = value.getHandle(js); |
| 109 | KJ_ASSERT(handle->IsUint8Array()); |
| 110 | auto view = handle.As<v8::Uint8Array>(); |
| 111 | KJ_ASSERT(kj::min(test.streamLength, test.bufferSize) == view->ByteLength()); |
| 112 | KJ_ASSERT(test.bufferSize == view->Buffer()->ByteLength()); |
| 113 | }))); |
| 114 | return kj::READY_NOW; |
| 115 | }); |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | KJ_TEST("PumpToReader regression") { |
| 120 | // If the promise holding the PumpToReader is dropped while the inner |
| 121 | // write to the sink is pending, the PumpToReader can free the sink. |
| 122 | // In some cases, this means that the sink can error because shutdownWrite |
| 123 | // is called while there is still a pending write promise. This test verifies |
| 124 | // that PumpToReader cancels any pending write promise when it is destroyed. |
| 125 | |
| 126 | struct TestSink final: public WritableStreamSink { |
| 127 | kj::TwoWayPipe pipe; |
| 128 | kj::PromiseFulfillerPair<void> paf; |
| 129 | kj::Vector<kj::String>& events; |
| 130 | TestSink(kj::Vector<kj::String>& events) |
| 131 | : pipe(kj::newTwoWayPipe()), |
| 132 | paf(kj::newPromiseAndFulfiller<void>()), |
| 133 | events(events) {} |
| 134 | |
| 135 | ~TestSink() { |
| 136 | events.add(kj::str("sink was destroyed")); |
| 137 | pipe.ends[0]->shutdownWrite(); |
| 138 | } |
| 139 | |
| 140 | kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override { |
| 141 | events.add(kj::str("got the write")); |
| 142 | |
| 143 | paf.fulfiller->fulfill(); |
| 144 | |
| 145 | return pipe.ends[0]->write(buffer).attach( |
| 146 | kj::defer([this] { events.add(kj::str("write promise was dropped")); })); |
| 147 | } |
| 148 | |
| 149 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override { |
| 150 | events.add(kj::str("got the write")); |
| 151 | paf.fulfiller->fulfill(); |
| 152 | // Concatenate pieces into a single buffer for the pipe write. |
| 153 | kj::Vector<byte> data; |
| 154 | for (auto& piece: pieces) { |
| 155 | data.addAll(piece); |
| 156 | } |
| 157 | auto arr = data.releaseAsArray(); |
| 158 | return pipe.ends[0]->write(arr).attach( |
| 159 | kj::mv(arr), kj::defer([this] { events.add(kj::str("write promise was dropped")); })); |
| 160 | } |
| 161 | |
| 162 | kj::Promise<void> end() override { |
| 163 | return kj::READY_NOW; |
| 164 | } |
| 165 | |
| 166 | void abort(kj::Exception reason) override {} |
| 167 | }; |
| 168 | |
| 169 | kj::Vector<kj::String> events; |
| 170 | capnp::MallocMessageBuilder flagsBuilder; |
| 171 | auto featureFlags = flagsBuilder.initRoot<CompatibilityFlags>(); |
| 172 | featureFlags.setStreamsJavaScriptControllers(true); |
| 173 | TestFixture testFixture({.featureFlags = featureFlags.asReader()}); |
| 174 | |
| 175 | testFixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 176 | auto& js = jsg::Lock::from(env.isolate); |
| 177 | auto stream = ReadableStream::constructor(js, |
| 178 | UnderlyingSource{.start = |
| 179 | [](jsg::Lock& js, auto controller) { |
| 180 | auto& c = KJ_REQUIRE_NONNULL( |
| 181 | controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>()); |
| 182 | c->enqueue(js, v8::ArrayBuffer::New(js.v8Isolate, 10)); |
| 183 | c->close(js); |
| 184 | return js.resolvedPromise(); |
| 185 | }}, |
| 186 | kj::none); |
| 187 | |
| 188 | auto sink = kj::heap<TestSink>(events); |
| 189 | auto writePromise = kj::mv(sink->paf.promise); |
| 190 | auto promise = stream->pumpTo(js, kj::mv(sink), true); |
| 191 | |
| 192 | return writePromise.attach(kj::mv(promise)); |
| 193 | }); |
| 194 | |
| 195 | KJ_ASSERT(events.size() == 3); |
| 196 | KJ_ASSERT(events[0] == "got the write"); |
| 197 | KJ_ASSERT(events[1] == "write promise was dropped"); |
| 198 | KJ_ASSERT(events[2] == "sink was destroyed"); |
| 199 | } |
| 200 | |
| 201 | } // namespace |
| 202 | } // namespace workerd::api |