File
Blob: src/workerd/tests/bench-pumpto.c++
| 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 | |
| 36 | namespace workerd::api::streams { |
| 37 | namespace { |
| 38 | |
| 39 | // ============================================================================= |
| 40 | // Stream configuration |
| 41 | // ============================================================================= |
| 42 | |
| 43 | enum 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 | |
| 49 | struct 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. |
| 58 | struct 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 | |
| 90 | static size_t benchChunkCounterStatic = 0; |
| 91 | |
| 92 | // Creates a JS-backed value ReadableStream that produces data synchronously in pull(). |
| 93 | jsg::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(). |
| 121 | jsg::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. |
| 152 | jsg::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 | |
| 191 | jsg::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 |
| 212 | static 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 | |
| 249 | static const StreamConfig VALUE_DEFAULT{.type = StreamType::VALUE}; |
| 250 | static const StreamConfig BYTE_DEFAULT{.type = StreamType::BYTE}; |
| 251 | static 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 |
| 261 | static void PumpTo_64B_Value(benchmark::State& state) { |
| 262 | benchPumpTo(state, 64, 16384, VALUE_DEFAULT); |
| 263 | } |
| 264 | static void PumpTo_256B_Value(benchmark::State& state) { |
| 265 | benchPumpTo(state, 256, 4096, VALUE_DEFAULT); |
| 266 | } |
| 267 | static void PumpTo_1KB_Value(benchmark::State& state) { |
| 268 | benchPumpTo(state, 1024, 1024, VALUE_DEFAULT); |
| 269 | } |
| 270 | static void PumpTo_4KB_Value(benchmark::State& state) { |
| 271 | benchPumpTo(state, 4096, 256, VALUE_DEFAULT); |
| 272 | } |
| 273 | static void PumpTo_16KB_Value(benchmark::State& state) { |
| 274 | benchPumpTo(state, 16384, 64, VALUE_DEFAULT); |
| 275 | } |
| 276 | static void PumpTo_64KB_Value(benchmark::State& state) { |
| 277 | benchPumpTo(state, 65536, 16, VALUE_DEFAULT); |
| 278 | } |
| 279 | |
| 280 | // Byte streams |
| 281 | static void PumpTo_64B_Byte(benchmark::State& state) { |
| 282 | benchPumpTo(state, 64, 16384, BYTE_DEFAULT); |
| 283 | } |
| 284 | static void PumpTo_256B_Byte(benchmark::State& state) { |
| 285 | benchPumpTo(state, 256, 4096, BYTE_DEFAULT); |
| 286 | } |
| 287 | static void PumpTo_1KB_Byte(benchmark::State& state) { |
| 288 | benchPumpTo(state, 1024, 1024, BYTE_DEFAULT); |
| 289 | } |
| 290 | static void PumpTo_4KB_Byte(benchmark::State& state) { |
| 291 | benchPumpTo(state, 4096, 256, BYTE_DEFAULT); |
| 292 | } |
| 293 | static void PumpTo_16KB_Byte(benchmark::State& state) { |
| 294 | benchPumpTo(state, 16384, 64, BYTE_DEFAULT); |
| 295 | } |
| 296 | static 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 | |
| 308 | static void PumpTo_256B_IoLatency(benchmark::State& state) { |
| 309 | benchPumpTo(state, 256, 256, IO_LATENCY_VALUE_DEFAULT); |
| 310 | } |
| 311 | static void PumpTo_4KB_IoLatency(benchmark::State& state) { |
| 312 | benchPumpTo(state, 4096, 16, IO_LATENCY_VALUE_DEFAULT); |
| 313 | } |
| 314 | static 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 | |
| 324 | static void PumpTo_64B_10MB_Value(benchmark::State& state) { |
| 325 | benchPumpTo(state, 64, 163840, VALUE_DEFAULT); |
| 326 | } |
| 327 | static void PumpTo_256B_10MB_Value(benchmark::State& state) { |
| 328 | benchPumpTo(state, 256, 40960, VALUE_DEFAULT); |
| 329 | } |
| 330 | static 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 |
| 339 | WD_BENCHMARK(PumpTo_64B_Value); |
| 340 | WD_BENCHMARK(PumpTo_256B_Value); |
| 341 | WD_BENCHMARK(PumpTo_1KB_Value); |
| 342 | WD_BENCHMARK(PumpTo_4KB_Value); |
| 343 | WD_BENCHMARK(PumpTo_16KB_Value); |
| 344 | WD_BENCHMARK(PumpTo_64KB_Value); |
| 345 | |
| 346 | // Sync 1 MiB — byte streams |
| 347 | WD_BENCHMARK(PumpTo_64B_Byte); |
| 348 | WD_BENCHMARK(PumpTo_256B_Byte); |
| 349 | WD_BENCHMARK(PumpTo_1KB_Byte); |
| 350 | WD_BENCHMARK(PumpTo_4KB_Byte); |
| 351 | WD_BENCHMARK(PumpTo_16KB_Byte); |
| 352 | WD_BENCHMARK(PumpTo_64KB_Byte); |
| 353 | |
| 354 | // I/O latency — 64 KiB (no-regression check) |
| 355 | WD_BENCHMARK(PumpTo_256B_IoLatency); |
| 356 | WD_BENCHMARK(PumpTo_4KB_IoLatency); |
| 357 | WD_BENCHMARK(PumpTo_64KB_IoLatency); |
| 358 | |
| 359 | // Large payload — 10 MiB value streams |
| 360 | WD_BENCHMARK(PumpTo_64B_10MB_Value); |
| 361 | WD_BENCHMARK(PumpTo_256B_10MB_Value); |
| 362 | WD_BENCHMARK(PumpTo_1KB_10MB_Value); |
| 363 | |
| 364 | } // namespace |
| 365 | } // namespace workerd::api::streams |