Skip to content
File

Blob: src/workerd/tests/bench-pumpto.c++

13.4 KB
1// Copyright (c) 2026 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// Benchmark for PumpToReader in standard.c++.
6//
7// Measures the performance of ReadableStream::pumpTo() which routes through
8// ReadableStreamJsController::pumpTo() → PumpToReader::pumpLoop().
9//
10// This benchmark establishes a baseline before the DrainingReader adoption,
11// then the same benchmarks are re-run after the change to quantify improvement.
12// This test was originally written to measure improvement from DrainingReader
13// adoption (deployed by an autogate), but remains broadly useful as a benchmark
14// even after we remove the autogate.
15//
16// Usage:
17// # Capture baseline (before changes):
18// bazel run --config=opt //src/workerd/tests:bench-pumpto \
19// -- --benchmark_format=json --benchmark_out=baseline.json
20//
21// # Capture comparison (after changes):
22// bazel run --config=opt //src/workerd/tests:bench-pumpto \
23// -- --benchmark_format=json --benchmark_out=after.json
24//
25// Key metrics:
26// - bytes_per_second: Primary throughput metric.
27// - WriteOps: Average sink write calls per iteration. Directly measures batching.
28// Before DrainingReader adoption: WriteOps ≈ numChunks (one write per chunk).
29// After: WriteOps ≪ numChunks (one vectored write per drain cycle).
30 
31#include <workerd/api/streams/standard.h>
32#include <workerd/api/system-streams.h>
33#include <workerd/tests/bench-tools.h>
34#include <workerd/tests/test-fixture.h>
35 
36namespace workerd::api::streams {
37namespace {
38 
39// =============================================================================
40// Stream configuration
41// =============================================================================
42 
43enum class StreamType {
44 VALUE, // Default ReadableStreamDefaultController
45 BYTE, // ReadableByteStreamController
46 IO_LATENCY_VALUE, // Value stream that yields to KJ event loop between chunks
47};
48 
49struct StreamConfig {
50 StreamType type = StreamType::VALUE;
51};
52 
53// =============================================================================
54// Test utilities
55// =============================================================================
56 
57// A discarding sink that counts bytes written and number of write operations.
58struct DiscardingSink final: public kj::AsyncOutputStream {
59 size_t bytesWritten = 0;
60 size_t writeCount = 0;
61 
62 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
63 writeCount++;
64 bytesWritten += buffer.size();
65 co_return;
66 }
67 
68 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
69 writeCount++;
70 for (auto piece: pieces) {
71 bytesWritten += piece.size();
72 }
73 co_return;
74 }
75 
76 kj::Promise<void> whenWriteDisconnected() override {
77 return kj::NEVER_DONE;
78 }
79 
80 void reset() {
81 bytesWritten = 0;
82 writeCount = 0;
83 }
84};
85 
86// =============================================================================
87// Stream creation helpers
88// =============================================================================
89 
90static size_t benchChunkCounterStatic = 0;
91 
92// Creates a JS-backed value ReadableStream that produces data synchronously in pull().
93jsg::Ref<ReadableStream> createValueStream(
94 jsg::Lock& js, size_t chunkSize, size_t numChunks, size_t* counter) {
95 return ReadableStream::constructor(js,
96 UnderlyingSource{
97 .pull =
98 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
99 auto& c =
100 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
101 
102 if ((*counter)++ < numChunks) {
103 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
104 jsg::BufferSource buffer(js, kj::mv(backing));
105 buffer.asArrayPtr().fill(0xAB);
106 c->enqueue(js, buffer.getHandle(js));
107 }
108 if (*counter == numChunks) {
109 c->close(js);
110 }
111 return js.resolvedPromise();
112 },
113 .expectedLength = chunkSize * numChunks,
114 },
115 StreamQueuingStrategy{
116 .highWaterMark = 0,
117 });
118}
119 
120// Creates a JS-backed byte ReadableStream that produces data synchronously in pull().
121jsg::Ref<ReadableStream> createByteStream(
122 jsg::Lock& js, size_t chunkSize, size_t numChunks, size_t* counter) {
123 return ReadableStream::constructor(js,
124 UnderlyingSource{
125 .type = kj::str("bytes"),
126 .pull =
127 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
128 auto& c =
129 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableByteStreamController>>());
130 
131 if ((*counter)++ < numChunks) {
132 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
133 jsg::BufferSource buffer(js, kj::mv(backing));
134 buffer.asArrayPtr().fill(0xAB);
135 c->enqueue(js, kj::mv(buffer));
136 }
137 if (*counter == numChunks) {
138 c->close(js);
139 }
140 return js.resolvedPromise();
141 },
142 .expectedLength = chunkSize * numChunks,
143 },
144 StreamQueuingStrategy{
145 .highWaterMark = 0,
146 });
147}
148 
149// Creates a value stream that yields to the KJ event loop between chunks.
150// Simulates a network stream where data arrives with real I/O latency.
151// Each chunk requires a KJ event loop iteration, so DrainingReader cannot batch them.
152jsg::Ref<ReadableStream> createIoLatencyValueStream(
153 jsg::Lock& js, size_t chunkSize, size_t numChunks, size_t* counter) {
154 return ReadableStream::constructor(js,
155 UnderlyingSource{
156 .pull =
157 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
158 auto& c =
159 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
160 
161 if (*counter >= numChunks) {
162 c->close(js);
163 return js.resolvedPromise();
164 }
165 
166 // Use IoContext.awaitIo() to wait for a KJ event loop yield.
167 // kj::evalLater() schedules on the next KJ event loop iteration.
168 auto& ioContext = IoContext::current();
169 auto cRef = c.addRef();
170 return ioContext.awaitIo(js, kj::evalLater([]() {}),
171 JSG_VISITABLE_LAMBDA(
172 (cRef = kj::mv(cRef), chunkSize, numChunks, counter), (cRef), (jsg::Lock & js) mutable {
173 if ((*counter)++ < numChunks) {
174 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
175 jsg::BufferSource buffer(js, kj::mv(backing));
176 buffer.asArrayPtr().fill(0xAB);
177 cRef->enqueue(js, buffer.getHandle(js));
178 }
179 if (*counter == numChunks) {
180 cRef->close(js);
181 }
182 }));
183 },
184 .expectedLength = chunkSize * numChunks,
185 },
186 StreamQueuingStrategy{
187 .highWaterMark = 0,
188 });
189}
190 
191jsg::Ref<ReadableStream> createConfiguredStream(
192 jsg::Lock& js, size_t chunkSize, size_t numChunks, const StreamConfig& config) {
193 benchChunkCounterStatic = 0;
194 size_t* counter = &benchChunkCounterStatic;
195 
196 switch (config.type) {
197 case StreamType::VALUE:
198 return createValueStream(js, chunkSize, numChunks, counter);
199 case StreamType::BYTE:
200 return createByteStream(js, chunkSize, numChunks, counter);
201 case StreamType::IO_LATENCY_VALUE:
202 return createIoLatencyValueStream(js, chunkSize, numChunks, counter);
203 }
204 KJ_UNREACHABLE;
205}
206 
207// =============================================================================
208// Core benchmark function
209// =============================================================================
210 
211// Exercises: ReadableStream::pumpTo() → ReadableStreamJsController::pumpTo() → PumpToReader
212static void benchPumpTo(
213 benchmark::State& state, size_t chunkSize, size_t numChunks, const StreamConfig& config) {
214 capnp::MallocMessageBuilder message;
215 auto flags = message.initRoot<CompatibilityFlags>();
216 flags.setStreamsJavaScriptControllers(true);
217 TestFixture fixture({.featureFlags = flags.asReader()});
218 
219 DiscardingSink sink;
220 size_t expectedBytes = chunkSize * numChunks;
221 
222 for (auto _: state) {
223 sink.reset();
224 
225 fixture.runInIoContext([&](const TestFixture::Environment& env) {
226 auto stream = createConfiguredStream(env.js, chunkSize, numChunks, config);
227 
228 // Wrap DiscardingSink as a WritableStreamSink via newSystemStream.
229 // This is the production path: PumpToReader receives a WritableStreamSink.
230 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
231 auto writableSink = newSystemStream(kj::mv(fakeOwn), StreamEncoding::IDENTITY, env.context);
232 
233 return env.context.waitForDeferredProxy(stream->pumpTo(env.js, kj::mv(writableSink), true));
234 });
235 
236 KJ_ASSERT(sink.bytesWritten == expectedBytes, "Expected", expectedBytes, "bytes but got",
237 sink.bytesWritten);
238 }
239 
240 state.SetBytesProcessed(state.iterations() * static_cast<int64_t>(expectedBytes));
241 state.counters["WriteOps"] =
242 benchmark::Counter(sink.writeCount, benchmark::Counter::kAvgIterations);
243}
244 
245// =============================================================================
246// Stream configs
247// =============================================================================
248 
249static const StreamConfig VALUE_DEFAULT{.type = StreamType::VALUE};
250static const StreamConfig BYTE_DEFAULT{.type = StreamType::BYTE};
251static const StreamConfig IO_LATENCY_VALUE_DEFAULT{.type = StreamType::IO_LATENCY_VALUE};
252 
253// =============================================================================
254// Synchronous streams — 1 MiB total payload
255// =============================================================================
256// These are the primary benchmarks. Data is produced synchronously in the pull
257// callback. DrainingReader (post-change) can drain all chunks in a single lock
258// acquisition, so small-chunk benchmarks should see large improvement.
259 
260// Value streams
261static void PumpTo_64B_Value(benchmark::State& state) {
262 benchPumpTo(state, 64, 16384, VALUE_DEFAULT);
263}
264static void PumpTo_256B_Value(benchmark::State& state) {
265 benchPumpTo(state, 256, 4096, VALUE_DEFAULT);
266}
267static void PumpTo_1KB_Value(benchmark::State& state) {
268 benchPumpTo(state, 1024, 1024, VALUE_DEFAULT);
269}
270static void PumpTo_4KB_Value(benchmark::State& state) {
271 benchPumpTo(state, 4096, 256, VALUE_DEFAULT);
272}
273static void PumpTo_16KB_Value(benchmark::State& state) {
274 benchPumpTo(state, 16384, 64, VALUE_DEFAULT);
275}
276static void PumpTo_64KB_Value(benchmark::State& state) {
277 benchPumpTo(state, 65536, 16, VALUE_DEFAULT);
278}
279 
280// Byte streams
281static void PumpTo_64B_Byte(benchmark::State& state) {
282 benchPumpTo(state, 64, 16384, BYTE_DEFAULT);
283}
284static void PumpTo_256B_Byte(benchmark::State& state) {
285 benchPumpTo(state, 256, 4096, BYTE_DEFAULT);
286}
287static void PumpTo_1KB_Byte(benchmark::State& state) {
288 benchPumpTo(state, 1024, 1024, BYTE_DEFAULT);
289}
290static void PumpTo_4KB_Byte(benchmark::State& state) {
291 benchPumpTo(state, 4096, 256, BYTE_DEFAULT);
292}
293static void PumpTo_16KB_Byte(benchmark::State& state) {
294 benchPumpTo(state, 16384, 64, BYTE_DEFAULT);
295}
296static void PumpTo_64KB_Byte(benchmark::State& state) {
297 benchPumpTo(state, 65536, 16, BYTE_DEFAULT);
298}
299 
300// =============================================================================
301// I/O latency streams — 64 KiB total payload
302// =============================================================================
303// Each chunk requires a KJ event loop yield, simulating real network I/O.
304// DrainingReader cannot batch these (at most 1 chunk per drain cycle).
305// These verify no regression from the PumpToReader change.
306// Smaller total payload because each chunk incurs real event loop overhead.
307 
308static void PumpTo_256B_IoLatency(benchmark::State& state) {
309 benchPumpTo(state, 256, 256, IO_LATENCY_VALUE_DEFAULT);
310}
311static void PumpTo_4KB_IoLatency(benchmark::State& state) {
312 benchPumpTo(state, 4096, 16, IO_LATENCY_VALUE_DEFAULT);
313}
314static void PumpTo_64KB_IoLatency(benchmark::State& state) {
315 benchPumpTo(state, 65536, 1, IO_LATENCY_VALUE_DEFAULT);
316}
317 
318// =============================================================================
319// Large payload — 10 MiB total, sync value streams
320// =============================================================================
321// Sustained throughput test with small chunks. More data amortizes fixture
322// setup cost, yielding more stable measurements.
323 
324static void PumpTo_64B_10MB_Value(benchmark::State& state) {
325 benchPumpTo(state, 64, 163840, VALUE_DEFAULT);
326}
327static void PumpTo_256B_10MB_Value(benchmark::State& state) {
328 benchPumpTo(state, 256, 40960, VALUE_DEFAULT);
329}
330static void PumpTo_1KB_10MB_Value(benchmark::State& state) {
331 benchPumpTo(state, 1024, 10240, VALUE_DEFAULT);
332}
333 
334// =============================================================================
335// Register benchmarks
336// =============================================================================
337 
338// Sync 1 MiB — value streams
339WD_BENCHMARK(PumpTo_64B_Value);
340WD_BENCHMARK(PumpTo_256B_Value);
341WD_BENCHMARK(PumpTo_1KB_Value);
342WD_BENCHMARK(PumpTo_4KB_Value);
343WD_BENCHMARK(PumpTo_16KB_Value);
344WD_BENCHMARK(PumpTo_64KB_Value);
345 
346// Sync 1 MiB — byte streams
347WD_BENCHMARK(PumpTo_64B_Byte);
348WD_BENCHMARK(PumpTo_256B_Byte);
349WD_BENCHMARK(PumpTo_1KB_Byte);
350WD_BENCHMARK(PumpTo_4KB_Byte);
351WD_BENCHMARK(PumpTo_16KB_Byte);
352WD_BENCHMARK(PumpTo_64KB_Byte);
353 
354// I/O latency — 64 KiB (no-regression check)
355WD_BENCHMARK(PumpTo_256B_IoLatency);
356WD_BENCHMARK(PumpTo_4KB_IoLatency);
357WD_BENCHMARK(PumpTo_64KB_IoLatency);
358 
359// Large payload — 10 MiB value streams
360WD_BENCHMARK(PumpTo_64B_10MB_Value);
361WD_BENCHMARK(PumpTo_256B_10MB_Value);
362WD_BENCHMARK(PumpTo_1KB_10MB_Value);
363 
364} // namespace
365} // namespace workerd::api::streams