Skip to content
File

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

7.3 KB
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 
8namespace workerd::api {
9 
10namespace {
11 
12class 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 
29KJ_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 
44KJ_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 
72KJ_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 
119KJ_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