File
Blob: src/workerd/api/streams/readable.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 "readable.h" |
| 6 | |
| 7 | #include "internal.h" |
| 8 | #include "writable.h" |
| 9 | |
| 10 | #include <workerd/api/system-streams.h> |
| 11 | #include <workerd/api/worker-rpc.h> |
| 12 | #include <workerd/io/features.h> |
| 13 | #include <workerd/jsg/jsg.h> |
| 14 | |
| 15 | namespace workerd::api { |
| 16 | |
| 17 | ReaderImpl::ReaderImpl(ReadableStreamController::Reader& reader) |
| 18 | : ioContext(tryGetIoContext()), |
| 19 | reader(reader), |
| 20 | state(ReaderState::create<Initial>()) {} |
| 21 | |
| 22 | ReaderImpl::~ReaderImpl() noexcept(false) { |
| 23 | KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { |
| 24 | attached.stream->getController().releaseReader(reader, kj::none); |
| 25 | } |
| 26 | } |
| 27 | |
| 28 | void ReaderImpl::attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) { |
| 29 | KJ_ASSERT(state.is<Initial>()); |
| 30 | state.transitionTo<Attached>(controller.addRef()); |
| 31 | this->closedPromise = kj::mv(closedPromise); |
| 32 | } |
| 33 | |
| 34 | void ReaderImpl::detach() { |
| 35 | // Only transition from Attached to Closed. |
| 36 | // All other states (Initial, Closed, Released) are no-ops. |
| 37 | if (state.isActive()) { |
| 38 | state.transitionTo<Closed>(); |
| 39 | } |
| 40 | } |
| 41 | |
| 42 | jsg::Promise<void> ReaderImpl::cancel( |
| 43 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 44 | assertAttachedOrTerminal(); |
| 45 | if (state.is<Released>()) { |
| 46 | return js.rejectedPromise<void>( |
| 47 | js.v8TypeError("This ReadableStream reader has been released."_kj)); |
| 48 | } |
| 49 | if (state.is<Closed>()) { |
| 50 | return js.resolvedPromise(); |
| 51 | } |
| 52 | auto& attached = state.requireActiveUnsafe(); |
| 53 | // In some edge cases, this reader is the last thing holding a strong |
| 54 | // reference to the stream. Calling cancel might cause the readers strong |
| 55 | // reference to be cleared, so let's make sure we keep a reference to |
| 56 | // the stream at least until the call to cancel completes. |
| 57 | auto ref = attached.stream.addRef(); |
| 58 | return attached.stream->getController().cancel(js, maybeReason); |
| 59 | } |
| 60 | |
| 61 | jsg::MemoizedIdentity<jsg::Promise<void>>& ReaderImpl::getClosed() { |
| 62 | // The closed promise should always be set after the object is created so this assert |
| 63 | // should always be safe. |
| 64 | return KJ_ASSERT_NONNULL(closedPromise); |
| 65 | } |
| 66 | |
| 67 | void ReaderImpl::lockToStream(jsg::Lock& js, ReadableStream& stream) { |
| 68 | KJ_ASSERT(!stream.isLocked()); |
| 69 | KJ_ASSERT(stream.getController().lockReader(js, reader)); |
| 70 | } |
| 71 | |
| 72 | jsg::Promise<ReadResult> ReaderImpl::read( |
| 73 | jsg::Lock& js, kj::Maybe<ReadableStreamController::ByobOptions> byobOptions) { |
| 74 | assertAttachedOrTerminal(); |
| 75 | if (state.is<Released>()) { |
| 76 | return js.rejectedPromise<ReadResult>( |
| 77 | js.v8TypeError("This ReadableStream reader has been released."_kj)); |
| 78 | } |
| 79 | if (state.is<Closed>()) { |
| 80 | return js.rejectedPromise<ReadResult>( |
| 81 | js.v8TypeError("This ReadableStream has been closed."_kj)); |
| 82 | } |
| 83 | auto& attached = state.requireActiveUnsafe(); |
| 84 | KJ_IF_SOME(options, byobOptions) { |
| 85 | // Per the spec, we must perform these checks before disturbing the stream. |
| 86 | size_t atLeast = options.atLeast.orDefault(1); |
| 87 | |
| 88 | if (options.byteLength == 0) { |
| 89 | return js.rejectedPromise<ReadResult>( |
| 90 | js.v8TypeError("You must call read() on a \"byob\" reader with a positive-sized " |
| 91 | "TypedArray object."_kj)); |
| 92 | } |
| 93 | if (atLeast == 0) { |
| 94 | return js.rejectedPromise<ReadResult>(js.v8TypeError( |
| 95 | kj::str("Requested invalid minimum number of bytes to read (", atLeast, ")."))); |
| 96 | } |
| 97 | |
| 98 | // Both read() and readAtLeast() pass atLeast in element count. |
| 99 | // Convert to bytes before validation and forwarding to the controller. |
| 100 | jsg::BufferSource source(js, options.bufferView.getHandle(js)); |
| 101 | auto elementSize = source.getElementSize(); |
| 102 | atLeast = atLeast * elementSize; |
| 103 | |
| 104 | if (atLeast > options.byteLength) { |
| 105 | return js.rejectedPromise<ReadResult>(js.v8TypeError(kj::str("Minimum bytes to read (", |
| 106 | atLeast, ") exceeds size of buffer (", options.byteLength, ")."))); |
| 107 | } |
| 108 | |
| 109 | options.atLeast = atLeast; |
| 110 | } |
| 111 | |
| 112 | return KJ_ASSERT_NONNULL(attached.stream->getController().read(js, kj::mv(byobOptions))); |
| 113 | } |
| 114 | |
| 115 | void ReaderImpl::releaseLock(jsg::Lock& js) { |
| 116 | // TODO(soon): Releasing the lock should cancel any pending reads. This is a recent |
| 117 | // modification to the spec that we have not yet implemented. |
| 118 | assertAttachedOrTerminal(); |
| 119 | // Closed and Released states are no-ops. |
| 120 | KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { |
| 121 | // In some edge cases, this reader is the last thing holding a strong |
| 122 | // reference to the stream. Calling releaseLock might cause the readers strong |
| 123 | // reference to be cleared, so let's make sure we keep a reference to |
| 124 | // the stream at least until the call to releaseLock completes. |
| 125 | auto ref = attached.stream.addRef(); |
| 126 | attached.stream->getController().releaseReader(reader, js); |
| 127 | state.transitionTo<Released>(); |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | void ReaderImpl::visitForGc(jsg::GcVisitor& visitor) { |
| 132 | KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { |
| 133 | visitor.visit(attached.stream); |
| 134 | } |
| 135 | visitor.visit(closedPromise); |
| 136 | } |
| 137 | |
| 138 | // ====================================================================================== |
| 139 | |
| 140 | ReadableStreamDefaultReader::ReadableStreamDefaultReader(): impl(*this) {} |
| 141 | |
| 142 | jsg::Ref<ReadableStreamDefaultReader> ReadableStreamDefaultReader::constructor( |
| 143 | jsg::Lock& js, jsg::Ref<ReadableStream> stream) { |
| 144 | JSG_REQUIRE( |
| 145 | !stream->isLocked(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 146 | auto reader = js.alloc<ReadableStreamDefaultReader>(); |
| 147 | reader->lockToStream(js, *stream); |
| 148 | return kj::mv(reader); |
| 149 | } |
| 150 | |
| 151 | void ReadableStreamDefaultReader::attach( |
| 152 | ReadableStreamController& controller, jsg::Promise<void> closedPromise) { |
| 153 | impl.attach(controller, kj::mv(closedPromise)); |
| 154 | } |
| 155 | |
| 156 | jsg::Promise<void> ReadableStreamDefaultReader::cancel( |
| 157 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 158 | return impl.cancel(js, kj::mv(maybeReason)); |
| 159 | } |
| 160 | |
| 161 | void ReadableStreamDefaultReader::detach() { |
| 162 | impl.detach(); |
| 163 | } |
| 164 | |
| 165 | jsg::MemoizedIdentity<jsg::Promise<void>>& ReadableStreamDefaultReader::getClosed() { |
| 166 | return impl.getClosed(); |
| 167 | } |
| 168 | |
| 169 | void ReadableStreamDefaultReader::lockToStream(jsg::Lock& js, ReadableStream& stream) { |
| 170 | impl.lockToStream(js, stream); |
| 171 | } |
| 172 | |
| 173 | jsg::Promise<ReadResult> ReadableStreamDefaultReader::read(jsg::Lock& js) { |
| 174 | return impl.read(js, kj::none); |
| 175 | } |
| 176 | |
| 177 | void ReadableStreamDefaultReader::releaseLock(jsg::Lock& js) { |
| 178 | impl.releaseLock(js); |
| 179 | } |
| 180 | |
| 181 | void ReadableStreamDefaultReader::visitForGc(jsg::GcVisitor& visitor) { |
| 182 | visitor.visit(impl); |
| 183 | } |
| 184 | |
| 185 | // ====================================================================================== |
| 186 | |
| 187 | ReadableStreamBYOBReader::ReadableStreamBYOBReader(): impl(*this) {} |
| 188 | |
| 189 | jsg::Ref<ReadableStreamBYOBReader> ReadableStreamBYOBReader::constructor( |
| 190 | jsg::Lock& js, jsg::Ref<ReadableStream> stream) { |
| 191 | JSG_REQUIRE( |
| 192 | !stream->isLocked(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 193 | |
| 194 | if (!stream->getController().isClosedOrErrored()) { |
| 195 | JSG_REQUIRE(stream->getController().isByteOriented(), TypeError, |
| 196 | "This ReadableStream does not support BYOB reads."); |
| 197 | } |
| 198 | |
| 199 | auto reader = js.alloc<ReadableStreamBYOBReader>(); |
| 200 | reader->lockToStream(js, *stream); |
| 201 | return kj::mv(reader); |
| 202 | } |
| 203 | |
| 204 | void ReadableStreamBYOBReader::attach( |
| 205 | ReadableStreamController& controller, jsg::Promise<void> closedPromise) { |
| 206 | impl.attach(controller, kj::mv(closedPromise)); |
| 207 | } |
| 208 | |
| 209 | jsg::Promise<void> ReadableStreamBYOBReader::cancel( |
| 210 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 211 | return impl.cancel(js, kj::mv(maybeReason)); |
| 212 | } |
| 213 | |
| 214 | void ReadableStreamBYOBReader::detach() { |
| 215 | impl.detach(); |
| 216 | } |
| 217 | |
| 218 | jsg::MemoizedIdentity<jsg::Promise<void>>& ReadableStreamBYOBReader::getClosed() { |
| 219 | return impl.getClosed(); |
| 220 | } |
| 221 | |
| 222 | void ReadableStreamBYOBReader::lockToStream(jsg::Lock& js, ReadableStream& stream) { |
| 223 | impl.lockToStream(js, stream); |
| 224 | } |
| 225 | |
| 226 | jsg::Promise<ReadResult> ReadableStreamBYOBReader::read(jsg::Lock& js, |
| 227 | v8::Local<v8::ArrayBufferView> byobBuffer, |
| 228 | jsg::Optional<ReadableStreamBYOBReaderReadOptions> maybeOptions) { |
| 229 | static const ReadableStreamBYOBReaderReadOptions defaultOptions{}; |
| 230 | auto options = ReadableStreamController::ByobOptions{ |
| 231 | .bufferView = js.v8Ref(byobBuffer), |
| 232 | .byteOffset = byobBuffer->ByteOffset(), |
| 233 | .byteLength = byobBuffer->ByteLength(), |
| 234 | .atLeast = maybeOptions.orDefault(defaultOptions).min.orDefault(1), |
| 235 | .detachBuffer = FeatureFlags::get(js).getStreamsByobReaderDetachesBuffer(), |
| 236 | }; |
| 237 | return impl.read(js, kj::mv(options)); |
| 238 | } |
| 239 | |
| 240 | jsg::Promise<ReadResult> ReadableStreamBYOBReader::readAtLeast( |
| 241 | jsg::Lock& js, int minElements, v8::Local<v8::ArrayBufferView> byobBuffer) { |
| 242 | auto options = ReadableStreamController::ByobOptions{ |
| 243 | .bufferView = js.v8Ref(byobBuffer), |
| 244 | .byteOffset = byobBuffer->ByteOffset(), |
| 245 | .byteLength = byobBuffer->ByteLength(), |
| 246 | .atLeast = minElements, |
| 247 | .detachBuffer = true, |
| 248 | }; |
| 249 | return impl.read(js, kj::mv(options)); |
| 250 | } |
| 251 | |
| 252 | void ReadableStreamBYOBReader::releaseLock(jsg::Lock& js) { |
| 253 | impl.releaseLock(js); |
| 254 | } |
| 255 | |
| 256 | void ReadableStreamBYOBReader::visitForGc(jsg::GcVisitor& visitor) { |
| 257 | visitor.visit(impl); |
| 258 | } |
| 259 | |
| 260 | // ====================================================================================== |
| 261 | // DrainingReader implementation |
| 262 | |
| 263 | DrainingReader::DrainingReader(): ioContext(tryGetIoContext()) {} |
| 264 | |
| 265 | DrainingReader::~DrainingReader() noexcept(false) { |
| 266 | KJ_IF_SOME(stream, state.tryGet<Attached>()) { |
| 267 | stream->getController().releaseReader(*this, kj::none); |
| 268 | } |
| 269 | } |
| 270 | |
| 271 | kj::Maybe<kj::Own<DrainingReader>> DrainingReader::create(jsg::Lock& js, ReadableStream& stream) { |
| 272 | if (stream.isLocked()) { |
| 273 | return kj::none; |
| 274 | } |
| 275 | auto reader = kj::heap<DrainingReader>(); |
| 276 | if (!stream.getController().lockReader(js, *reader)) { |
| 277 | return kj::none; |
| 278 | } |
| 279 | return kj::mv(reader); |
| 280 | } |
| 281 | |
| 282 | void DrainingReader::attach( |
| 283 | ReadableStreamController& controller, jsg::Promise<void> closedPromise) { |
| 284 | KJ_ASSERT(state.is<Initial>()); |
| 285 | state = controller.addRef(); |
| 286 | this->closedPromise = kj::mv(closedPromise); |
| 287 | } |
| 288 | |
| 289 | void DrainingReader::detach() { |
| 290 | KJ_SWITCH_ONEOF(state) { |
| 291 | KJ_CASE_ONEOF(i, Initial) { |
| 292 | return; |
| 293 | } |
| 294 | KJ_CASE_ONEOF(stream, Attached) { |
| 295 | state.init<StreamStates::Closed>(); |
| 296 | return; |
| 297 | } |
| 298 | KJ_CASE_ONEOF(c, StreamStates::Closed) { |
| 299 | return; |
| 300 | } |
| 301 | KJ_CASE_ONEOF(r, Released) { |
| 302 | return; |
| 303 | } |
| 304 | } |
| 305 | KJ_UNREACHABLE; |
| 306 | } |
| 307 | |
| 308 | jsg::Promise<DrainingReadResult> DrainingReader::read(jsg::Lock& js, size_t maxRead) { |
| 309 | KJ_SWITCH_ONEOF(state) { |
| 310 | KJ_CASE_ONEOF(i, Initial) { |
| 311 | KJ_FAIL_ASSERT("this reader was never attached"); |
| 312 | } |
| 313 | KJ_CASE_ONEOF(stream, Attached) { |
| 314 | auto& controller = stream->getController(); |
| 315 | KJ_IF_SOME(result, controller.drainingRead(js, maxRead)) { |
| 316 | return kj::mv(result); |
| 317 | } |
| 318 | return js.rejectedPromise<DrainingReadResult>( |
| 319 | js.v8TypeError("Unable to perform draining read on this stream."_kj)); |
| 320 | } |
| 321 | KJ_CASE_ONEOF(r, Released) { |
| 322 | return js.rejectedPromise<DrainingReadResult>( |
| 323 | js.v8TypeError("This ReadableStream reader has been released."_kj)); |
| 324 | } |
| 325 | KJ_CASE_ONEOF(c, StreamStates::Closed) { |
| 326 | return js.resolvedPromise(DrainingReadResult{ |
| 327 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 328 | .done = true, |
| 329 | }); |
| 330 | } |
| 331 | } |
| 332 | KJ_UNREACHABLE; |
| 333 | } |
| 334 | |
| 335 | jsg::Promise<void> DrainingReader::cancel( |
| 336 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 337 | KJ_SWITCH_ONEOF(state) { |
| 338 | KJ_CASE_ONEOF(i, Initial) { |
| 339 | KJ_FAIL_ASSERT("this reader was never attached"); |
| 340 | } |
| 341 | KJ_CASE_ONEOF(stream, Attached) { |
| 342 | auto ref = stream.addRef(); |
| 343 | return stream->getController().cancel(js, maybeReason); |
| 344 | } |
| 345 | KJ_CASE_ONEOF(r, Released) { |
| 346 | return js.rejectedPromise<void>( |
| 347 | js.v8TypeError("This ReadableStream reader has been released."_kj)); |
| 348 | } |
| 349 | KJ_CASE_ONEOF(c, StreamStates::Closed) { |
| 350 | return js.resolvedPromise(); |
| 351 | } |
| 352 | } |
| 353 | KJ_UNREACHABLE; |
| 354 | } |
| 355 | |
| 356 | void DrainingReader::releaseLock(jsg::Lock& js) { |
| 357 | KJ_SWITCH_ONEOF(state) { |
| 358 | KJ_CASE_ONEOF(i, Initial) { |
| 359 | KJ_FAIL_ASSERT("this reader was never attached"); |
| 360 | } |
| 361 | KJ_CASE_ONEOF(stream, Attached) { |
| 362 | auto ref = stream.addRef(); |
| 363 | stream->getController().releaseReader(*this, js); |
| 364 | state.init<Released>(); |
| 365 | return; |
| 366 | } |
| 367 | KJ_CASE_ONEOF(c, StreamStates::Closed) { |
| 368 | return; |
| 369 | } |
| 370 | KJ_CASE_ONEOF(r, Released) { |
| 371 | return; |
| 372 | } |
| 373 | } |
| 374 | KJ_UNREACHABLE; |
| 375 | } |
| 376 | |
| 377 | bool DrainingReader::isAttached() const { |
| 378 | return state.is<Attached>(); |
| 379 | } |
| 380 | |
| 381 | void DrainingReader::visitForGc(jsg::GcVisitor& visitor) { |
| 382 | KJ_IF_SOME(stream, state.tryGet<Attached>()) { |
| 383 | visitor.visit(stream); |
| 384 | } |
| 385 | visitor.visit(closedPromise); |
| 386 | } |
| 387 | |
| 388 | // ====================================================================================== |
| 389 | |
| 390 | ReadableStream::ReadableStream(IoContext& ioContext, kj::Own<ReadableStreamSource> source) |
| 391 | : ReadableStream(newReadableStreamInternalController(ioContext, kj::mv(source))) {} |
| 392 | |
| 393 | ReadableStream::ReadableStream(kj::Own<ReadableStreamController> controller) |
| 394 | : ioContext(tryGetIoContext()), |
| 395 | controller(kj::mv(controller)) { |
| 396 | getController().setOwnerRef(*this); |
| 397 | } |
| 398 | |
| 399 | void ReadableStream::visitForGc(jsg::GcVisitor& visitor) { |
| 400 | visitor.visit(getController()); |
| 401 | KJ_IF_SOME(pair, eofResolverPair) { |
| 402 | visitor.visit(pair.resolver); |
| 403 | visitor.visit(pair.promise); |
| 404 | } |
| 405 | } |
| 406 | |
| 407 | jsg::Ref<ReadableStream> ReadableStream::addRef() { |
| 408 | return JSG_THIS; |
| 409 | } |
| 410 | |
| 411 | bool ReadableStream::isDisturbed() { |
| 412 | return getController().isDisturbed(); |
| 413 | } |
| 414 | |
| 415 | bool ReadableStream::isLocked() { |
| 416 | return getController().isLockedToReader(); |
| 417 | } |
| 418 | |
| 419 | jsg::Promise<void> ReadableStream::onEof(jsg::Lock& js) { |
| 420 | eofResolverPair = js.newPromiseAndResolver<void>(); |
| 421 | return kj::mv(KJ_ASSERT_NONNULL(eofResolverPair).promise); |
| 422 | } |
| 423 | |
| 424 | void ReadableStream::signalEof(jsg::Lock& js) { |
| 425 | KJ_IF_SOME(pair, eofResolverPair) { |
| 426 | pair.resolver.resolve(js); |
| 427 | } |
| 428 | } |
| 429 | |
| 430 | ReadableStreamController& ReadableStream::getController() { |
| 431 | return *controller; |
| 432 | } |
| 433 | |
| 434 | jsg::Promise<void> ReadableStream::cancel( |
| 435 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 436 | if (isLocked()) { |
| 437 | return js.rejectedPromise<void>( |
| 438 | js.v8TypeError("This ReadableStream is currently locked to a reader."_kj)); |
| 439 | } |
| 440 | return getController().cancel(js, maybeReason); |
| 441 | } |
| 442 | |
| 443 | ReadableStream::Reader ReadableStream::getReader( |
| 444 | jsg::Lock& js, jsg::Optional<GetReaderOptions> options) { |
| 445 | JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 446 | |
| 447 | bool isByob = false; |
| 448 | KJ_IF_SOME(o, options) { |
| 449 | KJ_IF_SOME(mode, o.mode) { |
| 450 | JSG_REQUIRE( |
| 451 | mode == "byob", TypeError, "mode must be undefined or 'byob' in call to getReader()."); |
| 452 | // No need to check that the ReadableStream implementation is a byte stream: the first |
| 453 | // invocation of read() will do that for us and throw if necessary. Also, we should really |
| 454 | // just support reading non-byte streams with BYOB readers. |
| 455 | isByob = true; |
| 456 | } |
| 457 | } |
| 458 | |
| 459 | if (isByob) { |
| 460 | return ReadableStreamBYOBReader::constructor(js, JSG_THIS); |
| 461 | } |
| 462 | return ReadableStreamDefaultReader::constructor(js, JSG_THIS); |
| 463 | } |
| 464 | |
| 465 | jsg::Ref<ReadableStream::ReadableStreamAsyncIterator> ReadableStream::values( |
| 466 | jsg::Lock& js, jsg::Optional<ValuesOptions> options) { |
| 467 | static const auto defaultOptions = ValuesOptions{}; |
| 468 | return js.alloc<ReadableStreamAsyncIterator>(AsyncIteratorState{.ioContext = ioContext, |
| 469 | .reader = ReadableStreamDefaultReader::constructor(js, JSG_THIS), |
| 470 | .preventCancel = options.orDefault(defaultOptions).preventCancel.orDefault(false)}); |
| 471 | } |
| 472 | |
| 473 | jsg::Ref<ReadableStream> ReadableStream::pipeThrough( |
| 474 | jsg::Lock& js, Transform transform, jsg::Optional<PipeToOptions> maybeOptions) { |
| 475 | auto& controller = getController(); |
| 476 | |
| 477 | auto& destination = transform.writable->getController(); |
| 478 | JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 479 | JSG_REQUIRE(!destination.isLockedToWriter(), TypeError, |
| 480 | "This WritableStream is currently locked to a writer."); |
| 481 | |
| 482 | auto options = kj::mv(maybeOptions).orDefault({}); |
| 483 | options.pipeThrough = true; |
| 484 | // The lambda intentionally captures self as a visitable reference, ensuring |
| 485 | // JSG_THIS stays alive until the pipe promise resolves. |
| 486 | controller.pipeTo(js, destination, kj::mv(options)) |
| 487 | .then(js, |
| 488 | JSG_VISITABLE_LAMBDA( |
| 489 | (self = JSG_THIS), (self), (jsg::Lock& js) { return js.resolvedPromise(); })) |
| 490 | .markAsHandled(js); |
| 491 | return kj::mv(transform.readable); |
| 492 | } |
| 493 | |
| 494 | jsg::Promise<void> ReadableStream::pipeTo(jsg::Lock& js, |
| 495 | jsg::Ref<WritableStream> destination, |
| 496 | jsg::Optional<PipeToOptions> maybeOptions) { |
| 497 | if (isLocked()) { |
| 498 | return js.rejectedPromise<void>( |
| 499 | js.v8TypeError("This ReadableStream is currently locked to a reader."_kj)); |
| 500 | } |
| 501 | |
| 502 | if (destination->getController().isLockedToWriter()) { |
| 503 | return js.rejectedPromise<void>( |
| 504 | js.v8TypeError("This WritableStream is currently locked to a writer"_kj)); |
| 505 | } |
| 506 | |
| 507 | auto options = kj::mv(maybeOptions).orDefault({}); |
| 508 | return getController().pipeTo(js, destination->getController(), kj::mv(options)); |
| 509 | } |
| 510 | |
| 511 | kj::Array<jsg::Ref<ReadableStream>> ReadableStream::tee(jsg::Lock& js) { |
| 512 | JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader,"); |
| 513 | auto tee = getController().tee(js); |
| 514 | return kj::arr(kj::mv(tee.branch1), kj::mv(tee.branch2)); |
| 515 | } |
| 516 | |
| 517 | jsg::JsString ReadableStream::inspectState(jsg::Lock& js) { |
| 518 | if (controller->isClosedOrErrored()) { |
| 519 | return js.strIntern(controller->isClosed() ? "closed"_kj : "errored"_kj); |
| 520 | } else { |
| 521 | return js.strIntern("readable"_kj); |
| 522 | } |
| 523 | } |
| 524 | |
| 525 | bool ReadableStream::inspectSupportsBYOB() { |
| 526 | return controller->isByteOriented(); |
| 527 | } |
| 528 | |
| 529 | jsg::Optional<uint64_t> ReadableStream::inspectLength() { |
| 530 | return tryGetLength(StreamEncoding::IDENTITY); |
| 531 | } |
| 532 | |
| 533 | jsg::Promise<kj::Maybe<jsg::Value>> ReadableStream::nextFunction( |
| 534 | jsg::Lock& js, AsyncIteratorState& state) { |
| 535 | return state.reader->read(js).then( |
| 536 | js, [reader = state.reader.addRef()](jsg::Lock& js, ReadResult result) mutable { |
| 537 | if (result.done) { |
| 538 | reader->releaseLock(js); |
| 539 | return js.resolvedPromise(kj::Maybe<jsg::Value>(kj::none)); |
| 540 | } |
| 541 | return js.resolvedPromise<kj::Maybe<jsg::Value>>(kj::mv(result.value)); |
| 542 | }); |
| 543 | } |
| 544 | |
| 545 | jsg::Promise<void> ReadableStream::returnFunction( |
| 546 | jsg::Lock& js, AsyncIteratorState& state, jsg::Optional<jsg::Value>& value) { |
| 547 | if (state.reader.get() != nullptr) { |
| 548 | auto reader = kj::mv(state.reader); |
| 549 | if (!state.preventCancel) { |
| 550 | auto promise = reader->cancel(js, value.map([&](jsg::Value& v) { return v.getHandle(js); })); |
| 551 | reader->releaseLock(js); |
| 552 | auto result = promise.then(js, |
| 553 | JSG_VISITABLE_LAMBDA((reader = kj::mv(reader)), (reader), (jsg::Lock& js) { |
| 554 | // Ensure that the reader is not garbage collected until the cancel promise resolves. |
| 555 | return js.resolvedPromise(); |
| 556 | })); |
| 557 | // When the stream is already errored, cancel() returns a rejected promise |
| 558 | // that propagates through the .then() chain. Mark it as handled so V8 does |
| 559 | // not fire unhandledrejection events during iterator teardown. |
| 560 | result.markAsHandled(js); |
| 561 | return kj::mv(result); |
| 562 | } |
| 563 | |
| 564 | reader->releaseLock(js); |
| 565 | } |
| 566 | return js.resolvedPromise(); |
| 567 | } |
| 568 | |
| 569 | jsg::Ref<ReadableStream> ReadableStream::detach(jsg::Lock& js, bool ignoreDisturbed) { |
| 570 | JSG_REQUIRE( |
| 571 | !isDisturbed() || ignoreDisturbed, TypeError, "The ReadableStream has already been read."); |
| 572 | JSG_REQUIRE(!isLocked(), TypeError, "The ReadableStream has been locked to a reader."); |
| 573 | return js.alloc<ReadableStream>(getController().detach(js, ignoreDisturbed)); |
| 574 | } |
| 575 | |
| 576 | kj::Maybe<uint64_t> ReadableStream::tryGetLength(StreamEncoding encoding) { |
| 577 | return getController().tryGetLength(encoding); |
| 578 | } |
| 579 | |
| 580 | kj::Promise<DeferredProxy<void>> ReadableStream::pumpTo( |
| 581 | jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) { |
| 582 | JSG_REQUIRE( |
| 583 | IoContext::hasCurrent(), Error, "Unable to consume this ReadableStream outside of a request"); |
| 584 | JSG_REQUIRE(!isLocked(), TypeError, "The ReadableStream has been locked to a reader."); |
| 585 | return getController().pumpTo(js, kj::mv(sink), end); |
| 586 | } |
| 587 | |
| 588 | jsg::Ref<ReadableStream> ReadableStream::constructor(jsg::Lock& js, |
| 589 | jsg::Optional<UnderlyingSource> underlyingSource, |
| 590 | jsg::Optional<StreamQueuingStrategy> queuingStrategy) { |
| 591 | |
| 592 | JSG_REQUIRE(FeatureFlags::get(js).getStreamsJavaScriptControllers(), Error, |
| 593 | "To use the new ReadableStream() constructor, enable the " |
| 594 | "streams_enable_constructors compatibility flag. " |
| 595 | "Refer to the docs for more information: https://developers.cloudflare.com/workers/platform/compatibility-dates/#compatibility-flags"); |
| 596 | // We account for the memory usage of the ReadableStream and its controller together because their |
| 597 | // lifetimes are identical and memory accounting itself has a memory overhead. |
| 598 | auto controller = newReadableStreamJsController(); |
| 599 | auto stream = js.allocAccounted<ReadableStream>( |
| 600 | sizeof(ReadableStream) + controller->jsgGetMemorySelfSize(), kj::mv(controller)); |
| 601 | stream->getController().setup(js, kj::mv(underlyingSource), kj::mv(queuingStrategy)); |
| 602 | return kj::mv(stream); |
| 603 | } |
| 604 | |
| 605 | jsg::Optional<uint32_t> ByteLengthQueuingStrategy::size( |
| 606 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeValue) { |
| 607 | KJ_IF_SOME(value, maybeValue) { |
| 608 | if ((value)->IsArrayBuffer()) { |
| 609 | auto buffer = value.As<v8::ArrayBuffer>(); |
| 610 | return buffer->ByteLength(); |
| 611 | } else if ((value)->IsArrayBufferView()) { |
| 612 | auto view = value.As<v8::ArrayBufferView>(); |
| 613 | return view->ByteLength(); |
| 614 | } else { |
| 615 | // Per the WHATWG Streams spec, ByteLengthQueuingStrategy.size should return |
| 616 | // GetV(chunk, "byteLength"), which means getting the byteLength property |
| 617 | // from any object, not just ArrayBuffer/ArrayBufferView. |
| 618 | KJ_IF_SOME(obj, jsg::JsValue(value).tryCast<jsg::JsObject>()) { |
| 619 | auto byteLength = obj.get(js, "byteLength"_kj); |
| 620 | KJ_IF_SOME(num, byteLength.tryCast<jsg::JsNumber>()) { |
| 621 | KJ_IF_SOME(val, num.value(js)) { |
| 622 | return static_cast<uint32_t>(val); |
| 623 | } |
| 624 | } |
| 625 | } |
| 626 | } |
| 627 | } |
| 628 | return kj::none; |
| 629 | } |
| 630 | |
| 631 | namespace { |
| 632 | |
| 633 | // TODO(cleanup): These classes have been copied to external-pusher.c++. The copies here can be |
| 634 | // deleted as soon as we've switched from StreamSink to ExternalPusher and can delete all the |
| 635 | // StreamSink-related code. For now I'm not trying to avoid duplication. |
| 636 | |
| 637 | // HACK: We need as async pipe, like kj::newOneWayPipe(), except supporting explicit end(). So we |
| 638 | // wrap the two ends of the pipe in special adapters that track whether end() was called. |
| 639 | class ExplicitEndOutputPipeAdapter final: public capnp::ExplicitEndOutputStream { |
| 640 | public: |
| 641 | ExplicitEndOutputPipeAdapter( |
| 642 | kj::Own<kj::AsyncOutputStream> inner, kj::Own<kj::RefcountedWrapper<bool>> ended) |
| 643 | : inner(kj::mv(inner)), |
| 644 | ended(kj::mv(ended)) {} |
| 645 | |
| 646 | kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override { |
| 647 | return KJ_REQUIRE_NONNULL(inner)->write(buffer); |
| 648 | } |
| 649 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override { |
| 650 | return KJ_REQUIRE_NONNULL(inner)->write(pieces); |
| 651 | } |
| 652 | |
| 653 | kj::Maybe<kj::Promise<uint64_t>> tryPumpFrom( |
| 654 | kj::AsyncInputStream& input, uint64_t amount) override { |
| 655 | return KJ_REQUIRE_NONNULL(inner)->tryPumpFrom(input, amount); |
| 656 | } |
| 657 | |
| 658 | kj::Promise<void> whenWriteDisconnected() override { |
| 659 | return KJ_REQUIRE_NONNULL(inner)->whenWriteDisconnected(); |
| 660 | } |
| 661 | |
| 662 | kj::Promise<void> end() override { |
| 663 | // Signal to the other side that end() was actually called. |
| 664 | ended->getWrapped() = true; |
| 665 | inner = kj::none; |
| 666 | return kj::READY_NOW; |
| 667 | } |
| 668 | |
| 669 | private: |
| 670 | kj::Maybe<kj::Own<kj::AsyncOutputStream>> inner; |
| 671 | kj::Own<kj::RefcountedWrapper<bool>> ended; |
| 672 | }; |
| 673 | |
| 674 | class ExplicitEndInputPipeAdapter final: public kj::AsyncInputStream { |
| 675 | public: |
| 676 | ExplicitEndInputPipeAdapter(kj::Own<kj::AsyncInputStream> inner, |
| 677 | kj::Own<kj::RefcountedWrapper<bool>> ended, |
| 678 | kj::Maybe<uint64_t> expectedLength) |
| 679 | : inner(kj::mv(inner)), |
| 680 | ended(kj::mv(ended)), |
| 681 | expectedLength(expectedLength) {} |
| 682 | |
| 683 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 684 | size_t result = co_await inner->tryRead(buffer, minBytes, maxBytes); |
| 685 | |
| 686 | KJ_IF_SOME(l, expectedLength) { |
| 687 | KJ_ASSERT(result <= l); |
| 688 | l -= result; |
| 689 | if (l == 0) { |
| 690 | // If we got all the bytes we expected, we treat this as a successful end, because the |
| 691 | // underlying KJ pipe is not actually going to wait for the other side to drop. This is |
| 692 | // consistent with the behavior of Content-Length in HTTP anyway. |
| 693 | ended->getWrapped() = true; |
| 694 | } |
| 695 | } |
| 696 | |
| 697 | if (result < minBytes) { |
| 698 | // Verify that end() was called. |
| 699 | if (!ended->getWrapped()) { |
| 700 | JSG_FAIL_REQUIRE(Error, "ReadableStream received over RPC disconnected prematurely."); |
| 701 | } |
| 702 | } |
| 703 | co_return result; |
| 704 | } |
| 705 | |
| 706 | kj::Maybe<uint64_t> tryGetLength() override { |
| 707 | return inner->tryGetLength(); |
| 708 | } |
| 709 | |
| 710 | kj::Promise<uint64_t> pumpTo(kj::AsyncOutputStream& output, uint64_t amount) override { |
| 711 | return inner->pumpTo(output, amount); |
| 712 | } |
| 713 | |
| 714 | private: |
| 715 | kj::Own<kj::AsyncInputStream> inner; |
| 716 | kj::Own<kj::RefcountedWrapper<bool>> ended; |
| 717 | kj::Maybe<uint64_t> expectedLength; |
| 718 | }; |
| 719 | |
| 720 | // Wrapper around ReadableStreamSource that prevents deferred proxying. We need this for RPC |
| 721 | // streams because although they are "system streams", they become disconnected when the IoContext |
| 722 | // is destroyed, due to the JsRpcCustomEvent being canceled. |
| 723 | // |
| 724 | // TODO(someday): Devise a better way for RPC streams to extend the lifetime of the RPC session |
| 725 | // beyond the destruction of the IoContext, if it is being used for deferred proxying. |
| 726 | class NoDeferredProxyReadableStream final: public ReadableStreamSource { |
| 727 | public: |
| 728 | NoDeferredProxyReadableStream(kj::Own<ReadableStreamSource> inner, IoContext& ioctx) |
| 729 | : inner(kj::mv(inner)), |
| 730 | ioctx(ioctx) {} |
| 731 | |
| 732 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 733 | return inner->tryRead(buffer, minBytes, maxBytes); |
| 734 | } |
| 735 | |
| 736 | kj::Promise<DeferredProxy<void>> pumpTo(WritableStreamSink& output, bool end) override { |
| 737 | // Move the deferred proxy part of the task over to the non-deferred part. To do this, |
| 738 | // we use `ioctx.waitForDeferredProxy()`, which returns a single promise covering both parts |
| 739 | // (and, importantly, registering pending events where needed). Then, we add a noop deferred |
| 740 | // proxy to the end of that. |
| 741 | return addNoopDeferredProxy(ioctx.waitForDeferredProxy(inner->pumpTo(output, end))); |
| 742 | } |
| 743 | |
| 744 | StreamEncoding getPreferredEncoding() override { |
| 745 | return inner->getPreferredEncoding(); |
| 746 | } |
| 747 | |
| 748 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 749 | return inner->tryGetLength(encoding); |
| 750 | } |
| 751 | |
| 752 | void cancel(kj::Exception reason) override { |
| 753 | return inner->cancel(kj::mv(reason)); |
| 754 | } |
| 755 | |
| 756 | kj::Maybe<Tee> tryTee(uint64_t limit) override { |
| 757 | return inner->tryTee(limit).map([&](Tee tee) { |
| 758 | return Tee{.branches = { |
| 759 | kj::heap<NoDeferredProxyReadableStream>(kj::mv(tee.branches[0]), ioctx), |
| 760 | kj::heap<NoDeferredProxyReadableStream>(kj::mv(tee.branches[1]), ioctx), |
| 761 | }}; |
| 762 | }); |
| 763 | } |
| 764 | |
| 765 | private: |
| 766 | kj::Own<ReadableStreamSource> inner; |
| 767 | IoContext& ioctx; |
| 768 | }; |
| 769 | |
| 770 | } // namespace |
| 771 | |
| 772 | void ReadableStream::serialize(jsg::Lock& js, jsg::Serializer& serializer) { |
| 773 | // Serialize by effectively creating a `JsRpcStub` around this object and serializing that. |
| 774 | // Except we don't actually want to do _exactly_ that, because we do not want to actually create |
| 775 | // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by |
| 776 | // `JsRpcStub::serialize()`. |
| 777 | |
| 778 | auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError, |
| 779 | "ReadableStream can only be serialized for RPC."); |
| 780 | auto externalHandler = dynamic_cast<RpcSerializerExternalHandler*>(&handler); |
| 781 | JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError, |
| 782 | "ReadableStream can only be serialized for RPC."); |
| 783 | |
| 784 | // NOTE: We're counting on `pumpTo()`, below, to check that the stream is not locked or disturbed |
| 785 | // and other common checks. It's important that we don't modify the stream in any way before |
| 786 | // that call. |
| 787 | |
| 788 | IoContext& ioctx = IoContext::current(); |
| 789 | |
| 790 | auto& controller = getController(); |
| 791 | StreamEncoding encoding = controller.getPreferredEncoding(); |
| 792 | auto expectedLength = controller.tryGetLength(encoding); |
| 793 | |
| 794 | capnp::ByteStream::Client streamCap = [&]() { |
| 795 | KJ_IF_SOME(pusher, externalHandler->getExternalPusher()) { |
| 796 | auto req = pusher.pushByteStreamRequest(capnp::MessageSize{2, 0}); |
| 797 | KJ_IF_SOME(el, expectedLength) { |
| 798 | req.setLengthPlusOne(el + 1); |
| 799 | } |
| 800 | auto pipeline = req.sendForPipeline(); |
| 801 | |
| 802 | externalHandler->write([encoding, expectedLength, source = pipeline.getSource()]( |
| 803 | rpc::JsValue::External::Builder builder) mutable { |
| 804 | auto rs = builder.initReadableStream(); |
| 805 | rs.setStream(kj::mv(source)); |
| 806 | rs.setEncoding(encoding); |
| 807 | }); |
| 808 | |
| 809 | return pipeline.getSink(); |
| 810 | } else { |
| 811 | return externalHandler |
| 812 | ->writeStream( |
| 813 | [encoding, expectedLength](rpc::JsValue::External::Builder builder) mutable { |
| 814 | auto rs = builder.initReadableStream(); |
| 815 | rs.setEncoding(encoding); |
| 816 | KJ_IF_SOME(l, expectedLength) { |
| 817 | rs.getExpectedLength().setKnown(l); |
| 818 | } |
| 819 | }).castAs<capnp::ByteStream>(); |
| 820 | } |
| 821 | }(); |
| 822 | |
| 823 | kj::Own<capnp::ExplicitEndOutputStream> kjStream = |
| 824 | ioctx.getByteStreamFactory().capnpToKjExplicitEnd(kj::mv(streamCap)); |
| 825 | |
| 826 | auto sink = newSystemStream(kj::mv(kjStream), encoding, ioctx); |
| 827 | |
| 828 | ioctx.addTask( |
| 829 | ioctx.waitForDeferredProxy(pumpTo(js, kj::mv(sink), true)).catch_([](kj::Exception&& e) { |
| 830 | // Errors in pumpTo() are automatically propagated to the source and destination. We don't |
| 831 | // want to throw them from here since it'll cause an uncaught exception to be reported, even |
| 832 | // if the application actually does handle it! |
| 833 | })); |
| 834 | } |
| 835 | |
| 836 | jsg::Ref<ReadableStream> ReadableStream::deserialize( |
| 837 | jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) { |
| 838 | auto& handler = KJ_REQUIRE_NONNULL( |
| 839 | deserializer.getExternalHandler(), "got ReadableStream on non-RPC serialized object?"); |
| 840 | auto externalHandler = dynamic_cast<RpcDeserializerExternalHandler*>(&handler); |
| 841 | KJ_REQUIRE(externalHandler != nullptr, "got ReadableStream on non-RPC serialized object?"); |
| 842 | |
| 843 | auto reader = externalHandler->read(); |
| 844 | KJ_REQUIRE(reader.isReadableStream(), "external table slot type doesn't match serialization tag"); |
| 845 | |
| 846 | auto rs = reader.getReadableStream(); |
| 847 | auto encoding = rs.getEncoding(); |
| 848 | |
| 849 | KJ_REQUIRE( |
| 850 | static_cast<uint>(encoding) < capnp::Schema::from<StreamEncoding>().getEnumerants().size(), |
| 851 | "unknown StreamEncoding received from peer"); |
| 852 | |
| 853 | auto& ioctx = IoContext::current(); |
| 854 | |
| 855 | kj::Own<kj::AsyncInputStream> in; |
| 856 | if (rs.hasStream()) { |
| 857 | in = |
| 858 | ioctx.getExternalPusher()->unwrapStream(rs.getStream(), externalHandler->getDebugContext()); |
| 859 | } else { |
| 860 | kj::Maybe<uint64_t> expectedLength; |
| 861 | auto el = rs.getExpectedLength(); |
| 862 | if (el.isKnown()) { |
| 863 | expectedLength = el.getKnown(); |
| 864 | } |
| 865 | |
| 866 | auto pipe = kj::newOneWayPipe(expectedLength); |
| 867 | |
| 868 | auto endedFlag = kj::refcounted<kj::RefcountedWrapper<bool>>(false); |
| 869 | |
| 870 | auto out = kj::heap<ExplicitEndOutputPipeAdapter>(kj::mv(pipe.out), kj::addRef(*endedFlag)); |
| 871 | in = kj::heap<ExplicitEndInputPipeAdapter>(kj::mv(pipe.in), kj::mv(endedFlag), expectedLength); |
| 872 | |
| 873 | externalHandler->setLastStream(ioctx.getByteStreamFactory().kjToCapnp(kj::mv(out))); |
| 874 | } |
| 875 | |
| 876 | return js.alloc<ReadableStream>(ioctx, |
| 877 | kj::heap<NoDeferredProxyReadableStream>(newSystemStream(kj::mv(in), encoding, ioctx), ioctx)); |
| 878 | } |
| 879 | |
| 880 | kj::StringPtr ReaderImpl::jsgGetMemoryName() const { |
| 881 | return "ReaderImpl"_kjc; |
| 882 | } |
| 883 | |
| 884 | size_t ReaderImpl::jsgGetMemorySelfSize() const { |
| 885 | return sizeof(ReaderImpl); |
| 886 | } |
| 887 | |
| 888 | void ReaderImpl::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 889 | KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { |
| 890 | tracker.trackField("stream", attached.stream); |
| 891 | } |
| 892 | tracker.trackField("closedPromise", closedPromise); |
| 893 | } |
| 894 | |
| 895 | void ReadableStream::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 896 | tracker.trackField("controller", controller); |
| 897 | tracker.trackField("eofResolverPair", eofResolverPair); |
| 898 | } |
| 899 | |
| 900 | } // namespace workerd::api |