File
Blob: src/workerd/api/streams/standard.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 "standard.h" |
| 6 | |
| 7 | #include "readable.h" |
| 8 | #include "writable.h" |
| 9 | |
| 10 | #include <workerd/io/features.h> |
| 11 | #include <workerd/jsg/jsg.h> |
| 12 | #include <workerd/util/autogate.h> |
| 13 | #include <workerd/util/state-machine.h> |
| 14 | #include <workerd/util/weak-refs.h> |
| 15 | |
| 16 | #include <kj/debug.h> |
| 17 | #include <kj/vector.h> |
| 18 | |
| 19 | namespace workerd::api { |
| 20 | |
| 21 | using DefaultController = jsg::Ref<ReadableStreamDefaultController>; |
| 22 | using ByobController = jsg::Ref<ReadableByteStreamController>; |
| 23 | |
| 24 | namespace { |
| 25 | struct ValueReadable; |
| 26 | struct ByteReadable; |
| 27 | } // namespace |
| 28 | |
| 29 | // ======================================================================================= |
| 30 | // The Unlocked, Locked, ReaderLocked, and WriterLocked structs |
| 31 | // are used to track the current lock status of JavaScript-backed streams. |
| 32 | // All readable and writable streams begin in the Unlocked state. When a |
| 33 | // reader or writer are attached, the streams will transition into the |
| 34 | // ReaderLocked or WriterLocked state. When the reader is released, those |
| 35 | // will transition back to Unlocked. |
| 36 | // |
| 37 | // When a readable is piped to a writable, both will enter the PipeLocked state. |
| 38 | // (PipeLocked is defined within the ReadableLockImpl and WritableLockImpl classes |
| 39 | // below) When the pipe completes, both will transition back to Unlocked. |
| 40 | // |
| 41 | // When a ReadableStreamJsController is tee()'d, it will enter the locked state. |
| 42 | |
| 43 | namespace { |
| 44 | |
| 45 | // A utility class used by ReadableStreamJsController |
| 46 | // for implementing the reader lock in a consistent way (without duplicating any code). |
| 47 | template <typename Controller> |
| 48 | class ReadableLockImpl { |
| 49 | public: |
| 50 | using PipeController = ReadableStreamController::PipeController; |
| 51 | using Reader = ReadableStreamController::Reader; |
| 52 | |
| 53 | bool isLockedToReader() const { |
| 54 | return !state.template is<Unlocked>(); |
| 55 | } |
| 56 | |
| 57 | bool lockReader(jsg::Lock& js, Controller& self, Reader& reader); |
| 58 | |
| 59 | // See the comment for releaseReader in common.h for details on the use of maybeJs |
| 60 | void releaseReader(Controller& self, Reader& reader, kj::Maybe<jsg::Lock&> maybeJs); |
| 61 | |
| 62 | bool lock(); |
| 63 | |
| 64 | void onClose(jsg::Lock& js); |
| 65 | void onError(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 66 | |
| 67 | kj::Maybe<PipeController&> tryPipeLock(Controller& self); |
| 68 | |
| 69 | void visitForGc(jsg::GcVisitor& visitor); |
| 70 | |
| 71 | kj::StringPtr jsgGetMemoryName() const { |
| 72 | return "ReadableLockImpl"_kjc; |
| 73 | } |
| 74 | size_t jsgGetMemorySelfSize() const { |
| 75 | return sizeof(ReadableLockImpl); |
| 76 | } |
| 77 | void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 78 | KJ_SWITCH_ONEOF(state) { |
| 79 | KJ_CASE_ONEOF(locked, Locked) {} |
| 80 | KJ_CASE_ONEOF(unlocked, Unlocked) {} |
| 81 | KJ_CASE_ONEOF(pipeLocked, PipeLocked) {} |
| 82 | KJ_CASE_ONEOF(readerLocked, ReaderLocked) { |
| 83 | tracker.trackField("readerLocked", readerLocked); |
| 84 | } |
| 85 | } |
| 86 | } |
| 87 | |
| 88 | private: |
| 89 | class PipeLocked final: public PipeController { |
| 90 | public: |
| 91 | static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; |
| 92 | explicit PipeLocked(Controller& inner): inner(inner) {} |
| 93 | |
| 94 | bool isClosed() override { |
| 95 | return inner.state.template is<StreamStates::Closed>(); |
| 96 | } |
| 97 | |
| 98 | kj::Maybe<v8::Local<v8::Value>> tryGetErrored(jsg::Lock& js) override { |
| 99 | KJ_IF_SOME(errored, inner.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 100 | return errored.getHandle(js); |
| 101 | } |
| 102 | return kj::none; |
| 103 | } |
| 104 | |
| 105 | void cancel(jsg::Lock& js, v8::Local<v8::Value> reason) override { |
| 106 | // Cancel here returns a Promise but we do not need to propagate it. |
| 107 | // We can safely drop it on the floor here. |
| 108 | auto promise KJ_UNUSED = inner.cancel(js, reason); |
| 109 | } |
| 110 | |
| 111 | void close(jsg::Lock& js) override { |
| 112 | inner.doClose(js); |
| 113 | } |
| 114 | |
| 115 | void error(jsg::Lock& js, v8::Local<v8::Value> reason) override { |
| 116 | inner.doError(js, reason); |
| 117 | } |
| 118 | |
| 119 | void release(jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError = kj::none) override { |
| 120 | KJ_IF_SOME(error, maybeError) { |
| 121 | cancel(js, error); |
| 122 | } |
| 123 | inner.lock.state.template transitionTo<Unlocked>(); |
| 124 | } |
| 125 | |
| 126 | kj::Maybe<kj::Promise<void>> tryPumpTo(WritableStreamSink& sink, bool end) override; |
| 127 | |
| 128 | jsg::Promise<ReadResult> read(jsg::Lock& js) override; |
| 129 | |
| 130 | private: |
| 131 | Controller& inner; |
| 132 | |
| 133 | friend Controller; |
| 134 | }; |
| 135 | |
| 136 | // State machine for ReadableLockImpl: |
| 137 | // All states can transition to any other state (no terminal states). |
| 138 | // Unlocked -> Locked (lock() called for tee) |
| 139 | // Unlocked -> ReaderLocked (lockReader() called) |
| 140 | // Unlocked -> PipeLocked (tryPipeLock() called) |
| 141 | // ReaderLocked -> Unlocked (releaseReader() called) |
| 142 | // PipeLocked -> Unlocked (release() or onClose/onError called) |
| 143 | // Locked -> (remains until stream is done) |
| 144 | using LockState = StateMachine<Locked, PipeLocked, ReaderLocked, Unlocked>; |
| 145 | LockState state = LockState::template create<Unlocked>(); |
| 146 | friend Controller; |
| 147 | }; |
| 148 | |
| 149 | // A utility class used by WritableStreamJsController to implement the writer lock |
| 150 | // mechanism. Extracted for consistency with ReadableStreamJsController and to |
| 151 | // eventually allow it to be shared also with WritableStreamInternalController. |
| 152 | template <typename Controller> |
| 153 | class WritableLockImpl { |
| 154 | public: |
| 155 | using Writer = WritableStreamController::Writer; |
| 156 | |
| 157 | bool isLockedToWriter() const; |
| 158 | |
| 159 | bool lockWriter(jsg::Lock& js, Controller& self, Writer& writer); |
| 160 | |
| 161 | // See the comment for releaseWriter in common.h for details on the use of maybeJs |
| 162 | void releaseWriter(Controller& self, Writer& writer, kj::Maybe<jsg::Lock&> maybeJs); |
| 163 | |
| 164 | void visitForGc(jsg::GcVisitor& visitor); |
| 165 | |
| 166 | bool pipeLock(WritableStream& owner, jsg::Ref<ReadableStream> source, PipeToOptions& options); |
| 167 | void releasePipeLock(); |
| 168 | |
| 169 | JSG_MEMORY_INFO(WritableLockImpl) { |
| 170 | KJ_SWITCH_ONEOF(state) { |
| 171 | KJ_CASE_ONEOF(unlocked, Unlocked) {} |
| 172 | KJ_CASE_ONEOF(locked, Locked) {} |
| 173 | KJ_CASE_ONEOF(writerLocked, WriterLocked) { |
| 174 | tracker.trackField("writerLocked", writerLocked); |
| 175 | } |
| 176 | KJ_CASE_ONEOF(pipeLocked, PipeLocked) { |
| 177 | tracker.trackField("pipeLocked", pipeLocked); |
| 178 | } |
| 179 | } |
| 180 | } |
| 181 | |
| 182 | private: |
| 183 | struct PipeLocked { |
| 184 | static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; |
| 185 | ReadableStreamController::PipeController& source; |
| 186 | jsg::Ref<ReadableStream> readableStreamRef; |
| 187 | |
| 188 | kj::Maybe<jsg::Ref<AbortSignal>> maybeSignal; |
| 189 | |
| 190 | kj::Maybe<jsg::Promise<void>> checkSignal(jsg::Lock& js, Controller& self); |
| 191 | |
| 192 | struct Flags { |
| 193 | uint8_t preventAbort : 1 = 0; |
| 194 | uint8_t preventCancel : 1 = 0; |
| 195 | uint8_t preventClose : 1 = 0; |
| 196 | uint8_t pipeThrough : 1 = 0; |
| 197 | }; |
| 198 | Flags flags{}; |
| 199 | |
| 200 | JSG_MEMORY_INFO(PipeLocked) { |
| 201 | tracker.trackField("readableStreamRef", readableStreamRef); |
| 202 | tracker.trackField("signal", maybeSignal); |
| 203 | } |
| 204 | }; |
| 205 | |
| 206 | // State machine for WritableLockImpl: |
| 207 | // All states can transition to any other state (no terminal states). |
| 208 | // Unlocked -> Locked (not currently used) |
| 209 | // Unlocked -> WriterLocked (lockWriter() called) |
| 210 | // Unlocked -> PipeLocked (pipeLock() called) |
| 211 | // WriterLocked -> Unlocked (releaseWriter() called) |
| 212 | // PipeLocked -> Unlocked (releasePipeLock() called) |
| 213 | using LockState = StateMachine<Unlocked, Locked, WriterLocked, PipeLocked>; |
| 214 | LockState state = LockState::template create<Unlocked>(); |
| 215 | |
| 216 | inline kj::Maybe<PipeLocked&> tryGetPipe() { |
| 217 | KJ_IF_SOME(locked, state.template tryGetUnsafe<PipeLocked>()) { |
| 218 | return locked; |
| 219 | } |
| 220 | return kj::none; |
| 221 | } |
| 222 | |
| 223 | friend Controller; |
| 224 | }; |
| 225 | |
| 226 | // ====================================================================================== |
| 227 | |
| 228 | template <typename Controller> |
| 229 | bool ReadableLockImpl<Controller>::lock() { |
| 230 | if (isLockedToReader()) { |
| 231 | return false; |
| 232 | } |
| 233 | |
| 234 | state.template transitionTo<Locked>(); |
| 235 | return true; |
| 236 | } |
| 237 | |
| 238 | template <typename Controller> |
| 239 | bool ReadableLockImpl<Controller>::lockReader(jsg::Lock& js, Controller& self, Reader& reader) { |
| 240 | if (isLockedToReader()) { |
| 241 | return false; |
| 242 | } |
| 243 | |
| 244 | auto prp = js.newPromiseAndResolver<void>(); |
| 245 | prp.promise.markAsHandled(js); |
| 246 | |
| 247 | auto lock = ReaderLocked(reader, kj::mv(prp.resolver)); |
| 248 | |
| 249 | if (self.state.template is<StreamStates::Closed>()) { |
| 250 | maybeResolvePromise(js, lock.getClosedFulfiller()); |
| 251 | } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 252 | maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js)); |
| 253 | } |
| 254 | |
| 255 | state.template transitionTo<ReaderLocked>(kj::mv(lock)); |
| 256 | reader.attach(self, kj::mv(prp.promise)); |
| 257 | return true; |
| 258 | } |
| 259 | |
| 260 | template <typename Controller> |
| 261 | void ReadableLockImpl<Controller>::releaseReader( |
| 262 | Controller& self, Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) { |
| 263 | KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) { |
| 264 | KJ_ASSERT(&locked.getReader() == &reader); |
| 265 | |
| 266 | KJ_IF_SOME(js, maybeJs) { |
| 267 | auto reason = js.typeError("This ReadableStream reader has been released."_kj); |
| 268 | KJ_SWITCH_ONEOF(self.state) { |
| 269 | KJ_CASE_ONEOF(initial, typename Controller::Initial) {} |
| 270 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 271 | KJ_CASE_ONEOF(errored, StreamStates::Errored) {} |
| 272 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 273 | consumer->cancelPendingReads(js, reason); |
| 274 | } |
| 275 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 276 | consumer->cancelPendingReads(js, reason); |
| 277 | } |
| 278 | } |
| 279 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason); |
| 280 | } |
| 281 | |
| 282 | // Keep the locked.clear() after the isolate and hasPendingReadRequests check above. |
| 283 | // Clearing will release the references and we don't want to do that if the |
| 284 | // hasPendingReadRequests check fails. |
| 285 | locked.clear(); |
| 286 | |
| 287 | // When maybeJs is nullptr, that means releaseReader was called when the reader is |
| 288 | // being deconstructed and not as the result of explicitly calling releaseLock and |
| 289 | // we do not have an isolate lock. In that case, we don't want to change the lock |
| 290 | // state itself. Moving the lock above will free the lock state while keeping the |
| 291 | // ReadableStream marked as locked. |
| 292 | if (maybeJs != kj::none) { |
| 293 | state.template transitionTo<Unlocked>(); |
| 294 | } |
| 295 | } |
| 296 | } |
| 297 | |
| 298 | template <typename Controller> |
| 299 | kj::Maybe<ReadableStreamController::PipeController&> ReadableLockImpl<Controller>::tryPipeLock( |
| 300 | Controller& self) { |
| 301 | if (isLockedToReader()) { |
| 302 | return kj::none; |
| 303 | } |
| 304 | return state.template transitionTo<PipeLocked>(self); |
| 305 | } |
| 306 | |
| 307 | template <typename Controller> |
| 308 | void ReadableLockImpl<Controller>::visitForGc(jsg::GcVisitor& visitor) { |
| 309 | KJ_SWITCH_ONEOF(state) { |
| 310 | KJ_CASE_ONEOF(locked, Locked) {} |
| 311 | KJ_CASE_ONEOF(locked, Unlocked) {} |
| 312 | KJ_CASE_ONEOF(locked, PipeLocked) {} |
| 313 | KJ_CASE_ONEOF(locked, ReaderLocked) { |
| 314 | visitor.visit(locked); |
| 315 | } |
| 316 | } |
| 317 | } |
| 318 | |
| 319 | template <typename Controller> |
| 320 | void ReadableLockImpl<Controller>::onClose(jsg::Lock& js) { |
| 321 | KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) { |
| 322 | try { |
| 323 | maybeResolvePromise(js, locked.getClosedFulfiller()); |
| 324 | } catch (jsg::JsExceptionThrown&) { |
| 325 | // Resolving the promise could end up throwing an exception in some cases, |
| 326 | // causing a jsg::JsExceptionThrown to be thrown. At this point, however, |
| 327 | // we are already in the process of closing the stream and an error at this |
| 328 | // point is not recoverable. Log and move on. |
| 329 | LOG_NOSENTRY(ERROR, "Error resolving ReadableStream reader closed promise"); |
| 330 | }; |
| 331 | } else { |
| 332 | (void)state.template transitionFromTo<PipeLocked, Unlocked>(); |
| 333 | } |
| 334 | } |
| 335 | |
| 336 | template <typename Controller> |
| 337 | void ReadableLockImpl<Controller>::onError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 338 | KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) { |
| 339 | try { |
| 340 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason); |
| 341 | } catch (jsg::JsExceptionThrown&) { |
| 342 | // Rejecting the promise could end up throwing an exception in some cases, |
| 343 | // causing a jsg::JsExceptionThrown to be thrown. At this point, however, |
| 344 | // we are already in the process of closing the stream and an error at this |
| 345 | // point is not recoverable. Log and move on. |
| 346 | LOG_NOSENTRY(ERROR, "Error rejecting ReadableStream reader closed promise"); |
| 347 | } |
| 348 | } else { |
| 349 | (void)state.template transitionFromTo<PipeLocked, Unlocked>(); |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | template <typename Controller> |
| 354 | kj::Maybe<kj::Promise<void>> ReadableLockImpl<Controller>::PipeLocked::tryPumpTo( |
| 355 | WritableStreamSink& sink, bool end) { |
| 356 | // We return nullptr here because this controller does not support kj's pumpTo. |
| 357 | return kj::none; |
| 358 | } |
| 359 | |
| 360 | template <typename Controller> |
| 361 | jsg::Promise<ReadResult> ReadableLockImpl<Controller>::PipeLocked::read(jsg::Lock& js) { |
| 362 | return KJ_ASSERT_NONNULL(inner.read(js, kj::none)); |
| 363 | } |
| 364 | |
| 365 | // ====================================================================================== |
| 366 | |
| 367 | template <typename Controller> |
| 368 | bool WritableLockImpl<Controller>::isLockedToWriter() const { |
| 369 | return !state.template is<Unlocked>(); |
| 370 | } |
| 371 | |
| 372 | template <typename Controller> |
| 373 | bool WritableLockImpl<Controller>::lockWriter(jsg::Lock& js, Controller& self, Writer& writer) { |
| 374 | if (isLockedToWriter()) { |
| 375 | return false; |
| 376 | } |
| 377 | |
| 378 | auto closedPrp = js.newPromiseAndResolver<void>(); |
| 379 | closedPrp.promise.markAsHandled(js); |
| 380 | auto readyPrp = js.newPromiseAndResolver<void>(); |
| 381 | readyPrp.promise.markAsHandled(js); |
| 382 | |
| 383 | auto lock = WriterLocked(writer, kj::mv(closedPrp.resolver), kj::mv(readyPrp.resolver)); |
| 384 | |
| 385 | if (self.state.template is<StreamStates::Closed>()) { |
| 386 | maybeResolvePromise(js, lock.getClosedFulfiller()); |
| 387 | maybeResolvePromise(js, lock.getReadyFulfiller()); |
| 388 | } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 389 | maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js)); |
| 390 | maybeRejectPromise<void>(js, lock.getReadyFulfiller(), errored.getHandle(js)); |
| 391 | } else { |
| 392 | if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { |
| 393 | // Per spec (SetUpWritableStreamDefaultWriter step 4), the ready promise |
| 394 | // is resolved when the stream is writable and not experiencing backpressure, |
| 395 | // regardless of whether the start algorithm has completed. The backpressure |
| 396 | // state is set synchronously during SetUpWritableStreamDefaultController. |
| 397 | KJ_IF_SOME(erroring, self.isErroring(js)) { |
| 398 | maybeRejectPromise<void>(js, lock.getReadyFulfiller(), erroring); |
| 399 | } else if (!self.hasBackpressure()) { |
| 400 | maybeResolvePromise(js, lock.getReadyFulfiller()); |
| 401 | } |
| 402 | } else { |
| 403 | if (self.isStarted()) { |
| 404 | maybeResolvePromise(js, lock.getReadyFulfiller()); |
| 405 | } |
| 406 | } |
| 407 | } |
| 408 | |
| 409 | state.template transitionTo<WriterLocked>(kj::mv(lock)); |
| 410 | writer.attach(js, self, kj::mv(closedPrp.promise), kj::mv(readyPrp.promise)); |
| 411 | return true; |
| 412 | } |
| 413 | |
| 414 | template <typename Controller> |
| 415 | void WritableLockImpl<Controller>::releaseWriter( |
| 416 | Controller& self, Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) { |
| 417 | KJ_IF_SOME(locked, state.template tryGetUnsafe<WriterLocked>()) { |
| 418 | KJ_ASSERT(&locked.getWriter() == &writer); |
| 419 | KJ_IF_SOME(js, maybeJs) { |
| 420 | KJ_SWITCH_ONEOF(self.state) { |
| 421 | KJ_CASE_ONEOF(initial, typename Controller::Initial) {} |
| 422 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 423 | KJ_CASE_ONEOF(errored, StreamStates::Errored) {} |
| 424 | KJ_CASE_ONEOF(controller, jsg::Ref<WritableStreamDefaultController>) { |
| 425 | controller->cancelPendingWrites( |
| 426 | js, js.typeError("This WritableStream writer has been released."_kjc)); |
| 427 | } |
| 428 | } |
| 429 | |
| 430 | // Per spec (WritableStreamDefaultWriterRelease), both the ready and closed |
| 431 | // promises must be rejected when the writer is released. |
| 432 | auto releaseReason = js.v8TypeError("This WritableStream writer has been released."_kjc); |
| 433 | if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { |
| 434 | if (locked.getReadyFulfiller() != kj::none) { |
| 435 | maybeRejectPromise<void>(js, locked.getReadyFulfiller(), releaseReason); |
| 436 | } else { |
| 437 | // The ready fulfiller was already consumed (promise was resolved). |
| 438 | // Per spec (WritableStreamDefaultWriterEnsureReadyPromiseRejected), |
| 439 | // we must replace it with a new rejected promise. |
| 440 | auto prp = js.newPromiseAndResolver<void>(); |
| 441 | prp.promise.markAsHandled(js); |
| 442 | prp.resolver.reject(js, releaseReason); |
| 443 | locked.setReadyFulfiller(js, prp); |
| 444 | } |
| 445 | } else { |
| 446 | maybeRejectPromise<void>(js, locked.getReadyFulfiller(), releaseReason); |
| 447 | } |
| 448 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), releaseReason); |
| 449 | } |
| 450 | locked.clear(); |
| 451 | |
| 452 | // When maybeJs is nullptr, that means releaseWriter was called when the writer is |
| 453 | // being deconstructed and not as the result of explicitly calling releaseLock and |
| 454 | // we do not have an isolate lock. In that case, we don't want to change the lock |
| 455 | // state itself. Moving the lock above will free the lock state while keeping the |
| 456 | // WritableStream marked as locked. |
| 457 | if (maybeJs != kj::none) { |
| 458 | state.template transitionTo<Unlocked>(); |
| 459 | } |
| 460 | } |
| 461 | } |
| 462 | |
| 463 | template <typename Controller> |
| 464 | bool WritableLockImpl<Controller>::pipeLock( |
| 465 | WritableStream& owner, jsg::Ref<ReadableStream> source, PipeToOptions& options) { |
| 466 | if (isLockedToWriter()) { |
| 467 | return false; |
| 468 | } |
| 469 | |
| 470 | auto& sourceLock = KJ_ASSERT_NONNULL(source->getController().tryPipeLock()); |
| 471 | |
| 472 | state.template transitionTo<PipeLocked>(PipeLocked{ |
| 473 | .source = sourceLock, |
| 474 | .readableStreamRef = kj::mv(source), |
| 475 | .maybeSignal = kj::mv(options.signal), |
| 476 | .flags = |
| 477 | { |
| 478 | .preventAbort = options.preventAbort.orDefault(false), |
| 479 | .preventCancel = options.preventCancel.orDefault(false), |
| 480 | .preventClose = options.preventClose.orDefault(false), |
| 481 | .pipeThrough = options.pipeThrough, |
| 482 | }, |
| 483 | }); |
| 484 | return true; |
| 485 | } |
| 486 | |
| 487 | template <typename Controller> |
| 488 | void WritableLockImpl<Controller>::releasePipeLock() { |
| 489 | if (state.template is<PipeLocked>()) { |
| 490 | state.template transitionTo<Unlocked>(); |
| 491 | } |
| 492 | } |
| 493 | |
| 494 | template <typename Controller> |
| 495 | void WritableLockImpl<Controller>::visitForGc(jsg::GcVisitor& visitor) { |
| 496 | KJ_SWITCH_ONEOF(state) { |
| 497 | KJ_CASE_ONEOF(locked, Unlocked) {} |
| 498 | KJ_CASE_ONEOF(locked, Locked) {} |
| 499 | KJ_CASE_ONEOF(locked, WriterLocked) { |
| 500 | visitor.visit(locked); |
| 501 | } |
| 502 | KJ_CASE_ONEOF(locked, PipeLocked) { |
| 503 | visitor.visit(locked.readableStreamRef); |
| 504 | KJ_IF_SOME(signal, locked.maybeSignal) { |
| 505 | visitor.visit(signal); |
| 506 | } |
| 507 | } |
| 508 | } |
| 509 | } |
| 510 | |
| 511 | template <typename Controller> |
| 512 | kj::Maybe<jsg::Promise<void>> WritableLockImpl<Controller>::PipeLocked::checkSignal( |
| 513 | jsg::Lock& js, Controller& self) { |
| 514 | KJ_IF_SOME(signal, maybeSignal) { |
| 515 | if (signal->getAborted(js)) { |
| 516 | auto reason = signal->getReason(js); |
| 517 | if (!flags.preventCancel) { |
| 518 | source.release(js, v8::Local<v8::Value>(reason)); |
| 519 | } else { |
| 520 | source.release(js); |
| 521 | } |
| 522 | if (!flags.preventAbort) { |
| 523 | return self.abort(js, reason).then(js, JSG_VISITABLE_LAMBDA((this, reason = reason.addRef(js), ref = self.addRef()), (reason, ref), (jsg::Lock& js) { |
| 524 | return rejectedMaybeHandledPromise<void>(js, reason.getHandle(js), flags.pipeThrough); |
| 525 | })); |
| 526 | } |
| 527 | return rejectedMaybeHandledPromise<void>(js, reason, flags.pipeThrough); |
| 528 | } |
| 529 | } |
| 530 | return kj::none; |
| 531 | } |
| 532 | |
| 533 | auto maybeAddFunctor(jsg::Lock& js, auto promise, auto onSuccess, auto onFailure) { |
| 534 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 535 | return promise.then( |
| 536 | js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure))); |
| 537 | } else { |
| 538 | return promise.then(js, kj::mv(onSuccess), kj::mv(onFailure)); |
| 539 | } |
| 540 | } |
| 541 | |
| 542 | jsg::Promise<void> maybeRunAlgorithm( |
| 543 | jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) { |
| 544 | // The algorithm is a JavaScript function mapped through jsg::Function. |
| 545 | // It is expected to return a Promise mapped via jsg::Promise. If the |
| 546 | // function returns synchronously, the jsg::Promise wrapper ensures |
| 547 | // that it is properly mapped to a jsg::Promise, but if the Promise |
| 548 | // throws synchronously, we have to convert that synchronous throw |
| 549 | // into a proper rejected jsg::Promise. |
| 550 | KJ_IF_SOME(algorithm, maybeAlgorithm) { |
| 551 | // We need two layers of JSG_TRY here, unfortunately. The inner layer |
| 552 | // covers the algorithm implementation itself and is our typical error |
| 553 | // handling path. It ensures that if the algorithm throws an exception, |
| 554 | // that is properly converted in to a rejected promise that is *then* |
| 555 | // handled by the onFailure handler that is passed in. The outer JSG_TRY |
| 556 | // handles the rare and generally unexpected failure of the calls to |
| 557 | // .then() itself, which can throw JS exceptions synchronously in certain |
| 558 | // rare cases. For those we return a rejected promise but do not call the |
| 559 | // onFailure case since such errors are generally indicative of a fatal |
| 560 | // condition in the isolate (e.g. out of memory, other fatal exception, etc). |
| 561 | JSG_TRY(js) { |
| 562 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 563 | auto getInnerPromise = [&]() -> jsg::Promise<void> { |
| 564 | JSG_TRY(js) { |
| 565 | return algorithm(js, kj::fwd<decltype(args)>(args)...); |
| 566 | } |
| 567 | JSG_CATCH(exception) { |
| 568 | return js.rejectedPromise<void>(kj::mv(exception)); |
| 569 | } |
| 570 | }; |
| 571 | return getInnerPromise().then( |
| 572 | js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure))); |
| 573 | } else { |
| 574 | auto getInnerPromise = [&]() -> jsg::Promise<void> { |
| 575 | JSG_TRY(js) { |
| 576 | return algorithm(js, kj::fwd<decltype(args)>(args)...); |
| 577 | } |
| 578 | JSG_CATCH(exception) { |
| 579 | return js.rejectedPromise<void>(kj::mv(exception)); |
| 580 | } |
| 581 | }; |
| 582 | return getInnerPromise().then(js, kj::mv(onSuccess), kj::mv(onFailure)); |
| 583 | } |
| 584 | } |
| 585 | JSG_CATCH(exception) { |
| 586 | return js.rejectedPromise<void>(kj::mv(exception)); |
| 587 | } |
| 588 | } |
| 589 | |
| 590 | // If the algorithm does not exist, we just handle it as a success and move on. |
| 591 | onSuccess(js); |
| 592 | return js.resolvedPromise(); |
| 593 | } |
| 594 | |
| 595 | jsg::Promise<void> maybeRunAlgorithmAsync( |
| 596 | jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) { |
| 597 | // The algorithm is a JavaScript function mapped through jsg::Function. |
| 598 | // It is expected to return a Promise mapped via jsg::Promise. If the |
| 599 | // function returns synchronously, the jsg::Promise wrapper ensures |
| 600 | // that it is properly mapped to a jsg::Promise, but if the Promise |
| 601 | // throws synchronously, we have to convert that synchronous throw |
| 602 | // into a proper rejected jsg::Promise. |
| 603 | KJ_IF_SOME(algorithm, maybeAlgorithm) { |
| 604 | // We need two layers of tryCatch here, unfortunately. The inner layer |
| 605 | // covers the algorithm implementation itself and is our typical error |
| 606 | // handling path. It ensures that if the algorithm throws an exception, |
| 607 | // that is properly converted in to a rejected promise that is *then* |
| 608 | // handled by the onFailure handler that is passed in. The outer tryCatch |
| 609 | // handles the rare and generally unexpected failure of the calls to |
| 610 | // .then() itself, which can throw JS exceptions synchronously in certain |
| 611 | // rare cases. For those we return a rejected promise but do not call the |
| 612 | // onFailure case since such errors are generally indicative of a fatal |
| 613 | // condition in the isolate (e.g. out of memory, other fatal exception, etc). |
| 614 | return js.tryCatch([&] { |
| 615 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 616 | return js |
| 617 | .tryCatch([&] { return algorithm(js, kj::fwd<decltype(args)>(args)...); }, |
| 618 | [&](jsg::Value&& exception) { return js.rejectedPromise<void>(kj::mv(exception)); }) |
| 619 | .then(js, ioContext.addFunctor(kj::mv(onSuccess)), |
| 620 | ioContext.addFunctor(kj::mv(onFailure))); |
| 621 | } else { |
| 622 | return js |
| 623 | .tryCatch([&] { return algorithm(js, kj::fwd<decltype(args)>(args)...); }, |
| 624 | [&](jsg::Value&& exception) { |
| 625 | return js.rejectedPromise<void>(kj::mv(exception)); |
| 626 | }).then(js, kj::mv(onSuccess), kj::mv(onFailure)); |
| 627 | } |
| 628 | }, [&](jsg::Value&& exception) { return js.rejectedPromise<void>(kj::mv(exception)); }); |
| 629 | } |
| 630 | |
| 631 | // If the algorithm does not exist, we handle it as a success but ensure |
| 632 | // it runs asynchronously by scheduling via a resolved promise. |
| 633 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 634 | return js.resolvedPromise().then(js, ioContext.addFunctor(kj::mv(onSuccess))); |
| 635 | } else { |
| 636 | return js.resolvedPromise().then(js, kj::mv(onSuccess)); |
| 637 | } |
| 638 | } |
| 639 | |
| 640 | int getHighWaterMark( |
| 641 | const UnderlyingSource& underlyingSource, const StreamQueuingStrategy& queuingStrategy) { |
| 642 | bool isBytes = underlyingSource.type.map([](auto& s) { return s == "bytes"; }).orDefault(false); |
| 643 | return queuingStrategy.highWaterMark.orDefault(isBytes ? 0 : 1); |
| 644 | } |
| 645 | |
| 646 | } // namespace |
| 647 | |
| 648 | // It is possible for the controller state to be released synchronously while |
| 649 | // we are in the middle of a read. When that happens we need to defer the actual |
| 650 | // close/error state change until the read call is complete. deferControllerStateChange |
| 651 | // handles this for us by using the state machine's operation tracking to defer |
| 652 | // pending close/error transitions until the read is complete. |
| 653 | template <typename Controller> |
| 654 | jsg::Promise<ReadResult> deferControllerStateChange(jsg::Lock& js, |
| 655 | Controller& controller, |
| 656 | kj::FunctionParam<jsg::Promise<ReadResult>()> readCallback) { |
| 657 | bool endOperation = true; |
| 658 | // The readCallback and the controller.doClose(..) and controller.doError(...) |
| 659 | // methods, as well as the methods can trigger JavaScript errors to be thrown |
| 660 | // synchronously in some cases. We want to make sure non-fatal errors cause the |
| 661 | // stream to error and only fatal cases bubble up. |
| 662 | return js.tryCatch([&] { |
| 663 | controller.state.beginOperation(); |
| 664 | auto result = readCallback(); |
| 665 | endOperation = false; |
| 666 | |
| 667 | // endOperation() will automatically apply any pending state if this was the last operation. |
| 668 | // Returns true if a pending state was applied. |
| 669 | if (controller.state.endOperation()) { |
| 670 | // A pending state was applied. Call the appropriate callback. |
| 671 | // Skip callbacks if execution is being terminated (e.g., CPU time limit) since we can't |
| 672 | // safely execute JavaScript in that state. |
| 673 | if (!js.v8Isolate->IsExecutionTerminating()) { |
| 674 | if (controller.state.template is<StreamStates::Closed>()) { |
| 675 | controller.lock.onClose(js); |
| 676 | } else if (controller.state.template is<StreamStates::Errored>()) { |
| 677 | KJ_IF_SOME(err, controller.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 678 | controller.lock.onError(js, err.getHandle(js)); |
| 679 | } |
| 680 | } |
| 681 | } |
| 682 | } |
| 683 | |
| 684 | return kj::mv(result); |
| 685 | }, [&](jsg::Value exception) -> jsg::Promise<ReadResult> { |
| 686 | if (endOperation) { |
| 687 | // Clear any pending state since we're erroring |
| 688 | controller.state.clearPendingState(); |
| 689 | (void)controller.state.endOperation(); |
| 690 | } |
| 691 | controller.doError(js, exception.getHandle(js)); |
| 692 | return js.rejectedPromise<ReadResult>(kj::mv(exception)); |
| 693 | }); |
| 694 | } |
| 695 | |
| 696 | // The ReadableStreamJsController provides the implementation of custom |
| 697 | // ReadableStreams backed by a user-code provided Underlying Source. The implementation |
| 698 | // is fairly complicated and defined entirely by the streams specification. |
| 699 | // |
| 700 | // Another important thing to understand is that there are two types of JavaScript |
| 701 | // backed ReadableStreams: value-oriented, and byte-oriented. |
| 702 | // |
| 703 | // When user code uses the `new ReadableStream(underlyingSource)` constructor, the |
| 704 | // underlyingSource argument may have a `type` property, the value of which is either |
| 705 | // `undefined`, the empty string, or the string value `'bytes'`. If the underlyingSource |
| 706 | // argument is not given, the default value of `type` is `undefined`. If `type` is |
| 707 | // `undefined` or the empty string, the ReadableStream is value-oriented. If `type` is |
| 708 | // exactly equal to `'bytes'`, the ReadableStream is byte-oriented. |
| 709 | // |
| 710 | // For value-oriented streams, any JavaScript value can be pushed through the stream, |
| 711 | // and the stream will only support use of the ReadableStreamDefaultReader to consume |
| 712 | // the stream data. |
| 713 | // |
| 714 | // For byte-oriented streams, only byte data (as provided by `ArrayBufferView`s) can |
| 715 | // be pushed through the stream. All byte-oriented streams support using both |
| 716 | // ReadableStreamDefaultReader and ReadableStreamBYOBReader to consume the stream |
| 717 | // data. |
| 718 | // |
| 719 | // When the ReadableStreamJsController::setup() method is called the type |
| 720 | // of stream is determined, and the controller will create an instance of either |
| 721 | // jsg::Ref<ReadableStreamDefaultController> or jsg::Ref<ReadableByteStreamController>. |
| 722 | // These are the objects that are actually passed on to the user-code's Underlying Source |
| 723 | // implementation. |
| 724 | class ReadableStreamJsController final: public ReadableStreamController { |
| 725 | public: |
| 726 | using ReadableLockImpl = ReadableLockImpl<ReadableStreamJsController>; |
| 727 | |
| 728 | KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamJsController); |
| 729 | |
| 730 | explicit ReadableStreamJsController(); |
| 731 | explicit ReadableStreamJsController(StreamStates::Closed closed); |
| 732 | explicit ReadableStreamJsController(StreamStates::Errored errored); |
| 733 | explicit ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer); |
| 734 | explicit ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer); |
| 735 | |
| 736 | jsg::Ref<ReadableStream> addRef() override; |
| 737 | |
| 738 | void setup(jsg::Lock& js, |
| 739 | jsg::Optional<UnderlyingSource> maybeUnderlyingSource, |
| 740 | jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) override; |
| 741 | |
| 742 | // Signals that this ReadableStream is no longer interested in the underlying |
| 743 | // data source. Whether this cancels the underlying data source also depends |
| 744 | // on whether or not there are other ReadableStreams still attached to it. |
| 745 | // This operation is terminal. Once called, even while the returned Promise |
| 746 | // is still pending, the ReadableStream will be no longer usable and any |
| 747 | // data still in the queue will be dropped. Pending read requests will be |
| 748 | // rejected if a reason is given, or resolved with no data otherwise. |
| 749 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override; |
| 750 | |
| 751 | void doClose(jsg::Lock& js); |
| 752 | |
| 753 | void doError(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 754 | |
| 755 | bool canCloseOrEnqueue(); |
| 756 | bool hasBackpressure(); |
| 757 | |
| 758 | bool isByteOriented() const override; |
| 759 | |
| 760 | bool isDisturbed() override; |
| 761 | |
| 762 | bool isClosedOrErrored() const override; |
| 763 | |
| 764 | bool isClosed() const override; |
| 765 | |
| 766 | bool isLockedToReader() const override; |
| 767 | |
| 768 | bool lockReader(jsg::Lock& js, Reader& reader) override; |
| 769 | |
| 770 | kj::Maybe<v8::Local<v8::Value>> isErrored(jsg::Lock& js); |
| 771 | |
| 772 | kj::Maybe<int> getDesiredSize(); |
| 773 | |
| 774 | jsg::Promise<void> pipeTo( |
| 775 | jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) override; |
| 776 | |
| 777 | kj::Promise<DeferredProxy<void>> pumpTo( |
| 778 | jsg::Lock& js, kj::Own<WritableStreamSink>, bool end) override; |
| 779 | |
| 780 | kj::Maybe<jsg::Promise<ReadResult>> read( |
| 781 | jsg::Lock& js, kj::Maybe<ByobOptions> byobOptions) override; |
| 782 | |
| 783 | kj::Maybe<jsg::Promise<DrainingReadResult>> drainingRead( |
| 784 | jsg::Lock& js, size_t maxRead = kj::maxValue) override; |
| 785 | |
| 786 | // See the comment for releaseReader in common.h for details on the use of maybeJs |
| 787 | void releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) override; |
| 788 | |
| 789 | void setOwnerRef(ReadableStream& stream) override; |
| 790 | |
| 791 | Tee tee(jsg::Lock& js) override; |
| 792 | |
| 793 | kj::Maybe<PipeController&> tryPipeLock() override; |
| 794 | |
| 795 | void visitForGc(jsg::GcVisitor& visitor) override; |
| 796 | |
| 797 | kj::Maybe<kj::OneOf<DefaultController, ByobController>> getController(); |
| 798 | |
| 799 | jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit) override; |
| 800 | jsg::Promise<kj::String> readAllText(jsg::Lock& js, uint64_t limit) override; |
| 801 | |
| 802 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override; |
| 803 | |
| 804 | kj::Own<ReadableStreamController> detach(jsg::Lock& js, bool ignoreDisturbed) override; |
| 805 | |
| 806 | void setPendingClosure() override { |
| 807 | KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController"); |
| 808 | } |
| 809 | |
| 810 | kj::StringPtr jsgGetMemoryName() const override; |
| 811 | size_t jsgGetMemorySelfSize() const override; |
| 812 | void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const override; |
| 813 | |
| 814 | private: |
| 815 | // If the stream was created within the scope of a request, we want to treat it as I/O |
| 816 | // and make sure it is not advanced from the scope of a different request. |
| 817 | kj::Maybe<IoContext&> ioContext; |
| 818 | kj::Maybe<ReadableStream&> owner; |
| 819 | |
| 820 | // Initial state before setup() is called. |
| 821 | struct Initial { |
| 822 | static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; |
| 823 | }; |
| 824 | |
| 825 | // State machine for ReadableStreamJsController: |
| 826 | // Initial is the default state before setup() is called |
| 827 | // ValueReadable and ByteReadable are the active states (stream has data) |
| 828 | // Closed and Errored are terminal states (stream is done) |
| 829 | // Initial -> ValueReadable or ByteReadable (setup() called) |
| 830 | // Initial -> Closed (constructed with Closed) |
| 831 | // Initial -> Errored (constructed with Errored) |
| 832 | // ValueReadable -> Closed (doClose() or cancel() called) |
| 833 | // ValueReadable -> Errored (doError() called) |
| 834 | // ByteReadable -> Closed (doClose() or cancel() called) |
| 835 | // ByteReadable -> Errored (doError() called) |
| 836 | // Note: No single ActiveState since there are two active variants. |
| 837 | // PendingStates allows Closed/Errored transitions to be deferred during reads. |
| 838 | using State = StateMachine<TerminalStates<StreamStates::Closed>, |
| 839 | ErrorState<StreamStates::Errored>, |
| 840 | PendingStates<StreamStates::Closed, StreamStates::Errored>, |
| 841 | Initial, |
| 842 | StreamStates::Closed, |
| 843 | StreamStates::Errored, |
| 844 | kj::Own<ValueReadable>, |
| 845 | kj::Own<ByteReadable>>; |
| 846 | State state = State::create<Initial>(); |
| 847 | |
| 848 | kj::Maybe<uint64_t> expectedLength = kj::none; |
| 849 | bool canceling = false; |
| 850 | |
| 851 | // The lock state is separate because a closed or errored stream can still be locked. |
| 852 | ReadableLockImpl lock; |
| 853 | |
| 854 | bool disturbed = false; |
| 855 | |
| 856 | template <typename T> |
| 857 | jsg::Promise<T> readAll(jsg::Lock& js, uint64_t limit); |
| 858 | |
| 859 | friend ReadableLockImpl; |
| 860 | friend ReadableLockImpl::PipeLocked; |
| 861 | friend struct ValueReadable; |
| 862 | friend struct ByteReadable; |
| 863 | |
| 864 | template <typename Controller> |
| 865 | friend jsg::Promise<ReadResult> deferControllerStateChange(jsg::Lock& js, |
| 866 | Controller& controller, |
| 867 | kj::FunctionParam<jsg::Promise<ReadResult>()> readCallback); |
| 868 | }; |
| 869 | |
| 870 | // The WritableStreamJsController provides the implementation of custom |
| 871 | // WritableStream's backed by a user-code provided Underlying Sink. The implementation |
| 872 | // is fairly complicated and defined entirely by the streams specification. |
| 873 | class WritableStreamJsController final: public WritableStreamController { |
| 874 | public: |
| 875 | using WritableLockImpl = WritableLockImpl<WritableStreamJsController>; |
| 876 | |
| 877 | using Controller = jsg::Ref<WritableStreamDefaultController>; |
| 878 | |
| 879 | explicit WritableStreamJsController(); |
| 880 | |
| 881 | explicit WritableStreamJsController(StreamStates::Closed closed); |
| 882 | |
| 883 | explicit WritableStreamJsController(StreamStates::Errored errored); |
| 884 | |
| 885 | ~WritableStreamJsController() noexcept(false); |
| 886 | |
| 887 | KJ_DISALLOW_COPY_AND_MOVE(WritableStreamJsController); |
| 888 | |
| 889 | jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override; |
| 890 | |
| 891 | jsg::Ref<WritableStream> addRef() override; |
| 892 | |
| 893 | jsg::Promise<void> close(jsg::Lock& js, bool markAsHandled = false) override; |
| 894 | |
| 895 | jsg::Promise<void> flush(jsg::Lock& js, bool markAsHandled = false) override { |
| 896 | KJ_UNIMPLEMENTED("expected WritableStreamInternalController implementation to be enough"); |
| 897 | } |
| 898 | |
| 899 | void doClose(jsg::Lock& js); |
| 900 | |
| 901 | void doError(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 902 | |
| 903 | // Error through the underlying controller if available, going through the proper |
| 904 | // error transition (Erroring -> Errored). |
| 905 | void errorIfNeeded(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 906 | |
| 907 | kj::Maybe<int> getDesiredSize() override; |
| 908 | |
| 909 | kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js) override; |
| 910 | kj::Maybe<v8::Local<v8::Value>> isErroredOrErroring(jsg::Lock& js); |
| 911 | |
| 912 | bool isLocked() const; |
| 913 | |
| 914 | bool isLockedToWriter() const override; |
| 915 | |
| 916 | bool isStarted(); |
| 917 | |
| 918 | bool hasBackpressure(); |
| 919 | |
| 920 | inline bool isWritable() const { |
| 921 | return state.isActive(); |
| 922 | } |
| 923 | |
| 924 | bool lockWriter(jsg::Lock& js, Writer& writer) override; |
| 925 | |
| 926 | void maybeRejectReadyPromise(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 927 | |
| 928 | void maybeResolveReadyPromise(jsg::Lock& js); |
| 929 | |
| 930 | // See the comment for releaseWriter in common.h for details on the use of maybeJs |
| 931 | void releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) override; |
| 932 | |
| 933 | kj::Maybe<kj::Own<WritableStreamSink>> removeSink(jsg::Lock& js) override; |
| 934 | void detach(jsg::Lock& js) override; |
| 935 | |
| 936 | void setOwnerRef(WritableStream& stream) override; |
| 937 | |
| 938 | void setup(jsg::Lock& js, |
| 939 | jsg::Optional<UnderlyingSink> maybeUnderlyingSink, |
| 940 | jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) override; |
| 941 | |
| 942 | kj::Maybe<jsg::Promise<void>> tryPipeFrom( |
| 943 | jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) override; |
| 944 | |
| 945 | void updateBackpressure(jsg::Lock& js, bool backpressure); |
| 946 | |
| 947 | jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) override; |
| 948 | |
| 949 | void visitForGc(jsg::GcVisitor& visitor) override; |
| 950 | |
| 951 | bool isClosedOrClosing() override; |
| 952 | bool isErrored() override; |
| 953 | |
| 954 | inline bool isByteOriented() const override { |
| 955 | return false; |
| 956 | } |
| 957 | |
| 958 | void setPendingClosure() override { |
| 959 | KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController"); |
| 960 | } |
| 961 | |
| 962 | kj::StringPtr jsgGetMemoryName() const override; |
| 963 | size_t jsgGetMemorySelfSize() const override; |
| 964 | void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override; |
| 965 | |
| 966 | private: |
| 967 | jsg::Promise<void> pipeLoop(jsg::Lock& js); |
| 968 | |
| 969 | kj::Maybe<IoContext&> ioContext; |
| 970 | kj::Maybe<WritableStream&> owner; |
| 971 | |
| 972 | // Initial state before setup() is called. |
| 973 | struct Initial { |
| 974 | static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; |
| 975 | }; |
| 976 | |
| 977 | // State machine for WritableStreamJsController: |
| 978 | // Initial is the default state before setup() is called |
| 979 | // Controller is the active state (stream is writable) |
| 980 | // Closed is terminal, Errored is implicitly terminal via ErrorState |
| 981 | using State = StateMachine<TerminalStates<StreamStates::Closed>, |
| 982 | ErrorState<StreamStates::Errored>, |
| 983 | ActiveState<Controller>, |
| 984 | Initial, |
| 985 | StreamStates::Closed, |
| 986 | StreamStates::Errored, |
| 987 | Controller>; |
| 988 | State state = State::create<Initial>(); |
| 989 | |
| 990 | WritableLockImpl lock; |
| 991 | kj::Maybe<jsg::Promise<void>> maybeAbortPromise; |
| 992 | |
| 993 | friend WritableLockImpl; |
| 994 | }; |
| 995 | |
| 996 | kj::Own<ReadableStreamController> newReadableStreamJsController() { |
| 997 | return kj::heap<ReadableStreamJsController>(); |
| 998 | } |
| 999 | |
| 1000 | kj::Own<WritableStreamController> newWritableStreamJsController() { |
| 1001 | return kj::heap<WritableStreamJsController>(); |
| 1002 | } |
| 1003 | |
| 1004 | template <typename Self> |
| 1005 | ReadableImpl<Self>::ReadableImpl( |
| 1006 | UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) |
| 1007 | : state(State::template create<Queue>(getHighWaterMark(underlyingSource, queuingStrategy))), |
| 1008 | algorithms(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} |
| 1009 | |
| 1010 | template <typename Self> |
| 1011 | void ReadableImpl<Self>::start(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1012 | KJ_ASSERT(!flags.started && !flags.starting); |
| 1013 | flags.starting = true; |
| 1014 | |
| 1015 | // Per the streams spec, the size function should be called with `undefined` as `this`, |
| 1016 | // not as a method on the strategy object. |
| 1017 | KJ_IF_SOME(sizeFunc, algorithms.size) { |
| 1018 | sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined())); |
| 1019 | } |
| 1020 | |
| 1021 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1022 | flags.started = true; |
| 1023 | flags.starting = false; |
| 1024 | pullIfNeeded(js, kj::mv(self)); |
| 1025 | }); |
| 1026 | |
| 1027 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1028 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1029 | flags.started = true; |
| 1030 | flags.starting = false; |
| 1031 | doError(js, kj::mv(reason)); |
| 1032 | }); |
| 1033 | |
| 1034 | maybeRunAlgorithm(js, algorithms.start, kj::mv(onSuccess), kj::mv(onFailure), kj::mv(self)); |
| 1035 | algorithms.start = kj::none; |
| 1036 | } |
| 1037 | |
| 1038 | template <typename Self> |
| 1039 | size_t ReadableImpl<Self>::consumerCount() { |
| 1040 | return state.whenActiveOr([](Queue& q) { return q.getConsumerCount(); }, size_t{0}); |
| 1041 | } |
| 1042 | |
| 1043 | template <typename Self> |
| 1044 | jsg::Promise<void> ReadableImpl<Self>::cancel( |
| 1045 | jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1046 | if (state.template is<StreamStates::Closed>()) { |
| 1047 | // We are already closed. There's nothing to cancel. |
| 1048 | // This shouldn't happen but we handle the case anyway, just to be safe. |
| 1049 | return js.resolvedPromise(); |
| 1050 | } |
| 1051 | KJ_IF_SOME(errored, state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 1052 | // We are already errored. There's nothing to cancel. |
| 1053 | // This shouldn't happen but we handle the case anyway, just to be safe. |
| 1054 | return js.rejectedPromise<void>(errored.getHandle(js)); |
| 1055 | } |
| 1056 | |
| 1057 | auto& queue = state.template getUnsafe<Queue>(); |
| 1058 | size_t consumerCount = queue.getConsumerCount(); |
| 1059 | if (consumerCount > 1) { |
| 1060 | // If there is more than 1 consumer, then we just return here with an |
| 1061 | // immediately resolved promise. The consumer will remove itself, |
| 1062 | // canceling its interest in the underlying source but we do not yet |
| 1063 | // want to cancel the underlying source since there are still other |
| 1064 | // consumers that want data. |
| 1065 | return js.resolvedPromise(); |
| 1066 | } |
| 1067 | |
| 1068 | // Otherwise, there should be exactly one consumer at this point. |
| 1069 | KJ_ASSERT(consumerCount == 1); |
| 1070 | KJ_IF_SOME(pendingCancel, maybePendingCancel) { |
| 1071 | // If we're already waiting for cancel to complete, just return the |
| 1072 | // already existing pending promise. |
| 1073 | // This shouldn't happen but we handle the case anyway, just to be safe. |
| 1074 | return pendingCancel.promise.whenResolved(js); |
| 1075 | } |
| 1076 | |
| 1077 | auto prp = js.newPromiseAndResolver<void>(); |
| 1078 | maybePendingCancel = PendingCancel{ |
| 1079 | .fulfiller = kj::mv(prp.resolver), |
| 1080 | .promise = kj::mv(prp.promise), |
| 1081 | }; |
| 1082 | auto promise = KJ_ASSERT_NONNULL(maybePendingCancel).promise.whenResolved(js); |
| 1083 | doCancel(js, kj::mv(self), reason); |
| 1084 | return kj::mv(promise); |
| 1085 | } |
| 1086 | |
| 1087 | template <typename Self> |
| 1088 | bool ReadableImpl<Self>::canCloseOrEnqueue() { |
| 1089 | return state.isActive(); |
| 1090 | } |
| 1091 | |
| 1092 | // doCancel() is triggered by cancel() being called, which is an explicit signal from |
| 1093 | // the ReadableStream that we don't care about the data this controller provides any |
| 1094 | // more. We don't need to notify the consumers because we presume they already know |
| 1095 | // that they called cancel. What we do want to do here, tho, is close the implementation |
| 1096 | // and trigger the cancel algorithm. |
| 1097 | template <typename Self> |
| 1098 | void ReadableImpl<Self>::doCancel(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1099 | state.template transitionTo<StreamStates::Closed>(); |
| 1100 | |
| 1101 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1102 | doClose(js); |
| 1103 | KJ_IF_SOME(pendingCancel, maybePendingCancel) { |
| 1104 | maybeResolvePromise(js, pendingCancel.fulfiller); |
| 1105 | } else { |
| 1106 | // Else block to avert dangling else compiler warning. |
| 1107 | } |
| 1108 | }); |
| 1109 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1110 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1111 | // We do not call doError() here because there's really no point. Everything |
| 1112 | // that cares about the state of this controller impl has signaled that it |
| 1113 | // no longer cares and has gone away. |
| 1114 | doClose(js); |
| 1115 | KJ_IF_SOME(pendingCancel, maybePendingCancel) { |
| 1116 | maybeRejectPromise<void>(js, pendingCancel.fulfiller, reason.getHandle(js)); |
| 1117 | } else { |
| 1118 | // Else block to avert dangling else compiler warning. |
| 1119 | } |
| 1120 | }); |
| 1121 | |
| 1122 | maybeRunAlgorithm(js, algorithms.cancel, kj::mv(onSuccess), kj::mv(onFailure), reason); |
| 1123 | } |
| 1124 | |
| 1125 | template <typename Self> |
| 1126 | void ReadableImpl<Self>::enqueue(jsg::Lock& js, kj::Rc<Entry> entry, jsg::Ref<Self> self) { |
| 1127 | JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed."); |
| 1128 | KJ_DEFER(pullIfNeeded(js, kj::mv(self))); |
| 1129 | auto& queue = state.template getUnsafe<Queue>(); |
| 1130 | queue.push(js, kj::mv(entry)); |
| 1131 | } |
| 1132 | |
| 1133 | template <typename Self> |
| 1134 | void ReadableImpl<Self>::close(jsg::Lock& js) { |
| 1135 | JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed."); |
| 1136 | auto& queue = state.template getUnsafe<Queue>(); |
| 1137 | |
| 1138 | if (queue.hasPartiallyFulfilledRead()) { |
| 1139 | auto error = |
| 1140 | js.v8Ref(js.v8TypeError("This ReadableStream was closed with a partial read pending.")); |
| 1141 | doError(js, error.addRef(js)); |
| 1142 | js.throwException(kj::mv(error)); |
| 1143 | return; |
| 1144 | } |
| 1145 | |
| 1146 | queue.close(js); |
| 1147 | |
| 1148 | state.template transitionTo<StreamStates::Closed>(); |
| 1149 | doClose(js); |
| 1150 | } |
| 1151 | |
| 1152 | template <typename Self> |
| 1153 | void ReadableImpl<Self>::doClose(jsg::Lock& js) { |
| 1154 | // The state should have already been set to closed. |
| 1155 | KJ_ASSERT(state.template is<StreamStates::Closed>()); |
| 1156 | algorithms.clear(); |
| 1157 | } |
| 1158 | |
| 1159 | template <typename Self> |
| 1160 | void ReadableImpl<Self>::doError(jsg::Lock& js, jsg::Value reason) { |
| 1161 | // If already closed or errored, do nothing |
| 1162 | if (state.isInactive()) { |
| 1163 | return; |
| 1164 | } |
| 1165 | |
| 1166 | auto& queue = state.template getUnsafe<Queue>(); |
| 1167 | queue.error(js, reason.addRef(js)); |
| 1168 | state.template transitionTo<StreamStates::Errored>(kj::mv(reason)); |
| 1169 | algorithms.clear(); |
| 1170 | } |
| 1171 | |
| 1172 | template <typename Self> |
| 1173 | kj::Maybe<int> ReadableImpl<Self>::getDesiredSize() { |
| 1174 | if (state.template is<StreamStates::Closed>()) { |
| 1175 | return 0; |
| 1176 | } |
| 1177 | if (state.template is<StreamStates::Errored>()) { |
| 1178 | return kj::none; |
| 1179 | } |
| 1180 | return state.template getUnsafe<Queue>().desiredSize(); |
| 1181 | } |
| 1182 | |
| 1183 | // We should call pull if any of the consumers known to the queue have read requests or |
| 1184 | // we haven't yet signalled backpressure. |
| 1185 | template <typename Self> |
| 1186 | bool ReadableImpl<Self>::shouldCallPull() { |
| 1187 | return state.whenActiveOr( |
| 1188 | [this](Queue& q) { return q.wantsRead() || getDesiredSize().orDefault(0) > 0; }, false); |
| 1189 | } |
| 1190 | |
| 1191 | template <typename Self> |
| 1192 | void ReadableImpl<Self>::pullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1193 | // Determining if we need to pull is fairly complicated. All of the following |
| 1194 | // must hold true: |
| 1195 | if (!shouldCallPull()) { |
| 1196 | return; |
| 1197 | } |
| 1198 | |
| 1199 | if (flags.pulling) { |
| 1200 | flags.pullAgain = true; |
| 1201 | return; |
| 1202 | } |
| 1203 | KJ_ASSERT(!flags.pullAgain); |
| 1204 | flags.pulling = true; |
| 1205 | |
| 1206 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1207 | flags.pulling = false; |
| 1208 | if (flags.pullAgain) { |
| 1209 | flags.pullAgain = false; |
| 1210 | pullIfNeeded(js, kj::mv(self)); |
| 1211 | } |
| 1212 | }); |
| 1213 | |
| 1214 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1215 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1216 | flags.pulling = false; |
| 1217 | doError(js, kj::mv(reason)); |
| 1218 | }); |
| 1219 | |
| 1220 | maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); |
| 1221 | } |
| 1222 | |
| 1223 | template <typename Self> |
| 1224 | void ReadableImpl<Self>::forcePullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1225 | // Like pullIfNeeded but bypasses the shouldCallPull() check. Used for draining reads |
| 1226 | // which need to pull all available data regardless of backpressure settings. |
| 1227 | if (!canCloseOrEnqueue()) { |
| 1228 | return; |
| 1229 | } |
| 1230 | |
| 1231 | if (flags.pulling) { |
| 1232 | flags.pullAgain = true; |
| 1233 | return; |
| 1234 | } |
| 1235 | KJ_ASSERT(!flags.pullAgain); |
| 1236 | flags.pulling = true; |
| 1237 | |
| 1238 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1239 | flags.pulling = false; |
| 1240 | if (flags.pullAgain) { |
| 1241 | flags.pullAgain = false; |
| 1242 | // After a force pull, we go back to normal pullIfNeeded behavior. |
| 1243 | pullIfNeeded(js, kj::mv(self)); |
| 1244 | } |
| 1245 | }); |
| 1246 | |
| 1247 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1248 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1249 | flags.pulling = false; |
| 1250 | doError(js, kj::mv(reason)); |
| 1251 | }); |
| 1252 | |
| 1253 | maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); |
| 1254 | } |
| 1255 | |
| 1256 | template <typename Self> |
| 1257 | void ReadableImpl<Self>::visitForGc(jsg::GcVisitor& visitor) { |
| 1258 | state.visitForGc(visitor); |
| 1259 | KJ_IF_SOME(pendingCancel, maybePendingCancel) { |
| 1260 | visitor.visit(pendingCancel.fulfiller, pendingCancel.promise); |
| 1261 | } |
| 1262 | visitor.visit(algorithms); |
| 1263 | } |
| 1264 | |
| 1265 | template <typename Self> |
| 1266 | kj::Own<typename ReadableImpl<Self>::Consumer> ReadableImpl<Self>::getConsumer( |
| 1267 | kj::Maybe<ReadableImpl<Self>::StateListener&> listener) { |
| 1268 | auto& queue = state.template getUnsafe<Queue>(); |
| 1269 | return kj::heap<typename ReadableImpl<Self>::Consumer>(queue, listener); |
| 1270 | } |
| 1271 | |
| 1272 | // ====================================================================================== |
| 1273 | |
| 1274 | template <typename Self> |
| 1275 | WritableImpl<Self>::WritableImpl( |
| 1276 | jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal) |
| 1277 | : owner(owner.addWeakRef()), |
| 1278 | signal(kj::mv(abortSignal)) { |
| 1279 | flags.pedanticWpt = FeatureFlags::get(js).getPedanticWpt(); |
| 1280 | } |
| 1281 | |
| 1282 | template <typename Self> |
| 1283 | jsg::Promise<void> WritableImpl<Self>::abort( |
| 1284 | jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1285 | // Per the spec, the signal.reason should be a DOMException with name 'AbortError' |
| 1286 | // when no reason is provided, but the stored error should remain as the original reason. |
| 1287 | auto signalReason = [&]() -> jsg::JsValue { |
| 1288 | if (reason->IsUndefined() && FeatureFlags::get(js).getPedanticWpt()) { |
| 1289 | auto ex = js.domException( |
| 1290 | kj::str("AbortError"), kj::str("This writable stream has been aborted."), kj::none); |
| 1291 | return jsg::JsValue(KJ_ASSERT_NONNULL(ex.tryGetHandle(js))); |
| 1292 | } |
| 1293 | return jsg::JsValue(reason); |
| 1294 | }(); |
| 1295 | signal->triggerAbort(js, signalReason); |
| 1296 | |
| 1297 | // We have to check this again after the AbortSignal is triggered. |
| 1298 | if (state.isTerminal()) { |
| 1299 | return js.resolvedPromise(); |
| 1300 | } |
| 1301 | |
| 1302 | KJ_IF_SOME(pendingAbort, maybePendingAbort) { |
| 1303 | // Notice here that, per the spec, the reason given in this call of abort is |
| 1304 | // intentionally ignored if there is already an abort pending. |
| 1305 | return pendingAbort->whenResolved(js); |
| 1306 | } |
| 1307 | |
| 1308 | bool wasAlreadyErroring = false; |
| 1309 | if (state.template is<StreamStates::Erroring>()) { |
| 1310 | wasAlreadyErroring = true; |
| 1311 | reason = js.v8Undefined(); |
| 1312 | } |
| 1313 | |
| 1314 | KJ_DEFER(if (!wasAlreadyErroring) { startErroring(js, kj::mv(self), reason); }); |
| 1315 | |
| 1316 | maybePendingAbort = kj::heap<PendingAbort>(js, reason, wasAlreadyErroring); |
| 1317 | return KJ_ASSERT_NONNULL(maybePendingAbort)->whenResolved(js); |
| 1318 | } |
| 1319 | |
| 1320 | template <typename Self> |
| 1321 | kj::Maybe<WritableStreamJsController&> WritableImpl<Self>::tryGetOwner() { |
| 1322 | KJ_IF_SOME(o, owner) { |
| 1323 | return o->tryGet().map([](WritableStream& owner) -> WritableStreamJsController& { |
| 1324 | return static_cast<WritableStreamJsController&>(owner.getController()); |
| 1325 | }); |
| 1326 | } |
| 1327 | return kj::none; |
| 1328 | } |
| 1329 | |
| 1330 | template <typename Self> |
| 1331 | ssize_t WritableImpl<Self>::getDesiredSize() { |
| 1332 | return highWaterMark - amountBuffered; |
| 1333 | } |
| 1334 | |
| 1335 | template <typename Self> |
| 1336 | void WritableImpl<Self>::advanceQueueIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1337 | if (!flags.started || inFlightWrite != kj::none) { |
| 1338 | return; |
| 1339 | } |
| 1340 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1341 | |
| 1342 | if (state.template is<StreamStates::Erroring>()) { |
| 1343 | return finishErroring(js, kj::mv(self)); |
| 1344 | } |
| 1345 | |
| 1346 | if (writeRequests.empty()) { |
| 1347 | if (closeRequest != kj::none) { |
| 1348 | KJ_ASSERT(inFlightClose == kj::none); |
| 1349 | KJ_ASSERT_NONNULL(closeRequest); |
| 1350 | inFlightClose = kj::mv(closeRequest); |
| 1351 | |
| 1352 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), |
| 1353 | (jsg::Lock& js) { finishInFlightClose(js, kj::mv(self)); }); |
| 1354 | |
| 1355 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1356 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1357 | finishInFlightClose(js, kj::mv(self), reason.getHandle(js)); |
| 1358 | }); |
| 1359 | |
| 1360 | // Per the spec, the close algorithm should always run asynchronously, even if |
| 1361 | // there's no user-provided close handler. This ensures that releaseLock() can |
| 1362 | // reject the closed promise before the close completes. |
| 1363 | // The original maybeRunAlgorithm would call the onSuccess continuation |
| 1364 | // synchronously if algorithms.close is not specified. maybeRunAlgorithmAsync |
| 1365 | // always defers to a microtask. |
| 1366 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 1367 | maybeRunAlgorithmAsync(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure)); |
| 1368 | } else { |
| 1369 | maybeRunAlgorithm(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure)); |
| 1370 | } |
| 1371 | } |
| 1372 | return; |
| 1373 | } |
| 1374 | |
| 1375 | KJ_ASSERT(inFlightWrite == kj::none); |
| 1376 | auto req = dequeueWriteRequest(); |
| 1377 | auto value = req.value.addRef(js); |
| 1378 | auto size = req.size; |
| 1379 | inFlightWrite = kj::mv(req); |
| 1380 | |
| 1381 | auto onSuccess = |
| 1382 | JSG_VISITABLE_LAMBDA((this, self = self.addRef(), size), (self), (jsg::Lock& js) { |
| 1383 | amountBuffered -= size; |
| 1384 | finishInFlightWrite(js, self.addRef()); |
| 1385 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1386 | if (!isCloseQueuedOrInFlight() && isWritable()) { |
| 1387 | updateBackpressure(js); |
| 1388 | } |
| 1389 | if (state.template is<StreamStates::Erroring>() || writeRequests.empty()) { |
| 1390 | // In this case, we know advanceQueueIfNeeded won't recurse further, so we can |
| 1391 | // avoid the extra microtask hop. |
| 1392 | advanceQueueIfNeeded(js, kj::mv(self)); |
| 1393 | return js.resolvedPromise(); |
| 1394 | } |
| 1395 | // Here, however, let's avoid potentially deep recursion by hopping to a new |
| 1396 | // microtask to continue processing the queue. |
| 1397 | return js.resolvedPromise().then( |
| 1398 | js, JSG_VISITABLE_LAMBDA((this, self = kj::mv(self)), (self), (jsg::Lock & js) mutable { |
| 1399 | if (isWritable() || state.template is<StreamStates::Erroring>()) { |
| 1400 | advanceQueueIfNeeded(js, kj::mv(self)); |
| 1401 | } |
| 1402 | })); |
| 1403 | }); |
| 1404 | |
| 1405 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1406 | (this, self = self.addRef(), size), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1407 | amountBuffered -= size; |
| 1408 | finishInFlightWrite(js, kj::mv(self), reason.getHandle(js)); |
| 1409 | return js.resolvedPromise(); |
| 1410 | }); |
| 1411 | |
| 1412 | // Per the spec, the write algorithm should always run asynchronously, even if |
| 1413 | // there's no user-provided write handler. This ensures that backpressure changes |
| 1414 | // from the write don't resolve the ready promise synchronously, preserving correct |
| 1415 | // microtask ordering (e.g., ready rejects before closed on releaseLock). |
| 1416 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 1417 | maybeRunAlgorithmAsync(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure), |
| 1418 | value.getHandle(js), self.addRef()); |
| 1419 | } else { |
| 1420 | maybeRunAlgorithm(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure), |
| 1421 | value.getHandle(js), self.addRef()); |
| 1422 | } |
| 1423 | } |
| 1424 | |
| 1425 | template <typename Self> |
| 1426 | jsg::Promise<void> WritableImpl<Self>::close(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1427 | if (state.template is<StreamStates::Closed>()) { |
| 1428 | return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj)); |
| 1429 | } |
| 1430 | KJ_IF_SOME(errored, state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 1431 | return js.rejectedPromise<void>(errored.addRef(js)); |
| 1432 | } |
| 1433 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1434 | JSG_REQUIRE( |
| 1435 | !isCloseQueuedOrInFlight(), TypeError, "Cannot close a writer that is already being closed"); |
| 1436 | auto prp = js.newPromiseAndResolver<void>(); |
| 1437 | closeRequest = kj::mv(prp.resolver); |
| 1438 | |
| 1439 | if (flags.backpressure && isWritable()) { |
| 1440 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1441 | owner.maybeResolveReadyPromise(js); |
| 1442 | } |
| 1443 | } |
| 1444 | |
| 1445 | advanceQueueIfNeeded(js, kj::mv(self)); |
| 1446 | |
| 1447 | return kj::mv(prp.promise); |
| 1448 | } |
| 1449 | |
| 1450 | template <typename Self> |
| 1451 | void WritableImpl<Self>::dealWithRejection( |
| 1452 | jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1453 | if (isWritable()) { |
| 1454 | return startErroring(js, kj::mv(self), reason); |
| 1455 | } |
| 1456 | KJ_ASSERT(state.template is<StreamStates::Erroring>()); |
| 1457 | finishErroring(js, kj::mv(self)); |
| 1458 | } |
| 1459 | |
| 1460 | template <typename Self> |
| 1461 | WritableImpl<Self>::WriteRequest WritableImpl<Self>::dequeueWriteRequest() { |
| 1462 | auto write = kj::mv(writeRequests.front()); |
| 1463 | writeRequests.pop_front(); |
| 1464 | return kj::mv(write); |
| 1465 | } |
| 1466 | |
| 1467 | template <typename Self> |
| 1468 | void WritableImpl<Self>::doClose(jsg::Lock& js) { |
| 1469 | KJ_ASSERT(closeRequest == kj::none); |
| 1470 | KJ_ASSERT(inFlightClose == kj::none); |
| 1471 | KJ_ASSERT(inFlightWrite == kj::none); |
| 1472 | KJ_ASSERT(maybePendingAbort == kj::none); |
| 1473 | KJ_ASSERT(writeRequests.empty()); |
| 1474 | // State should have already been transitioned to Closed |
| 1475 | KJ_ASSERT(state.template is<StreamStates::Closed>()); |
| 1476 | algorithms.clear(); |
| 1477 | |
| 1478 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1479 | owner.doClose(js); |
| 1480 | } |
| 1481 | } |
| 1482 | |
| 1483 | template <typename Self> |
| 1484 | void WritableImpl<Self>::doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 1485 | KJ_ASSERT(closeRequest == kj::none); |
| 1486 | KJ_ASSERT(inFlightClose == kj::none); |
| 1487 | KJ_ASSERT(inFlightWrite == kj::none); |
| 1488 | KJ_ASSERT(maybePendingAbort == kj::none); |
| 1489 | KJ_ASSERT(writeRequests.empty()); |
| 1490 | // State should have already been transitioned to Errored |
| 1491 | KJ_ASSERT(state.template is<StreamStates::Errored>()); |
| 1492 | algorithms.clear(); |
| 1493 | |
| 1494 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1495 | owner.doError(js, reason); |
| 1496 | } |
| 1497 | } |
| 1498 | |
| 1499 | template <typename Self> |
| 1500 | void WritableImpl<Self>::error(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1501 | if (isWritable()) { |
| 1502 | algorithms.clear(); |
| 1503 | startErroring(js, kj::mv(self), reason); |
| 1504 | } |
| 1505 | } |
| 1506 | |
| 1507 | template <typename Self> |
| 1508 | void WritableImpl<Self>::finishErroring(jsg::Lock& js, jsg::Ref<Self> self) { |
| 1509 | auto erroring = kj::mv(KJ_ASSERT_NONNULL(state.template tryGetUnsafe<StreamStates::Erroring>())); |
| 1510 | auto reason = erroring.reason.getHandle(js); |
| 1511 | KJ_ASSERT(inFlightWrite == kj::none); |
| 1512 | KJ_ASSERT(inFlightClose == kj::none); |
| 1513 | state.template transitionTo<StreamStates::Errored>(kj::mv(erroring.reason)); |
| 1514 | |
| 1515 | while (!writeRequests.empty()) { |
| 1516 | dequeueWriteRequest().resolver.reject(js, reason); |
| 1517 | } |
| 1518 | KJ_ASSERT(writeRequests.empty()); |
| 1519 | |
| 1520 | KJ_IF_SOME(pendingAbort, maybePendingAbort) { |
| 1521 | if (pendingAbort->reject) { |
| 1522 | pendingAbort->fail(js, reason); |
| 1523 | return rejectCloseAndClosedPromiseIfNeeded(js); |
| 1524 | } |
| 1525 | |
| 1526 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1527 | auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort); |
| 1528 | pendingAbort->reject = false; |
| 1529 | pendingAbort->complete(js); |
| 1530 | rejectCloseAndClosedPromiseIfNeeded(js); |
| 1531 | }); |
| 1532 | |
| 1533 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1534 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1535 | auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort); |
| 1536 | pendingAbort->fail(js, reason.getHandle(js)); |
| 1537 | rejectCloseAndClosedPromiseIfNeeded(js); |
| 1538 | }); |
| 1539 | |
| 1540 | maybeRunAlgorithm(js, algorithms.abort, kj::mv(onSuccess), kj::mv(onFailure), reason); |
| 1541 | return; |
| 1542 | } |
| 1543 | rejectCloseAndClosedPromiseIfNeeded(js); |
| 1544 | } |
| 1545 | |
| 1546 | template <typename Self> |
| 1547 | void WritableImpl<Self>::finishInFlightClose( |
| 1548 | jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> maybeReason) { |
| 1549 | algorithms.clear(); |
| 1550 | KJ_ASSERT_NONNULL(inFlightClose); |
| 1551 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1552 | |
| 1553 | KJ_IF_SOME(reason, maybeReason) { |
| 1554 | maybeRejectPromise<void>(js, inFlightClose, reason); |
| 1555 | |
| 1556 | KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { |
| 1557 | pendingAbort->fail(js, reason); |
| 1558 | } |
| 1559 | |
| 1560 | return dealWithRejection(js, kj::mv(self), reason); |
| 1561 | } |
| 1562 | |
| 1563 | maybeResolvePromise(js, inFlightClose); |
| 1564 | |
| 1565 | if (state.template is<StreamStates::Erroring>()) { |
| 1566 | KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { |
| 1567 | pendingAbort->reject = false; |
| 1568 | pendingAbort->complete(js); |
| 1569 | } |
| 1570 | } |
| 1571 | KJ_ASSERT(maybePendingAbort == kj::none); |
| 1572 | |
| 1573 | state.template transitionTo<StreamStates::Closed>(); |
| 1574 | doClose(js); |
| 1575 | } |
| 1576 | |
| 1577 | template <typename Self> |
| 1578 | void WritableImpl<Self>::finishInFlightWrite( |
| 1579 | jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> maybeReason) { |
| 1580 | auto& write = KJ_ASSERT_NONNULL(inFlightWrite); |
| 1581 | |
| 1582 | KJ_IF_SOME(reason, maybeReason) { |
| 1583 | write.resolver.reject(js, reason); |
| 1584 | inFlightWrite = kj::none; |
| 1585 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1586 | return dealWithRejection(js, kj::mv(self), reason); |
| 1587 | } |
| 1588 | |
| 1589 | write.resolver.resolve(js); |
| 1590 | inFlightWrite = kj::none; |
| 1591 | } |
| 1592 | |
| 1593 | template <typename Self> |
| 1594 | bool WritableImpl<Self>::isCloseQueuedOrInFlight() { |
| 1595 | return closeRequest != kj::none || inFlightClose != kj::none; |
| 1596 | } |
| 1597 | |
| 1598 | template <typename Self> |
| 1599 | void WritableImpl<Self>::rejectCloseAndClosedPromiseIfNeeded(jsg::Lock& js) { |
| 1600 | algorithms.clear(); |
| 1601 | auto reason = |
| 1602 | KJ_ASSERT_NONNULL(state.template tryGetUnsafe<StreamStates::Errored>()).getHandle(js); |
| 1603 | maybeRejectPromise<void>(js, closeRequest, reason); |
| 1604 | PendingAbort::dequeue(maybePendingAbort); |
| 1605 | doError(js, reason); |
| 1606 | } |
| 1607 | |
| 1608 | template <typename Self> |
| 1609 | void WritableImpl<Self>::setup(jsg::Lock& js, |
| 1610 | jsg::Ref<Self> self, |
| 1611 | UnderlyingSink underlyingSink, |
| 1612 | StreamQueuingStrategy queuingStrategy) { |
| 1613 | KJ_ASSERT(!flags.started && !flags.starting); |
| 1614 | flags.starting = true; |
| 1615 | |
| 1616 | highWaterMark = queuingStrategy.highWaterMark.orDefault(1); |
| 1617 | auto startAlgorithm = kj::mv(underlyingSink.start); |
| 1618 | algorithms.write = kj::mv(underlyingSink.write); |
| 1619 | algorithms.close = kj::mv(underlyingSink.close); |
| 1620 | algorithms.abort = kj::mv(underlyingSink.abort); |
| 1621 | algorithms.size = kj::mv(queuingStrategy.size); |
| 1622 | // Per the streams spec, the size function should be called with `undefined` as `this`, |
| 1623 | // not as a method on the strategy object. |
| 1624 | KJ_IF_SOME(sizeFunc, algorithms.size) { |
| 1625 | sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined())); |
| 1626 | } |
| 1627 | |
| 1628 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { |
| 1629 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1630 | |
| 1631 | if (isWritable()) { |
| 1632 | // Only resolve the ready promise if an abort is not pending. |
| 1633 | // It will have been rejected already. |
| 1634 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1635 | owner.maybeResolveReadyPromise(js); |
| 1636 | } else { |
| 1637 | // Else block to avert dangling else compiler warning. |
| 1638 | } |
| 1639 | } |
| 1640 | |
| 1641 | flags.started = true; |
| 1642 | flags.starting = false; |
| 1643 | advanceQueueIfNeeded(js, kj::mv(self)); |
| 1644 | }); |
| 1645 | |
| 1646 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 1647 | (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { |
| 1648 | auto handle = reason.getHandle(js); |
| 1649 | KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>()); |
| 1650 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1651 | owner.maybeRejectReadyPromise(js, handle); |
| 1652 | } else { |
| 1653 | // Else block to avert dangling else compiler warning. |
| 1654 | } |
| 1655 | flags.started = true; |
| 1656 | flags.starting = false; |
| 1657 | dealWithRejection(js, kj::mv(self), handle); |
| 1658 | }); |
| 1659 | |
| 1660 | flags.backpressure = getDesiredSize() <= 0; |
| 1661 | |
| 1662 | maybeRunAlgorithm(js, startAlgorithm, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); |
| 1663 | } |
| 1664 | |
| 1665 | template <typename Self> |
| 1666 | void WritableImpl<Self>::startErroring( |
| 1667 | jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) { |
| 1668 | KJ_ASSERT(isWritable()); |
| 1669 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1670 | owner.maybeRejectReadyPromise(js, reason); |
| 1671 | } |
| 1672 | state.template transitionTo<StreamStates::Erroring>(js.v8Ref(reason)); |
| 1673 | if (inFlightWrite == kj::none && inFlightClose == kj::none && flags.started) { |
| 1674 | finishErroring(js, kj::mv(self)); |
| 1675 | } |
| 1676 | } |
| 1677 | |
| 1678 | template <typename Self> |
| 1679 | void WritableImpl<Self>::updateBackpressure(jsg::Lock& js) { |
| 1680 | KJ_ASSERT(isWritable()); |
| 1681 | KJ_ASSERT(!isCloseQueuedOrInFlight()); |
| 1682 | bool bp = getDesiredSize() <= 0; |
| 1683 | |
| 1684 | if (bp != flags.backpressure) { |
| 1685 | flags.backpressure = bp; |
| 1686 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1687 | owner.updateBackpressure(js, flags.backpressure); |
| 1688 | } |
| 1689 | } |
| 1690 | } |
| 1691 | |
| 1692 | template <typename Self> |
| 1693 | jsg::Promise<void> WritableImpl<Self>::write( |
| 1694 | jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> value) { |
| 1695 | |
| 1696 | size_t size = 1; |
| 1697 | KJ_IF_SOME(sizeFunc, algorithms.size) { |
| 1698 | kj::Maybe<jsg::Value> failure; |
| 1699 | JSG_TRY(js) { |
| 1700 | size = sizeFunc(js, value); |
| 1701 | } |
| 1702 | JSG_CATCH(exception) { |
| 1703 | startErroring(js, self.addRef(), exception.getHandle(js)); |
| 1704 | failure = kj::mv(exception); |
| 1705 | } |
| 1706 | KJ_IF_SOME(exception, failure) { |
| 1707 | return js.rejectedPromise<void>(kj::mv(exception)); |
| 1708 | } |
| 1709 | } |
| 1710 | |
| 1711 | // Per spec (WritableStreamDefaultWriterWrite step 5), after calling the size |
| 1712 | // algorithm, re-check that the stream is still locked to a writer. If |
| 1713 | // releaseLock() was called from within strategy.size(), the write must be |
| 1714 | // rejected. This check must occur before any state checks, as the stream |
| 1715 | // state may still appear writable even after the writer was released. |
| 1716 | if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { |
| 1717 | KJ_IF_SOME(owner, tryGetOwner()) { |
| 1718 | if (!owner.isLockedToWriter()) { |
| 1719 | return js.rejectedPromise<void>( |
| 1720 | js.v8TypeError("This WritableStream writer has been released."_kjc)); |
| 1721 | } |
| 1722 | } |
| 1723 | } |
| 1724 | |
| 1725 | KJ_IF_SOME(error, state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 1726 | return js.rejectedPromise<void>(error.addRef(js)); |
| 1727 | } |
| 1728 | |
| 1729 | if (isCloseQueuedOrInFlight() || state.template is<StreamStates::Closed>()) { |
| 1730 | return js.rejectedPromise<void>(js.v8TypeError("This ReadableStream is closed."_kj)); |
| 1731 | } |
| 1732 | |
| 1733 | KJ_IF_SOME(erroring, state.template tryGetUnsafe<StreamStates::Erroring>()) { |
| 1734 | return js.rejectedPromise<void>(erroring.reason.addRef(js)); |
| 1735 | } |
| 1736 | |
| 1737 | KJ_ASSERT(isWritable()); |
| 1738 | |
| 1739 | auto prp = js.newPromiseAndResolver<void>(); |
| 1740 | writeRequests.push_back(WriteRequest{ |
| 1741 | .resolver = kj::mv(prp.resolver), |
| 1742 | .value = js.v8Ref(value), |
| 1743 | .size = size, |
| 1744 | }); |
| 1745 | amountBuffered += size; |
| 1746 | |
| 1747 | updateBackpressure(js); |
| 1748 | advanceQueueIfNeeded(js, kj::mv(self)); |
| 1749 | return kj::mv(prp.promise); |
| 1750 | } |
| 1751 | |
| 1752 | template <typename Self> |
| 1753 | void WritableImpl<Self>::visitForGc(jsg::GcVisitor& visitor) { |
| 1754 | state.visitForGc(visitor); |
| 1755 | visitor.visit(inFlightWrite, inFlightClose, closeRequest, algorithms, signal); |
| 1756 | KJ_IF_SOME(pendingAbort, maybePendingAbort) { |
| 1757 | visitor.visit(*pendingAbort); |
| 1758 | } |
| 1759 | visitor.visitAll(writeRequests); |
| 1760 | } |
| 1761 | |
| 1762 | template <typename Self> |
| 1763 | bool WritableImpl<Self>::isWritable() const { |
| 1764 | return state.isActive(); |
| 1765 | } |
| 1766 | |
| 1767 | template <typename Self> |
| 1768 | void WritableImpl<Self>::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) { |
| 1769 | for (auto& write: writeRequests) { |
| 1770 | write.resolver.reject(js, reason); |
| 1771 | } |
| 1772 | writeRequests.clear(); |
| 1773 | } |
| 1774 | |
| 1775 | // ====================================================================================== |
| 1776 | |
| 1777 | namespace { |
| 1778 | template <typename Controller, typename Queue> |
| 1779 | struct ReadableState { |
| 1780 | Controller controller; |
| 1781 | kj::Own<typename Queue::Consumer> consumer; |
| 1782 | ReadableStreamJsController& owner; |
| 1783 | |
| 1784 | ReadableState(Controller controller, |
| 1785 | kj::Own<typename Queue::Consumer> consumer, |
| 1786 | ReadableStreamJsController& owner) |
| 1787 | : controller(kj::mv(controller)), |
| 1788 | consumer(kj::mv(consumer)), |
| 1789 | owner(owner) {} |
| 1790 | |
| 1791 | ReadableState(Controller controller, |
| 1792 | Queue::ConsumerImpl::StateListener& listener, |
| 1793 | ReadableStreamJsController& owner) |
| 1794 | : ReadableState(controller.addRef(), controller->getConsumer(listener), owner) {} |
| 1795 | |
| 1796 | ReadableState clone(jsg::Lock& js, |
| 1797 | Queue::ConsumerImpl::StateListener& listener, |
| 1798 | ReadableStreamJsController& owner) { |
| 1799 | return ReadableState(controller.addRef(), consumer->clone(js, listener), owner); |
| 1800 | } |
| 1801 | }; |
| 1802 | |
| 1803 | struct ValueReadable final: private api::ValueQueue::ConsumerImpl::StateListener { |
| 1804 | |
| 1805 | using State = ReadableState<DefaultController, ValueQueue>; |
| 1806 | kj::Maybe<State> state; |
| 1807 | bool reading = false; |
| 1808 | bool pendingCancel = false; |
| 1809 | |
| 1810 | JSG_MEMORY_INFO(ValueReadable) { |
| 1811 | KJ_IF_SOME(s, state) { |
| 1812 | tracker.trackField("controller", s.controller); |
| 1813 | tracker.trackField("consumer", s.consumer); |
| 1814 | } |
| 1815 | } |
| 1816 | |
| 1817 | void visitForGc(jsg::GcVisitor& visitor) { |
| 1818 | KJ_IF_SOME(s, state) { |
| 1819 | visitor.visit(s.controller, *s.consumer); |
| 1820 | } |
| 1821 | } |
| 1822 | |
| 1823 | ValueReadable(DefaultController controller, ReadableStreamJsController& owner) |
| 1824 | : state(State(kj::mv(controller), *this, owner)) {} |
| 1825 | |
| 1826 | ValueReadable(jsg::Lock& js, ReadableStreamJsController& owner, ValueReadable& other) |
| 1827 | : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)) {} |
| 1828 | |
| 1829 | KJ_DISALLOW_COPY_AND_MOVE(ValueReadable); |
| 1830 | |
| 1831 | void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { |
| 1832 | KJ_IF_SOME(s, state) { |
| 1833 | s.consumer->cancelPendingReads(js, reason); |
| 1834 | } |
| 1835 | } |
| 1836 | |
| 1837 | kj::Own<ValueReadable> clone(jsg::Lock& js, ReadableStreamJsController& owner) { |
| 1838 | // A single ReadableStreamDefaultController can have multiple consumers. |
| 1839 | // When the ValueReadable constructor is used, the new consumer is added |
| 1840 | // and starts to receive new data that becomes enqueued. When clone |
| 1841 | // is used, any state currently held by this consumer is copied to the |
| 1842 | // new consumer. |
| 1843 | return kj::heap<ValueReadable>(js, owner, *this); |
| 1844 | } |
| 1845 | |
| 1846 | jsg::Promise<ReadResult> read(jsg::Lock& js) { |
| 1847 | KJ_IF_SOME(s, state) { |
| 1848 | auto prp = js.newPromiseAndResolver<ReadResult>(); |
| 1849 | reading = true; |
| 1850 | s.consumer->read(js, |
| 1851 | ValueQueue::ReadRequest{ |
| 1852 | .resolver = kj::mv(prp.resolver), |
| 1853 | }); |
| 1854 | reading = false; |
| 1855 | if (pendingCancel) { |
| 1856 | // If we were canceled while reading, we need to drop our state now. |
| 1857 | state = kj::none; |
| 1858 | pendingCancel = false; |
| 1859 | } |
| 1860 | return kj::mv(prp.promise); |
| 1861 | } |
| 1862 | |
| 1863 | // We are canceled! There's nothing to do. |
| 1864 | return js.resolvedPromise(ReadResult{.done = true}); |
| 1865 | } |
| 1866 | |
| 1867 | jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead) { |
| 1868 | KJ_IF_SOME(s, state) { |
| 1869 | // Note: We do NOT call beginOperation()/endOperation() here. The caller |
| 1870 | // (ReadableStreamJsController::drainingRead) manages the operation scope |
| 1871 | // around both this call and the returned promise's lifetime. If we added |
| 1872 | // our own beginOperation/endOperation here, the endOperation would fire |
| 1873 | // before the caller's wrapDrainingRead could set up its .then() callbacks, |
| 1874 | // potentially destroying the Consumer while the returned promise still has |
| 1875 | // dangling this-capturing callbacks from consumer->drainingRead(). |
| 1876 | return s.consumer->drainingRead(js, maxRead); |
| 1877 | } |
| 1878 | |
| 1879 | // We are canceled! Return done with empty chunks. |
| 1880 | return js.resolvedPromise(DrainingReadResult{ |
| 1881 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 1882 | .done = true, |
| 1883 | }); |
| 1884 | } |
| 1885 | |
| 1886 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 1887 | // When a ReadableStream is canceled, the expected behavior is that the underlying |
| 1888 | // controller is notified and the cancel algorithm on the underlying source is |
| 1889 | // called. When there are multiple ReadableStreams sharing consumption of a |
| 1890 | // controller, however, it should act as a shared pointer of sorts, canceling |
| 1891 | // the underlying controller only when the last reader is canceled. |
| 1892 | // Here, we rely on the controller implementing the correct behavior since it owns |
| 1893 | // the queue that knows about all of the attached consumers. |
| 1894 | if (pendingCancel) return js.resolvedPromise(); |
| 1895 | KJ_IF_SOME(s, state) { |
| 1896 | // Check if there's a pending draining read before calling cancel, since cancel |
| 1897 | // will resolve the pending read and we need to know if we should defer destruction. |
| 1898 | bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead(); |
| 1899 | s.consumer->cancel(js, maybeReason); |
| 1900 | auto promise = s.controller->cancel(js, kj::mv(maybeReason)); |
| 1901 | // If we're currently in a read (sync or draining), we need to wait for that to |
| 1902 | // finish before dropping our state. For draining reads, the promise callbacks |
| 1903 | // capture 'this' (the Consumer) to clear hasPendingDrainingRead. If we destroy |
| 1904 | // the state now, those callbacks will UAF. |
| 1905 | if (reading || hasPendingDrainingRead) { |
| 1906 | pendingCancel = true; |
| 1907 | } else { |
| 1908 | state = kj::none; |
| 1909 | } |
| 1910 | return kj::mv(promise); |
| 1911 | } |
| 1912 | |
| 1913 | return js.resolvedPromise(); |
| 1914 | } |
| 1915 | |
| 1916 | void onConsumerClose(jsg::Lock& js) override { |
| 1917 | // Called by the consumer when a state change to closed happens. |
| 1918 | // We need to notify the owner. Note that the owner may drop this |
| 1919 | // readable in doClose so it is not safe to access anything on this |
| 1920 | // after calling doClose. |
| 1921 | KJ_IF_SOME(s, state) { |
| 1922 | s.owner.doClose(js); |
| 1923 | } |
| 1924 | } |
| 1925 | |
| 1926 | void onConsumerError(jsg::Lock& js, jsg::Value reason) override { |
| 1927 | // Called by the consumer when a state change to errored happens. |
| 1928 | // We need to notify the owner. Note that the owner may drop this |
| 1929 | // readable in doClose so it is not safe to access anything on this |
| 1930 | // after calling doError. |
| 1931 | KJ_IF_SOME(s, state) { |
| 1932 | s.owner.doError(js, reason.getHandle(js)); |
| 1933 | } |
| 1934 | } |
| 1935 | |
| 1936 | bool onConsumerWantsData(jsg::Lock& js) override { |
| 1937 | // Called by the consumer when it has a queued pending read and needs |
| 1938 | // data to be provided to fulfill it. We need to notify the controller |
| 1939 | // to initiate pulling to provide the data. |
| 1940 | // Returns true if the pull completed synchronously (meaning more pumping |
| 1941 | // might yield additional synchronous data), false otherwise. |
| 1942 | KJ_IF_SOME(s, state) { |
| 1943 | // Save a reference to the owner before calling pull. The pull callback |
| 1944 | // may trigger close/error which could destroy this ValueReadable. By |
| 1945 | // using beginOperation(), we ensure doClose/doError defers the |
| 1946 | // actual destruction until after we return. |
| 1947 | ReadableStreamJsController& owner = s.owner; |
| 1948 | owner.state.beginOperation(); |
| 1949 | |
| 1950 | // For draining reads, use forcePull to bypass backpressure checks. |
| 1951 | // This ensures we pull all available data regardless of highWaterMark. |
| 1952 | if (s.consumer->hasPendingDrainingRead()) { |
| 1953 | s.controller->forcePull(js); |
| 1954 | } else { |
| 1955 | s.controller->pull(js); |
| 1956 | } |
| 1957 | |
| 1958 | // Check if state is still valid BEFORE calling endOperation(), |
| 1959 | // because that call may destroy this ValueReadable if close was deferred. |
| 1960 | bool result = |
| 1961 | state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false); |
| 1962 | |
| 1963 | // Process any deferred close/error. This may destroy this ValueReadable. |
| 1964 | if (owner.state.endOperation()) { |
| 1965 | // A pending state was applied. Call the appropriate callback. |
| 1966 | if (owner.state.template is<StreamStates::Closed>()) { |
| 1967 | owner.lock.onClose(js); |
| 1968 | } else if (owner.state.template is<StreamStates::Errored>()) { |
| 1969 | KJ_IF_SOME(err, owner.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 1970 | owner.lock.onError(js, err.getHandle(js)); |
| 1971 | } |
| 1972 | } |
| 1973 | } |
| 1974 | |
| 1975 | return result; |
| 1976 | } |
| 1977 | return false; |
| 1978 | } |
| 1979 | |
| 1980 | kj::Maybe<int> getDesiredSize() { |
| 1981 | KJ_IF_SOME(s, state) { |
| 1982 | return s.controller->getDesiredSize(); |
| 1983 | } |
| 1984 | return kj::none; |
| 1985 | } |
| 1986 | |
| 1987 | bool canCloseOrEnqueue() { |
| 1988 | return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false); |
| 1989 | } |
| 1990 | |
| 1991 | kj::Maybe<DefaultController> getControllerRef() { |
| 1992 | return state.map([](State& s) { return s.controller.addRef(); }); |
| 1993 | } |
| 1994 | }; |
| 1995 | |
| 1996 | struct ByteReadable final: private api::ByteQueue::ConsumerImpl::StateListener { |
| 1997 | |
| 1998 | using State = ReadableState<ByobController, ByteQueue>; |
| 1999 | kj::Maybe<State> state; |
| 2000 | kj::Maybe<int> autoAllocateChunkSize; |
| 2001 | bool pendingCancel = false; |
| 2002 | |
| 2003 | JSG_MEMORY_INFO(ByteReadable) { |
| 2004 | KJ_IF_SOME(s, state) { |
| 2005 | tracker.trackField("controller", s.controller); |
| 2006 | tracker.trackField("consumer", s.consumer); |
| 2007 | } |
| 2008 | } |
| 2009 | |
| 2010 | void visitForGc(jsg::GcVisitor& visitor) { |
| 2011 | KJ_IF_SOME(s, state) { |
| 2012 | visitor.visit(s.controller, *s.consumer); |
| 2013 | } |
| 2014 | } |
| 2015 | |
| 2016 | ByteReadable(ByobController controller, |
| 2017 | ReadableStreamJsController& owner, |
| 2018 | kj::Maybe<int> autoAllocateChunkSize) |
| 2019 | : state(State(kj::mv(controller), *this, owner)), |
| 2020 | autoAllocateChunkSize(autoAllocateChunkSize) {} |
| 2021 | |
| 2022 | ByteReadable(jsg::Lock& js, ReadableStreamJsController& owner, ByteReadable& other) |
| 2023 | : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)), |
| 2024 | autoAllocateChunkSize(other.autoAllocateChunkSize) {} |
| 2025 | |
| 2026 | KJ_DISALLOW_COPY_AND_MOVE(ByteReadable); |
| 2027 | |
| 2028 | void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { |
| 2029 | KJ_IF_SOME(s, state) { |
| 2030 | s.consumer->cancelPendingReads(js, reason); |
| 2031 | } |
| 2032 | } |
| 2033 | |
| 2034 | // A single ReadableByteStreamController can have multiple consumers. |
| 2035 | // When the ByteReadable constructor is used, the new consumer is added |
| 2036 | // and starts to receive new data that becomes enqueued. When clone |
| 2037 | // is used, any state currently held by this consumer is copied to the |
| 2038 | // new consumer. |
| 2039 | kj::Own<ByteReadable> clone(jsg::Lock& js, ReadableStreamJsController& owner) { |
| 2040 | return kj::heap<ByteReadable>(js, owner, *this); |
| 2041 | } |
| 2042 | |
| 2043 | jsg::Promise<ReadResult> read( |
| 2044 | jsg::Lock& js, kj::Maybe<ReadableStreamController::ByobOptions> byobOptions) { |
| 2045 | KJ_IF_SOME(s, state) { |
| 2046 | auto prp = js.newPromiseAndResolver<ReadResult>(); |
| 2047 | |
| 2048 | KJ_IF_SOME(byob, byobOptions) { |
| 2049 | jsg::BufferSource source(js, byob.bufferView.getHandle(js)); |
| 2050 | // If atLeast is not given, then by default it is the element size of the view |
| 2051 | // that we were given. If atLeast is given, we make sure that it is aligned |
| 2052 | // with the element size. No matter what, atLeast cannot be less than 1. |
| 2053 | auto atLeast = kj::max(source.getElementSize(), byob.atLeast.orDefault(1)); |
| 2054 | atLeast = kj::max(1, atLeast - (atLeast % source.getElementSize())); |
| 2055 | s.consumer->read(js, |
| 2056 | ByteQueue::ReadRequest(kj::mv(prp.resolver), |
| 2057 | { |
| 2058 | .store = jsg::BufferSource(js, source.detach(js)), |
| 2059 | .atLeast = atLeast, |
| 2060 | .type = ByteQueue::ReadRequest::Type::BYOB, |
| 2061 | })); |
| 2062 | } else KJ_IF_SOME(chunkSize, autoAllocateChunkSize) { |
| 2063 | // autoAllocateChunkSize is set, so we allocate a buffer and do a BYOB read. |
| 2064 | // This makes the buffer available to the underlying source via controller.byobRequest. |
| 2065 | KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, chunkSize)) { |
| 2066 | // Ensure that the handle is created here so that the size of the buffer |
| 2067 | // is accounted for in the isolate memory tracking. |
| 2068 | s.consumer->read(js, |
| 2069 | ByteQueue::ReadRequest(kj::mv(prp.resolver), |
| 2070 | { |
| 2071 | .store = kj::mv(store), |
| 2072 | .type = ByteQueue::ReadRequest::Type::BYOB, |
| 2073 | })); |
| 2074 | } else { |
| 2075 | prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read.")); |
| 2076 | } |
| 2077 | } else { |
| 2078 | // autoAllocateChunkSize is not set. Per spec, we do a DEFAULT read which means |
| 2079 | // the underlying source's pull method won't get a byobRequest. It must use |
| 2080 | // controller.enqueue() to provide data instead. |
| 2081 | constexpr size_t kDefaultReadSize = 16384; // 16KB default buffer |
| 2082 | KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, kDefaultReadSize)) { |
| 2083 | s.consumer->read(js, |
| 2084 | ByteQueue::ReadRequest(kj::mv(prp.resolver), |
| 2085 | { |
| 2086 | .store = kj::mv(store), |
| 2087 | .type = ByteQueue::ReadRequest::Type::DEFAULT, |
| 2088 | })); |
| 2089 | } else { |
| 2090 | prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read.")); |
| 2091 | } |
| 2092 | } |
| 2093 | |
| 2094 | return kj::mv(prp.promise); |
| 2095 | } |
| 2096 | |
| 2097 | // We are canceled! There's nothing else to do. |
| 2098 | KJ_IF_SOME(byob, byobOptions) { |
| 2099 | // If a BYOB buffer was given, we need to give it back wrapped in a TypedArray |
| 2100 | // whose size is set to zero. |
| 2101 | jsg::BufferSource source(js, byob.bufferView.getHandle(js)); |
| 2102 | auto store = source.detach(js); |
| 2103 | store.consume(store.size()); |
| 2104 | return js.resolvedPromise(ReadResult{ |
| 2105 | .value = js.v8Ref(store.createHandle(js)), |
| 2106 | .done = true, |
| 2107 | }); |
| 2108 | } else { |
| 2109 | return js.resolvedPromise(ReadResult{.done = true}); |
| 2110 | } |
| 2111 | } |
| 2112 | |
| 2113 | jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead) { |
| 2114 | KJ_IF_SOME(s, state) { |
| 2115 | // Note: We do NOT call beginOperation()/endOperation() here. The caller |
| 2116 | // (ReadableStreamJsController::drainingRead) manages the operation scope |
| 2117 | // around both this call and the returned promise's lifetime. See the |
| 2118 | // comment in ValueReadable::drainingRead for the detailed explanation. |
| 2119 | return s.consumer->drainingRead(js, maxRead); |
| 2120 | } |
| 2121 | |
| 2122 | // We are canceled! Return done with empty chunks. |
| 2123 | return js.resolvedPromise(DrainingReadResult{ |
| 2124 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 2125 | .done = true, |
| 2126 | }); |
| 2127 | } |
| 2128 | |
| 2129 | // When a ReadableStream is canceled, the expected behavior is that the underlying |
| 2130 | // controller is notified and the cancel algorithm on the underlying source is |
| 2131 | // called. When there are multiple ReadableStreams sharing consumption of a |
| 2132 | // controller, however, it should act as a shared pointer of sorts, canceling |
| 2133 | // the underlying controller only when the last reader is canceled. |
| 2134 | // Here, we rely on the controller implementing the correct behavior since it owns |
| 2135 | // the queue that knows about all of the attached consumers. |
| 2136 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 2137 | if (pendingCancel) return js.resolvedPromise(); |
| 2138 | KJ_IF_SOME(s, state) { |
| 2139 | // Check if there's a pending draining read before calling cancel, since cancel |
| 2140 | // will resolve the pending read and we need to know if we should defer destruction. |
| 2141 | bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead(); |
| 2142 | s.consumer->cancel(js, maybeReason); |
| 2143 | auto promise = s.controller->cancel(js, kj::mv(maybeReason)); |
| 2144 | // If there's a pending draining read, we need to wait for it to finish before |
| 2145 | // dropping our state. The draining read's promise callbacks capture 'this' (the |
| 2146 | // Consumer) to clear hasPendingDrainingRead. If we destroy the state now, those |
| 2147 | // callbacks will UAF. |
| 2148 | if (hasPendingDrainingRead) { |
| 2149 | pendingCancel = true; |
| 2150 | } else { |
| 2151 | state = kj::none; |
| 2152 | } |
| 2153 | return kj::mv(promise); |
| 2154 | } |
| 2155 | |
| 2156 | return js.resolvedPromise(); |
| 2157 | } |
| 2158 | |
| 2159 | void onConsumerClose(jsg::Lock& js) override { |
| 2160 | // Note that the owner may drop this readable in doClose so it |
| 2161 | // is not safe to access anything on this after calling doClose. |
| 2162 | KJ_IF_SOME(s, state) { |
| 2163 | s.owner.doClose(js); |
| 2164 | } |
| 2165 | } |
| 2166 | |
| 2167 | void onConsumerError(jsg::Lock& js, jsg::Value reason) override { |
| 2168 | // Note that the owner may drop this readable in doClose so it |
| 2169 | // is not safe to access anything on this after calling doError. |
| 2170 | KJ_IF_SOME(s, state) { |
| 2171 | s.owner.doError(js, reason.getHandle(js)); |
| 2172 | }; |
| 2173 | } |
| 2174 | |
| 2175 | // Called by the consumer when it has a queued pending read and needs |
| 2176 | // data to be provided to fulfill it. We need to notify the controller |
| 2177 | // to initiate pulling to provide the data. |
| 2178 | // Returns true if the pull completed synchronously (meaning more pumping |
| 2179 | // might yield additional synchronous data), false otherwise. |
| 2180 | bool onConsumerWantsData(jsg::Lock& js) override { |
| 2181 | KJ_IF_SOME(s, state) { |
| 2182 | // Save a reference to the owner before calling pull. The pull callback |
| 2183 | // may trigger close/error which could destroy this ByteReadable. By |
| 2184 | // using beginOperation(), we ensure doClose/doError defers the |
| 2185 | // actual destruction until after we return. |
| 2186 | ReadableStreamJsController& owner = s.owner; |
| 2187 | owner.state.beginOperation(); |
| 2188 | |
| 2189 | // For draining reads, use forcePull to bypass backpressure checks. |
| 2190 | // This ensures we pull all available data regardless of highWaterMark. |
| 2191 | if (s.consumer->hasPendingDrainingRead()) { |
| 2192 | s.controller->forcePull(js); |
| 2193 | } else { |
| 2194 | s.controller->pull(js); |
| 2195 | } |
| 2196 | |
| 2197 | // Check if state is still valid BEFORE calling endOperation(), |
| 2198 | // because that call may destroy this ByteReadable if close was deferred. |
| 2199 | bool result = |
| 2200 | state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false); |
| 2201 | |
| 2202 | // Process any deferred close/error. This may destroy this ByteReadable. |
| 2203 | if (owner.state.endOperation()) { |
| 2204 | // A pending state was applied. Call the appropriate callback. |
| 2205 | if (owner.state.template is<StreamStates::Closed>()) { |
| 2206 | owner.lock.onClose(js); |
| 2207 | } else if (owner.state.template is<StreamStates::Errored>()) { |
| 2208 | KJ_IF_SOME(err, owner.state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 2209 | owner.lock.onError(js, err.getHandle(js)); |
| 2210 | } |
| 2211 | } |
| 2212 | } |
| 2213 | |
| 2214 | return result; |
| 2215 | } |
| 2216 | return false; |
| 2217 | } |
| 2218 | |
| 2219 | kj::Maybe<int> getDesiredSize() { |
| 2220 | KJ_IF_SOME(s, state) { |
| 2221 | return s.controller->getDesiredSize(); |
| 2222 | } |
| 2223 | return kj::none; |
| 2224 | } |
| 2225 | |
| 2226 | bool canCloseOrEnqueue() { |
| 2227 | return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false); |
| 2228 | } |
| 2229 | |
| 2230 | kj::Maybe<ByobController> getControllerRef() { |
| 2231 | return state.map([](State& state) { return state.controller.addRef(); }); |
| 2232 | } |
| 2233 | }; |
| 2234 | } // namespace |
| 2235 | |
| 2236 | // ======================================================================================= |
| 2237 | |
| 2238 | ReadableStreamDefaultController::ReadableStreamDefaultController( |
| 2239 | UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) |
| 2240 | : ioContext(tryGetIoContext()), |
| 2241 | impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} |
| 2242 | |
| 2243 | kj::Maybe<StreamStates::Errored> ReadableStreamDefaultController::getMaybeErrorState( |
| 2244 | jsg::Lock& js) { |
| 2245 | KJ_IF_SOME(errored, impl.state.tryGetUnsafe<StreamStates::Errored>()) { |
| 2246 | return errored.addRef(js); |
| 2247 | } |
| 2248 | return kj::none; |
| 2249 | } |
| 2250 | |
| 2251 | void ReadableStreamDefaultController::start(jsg::Lock& js) { |
| 2252 | impl.start(js, JSG_THIS); |
| 2253 | } |
| 2254 | |
| 2255 | bool ReadableStreamDefaultController::canCloseOrEnqueue() { |
| 2256 | return impl.canCloseOrEnqueue(); |
| 2257 | } |
| 2258 | |
| 2259 | bool ReadableStreamDefaultController::hasBackpressure() { |
| 2260 | return !impl.shouldCallPull(); |
| 2261 | } |
| 2262 | |
| 2263 | kj::Maybe<int> ReadableStreamDefaultController::getDesiredSize() { |
| 2264 | return impl.getDesiredSize(); |
| 2265 | } |
| 2266 | |
| 2267 | void ReadableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { |
| 2268 | visitor.visit(impl); |
| 2269 | } |
| 2270 | |
| 2271 | jsg::Promise<void> ReadableStreamDefaultController::cancel( |
| 2272 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 2273 | return impl.cancel(js, JSG_THIS, maybeReason.orDefault([&] { return js.v8Undefined(); })); |
| 2274 | } |
| 2275 | |
| 2276 | void ReadableStreamDefaultController::close(jsg::Lock& js) { |
| 2277 | impl.close(js); |
| 2278 | } |
| 2279 | |
| 2280 | void ReadableStreamDefaultController::enqueue( |
| 2281 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> chunk) { |
| 2282 | // Hold a strong reference to prevent this controller from being freed if the |
| 2283 | // user-provided size algorithm (below) re-enters JS and errors the controller |
| 2284 | // through a side-channel (e.g. TransformStreamDefaultController::error() |
| 2285 | // dropping all external jsg::Refs to this controller). |
| 2286 | auto self = JSG_THIS; |
| 2287 | auto value = chunk.orDefault(js.undefined()); |
| 2288 | |
| 2289 | JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "Unable to enqueue"); |
| 2290 | |
| 2291 | size_t size = 1; |
| 2292 | bool errored = false; |
| 2293 | KJ_IF_SOME(sizeFunc, impl.algorithms.size) { |
| 2294 | js.tryCatch([&] { size = sizeFunc(js, value); }, [&](jsg::Value exception) { |
| 2295 | impl.doError(js, kj::mv(exception)); |
| 2296 | errored = true; |
| 2297 | }); |
| 2298 | } |
| 2299 | |
| 2300 | // Re-check canCloseOrEnqueue: the size callback may have errored us without |
| 2301 | // throwing (e.g. by calling transformController.error()), in which case |
| 2302 | // `errored` is still false but the impl state has transitioned to Errored. |
| 2303 | if (!errored && impl.canCloseOrEnqueue()) { |
| 2304 | impl.enqueue(js, kj::rc<ValueQueue::Entry>(js.v8Ref(value), size), kj::mv(self)); |
| 2305 | } |
| 2306 | } |
| 2307 | |
| 2308 | void ReadableStreamDefaultController::error(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2309 | impl.doError(js, js.v8Ref(reason)); |
| 2310 | } |
| 2311 | |
| 2312 | // When a consumer receives a read request, but does not have the data available to |
| 2313 | // fulfill the request, the consumer will call pull on the controller to pull that |
| 2314 | // data if needed. |
| 2315 | void ReadableStreamDefaultController::pull(jsg::Lock& js) { |
| 2316 | impl.pullIfNeeded(js, JSG_THIS); |
| 2317 | } |
| 2318 | |
| 2319 | void ReadableStreamDefaultController::forcePull(jsg::Lock& js) { |
| 2320 | impl.forcePullIfNeeded(js, JSG_THIS); |
| 2321 | } |
| 2322 | |
| 2323 | kj::Own<ValueQueue::Consumer> ReadableStreamDefaultController::getConsumer( |
| 2324 | kj::Maybe<ValueQueue::ConsumerImpl::StateListener&> stateListener) { |
| 2325 | return impl.getConsumer(stateListener); |
| 2326 | } |
| 2327 | |
| 2328 | // ====================================================================================== |
| 2329 | |
| 2330 | ReadableStreamBYOBRequest::Impl::Impl(jsg::Lock& js, |
| 2331 | kj::Own<ByteQueue::ByobRequest> readRequest, |
| 2332 | kj::Rc<WeakRef<ReadableByteStreamController>> controller) |
| 2333 | : readRequest(kj::mv(readRequest)), |
| 2334 | controller(kj::mv(controller)), |
| 2335 | view(js.v8Ref(this->readRequest->getView(js))), |
| 2336 | originalBufferByteLength(this->readRequest->getOriginalBufferByteLength(js)), |
| 2337 | originalByteOffsetPlusBytesFilled(this->readRequest->getOriginalByteOffsetPlusBytesFilled()) { |
| 2338 | } |
| 2339 | |
| 2340 | void ReadableStreamBYOBRequest::Impl::updateView(jsg::Lock& js) { |
| 2341 | jsg::check(view.getHandle(js)->Buffer()->Detach(v8::Local<v8::Value>())); |
| 2342 | view = js.v8Ref(readRequest->getView(js)); |
| 2343 | } |
| 2344 | |
| 2345 | void ReadableStreamBYOBRequest::visitForGc(jsg::GcVisitor& visitor) { |
| 2346 | KJ_IF_SOME(impl, maybeImpl) { |
| 2347 | visitor.visit(impl.view); |
| 2348 | } |
| 2349 | } |
| 2350 | |
| 2351 | ReadableStreamBYOBRequest::ReadableStreamBYOBRequest(jsg::Lock& js, |
| 2352 | kj::Own<ByteQueue::ByobRequest> readRequest, |
| 2353 | kj::Rc<WeakRef<ReadableByteStreamController>> controller) |
| 2354 | : ioContext(tryGetIoContext()), |
| 2355 | maybeImpl(Impl(js, kj::mv(readRequest), kj::mv(controller))) {} |
| 2356 | |
| 2357 | kj::Maybe<int> ReadableStreamBYOBRequest::getAtLeast() { |
| 2358 | KJ_IF_SOME(impl, maybeImpl) { |
| 2359 | return impl.readRequest->getAtLeast(); |
| 2360 | } |
| 2361 | return kj::none; |
| 2362 | } |
| 2363 | |
| 2364 | kj::Maybe<jsg::V8Ref<v8::Uint8Array>> ReadableStreamBYOBRequest::getView(jsg::Lock& js) { |
| 2365 | KJ_IF_SOME(impl, maybeImpl) { |
| 2366 | return impl.view.addRef(js); |
| 2367 | } |
| 2368 | return kj::none; |
| 2369 | } |
| 2370 | |
| 2371 | void ReadableStreamBYOBRequest::invalidate(jsg::Lock& js) { |
| 2372 | KJ_IF_SOME(impl, maybeImpl) { |
| 2373 | // If the user code happened to have retained a reference to the view or |
| 2374 | // the buffer, we need to detach it so that those references cannot be used |
| 2375 | // to modify or observe modifications. |
| 2376 | jsg::check(impl.view.getHandle(js)->Buffer()->Detach(v8::Local<v8::Value>())); |
| 2377 | impl.controller->runIfAlive( |
| 2378 | [](ReadableByteStreamController& controller) { controller.maybeByobRequest = kj::none; }); |
| 2379 | } |
| 2380 | maybeImpl = kj::none; |
| 2381 | } |
| 2382 | |
| 2383 | void ReadableStreamBYOBRequest::respond(jsg::Lock& js, int bytesWritten) { |
| 2384 | auto& impl = JSG_REQUIRE_NONNULL( |
| 2385 | maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated."); |
| 2386 | JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid."); |
| 2387 | JSG_REQUIRE(impl.view.getHandle(js)->ByteLength() > 0, TypeError, |
| 2388 | "Cannot respond with a zero-length or detached view"); |
| 2389 | impl.controller->runIfAlive([&](ReadableByteStreamController& controller) { |
| 2390 | if (!controller.canCloseOrEnqueue()) { |
| 2391 | JSG_REQUIRE(bytesWritten == 0, TypeError, |
| 2392 | "The bytesWritten must be zero after the stream is closed."); |
| 2393 | KJ_ASSERT(impl.readRequest->isInvalidated()); |
| 2394 | invalidate(js); |
| 2395 | } else { |
| 2396 | bool shouldInvalidate = false; |
| 2397 | if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) { |
| 2398 | // While this particular request may be invalidated, there are still |
| 2399 | // other branches we can push the data to. Let's do so. |
| 2400 | jsg::BufferSource source(js, impl.view.getHandle(js)); |
| 2401 | auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, source.detach(js))); |
| 2402 | controller.impl.enqueue(js, kj::mv(entry), controller.getSelf()); |
| 2403 | } else { |
| 2404 | JSG_REQUIRE(bytesWritten > 0, TypeError, |
| 2405 | "The bytesWritten must be more than zero while the stream is open."); |
| 2406 | if (impl.readRequest->respond(js, bytesWritten)) { |
| 2407 | // The read request was fulfilled, we need to invalidate. |
| 2408 | shouldInvalidate = true; |
| 2409 | } else { |
| 2410 | // The response did not fulfill the minimum requirements of the read. |
| 2411 | // We do not want to invalidate the read request and we need to update the |
| 2412 | // view so that on the next read the view will be properly adjusted. |
| 2413 | impl.updateView(js); |
| 2414 | } |
| 2415 | } |
| 2416 | controller.pull(js); |
| 2417 | if (shouldInvalidate) { |
| 2418 | invalidate(js); |
| 2419 | } |
| 2420 | } |
| 2421 | }); |
| 2422 | } |
| 2423 | |
| 2424 | void ReadableStreamBYOBRequest::respondWithNewView(jsg::Lock& js, jsg::BufferSource view) { |
| 2425 | auto& impl = JSG_REQUIRE_NONNULL( |
| 2426 | maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated."); |
| 2427 | JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid."); |
| 2428 | impl.controller->runIfAlive([&](ReadableByteStreamController& controller) { |
| 2429 | if (!controller.canCloseOrEnqueue()) { |
| 2430 | JSG_REQUIRE(view.size() == 0, TypeError, |
| 2431 | "The view byte length must be zero after the stream is closed."); |
| 2432 | |
| 2433 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 2434 | // Per the spec, when the stream is closed: |
| 2435 | // 1. The view byte length must be zero (TypeError if not) |
| 2436 | // 2. The underlying buffer must not be detached (TypeError) |
| 2437 | // 3. The buffer byte length must not be zero (RangeError) |
| 2438 | // 4. The buffer byte length must match the original (RangeError) |
| 2439 | auto handle = view.getHandle(js); |
| 2440 | auto buffer = handle->IsArrayBuffer() ? handle.As<v8::ArrayBuffer>() |
| 2441 | : handle.As<v8::ArrayBufferView>()->Buffer(); |
| 2442 | JSG_REQUIRE( |
| 2443 | !buffer->WasDetached(), TypeError, "The underlying ArrayBuffer has been detached."); |
| 2444 | |
| 2445 | JSG_REQUIRE(view.canDetach(js), TypeError, "Unable to use non-detachable ArrayBuffer."); |
| 2446 | // Use the stored values since the ByobRequest may have been invalidated during close. |
| 2447 | auto actualBufferByteLength = buffer->ByteLength(); |
| 2448 | JSG_REQUIRE( |
| 2449 | actualBufferByteLength != 0, RangeError, "The underlying ArrayBuffer is zero-length."); |
| 2450 | JSG_REQUIRE(actualBufferByteLength == impl.originalBufferByteLength, RangeError, |
| 2451 | "The underlying ArrayBuffer is not the correct length."); |
| 2452 | // The view's byte offset must match the original byte offset plus bytes filled. |
| 2453 | auto viewByteOffset = |
| 2454 | handle->IsArrayBuffer() ? 0 : handle.As<v8::ArrayBufferView>()->ByteOffset(); |
| 2455 | JSG_REQUIRE(viewByteOffset == impl.originalByteOffsetPlusBytesFilled, RangeError, |
| 2456 | "The view has an invalid byte offset."); |
| 2457 | } else { |
| 2458 | KJ_ASSERT(impl.readRequest->isInvalidated()); |
| 2459 | } |
| 2460 | |
| 2461 | invalidate(js); |
| 2462 | } else { |
| 2463 | bool shouldInvalidate = false; |
| 2464 | if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) { |
| 2465 | // While this particular request may be invalidated, there are still |
| 2466 | // other branches we can push the data to. Let's do so. |
| 2467 | auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, view.detach(js))); |
| 2468 | controller.impl.enqueue(js, kj::mv(entry), controller.getSelf()); |
| 2469 | } else { |
| 2470 | JSG_REQUIRE(view.size() > 0, TypeError, |
| 2471 | "The view byte length must be more than zero while the stream is open."); |
| 2472 | if (impl.readRequest->respondWithNewView(js, kj::mv(view))) { |
| 2473 | // The read request was fulfilled, we need to invalidate. |
| 2474 | shouldInvalidate = true; |
| 2475 | } else { |
| 2476 | // The response did not fulfill the minimum requirements of the read. |
| 2477 | // We do not want to invalidate the read request and we need to update the |
| 2478 | // view so that on the next read the view will be properly adjusted. |
| 2479 | impl.updateView(js); |
| 2480 | } |
| 2481 | } |
| 2482 | |
| 2483 | controller.pull(js); |
| 2484 | if (shouldInvalidate) { |
| 2485 | invalidate(js); |
| 2486 | } |
| 2487 | } |
| 2488 | }); |
| 2489 | } |
| 2490 | |
| 2491 | bool ReadableStreamBYOBRequest::isPartiallyFulfilled() { |
| 2492 | KJ_IF_SOME(impl, maybeImpl) { |
| 2493 | return impl.readRequest->isPartiallyFulfilled(); |
| 2494 | } |
| 2495 | return false; |
| 2496 | } |
| 2497 | |
| 2498 | // ====================================================================================== |
| 2499 | |
| 2500 | ReadableByteStreamController::ReadableByteStreamController( |
| 2501 | UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) |
| 2502 | : weakSelf(kj::rc<WeakRef<ReadableByteStreamController>>( |
| 2503 | kj::Badge<ReadableByteStreamController>{}, *this)), |
| 2504 | ioContext(tryGetIoContext()), |
| 2505 | impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} |
| 2506 | |
| 2507 | ReadableByteStreamController::~ReadableByteStreamController() noexcept(false) { |
| 2508 | weakSelf->invalidate(); |
| 2509 | } |
| 2510 | |
| 2511 | void ReadableByteStreamController::start(jsg::Lock& js) { |
| 2512 | impl.start(js, JSG_THIS); |
| 2513 | } |
| 2514 | |
| 2515 | bool ReadableByteStreamController::canCloseOrEnqueue() { |
| 2516 | return impl.canCloseOrEnqueue(); |
| 2517 | } |
| 2518 | |
| 2519 | bool ReadableByteStreamController::hasBackpressure() { |
| 2520 | return !impl.shouldCallPull(); |
| 2521 | } |
| 2522 | |
| 2523 | kj::Maybe<int> ReadableByteStreamController::getDesiredSize() { |
| 2524 | return impl.getDesiredSize(); |
| 2525 | } |
| 2526 | |
| 2527 | void ReadableByteStreamController::visitForGc(jsg::GcVisitor& visitor) { |
| 2528 | visitor.visit(maybeByobRequest, impl); |
| 2529 | } |
| 2530 | |
| 2531 | jsg::Promise<void> ReadableByteStreamController::cancel( |
| 2532 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 2533 | KJ_IF_SOME(byobRequest, maybeByobRequest) { |
| 2534 | if (impl.consumerCount() == 1) { |
| 2535 | byobRequest->invalidate(js); |
| 2536 | } |
| 2537 | } |
| 2538 | return impl.cancel(js, JSG_THIS, maybeReason.orDefault(js.undefined())); |
| 2539 | } |
| 2540 | |
| 2541 | void ReadableByteStreamController::close(jsg::Lock& js) { |
| 2542 | KJ_IF_SOME(byobRequest, maybeByobRequest) { |
| 2543 | JSG_REQUIRE(!byobRequest->isPartiallyFulfilled(), TypeError, |
| 2544 | "This ReadableStream was closed with a partial read pending."); |
| 2545 | } else if (FeatureFlags::get(js).getPedanticWpt()) { |
| 2546 | // If maybeByobRequest is not set, check if there's a pending byob request. |
| 2547 | // If so, materialize it before closing so it remains accessible after |
| 2548 | // the state changes to Closed. This is required by the spec for proper |
| 2549 | // respondWithNewView() error handling in the closed state. |
| 2550 | // Only do this if the queue doesn't have a partially fulfilled read. |
| 2551 | KJ_IF_SOME(queue, impl.state.tryGetUnsafe<ByteQueue>()) { |
| 2552 | if (!queue.hasPartiallyFulfilledRead()) { |
| 2553 | getByobRequest(js); |
| 2554 | } |
| 2555 | } |
| 2556 | } |
| 2557 | impl.close(js); |
| 2558 | } |
| 2559 | |
| 2560 | void ReadableByteStreamController::enqueue(jsg::Lock& js, jsg::BufferSource chunk) { |
| 2561 | // Hold a strong reference up front. Operations below (invalidate, detach) touch |
| 2562 | // the JS heap and C++ argument evaluation order is unspecified, so JSG_THIS as a |
| 2563 | // function argument would not reliably precede chunk.detach(js). |
| 2564 | auto self = JSG_THIS; |
| 2565 | |
| 2566 | JSG_REQUIRE(chunk.size() > 0, TypeError, "Cannot enqueue a zero-length ArrayBuffer."); |
| 2567 | JSG_REQUIRE(chunk.canDetach(js), TypeError, "The provided ArrayBuffer must be detachable."); |
| 2568 | JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "This ReadableByteStreamController is closed."); |
| 2569 | |
| 2570 | KJ_IF_SOME(byobRequest, maybeByobRequest) { |
| 2571 | KJ_IF_SOME(view, byobRequest->getView(js)) { |
| 2572 | JSG_REQUIRE(view.getHandle(js)->ByteLength() > 0, TypeError, |
| 2573 | "The byobRequest.view is zero-length or was detached"); |
| 2574 | } |
| 2575 | byobRequest->invalidate(js); |
| 2576 | } |
| 2577 | |
| 2578 | impl.enqueue(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, chunk.detach(js))), kj::mv(self)); |
| 2579 | } |
| 2580 | |
| 2581 | void ReadableByteStreamController::error(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2582 | impl.doError(js, js.v8Ref(reason)); |
| 2583 | } |
| 2584 | |
| 2585 | kj::Maybe<jsg::Ref<ReadableStreamBYOBRequest>> ReadableByteStreamController::getByobRequest( |
| 2586 | jsg::Lock& js) { |
| 2587 | if (maybeByobRequest == kj::none) { |
| 2588 | KJ_IF_SOME(queue, impl.state.tryGetUnsafe<ByteQueue>()) { |
| 2589 | KJ_IF_SOME(pendingByob, queue.nextPendingByobReadRequest()) { |
| 2590 | maybeByobRequest = |
| 2591 | js.alloc<ReadableStreamBYOBRequest>(js, kj::mv(pendingByob), weakSelf.addRef()); |
| 2592 | } |
| 2593 | } else { |
| 2594 | return kj::none; |
| 2595 | } |
| 2596 | } |
| 2597 | |
| 2598 | return maybeByobRequest.map( |
| 2599 | [&](jsg::Ref<ReadableStreamBYOBRequest>& req) { return req.addRef(); }); |
| 2600 | } |
| 2601 | |
| 2602 | // When a consumer receives a read request, but does not have the data available to |
| 2603 | // fulfill the request, the consumer will call pull on the controller to pull that |
| 2604 | // data if needed. |
| 2605 | void ReadableByteStreamController::pull(jsg::Lock& js) { |
| 2606 | impl.pullIfNeeded(js, JSG_THIS); |
| 2607 | } |
| 2608 | |
| 2609 | void ReadableByteStreamController::forcePull(jsg::Lock& js) { |
| 2610 | impl.forcePullIfNeeded(js, JSG_THIS); |
| 2611 | } |
| 2612 | |
| 2613 | kj::Own<ByteQueue::Consumer> ReadableByteStreamController::getConsumer( |
| 2614 | kj::Maybe<ByteQueue::ConsumerImpl::StateListener&> stateListener) { |
| 2615 | return impl.getConsumer(stateListener); |
| 2616 | } |
| 2617 | |
| 2618 | // ====================================================================================== |
| 2619 | |
| 2620 | ReadableStreamJsController::ReadableStreamJsController(): ioContext(tryGetIoContext()) {} |
| 2621 | |
| 2622 | ReadableStreamJsController::ReadableStreamJsController(StreamStates::Closed closed) |
| 2623 | : ioContext(tryGetIoContext()) { |
| 2624 | state.transitionTo<StreamStates::Closed>(); |
| 2625 | } |
| 2626 | |
| 2627 | ReadableStreamJsController::ReadableStreamJsController(StreamStates::Errored errored) |
| 2628 | : ioContext(tryGetIoContext()) { |
| 2629 | state.transitionTo<StreamStates::Errored>(kj::mv(errored)); |
| 2630 | } |
| 2631 | |
| 2632 | ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer) |
| 2633 | : ioContext(tryGetIoContext()) { |
| 2634 | state.transitionTo<kj::Own<ValueReadable>>(consumer.clone(js, *this)); |
| 2635 | } |
| 2636 | |
| 2637 | ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer) |
| 2638 | : ioContext(tryGetIoContext()) { |
| 2639 | state.transitionTo<kj::Own<ByteReadable>>(consumer.clone(js, *this)); |
| 2640 | } |
| 2641 | |
| 2642 | jsg::Ref<ReadableStream> ReadableStreamJsController::addRef() { |
| 2643 | return KJ_REQUIRE_NONNULL(owner).addRef(); |
| 2644 | } |
| 2645 | |
| 2646 | jsg::Promise<void> ReadableStreamJsController::cancel( |
| 2647 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 2648 | disturbed = true; |
| 2649 | |
| 2650 | const auto doCancel = [&](auto& consumer) { |
| 2651 | auto reason = js.v8Ref(maybeReason.orDefault([&] { return js.v8Undefined(); })); |
| 2652 | KJ_DEFER(doClose(js)); |
| 2653 | return consumer->cancel(js, reason.getHandle(js)); |
| 2654 | }; |
| 2655 | |
| 2656 | // Check for pending state first (deferred close/error during a read operation) |
| 2657 | if (state.pendingStateIs<StreamStates::Closed>()) { |
| 2658 | return js.resolvedPromise(); |
| 2659 | } |
| 2660 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 2661 | return js.rejectedPromise<void>(pendingError.addRef(js)); |
| 2662 | } |
| 2663 | |
| 2664 | KJ_SWITCH_ONEOF(state) { |
| 2665 | KJ_CASE_ONEOF(initial, Initial) { |
| 2666 | // Stream not yet set up, treat as closed. |
| 2667 | return js.resolvedPromise(); |
| 2668 | } |
| 2669 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2670 | return js.resolvedPromise(); |
| 2671 | } |
| 2672 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2673 | return js.rejectedPromise<void>(errored.addRef(js)); |
| 2674 | } |
| 2675 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 2676 | if (canceling) return js.resolvedPromise(); |
| 2677 | canceling = true; |
| 2678 | return doCancel(consumer); |
| 2679 | } |
| 2680 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 2681 | if (canceling) return js.resolvedPromise(); |
| 2682 | canceling = true; |
| 2683 | return doCancel(consumer); |
| 2684 | } |
| 2685 | } |
| 2686 | |
| 2687 | KJ_UNREACHABLE; |
| 2688 | } |
| 2689 | |
| 2690 | // Finalizes the closed state of this ReadableStream. The connection to the underlying |
| 2691 | // controller is released with no further action. Importantly, this method is triggered |
| 2692 | // by the underlying controller as a result of that controller closing or being canceled. |
| 2693 | // We detach ourselves from the underlying controller by releasing the ValueReadable or |
| 2694 | // ByteReadable in the state and changing that to closed. |
| 2695 | // We also clean up other state here. |
| 2696 | void ReadableStreamJsController::doClose(jsg::Lock& js) { |
| 2697 | // If already in a terminal state, nothing to do. |
| 2698 | if (state.isTerminal()) return; |
| 2699 | |
| 2700 | // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately. |
| 2701 | // Returns true if transition happened immediately. |
| 2702 | if (state.deferTransitionTo<StreamStates::Closed>()) { |
| 2703 | lock.onClose(js); |
| 2704 | } |
| 2705 | // If deferred, lock.onClose will be called when the pending state is applied |
| 2706 | // via applyPendingState in deferControllerStateChange. |
| 2707 | } |
| 2708 | |
| 2709 | // As with doClose(), doError() finalizes the error state of this ReadableStream. |
| 2710 | // The connection to the underlying controller is released with no further action. |
| 2711 | // This method is triggered by the underlying controller as a result of that controller |
| 2712 | // erroring. We detach ourselves from the underlying controller by releasing the ValueReadable |
| 2713 | // or ByteReadable in the state and changing that to errored. |
| 2714 | // We also clean up other state here. |
| 2715 | void ReadableStreamJsController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2716 | // If already in a terminal state, nothing to do. |
| 2717 | if (state.isTerminal()) return; |
| 2718 | |
| 2719 | // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately. |
| 2720 | // Returns true if transition happened immediately. |
| 2721 | if (state.deferTransitionTo<StreamStates::Errored>(js.v8Ref(reason))) { |
| 2722 | lock.onError(js, reason); |
| 2723 | } |
| 2724 | // If deferred, lock.onError will be called when the pending state is applied |
| 2725 | // via applyPendingState in deferControllerStateChange. |
| 2726 | } |
| 2727 | |
| 2728 | bool ReadableStreamJsController::isByteOriented() const { |
| 2729 | return state.is<kj::Own<ByteReadable>>(); |
| 2730 | } |
| 2731 | |
| 2732 | bool ReadableStreamJsController::isClosedOrErrored() const { |
| 2733 | // Check if we're in a terminal state or have one pending |
| 2734 | return state.isTerminal() || state.hasPendingState(); |
| 2735 | } |
| 2736 | |
| 2737 | bool ReadableStreamJsController::isClosed() const { |
| 2738 | // Check current state first, then pending state |
| 2739 | if (state.is<StreamStates::Closed>()) return true; |
| 2740 | return state.pendingStateIs<StreamStates::Closed>(); |
| 2741 | } |
| 2742 | |
| 2743 | bool ReadableStreamJsController::isDisturbed() { |
| 2744 | return disturbed; |
| 2745 | } |
| 2746 | |
| 2747 | bool ReadableStreamJsController::isLockedToReader() const { |
| 2748 | return lock.isLockedToReader(); |
| 2749 | } |
| 2750 | |
| 2751 | bool ReadableStreamJsController::lockReader(jsg::Lock& js, Reader& reader) { |
| 2752 | return lock.lockReader(js, *this, reader); |
| 2753 | } |
| 2754 | |
| 2755 | jsg::Promise<void> ReadableStreamJsController::pipeTo( |
| 2756 | jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) { |
| 2757 | KJ_DASSERT(!isLockedToReader()); |
| 2758 | KJ_DASSERT(!destination.isLockedToWriter()); |
| 2759 | |
| 2760 | disturbed = true; |
| 2761 | KJ_IF_SOME(promise, destination.tryPipeFrom(js, addRef(), kj::mv(options))) { |
| 2762 | return kj::mv(promise); |
| 2763 | } |
| 2764 | |
| 2765 | return js.rejectedPromise<void>( |
| 2766 | js.v8TypeError("This ReadableStream cannot be piped to this WritableStream"_kj)); |
| 2767 | } |
| 2768 | |
| 2769 | kj::Maybe<jsg::Promise<ReadResult>> ReadableStreamJsController::read( |
| 2770 | jsg::Lock& js, kj::Maybe<ByobOptions> maybeByobOptions) { |
| 2771 | disturbed = true; |
| 2772 | |
| 2773 | KJ_IF_SOME(byobOptions, maybeByobOptions) { |
| 2774 | byobOptions.detachBuffer = true; |
| 2775 | auto view = byobOptions.bufferView.getHandle(js); |
| 2776 | if (!view->Buffer()->IsDetachable()) { |
| 2777 | return js.rejectedPromise<ReadResult>( |
| 2778 | js.v8TypeError("Unabled to use non-detachable ArrayBuffer."_kj)); |
| 2779 | } |
| 2780 | |
| 2781 | if (view->ByteLength() == 0 || view->Buffer()->ByteLength() == 0) { |
| 2782 | return js.rejectedPromise<ReadResult>( |
| 2783 | js.v8TypeError("Unable to use a zero-length ArrayBuffer."_kj)); |
| 2784 | } |
| 2785 | |
| 2786 | // Check for pending error first (deferred error during a prior read operation) |
| 2787 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 2788 | return js.rejectedPromise<ReadResult>(pendingError.addRef(js)); |
| 2789 | } |
| 2790 | |
| 2791 | if (state.is<StreamStates::Closed>() || state.pendingStateIs<StreamStates::Closed>()) { |
| 2792 | // If it is a BYOB read, then the spec requires that we return an empty |
| 2793 | // view of the same type provided, that uses the same backing memory |
| 2794 | // as that provided, but with zero-length. |
| 2795 | auto source = jsg::BufferSource(js, byobOptions.bufferView.getHandle(js)); |
| 2796 | auto store = source.detach(js); |
| 2797 | store.consume(store.size()); |
| 2798 | return js.resolvedPromise(ReadResult{ |
| 2799 | .value = js.v8Ref(store.createHandle(js)), |
| 2800 | .done = true, |
| 2801 | }); |
| 2802 | } |
| 2803 | } |
| 2804 | |
| 2805 | // Check for pending state (deferred close/error during a prior read operation) |
| 2806 | if (state.pendingStateIs<StreamStates::Closed>()) { |
| 2807 | // The closed state for BYOB reads is handled in the maybeByobOptions check above. |
| 2808 | KJ_ASSERT(maybeByobOptions == kj::none); |
| 2809 | return js.resolvedPromise(ReadResult{.done = true}); |
| 2810 | } |
| 2811 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 2812 | return js.rejectedPromise<ReadResult>(pendingError.addRef(js)); |
| 2813 | } |
| 2814 | |
| 2815 | KJ_SWITCH_ONEOF(state) { |
| 2816 | KJ_CASE_ONEOF(initial, Initial) { |
| 2817 | // Stream not yet set up, treat as closed. |
| 2818 | KJ_ASSERT(maybeByobOptions == kj::none); |
| 2819 | return js.resolvedPromise(ReadResult{.done = true}); |
| 2820 | } |
| 2821 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2822 | // The closed state for BYOB reads is handled in the maybeByobOptions check above. |
| 2823 | KJ_ASSERT(maybeByobOptions == kj::none); |
| 2824 | return js.resolvedPromise(ReadResult{.done = true}); |
| 2825 | } |
| 2826 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2827 | return js.rejectedPromise<ReadResult>(errored.addRef(js)); |
| 2828 | } |
| 2829 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 2830 | // The ReadableStreamDefaultController does not support ByobOptions. |
| 2831 | // It should never happen, but let's make sure. |
| 2832 | KJ_ASSERT(maybeByobOptions == kj::none); |
| 2833 | return deferControllerStateChange(js, *this, [&]() mutable { return consumer->read(js); }); |
| 2834 | } |
| 2835 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 2836 | return deferControllerStateChange( |
| 2837 | js, *this, [&]() mutable { return consumer->read(js, kj::mv(maybeByobOptions)); }); |
| 2838 | } |
| 2839 | } |
| 2840 | KJ_UNREACHABLE; |
| 2841 | } |
| 2842 | |
| 2843 | kj::Maybe<jsg::Promise<DrainingReadResult>> ReadableStreamJsController::drainingRead( |
| 2844 | jsg::Lock& js, size_t maxRead) { |
| 2845 | disturbed = true; |
| 2846 | |
| 2847 | // Check for pending state first (deferred close/error during a prior read operation) |
| 2848 | if (state.pendingStateIs<StreamStates::Closed>()) { |
| 2849 | return js.resolvedPromise(DrainingReadResult{ |
| 2850 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 2851 | .done = true, |
| 2852 | }); |
| 2853 | } |
| 2854 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 2855 | return js.rejectedPromise<DrainingReadResult>(pendingError.addRef(js)); |
| 2856 | } |
| 2857 | |
| 2858 | // Like deferControllerStateChange for regular reads, we need to prevent the controller |
| 2859 | // state from being destroyed while a draining read's promise callbacks are pending. |
| 2860 | // The drainingRead implementation captures `this` (the Consumer) in promise lambdas to |
| 2861 | // clear hasPendingDrainingRead. If the state is changed (destroying the Consumer) before |
| 2862 | // those callbacks run, we get a use-after-free. |
| 2863 | // |
| 2864 | // CRITICAL: state.beginOperation() MUST be called BEFORE consumer->drainingRead(), not |
| 2865 | // after. The consumer->drainingRead() call may trigger onConsumerWantsData -> forcePull |
| 2866 | // -> close/error, which calls deferTransitionTo. If no operation is in progress at that |
| 2867 | // point, the transition fires immediately, destroying the Consumer while we're still |
| 2868 | // inside its method and before the returned promise's .then() callbacks are set up. |
| 2869 | // The endOperation() happens in the .then() callbacks below, ensuring the deferred |
| 2870 | // state change only fires after the promise resolves/rejects and the Consumer's |
| 2871 | // this-capturing callbacks have already run. |
| 2872 | auto wrapDrainingRead = |
| 2873 | [this](jsg::Lock& js, |
| 2874 | jsg::Promise<DrainingReadResult> promise) -> jsg::Promise<DrainingReadResult> { |
| 2875 | return promise.then(js, [this](jsg::Lock& js, DrainingReadResult result) { |
| 2876 | if (state.endOperation()) { |
| 2877 | // A pending state was applied. Call the appropriate callback. |
| 2878 | if (state.template is<StreamStates::Closed>()) { |
| 2879 | lock.onClose(js); |
| 2880 | } else if (state.template is<StreamStates::Errored>()) { |
| 2881 | KJ_IF_SOME(err, state.template tryGetUnsafe<StreamStates::Errored>()) { |
| 2882 | lock.onError(js, err.getHandle(js)); |
| 2883 | // The error was applied during this operation โ the data we collected |
| 2884 | // may be invalid. Discard it and propagate the error rather than |
| 2885 | // silently returning possibly-corrupt data. |
| 2886 | js.throwException(err.addRef(js)); |
| 2887 | } |
| 2888 | } |
| 2889 | } |
| 2890 | return kj::mv(result); |
| 2891 | }, [this](jsg::Lock& js, jsg::Value exception) -> DrainingReadResult { |
| 2892 | state.clearPendingState(); |
| 2893 | (void)state.endOperation(); |
| 2894 | js.throwException(kj::mv(exception)); |
| 2895 | }); |
| 2896 | }; |
| 2897 | |
| 2898 | KJ_SWITCH_ONEOF(state) { |
| 2899 | KJ_CASE_ONEOF(initial, Initial) { |
| 2900 | // Stream not yet set up, treat as closed. |
| 2901 | return js.resolvedPromise(DrainingReadResult{ |
| 2902 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 2903 | .done = true, |
| 2904 | }); |
| 2905 | } |
| 2906 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2907 | return js.resolvedPromise(DrainingReadResult{ |
| 2908 | .chunks = kj::Array<kj::Array<kj::byte>>(), |
| 2909 | .done = true, |
| 2910 | }); |
| 2911 | } |
| 2912 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2913 | return js.rejectedPromise<DrainingReadResult>(errored.addRef(js)); |
| 2914 | } |
| 2915 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 2916 | // beginOperation MUST be before consumer->drainingRead() โ see comment above. |
| 2917 | state.beginOperation(); |
| 2918 | JSG_TRY(js) { |
| 2919 | return wrapDrainingRead(js, consumer->drainingRead(js, maxRead)); |
| 2920 | } |
| 2921 | JSG_CATCH(exception) { |
| 2922 | state.clearPendingState(); |
| 2923 | (void)state.endOperation(); |
| 2924 | doError(js, exception.getHandle(js)); |
| 2925 | return js.rejectedPromise<DrainingReadResult>(kj::mv(exception)); |
| 2926 | }; |
| 2927 | } |
| 2928 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 2929 | // beginOperation MUST be before consumer->drainingRead() โ see comment above. |
| 2930 | state.beginOperation(); |
| 2931 | JSG_TRY(js) { |
| 2932 | return wrapDrainingRead(js, consumer->drainingRead(js, maxRead)); |
| 2933 | } |
| 2934 | JSG_CATCH(exception) { |
| 2935 | state.clearPendingState(); |
| 2936 | (void)state.endOperation(); |
| 2937 | doError(js, exception.getHandle(js)); |
| 2938 | return js.rejectedPromise<DrainingReadResult>(kj::mv(exception)); |
| 2939 | }; |
| 2940 | } |
| 2941 | } |
| 2942 | KJ_UNREACHABLE; |
| 2943 | } |
| 2944 | |
| 2945 | void ReadableStreamJsController::releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) { |
| 2946 | lock.releaseReader(*this, reader, maybeJs); |
| 2947 | } |
| 2948 | |
| 2949 | ReadableStreamController::Tee ReadableStreamJsController::tee(jsg::Lock& js) { |
| 2950 | JSG_REQUIRE(!isLockedToReader(), TypeError, "This ReadableStream is locked to a reader."); |
| 2951 | lock.state.transitionTo<Locked>(); |
| 2952 | disturbed = true; |
| 2953 | |
| 2954 | // This will leave this stream locked, disturbed, and closed. |
| 2955 | |
| 2956 | // Check for pending state first (deferred close/error during a prior read operation) |
| 2957 | if (state.pendingStateIs<StreamStates::Closed>()) { |
| 2958 | return Tee{ |
| 2959 | .branch1 = |
| 2960 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2961 | .branch2 = |
| 2962 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2963 | }; |
| 2964 | } |
| 2965 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 2966 | return Tee{ |
| 2967 | .branch1 = |
| 2968 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(pendingError.addRef(js))), |
| 2969 | .branch2 = |
| 2970 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(pendingError.addRef(js))), |
| 2971 | }; |
| 2972 | } |
| 2973 | |
| 2974 | KJ_SWITCH_ONEOF(state) { |
| 2975 | KJ_CASE_ONEOF(initial, Initial) { |
| 2976 | // Stream not yet set up, treat as closed. |
| 2977 | return Tee{ |
| 2978 | .branch1 = |
| 2979 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2980 | .branch2 = |
| 2981 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2982 | }; |
| 2983 | } |
| 2984 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2985 | return Tee{ |
| 2986 | .branch1 = |
| 2987 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2988 | .branch2 = |
| 2989 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())), |
| 2990 | }; |
| 2991 | } |
| 2992 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2993 | return Tee{ |
| 2994 | .branch1 = |
| 2995 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(errored.addRef(js))), |
| 2996 | .branch2 = |
| 2997 | js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(errored.addRef(js))), |
| 2998 | }; |
| 2999 | } |
| 3000 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 3001 | KJ_DEFER(state.transitionTo<StreamStates::Closed>()); |
| 3002 | // We create two additional streams that clone this stream's consumer state, |
| 3003 | // then close this stream's consumer. |
| 3004 | return Tee{ |
| 3005 | .branch1 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)), |
| 3006 | .branch2 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)), |
| 3007 | }; |
| 3008 | } |
| 3009 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 3010 | KJ_DEFER(state.transitionTo<StreamStates::Closed>()); |
| 3011 | // We create two additional streams that clone this stream's consumer state, |
| 3012 | // then close this stream's consumer. |
| 3013 | return Tee{ |
| 3014 | .branch1 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)), |
| 3015 | .branch2 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)), |
| 3016 | }; |
| 3017 | } |
| 3018 | } |
| 3019 | KJ_UNREACHABLE; |
| 3020 | } |
| 3021 | |
| 3022 | void ReadableStreamJsController::setOwnerRef(ReadableStream& stream) { |
| 3023 | KJ_ASSERT(owner == kj::none); |
| 3024 | owner = &stream; |
| 3025 | } |
| 3026 | |
| 3027 | void ReadableStreamJsController::setup(jsg::Lock& js, |
| 3028 | jsg::Optional<UnderlyingSource> maybeUnderlyingSource, |
| 3029 | jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) { |
| 3030 | auto underlyingSource = kj::mv(maybeUnderlyingSource).orDefault({}); |
| 3031 | auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({}); |
| 3032 | auto type = underlyingSource.type.map([](kj::StringPtr s) { return s; }).orDefault(""_kj); |
| 3033 | |
| 3034 | expectedLength = underlyingSource.expectedLength; |
| 3035 | |
| 3036 | if (type == "bytes") { |
| 3037 | // Per spec, autoAllocateChunkSize should only be set if the user explicitly provides it. |
| 3038 | // If not set, the underlying source's pull method won't receive a byobRequest for |
| 3039 | // non-BYOB reads and must use controller.enqueue() instead. |
| 3040 | // |
| 3041 | // However, our original implementation always defaulted to 4096, so we need a compat flag |
| 3042 | // to control this behavior. Default to legacy behavior if flags aren't available. |
| 3043 | bool useSpecCompliantBehavior = false; |
| 3044 | KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) { |
| 3045 | useSpecCompliantBehavior = flags.getNoAutoAllocateChunkSize(); |
| 3046 | } |
| 3047 | |
| 3048 | kj::Maybe<int> autoAllocateChunkSize; |
| 3049 | if (useSpecCompliantBehavior) { |
| 3050 | // Spec-compliant: only set if user explicitly provides it |
| 3051 | autoAllocateChunkSize = |
| 3052 | underlyingSource.autoAllocateChunkSize.map([](int size) { return size; }); |
| 3053 | } else { |
| 3054 | // Legacy behavior: apply a default autoAllocateChunkSize if not provided. |
| 3055 | auto defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE; |
| 3056 | if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) { |
| 3057 | defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2; |
| 3058 | } |
| 3059 | autoAllocateChunkSize = underlyingSource.autoAllocateChunkSize.orDefault(defaultChunkSize); |
| 3060 | } |
| 3061 | |
| 3062 | auto controller = |
| 3063 | js.alloc<ReadableByteStreamController>(kj::mv(underlyingSource), kj::mv(queuingStrategy)); |
| 3064 | |
| 3065 | KJ_IF_SOME(chunkSize, autoAllocateChunkSize) { |
| 3066 | JSG_REQUIRE(chunkSize > 0, TypeError, "The autoAllocateChunkSize option cannot be zero."); |
| 3067 | } |
| 3068 | |
| 3069 | // We account for the memory usage of the ByteReadable and its controller together because |
| 3070 | // their lifetimes are identical (in practice) and memory accounting itself has a memory |
| 3071 | // overhead. The same applies to ValueReadable below. |
| 3072 | state.transitionTo<kj::Own<ByteReadable>>( |
| 3073 | kj::heap<ByteReadable>(controller.addRef(), *this, autoAllocateChunkSize) |
| 3074 | .attach(js.getExternalMemoryAdjustment( |
| 3075 | sizeof(ByteReadable) + sizeof(ReadableByteStreamController)))); |
| 3076 | controller->start(js); |
| 3077 | } else { |
| 3078 | JSG_REQUIRE( |
| 3079 | type == "", TypeError, kj::str("\"", type, "\" is not a valid type of ReadableStream.")); |
| 3080 | auto controller = js.alloc<ReadableStreamDefaultController>( |
| 3081 | kj::mv(underlyingSource), kj::mv(queuingStrategy)); |
| 3082 | state.transitionTo<kj::Own<ValueReadable>>( |
| 3083 | kj::heap<ValueReadable>(controller.addRef(), *this) |
| 3084 | .attach(js.getExternalMemoryAdjustment( |
| 3085 | sizeof(ValueReadable) + sizeof(ReadableStreamDefaultController)))); |
| 3086 | controller->start(js); |
| 3087 | } |
| 3088 | } |
| 3089 | |
| 3090 | kj::Maybe<ReadableStreamController::PipeController&> ReadableStreamJsController::tryPipeLock() { |
| 3091 | return lock.tryPipeLock(*this); |
| 3092 | } |
| 3093 | |
| 3094 | void ReadableStreamJsController::visitForGc(jsg::GcVisitor& visitor) { |
| 3095 | // Visit pending state if it's an error (Closed has no GC-traceable content) |
| 3096 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 3097 | visitor.visit(pendingError); |
| 3098 | } |
| 3099 | |
| 3100 | // Note: We cannot use state.visitForGc(visitor) here because the state machine's |
| 3101 | // visitForGc passes kj::Own<T>& to visitor.visit(), but GcVisitor expects T& for |
| 3102 | // types with visitForGc methods. We must dereference kj::Own manually. |
| 3103 | KJ_SWITCH_ONEOF(state) { |
| 3104 | KJ_CASE_ONEOF(initial, Initial) {} |
| 3105 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 3106 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 3107 | visitor.visit(error); |
| 3108 | } |
| 3109 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 3110 | visitor.visit(*consumer); |
| 3111 | } |
| 3112 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 3113 | visitor.visit(*consumer); |
| 3114 | } |
| 3115 | } |
| 3116 | visitor.visit(lock); |
| 3117 | } |
| 3118 | |
| 3119 | kj::Maybe<int> ReadableStreamJsController::getDesiredSize() { |
| 3120 | // If there's a pending state transition, return none |
| 3121 | if (state.hasPendingState()) { |
| 3122 | return kj::none; |
| 3123 | } |
| 3124 | |
| 3125 | KJ_SWITCH_ONEOF(state) { |
| 3126 | KJ_CASE_ONEOF(initial, Initial) { |
| 3127 | return kj::none; |
| 3128 | } |
| 3129 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3130 | return kj::none; |
| 3131 | } |
| 3132 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3133 | return kj::none; |
| 3134 | } |
| 3135 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 3136 | return consumer->getDesiredSize(); |
| 3137 | } |
| 3138 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 3139 | return consumer->getDesiredSize(); |
| 3140 | } |
| 3141 | } |
| 3142 | KJ_UNREACHABLE; |
| 3143 | } |
| 3144 | |
| 3145 | kj::Maybe<v8::Local<v8::Value>> ReadableStreamJsController::isErrored(jsg::Lock& js) { |
| 3146 | // Check for pending error first |
| 3147 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 3148 | return pendingError.getHandle(js); |
| 3149 | } |
| 3150 | // Pending Closed means not errored, so we can just check current state |
| 3151 | return state.tryGetUnsafe<StreamStates::Errored>().map( |
| 3152 | [&](jsg::Value& reason) { return reason.getHandle(js); }); |
| 3153 | } |
| 3154 | |
| 3155 | bool ReadableStreamJsController::canCloseOrEnqueue() { |
| 3156 | // If there's a pending state transition, can't close or enqueue |
| 3157 | if (state.hasPendingState()) { |
| 3158 | return false; |
| 3159 | } |
| 3160 | |
| 3161 | KJ_SWITCH_ONEOF(state) { |
| 3162 | KJ_CASE_ONEOF(initial, Initial) { |
| 3163 | return false; |
| 3164 | } |
| 3165 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3166 | return false; |
| 3167 | } |
| 3168 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3169 | return false; |
| 3170 | } |
| 3171 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 3172 | return consumer->canCloseOrEnqueue(); |
| 3173 | } |
| 3174 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 3175 | return consumer->canCloseOrEnqueue(); |
| 3176 | } |
| 3177 | } |
| 3178 | KJ_UNREACHABLE; |
| 3179 | } |
| 3180 | |
| 3181 | bool ReadableStreamJsController::hasBackpressure() { |
| 3182 | KJ_IF_SOME(size, getDesiredSize()) { |
| 3183 | return size <= 0; |
| 3184 | } |
| 3185 | return false; |
| 3186 | } |
| 3187 | |
| 3188 | kj::Maybe<kj::OneOf<DefaultController, ByobController>> ReadableStreamJsController:: |
| 3189 | getController() { |
| 3190 | // If there's a pending state transition, return none |
| 3191 | if (state.hasPendingState()) { |
| 3192 | return kj::none; |
| 3193 | } |
| 3194 | KJ_SWITCH_ONEOF(state) { |
| 3195 | KJ_CASE_ONEOF(initial, Initial) { |
| 3196 | return kj::none; |
| 3197 | } |
| 3198 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3199 | return kj::none; |
| 3200 | } |
| 3201 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3202 | return kj::none; |
| 3203 | } |
| 3204 | KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) { |
| 3205 | return consumer->getControllerRef(); |
| 3206 | } |
| 3207 | KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) { |
| 3208 | return consumer->getControllerRef(); |
| 3209 | } |
| 3210 | } |
| 3211 | KJ_UNREACHABLE; |
| 3212 | } |
| 3213 | |
| 3214 | namespace { |
| 3215 | // Consumes all bytes from a stream, buffering in memory, with the purpose |
| 3216 | // of producing either a single concatenated kj::Array<byte> or kj::String. |
| 3217 | class AllReader { |
| 3218 | public: |
| 3219 | using PartList = kj::Array<kj::ArrayPtr<byte>>; |
| 3220 | |
| 3221 | AllReader(jsg::Ref<ReadableStream> stream, uint64_t limit) |
| 3222 | : state(State::create<jsg::Ref<ReadableStream>>(kj::mv(stream))), |
| 3223 | limit(limit) {} |
| 3224 | KJ_DISALLOW_COPY_AND_MOVE(AllReader); |
| 3225 | |
| 3226 | jsg::Promise<jsg::BufferSource> allBytes(jsg::Lock& js) { |
| 3227 | return loop(js).then(js, [this](auto& js, PartList&& partPtrs) -> jsg::BufferSource { |
| 3228 | auto out = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, runningTotal); |
| 3229 | copyInto(out.asArrayPtr(), partPtrs.asPtr()); |
| 3230 | return jsg::BufferSource(js, kj::mv(out)); |
| 3231 | }); |
| 3232 | } |
| 3233 | |
| 3234 | jsg::Promise<kj::String> allText( |
| 3235 | jsg::Lock& js, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE) { |
| 3236 | return loop(js).then(js, [this, option](auto& js, PartList&& partPtrs) { |
| 3237 | // Strip UTF-8 BOM if requested |
| 3238 | if ((option & ReadAllTextOption::STRIP_BOM) && partPtrs.size() > 0 && |
| 3239 | hasUtf8Bom(partPtrs[0])) { |
| 3240 | partPtrs[0] = partPtrs[0].slice(UTF8_BOM_SIZE); |
| 3241 | runningTotal -= UTF8_BOM_SIZE; |
| 3242 | } |
| 3243 | |
| 3244 | JSG_REQUIRE(runningTotal <= v8::String::kMaxLength, RangeError, |
| 3245 | "String length exceeds v8::String::kMaxLength."); |
| 3246 | |
| 3247 | auto out = kj::heapArray<char>(runningTotal + 1); |
| 3248 | copyInto(out.first(out.size() - 1).asBytes(), partPtrs.asPtr()); |
| 3249 | out.back() = '\0'; |
| 3250 | return kj::String(kj::mv(out)); |
| 3251 | }); |
| 3252 | } |
| 3253 | |
| 3254 | void visitForGc(jsg::GcVisitor& visitor) { |
| 3255 | state.visitForGc(visitor); |
| 3256 | } |
| 3257 | |
| 3258 | private: |
| 3259 | // State machine for AllReader: |
| 3260 | // Closed is terminal, Errored is implicitly terminal via ErrorState. |
| 3261 | // jsg::Ref<ReadableStream> is the active state (still reading). |
| 3262 | using State = StateMachine<TerminalStates<StreamStates::Closed>, |
| 3263 | ErrorState<StreamStates::Errored>, |
| 3264 | ActiveState<jsg::Ref<ReadableStream>>, |
| 3265 | StreamStates::Closed, |
| 3266 | StreamStates::Errored, |
| 3267 | jsg::Ref<ReadableStream>>; |
| 3268 | State state; |
| 3269 | uint64_t limit; |
| 3270 | kj::Vector<jsg::BufferSource> parts; |
| 3271 | uint64_t runningTotal = 0; |
| 3272 | |
| 3273 | jsg::Promise<PartList> loop(jsg::Lock& js) { |
| 3274 | KJ_SWITCH_ONEOF(state) { |
| 3275 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3276 | return js.resolvedPromise(KJ_MAP(p, parts) { return p.asArrayPtr(); }); |
| 3277 | } |
| 3278 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3279 | return js.template rejectedPromise<PartList>(errored.getHandle(js)); |
| 3280 | } |
| 3281 | KJ_CASE_ONEOF(readable, jsg::Ref<ReadableStream>) { |
| 3282 | // Note that these nested lambda retain references to `this` and `readable` |
| 3283 | // and are passed into to promise returned by this method. It is the responsibility |
| 3284 | // of the caller to ensure that the AllReader instance is kept alive until the |
| 3285 | // promise is settled. |
| 3286 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, readable = readable.addRef()), (readable), |
| 3287 | (jsg::Lock & js, ReadResult result) mutable->jsg::Promise<PartList> { |
| 3288 | if (result.done) { |
| 3289 | state.template transitionTo<StreamStates::Closed>(); |
| 3290 | return loop(js); |
| 3291 | } |
| 3292 | |
| 3293 | // If we're not done, the result value must be interpretable as |
| 3294 | // bytes for the read to make any sense. |
| 3295 | auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js); |
| 3296 | if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) { |
| 3297 | auto error = js.v8TypeError("This ReadableStream did not return bytes."); |
| 3298 | state.template transitionTo<StreamStates::Errored>(js.v8Ref(error)); |
| 3299 | return readable->getController().cancel(js, error).then( |
| 3300 | js, [&](jsg::Lock& js) { return loop(js); }); |
| 3301 | } |
| 3302 | |
| 3303 | jsg::BufferSource bufferSource(js, handle); |
| 3304 | |
| 3305 | if (bufferSource.size() == 0) { |
| 3306 | // Weird but allowed, we'll skip it. |
| 3307 | return loop(js); |
| 3308 | } |
| 3309 | |
| 3310 | if ((runningTotal + bufferSource.size()) > limit) { |
| 3311 | auto error = js.v8TypeError("Memory limit exceeded before EOF."); |
| 3312 | state.template transitionTo<StreamStates::Errored>(js.v8Ref(error)); |
| 3313 | return readable->getController().cancel(js, error).then( |
| 3314 | js, [&](jsg::Lock& js) { return loop(js); }); |
| 3315 | } |
| 3316 | |
| 3317 | runningTotal += bufferSource.size(); |
| 3318 | parts.add(bufferSource.copy(js)); |
| 3319 | return loop(js); |
| 3320 | }); |
| 3321 | |
| 3322 | auto onFailure = [this](auto& js, jsg::Value exception) -> jsg::Promise<PartList> { |
| 3323 | // In this case the stream should already be errored. |
| 3324 | state.template transitionTo<StreamStates::Errored>(js.v8Ref(exception.getHandle(js))); |
| 3325 | return loop(js); |
| 3326 | }; |
| 3327 | |
| 3328 | return maybeAddFunctor(js, KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none)), |
| 3329 | kj::mv(onSuccess), kj::mv(onFailure)); |
| 3330 | } |
| 3331 | } |
| 3332 | KJ_UNREACHABLE; |
| 3333 | } |
| 3334 | |
| 3335 | void copyInto(kj::ArrayPtr<byte> out, kj::ArrayPtr<kj::ArrayPtr<byte>> in) { |
| 3336 | for (auto& part: in) { |
| 3337 | KJ_ASSERT(part.size() <= out.size()); |
| 3338 | out.first(part.size()).copyFrom(part); |
| 3339 | out = out.slice(part.size()); |
| 3340 | } |
| 3341 | } |
| 3342 | }; |
| 3343 | |
| 3344 | // PumpToReader implements the original JS promise-loop approach to pumping data from |
| 3345 | // a ReadableStream to a WritableStreamSink. It reads one chunk at a time using the |
| 3346 | // standard read() API, writes each chunk to the sink, and loops until done or errored. |
| 3347 | // This is the fallback path used when the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS |
| 3348 | // autogate is not enabled. |
| 3349 | class PumpToReader { |
| 3350 | public: |
| 3351 | PumpToReader(jsg::Ref<ReadableStream> stream, kj::Own<WritableStreamSink> sink, bool end) |
| 3352 | : ioContext(IoContext::current()), |
| 3353 | state(State::create<jsg::Ref<ReadableStream>>(kj::mv(stream))), |
| 3354 | sink(kj::mv(sink)), |
| 3355 | self(kj::refcounted<WeakRef<PumpToReader>>(kj::Badge<PumpToReader>{}, *this)), |
| 3356 | end(end) {} |
| 3357 | KJ_DISALLOW_COPY_AND_MOVE(PumpToReader); |
| 3358 | |
| 3359 | ~PumpToReader() noexcept(false) { |
| 3360 | self->invalidate(); |
| 3361 | // Ensure that if a write promise is pending it is proactively canceled. |
| 3362 | canceler.cancel("PumpToReader was destroyed"); |
| 3363 | } |
| 3364 | |
| 3365 | kj::Promise<void> pumpTo(jsg::Lock& js) { |
| 3366 | ioContext.requireCurrentOrThrowJs(); |
| 3367 | KJ_SWITCH_ONEOF(state) { |
| 3368 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 3369 | auto readable = stream.addRef(); |
| 3370 | state.template transitionTo<Pumping>(); |
| 3371 | return ioContext.awaitJs( |
| 3372 | js, pumpLoop(js, ioContext, kj::mv(readable), ioContext.addObject(self->addRef()))); |
| 3373 | } |
| 3374 | KJ_CASE_ONEOF(pumping, Pumping) { |
| 3375 | return KJ_EXCEPTION(FAILED, "pumping is already in progress"); |
| 3376 | } |
| 3377 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3378 | return KJ_EXCEPTION(FAILED, "stream has already been consumed"); |
| 3379 | } |
| 3380 | KJ_CASE_ONEOF(errored, kj::Exception) { |
| 3381 | return errored.clone(); |
| 3382 | } |
| 3383 | } |
| 3384 | KJ_UNREACHABLE; |
| 3385 | } |
| 3386 | |
| 3387 | private: |
| 3388 | struct Pumping { |
| 3389 | static constexpr kj::StringPtr NAME KJ_UNUSED = "pumping"_kj; |
| 3390 | }; |
| 3391 | IoContext& ioContext; |
| 3392 | |
| 3393 | using State = StateMachine<TerminalStates<StreamStates::Closed>, |
| 3394 | ErrorState<kj::Exception>, |
| 3395 | Pumping, |
| 3396 | StreamStates::Closed, |
| 3397 | kj::Exception, |
| 3398 | jsg::Ref<ReadableStream>>; |
| 3399 | State state; |
| 3400 | kj::Own<WritableStreamSink> sink; |
| 3401 | kj::Own<WeakRef<PumpToReader>> self; |
| 3402 | kj::Canceler canceler; |
| 3403 | bool end; |
| 3404 | |
| 3405 | bool isErroredOrClosed() { |
| 3406 | return state.isTerminal(); |
| 3407 | } |
| 3408 | |
| 3409 | jsg::Promise<void> pumpLoop(jsg::Lock& js, |
| 3410 | IoContext& ioContext, |
| 3411 | jsg::Ref<ReadableStream> readable, |
| 3412 | IoOwn<WeakRef<PumpToReader>> pumpToReader) { |
| 3413 | ioContext.requireCurrentOrThrowJs(); |
| 3414 | |
| 3415 | KJ_SWITCH_ONEOF(state) { |
| 3416 | KJ_CASE_ONEOF(ready, jsg::Ref<ReadableStream>) { |
| 3417 | KJ_UNREACHABLE; |
| 3418 | } |
| 3419 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3420 | return end ? ioContext.awaitIoLegacy(js, sink->end().attach(kj::mv(sink))) |
| 3421 | : js.resolvedPromise(); |
| 3422 | } |
| 3423 | KJ_CASE_ONEOF(errored, kj::Exception) { |
| 3424 | if (end) { |
| 3425 | sink->abort(errored.clone()); |
| 3426 | } |
| 3427 | return js.rejectedPromise<void>(errored.clone()); |
| 3428 | } |
| 3429 | KJ_CASE_ONEOF(pumping, Pumping) { |
| 3430 | using Result = kj::OneOf<Pumping, kj::Array<kj::byte>, StreamStates::Closed, jsg::Value>; |
| 3431 | |
| 3432 | return KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none)) |
| 3433 | .then(js, |
| 3434 | ioContext.addFunctor([byteStream = readable->getController().isByteOriented()]( |
| 3435 | auto& js, ReadResult result) mutable -> Result { |
| 3436 | if (result.done) { |
| 3437 | return StreamStates::Closed(); |
| 3438 | } |
| 3439 | |
| 3440 | auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js); |
| 3441 | if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) { |
| 3442 | return js.v8Ref(js.v8TypeError("This ReadableStream did not return bytes.")); |
| 3443 | } |
| 3444 | |
| 3445 | jsg::BufferSource bufferSource(js, handle); |
| 3446 | if (bufferSource.size() == 0) { |
| 3447 | return Pumping{}; |
| 3448 | } |
| 3449 | |
| 3450 | if (byteStream) { |
| 3451 | jsg::BackingStore backing = bufferSource.detach(js); |
| 3452 | return backing.asArrayPtr().attach(kj::mv(backing)); |
| 3453 | } |
| 3454 | return bufferSource.asArrayPtr().attach(kj::mv(bufferSource)); |
| 3455 | }), |
| 3456 | [](auto& js, jsg::Value exception) mutable -> Result { return kj::mv(exception); }) |
| 3457 | .then(js, ioContext.addFunctor( JSG_VISITABLE_LAMBDA((readable = kj::mv(readable), pumpToReader = kj::mv(pumpToReader)), (readable), (jsg::Lock & js, Result result) mutable { |
| 3458 | KJ_IF_SOME(reader, pumpToReader->tryGet()) { |
| 3459 | reader.ioContext.requireCurrentOrThrowJs(); |
| 3460 | auto& ioContext = IoContext::current(); |
| 3461 | KJ_SWITCH_ONEOF(result) { |
| 3462 | KJ_CASE_ONEOF(bytes, kj::Array<kj::byte>) { |
| 3463 | auto promise = reader.sink->write(bytes).attach(kj::mv(bytes)); |
| 3464 | return ioContext.awaitIo(js, reader.canceler.wrap(kj::mv(promise))) |
| 3465 | .then(js, |
| 3466 | [](jsg::Lock& js) -> kj::Maybe<jsg::Value> { |
| 3467 | return kj::Maybe<jsg::Value>(kj::none); |
| 3468 | }, |
| 3469 | [](jsg::Lock& js, jsg::Value exception) mutable -> kj::Maybe<jsg::Value> { |
| 3470 | return kj::mv(exception); |
| 3471 | }) |
| 3472 | .then(js, |
| 3473 | ioContext.addFunctor(JSG_VISITABLE_LAMBDA( |
| 3474 | (readable = readable.addRef(), pumpToReader = kj::mv(pumpToReader)), |
| 3475 | (readable), |
| 3476 | (jsg::Lock & js, kj::Maybe<jsg::Value> maybeException) mutable { |
| 3477 | KJ_IF_SOME(reader, pumpToReader->tryGet()) { |
| 3478 | auto& ioContext = reader.ioContext; |
| 3479 | ioContext.requireCurrentOrThrowJs(); |
| 3480 | KJ_IF_SOME(exception, maybeException) { |
| 3481 | if (!reader.isErroredOrClosed()) { |
| 3482 | reader.state.transitionTo<kj::Exception>( |
| 3483 | js.exceptionToKj(kj::mv(exception))); |
| 3484 | } |
| 3485 | } else { |
| 3486 | // Else block to avert dangling else compiler warning. |
| 3487 | } |
| 3488 | return reader.pumpLoop( |
| 3489 | js, ioContext, readable.addRef(), kj::mv(pumpToReader)); |
| 3490 | } else { |
| 3491 | return readable->getController().cancel(js, |
| 3492 | maybeException.map( |
| 3493 | [&](jsg::Value& ex) { return ex.getHandle(js); })); |
| 3494 | } |
| 3495 | }))); |
| 3496 | } |
| 3497 | KJ_CASE_ONEOF(pumping, Pumping) {} |
| 3498 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3499 | if (!reader.isErroredOrClosed()) { |
| 3500 | reader.state.transitionTo<StreamStates::Closed>(); |
| 3501 | } |
| 3502 | } |
| 3503 | KJ_CASE_ONEOF(exception, jsg::Value) { |
| 3504 | if (!reader.isErroredOrClosed()) { |
| 3505 | reader.state.transitionTo<kj::Exception>(js.exceptionToKj(kj::mv(exception))); |
| 3506 | } |
| 3507 | } |
| 3508 | } |
| 3509 | return reader.pumpLoop(js, ioContext, readable.addRef(), kj::mv(pumpToReader)); |
| 3510 | } else { |
| 3511 | KJ_SWITCH_ONEOF(result) { |
| 3512 | KJ_CASE_ONEOF(bytes, kj::Array<kj::byte>) { |
| 3513 | return readable->getController().cancel(js, kj::none); |
| 3514 | } |
| 3515 | KJ_CASE_ONEOF(pumping, Pumping) { |
| 3516 | return readable->getController().cancel(js, kj::none); |
| 3517 | } |
| 3518 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3519 | return js.resolvedPromise(); |
| 3520 | } |
| 3521 | KJ_CASE_ONEOF(exception, jsg::Value) { |
| 3522 | return readable->getController().cancel(js, exception.getHandle(js)); |
| 3523 | } |
| 3524 | } |
| 3525 | } |
| 3526 | KJ_UNREACHABLE; |
| 3527 | }))); |
| 3528 | } |
| 3529 | } |
| 3530 | KJ_UNREACHABLE; |
| 3531 | } |
| 3532 | }; |
| 3533 | |
| 3534 | // pumpToCoroutine uses a DrainingReader to efficiently pull all synchronously available |
| 3535 | // data from the stream in each iteration, then writes it to the sink using vectored |
| 3536 | // I/O. This minimizes isolate lock acquisitions by batching: each time the lock is |
| 3537 | // held, the stream's internal queue is fully drained and the JS pull callback is |
| 3538 | // pumped synchronously as many times as possible. |
| 3539 | // |
| 3540 | // The pump loop is a kj coroutine. Dropping the returned kj::Promise drops the |
| 3541 | // coroutine frame, which destroys the DrainingReader (releasing the stream lock) |
| 3542 | // and the sink. No WeakRef/IoOwn dance is needed because ownership is clear. |
| 3543 | // The coroutine that implements the pump loop takes ownership of the DrainingReader |
| 3544 | // and sink. The jsg::Ref<ReadableStream> is not passed into the coroutine because |
| 3545 | // jsg::Ref is disallowed in coroutine parameters; instead, the DrainingReader holds |
| 3546 | // a reference to the stream internally. |
| 3547 | kj::Promise<void> pumpToImpl(IoContext& ioContext, |
| 3548 | kj::Own<DrainingReader> reader, |
| 3549 | kj::Own<WritableStreamSink> sink, |
| 3550 | bool end) { |
| 3551 | |
| 3552 | bool writeFailed = false; |
| 3553 | |
| 3554 | KJ_TRY { |
| 3555 | while (true) { |
| 3556 | // Perform a draining read to get all synchronously available data if possible |
| 3557 | // or fall back to a regular read if not. |
| 3558 | DrainingReadResult result = co_await ioContext.run([&reader](jsg::Lock& js) mutable { |
| 3559 | auto& ioContext = IoContext::current(); |
| 3560 | // Use a 256KB limit to allow periodic yielding to the event loop, |
| 3561 | // preventing a fast producer from monopolizing the thread. |
| 3562 | constexpr size_t kMaxReadPerCycle = 256 * 1024; |
| 3563 | return ioContext.awaitJs(js, reader->read(js, kMaxReadPerCycle)); |
| 3564 | }); |
| 3565 | |
| 3566 | // Write all the chunks we received using vectored write for efficiency. |
| 3567 | if (result.chunks.size() > 0) { |
| 3568 | KJ_ON_SCOPE_FAILURE(writeFailed = true); |
| 3569 | auto pieces = |
| 3570 | KJ_MAP(chunk, result.chunks) -> kj::ArrayPtr<const kj::byte> { return chunk.asPtr(); }; |
| 3571 | co_await sink->write(pieces); |
| 3572 | } |
| 3573 | |
| 3574 | // If the stream is done, end the output if needed and exit. |
| 3575 | if (result.done) { |
| 3576 | KJ_ON_SCOPE_FAILURE(writeFailed = true); |
| 3577 | if (end) { |
| 3578 | co_await sink->end(); |
| 3579 | } |
| 3580 | co_return; |
| 3581 | } |
| 3582 | } |
| 3583 | } |
| 3584 | KJ_CATCH(exception) { |
| 3585 | if (!writeFailed) { |
| 3586 | sink->abort(exception.clone()); |
| 3587 | } |
| 3588 | |
| 3589 | co_await ioContext.run([&reader, ex = exception.clone()](jsg::Lock& js) mutable { |
| 3590 | auto& ioContext = IoContext::current(); |
| 3591 | auto error = js.exceptionToJsValue(kj::mv(ex)); |
| 3592 | return ioContext.awaitJs(js, reader->cancel(js, error.getHandle(js))); |
| 3593 | }); |
| 3594 | kj::throwFatalException(kj::mv(exception)); |
| 3595 | } |
| 3596 | } |
| 3597 | } // namespace |
| 3598 | |
| 3599 | template <typename T> |
| 3600 | jsg::Promise<T> ReadableStreamJsController::readAll(jsg::Lock& js, uint64_t limit) { |
| 3601 | if (isLockedToReader()) { |
| 3602 | return js.rejectedPromise<T>(KJ_EXCEPTION( |
| 3603 | FAILED, "jsg.TypeError: This ReadableStream is currently locked to a reader.")); |
| 3604 | } |
| 3605 | disturbed = true; |
| 3606 | |
| 3607 | bool stripBom = false; |
| 3608 | KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) { |
| 3609 | stripBom = flags.getStripBomInReadAllText(); |
| 3610 | } |
| 3611 | |
| 3612 | // This operation leaves the stream locked and disturbed. The loop will read until |
| 3613 | // the stream is closed or errored. If the limit is reached, the loop will error. |
| 3614 | |
| 3615 | const auto readAll = [this, limit, stripBom](auto& js) -> jsg::Promise<T> { |
| 3616 | KJ_ASSERT(lock.lock()); |
| 3617 | // The AllReader will hold a traceable reference to the ReadableStream. |
| 3618 | auto reader = kj::heap<AllReader>(addRef(), limit); |
| 3619 | |
| 3620 | auto promise = ([&js, &reader, stripBom]() -> jsg::Promise<T> { |
| 3621 | if constexpr (kj::isSameType<T, jsg::BufferSource>()) { |
| 3622 | (void)stripBom; // Unused in this branch. |
| 3623 | return reader->allBytes(js); |
| 3624 | } else { |
| 3625 | auto option = ReadAllTextOption::NULL_TERMINATE; |
| 3626 | if (stripBom) { |
| 3627 | option |= ReadAllTextOption::STRIP_BOM; |
| 3628 | } |
| 3629 | return reader->allText(js, option); |
| 3630 | } |
| 3631 | })(); |
| 3632 | |
| 3633 | return maybeAddFunctor(js, kj::mv(promise), |
| 3634 | // reader is a GC visitable type that holds a reference to either the stream |
| 3635 | // or an error. Accordingly, we wrap it in a visitable lambda attached as a |
| 3636 | // continuation on the promise to ensure that it is GC visited and kept alive until |
| 3637 | // the promise settles. |
| 3638 | JSG_VISITABLE_LAMBDA((reader = kj::mv(reader)), (reader), |
| 3639 | (jsg::Lock & js, T result)->jsg::Promise<T> { |
| 3640 | return js.resolvedPromise(kj::mv(result)); |
| 3641 | }), |
| 3642 | [](jsg::Lock& js, jsg::Value exception) -> jsg::Promise<T> { |
| 3643 | return js.rejectedPromise<T>(kj::mv(exception)); |
| 3644 | }); |
| 3645 | }; |
| 3646 | |
| 3647 | KJ_SWITCH_ONEOF(state) { |
| 3648 | KJ_CASE_ONEOF(initial, Initial) { |
| 3649 | // Stream not yet set up, treat as closed. |
| 3650 | if constexpr (kj::isSameType<T, jsg::BufferSource>()) { |
| 3651 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0); |
| 3652 | return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); |
| 3653 | } else { |
| 3654 | return js.resolvedPromise(T()); |
| 3655 | } |
| 3656 | } |
| 3657 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3658 | if constexpr (kj::isSameType<T, jsg::BufferSource>()) { |
| 3659 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0); |
| 3660 | return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); |
| 3661 | } else { |
| 3662 | return js.resolvedPromise(T()); |
| 3663 | } |
| 3664 | } |
| 3665 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3666 | return js.rejectedPromise<T>(errored.addRef(js)); |
| 3667 | } |
| 3668 | KJ_CASE_ONEOF(valueReadable, kj::Own<ValueReadable>) { |
| 3669 | return readAll(js); |
| 3670 | } |
| 3671 | KJ_CASE_ONEOF(byteReadable, kj::Own<ByteReadable>) { |
| 3672 | return readAll(js); |
| 3673 | } |
| 3674 | } |
| 3675 | KJ_UNREACHABLE; |
| 3676 | } |
| 3677 | |
| 3678 | jsg::Promise<jsg::BufferSource> ReadableStreamJsController::readAllBytes( |
| 3679 | jsg::Lock& js, uint64_t limit) { |
| 3680 | return readAll<jsg::BufferSource>(js, limit); |
| 3681 | } |
| 3682 | |
| 3683 | jsg::Promise<kj::String> ReadableStreamJsController::readAllText(jsg::Lock& js, uint64_t limit) { |
| 3684 | return readAll<kj::String>(js, limit); |
| 3685 | } |
| 3686 | |
| 3687 | kj::Own<ReadableStreamController> ReadableStreamJsController::detach( |
| 3688 | jsg::Lock& js, bool ignored /* unused */) { |
| 3689 | KJ_ASSERT(!isLockedToReader()); |
| 3690 | KJ_ASSERT(!isDisturbed()); |
| 3691 | KJ_ASSERT(!state.hasOperationInProgress(), "Unable to detach with read pending"); |
| 3692 | auto controller = kj::heap<ReadableStreamJsController>(); |
| 3693 | controller->expectedLength = expectedLength; |
| 3694 | disturbed = true; |
| 3695 | |
| 3696 | // Clones this streams state into a new ReadableStreamController, leaving this stream |
| 3697 | // locked, disturbed, and closed. |
| 3698 | |
| 3699 | // The controller starts in Initial state by default, so we can use regular transitionTo. |
| 3700 | KJ_SWITCH_ONEOF(state) { |
| 3701 | KJ_CASE_ONEOF(initial, Initial) { |
| 3702 | // Still in initial state, transition to closed |
| 3703 | controller->state.transitionTo<StreamStates::Closed>(); |
| 3704 | } |
| 3705 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3706 | controller->state.transitionTo<StreamStates::Closed>(); |
| 3707 | } |
| 3708 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3709 | controller->state.transitionTo<StreamStates::Errored>(errored.addRef(js)); |
| 3710 | } |
| 3711 | KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) { |
| 3712 | KJ_ASSERT(lock.lock()); |
| 3713 | controller->state.transitionTo<kj::Own<ValueReadable>>(readable->clone(js, *controller)); |
| 3714 | state.transitionTo<StreamStates::Closed>(); |
| 3715 | lock.onClose(js); |
| 3716 | } |
| 3717 | KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) { |
| 3718 | KJ_ASSERT(lock.lock()); |
| 3719 | controller->state.transitionTo<kj::Own<ByteReadable>>(readable->clone(js, *controller)); |
| 3720 | state.transitionTo<StreamStates::Closed>(); |
| 3721 | lock.onClose(js); |
| 3722 | } |
| 3723 | } |
| 3724 | |
| 3725 | return kj::mv(controller); |
| 3726 | } |
| 3727 | |
| 3728 | kj::Maybe<uint64_t> ReadableStreamJsController::tryGetLength(StreamEncoding encoding) { |
| 3729 | return expectedLength; |
| 3730 | } |
| 3731 | |
| 3732 | kj::Promise<DeferredProxy<void>> ReadableStreamJsController::pumpTo( |
| 3733 | jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) { |
| 3734 | KJ_ASSERT(IoContext::hasCurrent(), "Unable to consume this ReadableStream outside of a request"); |
| 3735 | KJ_REQUIRE(!isLockedToReader(), "This ReadableStream is currently locked to a reader."); |
| 3736 | disturbed = true; |
| 3737 | |
| 3738 | // This operation will leave the ReadableStream locked and disturbed. It will consume |
| 3739 | // the stream until it either closed or errors. |
| 3740 | // |
| 3741 | // When the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS autogate is enabled, uses the new |
| 3742 | // pumpToImpl coroutine with DrainingReader for batched reads and vectored writes. |
| 3743 | // Otherwise, falls back to the original PumpToReader JS promise loop that reads one |
| 3744 | // chunk at a time. |
| 3745 | |
| 3746 | const auto handlePump = [&] { |
| 3747 | if (util::Autogate::isEnabled(util::AutogateKey::ENABLE_DRAINING_READ_ON_STANDARD_STREAMS)) { |
| 3748 | auto reader = KJ_ASSERT_NONNULL(DrainingReader::create(js, *this->addRef()), |
| 3749 | "Failed to create DrainingReader โ stream should not be locked"); |
| 3750 | auto& ioContext = IoContext::current(); |
| 3751 | return addNoopDeferredProxy(pumpToImpl(ioContext, kj::mv(reader), kj::mv(sink), end)); |
| 3752 | } else { |
| 3753 | KJ_ASSERT(lock.lock()); |
| 3754 | auto reader = kj::heap<PumpToReader>(addRef(), kj::mv(sink), end); |
| 3755 | return addNoopDeferredProxy(reader->pumpTo(js).attach(kj::mv(reader))); |
| 3756 | } |
| 3757 | }; |
| 3758 | |
| 3759 | KJ_SWITCH_ONEOF(state) { |
| 3760 | KJ_CASE_ONEOF(initial, Initial) { |
| 3761 | // Stream not yet set up, treat as closed. |
| 3762 | return addNoopDeferredProxy(sink->end().attach(kj::mv(sink))); |
| 3763 | } |
| 3764 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3765 | return addNoopDeferredProxy(sink->end().attach(kj::mv(sink))); |
| 3766 | } |
| 3767 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3768 | return js.exceptionToKj(errored.addRef(js)); |
| 3769 | } |
| 3770 | KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) { |
| 3771 | return handlePump(); |
| 3772 | } |
| 3773 | KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) { |
| 3774 | return handlePump(); |
| 3775 | } |
| 3776 | } |
| 3777 | |
| 3778 | KJ_UNREACHABLE; |
| 3779 | } |
| 3780 | |
| 3781 | // ====================================================================================== |
| 3782 | |
| 3783 | WritableStreamDefaultController::WritableStreamDefaultController( |
| 3784 | jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal) |
| 3785 | : ioContext(tryGetIoContext()), |
| 3786 | impl(js, owner, kj::mv(abortSignal)) {} |
| 3787 | |
| 3788 | jsg::Promise<void> WritableStreamDefaultController::abort( |
| 3789 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 3790 | return impl.abort(js, JSG_THIS, reason); |
| 3791 | } |
| 3792 | |
| 3793 | void WritableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { |
| 3794 | visitor.visit(impl); |
| 3795 | } |
| 3796 | |
| 3797 | jsg::Promise<void> WritableStreamDefaultController::close(jsg::Lock& js) { |
| 3798 | return impl.close(js, JSG_THIS); |
| 3799 | } |
| 3800 | |
| 3801 | void WritableStreamDefaultController::error( |
| 3802 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) { |
| 3803 | impl.error(js, JSG_THIS, reason.orDefault(js.undefined())); |
| 3804 | } |
| 3805 | |
| 3806 | kj::Maybe<ssize_t> WritableStreamDefaultController::getDesiredSize() { |
| 3807 | // Per the spec, desiredSize should be null when the stream is erroring. |
| 3808 | if (impl.flags.pedanticWpt && isErroring()) { |
| 3809 | return kj::none; |
| 3810 | } |
| 3811 | return impl.getDesiredSize(); |
| 3812 | } |
| 3813 | |
| 3814 | jsg::Ref<AbortSignal> WritableStreamDefaultController::getSignal() { |
| 3815 | return impl.signal.addRef(); |
| 3816 | } |
| 3817 | |
| 3818 | kj::Maybe<v8::Local<v8::Value>> WritableStreamDefaultController::isErroring(jsg::Lock& js) { |
| 3819 | KJ_IF_SOME(erroring, impl.state.tryGetUnsafe<StreamStates::Erroring>()) { |
| 3820 | return erroring.reason.getHandle(js); |
| 3821 | } |
| 3822 | return kj::none; |
| 3823 | } |
| 3824 | |
| 3825 | void WritableStreamDefaultController::setup( |
| 3826 | jsg::Lock& js, UnderlyingSink underlyingSink, StreamQueuingStrategy queuingStrategy) { |
| 3827 | impl.setup(js, JSG_THIS, kj::mv(underlyingSink), kj::mv(queuingStrategy)); |
| 3828 | } |
| 3829 | |
| 3830 | jsg::Promise<void> WritableStreamDefaultController::write( |
| 3831 | jsg::Lock& js, v8::Local<v8::Value> value) { |
| 3832 | return impl.write(js, JSG_THIS, value); |
| 3833 | } |
| 3834 | |
| 3835 | void WritableStreamDefaultController::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) { |
| 3836 | impl.cancelPendingWrites(js, reason); |
| 3837 | } |
| 3838 | |
| 3839 | void WritableStreamDefaultController::clearAlgorithms() { |
| 3840 | impl.algorithms.clear(); |
| 3841 | } |
| 3842 | |
| 3843 | WritableStreamDefaultController::~WritableStreamDefaultController() noexcept(false) { |
| 3844 | // Clear algorithms in destructor to break circular references |
| 3845 | clearAlgorithms(); |
| 3846 | } |
| 3847 | |
| 3848 | // ====================================================================================== |
| 3849 | WritableStreamJsController::WritableStreamJsController(): ioContext(tryGetIoContext()) {} |
| 3850 | |
| 3851 | WritableStreamJsController::~WritableStreamJsController() noexcept(false) { |
| 3852 | // Clear algorithms to break circular references during destruction |
| 3853 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 3854 | controller->clearAlgorithms(); |
| 3855 | } |
| 3856 | // Clear the state to break the circular reference to the controller. |
| 3857 | // During destruction, we force the transition since the current state doesn't matter. |
| 3858 | state.forceTransitionTo<StreamStates::Closed>(); |
| 3859 | // Clear owner reference |
| 3860 | owner = kj::none; |
| 3861 | // Clear any pending abort promise |
| 3862 | maybeAbortPromise = kj::none; |
| 3863 | } |
| 3864 | |
| 3865 | WritableStreamJsController::WritableStreamJsController(StreamStates::Closed closed) |
| 3866 | : ioContext(tryGetIoContext()) { |
| 3867 | state.transitionTo<StreamStates::Closed>(); |
| 3868 | } |
| 3869 | |
| 3870 | WritableStreamJsController::WritableStreamJsController(StreamStates::Errored errored) |
| 3871 | : ioContext(tryGetIoContext()) { |
| 3872 | state.transitionTo<StreamStates::Errored>(kj::mv(errored)); |
| 3873 | } |
| 3874 | |
| 3875 | jsg::Promise<void> WritableStreamJsController::abort( |
| 3876 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) { |
| 3877 | // The spec requires that if abort is called multiple times, it is supposed to return the same |
| 3878 | // promise each time. That's a bit cumbersome here with jsg::Promise so we intentionally just |
| 3879 | // return a continuation branch off the same promise. |
| 3880 | KJ_IF_SOME(abortPromise, maybeAbortPromise) { |
| 3881 | return abortPromise.whenResolved(js); |
| 3882 | } |
| 3883 | KJ_SWITCH_ONEOF(state) { |
| 3884 | KJ_CASE_ONEOF(initial, Initial) { |
| 3885 | // Stream hasn't been set up yet - treat like closed for abort purposes |
| 3886 | maybeAbortPromise = js.resolvedPromise(); |
| 3887 | return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); |
| 3888 | } |
| 3889 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3890 | maybeAbortPromise = js.resolvedPromise(); |
| 3891 | return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); |
| 3892 | } |
| 3893 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3894 | // Per the spec, if the stream is errored, we are to return a resolved promise. |
| 3895 | maybeAbortPromise = js.resolvedPromise(); |
| 3896 | return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); |
| 3897 | } |
| 3898 | KJ_CASE_ONEOF(controller, Controller) { |
| 3899 | maybeAbortPromise = controller->abort(js, reason.orDefault(js.undefined())); |
| 3900 | return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); |
| 3901 | } |
| 3902 | } |
| 3903 | KJ_UNREACHABLE; |
| 3904 | } |
| 3905 | |
| 3906 | jsg::Ref<WritableStream> WritableStreamJsController::addRef() { |
| 3907 | return KJ_ASSERT_NONNULL(owner).addRef(); |
| 3908 | } |
| 3909 | |
| 3910 | bool WritableStreamJsController::isClosedOrClosing() { |
| 3911 | return state.is<StreamStates::Closed>(); |
| 3912 | } |
| 3913 | |
| 3914 | bool WritableStreamJsController::isErrored() { |
| 3915 | return state.isErrored(); |
| 3916 | } |
| 3917 | |
| 3918 | jsg::Promise<void> WritableStreamJsController::close(jsg::Lock& js, bool markAsHandled) { |
| 3919 | KJ_SWITCH_ONEOF(state) { |
| 3920 | KJ_CASE_ONEOF(initial, Initial) { |
| 3921 | return rejectedMaybeHandledPromise<void>( |
| 3922 | js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled); |
| 3923 | } |
| 3924 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 3925 | return rejectedMaybeHandledPromise<void>( |
| 3926 | js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled); |
| 3927 | } |
| 3928 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 3929 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 3930 | return rejectedMaybeHandledPromise<void>( |
| 3931 | js, js.v8TypeError("This WritableStream has been errored."_kj), markAsHandled); |
| 3932 | } |
| 3933 | return rejectedMaybeHandledPromise<void>(js, errored.getHandle(js), markAsHandled); |
| 3934 | } |
| 3935 | KJ_CASE_ONEOF(controller, Controller) { |
| 3936 | return controller->close(js); |
| 3937 | } |
| 3938 | } |
| 3939 | KJ_UNREACHABLE; |
| 3940 | } |
| 3941 | |
| 3942 | void WritableStreamJsController::doClose(jsg::Lock& js) { |
| 3943 | // If already in a terminal state, nothing to do. |
| 3944 | if (state.isTerminal()) return; |
| 3945 | |
| 3946 | // Clear algorithms to break circular references before changing state |
| 3947 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 3948 | controller->clearAlgorithms(); |
| 3949 | } |
| 3950 | |
| 3951 | state.transitionTo<StreamStates::Closed>(); |
| 3952 | KJ_IF_SOME(locked, lock.state.tryGetUnsafe<WriterLocked>()) { |
| 3953 | maybeResolvePromise(js, locked.getClosedFulfiller()); |
| 3954 | maybeResolvePromise(js, locked.getReadyFulfiller()); |
| 3955 | } else { |
| 3956 | (void)lock.state.transitionFromTo<WritableLockImpl::PipeLocked, Unlocked>(); |
| 3957 | } |
| 3958 | } |
| 3959 | |
| 3960 | void WritableStreamJsController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 3961 | // If already in a terminal state, nothing to do. |
| 3962 | if (state.isTerminal()) return; |
| 3963 | |
| 3964 | // Clear algorithms to break circular references before changing state |
| 3965 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 3966 | controller->clearAlgorithms(); |
| 3967 | } |
| 3968 | |
| 3969 | state.transitionTo<StreamStates::Errored>(js.v8Ref(reason)); |
| 3970 | KJ_IF_SOME(locked, lock.state.tryGetUnsafe<WriterLocked>()) { |
| 3971 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason); |
| 3972 | maybeResolvePromise(js, locked.getReadyFulfiller()); |
| 3973 | } else KJ_IF_SOME(pipeLocked, lock.state.tryGetUnsafe<WritableLockImpl::PipeLocked>()) { |
| 3974 | // When the writable side of a pipe errors, we need to release the source stream. |
| 3975 | // The pipeLoop may be waiting on a read from the source that will never complete, |
| 3976 | // so we need to proactively release the source here. |
| 3977 | if (!pipeLocked.flags.preventCancel) { |
| 3978 | pipeLocked.source.release(js, reason); |
| 3979 | } else { |
| 3980 | pipeLocked.source.release(js); |
| 3981 | } |
| 3982 | lock.state.transitionTo<Unlocked>(); |
| 3983 | } |
| 3984 | } |
| 3985 | |
| 3986 | void WritableStreamJsController::errorIfNeeded(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 3987 | // Error through the underlying controller if available, which goes through the proper |
| 3988 | // error transition (Erroring -> Errored). This allows close() to be called while the |
| 3989 | // stream is "erroring" and reject with the stored error. |
| 3990 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 3991 | controller->error(js, reason); |
| 3992 | } |
| 3993 | // If state is not Controller (already Closed or Errored), this is a no-op. |
| 3994 | } |
| 3995 | |
| 3996 | kj::Maybe<int> WritableStreamJsController::getDesiredSize() { |
| 3997 | KJ_SWITCH_ONEOF(state) { |
| 3998 | KJ_CASE_ONEOF(initial, Initial) { |
| 3999 | return 0; |
| 4000 | } |
| 4001 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 4002 | return 0; |
| 4003 | } |
| 4004 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 4005 | return kj::none; |
| 4006 | } |
| 4007 | KJ_CASE_ONEOF(controller, Controller) { |
| 4008 | return controller->getDesiredSize().map([](ssize_t size) -> int { return size; }); |
| 4009 | } |
| 4010 | } |
| 4011 | KJ_UNREACHABLE; |
| 4012 | } |
| 4013 | |
| 4014 | kj::Maybe<v8::Local<v8::Value>> WritableStreamJsController::isErroring(jsg::Lock& js) { |
| 4015 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 4016 | return controller->isErroring(js); |
| 4017 | } |
| 4018 | return kj::none; |
| 4019 | } |
| 4020 | |
| 4021 | bool WritableStreamDefaultController::isErroring() const { |
| 4022 | return impl.state.is<StreamStates::Erroring>(); |
| 4023 | } |
| 4024 | |
| 4025 | kj::Maybe<v8::Local<v8::Value>> WritableStreamJsController::isErroredOrErroring(jsg::Lock& js) { |
| 4026 | KJ_IF_SOME(err, state.tryGetErrorUnsafe()) { |
| 4027 | return err.getHandle(js); |
| 4028 | } |
| 4029 | return isErroring(js); |
| 4030 | } |
| 4031 | |
| 4032 | bool WritableStreamJsController::isStarted() { |
| 4033 | KJ_SWITCH_ONEOF(state) { |
| 4034 | KJ_CASE_ONEOF(initial, Initial) { |
| 4035 | return false; |
| 4036 | } |
| 4037 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 4038 | return true; |
| 4039 | } |
| 4040 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 4041 | return true; |
| 4042 | } |
| 4043 | KJ_CASE_ONEOF(controller, Controller) { |
| 4044 | return controller->isStarted(); |
| 4045 | } |
| 4046 | } |
| 4047 | KJ_UNREACHABLE; |
| 4048 | } |
| 4049 | |
| 4050 | bool WritableStreamJsController::hasBackpressure() { |
| 4051 | KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) { |
| 4052 | return controller->hasBackpressure(); |
| 4053 | } |
| 4054 | return false; |
| 4055 | } |
| 4056 | |
| 4057 | bool WritableStreamJsController::isLocked() const { |
| 4058 | return isLockedToWriter(); |
| 4059 | } |
| 4060 | |
| 4061 | bool WritableStreamJsController::isLockedToWriter() const { |
| 4062 | return !lock.state.is<Unlocked>(); |
| 4063 | } |
| 4064 | |
| 4065 | bool WritableStreamJsController::lockWriter(jsg::Lock& js, Writer& writer) { |
| 4066 | return lock.lockWriter(js, *this, writer); |
| 4067 | } |
| 4068 | |
| 4069 | void WritableStreamJsController::maybeRejectReadyPromise( |
| 4070 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 4071 | KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) { |
| 4072 | if (writerLock.getReadyFulfiller() != kj::none) { |
| 4073 | maybeRejectPromise<void>(js, writerLock.getReadyFulfiller(), reason); |
| 4074 | } else { |
| 4075 | auto prp = js.newPromiseAndResolver<void>(); |
| 4076 | prp.promise.markAsHandled(js); |
| 4077 | prp.resolver.reject(js, reason); |
| 4078 | writerLock.setReadyFulfiller(js, prp); |
| 4079 | } |
| 4080 | } |
| 4081 | } |
| 4082 | |
| 4083 | void WritableStreamJsController::maybeResolveReadyPromise(jsg::Lock& js) { |
| 4084 | KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) { |
| 4085 | maybeResolvePromise(js, writerLock.getReadyFulfiller()); |
| 4086 | } |
| 4087 | } |
| 4088 | |
| 4089 | void WritableStreamJsController::releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) { |
| 4090 | lock.releaseWriter(*this, writer, maybeJs); |
| 4091 | } |
| 4092 | |
| 4093 | kj::Maybe<kj::Own<WritableStreamSink>> WritableStreamJsController::removeSink(jsg::Lock& js) { |
| 4094 | return kj::none; |
| 4095 | } |
| 4096 | void WritableStreamJsController::detach(jsg::Lock& js) { |
| 4097 | KJ_UNIMPLEMENTED("WritableStreamJsController::detach is not implemented"); |
| 4098 | } |
| 4099 | |
| 4100 | void WritableStreamJsController::setOwnerRef(WritableStream& stream) { |
| 4101 | owner = stream; |
| 4102 | } |
| 4103 | |
| 4104 | void WritableStreamJsController::setup(jsg::Lock& js, |
| 4105 | jsg::Optional<UnderlyingSink> maybeUnderlyingSink, |
| 4106 | jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) { |
| 4107 | auto underlyingSink = kj::mv(maybeUnderlyingSink).orDefault({}); |
| 4108 | auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({}); |
| 4109 | |
| 4110 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4111 | // Per the spec, the type property for WritableStream's underlying sink must be undefined. |
| 4112 | // If it's anything else, throw a RangeError. |
| 4113 | JSG_REQUIRE(underlyingSink.type == kj::none, RangeError, |
| 4114 | "Invalid underlying sink type. Only undefined is valid."); |
| 4115 | } |
| 4116 | |
| 4117 | // We account for the memory usage of the WritableStreamDefaultController and AbortSignal together |
| 4118 | // because their lifetimes are identical and memory accounting itself has a memory overhead. |
| 4119 | auto controller = js.allocAccounted<WritableStreamDefaultController>( |
| 4120 | sizeof(WritableStreamDefaultController) + sizeof(AbortSignal), js, KJ_ASSERT_NONNULL(owner), |
| 4121 | js.alloc<AbortSignal>()); |
| 4122 | auto& controllerRef = *controller; |
| 4123 | state.transitionTo<Controller>(kj::mv(controller)); |
| 4124 | controllerRef.setup(js, kj::mv(underlyingSink), kj::mv(queuingStrategy)); |
| 4125 | } |
| 4126 | |
| 4127 | kj::Maybe<jsg::Promise<void>> WritableStreamJsController::tryPipeFrom( |
| 4128 | jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) { |
| 4129 | JSG_REQUIRE_NONNULL( |
| 4130 | ioContext, Error, "Unable to pipe to a WritableStream created outside of a request"); |
| 4131 | |
| 4132 | // The ReadableStream source here can be either a JavaScript-backed ReadableStream |
| 4133 | // or ReadableStreamSource-backed. In either case, however, this WritableStream is |
| 4134 | // JavaScript-based and must use a JavaScript promise-based data flow for piping data. |
| 4135 | // We'll treat all ReadableStreams as if they are JavaScript-backed. |
| 4136 | // |
| 4137 | // This method will return a JavaScript promise that is resolved when the pipe operation |
| 4138 | // completes, or is rejected if the pipe operation is aborted or errored. |
| 4139 | |
| 4140 | // Let's also acquire the destination pipe lock. |
| 4141 | lock.pipeLock(KJ_ASSERT_NONNULL(owner), kj::mv(source), options); |
| 4142 | |
| 4143 | return pipeLoop(js).then(js, JSG_VISITABLE_LAMBDA((ref = addRef()), (ref), (auto& js){})); |
| 4144 | } |
| 4145 | |
| 4146 | jsg::Promise<void> WritableStreamJsController::pipeLoop(jsg::Lock& js) { |
| 4147 | auto maybePipeLock = lock.tryGetPipe(); |
| 4148 | if (maybePipeLock == kj::none) return js.resolvedPromise(); |
| 4149 | auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock); |
| 4150 | |
| 4151 | auto preventAbort = pipeLock.flags.preventAbort; |
| 4152 | auto preventCancel = pipeLock.flags.preventCancel; |
| 4153 | auto preventClose = pipeLock.flags.preventClose; |
| 4154 | auto pipeThrough = pipeLock.flags.pipeThrough; |
| 4155 | auto& source = pipeLock.source; |
| 4156 | // At the start of each pipe step, we check to see if either the source or |
| 4157 | // the destination has closed or errored and propagate that on to the other. |
| 4158 | KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) { |
| 4159 | lock.releasePipeLock(); |
| 4160 | return kj::mv(promise); |
| 4161 | } |
| 4162 | |
| 4163 | KJ_IF_SOME(errored, pipeLock.source.tryGetErrored(js)) { |
| 4164 | source.release(js); |
| 4165 | lock.releasePipeLock(); |
| 4166 | if (!preventAbort) { |
| 4167 | auto onSuccess = JSG_VISITABLE_LAMBDA( |
| 4168 | (pipeThrough, reason = js.v8Ref(errored)), (reason), (jsg::Lock& js) { |
| 4169 | return rejectedMaybeHandledPromise<void>(js, reason.getHandle(js), pipeThrough); |
| 4170 | }); |
| 4171 | auto promise = abort(js, errored); |
| 4172 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 4173 | return promise.then(js, ioContext.addFunctor(kj::mv(onSuccess))); |
| 4174 | } else { |
| 4175 | return promise.then(js, kj::mv(onSuccess)); |
| 4176 | } |
| 4177 | } |
| 4178 | return rejectedMaybeHandledPromise<void>(js, errored, pipeThrough); |
| 4179 | } |
| 4180 | |
| 4181 | KJ_IF_SOME(errored, state.tryGetUnsafe<StreamStates::Errored>()) { |
| 4182 | lock.releasePipeLock(); |
| 4183 | auto reason = errored.getHandle(js); |
| 4184 | if (!preventCancel) { |
| 4185 | source.release(js, reason); |
| 4186 | } else { |
| 4187 | source.release(js); |
| 4188 | } |
| 4189 | return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough); |
| 4190 | } |
| 4191 | |
| 4192 | KJ_IF_SOME(erroring, isErroring(js)) { |
| 4193 | lock.releasePipeLock(); |
| 4194 | if (!preventCancel) { |
| 4195 | source.release(js, erroring); |
| 4196 | } else { |
| 4197 | source.release(js); |
| 4198 | } |
| 4199 | return rejectedMaybeHandledPromise<void>(js, erroring, pipeThrough); |
| 4200 | } |
| 4201 | |
| 4202 | if (source.isClosed()) { |
| 4203 | source.release(js); |
| 4204 | lock.releasePipeLock(); |
| 4205 | if (!preventClose) { |
| 4206 | auto promise = close(js); |
| 4207 | if (pipeThrough) { |
| 4208 | promise.markAsHandled(js); |
| 4209 | } |
| 4210 | return kj::mv(promise); |
| 4211 | } |
| 4212 | return js.resolvedPromise(); |
| 4213 | } |
| 4214 | |
| 4215 | if (state.is<StreamStates::Closed>()) { |
| 4216 | lock.releasePipeLock(); |
| 4217 | auto reason = js.v8TypeError("This destination writable stream is closed."_kj); |
| 4218 | if (!preventCancel) { |
| 4219 | source.release(js, reason); |
| 4220 | } else { |
| 4221 | source.release(js); |
| 4222 | } |
| 4223 | |
| 4224 | return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough); |
| 4225 | } |
| 4226 | |
| 4227 | // Assuming we get by that, we perform a read on the source. If the read errors, |
| 4228 | // we propagate the error to the destination, depending on options and reject |
| 4229 | // the pipe promise. If the read is successful then we'll get a ReadResult |
| 4230 | // back. If the ReadResult indicates done, then we close the destination |
| 4231 | // depending on options and resolve the pipe promise. If the ReadResult is |
| 4232 | // not done, we write the value on to the destination. If the write operation |
| 4233 | // fails, we reject the pipe promise and propagate the error back to the |
| 4234 | // source (again, depending on options). If the write operation is successful, |
| 4235 | // we call pipeLoop again to move on to the next iteration. |
| 4236 | |
| 4237 | auto onSuccess = JSG_VISITABLE_LAMBDA((this, ref = addRef(), preventCancel, pipeThrough), (ref), |
| 4238 | (jsg::Lock & js, ReadResult result)->jsg::Promise<void> { |
| 4239 | auto maybePipeLock = lock.tryGetPipe(); |
| 4240 | if (maybePipeLock == kj::none) return js.resolvedPromise(); |
| 4241 | auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock); |
| 4242 | |
| 4243 | KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) { |
| 4244 | lock.releasePipeLock(); |
| 4245 | return kj::mv(promise); |
| 4246 | } else { |
| 4247 | } // Trailing else() is squash compiler warning |
| 4248 | |
| 4249 | if (result.done) { |
| 4250 | // We'll handle the close at the start of the next iteration. |
| 4251 | return pipeLoop(js); |
| 4252 | } |
| 4253 | |
| 4254 | auto onSuccess = JSG_VISITABLE_LAMBDA( |
| 4255 | (this, ref=addRef()), (ref) , (jsg::Lock& js) { |
| 4256 | return pipeLoop(js); |
| 4257 | } ); |
| 4258 | |
| 4259 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 4260 | (this, ref=addRef(), preventCancel, pipeThrough), |
| 4261 | (ref) , (jsg::Lock& js, jsg::Value value) { |
| 4262 | // The write failed. We need to release the source if the pipe lock still exists. |
| 4263 | auto reason = value.getHandle(js); |
| 4264 | KJ_IF_SOME(pipeLock, lock.tryGetPipe()) { |
| 4265 | if (!preventCancel) { |
| 4266 | pipeLock.source.release(js, reason); |
| 4267 | } else { |
| 4268 | pipeLock.source.release(js); |
| 4269 | } |
| 4270 | } else {} // Trailing else() to squash compiler warning |
| 4271 | return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough); |
| 4272 | } ); |
| 4273 | |
| 4274 | auto promise = |
| 4275 | write(js, result.value.map([&](jsg::Value& value) { return value.getHandle(js); })); |
| 4276 | |
| 4277 | return maybeAddFunctor(js, kj::mv(promise), kj::mv(onSuccess), kj::mv(onFailure)); |
| 4278 | }); |
| 4279 | |
| 4280 | auto onFailure = |
| 4281 | JSG_VISITABLE_LAMBDA((this, ref = addRef()), (ref), (jsg::Lock& js, jsg::Value value) { |
| 4282 | // The read failed. We will handle the error at the start of the next iteration. |
| 4283 | return pipeLoop(js); |
| 4284 | }); |
| 4285 | |
| 4286 | return maybeAddFunctor(js, pipeLock.source.read(js), kj::mv(onSuccess), kj::mv(onFailure)); |
| 4287 | } |
| 4288 | |
| 4289 | void WritableStreamJsController::updateBackpressure(jsg::Lock& js, bool backpressure) { |
| 4290 | KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) { |
| 4291 | if (backpressure) { |
| 4292 | // Per the spec, when backpressure is updated and is true, we replace the existing |
| 4293 | // ready promise on the writer with a new pending promise, regardless of whether |
| 4294 | // the existing one is resolved or not. |
| 4295 | auto prp = js.newPromiseAndResolver<void>(); |
| 4296 | prp.promise.markAsHandled(js); |
| 4297 | return writerLock.setReadyFulfiller(js, prp); |
| 4298 | } |
| 4299 | |
| 4300 | // When backpressure is updated and is false, we resolve the ready promise on the writer |
| 4301 | maybeResolvePromise(js, writerLock.getReadyFulfiller()); |
| 4302 | } |
| 4303 | } |
| 4304 | |
| 4305 | jsg::Promise<void> WritableStreamJsController::write( |
| 4306 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) { |
| 4307 | KJ_SWITCH_ONEOF(state) { |
| 4308 | KJ_CASE_ONEOF(initial, Initial) { |
| 4309 | return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj)); |
| 4310 | } |
| 4311 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 4312 | return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj)); |
| 4313 | } |
| 4314 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 4315 | return js.rejectedPromise<void>(errored.addRef(js)); |
| 4316 | } |
| 4317 | KJ_CASE_ONEOF(controller, Controller) { |
| 4318 | return controller->write(js, value.orDefault([&] { return js.undefined(); })); |
| 4319 | } |
| 4320 | } |
| 4321 | KJ_UNREACHABLE; |
| 4322 | } |
| 4323 | |
| 4324 | void WritableStreamJsController::visitForGc(jsg::GcVisitor& visitor) { |
| 4325 | state.visitForGc(visitor); |
| 4326 | visitor.visit(maybeAbortPromise, lock); |
| 4327 | } |
| 4328 | |
| 4329 | // ======================================================================================= |
| 4330 | |
| 4331 | TransformStreamDefaultController::TransformStreamDefaultController(jsg::Lock& js) |
| 4332 | : ioContext(tryGetIoContext()), |
| 4333 | startPromise(js.newPromiseAndResolver<void>()) {} |
| 4334 | |
| 4335 | kj::Maybe<int> TransformStreamDefaultController::getDesiredSize() { |
| 4336 | KJ_IF_SOME(readableController, tryGetReadableController()) { |
| 4337 | return readableController.getDesiredSize(); |
| 4338 | } |
| 4339 | return kj::none; |
| 4340 | } |
| 4341 | |
| 4342 | void TransformStreamDefaultController::enqueue(jsg::Lock& js, v8::Local<v8::Value> chunk) { |
| 4343 | auto& readableController = JSG_REQUIRE_NONNULL(tryGetReadableController(), TypeError, |
| 4344 | "The readable side of this TransformStream is no longer readable."); |
| 4345 | // Hold a strong reference to the readable controller for the duration of this |
| 4346 | // method. The readableController.enqueue() call below invokes the user-provided |
| 4347 | // size algorithm, which can re-enter JS and call error() on this transform |
| 4348 | // controller, dropping the jsg::Ref held by this->readable and the one held by |
| 4349 | // the ReadableStreamJsController's ValueReadable. Without this ref the |
| 4350 | // ReadableStreamDefaultController would be freed while its enqueue() method is |
| 4351 | // still on the stack. |
| 4352 | auto readableControllerRef = kj::addRef(readableController); |
| 4353 | |
| 4354 | JSG_REQUIRE(readableController.canCloseOrEnqueue(), TypeError, |
| 4355 | "The readable side of this TransformStream is no longer readable."); |
| 4356 | js.tryCatch([&] { readableController.enqueue(js, chunk); }, [&](jsg::Value exception) { |
| 4357 | errorWritableAndUnblockWrite(js, exception.getHandle(js)); |
| 4358 | js.throwException(kj::mv(exception)); |
| 4359 | }); |
| 4360 | |
| 4361 | // If the controller was errored during the enqueue (e.g. by the size callback |
| 4362 | // calling error()), skip the backpressure update โ the stream is already torn down. |
| 4363 | if (!readableController.canCloseOrEnqueue()) { |
| 4364 | return; |
| 4365 | } |
| 4366 | |
| 4367 | bool newBackpressure = readableController.hasBackpressure(); |
| 4368 | if (newBackpressure != backpressure) { |
| 4369 | KJ_ASSERT(newBackpressure); |
| 4370 | // Unfortunately the original implementation forgot to actually set the backpressure |
| 4371 | // here so the backpressure signaling failed to work correctly. This is unfortunate |
| 4372 | // because applying the backpressure here could break existing code, so we need to |
| 4373 | // put the fix behind a compat flag. Doh! |
| 4374 | if (FeatureFlags::get(js).getFixupTransformStreamBackpressure()) { |
| 4375 | setBackpressure(js, true); |
| 4376 | } |
| 4377 | } |
| 4378 | } |
| 4379 | |
| 4380 | void TransformStreamDefaultController::error(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 4381 | KJ_IF_SOME(readableController, tryGetReadableController()) { |
| 4382 | readableController.error(js, reason); |
| 4383 | readable = kj::none; |
| 4384 | } |
| 4385 | errorWritableAndUnblockWrite(js, reason); |
| 4386 | } |
| 4387 | |
| 4388 | void TransformStreamDefaultController::terminate(jsg::Lock& js) { |
| 4389 | KJ_IF_SOME(readableController, tryGetReadableController()) { |
| 4390 | readableController.close(js); |
| 4391 | readable = kj::none; |
| 4392 | } |
| 4393 | errorWritableAndUnblockWrite(js, js.v8TypeError("The transform stream has been terminated"_kj)); |
| 4394 | } |
| 4395 | |
| 4396 | jsg::Promise<void> TransformStreamDefaultController::write( |
| 4397 | jsg::Lock& js, v8::Local<v8::Value> chunk) { |
| 4398 | KJ_IF_SOME(writableController, tryGetWritableController()) { |
| 4399 | KJ_IF_SOME(error, writableController.isErroredOrErroring(js)) { |
| 4400 | return js.rejectedPromise<void>(error); |
| 4401 | } |
| 4402 | |
| 4403 | KJ_ASSERT(writableController.isWritable()); |
| 4404 | |
| 4405 | if (backpressure) { |
| 4406 | auto chunkRef = js.v8Ref(chunk); |
| 4407 | return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js).then(js, |
| 4408 | JSG_VISITABLE_LAMBDA((chunkRef = kj::mv(chunkRef), ref=JSG_THIS), |
| 4409 | (chunkRef, ref), (jsg::Lock& js) mutable -> jsg::Promise<void> { |
| 4410 | KJ_IF_SOME(writableController, ref->tryGetWritableController()) { |
| 4411 | KJ_IF_SOME(error, writableController.isErroring(js)) { |
| 4412 | return js.rejectedPromise<void>(error); |
| 4413 | } else { |
| 4414 | // Else block to avert dangling else compiler warning. |
| 4415 | } |
| 4416 | } else { |
| 4417 | // Else block to avert dangling else compiler warning. |
| 4418 | } |
| 4419 | return ref->performTransform(js, chunkRef.getHandle(js)); |
| 4420 | })); |
| 4421 | } |
| 4422 | return performTransform(js, chunk); |
| 4423 | } else { |
| 4424 | return js.rejectedPromise<void>( |
| 4425 | KJ_EXCEPTION(FAILED, "jsg.TypeError: Writing to the TransformStream failed.")); |
| 4426 | } |
| 4427 | } |
| 4428 | |
| 4429 | jsg::Promise<void> TransformStreamDefaultController::abort( |
| 4430 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 4431 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4432 | // If a finish operation is already in progress, return the existing promise |
| 4433 | // or handle the case where we're being called synchronously from within another |
| 4434 | // finish operation. |
| 4435 | if (algorithms.finishStarted) { |
| 4436 | KJ_IF_SOME(finish, algorithms.maybeFinish) { |
| 4437 | return finish.whenResolved(js); |
| 4438 | } |
| 4439 | // finishStarted is true but maybeFinish is not set yet - this means we're being |
| 4440 | // called synchronously from within another finish operation (like cancel). |
| 4441 | // We need to error the stream with the abort reason so that both the current |
| 4442 | // operation and this abort reject with the abort reason. |
| 4443 | error(js, reason); |
| 4444 | return js.rejectedPromise<void>(js.v8Ref(reason)); |
| 4445 | } |
| 4446 | |
| 4447 | // Mark that we're starting a finish operation before running the algorithm. |
| 4448 | algorithms.finishStarted = true; |
| 4449 | } else { |
| 4450 | KJ_IF_SOME(finish, algorithms.maybeFinish) { |
| 4451 | return finish.whenResolved(js); |
| 4452 | } |
| 4453 | } |
| 4454 | |
| 4455 | return algorithms.maybeFinish |
| 4456 | .emplace(maybeRunAlgorithm(js, algorithms.cancel, |
| 4457 | JSG_VISITABLE_LAMBDA( |
| 4458 | (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason), |
| 4459 | (jsg::Lock & js)->jsg::Promise<void> { |
| 4460 | // If the readable side is errored, return a rejected promise with the stored error |
| 4461 | { |
| 4462 | KJ_IF_SOME(err, getReadableErrorState(js)) { |
| 4463 | return js.rejectedPromise<void>(kj::mv(err)); |
| 4464 | } else { |
| 4465 | // Else block to avert dangling else compiler warning. |
| 4466 | } |
| 4467 | } |
| 4468 | // Otherwise... error with the given reason and resolve the abort promise |
| 4469 | error(js, reason.getHandle(js)); |
| 4470 | return js.resolvedPromise(); |
| 4471 | }), |
| 4472 | JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref), |
| 4473 | (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> { |
| 4474 | error(js, reason.getHandle(js)); |
| 4475 | return js.rejectedPromise<void>(kj::mv(reason)); |
| 4476 | }), |
| 4477 | jsg::JsValue(reason))) |
| 4478 | .whenResolved(js); |
| 4479 | } |
| 4480 | |
| 4481 | jsg::Promise<void> TransformStreamDefaultController::close(jsg::Lock& js) { |
| 4482 | auto flags = FeatureFlags::get(js); |
| 4483 | if (flags.getPedanticWpt()) { |
| 4484 | // If a finish operation is already in progress (e.g., from cancel or abort), |
| 4485 | // we should not run flush. Per the WHATWG streams spec, close/flush should |
| 4486 | // coordinate with cancel to avoid calling both. |
| 4487 | if (algorithms.finishStarted) { |
| 4488 | KJ_IF_SOME(finish, algorithms.maybeFinish) { |
| 4489 | return finish.whenResolved(js); |
| 4490 | } |
| 4491 | // finishStarted is true but maybeFinish is not set yet - this means we're being |
| 4492 | // called synchronously from within another finish operation. If the stream was |
| 4493 | // errored during that operation, return a rejected promise with the error. |
| 4494 | KJ_IF_SOME(writableController, tryGetWritableController()) { |
| 4495 | KJ_IF_SOME(err, writableController.isErroredOrErroring(js)) { |
| 4496 | return js.rejectedPromise<void>(err); |
| 4497 | } |
| 4498 | } |
| 4499 | KJ_IF_SOME(err, getReadableErrorState(js)) { |
| 4500 | return js.rejectedPromise<void>(kj::mv(err)); |
| 4501 | } |
| 4502 | return js.resolvedPromise(); |
| 4503 | } |
| 4504 | |
| 4505 | // Mark that we're starting a finish operation before running the algorithm, |
| 4506 | // since the algorithm may synchronously call other finish operations. |
| 4507 | algorithms.finishStarted = true; |
| 4508 | } |
| 4509 | |
| 4510 | auto onSuccess = |
| 4511 | JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), (jsg::Lock & js)->jsg::Promise<void> { |
| 4512 | // If the stream was errored during the flush algorithm (e.g., by controller.error() |
| 4513 | // or by a parallel cancel() calling abort()), we should reject with that error. |
| 4514 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4515 | KJ_IF_SOME(err, ref->getReadableErrorState(js)) { |
| 4516 | return js.rejectedPromise<void>(kj::mv(err)); |
| 4517 | } else { |
| 4518 | // Else block to avert dangling else compiler warning. |
| 4519 | } |
| 4520 | } |
| 4521 | // Allows for a graceful close of the readable side. Close will |
| 4522 | // complete once all of the queued data is read or the stream |
| 4523 | // errors. Only close if the stream can still be closed (e.g., |
| 4524 | // it wasn't closed by a cancel operation from within flush). |
| 4525 | { |
| 4526 | KJ_IF_SOME(readableController, ref->tryGetReadableController()) { |
| 4527 | if (readableController.canCloseOrEnqueue()) { |
| 4528 | readableController.close(js); |
| 4529 | } |
| 4530 | } else { |
| 4531 | // Else block to avert dangling else compiler warning. |
| 4532 | } |
| 4533 | } |
| 4534 | return js.resolvedPromise(); |
| 4535 | }); |
| 4536 | |
| 4537 | auto onFailure = JSG_VISITABLE_LAMBDA( |
| 4538 | (ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> { |
| 4539 | ref->error(js, reason.getHandle(js)); |
| 4540 | return js.rejectedPromise<void>(kj::mv(reason)); |
| 4541 | }); |
| 4542 | |
| 4543 | if (flags.getPedanticWpt()) { |
| 4544 | return algorithms.maybeFinish |
| 4545 | .emplace( |
| 4546 | maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS)) |
| 4547 | .whenResolved(js); |
| 4548 | } |
| 4549 | |
| 4550 | return maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS); |
| 4551 | } |
| 4552 | |
| 4553 | jsg::Promise<void> TransformStreamDefaultController::pull(jsg::Lock& js) { |
| 4554 | KJ_ASSERT(backpressure); |
| 4555 | setBackpressure(js, false); |
| 4556 | return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js); |
| 4557 | } |
| 4558 | |
| 4559 | jsg::Promise<void> TransformStreamDefaultController::cancel( |
| 4560 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 4561 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4562 | // If a finish operation is already in progress, return the existing promise |
| 4563 | // or check for errors if we're being called synchronously from within another |
| 4564 | // finish operation. |
| 4565 | if (algorithms.finishStarted) { |
| 4566 | KJ_IF_SOME(finish, algorithms.maybeFinish) { |
| 4567 | return finish.whenResolved(js); |
| 4568 | } |
| 4569 | // finishStarted is true but maybeFinish is not set yet - check if the stream |
| 4570 | // was errored during that operation. |
| 4571 | KJ_IF_SOME(err, getReadableErrorState(js)) { |
| 4572 | return js.rejectedPromise<void>(kj::mv(err)); |
| 4573 | } |
| 4574 | return js.resolvedPromise(); |
| 4575 | } |
| 4576 | |
| 4577 | // Mark that we're starting a finish operation before running the algorithm. |
| 4578 | algorithms.finishStarted = true; |
| 4579 | } |
| 4580 | |
| 4581 | return algorithms.maybeFinish |
| 4582 | .emplace(maybeRunAlgorithm(js, algorithms.cancel, |
| 4583 | JSG_VISITABLE_LAMBDA( |
| 4584 | (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason), |
| 4585 | (jsg::Lock & js)->jsg::Promise<void> { |
| 4586 | // If the stream was errored during the cancel algorithm (e.g., by controller.error() |
| 4587 | // or by a parallel abort()), we should reject with that error. |
| 4588 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4589 | KJ_IF_SOME(err, getReadableErrorState(js)) { |
| 4590 | readable = kj::none; |
| 4591 | errorWritableAndUnblockWrite(js, reason.getHandle(js)); |
| 4592 | return js.rejectedPromise<void>(kj::mv(err)); |
| 4593 | } else { |
| 4594 | // Else block to avert dangling else compiler warning. |
| 4595 | } |
| 4596 | } |
| 4597 | readable = kj::none; |
| 4598 | errorWritableAndUnblockWrite(js, reason.getHandle(js)); |
| 4599 | return js.resolvedPromise(); |
| 4600 | }), |
| 4601 | JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref), |
| 4602 | (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> { |
| 4603 | readable = kj::none; |
| 4604 | errorWritableAndUnblockWrite(js, reason.getHandle(js)); |
| 4605 | return js.rejectedPromise<void>(kj::mv(reason)); |
| 4606 | }), |
| 4607 | jsg::JsValue(reason))) |
| 4608 | .whenResolved(js); |
| 4609 | } |
| 4610 | |
| 4611 | jsg::Promise<void> TransformStreamDefaultController::performTransform( |
| 4612 | jsg::Lock& js, v8::Local<v8::Value> chunk) { |
| 4613 | if (algorithms.transform != kj::none) { |
| 4614 | return maybeRunAlgorithm(js, algorithms.transform, |
| 4615 | [](jsg::Lock& js) -> jsg::Promise<void> { return js.resolvedPromise(); }, |
| 4616 | JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), |
| 4617 | (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> { |
| 4618 | ref->error(js, reason.getHandle(js)); |
| 4619 | return js.rejectedPromise<void>(kj::mv(reason)); |
| 4620 | }), |
| 4621 | chunk, JSG_THIS); |
| 4622 | } |
| 4623 | // If we got here, there is no transform algorithm. Per the spec, the default |
| 4624 | // behavior then is to just pass along the value untransformed. |
| 4625 | return js.tryCatch([&] { |
| 4626 | enqueue(js, chunk); |
| 4627 | return js.resolvedPromise(); |
| 4628 | }, [&](jsg::Value exception) { return js.rejectedPromise<void>(kj::mv(exception)); }); |
| 4629 | } |
| 4630 | |
| 4631 | void TransformStreamDefaultController::setBackpressure(jsg::Lock& js, bool newBackpressure) { |
| 4632 | KJ_ASSERT(newBackpressure != backpressure); |
| 4633 | KJ_IF_SOME(prp, maybeBackpressureChange) { |
| 4634 | prp.resolver.resolve(js); |
| 4635 | } |
| 4636 | maybeBackpressureChange = js.newPromiseAndResolver<void>(); |
| 4637 | KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.markAsHandled(js); |
| 4638 | backpressure = newBackpressure; |
| 4639 | } |
| 4640 | |
| 4641 | void TransformStreamDefaultController::errorWritableAndUnblockWrite( |
| 4642 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 4643 | algorithms.clear(); |
| 4644 | KJ_IF_SOME(writableController, tryGetWritableController()) { |
| 4645 | if (FeatureFlags::get(js).getPedanticWpt()) { |
| 4646 | // Use errorIfNeeded which goes through the proper error transition (Erroring -> Errored). |
| 4647 | // This allows close() to be called while the stream is "erroring" and reject with the |
| 4648 | // stored error, which is the expected behavior per the WHATWG streams spec. |
| 4649 | writableController.errorIfNeeded(js, reason); |
| 4650 | } else if (writableController.isWritable()) { |
| 4651 | writableController.doError(js, reason); |
| 4652 | } |
| 4653 | writable = kj::none; |
| 4654 | } |
| 4655 | if (backpressure) { |
| 4656 | setBackpressure(js, false); |
| 4657 | } |
| 4658 | } |
| 4659 | |
| 4660 | void TransformStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { |
| 4661 | KJ_IF_SOME(backpressureChange, maybeBackpressureChange) { |
| 4662 | visitor.visit(backpressureChange.promise, backpressureChange.resolver); |
| 4663 | } |
| 4664 | visitor.visit(writable, readable, startPromise.resolver, startPromise.promise, algorithms); |
| 4665 | } |
| 4666 | |
| 4667 | void TransformStreamDefaultController::init(jsg::Lock& js, |
| 4668 | jsg::Ref<ReadableStream>& readable, |
| 4669 | jsg::Ref<WritableStream>& writable, |
| 4670 | jsg::Optional<Transformer> maybeTransformer) { |
| 4671 | KJ_ASSERT(this->readable == kj::none); |
| 4672 | KJ_ASSERT(this->writable == kj::none); |
| 4673 | |
| 4674 | this->writable = writable.addRef(); |
| 4675 | |
| 4676 | // The TransformStreamDefaultController needs to have a reference to the underlying controller |
| 4677 | // and not just the readable because if the readable is teed, or passed off to source, etc, |
| 4678 | // the TransformStream has to make sure that it can continue to interface with the controller |
| 4679 | // to push data into it. |
| 4680 | auto& readableController = static_cast<ReadableStreamJsController&>(readable->getController()); |
| 4681 | auto readableRef = KJ_ASSERT_NONNULL(readableController.getController()); |
| 4682 | this->readable = KJ_ASSERT_NONNULL(readableRef.tryGet<DefaultController>()).addRef(); |
| 4683 | |
| 4684 | auto transformer = kj::mv(maybeTransformer).orDefault({}); |
| 4685 | |
| 4686 | // TODO(someday): The stream standard includes placeholders for supporting byte-oriented |
| 4687 | // TransformStreams but does not yet define them. For now, we are limiting our implementation |
| 4688 | // here to only support value-based transforms. |
| 4689 | JSG_REQUIRE(transformer.readableType == kj::none, TypeError, |
| 4690 | "transformer.readableType must be undefined."); |
| 4691 | JSG_REQUIRE(transformer.writableType == kj::none, TypeError, |
| 4692 | "transformer.writableType must be undefined."); |
| 4693 | |
| 4694 | KJ_IF_SOME(transform, transformer.transform) { |
| 4695 | algorithms.transform = kj::mv(transform); |
| 4696 | } |
| 4697 | |
| 4698 | KJ_IF_SOME(flush, transformer.flush) { |
| 4699 | algorithms.flush = kj::mv(flush); |
| 4700 | } |
| 4701 | |
| 4702 | KJ_IF_SOME(cancel, transformer.cancel) { |
| 4703 | algorithms.cancel = kj::mv(cancel); |
| 4704 | } |
| 4705 | |
| 4706 | setBackpressure(js, true); |
| 4707 | |
| 4708 | maybeRunAlgorithm(js, transformer.start, |
| 4709 | JSG_VISITABLE_LAMBDA( |
| 4710 | (ref = JSG_THIS), (ref), (jsg::Lock& js) { ref->startPromise.resolver.resolve(js); }), |
| 4711 | JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), |
| 4712 | (jsg::Lock& js, jsg::Value reason) { |
| 4713 | ref->startPromise.resolver.reject(js, reason.getHandle(js)); |
| 4714 | }), |
| 4715 | JSG_THIS); |
| 4716 | } |
| 4717 | |
| 4718 | kj::Maybe<ReadableStreamDefaultController&> TransformStreamDefaultController:: |
| 4719 | tryGetReadableController() { |
| 4720 | KJ_IF_SOME(controller, readable) { |
| 4721 | return *controller; |
| 4722 | } |
| 4723 | return kj::none; |
| 4724 | } |
| 4725 | |
| 4726 | kj::Maybe<WritableStreamJsController&> TransformStreamDefaultController:: |
| 4727 | tryGetWritableController() { |
| 4728 | KJ_IF_SOME(w, writable) { |
| 4729 | return static_cast<WritableStreamJsController&>(w->getController()); |
| 4730 | } |
| 4731 | return kj::none; |
| 4732 | } |
| 4733 | |
| 4734 | kj::Maybe<jsg::Value> TransformStreamDefaultController::getReadableErrorState(jsg::Lock& js) { |
| 4735 | KJ_IF_SOME(controller, tryGetReadableController()) { |
| 4736 | return controller.getMaybeErrorState(js); |
| 4737 | } |
| 4738 | return kj::none; |
| 4739 | } |
| 4740 | |
| 4741 | template <class Self> |
| 4742 | kj::StringPtr WritableImpl<Self>::jsgGetMemoryName() const { |
| 4743 | return "WritableImpl"_kjc; |
| 4744 | } |
| 4745 | |
| 4746 | template <class Self> |
| 4747 | size_t WritableImpl<Self>::jsgGetMemorySelfSize() const { |
| 4748 | return sizeof(WritableImpl<Self>); |
| 4749 | } |
| 4750 | |
| 4751 | template <class Self> |
| 4752 | void WritableImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4753 | tracker.trackField("signal", signal); |
| 4754 | |
| 4755 | KJ_SWITCH_ONEOF(state) { |
| 4756 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 4757 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 4758 | tracker.trackField("error", error); |
| 4759 | } |
| 4760 | KJ_CASE_ONEOF(erroring, StreamStates::Erroring) { |
| 4761 | tracker.trackField("erroring", erroring.reason); |
| 4762 | } |
| 4763 | KJ_CASE_ONEOF(writable, Writable) {} |
| 4764 | } |
| 4765 | |
| 4766 | tracker.trackField("abortAlgorithm", algorithms.abort); |
| 4767 | tracker.trackField("closeAlgorithm", algorithms.close); |
| 4768 | tracker.trackField("writeAlgorithm", algorithms.write); |
| 4769 | tracker.trackField("sizeAlgorithm", algorithms.size); |
| 4770 | |
| 4771 | for (auto& request: writeRequests) { |
| 4772 | tracker.trackField("pendingWrite", request); |
| 4773 | } |
| 4774 | |
| 4775 | tracker.trackField("inFlightWrite", inFlightWrite); |
| 4776 | tracker.trackField("inFlightClose", inFlightClose); |
| 4777 | tracker.trackField("closeRequest", closeRequest); |
| 4778 | tracker.trackField("maybePendingAbort", maybePendingAbort); |
| 4779 | } |
| 4780 | |
| 4781 | kj::StringPtr WritableStreamJsController::jsgGetMemoryName() const { |
| 4782 | return "WritableStreamJsController"_kjc; |
| 4783 | } |
| 4784 | |
| 4785 | size_t WritableStreamJsController::jsgGetMemorySelfSize() const { |
| 4786 | return sizeof(WritableStreamJsController); |
| 4787 | } |
| 4788 | |
| 4789 | void WritableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4790 | KJ_SWITCH_ONEOF(state) { |
| 4791 | KJ_CASE_ONEOF(initial, Initial) {} |
| 4792 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 4793 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 4794 | tracker.trackField("error", error); |
| 4795 | } |
| 4796 | KJ_CASE_ONEOF(controller, Controller) { |
| 4797 | tracker.trackField("controller", controller); |
| 4798 | } |
| 4799 | } |
| 4800 | tracker.trackField("lock", lock); |
| 4801 | tracker.trackField("maybeAbortPromise", maybeAbortPromise); |
| 4802 | } |
| 4803 | |
| 4804 | void WritableStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4805 | tracker.trackField("impl", impl); |
| 4806 | } |
| 4807 | |
| 4808 | kj::StringPtr ReadableStreamJsController::jsgGetMemoryName() const { |
| 4809 | return "ReadableStreamJsController"_kjc; |
| 4810 | } |
| 4811 | |
| 4812 | size_t ReadableStreamJsController::jsgGetMemorySelfSize() const { |
| 4813 | return sizeof(ReadableStreamJsController); |
| 4814 | } |
| 4815 | |
| 4816 | void ReadableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4817 | KJ_SWITCH_ONEOF(state) { |
| 4818 | KJ_CASE_ONEOF(initial, Initial) {} |
| 4819 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 4820 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 4821 | tracker.trackField("error", error); |
| 4822 | } |
| 4823 | KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) { |
| 4824 | tracker.trackField("readable", readable); |
| 4825 | } |
| 4826 | KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) { |
| 4827 | tracker.trackField("readable", readable); |
| 4828 | } |
| 4829 | } |
| 4830 | |
| 4831 | tracker.trackField("lock", lock); |
| 4832 | |
| 4833 | // Track pending error state if present (Closed has no trackable content) |
| 4834 | KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) { |
| 4835 | tracker.trackField("pendingError", pendingError); |
| 4836 | } |
| 4837 | } |
| 4838 | |
| 4839 | template <class Self> |
| 4840 | kj::StringPtr ReadableImpl<Self>::jsgGetMemoryName() const { |
| 4841 | return "ReadableImpl"_kjc; |
| 4842 | } |
| 4843 | |
| 4844 | template <class Self> |
| 4845 | size_t ReadableImpl<Self>::jsgGetMemorySelfSize() const { |
| 4846 | return sizeof(ReadableImpl); |
| 4847 | } |
| 4848 | |
| 4849 | template <class Self> |
| 4850 | void ReadableImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4851 | KJ_SWITCH_ONEOF(state) { |
| 4852 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 4853 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 4854 | tracker.trackField("error", error); |
| 4855 | } |
| 4856 | KJ_CASE_ONEOF(queue, Queue) { |
| 4857 | tracker.trackField("queue", queue); |
| 4858 | } |
| 4859 | } |
| 4860 | |
| 4861 | tracker.trackField("startAlgorithm", algorithms.start); |
| 4862 | tracker.trackField("pullAlgorithm", algorithms.pull); |
| 4863 | tracker.trackField("cancelAlgorithm", algorithms.cancel); |
| 4864 | tracker.trackField("sizeAlgorithm", algorithms.size); |
| 4865 | tracker.trackField("pendingCancel", maybePendingCancel); |
| 4866 | } |
| 4867 | |
| 4868 | void ReadableStreamBYOBRequest::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4869 | KJ_IF_SOME(impl, maybeImpl) { |
| 4870 | tracker.trackField("readRequest", impl.readRequest); |
| 4871 | tracker.trackField("view", impl.view); |
| 4872 | } |
| 4873 | } |
| 4874 | |
| 4875 | void TransformStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 4876 | tracker.trackField("startPromise", startPromise); |
| 4877 | tracker.trackField("maybeBackpressureChange", maybeBackpressureChange); |
| 4878 | tracker.trackField("transformAlgorithm", algorithms.transform); |
| 4879 | tracker.trackField("flushAlgorithm", algorithms.flush); |
| 4880 | tracker.trackField("writable", writable); |
| 4881 | tracker.trackField("readable", readable); |
| 4882 | } |
| 4883 | |
| 4884 | // ====================================================================================== |
| 4885 | |
| 4886 | jsg::Ref<ReadableStream> ReadableStream::from( |
| 4887 | jsg::Lock& js, jsg::AsyncGenerator<jsg::Value> generator) { |
| 4888 | |
| 4889 | // AsyncGenerator is not a refcounted type, so we need to wrap it in a refcounted |
| 4890 | // struct so that we can keep it alive through the various promise branches below. |
| 4891 | auto rcGenerator = |
| 4892 | kj::rc<kj::RefcountedWrapper<jsg::AsyncGenerator<jsg::Value>>>(kj::mv(generator)); |
| 4893 | |
| 4894 | // clang-format off |
| 4895 | return constructor(js, UnderlyingSource{ |
| 4896 | .pull = [generator = rcGenerator.addRef()](jsg::Lock& js, auto controller) mutable { |
| 4897 | auto& c = controller.template get<DefaultController>(); |
| 4898 | return generator->getWrapped().next(js).then(js, |
| 4899 | JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()), |
| 4900 | (controller), |
| 4901 | (jsg::Lock& js, kj::Maybe<jsg::Value> value) { |
| 4902 | KJ_IF_SOME(v, value) { |
| 4903 | auto handle = v.getHandle(js); |
| 4904 | // Per the ReadableStream.from spec, if the value is a promise, |
| 4905 | // the stream should wait for it to resolve and enqueue the |
| 4906 | // resolved value... |
| 4907 | // ... yes, this means that ReadableStream.from where the inputs |
| 4908 | // are promises will be slow, but that's the spec. |
| 4909 | if (handle->IsPromise()) { |
| 4910 | return js.toPromise(handle.As<v8::Promise>()).then(js, |
| 4911 | JSG_VISITABLE_LAMBDA( |
| 4912 | (controller=controller.addRef()), |
| 4913 | (controller), |
| 4914 | (jsg::Lock& js, jsg::Value val) mutable { |
| 4915 | controller->enqueue(js, val.getHandle(js)); |
| 4916 | return js.resolvedPromise(); |
| 4917 | })); |
| 4918 | } |
| 4919 | controller->enqueue(js, v.getHandle(js)); |
| 4920 | } else { |
| 4921 | controller->close(js); |
| 4922 | } |
| 4923 | return js.resolvedPromise(); |
| 4924 | }), |
| 4925 | JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()), |
| 4926 | (controller), (jsg::Lock& js, jsg::Value reason) { |
| 4927 | controller->error(js, reason.getHandle(js)); |
| 4928 | return js.rejectedPromise<void>(kj::mv(reason)); |
| 4929 | })); |
| 4930 | }, |
| 4931 | .cancel = [generator = rcGenerator.addRef()](jsg::Lock& js, auto reason) mutable { |
| 4932 | return generator->getWrapped().return_(js, js.v8Ref(reason)) |
| 4933 | .then(js, [generator = kj::mv(generator)](auto& lock, auto) { |
| 4934 | // The generator might produce a value on return and might even want to continue, |
| 4935 | // but the stream has been canceled at this point, so we stop here. |
| 4936 | }); |
| 4937 | }, |
| 4938 | }, StreamQueuingStrategy{ .highWaterMark = 0 }); |
| 4939 | // clang-format on |
| 4940 | } |
| 4941 | |
| 4942 | } // namespace workerd::api |