File
Blob: src/workerd/api/streams/internal-test.c++
| 1 | // Copyright (c) 2017-2022 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 | #include "internal.h" |
| 6 | #include "readable.h" |
| 7 | #include "standard.h" |
| 8 | #include "writable.h" |
| 9 | |
| 10 | #include <workerd/jsg/jsg-test.h> |
| 11 | #include <workerd/jsg/jsg.h> |
| 12 | #include <workerd/tests/test-fixture.h> |
| 13 | |
| 14 | #include <openssl/rand.h> |
| 15 | |
| 16 | namespace workerd::api { |
| 17 | namespace { |
| 18 | |
| 19 | // ====================================================================================== |
| 20 | // Shared test helpers |
| 21 | |
| 22 | // Simple source that returns EOF immediately |
| 23 | class EofSource final: public ReadableStreamSource { |
| 24 | public: |
| 25 | kj::Promise<size_t> tryRead(void*, size_t, size_t) override { |
| 26 | return static_cast<size_t>(0); |
| 27 | } |
| 28 | }; |
| 29 | |
| 30 | // Simple sink that accepts all writes |
| 31 | class NoopSink final: public WritableStreamSink { |
| 32 | public: |
| 33 | kj::Promise<void> write(kj::ArrayPtr<const byte>) override { |
| 34 | return kj::READY_NOW; |
| 35 | } |
| 36 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>>) override { |
| 37 | return kj::READY_NOW; |
| 38 | } |
| 39 | kj::Promise<void> end() override { |
| 40 | return kj::READY_NOW; |
| 41 | } |
| 42 | void abort(kj::Exception) override {} |
| 43 | }; |
| 44 | |
| 45 | // Creates a TestFixture with common flags for stream tests |
| 46 | TestFixture makeStreamTestFixture() { |
| 47 | capnp::MallocMessageBuilder message; |
| 48 | auto flags = message.initRoot<CompatibilityFlags>(); |
| 49 | flags.setStreamsJavaScriptControllers(true); |
| 50 | return TestFixture({.featureFlags = flags.asReader()}); |
| 51 | } |
| 52 | |
| 53 | // Creates a TestFixture with the abortClearsQueue flag for testing abort behavior |
| 54 | TestFixture makeAbortClearsQueueTestFixture() { |
| 55 | capnp::MallocMessageBuilder message; |
| 56 | auto flags = message.initRoot<CompatibilityFlags>(); |
| 57 | flags.setStreamsJavaScriptControllers(true); |
| 58 | flags.setInternalWritableStreamAbortClearsQueue(true); |
| 59 | return TestFixture({.featureFlags = flags.asReader()}); |
| 60 | } |
| 61 | |
| 62 | // Creates a BYOB-capable ReadableStream |
| 63 | jsg::Ref<ReadableStream> makeByteStream(jsg::Lock& js) { |
| 64 | auto rs = js.alloc<ReadableStream>(newReadableStreamJsController()); |
| 65 | rs->getController().setup( |
| 66 | js, UnderlyingSource{.type = kj::str("bytes")}, StreamQueuingStrategy{}); |
| 67 | return rs; |
| 68 | } |
| 69 | |
| 70 | // ====================================================================================== |
| 71 | // ReadableStreamSource test implementations |
| 72 | |
| 73 | template <int size> |
| 74 | class FooStream: public ReadableStreamSource { |
| 75 | public: |
| 76 | FooStream(): ptr(&data[0]), remaining_(size) { |
| 77 | KJ_ASSERT(RAND_bytes(data, size) == 1); |
| 78 | } |
| 79 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 80 | maxMaxBytesSeen_ = kj::max(maxMaxBytesSeen_, maxBytes); |
| 81 | numreads_++; |
| 82 | if (remaining_ == 0) return static_cast<size_t>(0); |
| 83 | KJ_ASSERT(minBytes == maxBytes); |
| 84 | auto amount = kj::min(remaining_, maxBytes); |
| 85 | memcpy(buffer, ptr, amount); |
| 86 | ptr += amount; |
| 87 | remaining_ -= amount; |
| 88 | return amount; |
| 89 | } |
| 90 | |
| 91 | kj::ArrayPtr<uint8_t> buf() { |
| 92 | return data; |
| 93 | } |
| 94 | |
| 95 | size_t remaining() { |
| 96 | return remaining_; |
| 97 | } |
| 98 | |
| 99 | size_t numreads() { |
| 100 | return numreads_; |
| 101 | } |
| 102 | |
| 103 | size_t maxMaxBytesSeen() { |
| 104 | return maxMaxBytesSeen_; |
| 105 | } |
| 106 | |
| 107 | private: |
| 108 | uint8_t data[size]; |
| 109 | uint8_t* ptr; |
| 110 | size_t remaining_; |
| 111 | size_t numreads_ = 0; |
| 112 | size_t maxMaxBytesSeen_ = 0; |
| 113 | }; |
| 114 | |
| 115 | template <int size> |
| 116 | class BarStream: public FooStream<size> { |
| 117 | public: |
| 118 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 119 | return size; |
| 120 | } |
| 121 | }; |
| 122 | |
| 123 | KJ_TEST("test") { |
| 124 | kj::EventLoop loop; |
| 125 | kj::WaitScope waitScope(loop); |
| 126 | |
| 127 | // In this first case, the stream does not report a length. The read size |
| 128 | // is min(limit, DEFAULT_BUFFER_CHUNK) = min(10001, 131072) = 10001, so the |
| 129 | // entire stream is consumed in a single read that returns a short read (10000 < 10001). |
| 130 | FooStream<10000> stream; |
| 131 | |
| 132 | stream.readAllBytes(10001) |
| 133 | .then([&](auto bytes) { |
| 134 | KJ_ASSERT(bytes.size() == 10000); |
| 135 | KJ_ASSERT(bytes == stream.buf().first(10000)); |
| 136 | }).wait(waitScope); |
| 137 | |
| 138 | KJ_ASSERT(stream.numreads() == 1); |
| 139 | KJ_ASSERT(stream.maxMaxBytesSeen() == 10001); |
| 140 | } |
| 141 | |
| 142 | KJ_TEST("test (text)") { |
| 143 | kj::EventLoop loop; |
| 144 | kj::WaitScope waitScope(loop); |
| 145 | |
| 146 | // In this first case, the stream does not report a length. The read size |
| 147 | // is min(limit, DEFAULT_BUFFER_CHUNK) = min(10001, 131072) = 10001, so the |
| 148 | // entire stream is consumed in a single read that returns a short read (10000 < 10001). |
| 149 | FooStream<10000> stream; |
| 150 | |
| 151 | stream.readAllText(10001) |
| 152 | .then([&](auto bytes) { |
| 153 | KJ_ASSERT(bytes.size() == 10000); |
| 154 | KJ_ASSERT(bytes.asBytes() == stream.buf().first(10000)); |
| 155 | }).wait(waitScope); |
| 156 | |
| 157 | KJ_ASSERT(stream.numreads() == 1); |
| 158 | KJ_ASSERT(stream.maxMaxBytesSeen() == 10001); |
| 159 | } |
| 160 | |
| 161 | KJ_TEST("test2") { |
| 162 | kj::EventLoop loop; |
| 163 | kj::WaitScope waitScope(loop); |
| 164 | |
| 165 | // In this second case, the stream does report a size, so we should see |
| 166 | // only one read. |
| 167 | BarStream<10000> stream; |
| 168 | |
| 169 | stream.readAllBytes(10001) |
| 170 | .then([&](auto bytes) { |
| 171 | KJ_ASSERT(bytes.size() == 10000); |
| 172 | KJ_ASSERT(bytes == stream.buf().first(10000)); |
| 173 | }).wait(waitScope); |
| 174 | |
| 175 | KJ_ASSERT(stream.numreads() == 2); |
| 176 | KJ_ASSERT(stream.maxMaxBytesSeen() == 10000); |
| 177 | } |
| 178 | |
| 179 | KJ_TEST("test2 (text)") { |
| 180 | kj::EventLoop loop; |
| 181 | kj::WaitScope waitScope(loop); |
| 182 | |
| 183 | // In this second case, the stream does report a size, so we should see |
| 184 | // only one read. |
| 185 | BarStream<10000> stream; |
| 186 | |
| 187 | stream.readAllText(10001) |
| 188 | .then([&](auto bytes) { |
| 189 | KJ_ASSERT(bytes.size() == 10000); |
| 190 | KJ_ASSERT(bytes.asBytes() == stream.buf().first(10000)); |
| 191 | }).wait(waitScope); |
| 192 | |
| 193 | KJ_ASSERT(stream.numreads() == 2); |
| 194 | KJ_ASSERT(stream.maxMaxBytesSeen() == 10000); |
| 195 | } |
| 196 | |
| 197 | KJ_TEST("zero-length stream") { |
| 198 | kj::EventLoop loop; |
| 199 | kj::WaitScope waitScope(loop); |
| 200 | |
| 201 | class Zero: public ReadableStreamSource { |
| 202 | public: |
| 203 | kj::Promise<size_t> tryRead(void*, size_t, size_t) override { |
| 204 | return static_cast<size_t>(0); |
| 205 | } |
| 206 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 207 | return static_cast<size_t>(0); |
| 208 | } |
| 209 | }; |
| 210 | |
| 211 | Zero zero; |
| 212 | zero.readAllBytes(10).then([&](kj::Array<kj::byte> bytes) { |
| 213 | KJ_ASSERT(bytes.size() == 0); |
| 214 | }).wait(waitScope); |
| 215 | } |
| 216 | |
| 217 | KJ_TEST("lying stream") { |
| 218 | kj::EventLoop loop; |
| 219 | kj::WaitScope waitScope(loop); |
| 220 | |
| 221 | class Dishonest: public FooStream<10000> { |
| 222 | public: |
| 223 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 224 | return static_cast<size_t>(10); |
| 225 | } |
| 226 | }; |
| 227 | |
| 228 | Dishonest stream; |
| 229 | stream.readAllBytes(10001) |
| 230 | .then([&](kj::Array<kj::byte> bytes) { |
| 231 | // The stream lies! it says there are only 10 bytes but there are more. |
| 232 | // oh well, we at least make sure we get the right result. |
| 233 | KJ_ASSERT(bytes.size() == 10000); |
| 234 | }).wait(waitScope); |
| 235 | |
| 236 | KJ_ASSERT(stream.numreads() == 1001); |
| 237 | KJ_ASSERT(stream.maxMaxBytesSeen() == 10); |
| 238 | } |
| 239 | |
| 240 | KJ_TEST("honest small stream") { |
| 241 | kj::EventLoop loop; |
| 242 | kj::WaitScope waitScope(loop); |
| 243 | |
| 244 | class HonestSmall: public FooStream<100> { |
| 245 | public: |
| 246 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 247 | return static_cast<size_t>(100); |
| 248 | } |
| 249 | }; |
| 250 | |
| 251 | HonestSmall stream; |
| 252 | stream.readAllBytes(1001).then([&](kj::Array<kj::byte> bytes) { |
| 253 | KJ_ASSERT(bytes.size() == 100); |
| 254 | }).wait(waitScope); |
| 255 | |
| 256 | KJ_ASSERT(stream.numreads() == 2); |
| 257 | KJ_ASSERT(stream.maxMaxBytesSeen(), 100); |
| 258 | } |
| 259 | |
| 260 | KJ_TEST("WritableStreamInternalController queue size assertion") { |
| 261 | auto fixture = makeStreamTestFixture(); |
| 262 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 263 | // Make sure that while an internal sink is being piped into, no other writes are |
| 264 | // allowed to be queued. |
| 265 | |
| 266 | jsg::Ref<ReadableStream> source = ReadableStream::constructor(env.js, kj::none, kj::none); |
| 267 | jsg::Ref<WritableStream> sink = |
| 268 | env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none); |
| 269 | |
| 270 | auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{.preventClose = true}); |
| 271 | |
| 272 | KJ_ASSERT(sink->isLocked()); |
| 273 | try { |
| 274 | sink->getWriter(env.js); |
| 275 | KJ_FAIL_ASSERT("Expected getWriter to throw"); |
| 276 | } catch (...) { |
| 277 | auto ex = kj::getCaughtExceptionAsKj(); |
| 278 | KJ_ASSERT(ex.getDescription() == |
| 279 | "expected !stream->isLocked(); jsg.TypeError: This WritableStream " |
| 280 | "is currently locked to a writer."); |
| 281 | } |
| 282 | |
| 283 | auto buffersource = env.js.bytes(kj::heapArray<kj::byte>(10)); |
| 284 | |
| 285 | bool writeFailed = false; |
| 286 | |
| 287 | auto write = sink->getController() |
| 288 | .write(env.js, buffersource.getHandle(env.js)) |
| 289 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value value) { |
| 290 | writeFailed = true; |
| 291 | auto ex = js.exceptionToKj(kj::mv(value)); |
| 292 | KJ_ASSERT( |
| 293 | ex.getDescription() == "jsg.TypeError: This WritableStream is currently being piped to."); |
| 294 | }); |
| 295 | |
| 296 | source->getController().cancel(env.js, kj::none); |
| 297 | |
| 298 | env.js.runMicrotasks(); |
| 299 | |
| 300 | KJ_ASSERT(!sink->isLocked()); |
| 301 | KJ_ASSERT(!sink->getController().isClosedOrClosing()); |
| 302 | KJ_ASSERT(!sink->getController().isErrored()); |
| 303 | KJ_ASSERT(sink->getController().isErroring(env.js) == kj::none); |
| 304 | |
| 305 | // Getting a writer at this point does not throw... |
| 306 | sink->getWriter(env.js); |
| 307 | }); |
| 308 | } |
| 309 | |
| 310 | KJ_TEST("WritableStreamInternalController operations reject when piped to") { |
| 311 | // Tests that close/flush/tryPipeFrom reject with "currently being piped to" |
| 312 | // during an active pipe operation. |
| 313 | auto fixture = makeStreamTestFixture(); |
| 314 | |
| 315 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 316 | auto source = ReadableStream::constructor(env.js, kj::none, kj::none); |
| 317 | auto source2 = ReadableStream::constructor(env.js, kj::none, kj::none); |
| 318 | auto sink = env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none); |
| 319 | |
| 320 | auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{.preventClose = true}); |
| 321 | KJ_ASSERT(sink->isLocked()); |
| 322 | |
| 323 | constexpr auto expectedError = |
| 324 | "jsg.TypeError: This WritableStream is currently being piped to."_kj; |
| 325 | |
| 326 | auto expectReject = [&](jsg::Promise<void> promise, bool& flag) { |
| 327 | promise.catch_(env.js, [&](jsg::Lock& js, jsg::Value value) { |
| 328 | flag = true; |
| 329 | KJ_ASSERT(js.exceptionToKj(kj::mv(value)).getDescription() == expectedError); |
| 330 | }); |
| 331 | }; |
| 332 | |
| 333 | bool closeFailed = false, flushFailed = false, pipeFailed = false; |
| 334 | |
| 335 | expectReject(sink->getController().close(env.js), closeFailed); |
| 336 | expectReject(sink->getController().flush(env.js), flushFailed); |
| 337 | |
| 338 | KJ_IF_SOME(secondPipe, |
| 339 | sink->getController().tryPipeFrom( |
| 340 | env.js, source2.addRef(), PipeToOptions{.preventClose = true})) { |
| 341 | expectReject(kj::mv(secondPipe), pipeFailed); |
| 342 | } else { |
| 343 | KJ_FAIL_ASSERT("Expected tryPipeFrom to return a promise"); |
| 344 | } |
| 345 | |
| 346 | source->getController().cancel(env.js, kj::none); |
| 347 | env.js.runMicrotasks(); |
| 348 | |
| 349 | KJ_ASSERT(closeFailed); |
| 350 | KJ_ASSERT(flushFailed); |
| 351 | KJ_ASSERT(pipeFailed); |
| 352 | }); |
| 353 | } |
| 354 | |
| 355 | KJ_TEST("WritableStreamInternalController observability") { |
| 356 | auto fixture = makeStreamTestFixture(); |
| 357 | |
| 358 | class MyObserver final: public ByteStreamObserver { |
| 359 | public: |
| 360 | void onChunkEnqueued(size_t bytes) override { |
| 361 | ++queueSize; |
| 362 | queueSizeBytes += bytes; |
| 363 | }; |
| 364 | void onChunkDequeued(size_t bytes) override { |
| 365 | queueSizeBytes -= bytes; |
| 366 | --queueSize; |
| 367 | }; |
| 368 | uint64_t queueSize = 0; |
| 369 | uint64_t queueSizeBytes = 0; |
| 370 | }; |
| 371 | |
| 372 | auto myObserver = kj::heap<MyObserver>(); |
| 373 | auto& observer = *myObserver; |
| 374 | kj::Maybe<jsg::Ref<WritableStream>> stream; |
| 375 | fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 376 | stream = env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::mv(myObserver)); |
| 377 | |
| 378 | auto write = [&](size_t size) { |
| 379 | auto buffersource = env.js.bytes(kj::heapArray<kj::byte>(size)); |
| 380 | return env.context.awaitJs(env.js, |
| 381 | KJ_ASSERT_NONNULL(stream)->getController().write(env.js, buffersource.getHandle(env.js))); |
| 382 | }; |
| 383 | |
| 384 | KJ_ASSERT(observer.queueSize == 0); |
| 385 | KJ_ASSERT(observer.queueSizeBytes == 0); |
| 386 | |
| 387 | auto builder = kj::heapArrayBuilder<kj::Promise<void>>(2); |
| 388 | builder.add(write(1)); |
| 389 | |
| 390 | KJ_ASSERT(observer.queueSize == 1); |
| 391 | KJ_ASSERT(observer.queueSizeBytes == 1); |
| 392 | |
| 393 | builder.add(write(10)); |
| 394 | |
| 395 | KJ_ASSERT(observer.queueSize == 2); |
| 396 | KJ_ASSERT(observer.queueSizeBytes == 11); |
| 397 | |
| 398 | return kj::joinPromises(builder.finish()); |
| 399 | }); |
| 400 | |
| 401 | KJ_ASSERT(observer.queueSize == 0); |
| 402 | KJ_ASSERT(observer.queueSizeBytes == 0); |
| 403 | } |
| 404 | |
| 405 | // Test for use-after-free fix in pipeLoop when abort is called during pending read. |
| 406 | // The fix ensures the Pipe::State is ref-counted and survives until all callbacks complete. |
| 407 | KJ_TEST("WritableStreamInternalController pipeLoop abort during pending read") { |
| 408 | auto fixture = makeAbortClearsQueueTestFixture(); |
| 409 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 410 | // Create a JavaScript-backed ReadableStream. |
| 411 | // The pull function will be called when the pipe tries to read. |
| 412 | // We use a JS-backed stream so that pipeLoop is used (not the kj pipe path). |
| 413 | // |
| 414 | // We need to simulate: |
| 415 | // 1. First read succeeds with some data |
| 416 | // 2. Second read is pending (the promise from pull is not resolved) |
| 417 | // 3. While pending, we abort the writable stream |
| 418 | // |
| 419 | // Using an UnderlyingSource with a pull callback that enqueues data once, |
| 420 | // then on the second call returns without enqueuing (leaving the read pending). |
| 421 | |
| 422 | int pullCount = 0; |
| 423 | jsg::Ref<ReadableStream> source = ReadableStream::constructor(env.js, |
| 424 | UnderlyingSource{.pull = |
| 425 | [&pullCount](jsg::Lock& js, UnderlyingSource::Controller controller) { |
| 426 | pullCount++; |
| 427 | auto& c = KJ_ASSERT_NONNULL(controller.tryGet<jsg::Ref<ReadableStreamDefaultController>>()); |
| 428 | if (pullCount == 1) { |
| 429 | // First pull: enqueue some data so the pipe loop can make progress |
| 430 | auto data = js.bytes(kj::heapArray<kj::byte>({1, 2, 3, 4})); |
| 431 | c->enqueue(js, data.getHandle(js)); |
| 432 | } |
| 433 | // Second pull onwards: don't enqueue anything, leaving the read pending. |
| 434 | // This simulates an async data source that hasn't received data yet. |
| 435 | // The promise returned by read() will be pending. |
| 436 | return js.resolvedPromise(); |
| 437 | }}, |
| 438 | kj::none); |
| 439 | |
| 440 | jsg::Ref<WritableStream> sink = |
| 441 | env.js.alloc<WritableStream>(env.context, kj::heap<NoopSink>(), kj::none); |
| 442 | |
| 443 | auto pipeTo = source->pipeTo(env.js, sink.addRef(), PipeToOptions{}); |
| 444 | pipeTo.markAsHandled(env.js); |
| 445 | env.js.runMicrotasks(); |
| 446 | |
| 447 | // Abort while pipeLoop is waiting for a pending read |
| 448 | auto abortPromise = sink->getController().abort(env.js, env.js.v8TypeError("Test abort"_kj)); |
| 449 | abortPromise.markAsHandled(env.js); |
| 450 | env.js.runMicrotasks(); |
| 451 | |
| 452 | // If we get here without crashing, the test passes |
| 453 | KJ_ASSERT(pullCount >= 1); |
| 454 | }); |
| 455 | } |
| 456 | |
| 457 | // ====================================================================================== |
| 458 | // DrainingReader tests for internal streams |
| 459 | // |
| 460 | // The internal stream implementation's drainingRead() behaves like a normal read() - |
| 461 | // it returns at most one chunk at a time rather than draining all buffered data. |
| 462 | // This is because internal streams are backed by kj I/O which is inherently async |
| 463 | // and doesn't have internal JS-side buffering. |
| 464 | |
| 465 | KJ_TEST("DrainingReader basic creation and locking (internal stream)") { |
| 466 | auto fixture = makeStreamTestFixture(); |
| 467 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 468 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>()); |
| 469 | KJ_ASSERT(!rs->isLocked()); |
| 470 | |
| 471 | KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) { |
| 472 | KJ_ASSERT(rs->isLocked()); |
| 473 | KJ_ASSERT(reader->isAttached()); |
| 474 | |
| 475 | reader->releaseLock(env.js); |
| 476 | KJ_ASSERT(!rs->isLocked()); |
| 477 | KJ_ASSERT(!reader->isAttached()); |
| 478 | } else { |
| 479 | KJ_FAIL_ASSERT("Failed to create DrainingReader"); |
| 480 | } |
| 481 | }); |
| 482 | } |
| 483 | |
| 484 | KJ_TEST("DrainingReader cannot be created on locked internal stream") { |
| 485 | auto fixture = makeStreamTestFixture(); |
| 486 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 487 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>()); |
| 488 | |
| 489 | KJ_IF_SOME(reader1, DrainingReader::create(env.js, *rs)) { |
| 490 | KJ_ASSERT(rs->isLocked()); |
| 491 | KJ_ASSERT(DrainingReader::create(env.js, *rs) == kj::none); |
| 492 | reader1->releaseLock(env.js); |
| 493 | } else { |
| 494 | KJ_FAIL_ASSERT("Failed to create first DrainingReader"); |
| 495 | } |
| 496 | }); |
| 497 | } |
| 498 | |
| 499 | KJ_TEST("DrainingReader read after releaseLock rejects (internal stream)") { |
| 500 | auto fixture = makeStreamTestFixture(); |
| 501 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 502 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>()); |
| 503 | |
| 504 | KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) { |
| 505 | reader->releaseLock(env.js); |
| 506 | |
| 507 | bool rejected = false; |
| 508 | reader->read(env.js).catch_(env.js, [&](jsg::Lock&, jsg::Value) -> DrainingReadResult { |
| 509 | rejected = true; |
| 510 | return {.done = true}; |
| 511 | }); |
| 512 | env.js.runMicrotasks(); |
| 513 | KJ_ASSERT(rejected); |
| 514 | } else { |
| 515 | KJ_FAIL_ASSERT("Failed to create DrainingReader"); |
| 516 | } |
| 517 | }); |
| 518 | } |
| 519 | |
| 520 | KJ_TEST("DrainingReader with maxRead parameter (internal stream)") { |
| 521 | // Test that the maxRead parameter is respected for internal streams |
| 522 | capnp::MallocMessageBuilder message; |
| 523 | auto flags = message.initRoot<CompatibilityFlags>(); |
| 524 | flags.setStreamsJavaScriptControllers(true); |
| 525 | |
| 526 | TestFixture fixture({.featureFlags = flags.asReader()}); |
| 527 | |
| 528 | bool testCompleted = false; |
| 529 | size_t lastMaxBytes = 0; |
| 530 | |
| 531 | fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 532 | class TestSource final: public ReadableStreamSource { |
| 533 | public: |
| 534 | explicit TestSource(size_t& maxBytesOut): lastMaxBytesOut(maxBytesOut) {} |
| 535 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 536 | readCount++; |
| 537 | // Note: maxBytes should be limited by the maxRead parameter |
| 538 | lastMaxBytesOut = maxBytes; |
| 539 | if (readCount == 1) { |
| 540 | // Return less than maxBytes |
| 541 | auto toWrite = kj::min(maxBytes, static_cast<size_t>(100)); |
| 542 | memset(buffer, 'x', toWrite); |
| 543 | return toWrite; |
| 544 | } |
| 545 | return static_cast<size_t>(0); // EOF |
| 546 | } |
| 547 | uint readCount = 0; |
| 548 | size_t& lastMaxBytesOut; |
| 549 | }; |
| 550 | |
| 551 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>(lastMaxBytes)); |
| 552 | |
| 553 | auto maybeReader = DrainingReader::create(env.js, *rs); |
| 554 | KJ_ASSERT(maybeReader != kj::none, "Failed to create DrainingReader"); |
| 555 | auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader)); |
| 556 | |
| 557 | // Pass a small maxRead value |
| 558 | auto readPromise = reader->read(env.js, 50); |
| 559 | |
| 560 | return env.context.awaitJs(env.js, |
| 561 | kj::mv(readPromise) |
| 562 | .then(env.js, |
| 563 | [&testCompleted, reader = kj::mv(reader)]( |
| 564 | jsg::Lock& js, DrainingReadResult&& result) mutable { |
| 565 | KJ_ASSERT(result.chunks.size() == 1); |
| 566 | // The internal implementation uses maxRead to allocate the buffer |
| 567 | KJ_ASSERT(result.chunks[0].size() <= 50); |
| 568 | KJ_ASSERT(!result.done); |
| 569 | reader->releaseLock(js); |
| 570 | testCompleted = true; |
| 571 | })); |
| 572 | }); |
| 573 | |
| 574 | KJ_ASSERT(testCompleted); |
| 575 | // Verify maxBytes was limited |
| 576 | KJ_ASSERT(lastMaxBytes == 50); |
| 577 | } |
| 578 | |
| 579 | KJ_TEST("DrainingReader with maxRead = 0 (internal stream)") { |
| 580 | // Test that the maxRead = 0 parameter is respected for internal streams |
| 581 | capnp::MallocMessageBuilder message; |
| 582 | auto flags = message.initRoot<CompatibilityFlags>(); |
| 583 | flags.setStreamsJavaScriptControllers(true); |
| 584 | |
| 585 | TestFixture fixture({.featureFlags = flags.asReader()}); |
| 586 | |
| 587 | bool testCompleted = false; |
| 588 | |
| 589 | fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 590 | class TestSource final: public ReadableStreamSource { |
| 591 | public: |
| 592 | explicit TestSource() = default; |
| 593 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 594 | KJ_FAIL_ASSERT("tryRead should not be called when maxRead = 0"); |
| 595 | } |
| 596 | }; |
| 597 | |
| 598 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>()); |
| 599 | |
| 600 | auto maybeReader = DrainingReader::create(env.js, *rs); |
| 601 | KJ_ASSERT(maybeReader != kj::none, "Failed to create DrainingReader"); |
| 602 | auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader)); |
| 603 | |
| 604 | auto readPromise = reader->read(env.js, 0); |
| 605 | |
| 606 | return env.context.awaitJs(env.js, |
| 607 | kj::mv(readPromise) |
| 608 | .then(env.js, |
| 609 | [&testCompleted, reader = kj::mv(reader)]( |
| 610 | jsg::Lock& js, DrainingReadResult&& result) mutable { |
| 611 | KJ_ASSERT(result.chunks.size() == 0); |
| 612 | KJ_ASSERT(!result.done); |
| 613 | reader->releaseLock(js); |
| 614 | testCompleted = true; |
| 615 | })); |
| 616 | }); |
| 617 | |
| 618 | KJ_ASSERT(testCompleted); |
| 619 | } |
| 620 | |
| 621 | KJ_TEST("DrainingReader on stream with pending closure (internal stream)") { |
| 622 | auto fixture = makeStreamTestFixture(); |
| 623 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 624 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>()); |
| 625 | rs->getController().setPendingClosure(); |
| 626 | |
| 627 | KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) { |
| 628 | bool rejected = false; |
| 629 | reader->read(env.js).catch_( |
| 630 | env.js, [&](jsg::Lock& js, jsg::Value reason) -> DrainingReadResult { |
| 631 | rejected = true; |
| 632 | auto msg = kj::str(reason.getHandle(js)); |
| 633 | KJ_ASSERT(msg.contains("closing"), msg); |
| 634 | return {.done = true}; |
| 635 | }); |
| 636 | env.js.runMicrotasks(); |
| 637 | KJ_ASSERT(rejected); |
| 638 | } else { |
| 639 | KJ_FAIL_ASSERT("Failed to create DrainingReader"); |
| 640 | } |
| 641 | }); |
| 642 | } |
| 643 | |
| 644 | KJ_TEST("DrainingReader on closed stream (internal stream)") { |
| 645 | auto fixture = makeStreamTestFixture(); |
| 646 | |
| 647 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 648 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<EofSource>()); |
| 649 | rs->getController().cancel(env.js, kj::none); |
| 650 | |
| 651 | KJ_IF_SOME(reader, DrainingReader::create(env.js, *rs)) { |
| 652 | bool done = false; |
| 653 | reader->read(env.js).then(env.js, [&](jsg::Lock&, DrainingReadResult&& result) { |
| 654 | done = true; |
| 655 | KJ_ASSERT(result.done); |
| 656 | KJ_ASSERT(result.chunks.size() == 0); |
| 657 | }); |
| 658 | env.js.runMicrotasks(); |
| 659 | KJ_ASSERT(done); |
| 660 | } else { |
| 661 | KJ_FAIL_ASSERT("Failed to create DrainingReader"); |
| 662 | } |
| 663 | }); |
| 664 | } |
| 665 | |
| 666 | KJ_TEST("DrainingReader error propagation (internal stream)") { |
| 667 | // Tests that I/O errors from tryRead are propagated and put the stream in errored state. |
| 668 | auto fixture = makeStreamTestFixture(); |
| 669 | |
| 670 | bool testCompleted = false; |
| 671 | KJ_EXPECT_LOG(ERROR, "Simulated I/O error"); |
| 672 | |
| 673 | fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 674 | class TestSource final: public ReadableStreamSource { |
| 675 | public: |
| 676 | kj::Promise<size_t> tryRead(void*, size_t, size_t) override { |
| 677 | return KJ_EXCEPTION(FAILED, "Simulated I/O error"); |
| 678 | } |
| 679 | }; |
| 680 | |
| 681 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>()); |
| 682 | auto maybeReader = DrainingReader::create(env.js, *rs); |
| 683 | auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader)); |
| 684 | |
| 685 | return env.context.awaitJs(env.js, |
| 686 | reader->read(env.js) |
| 687 | .catch_(env.js, |
| 688 | [&, reader = kj::mv(reader)]( |
| 689 | jsg::Lock& js, jsg::Value) mutable -> DrainingReadResult { |
| 690 | reader->releaseLock(js); |
| 691 | testCompleted = true; |
| 692 | return {.done = true}; |
| 693 | }).then(env.js, [](jsg::Lock&, DrainingReadResult&&) {})); |
| 694 | }); |
| 695 | |
| 696 | KJ_ASSERT(testCompleted); |
| 697 | } |
| 698 | |
| 699 | KJ_TEST("DrainingReader concurrent read rejection (internal stream)") { |
| 700 | auto fixture = makeStreamTestFixture(); |
| 701 | bool testCompleted = false; |
| 702 | |
| 703 | fixture.runInIoContext([&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 704 | auto paf = kj::newPromiseAndFulfiller<size_t>(); |
| 705 | |
| 706 | class TestSource final: public ReadableStreamSource { |
| 707 | public: |
| 708 | explicit TestSource(kj::Promise<size_t> p): promise(kj::mv(p)) {} |
| 709 | kj::Promise<size_t> tryRead(void*, size_t, size_t) override { |
| 710 | return kj::mv(promise); |
| 711 | } |
| 712 | kj::Promise<size_t> promise; |
| 713 | }; |
| 714 | |
| 715 | auto rs = env.js.alloc<ReadableStream>(env.context, kj::heap<TestSource>(kj::mv(paf.promise))); |
| 716 | auto maybeReader = DrainingReader::create(env.js, *rs); |
| 717 | auto reader = kj::mv(KJ_ASSERT_NONNULL(maybeReader)); |
| 718 | |
| 719 | auto firstRead = reader->read(env.js); |
| 720 | |
| 721 | // Second read while first is pending should reject synchronously |
| 722 | bool rejected = false; |
| 723 | reader->read(env.js).catch_( |
| 724 | env.js, [&](jsg::Lock& js, jsg::Value reason) -> DrainingReadResult { |
| 725 | rejected = true; |
| 726 | auto msg = kj::str(reason.getHandle(js)); |
| 727 | KJ_ASSERT(msg.contains("single pending read"), msg); |
| 728 | return {.done = true}; |
| 729 | }); |
| 730 | env.js.runMicrotasks(); |
| 731 | KJ_ASSERT(rejected); |
| 732 | |
| 733 | paf.fulfiller->fulfill(0); // Complete first read with EOF |
| 734 | |
| 735 | return env.context.awaitJs(env.js, |
| 736 | kj::mv(firstRead).then( |
| 737 | env.js, [&, reader = kj::mv(reader)](jsg::Lock& js, DrainingReadResult&&) mutable { |
| 738 | reader->releaseLock(js); |
| 739 | testCompleted = true; |
| 740 | })); |
| 741 | }); |
| 742 | |
| 743 | KJ_ASSERT(testCompleted); |
| 744 | } |
| 745 | |
| 746 | // ====================================================================================== |
| 747 | // ReadableStreamBYOBReader validation tests |
| 748 | |
| 749 | KJ_TEST("ReadableStreamBYOBReader rejects read with zero-sized buffer") { |
| 750 | KJ_EXPECT_LOG(ERROR, "read() on a BYOB reader requires a positive-sized TypedArray"); |
| 751 | auto fixture = makeStreamTestFixture(); |
| 752 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 753 | auto rs = makeByteStream(env.js); |
| 754 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 755 | |
| 756 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 0); |
| 757 | auto view = v8::Uint8Array::New(buffer, 0, 0); |
| 758 | |
| 759 | bool rejected = false; |
| 760 | reader->read(env.js, view, kj::none) |
| 761 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 762 | rejected = true; |
| 763 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 764 | KJ_ASSERT(ex.getDescription().contains( |
| 765 | "read() on a BYOB reader requires a positive-sized TypedArray"), |
| 766 | ex); |
| 767 | return {.done = true}; |
| 768 | }); |
| 769 | env.js.runMicrotasks(); |
| 770 | KJ_ASSERT(rejected, "Expected read() to reject with zero-sized buffer"); |
| 771 | }); |
| 772 | } |
| 773 | |
| 774 | KJ_TEST("ReadableStreamBYOBReader rejects read with atLeast=0") { |
| 775 | auto fixture = makeStreamTestFixture(); |
| 776 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 777 | auto rs = makeByteStream(env.js); |
| 778 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 779 | |
| 780 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10); |
| 781 | auto view = v8::Uint8Array::New(buffer, 0, 10); |
| 782 | |
| 783 | bool rejected = false; |
| 784 | reader->readAtLeast(env.js, 0, view) |
| 785 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 786 | rejected = true; |
| 787 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 788 | KJ_ASSERT( |
| 789 | ex.getDescription().contains("Requested invalid minimum number of bytes to read (0)"), |
| 790 | ex); |
| 791 | return {.done = true}; |
| 792 | }); |
| 793 | env.js.runMicrotasks(); |
| 794 | KJ_ASSERT(rejected, "Expected readAtLeast() to reject with atLeast=0"); |
| 795 | }); |
| 796 | } |
| 797 | |
| 798 | KJ_TEST("ReadableStreamBYOBReader rejects read when atLeast exceeds buffer size") { |
| 799 | auto fixture = makeStreamTestFixture(); |
| 800 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 801 | auto rs = makeByteStream(env.js); |
| 802 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 803 | |
| 804 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10); |
| 805 | auto view = v8::Uint8Array::New(buffer, 0, 10); |
| 806 | |
| 807 | bool rejected = false; |
| 808 | reader->readAtLeast(env.js, 20, view) |
| 809 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 810 | rejected = true; |
| 811 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 812 | KJ_ASSERT( |
| 813 | ex.getDescription().contains("Minimum bytes to read (20) exceeds size of buffer (10)"), |
| 814 | ex); |
| 815 | return {.done = true}; |
| 816 | }); |
| 817 | env.js.runMicrotasks(); |
| 818 | KJ_ASSERT(rejected, "Expected readAtLeast() to reject when atLeast exceeds buffer size"); |
| 819 | }); |
| 820 | } |
| 821 | |
| 822 | KJ_TEST("ReadableStreamBYOBReader readAtLeast with element count within capacity succeeds") { |
| 823 | // readAtLeast() treats its first argument as an element count. |
| 824 | // readAtLeast(10, Uint32Array(10)) → 10 elements * 4 bytes = 40 bytes, buffer is 40 → OK. |
| 825 | auto fixture = makeStreamTestFixture(); |
| 826 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 827 | auto rs = makeByteStream(env.js); |
| 828 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 829 | |
| 830 | // Uint32Array: element size 4, byteLength 40, length 10 |
| 831 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40); |
| 832 | auto view = v8::Uint32Array::New(buffer, 0, 10); |
| 833 | |
| 834 | bool rejected = false; |
| 835 | reader->readAtLeast(env.js, 10, view) |
| 836 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 837 | rejected = true; |
| 838 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 839 | KJ_FAIL_ASSERT("readAtLeast(10) on 10-element Uint32Array should not reject", ex); |
| 840 | return {.done = true}; |
| 841 | }); |
| 842 | env.js.runMicrotasks(); |
| 843 | KJ_ASSERT(!rejected); |
| 844 | }); |
| 845 | } |
| 846 | |
| 847 | KJ_TEST("ReadableStreamBYOBReader readAtLeast rejects when element count exceeds capacity") { |
| 848 | // Regression test for the bug: before the fix, validation compared |
| 849 | // the raw element count against byteLength (mixed units), so large element |
| 850 | // counts slipped through and produced an impossible (minBytes > maxBytes) pair |
| 851 | // that caused stack overflow in decoders like brotli. |
| 852 | // readAtLeast(11, Uint32Array(10)) → 11 * 4 = 44 bytes > 40 byte buffer → reject. |
| 853 | auto fixture = makeStreamTestFixture(); |
| 854 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 855 | auto rs = makeByteStream(env.js); |
| 856 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 857 | |
| 858 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40); |
| 859 | auto view = v8::Uint32Array::New(buffer, 0, 10); |
| 860 | |
| 861 | bool rejected = false; |
| 862 | reader->readAtLeast(env.js, 11, view) |
| 863 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 864 | rejected = true; |
| 865 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 866 | KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex); |
| 867 | return {.done = true}; |
| 868 | }); |
| 869 | env.js.runMicrotasks(); |
| 870 | KJ_ASSERT(rejected, "readAtLeast(11) on 10-element Uint32Array should reject"); |
| 871 | }); |
| 872 | } |
| 873 | |
| 874 | KJ_TEST("ReadableStreamBYOBReader readAtLeast rejects byteLength as element count") { |
| 875 | // readAtLeast(view.byteLength, view) with a Uint32Array. |
| 876 | // byteLength=4096, element count interpretation → 4096 * 4 = 16384 > 4096 → must reject. |
| 877 | auto fixture = makeStreamTestFixture(); |
| 878 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 879 | auto rs = makeByteStream(env.js); |
| 880 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 881 | |
| 882 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 4096); |
| 883 | auto view = v8::Uint32Array::New(buffer, 0, 1024); |
| 884 | |
| 885 | bool rejected = false; |
| 886 | reader->readAtLeast(env.js, 4096, view) |
| 887 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 888 | rejected = true; |
| 889 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 890 | KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex); |
| 891 | return {.done = true}; |
| 892 | }); |
| 893 | env.js.runMicrotasks(); |
| 894 | KJ_ASSERT(rejected, "readAtLeast(4096) on 1024-element Uint32Array must reject"); |
| 895 | }); |
| 896 | } |
| 897 | |
| 898 | KJ_TEST("ReadableStreamBYOBReader read() with min exceeding element capacity rejects") { |
| 899 | // read(view, {min: N}) where N is in elements. Before the fix, validation |
| 900 | // compared element count against byte length (mixed units), so large values |
| 901 | // slipped through. |
| 902 | // min=11 elements on Uint32Array(10) → 44 bytes > 40 → reject. |
| 903 | auto fixture = makeStreamTestFixture(); |
| 904 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 905 | auto rs = makeByteStream(env.js); |
| 906 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 907 | |
| 908 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 40); |
| 909 | auto view = v8::Uint32Array::New(buffer, 0, 10); |
| 910 | |
| 911 | ReadableStreamBYOBReader::ReadableStreamBYOBReaderReadOptions opts; |
| 912 | opts.min = 11; |
| 913 | bool rejected = false; |
| 914 | reader->read(env.js, view, kj::mv(opts)) |
| 915 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 916 | rejected = true; |
| 917 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 918 | KJ_ASSERT(ex.getDescription().contains("exceeds size of buffer"), ex); |
| 919 | return {.done = true}; |
| 920 | }); |
| 921 | env.js.runMicrotasks(); |
| 922 | KJ_ASSERT(rejected, "read() with min=11 on 10-element Uint32Array should reject"); |
| 923 | }); |
| 924 | } |
| 925 | |
| 926 | KJ_TEST("ReadableStreamBYOBReader rejects read after releaseLock") { |
| 927 | auto fixture = makeStreamTestFixture(); |
| 928 | fixture.runInIoContext([&](const TestFixture::Environment& env) { |
| 929 | auto rs = makeByteStream(env.js); |
| 930 | auto reader = ReadableStreamBYOBReader::constructor(env.js, rs.addRef()); |
| 931 | reader->releaseLock(env.js); |
| 932 | |
| 933 | auto buffer = v8::ArrayBuffer::New(env.js.v8Isolate, 10); |
| 934 | auto view = v8::Uint8Array::New(buffer, 0, 10); |
| 935 | |
| 936 | bool rejected = false; |
| 937 | reader->read(env.js, view, kj::none) |
| 938 | .catch_(env.js, [&](jsg::Lock& js, jsg::Value reason) -> ReadResult { |
| 939 | rejected = true; |
| 940 | auto ex = js.exceptionToKj(kj::mv(reason)); |
| 941 | KJ_ASSERT(ex.getDescription().contains("This ReadableStream reader has been released"), ex); |
| 942 | return {.done = true}; |
| 943 | }); |
| 944 | env.js.runMicrotasks(); |
| 945 | KJ_ASSERT(rejected, "Expected read() to reject after releaseLock"); |
| 946 | }); |
| 947 | } |
| 948 | |
| 949 | } // namespace |
| 950 | } // namespace workerd::api |