Skip to content
File

Blob: src/workerd/tests/bench-stream-piping.c++

48.5 KB
1// Copyright (c) 2024 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 to compare stream piping implementations:
6// 1. Existing approach (ReadableStream::pumpTo via PumpToReader) - uses JS promise-based loop
7// 2. New approach (ReadableSourceKjAdapter::pumpTo) - uses DrainingReader to pull all
8// synchronously available data at once, then writes with vectored I/O
9//
10// Run with: bazel run --config=opt //src/workerd/tests:bench-stream-piping
11 
12#include <workerd/api/streams/readable-source-adapter.h>
13#include <workerd/api/streams/standard.h>
14#include <workerd/api/streams/writable-sink.h>
15#include <workerd/api/system-streams.h>
16#include <workerd/tests/bench-tools.h>
17#include <workerd/tests/test-fixture.h>
18 
19#include <kj/compat/http.h>
20 
21namespace workerd::api::streams {
22namespace {
23 
24// =============================================================================
25// Stream configuration types
26// =============================================================================
27 
28enum class StreamType {
29 VALUE, // Default ReadableStreamDefaultController
30 BYTE, // ReadableByteStreamController
31 SLOW_VALUE, // Value stream that produces one chunk per microtask (async)
32 IO_LATENCY_VALUE, // Value stream that yields to KJ event loop between chunks
33 IO_LATENCY_BYTE, // Byte stream that yields to KJ event loop between chunks
34 TIMED_VALUE, // Value stream with configurable timer delay between chunks
35};
36 
37struct StreamConfig {
38 StreamType type = StreamType::VALUE;
39 kj::Maybe<size_t> autoAllocateChunkSize; // Only valid for BYTE streams
40 kj::Duration chunkDelay = 0 * kj::MILLISECONDS; // Delay between chunks for TIMED_* streams
41 double highWaterMark = 0; // 0 means default (pull on demand)
42};
43 
44// =============================================================================
45// Test utilities
46// =============================================================================
47 
48// A discarding sink that just counts bytes written (more representative of real network I/O).
49struct DiscardingSink final: public kj::AsyncOutputStream {
50 size_t bytesWritten = 0;
51 size_t writeCount = 0;
52 
53 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
54 writeCount++;
55 bytesWritten += buffer.size();
56 co_return;
57 }
58 
59 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
60 writeCount++;
61 for (auto piece: pieces) {
62 bytesWritten += piece.size();
63 }
64 co_return;
65 }
66 
67 kj::Promise<void> whenWriteDisconnected() override {
68 return kj::NEVER_DONE;
69 }
70 
71 void reset() {
72 bytesWritten = 0;
73 writeCount = 0;
74 }
75};
76 
77// A sink that simulates network backpressure with configurable latency per write.
78// This represents real-world scenarios where the downstream connection is slower
79// than the upstream source (e.g., slow client, congested network).
80struct LatencySink final: public kj::AsyncOutputStream {
81 kj::Timer& timer;
82 kj::Duration writeLatency;
83 size_t bytesWritten = 0;
84 size_t writeCount = 0;
85 
86 LatencySink(kj::Timer& timer, kj::Duration writeLatency)
87 : timer(timer),
88 writeLatency(writeLatency) {}
89 
90 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
91 writeCount++;
92 bytesWritten += buffer.size();
93 if (writeLatency > 0 * kj::MILLISECONDS) {
94 co_await timer.afterDelay(writeLatency);
95 }
96 co_return;
97 }
98 
99 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
100 writeCount++;
101 for (auto piece: pieces) {
102 bytesWritten += piece.size();
103 }
104 if (writeLatency > 0 * kj::MILLISECONDS) {
105 co_await timer.afterDelay(writeLatency);
106 }
107 co_return;
108 }
109 
110 kj::Promise<void> whenWriteDisconnected() override {
111 return kj::NEVER_DONE;
112 }
113 
114 void reset() {
115 bytesWritten = 0;
116 writeCount = 0;
117 }
118};
119 
120// Creates a JS-backed ReadableStream with the specified configuration.
121// Uses a counter pointer similar to the unit tests in readable-source-adapter-test.c++.
122static size_t benchChunkCounterStatic = 0;
123 
124jsg::Ref<ReadableStream> createValueStream(
125 jsg::Lock& js, size_t chunkSize, size_t numChunks, double highWaterMark, size_t* counter) {
126 return ReadableStream::constructor(js,
127 UnderlyingSource{
128 .pull =
129 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
130 auto& c =
131 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
132 
133 if ((*counter)++ < numChunks) {
134 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
135 jsg::BufferSource buffer(js, kj::mv(backing));
136 buffer.asArrayPtr().fill(0xAB);
137 c->enqueue(js, buffer.getHandle(js));
138 }
139 if (*counter == numChunks) {
140 c->close(js);
141 }
142 return js.resolvedPromise();
143 },
144 .expectedLength = chunkSize * numChunks,
145 },
146 StreamQueuingStrategy{
147 .highWaterMark = highWaterMark,
148 });
149}
150 
151jsg::Ref<ReadableStream> createByteStream(jsg::Lock& js,
152 size_t chunkSize,
153 size_t numChunks,
154 kj::Maybe<size_t> autoAllocateChunkSize,
155 double highWaterMark,
156 size_t* counter) {
157 return ReadableStream::constructor(js,
158 UnderlyingSource{
159 .type = kj::str("bytes"),
160 .autoAllocateChunkSize = autoAllocateChunkSize,
161 .pull =
162 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
163 auto& c =
164 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableByteStreamController>>());
165 
166 if ((*counter)++ < numChunks) {
167 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
168 jsg::BufferSource buffer(js, kj::mv(backing));
169 buffer.asArrayPtr().fill(0xAB);
170 c->enqueue(js, kj::mv(buffer));
171 }
172 if (*counter == numChunks) {
173 c->close(js);
174 }
175 return js.resolvedPromise();
176 },
177 .expectedLength = chunkSize * numChunks,
178 },
179 StreamQueuingStrategy{
180 .highWaterMark = highWaterMark,
181 });
182}
183 
184// Creates a "slow" value stream that produces one chunk per microtask.
185// This simulates a stream where pull() has async work to do before data is ready.
186// The pull() function returns a promise that resolves on the next microtask,
187// and only enqueues data WHEN the promise resolves.
188//
189// NOTE: This does NOT prevent batching or trigger the adaptive read policy!
190// Microtask delays execute synchronously within the JS event loop turn, so
191// readInternal's promise chain runs to completion before returning to KJ.
192// The buffer still fills completely, achieving full batching (100 chunks → 1 write).
193// See PUMP_PERFORMANCE_ANALYSIS.md section 9 for detailed analysis.
194jsg::Ref<ReadableStream> createSlowValueStream(
195 jsg::Lock& js, size_t chunkSize, size_t numChunks, double highWaterMark, size_t* counter) {
196 return ReadableStream::constructor(js,
197 UnderlyingSource{
198 .pull =
199 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
200 auto& c =
201 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
202 
203 if (*counter >= numChunks) {
204 c->close(js);
205 return js.resolvedPromise();
206 }
207 
208 // Return a promise that enqueues data on the next microtask.
209 // This adds a tiny delay per chunk, but does NOT prevent batching -
210 // the entire promise chain still runs within one JS event loop turn.
211 auto cRef = c.addRef();
212 return js.resolvedPromise().then(js,
213 JSG_VISITABLE_LAMBDA(
214 (cRef = kj::mv(cRef), chunkSize, numChunks, counter), (cRef), (jsg::Lock & js) mutable {
215 if ((*counter)++ < numChunks) {
216 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
217 jsg::BufferSource buffer(js, kj::mv(backing));
218 buffer.asArrayPtr().fill(0xAB);
219 cRef->enqueue(js, buffer.getHandle(js));
220 }
221 if (*counter == numChunks) {
222 cRef->close(js);
223 }
224 return js.resolvedPromise();
225 }));
226 },
227 .expectedLength = chunkSize * numChunks,
228 },
229 StreamQueuingStrategy{
230 .highWaterMark = highWaterMark,
231 });
232}
233 
234// Creates a value stream that yields to the KJ event loop between chunks.
235// This simulates a network stream (like fetch response body) where data arrives with real
236// I/O latency. Unlike the "slow" stream that uses microtask delays, this stream's pull()
237// returns a promise that only resolves after a KJ event loop iteration.
238//
239// This WILL cause pumpReadImpl to return early, potentially triggering the adaptive read policy.
240// See PUMP_PERFORMANCE_ANALYSIS.md section 9 for why this is different from microtask delays.
241jsg::Ref<ReadableStream> createIoLatencyValueStream(
242 jsg::Lock& js, size_t chunkSize, size_t numChunks, double highWaterMark, size_t* counter) {
243 return ReadableStream::constructor(js,
244 UnderlyingSource{
245 .pull =
246 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
247 auto& c =
248 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
249 
250 if (*counter >= numChunks) {
251 c->close(js);
252 return js.resolvedPromise();
253 }
254 
255 // Use IoContext.awaitIo() to wait for a KJ event loop yield.
256 // This simulates real network I/O latency where we yield to KJ between chunks.
257 // kj::evalLater() schedules on the next KJ event loop iteration.
258 auto& ioContext = IoContext::current();
259 auto cRef = c.addRef();
260 return ioContext.awaitIo(js, kj::evalLater([]() {}),
261 JSG_VISITABLE_LAMBDA(
262 (cRef = kj::mv(cRef), chunkSize, numChunks, counter), (cRef), (jsg::Lock & js) mutable {
263 if ((*counter)++ < numChunks) {
264 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
265 jsg::BufferSource buffer(js, kj::mv(backing));
266 buffer.asArrayPtr().fill(0xAB);
267 cRef->enqueue(js, buffer.getHandle(js));
268 }
269 if (*counter == numChunks) {
270 cRef->close(js);
271 }
272 }));
273 },
274 .expectedLength = chunkSize * numChunks,
275 },
276 StreamQueuingStrategy{
277 .highWaterMark = highWaterMark,
278 });
279}
280 
281// Creates a byte stream that yields to the KJ event loop between chunks.
282// Same as createIoLatencyValueStream but uses ReadableByteStreamController.
283jsg::Ref<ReadableStream> createIoLatencyByteStream(
284 jsg::Lock& js, size_t chunkSize, size_t numChunks, double highWaterMark, size_t* counter) {
285 return ReadableStream::constructor(js,
286 UnderlyingSource{
287 .type = kj::str("bytes"),
288 .pull =
289 [chunkSize, numChunks, counter](jsg::Lock& js, auto controller) {
290 auto& c =
291 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableByteStreamController>>());
292 
293 if (*counter >= numChunks) {
294 c->close(js);
295 return js.resolvedPromise();
296 }
297 
298 auto& ioContext = IoContext::current();
299 auto cRef = c.addRef();
300 return ioContext.awaitIo(js, kj::evalLater([]() {}),
301 JSG_VISITABLE_LAMBDA(
302 (cRef = kj::mv(cRef), chunkSize, numChunks, counter), (cRef), (jsg::Lock & js) mutable {
303 if ((*counter)++ < numChunks) {
304 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
305 jsg::BufferSource buffer(js, kj::mv(backing));
306 buffer.asArrayPtr().fill(0xAB);
307 cRef->enqueue(js, kj::mv(buffer));
308 }
309 if (*counter == numChunks) {
310 cRef->close(js);
311 }
312 }));
313 },
314 .expectedLength = chunkSize * numChunks,
315 },
316 StreamQueuingStrategy{
317 .highWaterMark = highWaterMark,
318 });
319}
320 
321// Creates a value stream with actual timer-based delays between chunks.
322// This simulates real network I/O where data arrives with measurable latency.
323// Unlike evalLater() which resumes immediately, timer delays represent real wall-clock time.
324//
325// With delays, we can observe:
326// 1. How throughput scales with I/O latency
327// 2. Whether double-buffering provides real overlap benefit
328// 3. The true cost of per-chunk I/O operations
329jsg::Ref<ReadableStream> createTimedValueStream(jsg::Lock& js,
330 size_t chunkSize,
331 size_t numChunks,
332 double highWaterMark,
333 kj::Duration delay,
334 size_t* counter) {
335 return ReadableStream::constructor(js,
336 UnderlyingSource{
337 .pull =
338 [chunkSize, numChunks, delay, counter](jsg::Lock& js, auto controller) {
339 auto& c =
340 KJ_ASSERT_NONNULL(controller.template tryGet<jsg::Ref<ReadableStreamDefaultController>>());
341 
342 if (*counter >= numChunks) {
343 c->close(js);
344 return js.resolvedPromise();
345 }
346 
347 // Use afterLimitTimeout for actual timer-based delay
348 auto& ioContext = IoContext::current();
349 auto cRef = c.addRef();
350 return ioContext.awaitIo(js, ioContext.afterLimitTimeout(delay),
351 JSG_VISITABLE_LAMBDA(
352 (cRef = kj::mv(cRef), chunkSize, numChunks, counter), (cRef), (jsg::Lock & js) mutable {
353 if ((*counter)++ < numChunks) {
354 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, chunkSize);
355 jsg::BufferSource buffer(js, kj::mv(backing));
356 buffer.asArrayPtr().fill(0xAB);
357 cRef->enqueue(js, buffer.getHandle(js));
358 }
359 if (*counter == numChunks) {
360 cRef->close(js);
361 }
362 }));
363 },
364 .expectedLength = chunkSize * numChunks,
365 },
366 StreamQueuingStrategy{
367 .highWaterMark = highWaterMark,
368 });
369}
370 
371jsg::Ref<ReadableStream> createConfiguredStream(
372 jsg::Lock& js, size_t chunkSize, size_t numChunks, const StreamConfig& config) {
373 benchChunkCounterStatic = 0;
374 size_t* counter = &benchChunkCounterStatic;
375 
376 switch (config.type) {
377 case StreamType::VALUE:
378 return createValueStream(js, chunkSize, numChunks, config.highWaterMark, counter);
379 case StreamType::BYTE:
380 return createByteStream(
381 js, chunkSize, numChunks, config.autoAllocateChunkSize, config.highWaterMark, counter);
382 case StreamType::SLOW_VALUE:
383 return createSlowValueStream(js, chunkSize, numChunks, config.highWaterMark, counter);
384 case StreamType::IO_LATENCY_VALUE:
385 return createIoLatencyValueStream(js, chunkSize, numChunks, config.highWaterMark, counter);
386 case StreamType::IO_LATENCY_BYTE:
387 return createIoLatencyByteStream(js, chunkSize, numChunks, config.highWaterMark, counter);
388 case StreamType::TIMED_VALUE:
389 return createTimedValueStream(
390 js, chunkSize, numChunks, config.highWaterMark, config.chunkDelay, counter);
391 }
392 KJ_UNREACHABLE;
393}
394 
395// =============================================================================
396// Benchmark: New approach using ReadableSourceKjAdapter::pumpTo
397// =============================================================================
398 
399static void benchNewApproachPumpTo(
400 benchmark::State& state, size_t chunkSize, size_t numChunks, const StreamConfig& config) {
401 capnp::MallocMessageBuilder message;
402 auto flags = message.initRoot<CompatibilityFlags>();
403 flags.setStreamsJavaScriptControllers(true);
404 // Always enable spec-compliant autoAllocateChunkSize behavior for the "New" approach.
405 // This is the behavior we're adopting going forward. When enabled, byte streams without
406 // explicit autoAllocateChunkSize will use DEFAULT reads (16KB buffer, no ByobRequest)
407 // instead of the legacy BYOB reads (4KB buffer, ByobRequest available).
408 // The "Existing" benchmarks keep the legacy behavior for comparison.
409 flags.setNoAutoAllocateChunkSize(true);
410 // Enable real timers for streams that need actual timer functionality (e.g., TIMED_VALUE).
411 bool needsRealTimers = config.type == StreamType::TIMED_VALUE;
412 TestFixture fixture({.featureFlags = flags.asReader(), .useRealTimers = needsRealTimers});
413 
414 DiscardingSink sink;
415 size_t expectedBytes = chunkSize * numChunks;
416 
417 for (auto _: state) {
418 sink.reset();
419 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
420 auto writableSink = newWritableSink(kj::mv(fakeOwn));
421 
422 fixture.runInIoContext([&](const TestFixture::Environment& env) {
423 auto stream = createConfiguredStream(env.js, chunkSize, numChunks, config);
424 auto adapter = kj::heap<ReadableSourceKjAdapter>(env.js, env.context, stream.addRef());
425 return adapter->pumpTo(*writableSink, EndAfterPump::YES).attach(kj::mv(adapter));
426 });
427 
428 // Verify all expected bytes were written
429 KJ_ASSERT(sink.bytesWritten == expectedBytes, "New approach: expected", expectedBytes,
430 "bytes but got", sink.bytesWritten);
431 }
432 
433 state.SetBytesProcessed(
434 state.iterations() * static_cast<int64_t>(chunkSize) * static_cast<int64_t>(numChunks));
435 state.counters["WriteOps"] =
436 benchmark::Counter(sink.writeCount, benchmark::Counter::kAvgIterations);
437}
438 
439// =============================================================================
440// Benchmark: Existing approach using ReadableStream::pumpTo (PumpToReader)
441// =============================================================================
442 
443static void benchExistingApproachPumpTo(
444 benchmark::State& state, size_t chunkSize, size_t numChunks, const StreamConfig& config) {
445 capnp::MallocMessageBuilder message;
446 auto flags = message.initRoot<CompatibilityFlags>();
447 flags.setStreamsJavaScriptControllers(true);
448 // Keep legacy autoAllocateChunkSize behavior for "Existing" benchmarks.
449 // This uses 4KB BYOB reads with ByobRequest available, which is the original workerd behavior.
450 // The "New" benchmarks use spec-compliant behavior for comparison.
451 flags.setNoAutoAllocateChunkSize(false);
452 // Enable real timers for streams that need actual timer functionality (e.g., TIMED_VALUE).
453 bool needsRealTimers = config.type == StreamType::TIMED_VALUE;
454 TestFixture fixture({.featureFlags = flags.asReader(), .useRealTimers = needsRealTimers});
455 
456 DiscardingSink sink;
457 size_t expectedBytes = chunkSize * numChunks;
458 
459 for (auto _: state) {
460 sink.reset();
461 
462 fixture.runInIoContext([&](const TestFixture::Environment& env) {
463 auto stream = createConfiguredStream(env.js, chunkSize, numChunks, config);
464 
465 kj::Own<kj::AsyncOutputStream> fakeOwn(&sink, kj::NullDisposer::instance);
466 auto writableSink = newSystemStream(kj::mv(fakeOwn), StreamEncoding::IDENTITY, env.context);
467 
468 return env.context.waitForDeferredProxy(stream->pumpTo(env.js, kj::mv(writableSink), true));
469 });
470 
471 // Verify all expected bytes were written
472 KJ_ASSERT(sink.bytesWritten == expectedBytes, "Existing approach: expected", expectedBytes,
473 "bytes but got", sink.bytesWritten);
474 }
475 
476 state.SetBytesProcessed(
477 state.iterations() * static_cast<int64_t>(chunkSize) * static_cast<int64_t>(numChunks));
478 state.counters["WriteOps"] =
479 benchmark::Counter(sink.writeCount, benchmark::Counter::kAvgIterations);
480}
481 
482// =============================================================================
483// Stream configurations to benchmark
484// =============================================================================
485 
486// Value stream with default highWaterMark (0)
487static const StreamConfig VALUE_DEFAULT{
488 .type = StreamType::VALUE,
489 .autoAllocateChunkSize = kj::none,
490 .highWaterMark = 0,
491};
492 
493// Value stream with 16KB highWaterMark
494static const StreamConfig VALUE_HWM_16K{
495 .type = StreamType::VALUE,
496 .autoAllocateChunkSize = kj::none,
497 .highWaterMark = 16 * 1024,
498};
499 
500// Byte stream without autoAllocateChunkSize, default highWaterMark
501static const StreamConfig BYTE_DEFAULT{
502 .type = StreamType::BYTE,
503 .autoAllocateChunkSize = kj::none,
504 .highWaterMark = 0,
505};
506 
507// Byte stream with autoAllocateChunkSize=64KB (fixed), default highWaterMark
508static const StreamConfig BYTE_AUTO_64K{
509 .type = StreamType::BYTE,
510 .autoAllocateChunkSize = 65536,
511 .highWaterMark = 0,
512};
513 
514// Byte stream without autoAllocateChunkSize, 16KB highWaterMark
515static const StreamConfig BYTE_HWM_16K{
516 .type = StreamType::BYTE,
517 .autoAllocateChunkSize = kj::none,
518 .highWaterMark = 16 * 1024,
519};
520 
521// Byte stream with autoAllocateChunkSize=64KB, 16KB highWaterMark
522static const StreamConfig BYTE_AUTO_64K_HWM_16K{
523 .type = StreamType::BYTE,
524 .autoAllocateChunkSize = 65536,
525 .highWaterMark = 16 * 1024,
526};
527 
528// Slow value stream (async, one chunk per microtask) - does NOT trigger adaptive read policy
529// because microtasks execute synchronously within the JS event loop turn.
530static const StreamConfig SLOW_VALUE_DEFAULT{
531 .type = StreamType::SLOW_VALUE,
532 .autoAllocateChunkSize = kj::none,
533 .highWaterMark = 0,
534};
535 
536// I/O latency value stream - yields to KJ event loop between chunks, simulating network I/O.
537// This DOES trigger early returns from pumpReadImpl and may activate the adaptive policy.
538static const StreamConfig IO_LATENCY_VALUE_DEFAULT{
539 .type = StreamType::IO_LATENCY_VALUE,
540 .autoAllocateChunkSize = kj::none,
541 .highWaterMark = 0,
542};
543 
544// I/O latency byte stream - same as above but using ReadableByteStreamController.
545// Tests how byte streams interact with I/O latency patterns.
546static const StreamConfig IO_LATENCY_BYTE_DEFAULT{
547 .type = StreamType::IO_LATENCY_BYTE,
548 .autoAllocateChunkSize = kj::none,
549 .highWaterMark = 0,
550};
551 
552// Timed value streams - actual timer-based delays between chunks.
553// These simulate real network I/O with measurable latency.
554// The delay represents the time waiting for the next chunk from the network.
555 
556// 10μs delay - fast network, minimal latency (e.g., local network)
557static const StreamConfig TIMED_VALUE_10US{
558 .type = StreamType::TIMED_VALUE,
559 .autoAllocateChunkSize = kj::none,
560 .chunkDelay = 10 * kj::MICROSECONDS,
561 .highWaterMark = 0,
562};
563 
564// 100μs delay - typical datacenter latency
565static const StreamConfig TIMED_VALUE_100US{
566 .type = StreamType::TIMED_VALUE,
567 .autoAllocateChunkSize = kj::none,
568 .chunkDelay = 100 * kj::MICROSECONDS,
569 .highWaterMark = 0,
570};
571 
572// 1ms delay - typical internet latency / slow upstream
573static const StreamConfig TIMED_VALUE_1MS{
574 .type = StreamType::TIMED_VALUE,
575 .autoAllocateChunkSize = kj::none,
576 .chunkDelay = 1 * kj::MILLISECONDS,
577 .highWaterMark = 0,
578};
579 
580// =============================================================================
581// Chunk size configurations
582// =============================================================================
583 
584// Tiny chunks (worst case for JS overhead): 64 * 256 = 16,384 bytes
585static constexpr size_t TINY_CHUNK_SIZE = 64;
586static constexpr size_t TINY_NUM_CHUNKS = 256;
587 
588// Small chunks (chatty protocol pattern): 256 * 100 = 25,600 bytes
589static constexpr size_t SMALL_CHUNK_SIZE = 256;
590static constexpr size_t SMALL_NUM_CHUNKS = 100;
591 
592// Medium chunks (typical HTTP response): 4096 * 100 = 409,600 bytes (~400KB)
593static constexpr size_t MEDIUM_CHUNK_SIZE = 4096;
594static constexpr size_t MEDIUM_NUM_CHUNKS = 100;
595 
596// Large chunks (file transfer pattern): 65536 * 16 = 1,048,576 bytes (1MB)
597static constexpr size_t LARGE_CHUNK_SIZE = 65536;
598static constexpr size_t LARGE_NUM_CHUNKS = 16;
599 
600// =============================================================================
601// Benchmark functions - Value streams
602// =============================================================================
603 
604// Value stream, default HWM
605static void New_Tiny_Value(benchmark::State& state) {
606 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, VALUE_DEFAULT);
607}
608static void Existing_Tiny_Value(benchmark::State& state) {
609 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, VALUE_DEFAULT);
610}
611static void New_Small_Value(benchmark::State& state) {
612 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, VALUE_DEFAULT);
613}
614static void Existing_Small_Value(benchmark::State& state) {
615 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, VALUE_DEFAULT);
616}
617static void New_Medium_Value(benchmark::State& state) {
618 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, VALUE_DEFAULT);
619}
620static void Existing_Medium_Value(benchmark::State& state) {
621 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, VALUE_DEFAULT);
622}
623static void New_Large_Value(benchmark::State& state) {
624 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, VALUE_DEFAULT);
625}
626static void Existing_Large_Value(benchmark::State& state) {
627 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, VALUE_DEFAULT);
628}
629 
630// Value stream, 16KB HWM
631static void New_Tiny_Value_HWM16K(benchmark::State& state) {
632 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, VALUE_HWM_16K);
633}
634static void Existing_Tiny_Value_HWM16K(benchmark::State& state) {
635 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, VALUE_HWM_16K);
636}
637static void New_Small_Value_HWM16K(benchmark::State& state) {
638 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, VALUE_HWM_16K);
639}
640static void Existing_Small_Value_HWM16K(benchmark::State& state) {
641 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, VALUE_HWM_16K);
642}
643static void New_Medium_Value_HWM16K(benchmark::State& state) {
644 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, VALUE_HWM_16K);
645}
646static void Existing_Medium_Value_HWM16K(benchmark::State& state) {
647 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, VALUE_HWM_16K);
648}
649static void New_Large_Value_HWM16K(benchmark::State& state) {
650 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, VALUE_HWM_16K);
651}
652static void Existing_Large_Value_HWM16K(benchmark::State& state) {
653 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, VALUE_HWM_16K);
654}
655 
656// =============================================================================
657// Benchmark functions - Byte streams, no autoAllocate
658// =============================================================================
659 
660// Byte stream, default HWM, no autoAllocate
661static void New_Tiny_Byte(benchmark::State& state) {
662 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_DEFAULT);
663}
664static void Existing_Tiny_Byte(benchmark::State& state) {
665 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_DEFAULT);
666}
667static void New_Small_Byte(benchmark::State& state) {
668 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_DEFAULT);
669}
670static void Existing_Small_Byte(benchmark::State& state) {
671 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_DEFAULT);
672}
673static void New_Medium_Byte(benchmark::State& state) {
674 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_DEFAULT);
675}
676static void Existing_Medium_Byte(benchmark::State& state) {
677 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_DEFAULT);
678}
679static void New_Large_Byte(benchmark::State& state) {
680 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_DEFAULT);
681}
682static void Existing_Large_Byte(benchmark::State& state) {
683 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_DEFAULT);
684}
685 
686// Byte stream, 16KB HWM, no autoAllocate
687static void New_Tiny_Byte_HWM16K(benchmark::State& state) {
688 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_HWM_16K);
689}
690static void Existing_Tiny_Byte_HWM16K(benchmark::State& state) {
691 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_HWM_16K);
692}
693static void New_Small_Byte_HWM16K(benchmark::State& state) {
694 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_HWM_16K);
695}
696static void Existing_Small_Byte_HWM16K(benchmark::State& state) {
697 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_HWM_16K);
698}
699static void New_Medium_Byte_HWM16K(benchmark::State& state) {
700 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_HWM_16K);
701}
702static void Existing_Medium_Byte_HWM16K(benchmark::State& state) {
703 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_HWM_16K);
704}
705static void New_Large_Byte_HWM16K(benchmark::State& state) {
706 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_HWM_16K);
707}
708static void Existing_Large_Byte_HWM16K(benchmark::State& state) {
709 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_HWM_16K);
710}
711 
712// =============================================================================
713// Benchmark functions - Byte streams with autoAllocate=64KB
714// =============================================================================
715 
716// Byte stream, default HWM, autoAllocate=64KB
717static void New_Tiny_Byte_Auto64K(benchmark::State& state) {
718 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_AUTO_64K);
719}
720static void Existing_Tiny_Byte_Auto64K(benchmark::State& state) {
721 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_AUTO_64K);
722}
723static void New_Small_Byte_Auto64K(benchmark::State& state) {
724 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_AUTO_64K);
725}
726static void Existing_Small_Byte_Auto64K(benchmark::State& state) {
727 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_AUTO_64K);
728}
729static void New_Medium_Byte_Auto64K(benchmark::State& state) {
730 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_AUTO_64K);
731}
732static void Existing_Medium_Byte_Auto64K(benchmark::State& state) {
733 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_AUTO_64K);
734}
735static void New_Large_Byte_Auto64K(benchmark::State& state) {
736 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_AUTO_64K);
737}
738static void Existing_Large_Byte_Auto64K(benchmark::State& state) {
739 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_AUTO_64K);
740}
741 
742// Byte stream, 16KB HWM, autoAllocate=64KB
743static void New_Tiny_Byte_Auto64K_HWM16K(benchmark::State& state) {
744 benchNewApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
745}
746static void Existing_Tiny_Byte_Auto64K_HWM16K(benchmark::State& state) {
747 benchExistingApproachPumpTo(state, TINY_CHUNK_SIZE, TINY_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
748}
749static void New_Small_Byte_Auto64K_HWM16K(benchmark::State& state) {
750 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
751}
752static void Existing_Small_Byte_Auto64K_HWM16K(benchmark::State& state) {
753 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
754}
755static void New_Medium_Byte_Auto64K_HWM16K(benchmark::State& state) {
756 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
757}
758static void Existing_Medium_Byte_Auto64K_HWM16K(benchmark::State& state) {
759 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
760}
761static void New_Large_Byte_Auto64K_HWM16K(benchmark::State& state) {
762 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
763}
764static void Existing_Large_Byte_Auto64K_HWM16K(benchmark::State& state) {
765 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, BYTE_AUTO_64K_HWM_16K);
766}
767 
768// =============================================================================
769// Benchmark functions - Slow value streams (async with microtask delays)
770// =============================================================================
771 
772// Slow value stream - these produce one chunk per microtask, adding processing overhead.
773// Note: This does NOT trigger the adaptive read policy because microtask delays don't
774// cause early returns from pumpReadImpl. The policy would only activate with real I/O
775// latency that blocks the KJ event loop. See PUMP_PERFORMANCE_ANALYSIS.md section 9.
776static void New_Small_SlowValue(benchmark::State& state) {
777 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, SLOW_VALUE_DEFAULT);
778}
779static void Existing_Small_SlowValue(benchmark::State& state) {
780 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, SLOW_VALUE_DEFAULT);
781}
782static void New_Medium_SlowValue(benchmark::State& state) {
783 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, SLOW_VALUE_DEFAULT);
784}
785static void Existing_Medium_SlowValue(benchmark::State& state) {
786 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, SLOW_VALUE_DEFAULT);
787}
788 
789// =============================================================================
790// Benchmark functions - I/O latency streams (real KJ event loop yields)
791// =============================================================================
792 
793// I/O latency streams yield to the KJ event loop between chunks, simulating real network I/O.
794// This tests how the adaptive read policy responds to streams with actual I/O latency,
795// unlike microtask-based "slow" streams which complete within one JS event loop turn.
796//
797// Key differences from SlowValue:
798// - Each chunk requires a KJ event loop iteration (not just a microtask)
799// - pumpReadImpl returns early after each chunk
800// - Adaptive policy may switch to IMMEDIATE mode after observing small reads
801 
802// I/O latency value streams
803static void New_Small_IoLatencyValue(benchmark::State& state) {
804 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
805}
806static void Existing_Small_IoLatencyValue(benchmark::State& state) {
807 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
808}
809static void New_Medium_IoLatencyValue(benchmark::State& state) {
810 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
811}
812static void Existing_Medium_IoLatencyValue(benchmark::State& state) {
813 benchExistingApproachPumpTo(
814 state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
815}
816static void New_Large_IoLatencyValue(benchmark::State& state) {
817 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
818}
819static void Existing_Large_IoLatencyValue(benchmark::State& state) {
820 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, IO_LATENCY_VALUE_DEFAULT);
821}
822 
823// I/O latency byte streams
824static void New_Small_IoLatencyByte(benchmark::State& state) {
825 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
826}
827static void Existing_Small_IoLatencyByte(benchmark::State& state) {
828 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
829}
830static void New_Medium_IoLatencyByte(benchmark::State& state) {
831 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
832}
833static void Existing_Medium_IoLatencyByte(benchmark::State& state) {
834 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
835}
836static void New_Large_IoLatencyByte(benchmark::State& state) {
837 benchNewApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
838}
839static void Existing_Large_IoLatencyByte(benchmark::State& state) {
840 benchExistingApproachPumpTo(state, LARGE_CHUNK_SIZE, LARGE_NUM_CHUNKS, IO_LATENCY_BYTE_DEFAULT);
841}
842 
843// =============================================================================
844// Benchmark functions - Timed value streams (real timer-based delays)
845// =============================================================================
846 
847// These benchmarks use actual timer delays to simulate real network I/O.
848// Unlike evalLater() which resumes immediately, these represent real wall-clock time.
849// We test with small chunks to see how latency affects batching behavior.
850 
851// 10μs delay per chunk (1ms total for 100 chunks)
852static void New_Small_Timed10us(benchmark::State& state) {
853 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_10US);
854}
855static void Existing_Small_Timed10us(benchmark::State& state) {
856 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_10US);
857}
858 
859// 100μs delay per chunk (10ms total for 100 chunks)
860static void New_Small_Timed100us(benchmark::State& state) {
861 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_100US);
862}
863static void Existing_Small_Timed100us(benchmark::State& state) {
864 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_100US);
865}
866 
867// 1ms delay per chunk (100ms total for 100 chunks) - very slow, representative of slow network
868static void New_Small_Timed1ms(benchmark::State& state) {
869 benchNewApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_1MS);
870}
871static void Existing_Small_Timed1ms(benchmark::State& state) {
872 benchExistingApproachPumpTo(state, SMALL_CHUNK_SIZE, SMALL_NUM_CHUNKS, TIMED_VALUE_1MS);
873}
874 
875// Medium chunks with 100μs delay - tests larger chunk batching with latency
876static void New_Medium_Timed100us(benchmark::State& state) {
877 benchNewApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, TIMED_VALUE_100US);
878}
879static void Existing_Medium_Timed100us(benchmark::State& state) {
880 benchExistingApproachPumpTo(state, MEDIUM_CHUNK_SIZE, MEDIUM_NUM_CHUNKS, TIMED_VALUE_100US);
881}
882 
883// =============================================================================
884// Register benchmarks - organized by chunk size for easy comparison
885// =============================================================================
886 
887// Tiny chunks - all configurations
888WD_BENCHMARK(New_Tiny_Value);
889WD_BENCHMARK(Existing_Tiny_Value);
890WD_BENCHMARK(New_Tiny_Value_HWM16K);
891WD_BENCHMARK(Existing_Tiny_Value_HWM16K);
892WD_BENCHMARK(New_Tiny_Byte);
893WD_BENCHMARK(Existing_Tiny_Byte);
894WD_BENCHMARK(New_Tiny_Byte_HWM16K);
895WD_BENCHMARK(Existing_Tiny_Byte_HWM16K);
896WD_BENCHMARK(New_Tiny_Byte_Auto64K);
897WD_BENCHMARK(Existing_Tiny_Byte_Auto64K);
898WD_BENCHMARK(New_Tiny_Byte_Auto64K_HWM16K);
899WD_BENCHMARK(Existing_Tiny_Byte_Auto64K_HWM16K);
900 
901// Small chunks - all configurations
902WD_BENCHMARK(New_Small_Value);
903WD_BENCHMARK(Existing_Small_Value);
904WD_BENCHMARK(New_Small_Value_HWM16K);
905WD_BENCHMARK(Existing_Small_Value_HWM16K);
906WD_BENCHMARK(New_Small_Byte);
907WD_BENCHMARK(Existing_Small_Byte);
908WD_BENCHMARK(New_Small_Byte_HWM16K);
909WD_BENCHMARK(Existing_Small_Byte_HWM16K);
910WD_BENCHMARK(New_Small_Byte_Auto64K);
911WD_BENCHMARK(Existing_Small_Byte_Auto64K);
912WD_BENCHMARK(New_Small_Byte_Auto64K_HWM16K);
913WD_BENCHMARK(Existing_Small_Byte_Auto64K_HWM16K);
914 
915// Medium chunks - all configurations
916WD_BENCHMARK(New_Medium_Value);
917WD_BENCHMARK(Existing_Medium_Value);
918WD_BENCHMARK(New_Medium_Value_HWM16K);
919WD_BENCHMARK(Existing_Medium_Value_HWM16K);
920WD_BENCHMARK(New_Medium_Byte);
921WD_BENCHMARK(Existing_Medium_Byte);
922WD_BENCHMARK(New_Medium_Byte_HWM16K);
923WD_BENCHMARK(Existing_Medium_Byte_HWM16K);
924WD_BENCHMARK(New_Medium_Byte_Auto64K);
925WD_BENCHMARK(Existing_Medium_Byte_Auto64K);
926WD_BENCHMARK(New_Medium_Byte_Auto64K_HWM16K);
927WD_BENCHMARK(Existing_Medium_Byte_Auto64K_HWM16K);
928 
929// Large chunks - all configurations
930WD_BENCHMARK(New_Large_Value);
931WD_BENCHMARK(Existing_Large_Value);
932WD_BENCHMARK(New_Large_Value_HWM16K);
933WD_BENCHMARK(Existing_Large_Value_HWM16K);
934WD_BENCHMARK(New_Large_Byte);
935WD_BENCHMARK(Existing_Large_Byte);
936WD_BENCHMARK(New_Large_Byte_HWM16K);
937WD_BENCHMARK(Existing_Large_Byte_HWM16K);
938WD_BENCHMARK(New_Large_Byte_Auto64K);
939WD_BENCHMARK(Existing_Large_Byte_Auto64K);
940WD_BENCHMARK(New_Large_Byte_Auto64K_HWM16K);
941WD_BENCHMARK(Existing_Large_Byte_Auto64K_HWM16K);
942 
943// Slow value stream - async streams with microtask delays (tests batching overhead)
944WD_BENCHMARK(New_Small_SlowValue);
945WD_BENCHMARK(Existing_Small_SlowValue);
946WD_BENCHMARK(New_Medium_SlowValue);
947WD_BENCHMARK(Existing_Medium_SlowValue);
948 
949// I/O latency streams - real KJ event loop yields (simulates network I/O)
950// These test how the adaptive read policy behaves with actual I/O latency
951WD_BENCHMARK(New_Small_IoLatencyValue);
952WD_BENCHMARK(Existing_Small_IoLatencyValue);
953WD_BENCHMARK(New_Medium_IoLatencyValue);
954WD_BENCHMARK(Existing_Medium_IoLatencyValue);
955WD_BENCHMARK(New_Large_IoLatencyValue);
956WD_BENCHMARK(Existing_Large_IoLatencyValue);
957WD_BENCHMARK(New_Small_IoLatencyByte);
958WD_BENCHMARK(Existing_Small_IoLatencyByte);
959WD_BENCHMARK(New_Medium_IoLatencyByte);
960WD_BENCHMARK(Existing_Medium_IoLatencyByte);
961WD_BENCHMARK(New_Large_IoLatencyByte);
962WD_BENCHMARK(Existing_Large_IoLatencyByte);
963 
964// Timed stream benchmarks - uses real timers via useRealTimers=true in SetupParams.
965// These simulate actual blocking I/O with timer delays between chunks.
966WD_BENCHMARK(New_Small_Timed10us);
967WD_BENCHMARK(Existing_Small_Timed10us);
968WD_BENCHMARK(New_Small_Timed100us);
969WD_BENCHMARK(Existing_Small_Timed100us);
970WD_BENCHMARK(New_Small_Timed1ms);
971WD_BENCHMARK(Existing_Small_Timed1ms);
972WD_BENCHMARK(New_Medium_Timed100us);
973WD_BENCHMARK(Existing_Medium_Timed100us);
974 
975// =============================================================================
976// maxRead limit benchmarks - tests maxRead guard effectiveness
977// =============================================================================
978 
979// These benchmarks test the maxRead limit with large finite streams.
980// They demonstrate that maxRead properly limits how much data is read in a single
981// draining read operation, even when more data is synchronously available.
982//
983// The benchmarks compare different maxRead limits to show:
984// 1. How maxRead affects batching behavior and write coalescing
985// 2. Whether smaller maxRead values add overhead
986// 3. The trade-off between limiting reads vs throughput
987 
988// Configuration for streams with more chunks than we want to read
989static constexpr size_t LARGE_STREAM_CHUNKS = 10000; // 10K chunks available
990static constexpr size_t LARGE_STREAM_CHUNK_SIZE = 256; // 256 bytes per chunk = 2.56MB total
991 
992// Stream config for a large finite value stream
993static const StreamConfig LARGE_VALUE_STREAM{
994 .type = StreamType::VALUE,
995 .autoAllocateChunkSize = kj::none,
996 .highWaterMark = 0,
997};
998 
999// Benchmark using ReadableSourceKjAdapter::pumpTo with large streams (unlimited maxRead).
1000static void New_LargeStream_Value(benchmark::State& state) {
1001 benchNewApproachPumpTo(state, LARGE_STREAM_CHUNK_SIZE, LARGE_STREAM_CHUNKS, LARGE_VALUE_STREAM);
1002}
1003static void Existing_LargeStream_Value(benchmark::State& state) {
1004 benchExistingApproachPumpTo(
1005 state, LARGE_STREAM_CHUNK_SIZE, LARGE_STREAM_CHUNKS, LARGE_VALUE_STREAM);
1006}
1007 
1008// =============================================================================
1009// DrainingReader maxRead limit benchmarks
1010// =============================================================================
1011 
1012// These benchmarks directly test DrainingReader::read with different maxRead limits.
1013// Each iteration performs multiple reads until the stream is exhausted, allowing us
1014// to measure the overhead of different maxRead limits.
1015 
1016static void benchDrainingReaderMaxRead(benchmark::State& state,
1017 size_t chunkSize,
1018 size_t numChunks,
1019 size_t maxReadLimit,
1020 bool byteStream) {
1021 capnp::MallocMessageBuilder message;
1022 auto flags = message.initRoot<CompatibilityFlags>();
1023 flags.setStreamsJavaScriptControllers(true);
1024 // Use spec-compliant autoAllocateChunkSize behavior for DrainingReader benchmarks.
1025 flags.setNoAutoAllocateChunkSize(true);
1026 TestFixture fixture({.featureFlags = flags.asReader()});
1027 
1028 size_t totalBytes = chunkSize * numChunks;
1029 size_t readCount = 0;
1030 
1031 for (auto _: state) {
1032 benchChunkCounterStatic = 0;
1033 size_t bytesRead = 0;
1034 readCount = 0;
1035 
1036 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1037 auto stream = byteStream
1038 ? createByteStream(env.js, chunkSize, numChunks, kj::none, 0, &benchChunkCounterStatic)
1039 : createValueStream(env.js, chunkSize, numChunks, 0, &benchChunkCounterStatic);
1040 
1041 auto maybeReader = DrainingReader::create(env.js, *stream);
1042 KJ_ASSERT(maybeReader != kj::none, "Failed to create DrainingReader");
1043 auto& reader = KJ_ASSERT_NONNULL(maybeReader);
1044 
1045 // Read in chunks limited by maxReadLimit until stream is done
1046 bool done = false;
1047 while (!done) {
1048 auto jsPromise = reader->read(env.js, maxReadLimit)
1049 .then(env.js, [&](jsg::Lock& js, DrainingReadResult&& result) {
1050 for (auto& chunk: result.chunks) {
1051 bytesRead += chunk.size();
1052 }
1053 done = result.done;
1054 readCount++;
1055 });
1056 // Run the promise synchronously
1057 env.js.runMicrotasks();
1058 }
1059 
1060 reader->releaseLock(env.js);
1061 return kj::READY_NOW;
1062 });
1063 
1064 KJ_ASSERT(bytesRead == totalBytes, "Expected", totalBytes, "bytes but got", bytesRead);
1065 }
1066 
1067 state.SetBytesProcessed(state.iterations() * static_cast<int64_t>(totalBytes));
1068 state.counters["ReadCalls"] = benchmark::Counter(readCount, benchmark::Counter::kAvgIterations);
1069 state.counters["BytesPerRead"] =
1070 benchmark::Counter(totalBytes / readCount, benchmark::Counter::kAvgIterations);
1071}
1072 
1073// Value stream benchmarks with different maxRead limits
1074// Using 1000 chunks of 1KB each = 1MB total data
1075 
1076// maxRead = 16KB (small limit, many read calls)
1077static void DrainingRead_Value_MaxRead_16KB(benchmark::State& state) {
1078 benchDrainingReaderMaxRead(state, 1024, 1000, 16 * 1024, false);
1079}
1080 
1081// maxRead = 64KB
1082static void DrainingRead_Value_MaxRead_64KB(benchmark::State& state) {
1083 benchDrainingReaderMaxRead(state, 1024, 1000, 64 * 1024, false);
1084}
1085 
1086// maxRead = 256KB
1087static void DrainingRead_Value_MaxRead_256KB(benchmark::State& state) {
1088 benchDrainingReaderMaxRead(state, 1024, 1000, 256 * 1024, false);
1089}
1090 
1091// maxRead = 1MB (equals total data size)
1092static void DrainingRead_Value_MaxRead_1MB(benchmark::State& state) {
1093 benchDrainingReaderMaxRead(state, 1024, 1000, 1024 * 1024, false);
1094}
1095 
1096// maxRead = unlimited (default behavior)
1097static void DrainingRead_Value_MaxRead_Unlimited(benchmark::State& state) {
1098 benchDrainingReaderMaxRead(state, 1024, 1000, kj::maxValue, false);
1099}
1100 
1101// Byte stream benchmarks with different maxRead limits
1102static void DrainingRead_Byte_MaxRead_16KB(benchmark::State& state) {
1103 benchDrainingReaderMaxRead(state, 1024, 1000, 16 * 1024, true);
1104}
1105 
1106static void DrainingRead_Byte_MaxRead_64KB(benchmark::State& state) {
1107 benchDrainingReaderMaxRead(state, 1024, 1000, 64 * 1024, true);
1108}
1109 
1110static void DrainingRead_Byte_MaxRead_256KB(benchmark::State& state) {
1111 benchDrainingReaderMaxRead(state, 1024, 1000, 256 * 1024, true);
1112}
1113 
1114static void DrainingRead_Byte_MaxRead_1MB(benchmark::State& state) {
1115 benchDrainingReaderMaxRead(state, 1024, 1000, 1024 * 1024, true);
1116}
1117 
1118static void DrainingRead_Byte_MaxRead_Unlimited(benchmark::State& state) {
1119 benchDrainingReaderMaxRead(state, 1024, 1000, kj::maxValue, true);
1120}
1121 
1122// Small chunk benchmarks (64 bytes) - tests overhead with many small chunks
1123// 16000 chunks of 64 bytes = 1MB total
1124static void DrainingRead_SmallChunks_MaxRead_16KB(benchmark::State& state) {
1125 benchDrainingReaderMaxRead(state, 64, 16000, 16 * 1024, false);
1126}
1127 
1128static void DrainingRead_SmallChunks_MaxRead_64KB(benchmark::State& state) {
1129 benchDrainingReaderMaxRead(state, 64, 16000, 64 * 1024, false);
1130}
1131 
1132static void DrainingRead_SmallChunks_MaxRead_Unlimited(benchmark::State& state) {
1133 benchDrainingReaderMaxRead(state, 64, 16000, kj::maxValue, false);
1134}
1135 
1136// =============================================================================
1137// Large stream benchmarks - 10MB total to better exercise maxRead limits
1138// =============================================================================
1139 
1140// 10000 chunks of 1KB = 10MB total
1141static void DrainingRead_Large_MaxRead_64KB(benchmark::State& state) {
1142 benchDrainingReaderMaxRead(state, 1024, 10000, 64 * 1024, false);
1143}
1144 
1145static void DrainingRead_Large_MaxRead_256KB(benchmark::State& state) {
1146 benchDrainingReaderMaxRead(state, 1024, 10000, 256 * 1024, false);
1147}
1148 
1149static void DrainingRead_Large_MaxRead_1MB(benchmark::State& state) {
1150 benchDrainingReaderMaxRead(state, 1024, 10000, 1024 * 1024, false);
1151}
1152 
1153static void DrainingRead_Large_MaxRead_Unlimited(benchmark::State& state) {
1154 benchDrainingReaderMaxRead(state, 1024, 10000, kj::maxValue, false);
1155}
1156 
1157// Register large stream benchmarks
1158WD_BENCHMARK(New_LargeStream_Value);
1159WD_BENCHMARK(Existing_LargeStream_Value);
1160 
1161// Register maxRead limit benchmarks - value streams (1MB total, sync)
1162WD_BENCHMARK(DrainingRead_Value_MaxRead_16KB);
1163WD_BENCHMARK(DrainingRead_Value_MaxRead_64KB);
1164WD_BENCHMARK(DrainingRead_Value_MaxRead_256KB);
1165WD_BENCHMARK(DrainingRead_Value_MaxRead_1MB);
1166WD_BENCHMARK(DrainingRead_Value_MaxRead_Unlimited);
1167 
1168// Register maxRead limit benchmarks - byte streams (1MB total, sync)
1169WD_BENCHMARK(DrainingRead_Byte_MaxRead_16KB);
1170WD_BENCHMARK(DrainingRead_Byte_MaxRead_64KB);
1171WD_BENCHMARK(DrainingRead_Byte_MaxRead_256KB);
1172WD_BENCHMARK(DrainingRead_Byte_MaxRead_1MB);
1173WD_BENCHMARK(DrainingRead_Byte_MaxRead_Unlimited);
1174 
1175// Register small chunk benchmarks (1MB total, 64-byte chunks, sync)
1176WD_BENCHMARK(DrainingRead_SmallChunks_MaxRead_16KB);
1177WD_BENCHMARK(DrainingRead_SmallChunks_MaxRead_64KB);
1178WD_BENCHMARK(DrainingRead_SmallChunks_MaxRead_Unlimited);
1179 
1180// Register large stream benchmarks (10MB total, sync)
1181WD_BENCHMARK(DrainingRead_Large_MaxRead_64KB);
1182WD_BENCHMARK(DrainingRead_Large_MaxRead_256KB);
1183WD_BENCHMARK(DrainingRead_Large_MaxRead_1MB);
1184WD_BENCHMARK(DrainingRead_Large_MaxRead_Unlimited);
1185 
1186} // namespace
1187} // namespace workerd::api::streams