File
Blob: src/workerd/api/streams/readable-source-test.c++
| 1 | #include "readable-source.h" |
| 2 | |
| 3 | #include <workerd/api/streams/writable-sink.h> |
| 4 | #include <workerd/jsg/jsg-test.h> |
| 5 | #include <workerd/tests/test-fixture.h> |
| 6 | #include <workerd/util/own-util.h> |
| 7 | #include <workerd/util/stream-utils.h> |
| 8 | |
| 9 | #include <kj/async-io.h> |
| 10 | #include <kj/test.h> |
| 11 | |
| 12 | // We Thank Claude for Tests. |
| 13 | |
| 14 | namespace workerd::api::streams { |
| 15 | namespace { |
| 16 | |
| 17 | // Mock WritableSink for testing pumpTo functionality |
| 18 | class MockWritableSink final: public WritableSink { |
| 19 | public: |
| 20 | MockWritableSink() = default; |
| 21 | ~MockWritableSink() = default; |
| 22 | |
| 23 | kj::Promise<void> write(kj::ArrayPtr<const kj::byte> buffer) override { |
| 24 | writeCallCount++; |
| 25 | totalBytesWritten += buffer.size(); |
| 26 | writtenData.addAll(buffer); |
| 27 | |
| 28 | if (shouldFailWrite) { |
| 29 | KJ_FAIL_REQUIRE("Expected failure"); |
| 30 | } |
| 31 | |
| 32 | co_return; |
| 33 | } |
| 34 | |
| 35 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) override { |
| 36 | multiWriteCallCount++; |
| 37 | for (auto piece: pieces) { |
| 38 | totalBytesWritten += piece.size(); |
| 39 | writtenData.addAll(piece); |
| 40 | } |
| 41 | |
| 42 | if (shouldFailWrite) { |
| 43 | return KJ_EXCEPTION(FAILED, "Mock multi-write failure"); |
| 44 | } |
| 45 | |
| 46 | return kj::READY_NOW; |
| 47 | } |
| 48 | |
| 49 | kj::Promise<void> end() override { |
| 50 | endCallCount++; |
| 51 | isEnded = true; |
| 52 | |
| 53 | if (shouldFailEnd) { |
| 54 | KJ_FAIL_REQUIRE("Expected failure"); |
| 55 | } |
| 56 | |
| 57 | co_return; |
| 58 | } |
| 59 | |
| 60 | void abort(kj::Exception reason) override { |
| 61 | abortCallCount++; |
| 62 | isAborted = true; |
| 63 | abortReason = kj::mv(reason); |
| 64 | } |
| 65 | |
| 66 | rpc::StreamEncoding disownEncodingResponsibility() override { |
| 67 | disownCallCount++; |
| 68 | return encoding; |
| 69 | } |
| 70 | |
| 71 | rpc::StreamEncoding getEncoding() override { |
| 72 | return encoding; |
| 73 | } |
| 74 | |
| 75 | // Test state |
| 76 | uint32_t writeCallCount = 0; |
| 77 | uint32_t multiWriteCallCount = 0; |
| 78 | uint32_t endCallCount = 0; |
| 79 | uint32_t abortCallCount = 0; |
| 80 | uint32_t disownCallCount = 0; |
| 81 | |
| 82 | size_t totalBytesWritten = 0; |
| 83 | bool isEnded = false; |
| 84 | bool isAborted = false; |
| 85 | kj::Maybe<kj::Exception> abortReason; |
| 86 | |
| 87 | kj::Vector<kj::byte> writtenData; |
| 88 | rpc::StreamEncoding encoding = rpc::StreamEncoding::IDENTITY; |
| 89 | |
| 90 | // Control behavior |
| 91 | bool shouldFailWrite = false; |
| 92 | bool shouldFailEnd = false; |
| 93 | }; |
| 94 | |
| 95 | // Memory-based AsyncInputStream for factory function tests |
| 96 | class MemoryAsyncInputStream: public kj::AsyncInputStream { |
| 97 | public: |
| 98 | MemoryAsyncInputStream(kj::ArrayPtr<const kj::byte> data): data_(data) {} |
| 99 | |
| 100 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 101 | auto dest = kj::arrayPtr(static_cast<kj::byte*>(buffer), maxBytes); |
| 102 | size_t amount = kj::min(dest.size(), data_.size()); |
| 103 | dest.first(amount).copyFrom(data_.first(amount)); |
| 104 | data_ = data_.slice(amount); |
| 105 | return amount; |
| 106 | } |
| 107 | |
| 108 | kj::Maybe<uint64_t> tryGetLength() override { |
| 109 | return data_.size(); |
| 110 | } |
| 111 | |
| 112 | private: |
| 113 | kj::ArrayPtr<const kj::byte> data_; |
| 114 | }; |
| 115 | |
| 116 | class MemoryAsyncOutputStream final: public kj::AsyncOutputStream { |
| 117 | public: |
| 118 | kj::Promise<void> write(kj::ArrayPtr<const kj::byte> buffer) override { |
| 119 | data.addAll(buffer); |
| 120 | |
| 121 | if (writeShouldError) { |
| 122 | KJ_FAIL_REQUIRE("Expected failure"); |
| 123 | } |
| 124 | |
| 125 | co_return; |
| 126 | } |
| 127 | |
| 128 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) override { |
| 129 | for (auto piece: pieces) { |
| 130 | data.addAll(piece); |
| 131 | } |
| 132 | |
| 133 | if (writeShouldError) { |
| 134 | KJ_FAIL_REQUIRE("Expected failure"); |
| 135 | } |
| 136 | |
| 137 | co_return; |
| 138 | } |
| 139 | |
| 140 | kj::Promise<void> whenWriteDisconnected() override { |
| 141 | return kj::NEVER_DONE; |
| 142 | } |
| 143 | |
| 144 | bool writeShouldError = false; |
| 145 | |
| 146 | kj::Vector<kj::byte> data; |
| 147 | }; |
| 148 | |
| 149 | // ====================================================================================== |
| 150 | // Core ReadableSource Interface Tests |
| 151 | |
| 152 | KJ_TEST("ReadableSource basic read operations (full)") { |
| 153 | TestFixture fixture; |
| 154 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 155 | |
| 156 | MemoryAsyncInputStream input(testData); |
| 157 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 158 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 159 | |
| 160 | KJ_ASSERT(source->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 161 | |
| 162 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 163 | kj::FixedArray<kj::byte, 15> buffer; |
| 164 | |
| 165 | KJ_ASSERT(KJ_ASSERT_NONNULL(source->tryGetLength(rpc::StreamEncoding::IDENTITY)) == 10); |
| 166 | KJ_ASSERT(source->tryGetLength(rpc::StreamEncoding::GZIP) == kj::none); |
| 167 | |
| 168 | // Read at least 5 bytes, at most 15. |
| 169 | auto bytesRead = co_await source->read(buffer, 5); |
| 170 | KJ_ASSERT(bytesRead == 10); // Should read all available data |
| 171 | |
| 172 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == testData); |
| 173 | |
| 174 | // Next read should return nothing |
| 175 | bytesRead = co_await source->read(buffer, 1); |
| 176 | KJ_ASSERT(bytesRead == 0); // EOF |
| 177 | }); |
| 178 | } |
| 179 | |
| 180 | KJ_TEST("ReadableSource basic read operations (partial)") { |
| 181 | TestFixture fixture; |
| 182 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 183 | |
| 184 | MemoryAsyncInputStream input(testData); |
| 185 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 186 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 187 | |
| 188 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 189 | kj::FixedArray<kj::byte, 5> buffer; |
| 190 | |
| 191 | KJ_ASSERT(KJ_ASSERT_NONNULL(source->tryGetLength(rpc::StreamEncoding::IDENTITY)) == 10); |
| 192 | KJ_ASSERT(source->tryGetLength(rpc::StreamEncoding::GZIP) == kj::none); |
| 193 | |
| 194 | // Read at most 5 bytes |
| 195 | auto bytesRead = co_await source->read(buffer, 5); |
| 196 | KJ_ASSERT(bytesRead == 5); // Should read all available data |
| 197 | |
| 198 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == kj::arrayPtr(testData, 5)); |
| 199 | |
| 200 | // Next read should return 5 byte |
| 201 | bytesRead = co_await source->read(buffer, 1); |
| 202 | KJ_ASSERT(bytesRead == 5); // EOF |
| 203 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == kj::arrayPtr(testData + 5, 5)); |
| 204 | |
| 205 | // Next read should return nothing |
| 206 | bytesRead = co_await source->read(buffer, 1); |
| 207 | KJ_ASSERT(bytesRead == 0); // EOF |
| 208 | }); |
| 209 | } |
| 210 | |
| 211 | KJ_TEST("ReadableSource concurrent reads forbidden") { |
| 212 | TestFixture fixture; |
| 213 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 214 | |
| 215 | MemoryAsyncInputStream input(testData); |
| 216 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 217 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 218 | |
| 219 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 220 | kj::FixedArray<kj::byte, 5> buffer; |
| 221 | |
| 222 | auto readPromise1 = source->read(buffer, 5); |
| 223 | auto readPromise2 = source->read(buffer, 5); |
| 224 | |
| 225 | try { |
| 226 | co_await readPromise2; |
| 227 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 228 | } catch (...) { |
| 229 | auto exception = kj::getCaughtExceptionAsKj(); |
| 230 | KJ_ASSERT(exception.getDescription().contains("already being read")); |
| 231 | } |
| 232 | |
| 233 | // But the first read should still succeed. |
| 234 | auto bytesRead = co_await readPromise1; |
| 235 | KJ_ASSERT(bytesRead == 5); |
| 236 | }); |
| 237 | } |
| 238 | |
| 239 | // ====================================================================================== |
| 240 | // PumpTo Tests |
| 241 | |
| 242 | KJ_TEST("ReadableSource pumpTo with end") { |
| 243 | TestFixture fixture; |
| 244 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 245 | |
| 246 | MemoryAsyncInputStream input(testData); |
| 247 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 248 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 249 | |
| 250 | MockWritableSink sink; |
| 251 | |
| 252 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 253 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 254 | KJ_ASSERT(sink.totalBytesWritten == 10); |
| 255 | KJ_ASSERT(sink.isEnded); |
| 256 | KJ_ASSERT(sink.writtenData == testData); |
| 257 | }); |
| 258 | } |
| 259 | |
| 260 | KJ_TEST("ReadableSource pumpTo without end") { |
| 261 | TestFixture fixture; |
| 262 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 263 | |
| 264 | MemoryAsyncInputStream input(testData); |
| 265 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 266 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 267 | |
| 268 | MockWritableSink sink; |
| 269 | |
| 270 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 271 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::NO)); |
| 272 | KJ_ASSERT(sink.totalBytesWritten == 10); |
| 273 | KJ_ASSERT(!sink.isEnded); |
| 274 | KJ_ASSERT(sink.writtenData == testData); |
| 275 | }); |
| 276 | } |
| 277 | |
| 278 | KJ_TEST("ReadableSource large pumpTo with end") { |
| 279 | TestFixture fixture; |
| 280 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 281 | testData.asPtr().fill(42); |
| 282 | |
| 283 | MemoryAsyncInputStream input(testData); |
| 284 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 285 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 286 | |
| 287 | MockWritableSink sink; |
| 288 | |
| 289 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 290 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 291 | KJ_ASSERT(sink.totalBytesWritten == 52 * 1024); |
| 292 | KJ_ASSERT(sink.isEnded); |
| 293 | KJ_ASSERT(sink.writtenData == testData); |
| 294 | }); |
| 295 | } |
| 296 | |
| 297 | KJ_TEST("ReadableSource large pumpTo canceled") { |
| 298 | TestFixture fixture; |
| 299 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 300 | testData.asPtr().fill(42); |
| 301 | |
| 302 | MemoryAsyncInputStream input(testData); |
| 303 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 304 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 305 | |
| 306 | MockWritableSink sink; |
| 307 | |
| 308 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 309 | auto promise = |
| 310 | environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 311 | source->cancel(KJ_EXCEPTION(FAILED, "test abort")); |
| 312 | try { |
| 313 | co_await promise; |
| 314 | } catch (...) { |
| 315 | auto exception = kj::getCaughtExceptionAsKj(); |
| 316 | KJ_ASSERT(exception.getDescription().contains("test abort")); |
| 317 | } |
| 318 | }); |
| 319 | } |
| 320 | |
| 321 | KJ_TEST("ReadableSource large pumpTo canceled before") { |
| 322 | TestFixture fixture; |
| 323 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 324 | testData.asPtr().fill(42); |
| 325 | |
| 326 | MemoryAsyncInputStream input(testData); |
| 327 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 328 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 329 | |
| 330 | MockWritableSink sink; |
| 331 | |
| 332 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 333 | source->cancel(KJ_EXCEPTION(FAILED, "test abort")); |
| 334 | auto promise = |
| 335 | environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 336 | try { |
| 337 | co_await promise; |
| 338 | } catch (...) { |
| 339 | auto exception = kj::getCaughtExceptionAsKj(); |
| 340 | KJ_ASSERT(exception.getDescription().contains("test abort")); |
| 341 | } |
| 342 | }); |
| 343 | } |
| 344 | |
| 345 | KJ_TEST("ReadableSource large pumpTo closed") { |
| 346 | TestFixture fixture; |
| 347 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 348 | testData.asPtr().fill(42); |
| 349 | |
| 350 | MemoryAsyncInputStream input(testData); |
| 351 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 352 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 353 | |
| 354 | MockWritableSink sink; |
| 355 | |
| 356 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 357 | auto& context = environment.context; |
| 358 | co_await source->readAllBytes(kj::maxValue); |
| 359 | co_await context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 360 | KJ_ASSERT(sink.totalBytesWritten == 0); |
| 361 | }); |
| 362 | } |
| 363 | |
| 364 | KJ_TEST("ReadableSource large pumpTo, concurrent read fails") { |
| 365 | TestFixture fixture; |
| 366 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 367 | testData.asPtr().fill(42); |
| 368 | |
| 369 | MemoryAsyncInputStream input(testData); |
| 370 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 371 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 372 | |
| 373 | MockWritableSink sink; |
| 374 | |
| 375 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 376 | auto promise = |
| 377 | environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 378 | |
| 379 | // Concurrent read should fail. |
| 380 | try { |
| 381 | co_await source->readAllBytes(kj::maxValue); |
| 382 | } catch (...) { |
| 383 | auto exception = kj::getCaughtExceptionAsKj(); |
| 384 | KJ_ASSERT(exception.getDescription().contains("already being read")); |
| 385 | } |
| 386 | |
| 387 | // But the pump should still succeed. |
| 388 | co_await promise; |
| 389 | KJ_ASSERT(sink.totalBytesWritten == 52 * 1024); |
| 390 | |
| 391 | // And we can read again afterwards, but will be at EOF. |
| 392 | auto allBytes = co_await source->readAllBytes(kj::maxValue); |
| 393 | KJ_ASSERT(allBytes.size() == 0); |
| 394 | }); |
| 395 | } |
| 396 | |
| 397 | // ====================================================================================== |
| 398 | // Read all tests |
| 399 | |
| 400 | KJ_TEST("ReadableSource read all bytes (small)") { |
| 401 | TestFixture fixture; |
| 402 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 403 | |
| 404 | MemoryAsyncInputStream input(testData); |
| 405 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 406 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 407 | |
| 408 | KJ_ASSERT(source->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 409 | |
| 410 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 411 | auto allBytes = co_await source->readAllBytes(100); |
| 412 | KJ_ASSERT(allBytes.size() == 10); |
| 413 | KJ_ASSERT(allBytes.asPtr() == kj::ArrayPtr<const kj::byte>(testData, 10)); |
| 414 | }); |
| 415 | } |
| 416 | |
| 417 | KJ_TEST("ReadableSource read all bytes (large)") { |
| 418 | TestFixture fixture; |
| 419 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 420 | testData.asPtr().fill(42); |
| 421 | |
| 422 | MemoryAsyncInputStream input(testData); |
| 423 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 424 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 425 | |
| 426 | KJ_ASSERT(source->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 427 | |
| 428 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 429 | auto allBytes = co_await source->readAllBytes((52 * 1024) + 1); |
| 430 | KJ_ASSERT(allBytes.size() == 52 * 1024); |
| 431 | KJ_ASSERT(allBytes.asPtr() == testData); |
| 432 | }); |
| 433 | } |
| 434 | |
| 435 | KJ_TEST("ReadableSource read all text (small)") { |
| 436 | TestFixture fixture; |
| 437 | kj::byte testData[] = {'a', 'b', 'c', 'd', 'e', 'f', 'g', 'h', 'i', 'j'}; |
| 438 | |
| 439 | MemoryAsyncInputStream input(testData); |
| 440 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 441 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 442 | |
| 443 | KJ_ASSERT(source->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 444 | |
| 445 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 446 | auto allText = co_await source->readAllBytes(100); |
| 447 | KJ_ASSERT(allText.size() == 10); |
| 448 | KJ_ASSERT(allText.asPtr() == kj::ArrayPtr<const kj::byte>(testData, 10)); |
| 449 | }); |
| 450 | } |
| 451 | |
| 452 | KJ_TEST("ReadableSource read all text (large)") { |
| 453 | TestFixture fixture; |
| 454 | auto testData = kj::heapArray<char>(52 * 1024); |
| 455 | testData.asPtr().fill('a'); |
| 456 | |
| 457 | MemoryAsyncInputStream input(testData.asBytes()); |
| 458 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 459 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 460 | |
| 461 | KJ_ASSERT(source->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 462 | |
| 463 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 464 | auto allText = co_await source->readAllText((52 * 1024) + 1); |
| 465 | KJ_ASSERT(allText.size() == 52 * 1024); |
| 466 | KJ_ASSERT(allText.asPtr() == testData); |
| 467 | }); |
| 468 | } |
| 469 | |
| 470 | KJ_TEST("ReadableSource read all aborted (after read)") { |
| 471 | TestFixture fixture; |
| 472 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 473 | testData.asPtr().fill(42); |
| 474 | |
| 475 | MemoryAsyncInputStream input(testData); |
| 476 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 477 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 478 | |
| 479 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 480 | auto promise = source->readAllBytes(52 * 1024); |
| 481 | source->cancel(KJ_EXCEPTION(FAILED, "test abort")); |
| 482 | try { |
| 483 | co_await promise; |
| 484 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 485 | } catch (...) { |
| 486 | auto exception = kj::getCaughtExceptionAsKj(); |
| 487 | KJ_ASSERT(exception.getDescription().contains("test abort")); |
| 488 | } |
| 489 | }); |
| 490 | } |
| 491 | |
| 492 | KJ_TEST("ReadableSource read all aborted (prior to read)") { |
| 493 | TestFixture fixture; |
| 494 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 495 | testData.asPtr().fill(42); |
| 496 | |
| 497 | MemoryAsyncInputStream input(testData); |
| 498 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 499 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 500 | |
| 501 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 502 | source->cancel(KJ_EXCEPTION(FAILED, "test abort")); |
| 503 | auto promise = source->readAllBytes(52 * 1024); |
| 504 | try { |
| 505 | co_await promise; |
| 506 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 507 | } catch (...) { |
| 508 | auto exception = kj::getCaughtExceptionAsKj(); |
| 509 | KJ_ASSERT(exception.getDescription().contains("test abort")); |
| 510 | } |
| 511 | }); |
| 512 | } |
| 513 | |
| 514 | KJ_TEST("ReadableSource read all aborted (dropped)") { |
| 515 | TestFixture fixture; |
| 516 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 517 | testData.asPtr().fill(42); |
| 518 | |
| 519 | MemoryAsyncInputStream input(testData); |
| 520 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 521 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 522 | |
| 523 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 524 | auto promise = source->readAllBytes(52 * 1024); |
| 525 | { auto freeme = kj::mv(source); } |
| 526 | try { |
| 527 | co_await promise; |
| 528 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 529 | } catch (...) { |
| 530 | auto exception = kj::getCaughtExceptionAsKj(); |
| 531 | KJ_ASSERT(exception.getDescription().contains("stream was dropped")); |
| 532 | } |
| 533 | }); |
| 534 | } |
| 535 | |
| 536 | // ====================================================================================== |
| 537 | // Tee tests |
| 538 | |
| 539 | KJ_TEST("ReadableSource tee (small, no limit)") { |
| 540 | TestFixture fixture; |
| 541 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 542 | |
| 543 | MemoryAsyncInputStream input(testData); |
| 544 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 545 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 546 | |
| 547 | auto tee = source->tee(200); |
| 548 | auto branch1 = kj::mv(tee.branch1); |
| 549 | auto branch2 = kj::mv(tee.branch2); |
| 550 | |
| 551 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 552 | auto allBytes1 = co_await branch1->readAllBytes(100); |
| 553 | auto allBytes2 = co_await branch2->readAllBytes(100); |
| 554 | KJ_ASSERT(allBytes1.size() == 10); |
| 555 | KJ_ASSERT(allBytes1 == testData); |
| 556 | KJ_ASSERT(allBytes2.size() == 10); |
| 557 | KJ_ASSERT(allBytes2 == testData); |
| 558 | |
| 559 | // Original source should be closed and return EOF |
| 560 | kj::FixedArray<kj::byte, 10> buffer; |
| 561 | auto bytesRead = co_await source->read(buffer, 1); |
| 562 | KJ_ASSERT(bytesRead == 0); // EOF |
| 563 | }); |
| 564 | } |
| 565 | |
| 566 | KJ_TEST("ReadableSource tee (small, no limit, independent)") { |
| 567 | TestFixture fixture; |
| 568 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 569 | |
| 570 | MemoryAsyncInputStream input(testData); |
| 571 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 572 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 573 | |
| 574 | auto tee = source->tee(200); |
| 575 | auto branch1 = kj::mv(tee.branch1); |
| 576 | auto branch2 = kj::mv(tee.branch2); |
| 577 | branch2->cancel(KJ_EXCEPTION(FAILED, "test abort")); |
| 578 | |
| 579 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 580 | auto allBytes1 = co_await branch1->readAllBytes(100); |
| 581 | |
| 582 | try { |
| 583 | co_await branch2->readAllBytes(100); |
| 584 | } catch (...) { |
| 585 | auto exception = kj::getCaughtExceptionAsKj(); |
| 586 | KJ_ASSERT(exception.getDescription().contains("test abort")); |
| 587 | } |
| 588 | KJ_ASSERT(allBytes1.size() == 10); |
| 589 | KJ_ASSERT(allBytes1.asPtr() == testData); |
| 590 | |
| 591 | // Original source should be closed and return EOF |
| 592 | kj::FixedArray<kj::byte, 10> buffer; |
| 593 | auto bytesRead = co_await source->read(buffer, 1); |
| 594 | KJ_ASSERT(bytesRead == 0); // EOF |
| 595 | }); |
| 596 | } |
| 597 | |
| 598 | KJ_TEST("ReadableSource tee (large, no limit)") { |
| 599 | TestFixture fixture; |
| 600 | auto testData = kj::heapArray<kj::byte>(52 * 1024); |
| 601 | testData.asPtr().fill(42); |
| 602 | |
| 603 | MemoryAsyncInputStream input(testData); |
| 604 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 605 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 606 | |
| 607 | auto tee = source->tee(0xffffffff); |
| 608 | auto branch1 = kj::mv(tee.branch1); |
| 609 | auto branch2 = kj::mv(tee.branch2); |
| 610 | |
| 611 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 612 | auto allBytes1 = co_await branch1->readAllBytes(kj::maxValue); |
| 613 | auto allBytes2 = co_await branch2->readAllBytes(kj::maxValue); |
| 614 | KJ_ASSERT(allBytes1.size() == 52 * 1024); |
| 615 | KJ_ASSERT(allBytes1 == testData); |
| 616 | KJ_ASSERT(allBytes2.size() == 52 * 1024); |
| 617 | KJ_ASSERT(allBytes2 == testData); |
| 618 | |
| 619 | // Original source should be closed and return EOF |
| 620 | kj::FixedArray<kj::byte, 10> buffer; |
| 621 | auto bytesRead = co_await source->read(buffer, 1); |
| 622 | KJ_ASSERT(bytesRead == 0); // EOF |
| 623 | }); |
| 624 | } |
| 625 | |
| 626 | KJ_TEST("ReadableSource tee (large, buffer limit)") { |
| 627 | TestFixture fixture; |
| 628 | auto testData = kj::heapArray<kj::byte>(1024); |
| 629 | testData.asPtr().fill(42); |
| 630 | |
| 631 | MemoryAsyncInputStream input(testData); |
| 632 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 633 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 634 | |
| 635 | auto tee = source->tee(100); |
| 636 | auto branch1 = kj::mv(tee.branch1); |
| 637 | auto branch2 = kj::mv(tee.branch2); |
| 638 | |
| 639 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 640 | try { |
| 641 | co_await branch1->readAllBytes(kj::maxValue); |
| 642 | } catch (...) { |
| 643 | auto exception = kj::getCaughtExceptionAsKj(); |
| 644 | KJ_ASSERT(exception.getDescription().contains("buffer limit exceeded")); |
| 645 | } |
| 646 | }); |
| 647 | } |
| 648 | |
| 649 | KJ_TEST("ReadableSource after read") { |
| 650 | TestFixture fixture; |
| 651 | auto testData = kj::heapArray<kj::byte>(1024); |
| 652 | testData.asPtr().fill(42); |
| 653 | |
| 654 | MemoryAsyncInputStream input(testData); |
| 655 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 656 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 657 | |
| 658 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 659 | kj::FixedArray<kj::byte, 512> buffer; |
| 660 | buffer.asPtr().first(512).fill(0); |
| 661 | buffer.asPtr().slice(512).fill(1); |
| 662 | auto bytesRead = co_await source->read(buffer, 1); |
| 663 | KJ_ASSERT(bytesRead == 512); |
| 664 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == testData.asPtr().first(bytesRead)); |
| 665 | |
| 666 | auto tee = source->tee(0xffffffff); |
| 667 | auto branch1 = kj::mv(tee.branch1); |
| 668 | auto branch2 = kj::mv(tee.branch2); |
| 669 | |
| 670 | // Each branch should get the remaining data |
| 671 | auto allBytes1 = co_await branch1->readAllBytes(kj::maxValue); |
| 672 | auto allBytes2 = co_await branch2->readAllBytes(kj::maxValue); |
| 673 | KJ_ASSERT(allBytes1.size() == 512); |
| 674 | KJ_ASSERT(allBytes1 == testData.asPtr().slice(512)); |
| 675 | KJ_ASSERT(allBytes2.size() == 512); |
| 676 | KJ_ASSERT(allBytes2 == testData.asPtr().slice(512)); |
| 677 | }); |
| 678 | } |
| 679 | |
| 680 | // ====================================================================================== |
| 681 | // ReadableSourceWrapper Tests |
| 682 | |
| 683 | KJ_TEST("ReadableSourceWrapper delegation") { |
| 684 | TestFixture fixture; |
| 685 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 686 | |
| 687 | MemoryAsyncInputStream input(testData); |
| 688 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 689 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 690 | |
| 691 | class TestWrapper: public ReadableSourceWrapper { |
| 692 | public: |
| 693 | TestWrapper(kj::Own<ReadableSource> inner): ReadableSourceWrapper(kj::mv(inner)) {} |
| 694 | }; |
| 695 | auto wrapper = kj::heap<TestWrapper>(kj::mv(source)); |
| 696 | |
| 697 | KJ_ASSERT(wrapper->getEncoding() == rpc::StreamEncoding::IDENTITY); |
| 698 | |
| 699 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 700 | kj::FixedArray<kj::byte, 15> buffer; |
| 701 | |
| 702 | KJ_ASSERT(KJ_ASSERT_NONNULL(wrapper->tryGetLength(rpc::StreamEncoding::IDENTITY)) == 10); |
| 703 | KJ_ASSERT(wrapper->tryGetLength(rpc::StreamEncoding::GZIP) == kj::none); |
| 704 | |
| 705 | // Read at least 5 bytes, at most 15. |
| 706 | auto bytesRead = co_await wrapper->read(buffer, 5); |
| 707 | KJ_ASSERT(bytesRead == 10); // Should read all available data |
| 708 | |
| 709 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == testData); |
| 710 | |
| 711 | // Next read should return nothing |
| 712 | bytesRead = co_await wrapper->read(buffer, 1); |
| 713 | KJ_ASSERT(bytesRead == 0); // EOF |
| 714 | }); |
| 715 | } |
| 716 | |
| 717 | KJ_TEST("ReadableSourceWrapper release") { |
| 718 | TestFixture fixture; |
| 719 | kj::byte testData[] = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}; |
| 720 | |
| 721 | MemoryAsyncInputStream input(testData); |
| 722 | auto fakeOwn = kj::Own<MemoryAsyncInputStream>(&input, kj::NullDisposer::instance); |
| 723 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 724 | |
| 725 | class TestWrapper: public ReadableSourceWrapper { |
| 726 | public: |
| 727 | TestWrapper(kj::Own<ReadableSource> inner): ReadableSourceWrapper(kj::mv(inner)) {} |
| 728 | }; |
| 729 | auto wrapper = kj::heap<TestWrapper>(kj::mv(source)); |
| 730 | |
| 731 | source = wrapper->release(); |
| 732 | |
| 733 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 734 | kj::FixedArray<kj::byte, 15> buffer; |
| 735 | |
| 736 | // Using the wrapper shuld fail. |
| 737 | try { |
| 738 | co_await wrapper->read(buffer, 1); |
| 739 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 740 | } catch (...) { |
| 741 | auto exception = kj::getCaughtExceptionAsKj(); |
| 742 | KJ_ASSERT(exception.getDescription().contains("inner != nullptr")); |
| 743 | } |
| 744 | |
| 745 | KJ_ASSERT(KJ_ASSERT_NONNULL(source->tryGetLength(rpc::StreamEncoding::IDENTITY)) == 10); |
| 746 | KJ_ASSERT(source->tryGetLength(rpc::StreamEncoding::GZIP) == kj::none); |
| 747 | |
| 748 | // Read at least 5 bytes, at most 15. |
| 749 | auto bytesRead = co_await source->read(buffer, 5); |
| 750 | KJ_ASSERT(bytesRead == 10); // Should read all available data |
| 751 | |
| 752 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == testData); |
| 753 | |
| 754 | // Next read should return nothing |
| 755 | bytesRead = co_await source->read(buffer, 1); |
| 756 | KJ_ASSERT(bytesRead == 0); // EOF |
| 757 | }); |
| 758 | } |
| 759 | |
| 760 | // ====================================================================================== |
| 761 | // Factory Function Tests |
| 762 | |
| 763 | KJ_TEST("newClosedReadableSource") { |
| 764 | TestFixture fixture; |
| 765 | auto source = newClosedReadableSource(); |
| 766 | |
| 767 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 768 | kj::byte buffer[10]; |
| 769 | auto bytesRead = co_await source->read(buffer, 1); |
| 770 | KJ_ASSERT(bytesRead == 0); // EOF |
| 771 | }); |
| 772 | } |
| 773 | |
| 774 | KJ_TEST("newErroredReadableSource") { |
| 775 | TestFixture fixture; |
| 776 | auto exception = KJ_EXCEPTION(FAILED, "test error"); |
| 777 | auto source = newErroredReadableSource(exception.clone()); |
| 778 | |
| 779 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 780 | kj::byte buffer[10]; |
| 781 | try { |
| 782 | co_await source->read(buffer, 1); |
| 783 | KJ_FAIL_REQUIRE("was expected to throw"); |
| 784 | } catch (...) { |
| 785 | auto caught = kj::getCaughtExceptionAsKj(); |
| 786 | KJ_ASSERT(caught.getDescription().contains("test error")); |
| 787 | } |
| 788 | }); |
| 789 | } |
| 790 | |
| 791 | KJ_TEST("newReadableSourceFromBytes (copy)") { |
| 792 | TestFixture fixture; |
| 793 | kj::byte testData[] = {1, 2, 3, 4, 5}; |
| 794 | auto source = newReadableSourceFromBytes(kj::ArrayPtr<const kj::byte>(testData, 5)); |
| 795 | testData[0] = 42; // Modify original to ensure copy was made. |
| 796 | |
| 797 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 798 | kj::FixedArray<kj::byte, 5> buffer; |
| 799 | auto bytesRead = co_await source->read(buffer, 1); |
| 800 | KJ_ASSERT(bytesRead == 5); |
| 801 | KJ_ASSERT(buffer[0] == 1); // Original data |
| 802 | |
| 803 | // Next read should return nothing |
| 804 | bytesRead = co_await source->read(buffer, 1); |
| 805 | KJ_ASSERT(bytesRead == 0); // EOF |
| 806 | }); |
| 807 | } |
| 808 | |
| 809 | KJ_TEST("newReadableSourceFromBytes (owned)") { |
| 810 | TestFixture fixture; |
| 811 | auto ownedData = kj::heapArray<kj::byte>(5); |
| 812 | ownedData.asPtr().fill(0); |
| 813 | auto ptr = ownedData.asPtr(); |
| 814 | auto source = newReadableSourceFromBytes(ptr, kj::heap(kj::mv(ownedData))); |
| 815 | ptr[0] = 42; // Modify original to ensure copy was made. |
| 816 | |
| 817 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 818 | kj::FixedArray<kj::byte, 5> buffer; |
| 819 | auto bytesRead = co_await source->read(buffer, 1); |
| 820 | KJ_ASSERT(bytesRead == 5); |
| 821 | KJ_ASSERT(buffer[0] == 42); // Modified data |
| 822 | |
| 823 | // Next read should return nothing |
| 824 | bytesRead = co_await source->read(buffer, 1); |
| 825 | KJ_ASSERT(bytesRead == 0); // EOF |
| 826 | }); |
| 827 | } |
| 828 | |
| 829 | KJ_TEST("newReadableSourceFromDelegate") { |
| 830 | TestFixture fixture; |
| 831 | kj::byte testData[] = {1, 2, 3, 4, 5}; |
| 832 | size_t position = 0; |
| 833 | |
| 834 | auto producer = [&](kj::ArrayPtr<kj::byte> buffer, size_t minBytes) -> kj::Promise<size_t> { |
| 835 | if (position >= 5) { |
| 836 | return static_cast<size_t>(0); // EOF |
| 837 | } |
| 838 | |
| 839 | size_t available = 5 - position; |
| 840 | size_t toRead = kj::min(available, buffer.size()); |
| 841 | |
| 842 | if (toRead > 0) { |
| 843 | memcpy(buffer.begin(), testData + position, toRead); |
| 844 | position += toRead; |
| 845 | } |
| 846 | |
| 847 | return toRead; |
| 848 | }; |
| 849 | |
| 850 | auto source = newReadableSourceFromProducer(kj::mv(producer), 5); |
| 851 | |
| 852 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 853 | kj::FixedArray<kj::byte, 5> buffer; |
| 854 | auto bytesRead = co_await source->read(buffer, 1); |
| 855 | KJ_ASSERT(bytesRead == 5); |
| 856 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == kj::ArrayPtr<const kj::byte>(testData, 5)); |
| 857 | |
| 858 | // Next read should return nothing |
| 859 | bytesRead = co_await source->read(buffer, 1); |
| 860 | KJ_ASSERT(bytesRead == 0); // EOF |
| 861 | }); |
| 862 | } |
| 863 | |
| 864 | KJ_TEST("newReadableSourceFromDelegate (not enough bytes)") { |
| 865 | TestFixture fixture; |
| 866 | kj::byte testData[] = {1, 2, 3, 4, 5}; |
| 867 | size_t position = 0; |
| 868 | |
| 869 | auto producer = [&](kj::ArrayPtr<kj::byte> buffer, size_t minBytes) -> kj::Promise<size_t> { |
| 870 | if (position >= 5) { |
| 871 | return static_cast<size_t>(0); // EOF |
| 872 | } |
| 873 | |
| 874 | size_t available = 5 - position; |
| 875 | size_t toRead = kj::min(available, buffer.size()); |
| 876 | |
| 877 | if (toRead > 0) { |
| 878 | memcpy(buffer.begin(), testData + position, toRead); |
| 879 | position += toRead; |
| 880 | } |
| 881 | |
| 882 | return toRead; |
| 883 | }; |
| 884 | |
| 885 | auto source = newReadableSourceFromProducer(kj::mv(producer), 10); |
| 886 | |
| 887 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 888 | kj::FixedArray<kj::byte, 5> buffer; |
| 889 | auto bytesRead = co_await source->read(buffer, 1); |
| 890 | KJ_ASSERT(bytesRead == 5); |
| 891 | KJ_ASSERT(buffer.asPtr().first(bytesRead) == kj::ArrayPtr<const kj::byte>(testData, 5)); |
| 892 | |
| 893 | // Next read should fail since producer did not produce the expected number of bytes |
| 894 | try { |
| 895 | co_await source->read(buffer, 1); |
| 896 | } catch (...) { |
| 897 | auto exception = kj::getCaughtExceptionAsKj(); |
| 898 | KJ_ASSERT(exception.getDescription().contains("ended stream early")); |
| 899 | } |
| 900 | }); |
| 901 | } |
| 902 | |
| 903 | // ====================================================================================== |
| 904 | // Gzip encoding |
| 905 | |
| 906 | KJ_TEST("Gzip encoded stream") { |
| 907 | TestFixture fixture; |
| 908 | static constexpr kj::byte data[] = {31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 43, 206, 207, 77, 85, 72, 73, |
| 909 | 44, 73, 84, 40, 201, 87, 72, 175, 202, 44, 0, 0, 40, 58, 113, 128, 17, 0, 0, 0}; |
| 910 | auto inner = newMemoryInputStream(data); |
| 911 | auto source = newEncodedReadableSource(rpc::StreamEncoding::GZIP, kj::mv(inner)); |
| 912 | |
| 913 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 914 | // Should decompress on read all... |
| 915 | auto allBytes = co_await source->readAllBytes(kj::maxValue); |
| 916 | KJ_ASSERT(allBytes == "some data to gzip"_kjb); |
| 917 | }); |
| 918 | } |
| 919 | |
| 920 | KJ_TEST("Gzip encoded stream (pumpTo)") { |
| 921 | TestFixture fixture; |
| 922 | static constexpr kj::byte data[] = {31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 43, 206, 207, 77, 85, 72, 73, |
| 923 | 44, 73, 84, 40, 201, 87, 72, 175, 202, 44, 0, 0, 40, 58, 113, 128, 17, 0, 0, 0}; |
| 924 | auto inner = newMemoryInputStream(data); |
| 925 | auto source = newEncodedReadableSource(rpc::StreamEncoding::GZIP, kj::mv(inner)); |
| 926 | |
| 927 | MockWritableSink sink; |
| 928 | |
| 929 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 930 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 931 | }); |
| 932 | |
| 933 | KJ_ASSERT(sink.writtenData == "some data to gzip"_kjb); |
| 934 | } |
| 935 | |
| 936 | KJ_TEST("Gzip encoded stream (pumpTo same encoding)") { |
| 937 | TestFixture fixture; |
| 938 | static const kj::byte data[] = {31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 43, 206, 207, 77, 85, 72, 73, 44, |
| 939 | 73, 84, 40, 201, 87, 72, 175, 202, 44, 0, 0, 40, 58, 113, 128, 17, 0, 0, 0}; |
| 940 | auto in = newMemoryInputStream(data); |
| 941 | auto source = newEncodedReadableSource(rpc::StreamEncoding::GZIP, kj::mv(in)); |
| 942 | |
| 943 | MemoryAsyncOutputStream inner; |
| 944 | auto fakeOwn = kj::Own<MemoryAsyncOutputStream>(&inner, kj::NullDisposer::instance); |
| 945 | auto sink = newEncodedWritableSink(rpc::StreamEncoding::GZIP, kj::mv(fakeOwn)); |
| 946 | |
| 947 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 948 | co_await environment.context.waitForDeferredProxy(source->pumpTo(*sink, EndAfterPump::YES)); |
| 949 | }); |
| 950 | |
| 951 | // The data should pass through unchanged. |
| 952 | KJ_ASSERT(inner.data == data); |
| 953 | } |
| 954 | |
| 955 | KJ_TEST("Gzip encoded stream (pumpTo different encoding)") { |
| 956 | TestFixture fixture; |
| 957 | static const kj::byte data[] = {31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 43, 206, 207, 77, 85, 72, 73, 44, |
| 958 | 73, 84, 40, 201, 87, 72, 175, 202, 44, 0, 0, 40, 58, 113, 128, 17, 0, 0, 0}; |
| 959 | auto in = newMemoryInputStream(data); |
| 960 | auto source = newEncodedReadableSource(rpc::StreamEncoding::GZIP, kj::mv(in)); |
| 961 | |
| 962 | MemoryAsyncOutputStream inner; |
| 963 | auto fakeOwn = kj::Own<MemoryAsyncOutputStream>(&inner, kj::NullDisposer::instance); |
| 964 | auto sink = newEncodedWritableSink(rpc::StreamEncoding::BROTLI, kj::mv(fakeOwn)); |
| 965 | |
| 966 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 967 | co_await environment.context.waitForDeferredProxy(source->pumpTo(*sink, EndAfterPump::YES)); |
| 968 | }); |
| 969 | |
| 970 | // The data shuld be brotli compressed. |
| 971 | static const kj::byte expected[] = { |
| 972 | 5, 8, 128, 115, 111, 109, 101, 32, 100, 97, 116, 97, 32, 116, 111, 32, 103, 122, 105, 112, 3}; |
| 973 | KJ_ASSERT(inner.data == expected); |
| 974 | } |
| 975 | |
| 976 | // ====================================================================================== |
| 977 | // Adaptive Pump Behavior Tests |
| 978 | // These tests verify the adaptive pump heuristics without relying on timing. |
| 979 | |
| 980 | // Mock AsyncInputStream that tracks tryRead() calls and their parameters |
| 981 | class AdaptiveTestInputStream final: public kj::AsyncInputStream { |
| 982 | public: |
| 983 | enum class FillBehavior { |
| 984 | ALWAYS_FILL_COMPLETELY, // Always fill the buffer completely |
| 985 | PARTIAL_FILLS, // Always return partial fills |
| 986 | MIXED, // Alternate between full and partial |
| 987 | }; |
| 988 | |
| 989 | AdaptiveTestInputStream(size_t totalSize, FillBehavior behavior, size_t chunkSize = 0) |
| 990 | : totalSize_(totalSize), |
| 991 | position_(0), |
| 992 | behavior_(behavior), |
| 993 | chunkSize_(chunkSize), |
| 994 | readCount_(0) {} |
| 995 | |
| 996 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 997 | readCount_++; |
| 998 | |
| 999 | // Track the minBytes parameter on each call |
| 1000 | minBytesHistory_.add(minBytes); |
| 1001 | maxBytesHistory_.add(maxBytes); |
| 1002 | |
| 1003 | if (position_ >= totalSize_) { |
| 1004 | co_return 0; // EOF |
| 1005 | } |
| 1006 | |
| 1007 | size_t remaining = totalSize_ - position_; |
| 1008 | size_t toRead = 0; |
| 1009 | |
| 1010 | switch (behavior_) { |
| 1011 | case FillBehavior::ALWAYS_FILL_COMPLETELY: |
| 1012 | // Fill the buffer completely up to maxBytes |
| 1013 | toRead = kj::min(remaining, maxBytes); |
| 1014 | break; |
| 1015 | |
| 1016 | case FillBehavior::PARTIAL_FILLS: |
| 1017 | // Return partial fills - less than maxBytes but at least minBytes |
| 1018 | // (unless at EOF). This simulates a stream with natural boundaries. |
| 1019 | if (remaining >= minBytes) { |
| 1020 | // We have enough data to satisfy minBytes |
| 1021 | if (chunkSize_ > 0 && chunkSize_ >= minBytes) { |
| 1022 | // Use chunkSize if it's large enough |
| 1023 | toRead = kj::min(remaining, chunkSize_); |
| 1024 | } else { |
| 1025 | // Otherwise, use minBytes to avoid triggering EOF |
| 1026 | toRead = minBytes; |
| 1027 | } |
| 1028 | } else { |
| 1029 | // At the end, return what's left (even if less than minBytes) |
| 1030 | toRead = remaining; |
| 1031 | } |
| 1032 | break; |
| 1033 | |
| 1034 | case FillBehavior::MIXED: |
| 1035 | // Alternate between full and partial fills |
| 1036 | if (readCount_ % 2 == 1) { |
| 1037 | toRead = kj::min(remaining, maxBytes); |
| 1038 | } else { |
| 1039 | toRead = kj::min(remaining, minBytes); |
| 1040 | } |
| 1041 | break; |
| 1042 | } |
| 1043 | |
| 1044 | // Fill buffer with predictable data |
| 1045 | auto dest = static_cast<kj::byte*>(buffer); |
| 1046 | for (size_t i = 0; i < toRead; i++) { |
| 1047 | dest[i] = static_cast<kj::byte>((position_ + i) & 0xFF); |
| 1048 | } |
| 1049 | |
| 1050 | position_ += toRead; |
| 1051 | co_return toRead; |
| 1052 | } |
| 1053 | |
| 1054 | kj::Maybe<uint64_t> tryGetLength() override { |
| 1055 | return totalSize_; |
| 1056 | } |
| 1057 | |
| 1058 | // Test accessors |
| 1059 | const kj::ArrayPtr<const size_t> getMinBytesHistory() const { |
| 1060 | return minBytesHistory_.asPtr(); |
| 1061 | } |
| 1062 | const kj::ArrayPtr<const size_t> getMaxBytesHistory() const { |
| 1063 | return maxBytesHistory_.asPtr(); |
| 1064 | } |
| 1065 | size_t getReadCount() const { |
| 1066 | return readCount_; |
| 1067 | } |
| 1068 | |
| 1069 | private: |
| 1070 | size_t totalSize_; |
| 1071 | size_t position_; |
| 1072 | FillBehavior behavior_; |
| 1073 | size_t chunkSize_; |
| 1074 | size_t readCount_; |
| 1075 | kj::Vector<size_t> minBytesHistory_; |
| 1076 | kj::Vector<size_t> maxBytesHistory_; |
| 1077 | }; |
| 1078 | |
| 1079 | // Mock WritableSink that tracks write patterns |
| 1080 | class AdaptiveTestSink final: public WritableSink { |
| 1081 | public: |
| 1082 | AdaptiveTestSink() = default; |
| 1083 | ~AdaptiveTestSink() = default; |
| 1084 | |
| 1085 | kj::Promise<void> write(kj::ArrayPtr<const kj::byte> buffer) override { |
| 1086 | writeCallCount_++; |
| 1087 | writeSizes_.add(buffer.size()); |
| 1088 | totalBytesWritten_ += buffer.size(); |
| 1089 | co_return; |
| 1090 | } |
| 1091 | |
| 1092 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) override { |
| 1093 | KJ_FAIL_ASSERT("Should not be called in these tests"); |
| 1094 | } |
| 1095 | |
| 1096 | kj::Promise<void> end() override { |
| 1097 | endCallCount_++; |
| 1098 | co_return; |
| 1099 | } |
| 1100 | |
| 1101 | void abort(kj::Exception reason) override { |
| 1102 | abortCallCount_++; |
| 1103 | } |
| 1104 | |
| 1105 | rpc::StreamEncoding disownEncodingResponsibility() override { |
| 1106 | return rpc::StreamEncoding::IDENTITY; |
| 1107 | } |
| 1108 | |
| 1109 | rpc::StreamEncoding getEncoding() override { |
| 1110 | return rpc::StreamEncoding::IDENTITY; |
| 1111 | } |
| 1112 | |
| 1113 | // Test accessors |
| 1114 | uint32_t getWriteCallCount() const { |
| 1115 | return writeCallCount_; |
| 1116 | } |
| 1117 | const kj::ArrayPtr<const size_t> getWriteSizes() const { |
| 1118 | return writeSizes_.asPtr(); |
| 1119 | } |
| 1120 | size_t getTotalBytesWritten() const { |
| 1121 | return totalBytesWritten_; |
| 1122 | } |
| 1123 | |
| 1124 | private: |
| 1125 | uint32_t writeCallCount_ = 0; |
| 1126 | uint32_t endCallCount_ = 0; |
| 1127 | uint32_t abortCallCount_ = 0; |
| 1128 | size_t totalBytesWritten_ = 0; |
| 1129 | kj::Vector<size_t> writeSizes_; |
| 1130 | }; |
| 1131 | |
| 1132 | KJ_TEST("Adaptive pump: verify mock stream is called") { |
| 1133 | TestFixture fixture; |
| 1134 | |
| 1135 | // Simple test to verify the mock tracking actually works |
| 1136 | AdaptiveTestInputStream input( |
| 1137 | 100 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1138 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1139 | |
| 1140 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1141 | AdaptiveTestSink sink; |
| 1142 | |
| 1143 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1144 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1145 | KJ_ASSERT(sink.getTotalBytesWritten() == 100 * 1024); |
| 1146 | |
| 1147 | // Verify that our mock was actually called |
| 1148 | // Note: Vector tracking doesn't work reliably in coroutine context, but readCount does |
| 1149 | KJ_ASSERT(input.getReadCount() > 0, "Mock stream was never read from"); |
| 1150 | }); |
| 1151 | } |
| 1152 | |
| 1153 | KJ_TEST("Adaptive pump: buffer sizing for small stream (2KB)") { |
| 1154 | TestFixture fixture; |
| 1155 | |
| 1156 | // Small stream should use a small buffer (power of 2, clamped to range) |
| 1157 | // For 2KB, expect buffer size of 2KB (next power of 2), clamped to MED_BUFFER_SIZE (64KB) |
| 1158 | AdaptiveTestInputStream input( |
| 1159 | 2 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1160 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1161 | |
| 1162 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1163 | AdaptiveTestSink sink; |
| 1164 | |
| 1165 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1166 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1167 | |
| 1168 | // Verify the stream was read efficiently |
| 1169 | KJ_ASSERT(sink.getTotalBytesWritten() == 2 * 1024); |
| 1170 | |
| 1171 | // For small streams that fill completely, we should see efficient buffer usage |
| 1172 | // The buffer should be sized appropriately (power of 2, at least MIN_BUFFER_SIZE) |
| 1173 | auto maxBytesHistory = input.getMaxBytesHistory(); |
| 1174 | if (maxBytesHistory.size() > 0) { |
| 1175 | // First read should use a buffer size that's a power of 2 and >= 2KB |
| 1176 | size_t firstBufferSize = maxBytesHistory[0]; |
| 1177 | KJ_ASSERT(firstBufferSize >= 2 * 1024, firstBufferSize); |
| 1178 | KJ_ASSERT( |
| 1179 | (firstBufferSize & (firstBufferSize - 1)) == 0, "Should be power of 2", firstBufferSize); |
| 1180 | } |
| 1181 | // For very small streams, there might be optimizations that bypass our tracking |
| 1182 | }); |
| 1183 | } |
| 1184 | |
| 1185 | KJ_TEST("Adaptive pump: buffer sizing for medium stream (500KB)") { |
| 1186 | TestFixture fixture; |
| 1187 | |
| 1188 | // Medium stream should be read efficiently in a reasonable number of chunks |
| 1189 | AdaptiveTestInputStream input( |
| 1190 | 500 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1191 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1192 | |
| 1193 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1194 | AdaptiveTestSink sink; |
| 1195 | |
| 1196 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1197 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1198 | |
| 1199 | KJ_ASSERT(sink.getTotalBytesWritten() == 500 * 1024); |
| 1200 | |
| 1201 | // For a 500KB stream, with reasonable buffer sizing (likely 64KB), |
| 1202 | // we should see around 8-10 reads |
| 1203 | auto readCount = input.getReadCount(); |
| 1204 | KJ_ASSERT(readCount >= 4 && readCount <= 20, "Expected 4-20 reads for 500KB stream", readCount); |
| 1205 | }); |
| 1206 | } |
| 1207 | |
| 1208 | KJ_TEST("Adaptive pump: buffer sizing for large stream (2MB)") { |
| 1209 | TestFixture fixture; |
| 1210 | |
| 1211 | // Large stream (>1MB) should use MAX_BUFFER_SIZE and complete efficiently |
| 1212 | AdaptiveTestInputStream input( |
| 1213 | 2 * 1024 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1214 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1215 | |
| 1216 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1217 | AdaptiveTestSink sink; |
| 1218 | |
| 1219 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1220 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1221 | |
| 1222 | KJ_ASSERT(sink.getTotalBytesWritten() == 2 * 1024 * 1024); |
| 1223 | |
| 1224 | // For a 2MB stream with MAX_BUFFER_SIZE (128KB), we should see around 16-18 reads |
| 1225 | auto readCount = input.getReadCount(); |
| 1226 | KJ_ASSERT(readCount >= 10 && readCount <= 30, "Expected 10-30 reads for 2MB stream", readCount); |
| 1227 | }); |
| 1228 | } |
| 1229 | |
| 1230 | KJ_TEST("Adaptive pump: fast-filling stream efficiency") { |
| 1231 | TestFixture fixture; |
| 1232 | |
| 1233 | // Stream that always fills the buffer completely should be read efficiently |
| 1234 | AdaptiveTestInputStream input( |
| 1235 | 200 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1236 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1237 | |
| 1238 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1239 | AdaptiveTestSink sink; |
| 1240 | |
| 1241 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1242 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1243 | |
| 1244 | KJ_ASSERT(sink.getTotalBytesWritten() == 200 * 1024); |
| 1245 | |
| 1246 | // Fast-filling streams should complete in relatively few iterations |
| 1247 | auto readCount = input.getReadCount(); |
| 1248 | KJ_ASSERT(readCount >= 2 && readCount <= 10, |
| 1249 | "Expected 2-10 reads for 200KB fast-filling stream", readCount); |
| 1250 | |
| 1251 | // Write count should be close to read count (double buffering) |
| 1252 | auto writeCount = sink.getWriteCallCount(); |
| 1253 | KJ_ASSERT(writeCount >= readCount - 2 && writeCount <= readCount + 2, |
| 1254 | "Write count should be close to read count", writeCount, readCount); |
| 1255 | }); |
| 1256 | } |
| 1257 | |
| 1258 | KJ_TEST("Adaptive pump: partial-filling stream behavior") { |
| 1259 | TestFixture fixture; |
| 1260 | |
| 1261 | // Stream that returns partial fills (32KB chunks) |
| 1262 | // Should require more iterations than a fast-filling stream |
| 1263 | AdaptiveTestInputStream input( |
| 1264 | 200 * 1024, AdaptiveTestInputStream::FillBehavior::PARTIAL_FILLS, 32 * 1024); |
| 1265 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1266 | |
| 1267 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1268 | AdaptiveTestSink sink; |
| 1269 | |
| 1270 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1271 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1272 | |
| 1273 | KJ_ASSERT(sink.getTotalBytesWritten() == 200 * 1024); |
| 1274 | |
| 1275 | // Partial-filling streams require more reads than fast-filling streams |
| 1276 | // 200KB / 32KB chunks = ~7 reads minimum |
| 1277 | auto readCount = input.getReadCount(); |
| 1278 | KJ_ASSERT(readCount >= 5, "Expected at least 5 reads for partial-fill stream", readCount); |
| 1279 | }); |
| 1280 | } |
| 1281 | |
| 1282 | KJ_TEST("Adaptive pump: large stream efficiency") { |
| 1283 | TestFixture fixture; |
| 1284 | |
| 1285 | // Large streams should complete efficiently with appropriate buffer sizing |
| 1286 | AdaptiveTestInputStream input( |
| 1287 | 2 * 1024 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1288 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1289 | |
| 1290 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1291 | AdaptiveTestSink sink; |
| 1292 | |
| 1293 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1294 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1295 | |
| 1296 | KJ_ASSERT(sink.getTotalBytesWritten() == 2 * 1024 * 1024); |
| 1297 | |
| 1298 | // Should complete in a reasonable number of reads |
| 1299 | // With 128KB buffers: 2MB / 128KB = ~16 reads |
| 1300 | auto readCount = input.getReadCount(); |
| 1301 | KJ_ASSERT(readCount >= 10 && readCount <= 30, "Expected 10-30 reads for 2MB stream", readCount); |
| 1302 | }); |
| 1303 | } |
| 1304 | |
| 1305 | KJ_TEST("Adaptive pump: mixed behavior stream") { |
| 1306 | TestFixture fixture; |
| 1307 | |
| 1308 | // Stream that alternates between full and partial fills |
| 1309 | // Should still complete reasonably efficiently |
| 1310 | AdaptiveTestInputStream input(1024 * 1024, AdaptiveTestInputStream::FillBehavior::MIXED); |
| 1311 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1312 | |
| 1313 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1314 | AdaptiveTestSink sink; |
| 1315 | |
| 1316 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1317 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1318 | |
| 1319 | KJ_ASSERT(sink.getTotalBytesWritten() == 1024 * 1024); |
| 1320 | |
| 1321 | // Mixed behavior should still complete in a reasonable number of reads |
| 1322 | auto readCount = input.getReadCount(); |
| 1323 | KJ_ASSERT( |
| 1324 | readCount >= 5 && readCount <= 40, "Expected 5-40 reads for 1MB mixed stream", readCount); |
| 1325 | }); |
| 1326 | } |
| 1327 | |
| 1328 | KJ_TEST("Adaptive pump: double buffering behavior") { |
| 1329 | TestFixture fixture; |
| 1330 | |
| 1331 | // Verify that the pump uses double buffering effectively |
| 1332 | // We can observe this by checking write patterns match read patterns |
| 1333 | AdaptiveTestInputStream input( |
| 1334 | 100 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1335 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1336 | |
| 1337 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1338 | AdaptiveTestSink sink; |
| 1339 | |
| 1340 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1341 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1342 | |
| 1343 | KJ_ASSERT(sink.getTotalBytesWritten() == 100 * 1024); |
| 1344 | |
| 1345 | // With double buffering, the number of writes should be close to the number of reads |
| 1346 | // (minus one since the last read returns EOF) |
| 1347 | auto readCount = input.getReadCount(); |
| 1348 | auto writeCount = sink.getWriteCallCount(); |
| 1349 | |
| 1350 | // Reads should be at least as many as writes (or equal for small streams) |
| 1351 | KJ_ASSERT(readCount >= writeCount, readCount, writeCount); |
| 1352 | |
| 1353 | // For properly pipelined operation, reads and writes should be close |
| 1354 | // The difference should be small (typically 0-1 for good pipelining) |
| 1355 | KJ_ASSERT(readCount - writeCount <= 2, "Pipelining gap too large", readCount, writeCount); |
| 1356 | }); |
| 1357 | } |
| 1358 | |
| 1359 | KJ_TEST("Adaptive pump: verify heuristics optimize for throughput") { |
| 1360 | TestFixture fixture; |
| 1361 | |
| 1362 | // Large stream with consistent full fills should optimize for throughput |
| 1363 | // by using large buffers and appropriate minBytes |
| 1364 | AdaptiveTestInputStream input( |
| 1365 | 1024 * 1024, AdaptiveTestInputStream::FillBehavior::ALWAYS_FILL_COMPLETELY); |
| 1366 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1367 | |
| 1368 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1369 | AdaptiveTestSink sink; |
| 1370 | |
| 1371 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1372 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1373 | |
| 1374 | KJ_ASSERT(sink.getTotalBytesWritten() == 1024 * 1024); |
| 1375 | |
| 1376 | auto writeSizes = sink.getWriteSizes(); |
| 1377 | auto readCount = input.getReadCount(); |
| 1378 | |
| 1379 | // For a 1MB stream with fast fills, we should see efficient large writes |
| 1380 | // The number of iterations should be relatively small |
| 1381 | KJ_ASSERT(readCount <= 20, "Too many iterations for 1MB stream", readCount); |
| 1382 | |
| 1383 | // Most writes should be using the full buffer |
| 1384 | size_t largeWrites = 0; |
| 1385 | for (size_t size: writeSizes) { |
| 1386 | if (size >= 32 * 1024) { // Reasonably large writes |
| 1387 | largeWrites++; |
| 1388 | } |
| 1389 | } |
| 1390 | |
| 1391 | // Most writes should be large for throughput optimization |
| 1392 | KJ_ASSERT(largeWrites >= writeSizes.size() / 2, "Expected mostly large writes for throughput", |
| 1393 | largeWrites, writeSizes.size()); |
| 1394 | }); |
| 1395 | } |
| 1396 | |
| 1397 | KJ_TEST("Adaptive pump: verify heuristics optimize for responsiveness") { |
| 1398 | TestFixture fixture; |
| 1399 | |
| 1400 | // Stream with medium chunks should optimize for responsiveness |
| 1401 | // Using 16KB chunks which will not fill larger buffers |
| 1402 | AdaptiveTestInputStream input( |
| 1403 | 256 * 1024, AdaptiveTestInputStream::FillBehavior::PARTIAL_FILLS, 16 * 1024); |
| 1404 | auto fakeOwn = kj::Own<AdaptiveTestInputStream>(&input, kj::NullDisposer::instance); |
| 1405 | |
| 1406 | auto source = newReadableSource(kj::mv(fakeOwn)); |
| 1407 | AdaptiveTestSink sink; |
| 1408 | |
| 1409 | fixture.runInIoContext([&](const auto& environment) -> kj::Promise<void> { |
| 1410 | co_await environment.context.waitForDeferredProxy(source->pumpTo(sink, EndAfterPump::YES)); |
| 1411 | |
| 1412 | KJ_ASSERT(sink.getTotalBytesWritten() == 256 * 1024); |
| 1413 | |
| 1414 | auto writeSizes = sink.getWriteSizes(); |
| 1415 | |
| 1416 | // For partial-fill streams, writes should match the stream's natural chunk size |
| 1417 | // We should see multiple writes rather than trying to accumulate into large ones |
| 1418 | KJ_ASSERT(writeSizes.size() >= 4, "Expected multiple writes for partial-fill stream", |
| 1419 | writeSizes.size()); |
| 1420 | |
| 1421 | // The write pattern should reflect the stream's behavior |
| 1422 | // Most writes should be around the chunk size (16KB) or minBytes |
| 1423 | size_t mediumWrites = 0; |
| 1424 | for (size_t size: writeSizes) { |
| 1425 | if (size >= 8 * 1024 && size <= 32 * 1024) { // Medium chunks |
| 1426 | mediumWrites++; |
| 1427 | } |
| 1428 | } |
| 1429 | |
| 1430 | // Should have multiple medium-sized writes reflecting the partial-fill pattern |
| 1431 | KJ_ASSERT(mediumWrites >= 2, "Expected some medium writes for responsive stream", mediumWrites); |
| 1432 | }); |
| 1433 | } |
| 1434 | |
| 1435 | } // namespace |
| 1436 | } // namespace workerd::api::streams |