File
Blob: src/workerd/api/streams/internal.c++
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #include "internal.h" |
| 6 | |
| 7 | #include "identity-transform-stream.h" |
| 8 | #include "readable.h" |
| 9 | #include "writable.h" |
| 10 | |
| 11 | #include <workerd/api/util.h> |
| 12 | #include <workerd/io/features.h> |
| 13 | #include <workerd/jsg/jsg.h> |
| 14 | #include <workerd/util/autogate.h> |
| 15 | #include <workerd/util/string-buffer.h> |
| 16 | |
| 17 | #include <kj/vector.h> |
| 18 | |
| 19 | namespace workerd::api { |
| 20 | |
| 21 | namespace { |
| 22 | // Use this in places where the exception thrown would cause finalizers to run. Your exception |
| 23 | // will not go anywhere, but we'll log the exception message to the console until the problem this |
| 24 | // papers over is fixed. |
| 25 | [[noreturn]] void throwTypeErrorAndConsoleWarn(kj::StringPtr message) { |
| 26 | KJ_IF_SOME(context, IoContext::tryCurrent()) { |
| 27 | if (context.hasWarningHandler()) { |
| 28 | context.logWarning(message); |
| 29 | } |
| 30 | } |
| 31 | |
| 32 | kj::throwFatalException(kj::Exception(kj::Exception::Type::FAILED, __FILE__, __LINE__, |
| 33 | kj::str(JSG_EXCEPTION(TypeError) ": ", message))); |
| 34 | } |
| 35 | |
| 36 | kj::Promise<void> pumpTo(ReadableStreamSource& input, WritableStreamSink& output, bool end) { |
| 37 | kj::byte buffer[65536]{}; |
| 38 | |
| 39 | while (true) { |
| 40 | auto amount = co_await input.tryRead(buffer, 1, kj::size(buffer)); |
| 41 | |
| 42 | if (amount == 0) { |
| 43 | if (end) { |
| 44 | co_await output.end(); |
| 45 | } |
| 46 | co_return; |
| 47 | } |
| 48 | |
| 49 | co_await output.write(kj::arrayPtr(buffer, amount)); |
| 50 | } |
| 51 | } |
| 52 | |
| 53 | // Modified from AllReader in kj/async-io.c++. |
| 54 | class AllReader final { |
| 55 | public: |
| 56 | explicit AllReader(ReadableStreamSource& input, uint64_t limit): input(input), limit(limit) { |
| 57 | JSG_REQUIRE(limit > 0, TypeError, "Memory limit exceeded before EOF."); |
| 58 | KJ_IF_SOME(length, input.tryGetLength(StreamEncoding::IDENTITY)) { |
| 59 | // Oh hey, we might be able to bail early. |
| 60 | JSG_REQUIRE(length < limit, TypeError, "Memory limit would be exceeded before EOF."); |
| 61 | } |
| 62 | } |
| 63 | KJ_DISALLOW_COPY_AND_MOVE(AllReader); |
| 64 | |
| 65 | kj::Promise<kj::Array<kj::byte>> readAllBytes() { |
| 66 | return read<kj::byte>(); |
| 67 | } |
| 68 | |
| 69 | kj::Promise<kj::String> readAllText( |
| 70 | ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE) { |
| 71 | auto data = co_await read<char>(option); |
| 72 | co_return kj::String(kj::mv(data)); |
| 73 | } |
| 74 | |
| 75 | private: |
| 76 | ReadableStreamSource& input; |
| 77 | uint64_t limit; |
| 78 | |
| 79 | template <typename T> |
| 80 | kj::Promise<kj::Array<T>> read(ReadAllTextOption option = ReadAllTextOption::NONE) { |
| 81 | // There are a few complexities in this operation that make it difficult to completely |
| 82 | // optimize. The most important is that even if a stream reports an expected length |
| 83 | // using tryGetLength, we really don't know how much data the stream will produce until |
| 84 | // we try to read it. The only signal we have that the stream is done producing data |
| 85 | // is a zero-length result from tryRead. Unfortunately, we have to allocate a buffer |
| 86 | // in advance of calling tryRead so we have to guess a bit at the size of the buffer |
| 87 | // to allocate. |
| 88 | // |
| 89 | // In the previous implementation of this method, we would just blindly allocate a |
| 90 | // 4096 byte buffer on every allocation, limiting each read iteration to a maximum |
| 91 | // of 4096 bytes. This works fine for streams producing a small amount of data but |
| 92 | // risks requiring a greater number of loop iterations and small allocations for streams |
| 93 | // that produce larger amounts of data. Also in the previous implementation, every |
| 94 | // loop iteration would allocate a new buffer regardless of how much of the previous |
| 95 | // allocation was actually used -- so a stream that produces only 4000 bytes total |
| 96 | // but only provides 10 bytes per iteration would end up with 400 reads and 400 4096 |
| 97 | // byte allocations. Doh! Fortunately our stream implementations tend to be a bit |
| 98 | // smarter than that but it's still a worst case possibility that it's likely better |
| 99 | // to avoid. |
| 100 | // |
| 101 | // So this implementation does things a bit differently. |
| 102 | // First, we check to see if the stream can give an estimate on how much data it |
| 103 | // expects to produce. If that length is within a given threshold, then best case |
| 104 | // is we can perform the entire read with at most two allocations and two calls to |
| 105 | // tryRead. The first allocation will be for the entire expected size of the stream, |
| 106 | // which the first tryRead will attempt to fulfill completely. In the best case the |
| 107 | // stream provides all of the data. The next allocation would be smaller and would |
| 108 | // end up resulting in a zero-length read signaling that we are done. Hooray! |
| 109 | // |
| 110 | // Not everything can be best case scenario tho, unfortunately. If our first tryRead |
| 111 | // does not fully consume the stream or fully fill the destination buffer, we're |
| 112 | // going to need to try again. It is possible that the new allocation in the next |
| 113 | // iteration will be wasted if the stream doesn't have any more data so it's important |
| 114 | // for us to try to be conservative with the allocation. If the running total of data |
| 115 | // we've seen so far is equal to or greater than the expected total length of the stream, |
| 116 | // then the most likely case is that the next read will be zero-length -- but unfortunately |
| 117 | // we can't know for sure! So for this we will fall back to a more conservative allocation |
| 118 | // which is either MIN_BUFFER_CHUNK or the calculated amountToRead, whichever is the lower |
| 119 | // number. |
| 120 | // |
| 121 | // The chunk sizes here are intentionally large to avoid pathological allocation patterns |
| 122 | // when reading from tee'd streams. The KJ tee buffers data in 16KB chunks; if we read |
| 123 | // with smaller buffers, each partial consume of a tee chunk allocates a new heap array |
| 124 | // for the remainder. With many green threads contending on tcmalloc, the cumulative |
| 125 | // allocation overhead can exceed watchdog timeouts. Using 128KB default reads ensures |
| 126 | // tee chunks are consumed whole, eliminating the amplification entirely. |
| 127 | |
| 128 | kj::Vector<kj::Array<T>> parts; |
| 129 | uint64_t runningTotal = 0; |
| 130 | static constexpr uint64_t MIN_BUFFER_CHUNK = 65536; // 64KB |
| 131 | static constexpr uint64_t DEFAULT_BUFFER_CHUNK = 131072; // 128KB |
| 132 | static constexpr uint64_t MAX_BUFFER_CHUNK = DEFAULT_BUFFER_CHUNK * 4; // 512KB |
| 133 | |
| 134 | // If we know in advance how much data we'll be reading, then we can attempt to |
| 135 | // optimize the loop here by setting the value specifically so we are only |
| 136 | // allocating at most twice. But, to be safe, let's enforce an upper bound on each |
| 137 | // allocation even if we do know the total. |
| 138 | kj::Maybe<uint64_t> maybeLength = input.tryGetLength(StreamEncoding::IDENTITY); |
| 139 | |
| 140 | // The amountToRead is the regular allocation size we'll use right up until we've |
| 141 | // read the number of expected bytes (if known). This number is calculated as the |
| 142 | // minimum of (limit, MAX_BUFFER_CHUNK, maybeLength or DEFAULT_BUFFER_CHUNK). In |
| 143 | // the best case scenario, this number is calculated such that we can read the |
| 144 | // entire stream in one go if the amount of data is small enough and the stream |
| 145 | // is well behaved. |
| 146 | // If the stream does report a length, once we've read that number of bytes, we'll |
| 147 | // fallback to the conservativeAllocation. |
| 148 | uint64_t amountToRead = |
| 149 | kj::min(limit, kj::min(MAX_BUFFER_CHUNK, maybeLength.orDefault(DEFAULT_BUFFER_CHUNK))); |
| 150 | // amountToRead can be zero if the stream reported a zero-length. While the stream could |
| 151 | // be lying about its length, let's skip reading anything in this case. |
| 152 | if (amountToRead > 0) { |
| 153 | for (;;) { |
| 154 | auto bytes = kj::heapArray<T>(amountToRead); |
| 155 | // Note that we're passing amountToRead as the *minBytes* here so the tryRead should |
| 156 | // attempt to fill the entire buffer. If it doesn't, the implication is that we read |
| 157 | // everything. |
| 158 | uint64_t amount = co_await input.tryRead(bytes.begin(), amountToRead, amountToRead); |
| 159 | KJ_DASSERT(amount <= amountToRead); |
| 160 | |
| 161 | runningTotal += amount; |
| 162 | JSG_REQUIRE(runningTotal < limit, TypeError, "Memory limit exceeded before EOF."); |
| 163 | |
| 164 | if (amount < amountToRead) { |
| 165 | // The stream has indicated that we're all done by returning a value less than the |
| 166 | // full buffer length. |
| 167 | // It is possible/likely that at least some amount of data was written to the buffer. |
| 168 | // In which case we want to add that subset to the parts list here before we exit |
| 169 | // the loop. |
| 170 | if (amount > 0) { |
| 171 | parts.add(bytes.first(amount).attach(kj::mv(bytes))); |
| 172 | } |
| 173 | break; |
| 174 | } |
| 175 | |
| 176 | // Because we specify minSize equal to maxSize in the tryRead above, we should only |
| 177 | // get here if the buffer was completely filled by the read. If it wasn't completely |
| 178 | // filled, that is an indication that the stream is complete which is handled above. |
| 179 | KJ_DASSERT(amount == bytes.size()); |
| 180 | parts.add(kj::mv(bytes)); |
| 181 | |
| 182 | // If the stream provided an expected length and our running total is equal to |
| 183 | // or greater than that length then we assume we're done. |
| 184 | KJ_IF_SOME(length, maybeLength) { |
| 185 | if (runningTotal >= length) { |
| 186 | // We've read everything we expect to read but some streams need to be read |
| 187 | // completely in order to properly finish and other streams might lie (although |
| 188 | // they shouldn't). Sigh. So we're going to make the next allocation potentially |
| 189 | // smaller and keep reading until we get a zero length. In the best case, the next |
| 190 | // read is going to be zero length but we have to try which will require at least |
| 191 | // one additional (potentially wasted) allocation. (If we don't there are multiple |
| 192 | // test failures). |
| 193 | amountToRead = kj::min(MIN_BUFFER_CHUNK, amountToRead); |
| 194 | continue; |
| 195 | } |
| 196 | } |
| 197 | } |
| 198 | } |
| 199 | |
| 200 | KJ_IF_SOME(length, maybeLength) { |
| 201 | if (runningTotal > length) { |
| 202 | // Realistically runningTotal should never be more than length so we'll emit |
| 203 | // a warning if it is just so we know. It would be indicative of a bug somewhere |
| 204 | // in the implementation. |
| 205 | KJ_LOG(WARNING, "ReadableStream provided more data than advertised", runningTotal, length); |
| 206 | } |
| 207 | } |
| 208 | |
| 209 | // Strip UTF-8 BOM if requested |
| 210 | size_t skipBytes = 0; |
| 211 | if ((option & ReadAllTextOption::STRIP_BOM) && parts.size() > 0 && |
| 212 | hasUtf8Bom(parts[0].asBytes())) { |
| 213 | skipBytes = UTF8_BOM_SIZE; |
| 214 | runningTotal -= UTF8_BOM_SIZE; |
| 215 | } |
| 216 | |
| 217 | if (option & ReadAllTextOption::NULL_TERMINATE) { |
| 218 | auto out = kj::heapArray<T>(runningTotal + 1); |
| 219 | out[runningTotal] = '\0'; |
| 220 | copyInto<T>(out, parts.asPtr(), skipBytes); |
| 221 | co_return kj::mv(out); |
| 222 | } |
| 223 | |
| 224 | // As an optimization, if there's only a single part in the list, we can avoid |
| 225 | // further copies. |
| 226 | if (parts.size() == 1) { |
| 227 | co_return kj::mv(parts[0]); |
| 228 | } |
| 229 | |
| 230 | auto out = kj::heapArray<T>(runningTotal); |
| 231 | copyInto<T>(out, parts.asPtr()); |
| 232 | co_return kj::mv(out); |
| 233 | } |
| 234 | |
| 235 | template <typename T> |
| 236 | void copyInto(kj::ArrayPtr<T> out, kj::ArrayPtr<kj::Array<T>> in, size_t skipBytes = 0) { |
| 237 | for (auto& part: in) { |
| 238 | if (out.size() == 0) { |
| 239 | break; |
| 240 | } |
| 241 | // The skipBytes are used to skip the BOM on the first part only. |
| 242 | KJ_DASSERT(skipBytes <= part.size()); |
| 243 | auto slicedPart = skipBytes ? part.slice(skipBytes) : part; |
| 244 | skipBytes = 0; |
| 245 | if (slicedPart.size() == 0) { |
| 246 | continue; |
| 247 | } |
| 248 | KJ_DASSERT(slicedPart.size() <= out.size()); |
| 249 | out.first(slicedPart.size()).copyFrom(slicedPart); |
| 250 | out = out.slice(slicedPart.size()); |
| 251 | } |
| 252 | } |
| 253 | }; |
| 254 | |
| 255 | kj::Exception reasonToException(jsg::Lock& js, |
| 256 | jsg::Optional<v8::Local<v8::Value>> maybeReason, |
| 257 | kj::String defaultDescription = kj::str(JSG_EXCEPTION(Error) ": Stream was cancelled.")) { |
| 258 | KJ_IF_SOME(reason, maybeReason) { |
| 259 | return js.exceptionToKj(js.v8Ref(reason)); |
| 260 | } else { |
| 261 | // We get here if the caller is something like `r.cancel()` (or `r.cancel(undefined)`). |
| 262 | return kj::Exception( |
| 263 | kj::Exception::Type::FAILED, __FILE__, __LINE__, kj::mv(defaultDescription)); |
| 264 | } |
| 265 | } |
| 266 | |
| 267 | // ======================================================================================= |
| 268 | |
| 269 | // Adapt ReadableStreamSource to kj::AsyncInputStream's interface for use with `kj::newTee()`. |
| 270 | class TeeAdapter final: public kj::AsyncInputStream { |
| 271 | public: |
| 272 | explicit TeeAdapter(kj::Own<ReadableStreamSource> inner): inner(kj::mv(inner)) {} |
| 273 | |
| 274 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 275 | return inner->tryRead(buffer, minBytes, maxBytes); |
| 276 | } |
| 277 | |
| 278 | kj::Maybe<uint64_t> tryGetLength() override { |
| 279 | return inner->tryGetLength(StreamEncoding::IDENTITY); |
| 280 | } |
| 281 | |
| 282 | private: |
| 283 | kj::Own<ReadableStreamSource> inner; |
| 284 | }; |
| 285 | |
| 286 | class TeeBranch final: public ReadableStreamSource { |
| 287 | public: |
| 288 | explicit TeeBranch(kj::Own<kj::AsyncInputStream> inner): inner(kj::mv(inner)) {} |
| 289 | |
| 290 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 291 | return inner->tryRead(buffer, minBytes, maxBytes); |
| 292 | } |
| 293 | |
| 294 | kj::Promise<DeferredProxy<void>> pumpTo(WritableStreamSink& output, bool end) override { |
| 295 | #ifdef KJ_NO_RTTI |
| 296 | // Yes, I'm paranoid. |
| 297 | static_assert(!KJ_NO_RTTI, "Need RTTI for correctness"); |
| 298 | #endif |
| 299 | |
| 300 | // HACK: If `output` is another TransformStream, we don't allow pumping to it, in order to |
| 301 | // guarantee that we can't create cycles. Note that currently TeeBranch only ever wraps |
| 302 | // TransformStreams, never system streams. |
| 303 | JSG_REQUIRE(!isIdentityTransformStream(output), TypeError, |
| 304 | "Inter-TransformStream ReadableStream.pipeTo() is not implemented."); |
| 305 | |
| 306 | // It is important we actually call `inner->pumpTo()` so that `kj::newTee()` is aware of this |
| 307 | // pump operation's backpressure. So we can't use the default `ReadableStreamSource::pumpTo()` |
| 308 | // implementation, and have to implement our own. |
| 309 | |
| 310 | PumpAdapter outputAdapter(output); |
| 311 | co_await inner->pumpTo(outputAdapter); |
| 312 | |
| 313 | if (end) { |
| 314 | co_await output.end(); |
| 315 | } |
| 316 | |
| 317 | // We only use `TeeBranch` when a locally-sourced stream was tee'd (because system streams |
| 318 | // implement `tryTee()` in a different way that doesn't use `TeeBranch`). So, we know that |
| 319 | // none of the pump can be performed without the IoContext active, and thus we do not |
| 320 | // `KJ_CO_MAGIC BEGIN_DEFERRED_PROXYING`. |
| 321 | co_return; |
| 322 | } |
| 323 | |
| 324 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 325 | if (encoding == StreamEncoding::IDENTITY) { |
| 326 | return inner->tryGetLength(); |
| 327 | } else { |
| 328 | return kj::none; |
| 329 | } |
| 330 | } |
| 331 | |
| 332 | kj::Maybe<Tee> tryTee(uint64_t limit) override { |
| 333 | KJ_IF_SOME(t, inner->tryTee(limit)) { |
| 334 | auto branch = kj::heap<TeeBranch>(newTeeErrorAdapter(kj::mv(t))); |
| 335 | auto consumed = kj::heap<TeeBranch>(kj::mv(inner)); |
| 336 | return Tee{kj::mv(branch), kj::mv(consumed)}; |
| 337 | } |
| 338 | |
| 339 | return kj::none; |
| 340 | } |
| 341 | |
| 342 | void cancel(kj::Exception reason) override { |
| 343 | // TODO(someday): What to do? |
| 344 | } |
| 345 | |
| 346 | private: |
| 347 | // Adapt WritableStreamSink to kj::AsyncOutputStream's interface for use in |
| 348 | // `TeeBranch::pumpTo()`. If you squint, the write logic looks very similar to TeeAdapter's |
| 349 | // read logic. |
| 350 | class PumpAdapter final: public kj::AsyncOutputStream { |
| 351 | public: |
| 352 | explicit PumpAdapter(WritableStreamSink& inner): inner(inner) {} |
| 353 | |
| 354 | kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override { |
| 355 | return inner.write(buffer); |
| 356 | } |
| 357 | |
| 358 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override { |
| 359 | return inner.write(pieces); |
| 360 | } |
| 361 | |
| 362 | kj::Promise<void> whenWriteDisconnected() override { |
| 363 | KJ_UNIMPLEMENTED("whenWriteDisconnected() not expected on PumpAdapter"); |
| 364 | } |
| 365 | |
| 366 | WritableStreamSink& inner; |
| 367 | }; |
| 368 | |
| 369 | kj::Own<kj::AsyncInputStream> inner; |
| 370 | }; |
| 371 | } // namespace |
| 372 | |
| 373 | // ======================================================================================= |
| 374 | |
| 375 | kj::Promise<DeferredProxy<void>> ReadableStreamSource::pumpTo( |
| 376 | WritableStreamSink& output, bool end) { |
| 377 | KJ_IF_SOME(p, output.tryPumpFrom(*this, end)) { |
| 378 | return kj::mv(p); |
| 379 | } |
| 380 | |
| 381 | // Non-optimized pumpTo() is presumed to require the IoContext to remain live, so don't do |
| 382 | // anything in the deferred proxy part. |
| 383 | return addNoopDeferredProxy(api::pumpTo(*this, output, end)); |
| 384 | } |
| 385 | |
| 386 | kj::Maybe<uint64_t> ReadableStreamSource::tryGetLength(StreamEncoding encoding) { |
| 387 | return kj::none; |
| 388 | } |
| 389 | |
| 390 | kj::Promise<kj::Array<byte>> ReadableStreamSource::readAllBytes(uint64_t limit) { |
| 391 | try { |
| 392 | AllReader allReader(*this, limit); |
| 393 | co_return co_await allReader.readAllBytes(); |
| 394 | } catch (...) { |
| 395 | // TODO(soon): Temporary logging. |
| 396 | auto ex = kj::getCaughtExceptionAsKj(); |
| 397 | if (ex.getDescription().endsWith("exceeded before EOF.")) { |
| 398 | LOG_WARNING_PERIODICALLY("NOSENTRY Internal Stream readAllBytes - Exceeded limit"); |
| 399 | } |
| 400 | kj::throwFatalException(kj::mv(ex)); |
| 401 | } |
| 402 | } |
| 403 | |
| 404 | kj::Promise<kj::String> ReadableStreamSource::readAllText( |
| 405 | uint64_t limit, ReadAllTextOption option) { |
| 406 | try { |
| 407 | AllReader allReader(*this, limit); |
| 408 | co_return co_await allReader.readAllText(option); |
| 409 | } catch (...) { |
| 410 | // TODO(soon): Temporary logging. |
| 411 | auto ex = kj::getCaughtExceptionAsKj(); |
| 412 | if (ex.getDescription().endsWith("exceeded before EOF.")) { |
| 413 | LOG_WARNING_PERIODICALLY("NOSENTRY Internal Stream readAllText - Exceeded limit"); |
| 414 | } |
| 415 | kj::throwFatalException(kj::mv(ex)); |
| 416 | } |
| 417 | } |
| 418 | |
| 419 | void ReadableStreamSource::cancel(kj::Exception reason) {} |
| 420 | |
| 421 | kj::Maybe<ReadableStreamSource::Tee> ReadableStreamSource::tryTee(uint64_t limit) { |
| 422 | return kj::none; |
| 423 | } |
| 424 | |
| 425 | kj::Maybe<kj::Promise<DeferredProxy<void>>> WritableStreamSink::tryPumpFrom( |
| 426 | ReadableStreamSource& input, bool end) { |
| 427 | return kj::none; |
| 428 | } |
| 429 | |
| 430 | // ======================================================================================= |
| 431 | |
| 432 | ReadableStreamInternalController::~ReadableStreamInternalController() noexcept(false) { |
| 433 | if (readState.is<ReaderLocked>()) { |
| 434 | readState.transitionTo<Unlocked>(); |
| 435 | } |
| 436 | } |
| 437 | |
| 438 | jsg::Ref<ReadableStream> ReadableStreamInternalController::addRef() { |
| 439 | return KJ_ASSERT_NONNULL(owner).addRef(); |
| 440 | } |
| 441 | |
| 442 | kj::Maybe<jsg::Promise<ReadResult>> ReadableStreamInternalController::read( |
| 443 | jsg::Lock& js, kj::Maybe<ByobOptions> maybeByobOptions) { |
| 444 | |
| 445 | if (isPendingClosure) { |
| 446 | return js.rejectedPromise<ReadResult>( |
| 447 | js.v8TypeError("This ReadableStream belongs to an object that is closing."_kj)); |
| 448 | } |
| 449 | |
| 450 | v8::Local<v8::ArrayBuffer> store; |
| 451 | size_t byteLength = 0; |
| 452 | size_t byteOffset = 0; |
| 453 | size_t atLeast = 1; |
| 454 | |
| 455 | KJ_IF_SOME(byobOptions, maybeByobOptions) { |
| 456 | store = byobOptions.bufferView.getHandle(js)->Buffer(); |
| 457 | byteOffset = byobOptions.byteOffset; |
| 458 | byteLength = byobOptions.byteLength; |
| 459 | atLeast = byobOptions.atLeast.orDefault(atLeast); |
| 460 | if (byobOptions.detachBuffer) { |
| 461 | if (!store->IsDetachable()) { |
| 462 | return js.rejectedPromise<ReadResult>( |
| 463 | js.v8TypeError("Unable to use non-detachable ArrayBuffer"_kj)); |
| 464 | } |
| 465 | auto backing = store->GetBackingStore(); |
| 466 | jsg::check(store->Detach(v8::Local<v8::Value>())); |
| 467 | store = v8::ArrayBuffer::New(js.v8Isolate, kj::mv(backing)); |
| 468 | } |
| 469 | } |
| 470 | |
| 471 | auto getOrInitStore = [&](bool errorCase = false) { |
| 472 | if (store.IsEmpty()) { |
| 473 | if (errorCase) { |
| 474 | byteLength = 0; |
| 475 | } else if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) { |
| 476 | byteLength = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2; |
| 477 | } else { |
| 478 | byteLength = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE; |
| 479 | } |
| 480 | |
| 481 | if (!v8::ArrayBuffer::MaybeNew(js.v8Isolate, byteLength).ToLocal(&store)) { |
| 482 | return v8::Local<v8::ArrayBuffer>(); |
| 483 | } |
| 484 | } |
| 485 | return store; |
| 486 | }; |
| 487 | |
| 488 | disturbed = true; |
| 489 | |
| 490 | KJ_SWITCH_ONEOF(state) { |
| 491 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 492 | if (maybeByobOptions != kj::none && FeatureFlags::get(js).getInternalStreamByobReturn()) { |
| 493 | // When using the BYOB reader, we must return a sized-0 Uint8Array that is backed |
| 494 | // by the ArrayBuffer passed in the options. |
| 495 | auto theStore = getOrInitStore(true); |
| 496 | if (theStore.IsEmpty()) { |
| 497 | return js.rejectedPromise<ReadResult>( |
| 498 | js.v8TypeError("Unable to allocate memory for read"_kj)); |
| 499 | } |
| 500 | return js.resolvedPromise(ReadResult{ |
| 501 | .value = js.v8Ref(v8::Uint8Array::New(theStore, 0, 0).As<v8::Value>()), |
| 502 | .done = true, |
| 503 | }); |
| 504 | } |
| 505 | return js.resolvedPromise(ReadResult{.done = true}); |
| 506 | } |
| 507 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 508 | return js.rejectedPromise<ReadResult>(errored.addRef(js)); |
| 509 | } |
| 510 | KJ_CASE_ONEOF(readable, Readable) { |
| 511 | // TODO(conform): Requiring serialized read requests is non-conformant, but we've never had a |
| 512 | // use case for them. At one time, our implementation of TransformStream supported multiple |
| 513 | // simultaneous read requests, but it is highly unlikely that anyone relied on this. Our |
| 514 | // ReadableStream implementation that wraps native streams has never supported them, our |
| 515 | // TransformStream implementation is primarily (only?) used for constructing manually |
| 516 | // streamed Responses, and no teed ReadableStream has ever supported them. |
| 517 | if (readPending) { |
| 518 | return js.rejectedPromise<ReadResult>(js.v8TypeError( |
| 519 | "This ReadableStream only supports a single pending read request at a time."_kj)); |
| 520 | } |
| 521 | readPending = true; |
| 522 | |
| 523 | auto theStore = getOrInitStore(); |
| 524 | if (theStore.IsEmpty()) { |
| 525 | return js.rejectedPromise<ReadResult>( |
| 526 | js.v8TypeError("Unable to allocate memory for read"_kj)); |
| 527 | } |
| 528 | |
| 529 | // In the case the ArrayBuffer is detached/transfered while the read is pending, we |
| 530 | // need to make sure that the ptr remains stable, so we grab a shared ptr to the |
| 531 | // backing store and use that to get the pointer to the data. If the buffer is detached |
| 532 | // while the read is pending, this does mean that the read data will end up being lost, |
| 533 | // but there's not really a better option. The best we can do here is warn the user |
| 534 | // that this is happening so they can avoid doing it in the future. |
| 535 | // Also, the user really shouldn't do this because the read will end up completing into |
| 536 | // the detached backing store still which could cause issues with whatever code now actually |
| 537 | // owns the transfered buffer. Below we'll warn the user about this if it happens so they |
| 538 | // can avoid doing it in the future. |
| 539 | auto backing = theStore->GetBackingStore(); |
| 540 | |
| 541 | // For resizable ArrayBuffers, the buffer may be resized while the read is |
| 542 | // pending, decommitting memory pages and making the pointer invalid (SIGSEGV). |
| 543 | // We read into a temporary buffer and copy the data back in the .then() |
| 544 | // callback, where we can validate the buffer is still large enough. |
| 545 | bool isResizable = theStore->IsResizableByUserJavaScript(); |
| 546 | |
| 547 | kj::Array<kj::byte> tempBuffer; |
| 548 | kj::byte* readPtr; |
| 549 | if (isResizable) { |
| 550 | auto currentByteLength = theStore->ByteLength(); |
| 551 | if (byteOffset >= currentByteLength) { |
| 552 | readPending = false; |
| 553 | return js.resolvedPromise(ReadResult{ |
| 554 | .value = js.v8Ref(v8::Uint8Array::New(theStore, 0, 0).As<v8::Value>()), |
| 555 | .done = false, |
| 556 | }); |
| 557 | } |
| 558 | if (byteOffset + byteLength > currentByteLength) { |
| 559 | byteLength = currentByteLength - byteOffset; |
| 560 | if (atLeast > byteLength) { |
| 561 | atLeast = byteLength > 0 ? byteLength : 1; |
| 562 | } |
| 563 | } |
| 564 | tempBuffer = kj::heapArray<kj::byte>(byteLength); |
| 565 | readPtr = tempBuffer.begin(); |
| 566 | } else { |
| 567 | auto ptr = static_cast<kj::byte*>(backing->Data()); |
| 568 | readPtr = ptr + byteOffset; |
| 569 | } |
| 570 | auto bytes = kj::arrayPtr(readPtr, byteLength); |
| 571 | |
| 572 | KJ_ASSERT(atLeast <= bytes.size(), "minBytes must not exceed maxBytes in tryRead"); |
| 573 | |
| 574 | auto promise = kj::evalNow([&] { |
| 575 | return readable->tryRead(bytes.begin(), atLeast, bytes.size()).attach(kj::mv(backing)); |
| 576 | }); |
| 577 | KJ_IF_SOME(readerLock, readState.tryGetUnsafe<ReaderLocked>()) { |
| 578 | promise = KJ_ASSERT_NONNULL(readerLock.getCanceler())->wrap(kj::mv(promise)); |
| 579 | } |
| 580 | |
| 581 | // TODO(soon): We use awaitIoLegacy() here because if the stream terminates in JavaScript in |
| 582 | // this same isolate, then the promise may actually be waiting on JavaScript to do something, |
| 583 | // and so should not be considered waiting on external I/O. We will need to use |
| 584 | // registerPendingEvent() manually when reading from an external stream. Ideally, we would |
| 585 | // refactor the implementation so that when waiting on a JavaScript stream, we strictly use |
| 586 | // jsg::Promises and not kj::Promises, so that it doesn't look like I/O at all, and there's |
| 587 | // no need to drop the isolate lock and take it again every time some data is read/written. |
| 588 | // That's a larger refactor, though. |
| 589 | auto& ioContext = IoContext::current(); |
| 590 | return ioContext.awaitIoLegacy(js, kj::mv(promise)) |
| 591 | .then(js, ioContext.addFunctor(JSG_VISITABLE_LAMBDA( |
| 592 | (this, ref = addRef(), store = js.v8Ref(store), |
| 593 | byteOffset, byteLength, isByob = maybeByobOptions != kj::none, |
| 594 | isResizable, readPtr, tempBuffer = kj::mv(tempBuffer)), |
| 595 | (ref), |
| 596 | (jsg::Lock& js, size_t amount) mutable -> jsg::Promise<ReadResult> { |
| 597 | readPending = false; |
| 598 | KJ_ASSERT(amount <= byteLength); |
| 599 | if (amount == 0) { |
| 600 | if (!state.is<StreamStates::Errored>()) { |
| 601 | doClose(js); |
| 602 | } |
| 603 | KJ_IF_SOME(o, owner) { |
| 604 | o.signalEof(js); |
| 605 | } else {} |
| 606 | if (isByob && FeatureFlags::get(js).getInternalStreamByobReturn()) { |
| 607 | // When using the BYOB reader, we must return a sized-0 Uint8Array that is backed |
| 608 | // by the ArrayBuffer passed in the options. |
| 609 | auto u8 = v8::Uint8Array::New(store.getHandle(js), 0, 0); |
| 610 | return js.resolvedPromise(ReadResult{ |
| 611 | .value = js.v8Ref(u8.As<v8::Value>()), |
| 612 | .done = true, |
| 613 | }); |
| 614 | } |
| 615 | return js.resolvedPromise(ReadResult{.done = true}); |
| 616 | } |
| 617 | // Return a slice so the script can see how many bytes were read. |
| 618 | |
| 619 | // We have to check to see if the store was detached or resized while we were waiting |
| 620 | // for the read to complete. |
| 621 | auto handle = store.getHandle(js); |
| 622 | if (handle->WasDetached()) { |
| 623 | // If the buffer was detached, we resolve with a new zero-length ArrayBuffer. |
| 624 | // The bytes that were read are lost, but this is a valid result. |
| 625 | |
| 626 | // Silly user, trix are for kids. |
| 627 | IoContext::current().logWarningOnce( |
| 628 | "A buffer that was being used for a read operation on a ReadableStream was detached " |
| 629 | "while the read was pending. The read completed with a zero-length buffer and the data " |
| 630 | "that was read is lost. Avoid detaching buffers that are being used for active read " |
| 631 | "operations on streams, or use the streams_byob_reader_detaches_buffer compatibility " |
| 632 | "flag, to prevent this from happening."_kj); |
| 633 | |
| 634 | auto buffer = v8::ArrayBuffer::New(js.v8Isolate, 0); |
| 635 | return js.resolvedPromise(ReadResult{ |
| 636 | .value = js.v8Ref(v8::Uint8Array::New(buffer, 0, 0).As<v8::Value>()), |
| 637 | .done = false, |
| 638 | }); |
| 639 | } |
| 640 | |
| 641 | if (byteOffset + amount > handle->ByteLength()) { |
| 642 | // If the buffer was resized smaller, we return a truncated result. |
| 643 | |
| 644 | IoContext::current().logWarningOnce( |
| 645 | "A buffer that was being used for a read operation on a ReadableStream was resized " |
| 646 | "smaller while the read was pending. The read completed with a truncated buffer " |
| 647 | "containing only the bytes that fit within the new size. Avoid resizing buffers that " |
| 648 | "are being used for active read operations on streams, or use the " |
| 649 | "streams_byob_reader_detaches_buffer compatibility flag, to prevent this from " |
| 650 | "happening."_kj); |
| 651 | |
| 652 | if (byteOffset >= handle->ByteLength()) { |
| 653 | return js.resolvedPromise(ReadResult{ |
| 654 | .value = js.v8Ref(v8::Uint8Array::New(store.getHandle(js), 0, 0).As<v8::Value>()), |
| 655 | .done = false, |
| 656 | }); |
| 657 | } |
| 658 | amount = handle->ByteLength() - byteOffset; |
| 659 | } |
| 660 | |
| 661 | if (isResizable && byteOffset + amount <= handle->ByteLength()) { |
| 662 | // For resizable buffers, the data was read into a temporary buffer. |
| 663 | // Copy it back into the user's (still valid) buffer region. |
| 664 | auto destPtr = static_cast<kj::byte*>(handle->GetBackingStore()->Data()); |
| 665 | memcpy(destPtr + byteOffset, readPtr, amount); |
| 666 | } |
| 667 | |
| 668 | return js.resolvedPromise(ReadResult{ |
| 669 | .value = js.v8Ref( |
| 670 | v8::Uint8Array::New(store.getHandle(js), byteOffset, amount).As<v8::Value>()), |
| 671 | .done = false, |
| 672 | }); |
| 673 | })), |
| 674 | ioContext.addFunctor(JSG_VISITABLE_LAMBDA( |
| 675 | (this, ref = addRef()), |
| 676 | (ref), |
| 677 | (jsg::Lock& js, jsg::Value reason) -> jsg::Promise<ReadResult> { |
| 678 | readPending = false; |
| 679 | if (!state.is<StreamStates::Errored>()) { |
| 680 | doError(js, reason.getHandle(js)); |
| 681 | } |
| 682 | return js.rejectedPromise<ReadResult>(kj::mv(reason)); |
| 683 | }))); |
| 684 | } |
| 685 | } |
| 686 | KJ_UNREACHABLE; |
| 687 | } |
| 688 | |
| 689 | kj::Maybe<jsg::Promise<DrainingReadResult>> ReadableStreamInternalController::drainingRead( |
| 690 | jsg::Lock& js, size_t maxRead) { |
| 691 | // InternalController does not support draining reads fully since all reads are |
| 692 | // async. We implement a simplified version that just performs a normal read |
| 693 | // like read(). The significant difference is that with JS-backed streams, a draining |
| 694 | // read will pull any already enqueued data from the stream buffer and try synchronously |
| 695 | // pumping the stream for more data until either maxRead is satisfied or the stream |
| 696 | // indicates EOF, error, or that it needs to wait for more data. Internal streams have |
| 697 | // no such internal buffering and never provide data synchronously so drainingRead |
| 698 | // is effectively the same as read(). |
| 699 | |
| 700 | if (isPendingClosure) { |
| 701 | return js.rejectedPromise<DrainingReadResult>( |
| 702 | js.v8TypeError("This ReadableStream belongs to an object that is closing."_kj)); |
| 703 | } |
| 704 | |
| 705 | static constexpr size_t kAtLeast = 1; |
| 706 | |
| 707 | disturbed = true; |
| 708 | |
| 709 | KJ_SWITCH_ONEOF(state) { |
| 710 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 711 | return js.resolvedPromise(DrainingReadResult{.done = true}); |
| 712 | } |
| 713 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 714 | return js.rejectedPromise<DrainingReadResult>(errored.addRef(js)); |
| 715 | } |
| 716 | KJ_CASE_ONEOF(readable, Readable) { |
| 717 | if (readPending) { |
| 718 | return js.rejectedPromise<DrainingReadResult>(js.v8TypeError( |
| 719 | "This ReadableStream only supports a single pending read request at a time."_kj)); |
| 720 | } |
| 721 | readPending = true; |
| 722 | |
| 723 | // TODO(later): In the case that maxRead is large, we may consider splitting this into |
| 724 | // multiple reads to avoid allocating too large of a buffer at once. The draining read |
| 725 | // result can handle multiple chunks so this would be feasible at the cost of more |
| 726 | // read calls. For now we just do a single read up to maxRead. |
| 727 | // At the very least, we cap maxRead to some reasonable limit to avoid |
| 728 | // potential OOM issues. |
| 729 | static constexpr size_t kMaxReadCap = 1 * 1024 * 1024; // 1 MB |
| 730 | maxRead = kj::min(maxRead, kMaxReadCap); |
| 731 | |
| 732 | if (maxRead == 0) { |
| 733 | // No data requested, return empty result. |
| 734 | // This really shouldn't ever happen but let's handle it gracefully. |
| 735 | readPending = false; |
| 736 | return js.resolvedPromise(DrainingReadResult{ |
| 737 | .chunks = nullptr, |
| 738 | .done = false, |
| 739 | }); |
| 740 | } |
| 741 | |
| 742 | auto store = kj::heapArray<kj::byte>(maxRead); |
| 743 | |
| 744 | auto promise = |
| 745 | kj::evalNow([&] { return readable->tryRead(store.begin(), kAtLeast, store.size()); }); |
| 746 | KJ_IF_SOME(readerLock, readState.tryGetUnsafe<ReaderLocked>()) { |
| 747 | promise = KJ_ASSERT_NONNULL(readerLock.getCanceler())->wrap(kj::mv(promise)); |
| 748 | } |
| 749 | |
| 750 | auto& ioContext = IoContext::current(); |
| 751 | return ioContext.awaitIoLegacy(js, kj::mv(promise)) |
| 752 | .then(js, ioContext.addFunctor(JSG_VISITABLE_LAMBDA( |
| 753 | (this, ref = addRef(), store = kj::mv(store)), |
| 754 | (ref), |
| 755 | (jsg::Lock& js, size_t amount) mutable -> jsg::Promise<DrainingReadResult> { |
| 756 | readPending = false; |
| 757 | KJ_ASSERT(amount <= store.size()); |
| 758 | if (amount == 0) { |
| 759 | if (!state.is<StreamStates::Errored>()) { |
| 760 | doClose(js); |
| 761 | } |
| 762 | KJ_IF_SOME(o, owner) { |
| 763 | o.signalEof(js); |
| 764 | } else {} |
| 765 | return js.resolvedPromise(DrainingReadResult{.done = true}); |
| 766 | } |
| 767 | // Return a slice so the script can see how many bytes were read. |
| 768 | return js.resolvedPromise(DrainingReadResult{ |
| 769 | .chunks = kj::arr(store.slice(0, amount).attach(kj::mv(store))), .done = false}); |
| 770 | })), |
| 771 | ioContext.addFunctor(JSG_VISITABLE_LAMBDA( |
| 772 | (this, ref = addRef()), |
| 773 | (ref), |
| 774 | (jsg::Lock& js, jsg::Value reason) -> jsg::Promise<DrainingReadResult> { |
| 775 | readPending = false; |
| 776 | if (!state.is<StreamStates::Errored>()) { |
| 777 | doError(js, reason.getHandle(js)); |
| 778 | } |
| 779 | return js.rejectedPromise<DrainingReadResult>(kj::mv(reason)); |
| 780 | }))); |
| 781 | } |
| 782 | } |
| 783 | KJ_UNREACHABLE; |
| 784 | } |
| 785 | |
| 786 | jsg::Promise<void> ReadableStreamInternalController::pipeTo( |
| 787 | jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) { |
| 788 | |
| 789 | KJ_DASSERT(!isLockedToReader()); |
| 790 | KJ_DASSERT(!destination.isLockedToWriter()); |
| 791 | |
| 792 | if (isPendingClosure) { |
| 793 | return js.rejectedPromise<void>( |
| 794 | js.v8TypeError("This ReadableStream belongs to an object that is closing."_kj)); |
| 795 | } |
| 796 | |
| 797 | disturbed = true; |
| 798 | KJ_IF_SOME(promise, |
| 799 | destination.tryPipeFrom(js, KJ_ASSERT_NONNULL(owner).addRef(), kj::mv(options))) { |
| 800 | return kj::mv(promise); |
| 801 | } |
| 802 | |
| 803 | return js.rejectedPromise<void>( |
| 804 | js.v8TypeError("This ReadableStream cannot be piped to this WritableStream."_kj)); |
| 805 | } |
| 806 | |
| 807 | jsg::Promise<void> ReadableStreamInternalController::cancel( |
| 808 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 809 | disturbed = true; |
| 810 | |
| 811 | KJ_IF_SOME(errored, state.tryGetUnsafe<StreamStates::Errored>()) { |
| 812 | return js.rejectedPromise<void>(errored.getHandle(js)); |
| 813 | } |
| 814 | |
| 815 | doCancel(js, maybeReason); |
| 816 | |
| 817 | return js.resolvedPromise(); |
| 818 | } |
| 819 | |
| 820 | void ReadableStreamInternalController::doCancel( |
| 821 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 822 | auto exception = reasonToException(js, maybeReason); |
| 823 | KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) { |
| 824 | KJ_IF_SOME(canceler, locked.getCanceler()) { |
| 825 | canceler->cancel(exception.clone()); |
| 826 | } |
| 827 | } |
| 828 | KJ_IF_SOME(readable, state.tryGetUnsafe<Readable>()) { |
| 829 | readable->cancel(kj::mv(exception)); |
| 830 | doClose(js); |
| 831 | } |
| 832 | } |
| 833 | |
| 834 | void ReadableStreamInternalController::doClose(jsg::Lock& js) { |
| 835 | // If already in a terminal state, nothing to do. |
| 836 | if (state.isTerminal()) return; |
| 837 | |
| 838 | state.transitionTo<StreamStates::Closed>(); |
| 839 | KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) { |
| 840 | maybeResolvePromise(js, locked.getClosedFulfiller()); |
| 841 | } else { |
| 842 | (void)readState.transitionFromTo<PipeLocked, Unlocked>(); |
| 843 | } |
| 844 | } |
| 845 | |
| 846 | void ReadableStreamInternalController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 847 | // If already in a terminal state, nothing to do. |
| 848 | if (state.isTerminal()) return; |
| 849 | |
| 850 | state.transitionTo<StreamStates::Errored>(js.v8Ref(reason)); |
| 851 | KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) { |
| 852 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason); |
| 853 | } else { |
| 854 | (void)readState.transitionFromTo<PipeLocked, Unlocked>(); |
| 855 | } |
| 856 | } |
| 857 | |
| 858 | ReadableStreamController::Tee ReadableStreamInternalController::tee(jsg::Lock& js) { |
| 859 | JSG_REQUIRE( |
| 860 | !isLockedToReader(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 861 | JSG_REQUIRE( |
| 862 | !isPendingClosure, TypeError, "This ReadableStream belongs to an object that is closing."); |
| 863 | readState.transitionTo<Locked>(); |
| 864 | disturbed = true; |
| 865 | KJ_SWITCH_ONEOF(state) { |
| 866 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 867 | // Create two closed ReadableStreams. |
| 868 | return Tee{ |
| 869 | .branch1 = js.alloc<ReadableStream>(kj::heap<ReadableStreamInternalController>(closed)), |
| 870 | .branch2 = js.alloc<ReadableStream>(kj::heap<ReadableStreamInternalController>(closed)), |
| 871 | }; |
| 872 | } |
| 873 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 874 | // Create two errored ReadableStreams. |
| 875 | return Tee{ |
| 876 | .branch1 = js.alloc<ReadableStream>( |
| 877 | kj::heap<ReadableStreamInternalController>(errored.addRef(js))), |
| 878 | .branch2 = js.alloc<ReadableStream>( |
| 879 | kj::heap<ReadableStreamInternalController>(errored.addRef(js))), |
| 880 | }; |
| 881 | } |
| 882 | KJ_CASE_ONEOF(readable, Readable) { |
| 883 | auto& ioContext = IoContext::current(); |
| 884 | |
| 885 | auto makeTee = [&](kj::Own<ReadableStreamSource> b1, |
| 886 | kj::Own<ReadableStreamSource> b2) -> Tee { |
| 887 | doClose(js); |
| 888 | return Tee{ |
| 889 | .branch1 = js.alloc<ReadableStream>(ioContext, kj::mv(b1)), |
| 890 | .branch2 = js.alloc<ReadableStream>(ioContext, kj::mv(b2)), |
| 891 | }; |
| 892 | }; |
| 893 | |
| 894 | auto bufferLimit = ioContext.getLimitEnforcer().getBufferingLimit(); |
| 895 | KJ_IF_SOME(tee, readable->tryTee(bufferLimit)) { |
| 896 | // This ReadableStreamSource has an optimized tee implementation. |
| 897 | return makeTee(kj::mv(tee.branches[0]), kj::mv(tee.branches[1])); |
| 898 | } |
| 899 | |
| 900 | auto tee = kj::newTee(kj::heap<TeeAdapter>(kj::mv(readable)), bufferLimit); |
| 901 | |
| 902 | return makeTee(kj::heap<TeeBranch>(newTeeErrorAdapter(kj::mv(tee.branches[0]))), |
| 903 | kj::heap<TeeBranch>(newTeeErrorAdapter(kj::mv(tee.branches[1])))); |
| 904 | } |
| 905 | } |
| 906 | |
| 907 | KJ_UNREACHABLE; |
| 908 | } |
| 909 | |
| 910 | kj::Maybe<kj::Own<ReadableStreamSource>> ReadableStreamInternalController::removeSource( |
| 911 | jsg::Lock& js, bool ignoreDisturbed) { |
| 912 | JSG_REQUIRE( |
| 913 | !isLockedToReader(), TypeError, "This ReadableStream is currently locked to a reader."); |
| 914 | JSG_REQUIRE(!disturbed || ignoreDisturbed, TypeError, "This ReadableStream is disturbed."); |
| 915 | |
| 916 | readState.transitionTo<Locked>(); |
| 917 | disturbed = true; |
| 918 | |
| 919 | KJ_SWITCH_ONEOF(state) { |
| 920 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 921 | class NullSource final: public ReadableStreamSource { |
| 922 | public: |
| 923 | kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override { |
| 924 | return static_cast<size_t>(0); |
| 925 | } |
| 926 | |
| 927 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override { |
| 928 | return static_cast<uint64_t>(0); |
| 929 | } |
| 930 | }; |
| 931 | |
| 932 | return kj::heap<NullSource>(); |
| 933 | } |
| 934 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 935 | kj::throwFatalException(js.exceptionToKj(errored.addRef(js))); |
| 936 | } |
| 937 | KJ_CASE_ONEOF(readable, Readable) { |
| 938 | auto result = kj::mv(readable); |
| 939 | state.transitionTo<StreamStates::Closed>(); |
| 940 | return kj::Maybe<kj::Own<ReadableStreamSource>>(kj::mv(result)); |
| 941 | } |
| 942 | } |
| 943 | |
| 944 | KJ_UNREACHABLE; |
| 945 | } |
| 946 | |
| 947 | bool ReadableStreamInternalController::lockReader(jsg::Lock& js, Reader& reader) { |
| 948 | if (isLockedToReader()) { |
| 949 | return false; |
| 950 | } |
| 951 | |
| 952 | auto prp = js.newPromiseAndResolver<void>(); |
| 953 | prp.promise.markAsHandled(js); |
| 954 | |
| 955 | auto lock = ReaderLocked( |
| 956 | reader, kj::mv(prp.resolver), IoContext::current().addObject(kj::heap<kj::Canceler>())); |
| 957 | |
| 958 | KJ_SWITCH_ONEOF(state) { |
| 959 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 960 | maybeResolvePromise(js, lock.getClosedFulfiller()); |
| 961 | } |
| 962 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 963 | maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js)); |
| 964 | } |
| 965 | KJ_CASE_ONEOF(readable, Readable) { |
| 966 | // Nothing to do. |
| 967 | } |
| 968 | } |
| 969 | |
| 970 | readState.transitionTo<ReaderLocked>(kj::mv(lock)); |
| 971 | reader.attach(*this, kj::mv(prp.promise)); |
| 972 | return true; |
| 973 | } |
| 974 | |
| 975 | void ReadableStreamInternalController::releaseReader( |
| 976 | Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) { |
| 977 | KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) { |
| 978 | KJ_ASSERT(&locked.getReader() == &reader); |
| 979 | KJ_IF_SOME(js, maybeJs) { |
| 980 | KJ_IF_SOME(canceler, locked.getCanceler()) { |
| 981 | JSG_REQUIRE(canceler->isEmpty(), TypeError, |
| 982 | "Cannot call releaseLock() on a reader with outstanding read promises."); |
| 983 | } |
| 984 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), |
| 985 | js.v8TypeError("This ReadableStream reader has been released."_kj)); |
| 986 | } |
| 987 | locked.clear(); |
| 988 | |
| 989 | // When maybeJs is nullptr, that means releaseReader was called when the reader is |
| 990 | // being deconstructed and not as the result of explicitly calling releaseLock. In |
| 991 | // that case, we don't want to change the lock state itself because we do not have |
| 992 | // an isolate lock. Clearing the lock above will free the lock state while keeping the |
| 993 | // ReadableStream marked as locked. |
| 994 | if (maybeJs != kj::none) { |
| 995 | readState.transitionTo<Unlocked>(); |
| 996 | } |
| 997 | } |
| 998 | } |
| 999 | |
| 1000 | void WritableStreamInternalController::Writable::abort(kj::Exception&& ex) { |
| 1001 | canceler.cancel(ex.clone()); |
| 1002 | sink->abort(kj::mv(ex)); |
| 1003 | } |
| 1004 | |
| 1005 | WritableStreamInternalController::~WritableStreamInternalController() noexcept(false) { |
| 1006 | if (writeState.is<WriterLocked>()) { |
| 1007 | writeState.transitionTo<Unlocked>(); |
| 1008 | } |
| 1009 | } |
| 1010 | |
| 1011 | jsg::Ref<WritableStream> WritableStreamInternalController::addRef() { |
| 1012 | return KJ_ASSERT_NONNULL(owner).addRef(); |
| 1013 | } |
| 1014 | |
| 1015 | jsg::Promise<void> WritableStreamInternalController::write( |
| 1016 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) { |
| 1017 | if (isPendingClosure) { |
| 1018 | return js.rejectedPromise<void>( |
| 1019 | js.v8TypeError("This WritableStream belongs to an object that is closing."_kj)); |
| 1020 | } |
| 1021 | if (isClosedOrClosing()) { |
| 1022 | return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj)); |
| 1023 | } |
| 1024 | if (isPiping()) { |
| 1025 | return js.rejectedPromise<void>( |
| 1026 | js.v8TypeError("This WritableStream is currently being piped to."_kj)); |
| 1027 | } |
| 1028 | |
| 1029 | KJ_SWITCH_ONEOF(state) { |
| 1030 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1031 | // Handled by isClosedOrClosing(). |
| 1032 | KJ_UNREACHABLE; |
| 1033 | } |
| 1034 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1035 | return js.rejectedPromise<void>(errored.addRef(js)); |
| 1036 | } |
| 1037 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1038 | if (value == kj::none) { |
| 1039 | return js.resolvedPromise(); |
| 1040 | } |
| 1041 | auto chunk = KJ_ASSERT_NONNULL(value); |
| 1042 | |
| 1043 | std::shared_ptr<v8::BackingStore> store; |
| 1044 | size_t byteLength = 0; |
| 1045 | size_t byteOffset = 0; |
| 1046 | if (chunk->IsArrayBuffer()) { |
| 1047 | auto buffer = chunk.As<v8::ArrayBuffer>(); |
| 1048 | store = buffer->GetBackingStore(); |
| 1049 | byteLength = buffer->ByteLength(); |
| 1050 | } else if (chunk->IsArrayBufferView()) { |
| 1051 | auto view = chunk.As<v8::ArrayBufferView>(); |
| 1052 | store = view->Buffer()->GetBackingStore(); |
| 1053 | byteLength = view->ByteLength(); |
| 1054 | byteOffset = view->ByteOffset(); |
| 1055 | } else if (chunk->IsString()) { |
| 1056 | // TODO(later): This really ought to return a rejected promise and not a sync throw. |
| 1057 | // This case caused me a moment of confusion during testing, so I think it's worth |
| 1058 | // a specific error message. |
| 1059 | throwTypeErrorAndConsoleWarn( |
| 1060 | "This TransformStream is being used as a byte stream, but received a string on its " |
| 1061 | "writable side. If you wish to write a string, you'll probably want to explicitly " |
| 1062 | "UTF-8-encode it with TextEncoder."); |
| 1063 | } else { |
| 1064 | // TODO(later): This really ought to return a rejected promise and not a sync throw. |
| 1065 | throwTypeErrorAndConsoleWarn( |
| 1066 | "This TransformStream is being used as a byte stream, but received an object of " |
| 1067 | "non-ArrayBuffer/ArrayBufferView type on its writable side."); |
| 1068 | } |
| 1069 | |
| 1070 | if (byteLength == 0) { |
| 1071 | return js.resolvedPromise(); |
| 1072 | } |
| 1073 | |
| 1074 | auto prp = js.newPromiseAndResolver<void>(); |
| 1075 | adjustWriteBufferSize(js, byteLength); |
| 1076 | KJ_IF_SOME(o, observer) { |
| 1077 | o->onChunkEnqueued(byteLength); |
| 1078 | } |
| 1079 | |
| 1080 | auto src = kj::arrayPtr(static_cast<kj::byte*>(store->Data()) + byteOffset, byteLength); |
| 1081 | auto data = kj::heapArray<kj::byte>(src.size()); |
| 1082 | data.asPtr().copyFrom(src); |
| 1083 | auto ptr = data.asPtr(); |
| 1084 | queue.push_back( |
| 1085 | WriteEvent{.outputLock = IoContext::current().waitForOutputLocksIfNecessaryIoOwn(), |
| 1086 | .event = kj::heap<Write>({ |
| 1087 | .promise = kj::mv(prp.resolver), |
| 1088 | .totalBytes = store->ByteLength(), |
| 1089 | .ownBytes = kj::mv(data), |
| 1090 | .bytes = ptr, |
| 1091 | })}); |
| 1092 | |
| 1093 | ensureWriting(js); |
| 1094 | return kj::mv(prp.promise); |
| 1095 | } |
| 1096 | } |
| 1097 | |
| 1098 | KJ_UNREACHABLE; |
| 1099 | } |
| 1100 | |
| 1101 | void WritableStreamInternalController::adjustWriteBufferSize(jsg::Lock& js, int64_t amount) { |
| 1102 | KJ_DASSERT(amount >= 0 || std::abs(amount) <= currentWriteBufferSize); |
| 1103 | currentWriteBufferSize += amount; |
| 1104 | KJ_IF_SOME(highWaterMark, maybeHighWaterMark) { |
| 1105 | int64_t desiredSize = highWaterMark - currentWriteBufferSize; |
| 1106 | updateBackpressure(js, desiredSize <= 0); |
| 1107 | } |
| 1108 | } |
| 1109 | |
| 1110 | void WritableStreamInternalController::updateBackpressure(jsg::Lock& js, bool backpressure) { |
| 1111 | KJ_IF_SOME(writerLock, writeState.tryGetUnsafe<WriterLocked>()) { |
| 1112 | if (backpressure) { |
| 1113 | // Per the spec, when backpressure is updated and is true, we replace the existing |
| 1114 | // ready promise on the writer with a new pending promise, regardless of whether |
| 1115 | // the existing one is resolved or not. |
| 1116 | auto prp = js.newPromiseAndResolver<void>(); |
| 1117 | prp.promise.markAsHandled(js); |
| 1118 | writerLock.setReadyFulfiller(js, prp); |
| 1119 | return; |
| 1120 | } |
| 1121 | |
| 1122 | // When backpressure is updated and is false, we resolve the ready promise on the writer |
| 1123 | maybeResolvePromise(js, writerLock.getReadyFulfiller()); |
| 1124 | } |
| 1125 | } |
| 1126 | |
| 1127 | void WritableStreamInternalController::setHighWaterMark(uint64_t highWaterMark) { |
| 1128 | maybeHighWaterMark = highWaterMark; |
| 1129 | } |
| 1130 | |
| 1131 | jsg::Promise<void> WritableStreamInternalController::closeImpl(jsg::Lock& js, bool markAsHandled) { |
| 1132 | if (isClosedOrClosing()) { |
| 1133 | return js.resolvedPromise(); |
| 1134 | } |
| 1135 | if (isPiping()) { |
| 1136 | auto reason = js.v8TypeError("This WritableStream is currently being piped to."_kj); |
| 1137 | return rejectedMaybeHandledPromise<void>(js, reason, markAsHandled); |
| 1138 | } |
| 1139 | |
| 1140 | KJ_SWITCH_ONEOF(state) { |
| 1141 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1142 | // Handled by isClosedOrClosing(). |
| 1143 | KJ_UNREACHABLE; |
| 1144 | } |
| 1145 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1146 | auto reason = errored.getHandle(js); |
| 1147 | return rejectedMaybeHandledPromise<void>(js, reason, markAsHandled); |
| 1148 | } |
| 1149 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1150 | auto prp = js.newPromiseAndResolver<void>(); |
| 1151 | if (markAsHandled) { |
| 1152 | prp.promise.markAsHandled(js); |
| 1153 | } |
| 1154 | queue.push_back( |
| 1155 | WriteEvent{.outputLock = IoContext::current().waitForOutputLocksIfNecessaryIoOwn(), |
| 1156 | .event = kj::heap<Close>({.promise = kj::mv(prp.resolver)})}); |
| 1157 | ensureWriting(js); |
| 1158 | return kj::mv(prp.promise); |
| 1159 | } |
| 1160 | } |
| 1161 | |
| 1162 | KJ_UNREACHABLE; |
| 1163 | } |
| 1164 | |
| 1165 | jsg::Promise<void> WritableStreamInternalController::close(jsg::Lock& js, bool markAsHandled) { |
| 1166 | KJ_IF_SOME(closureWaitable, maybeClosureWaitable) { |
| 1167 | // If we're already waiting on the closure waitable, then we do not want to try scheduling |
| 1168 | // it again, let's just wait for the existing one to be resolved. |
| 1169 | if (waitingOnClosureWritableAlready) { |
| 1170 | return closureWaitable.whenResolved(js); |
| 1171 | } |
| 1172 | waitingOnClosureWritableAlready = true; |
| 1173 | auto promise = closureWaitable.then(js, [markAsHandled, this](jsg::Lock& js) { |
| 1174 | return closeImpl(js, markAsHandled); |
| 1175 | }, [](jsg::Lock& js, jsg::Value) { |
| 1176 | // Ignore rejection as it will be reported in the Socket's `closed`/`opened` promises |
| 1177 | // instead. |
| 1178 | return js.resolvedPromise(); |
| 1179 | }); |
| 1180 | maybeClosureWaitable = promise.whenResolved(js); |
| 1181 | return kj::mv(promise); |
| 1182 | } else { |
| 1183 | return closeImpl(js, markAsHandled); |
| 1184 | } |
| 1185 | } |
| 1186 | |
| 1187 | jsg::Promise<void> WritableStreamInternalController::flush(jsg::Lock& js, bool markAsHandled) { |
| 1188 | if (isClosedOrClosing()) { |
| 1189 | auto reason = js.v8TypeError("This WritableStream has been closed."_kj); |
| 1190 | return rejectedMaybeHandledPromise<void>(js, reason, markAsHandled); |
| 1191 | } |
| 1192 | if (isPiping()) { |
| 1193 | auto reason = js.v8TypeError("This WritableStream is currently being piped to."_kj); |
| 1194 | return rejectedMaybeHandledPromise<void>(js, reason, markAsHandled); |
| 1195 | } |
| 1196 | |
| 1197 | KJ_SWITCH_ONEOF(state) { |
| 1198 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1199 | // Handled by isClosedOrClosing(). |
| 1200 | KJ_UNREACHABLE; |
| 1201 | } |
| 1202 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1203 | auto reason = errored.getHandle(js); |
| 1204 | return rejectedMaybeHandledPromise<void>(js, reason, markAsHandled); |
| 1205 | } |
| 1206 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1207 | auto prp = js.newPromiseAndResolver<void>(); |
| 1208 | if (markAsHandled) { |
| 1209 | prp.promise.markAsHandled(js); |
| 1210 | } |
| 1211 | queue.push_back( |
| 1212 | WriteEvent{.outputLock = IoContext::current().waitForOutputLocksIfNecessaryIoOwn(), |
| 1213 | .event = kj::heap<Flush>({.promise = kj::mv(prp.resolver)})}); |
| 1214 | ensureWriting(js); |
| 1215 | return kj::mv(prp.promise); |
| 1216 | } |
| 1217 | } |
| 1218 | |
| 1219 | KJ_UNREACHABLE; |
| 1220 | } |
| 1221 | |
| 1222 | jsg::Promise<void> WritableStreamInternalController::abort( |
| 1223 | jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 1224 | // While it may be confusing to users to throw `undefined` rather than a more helpful Error here, |
| 1225 | // doing so is required by the relevant spec: |
| 1226 | // https://streams.spec.whatwg.org/#writable-stream-abort |
| 1227 | return doAbort(js, maybeReason.orDefault(js.v8Undefined())); |
| 1228 | } |
| 1229 | |
| 1230 | jsg::Promise<void> WritableStreamInternalController::doAbort( |
| 1231 | jsg::Lock& js, v8::Local<v8::Value> reason, AbortOptions options) { |
| 1232 | // If maybePendingAbort is set, then the returned abort promise will be rejected |
| 1233 | // with the specified error once the abort is completed, otherwise the promise will |
| 1234 | // be resolved with undefined. |
| 1235 | |
| 1236 | // If there is already an abort pending, return that pending promise |
| 1237 | // instead of trying to schedule another. |
| 1238 | KJ_IF_SOME(pendingAbort, maybePendingAbort) { |
| 1239 | pendingAbort->reject = options.reject; |
| 1240 | auto promise = pendingAbort->whenResolved(js); |
| 1241 | if (options.handled) { |
| 1242 | promise.markAsHandled(js); |
| 1243 | } |
| 1244 | return kj::mv(promise); |
| 1245 | } |
| 1246 | |
| 1247 | KJ_IF_SOME(writable, state.tryGetUnsafe<IoOwn<Writable>>()) { |
| 1248 | auto exception = js.exceptionToKj(js.v8Ref(reason)); |
| 1249 | |
| 1250 | if (FeatureFlags::get(js).getInternalWritableStreamAbortClearsQueue()) { |
| 1251 | // If this flag is set, we will clear the queue proactively and immediately |
| 1252 | // error the stream rather than handling the abort lazily. In this case, the |
| 1253 | // stream will be put into an errored state immediately after draining the |
| 1254 | // queue. All pending writes and other operations in the queue will be rejected |
| 1255 | // immediately and an immediately resolved or rejected promise will be returned. |
| 1256 | writable->abort(exception.clone()); |
| 1257 | drain(js, reason); |
| 1258 | return options.reject ? rejectedMaybeHandledPromise<void>(js, reason, options.handled) |
| 1259 | : js.resolvedPromise(); |
| 1260 | } |
| 1261 | |
| 1262 | if (queue.empty()) { |
| 1263 | writable->abort(exception.clone()); |
| 1264 | doError(js, reason); |
| 1265 | return options.reject ? rejectedMaybeHandledPromise<void>(js, reason, options.handled) |
| 1266 | : js.resolvedPromise(); |
| 1267 | } |
| 1268 | |
| 1269 | maybePendingAbort = kj::heap<PendingAbort>(js, reason, options.reject); |
| 1270 | auto promise = KJ_ASSERT_NONNULL(maybePendingAbort)->whenResolved(js); |
| 1271 | if (options.handled) { |
| 1272 | promise.markAsHandled(js); |
| 1273 | } |
| 1274 | return kj::mv(promise); |
| 1275 | } |
| 1276 | |
| 1277 | return options.reject ? rejectedMaybeHandledPromise<void>(js, reason, options.handled) |
| 1278 | : js.resolvedPromise(); |
| 1279 | } |
| 1280 | |
| 1281 | kj::Maybe<jsg::Promise<void>> WritableStreamInternalController::tryPipeFrom( |
| 1282 | jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) { |
| 1283 | |
| 1284 | // The ReadableStream source here can be either a JavaScript-backed ReadableStream |
| 1285 | // or ReadableStreamSource-backed. |
| 1286 | // |
| 1287 | // If the source is ReadableStreamSource-backed, then we can use kj's low level mechanisms |
| 1288 | // for piping the data. If the source is JavaScript-backed, then we need to rely on the |
| 1289 | // JavaScript-based Promise API for piping the data. |
| 1290 | |
| 1291 | auto preventAbort = options.preventAbort.orDefault(false); |
| 1292 | auto preventClose = options.preventClose.orDefault(false); |
| 1293 | auto preventCancel = options.preventCancel.orDefault(false); |
| 1294 | auto pipeThrough = options.pipeThrough; |
| 1295 | |
| 1296 | if (isPiping()) { |
| 1297 | auto reason = js.v8TypeError("This WritableStream is currently being piped to."_kj); |
| 1298 | return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough); |
| 1299 | } |
| 1300 | |
| 1301 | // If a signal is provided, we need to check that it is not already triggered. If it |
| 1302 | // is, we return a rejected promise using the signal's reason. |
| 1303 | KJ_IF_SOME(signal, options.signal) { |
| 1304 | if (signal->getAborted(js)) { |
| 1305 | return rejectedMaybeHandledPromise<void>(js, signal->getReason(js), pipeThrough); |
| 1306 | } |
| 1307 | } |
| 1308 | |
| 1309 | // With either type of source, our first step is to acquire the source pipe lock. This |
| 1310 | // will help abstract most of the details of which type of source we're working with. |
| 1311 | auto& sourceLock = KJ_ASSERT_NONNULL(source->getController().tryPipeLock()); |
| 1312 | |
| 1313 | // Let's also acquire the destination pipe lock. |
| 1314 | writeState.transitionTo<PipeLocked>(*source); |
| 1315 | |
| 1316 | // If the source has errored, the spec requires us to reject the pipe promise and, if preventAbort |
| 1317 | // is false, error the destination (Propagate error forward). The errored source will be unlocked |
| 1318 | // immediately. The destination will be unlocked once the abort completes. |
| 1319 | KJ_IF_SOME(errored, sourceLock.tryGetErrored(js)) { |
| 1320 | sourceLock.release(js); |
| 1321 | if (!preventAbort) { |
| 1322 | if (state.tryGetUnsafe<IoOwn<Writable>>() != kj::none) { |
| 1323 | return doAbort(js, errored, {.reject = true, .handled = pipeThrough}); |
| 1324 | } |
| 1325 | } |
| 1326 | |
| 1327 | // If preventAbort was true, we're going to unlock the destination now. |
| 1328 | writeState.transitionTo<Unlocked>(); |
| 1329 | return rejectedMaybeHandledPromise<void>(js, errored, pipeThrough); |
| 1330 | } |
| 1331 | |
| 1332 | // If the destination has errored, the spec requires us to reject the pipe promise and, if |
| 1333 | // preventCancel is false, error the source (Propagate error backward). The errored destination |
| 1334 | // will be unlocked immediately. |
| 1335 | KJ_IF_SOME(errored, state.tryGetUnsafe<StreamStates::Errored>()) { |
| 1336 | writeState.transitionTo<Unlocked>(); |
| 1337 | if (!preventCancel) { |
| 1338 | sourceLock.release(js, errored.getHandle(js)); |
| 1339 | } else { |
| 1340 | sourceLock.release(js); |
| 1341 | } |
| 1342 | return rejectedMaybeHandledPromise<void>(js, errored.getHandle(js), pipeThrough); |
| 1343 | } |
| 1344 | |
| 1345 | // If the source has closed, the spec requires us to close the destination if preventClose |
| 1346 | // is false (Propagate closing forward). The source is unlocked immediately. The destination |
| 1347 | // will be unlocked as soon as the close completes. |
| 1348 | if (sourceLock.isClosed()) { |
| 1349 | sourceLock.release(js); |
| 1350 | if (!preventClose) { |
| 1351 | // The spec would have us check to see if `destination` is errored and, if so, return its |
| 1352 | // stored error. But if `destination` were errored, we would already have caught that case |
| 1353 | // above. The spec is probably concerned about cases where the readable and writable sides |
| 1354 | // transition to such states in a racey way. But our pump implementation will take care of |
| 1355 | // this naively. |
| 1356 | KJ_ASSERT(!state.is<StreamStates::Errored>()); |
| 1357 | if (!isClosedOrClosing()) { |
| 1358 | return close(js); |
| 1359 | } |
| 1360 | } |
| 1361 | writeState.transitionTo<Unlocked>(); |
| 1362 | return js.resolvedPromise(); |
| 1363 | } |
| 1364 | |
| 1365 | // If the destination has closed, the spec requires us to close the source if |
| 1366 | // preventCancel is false (Propagate closing backward). |
| 1367 | if (isClosedOrClosing()) { |
| 1368 | auto destClosed = js.v8TypeError("This destination writable stream is closed."_kj); |
| 1369 | writeState.transitionTo<Unlocked>(); |
| 1370 | |
| 1371 | if (!preventCancel) { |
| 1372 | sourceLock.release(js, destClosed); |
| 1373 | } else { |
| 1374 | sourceLock.release(js); |
| 1375 | } |
| 1376 | |
| 1377 | return rejectedMaybeHandledPromise<void>(js, destClosed, pipeThrough); |
| 1378 | } |
| 1379 | |
| 1380 | // The pipe will continue until either the source closes or errors, or until the destination |
| 1381 | // closes or errors. In either case, both will end up being closed or errored, which will |
| 1382 | // release the locks on both. |
| 1383 | // |
| 1384 | // For either type of source, our next step is to wait for the write loop to process the |
| 1385 | // pending Pipe event we queue below. |
| 1386 | auto prp = js.newPromiseAndResolver<void>(); |
| 1387 | if (pipeThrough) { |
| 1388 | prp.promise.markAsHandled(js); |
| 1389 | } |
| 1390 | queue.push_back(WriteEvent{ |
| 1391 | .outputLock = IoContext::current().waitForOutputLocksIfNecessaryIoOwn(), |
| 1392 | .event = kj::heap<Pipe>(*this, sourceLock, kj::mv(prp.resolver), preventAbort, preventClose, |
| 1393 | preventCancel, kj::mv(options.signal)), |
| 1394 | }); |
| 1395 | ensureWriting(js); |
| 1396 | return kj::mv(prp.promise); |
| 1397 | } |
| 1398 | |
| 1399 | kj::Maybe<kj::Own<WritableStreamSink>> WritableStreamInternalController::removeSink(jsg::Lock& js) { |
| 1400 | JSG_REQUIRE( |
| 1401 | !isLockedToWriter(), TypeError, "This WritableStream is currently locked to a writer."); |
| 1402 | JSG_REQUIRE(!isClosedOrClosing(), TypeError, "This WritableStream is closed."); |
| 1403 | |
| 1404 | writeState.transitionTo<Locked>(); |
| 1405 | |
| 1406 | KJ_SWITCH_ONEOF(state) { |
| 1407 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1408 | // Handled by the isClosedOrClosing() check above; |
| 1409 | KJ_UNREACHABLE; |
| 1410 | } |
| 1411 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1412 | kj::throwFatalException(js.exceptionToKj(errored.addRef(js))); |
| 1413 | } |
| 1414 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1415 | auto result = kj::mv(writable->sink); |
| 1416 | state.transitionTo<StreamStates::Closed>(); |
| 1417 | return kj::Maybe<kj::Own<WritableStreamSink>>(kj::mv(result)); |
| 1418 | } |
| 1419 | } |
| 1420 | |
| 1421 | KJ_UNREACHABLE; |
| 1422 | } |
| 1423 | |
| 1424 | void WritableStreamInternalController::detach(jsg::Lock& js) { |
| 1425 | JSG_REQUIRE( |
| 1426 | !isLockedToWriter(), TypeError, "This WritableStream is currently locked to a writer."); |
| 1427 | JSG_REQUIRE(!isClosedOrClosing(), TypeError, "This WritableStream is closed."); |
| 1428 | |
| 1429 | writeState.transitionTo<Locked>(); |
| 1430 | |
| 1431 | KJ_SWITCH_ONEOF(state) { |
| 1432 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1433 | // Handled by the isClosedOrClosing() check above; |
| 1434 | KJ_UNREACHABLE; |
| 1435 | } |
| 1436 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1437 | kj::throwFatalException(js.exceptionToKj(errored.addRef(js))); |
| 1438 | } |
| 1439 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1440 | state.transitionTo<StreamStates::Closed>(); |
| 1441 | return; |
| 1442 | } |
| 1443 | } |
| 1444 | |
| 1445 | KJ_UNREACHABLE; |
| 1446 | } |
| 1447 | |
| 1448 | kj::Maybe<int> WritableStreamInternalController::getDesiredSize() { |
| 1449 | KJ_SWITCH_ONEOF(state) { |
| 1450 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1451 | return 0; |
| 1452 | } |
| 1453 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1454 | return kj::none; |
| 1455 | } |
| 1456 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1457 | KJ_IF_SOME(highWaterMark, maybeHighWaterMark) { |
| 1458 | return highWaterMark - currentWriteBufferSize; |
| 1459 | } |
| 1460 | return 1; |
| 1461 | } |
| 1462 | } |
| 1463 | |
| 1464 | KJ_UNREACHABLE; |
| 1465 | } |
| 1466 | |
| 1467 | bool WritableStreamInternalController::lockWriter(jsg::Lock& js, Writer& writer) { |
| 1468 | if (isLockedToWriter()) { |
| 1469 | return false; |
| 1470 | } |
| 1471 | |
| 1472 | auto closedPrp = js.newPromiseAndResolver<void>(); |
| 1473 | closedPrp.promise.markAsHandled(js); |
| 1474 | |
| 1475 | auto readyPrp = js.newPromiseAndResolver<void>(); |
| 1476 | readyPrp.promise.markAsHandled(js); |
| 1477 | |
| 1478 | auto lock = WriterLocked(writer, kj::mv(closedPrp.resolver), kj::mv(readyPrp.resolver)); |
| 1479 | |
| 1480 | KJ_SWITCH_ONEOF(state) { |
| 1481 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 1482 | maybeResolvePromise(js, lock.getClosedFulfiller()); |
| 1483 | maybeResolvePromise(js, lock.getReadyFulfiller()); |
| 1484 | } |
| 1485 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 1486 | maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js)); |
| 1487 | maybeRejectPromise<void>(js, lock.getReadyFulfiller(), errored.getHandle(js)); |
| 1488 | } |
| 1489 | KJ_CASE_ONEOF(writable, IoOwn<Writable>) { |
| 1490 | maybeResolvePromise(js, lock.getReadyFulfiller()); |
| 1491 | } |
| 1492 | } |
| 1493 | |
| 1494 | writeState.transitionTo<WriterLocked>(kj::mv(lock)); |
| 1495 | writer.attach(js, *this, kj::mv(closedPrp.promise), kj::mv(readyPrp.promise)); |
| 1496 | return true; |
| 1497 | } |
| 1498 | |
| 1499 | void WritableStreamInternalController::releaseWriter( |
| 1500 | Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) { |
| 1501 | KJ_IF_SOME(locked, writeState.tryGetUnsafe<WriterLocked>()) { |
| 1502 | KJ_ASSERT(&locked.getWriter() == &writer); |
| 1503 | KJ_IF_SOME(js, maybeJs) { |
| 1504 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), |
| 1505 | js.v8TypeError("This WritableStream writer has been released."_kj)); |
| 1506 | } |
| 1507 | locked.clear(); |
| 1508 | |
| 1509 | // When maybeJs is nullptr, that means releaseWriter was called when the writer is |
| 1510 | // being deconstructed and not as the result of explicitly calling releaseLock and |
| 1511 | // we do not have an isolate lock. In that case, we don't want to change the lock |
| 1512 | // state itself. Clearing the lock above will free the lock state while keeping the |
| 1513 | // WritableStream marked as locked. |
| 1514 | if (maybeJs != kj::none) { |
| 1515 | writeState.transitionTo<Unlocked>(); |
| 1516 | } |
| 1517 | } |
| 1518 | } |
| 1519 | |
| 1520 | bool WritableStreamInternalController::isClosedOrClosing() { |
| 1521 | |
| 1522 | bool isClosing = !queue.empty() && queue.back().event.is<kj::Own<Close>>(); |
| 1523 | bool isFlushing = !queue.empty() && queue.back().event.is<kj::Own<Flush>>(); |
| 1524 | return state.is<StreamStates::Closed>() || isClosing || isFlushing; |
| 1525 | } |
| 1526 | |
| 1527 | bool WritableStreamInternalController::isPiping() { |
| 1528 | return state.is<IoOwn<Writable>>() && !queue.empty() && queue.back().event.is<kj::Own<Pipe>>(); |
| 1529 | } |
| 1530 | |
| 1531 | bool WritableStreamInternalController::isErrored() { |
| 1532 | return state.is<StreamStates::Errored>(); |
| 1533 | } |
| 1534 | |
| 1535 | void WritableStreamInternalController::doClose(jsg::Lock& js) { |
| 1536 | // If already in a terminal state, nothing to do. |
| 1537 | if (state.isTerminal()) return; |
| 1538 | |
| 1539 | state.transitionTo<StreamStates::Closed>(); |
| 1540 | KJ_IF_SOME(locked, writeState.tryGetUnsafe<WriterLocked>()) { |
| 1541 | maybeResolvePromise(js, locked.getClosedFulfiller()); |
| 1542 | maybeResolvePromise(js, locked.getReadyFulfiller()); |
| 1543 | writeState.transitionTo<Locked>(); |
| 1544 | } else { |
| 1545 | (void)writeState.transitionFromTo<PipeLocked, Unlocked>(); |
| 1546 | } |
| 1547 | PendingAbort::dequeue(maybePendingAbort); |
| 1548 | } |
| 1549 | |
| 1550 | void WritableStreamInternalController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 1551 | // If already in a terminal state, nothing to do. |
| 1552 | if (state.isTerminal()) return; |
| 1553 | |
| 1554 | state.transitionTo<StreamStates::Errored>(js.v8Ref(reason)); |
| 1555 | KJ_IF_SOME(locked, writeState.tryGetUnsafe<WriterLocked>()) { |
| 1556 | maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason); |
| 1557 | maybeResolvePromise(js, locked.getReadyFulfiller()); |
| 1558 | writeState.transitionTo<Locked>(); |
| 1559 | } else { |
| 1560 | (void)writeState.transitionFromTo<PipeLocked, Unlocked>(); |
| 1561 | } |
| 1562 | PendingAbort::dequeue(maybePendingAbort); |
| 1563 | } |
| 1564 | |
| 1565 | void WritableStreamInternalController::ensureWriting(jsg::Lock& js) { |
| 1566 | auto& ioContext = IoContext::current(); |
| 1567 | if (queue.size() == 1) { |
| 1568 | ioContext.addTask(ioContext.awaitJs(js, writeLoop(js, ioContext)).attach(addRef())); |
| 1569 | } |
| 1570 | } |
| 1571 | |
| 1572 | jsg::Promise<void> WritableStreamInternalController::writeLoop( |
| 1573 | jsg::Lock& js, IoContext& ioContext) { |
| 1574 | if (queue.empty()) { |
| 1575 | return js.resolvedPromise(); |
| 1576 | } else KJ_IF_SOME(promise, queue.front().outputLock) { |
| 1577 | return ioContext.awaitIo(js, kj::mv(*promise), |
| 1578 | [this](jsg::Lock& js) -> jsg::Promise<void> { return writeLoopAfterFrontOutputLock(js); }); |
| 1579 | } else { |
| 1580 | return writeLoopAfterFrontOutputLock(js); |
| 1581 | } |
| 1582 | } |
| 1583 | |
| 1584 | void WritableStreamInternalController::finishClose(jsg::Lock& js) { |
| 1585 | KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { |
| 1586 | pendingAbort->complete(js); |
| 1587 | } |
| 1588 | |
| 1589 | doClose(js); |
| 1590 | } |
| 1591 | |
| 1592 | void WritableStreamInternalController::finishError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 1593 | KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { |
| 1594 | // In this case, and only this case, we ignore any pending rejection |
| 1595 | // that may be stored in the pendingAbort. The current exception takes |
| 1596 | // precedence. |
| 1597 | pendingAbort->fail(js, reason); |
| 1598 | } |
| 1599 | |
| 1600 | doError(js, reason); |
| 1601 | } |
| 1602 | |
| 1603 | jsg::Promise<void> WritableStreamInternalController::writeLoopAfterFrontOutputLock(jsg::Lock& js) { |
| 1604 | auto& ioContext = IoContext::current(); |
| 1605 | |
| 1606 | // This helper function is just used to enhance the assert logging when checking |
| 1607 | // that the request in flight is the one we expect. |
| 1608 | static constexpr auto inspectQueue = [](auto& queue, kj::StringPtr name) { |
| 1609 | if (queue.size() > 1) { |
| 1610 | kj::Vector<kj::String> events; |
| 1611 | for (auto& event: queue) { |
| 1612 | KJ_SWITCH_ONEOF(event.event) { |
| 1613 | KJ_CASE_ONEOF(write, kj::Own<Write>) { |
| 1614 | events.add(kj::str("Write")); |
| 1615 | } |
| 1616 | KJ_CASE_ONEOF(flush, kj::Own<Flush>) { |
| 1617 | events.add(kj::str("Flush")); |
| 1618 | } |
| 1619 | KJ_CASE_ONEOF(close, kj::Own<Close>) { |
| 1620 | events.add(kj::str("Close")); |
| 1621 | } |
| 1622 | KJ_CASE_ONEOF(pipe, kj::Own<Pipe>) { |
| 1623 | events.add(kj::str("Pipe")); |
| 1624 | } |
| 1625 | } |
| 1626 | } |
| 1627 | return kj::str("Too many events in internal writablestream queue: ", |
| 1628 | kj::delimited(kj::mv(events), ", ")); |
| 1629 | } |
| 1630 | return kj::String(); |
| 1631 | }; |
| 1632 | |
| 1633 | const auto makeChecker = [this]() { |
| 1634 | // Make a helper function that asserts that the queue did not change state during a write/close |
| 1635 | // operation. We normally only pop/drain the queue after write/close completion. We drain the |
| 1636 | // queue concurrently during finalization, but finalization would also have canceled our |
| 1637 | // write/close promise. The helper function also helpfully returns a reference to the current |
| 1638 | // request in flight. |
| 1639 | // |
| 1640 | // We capture the current generation and verify it hasn't changed, rather than using pointer |
| 1641 | // comparison, because RingBuffer may relocate elements when it grows. |
| 1642 | |
| 1643 | return [this, expectedGeneration = queue.currentGeneration()]<typename Request>() -> Request& { |
| 1644 | if constexpr (kj::isSameType<Request, Write>() || kj::isSameType<Request, Flush>()) { |
| 1645 | // Write and flush requests can have any number of requests backed up after them. |
| 1646 | KJ_ASSERT(!queue.empty()); |
| 1647 | } else if constexpr (kj::isSameType<Request, Close>()) { |
| 1648 | // Pipe and Close requests are always the last one in the queue. |
| 1649 | KJ_ASSERT(queue.size() == 1, queue.size(), inspectQueue(queue, "Pipe")); |
| 1650 | } else if constexpr (kj::isSameType<Request, Pipe>()) { |
| 1651 | // Pipe and Close requests are always the last one in the queue. |
| 1652 | KJ_ASSERT(queue.size() == 1, queue.size(), inspectQueue(queue, "Pipe")); |
| 1653 | } |
| 1654 | |
| 1655 | // Verify nothing was popped from the queue while we were waiting. |
| 1656 | KJ_ASSERT(queue.currentGeneration() == expectedGeneration); |
| 1657 | |
| 1658 | return *queue.front().event.get<kj::Own<Request>>(); |
| 1659 | }; |
| 1660 | }; |
| 1661 | |
| 1662 | const auto maybeAbort = [this](jsg::Lock& js) -> bool { |
| 1663 | auto& writable = KJ_ASSERT_NONNULL(state.tryGetUnsafe<IoOwn<Writable>>()); |
| 1664 | KJ_IF_SOME(pendingAbort, WritableStreamController::PendingAbort::dequeue(maybePendingAbort)) { |
| 1665 | auto ex = js.exceptionToKj(pendingAbort->reason.addRef(js)); |
| 1666 | writable->abort(kj::mv(ex)); |
| 1667 | drain(js, pendingAbort->reason.getHandle(js)); |
| 1668 | pendingAbort->complete(js); |
| 1669 | return true; |
| 1670 | } |
| 1671 | return false; |
| 1672 | }; |
| 1673 | |
| 1674 | // Do we have anything left to do? |
| 1675 | if (queue.empty()) return js.resolvedPromise(); |
| 1676 | |
| 1677 | KJ_SWITCH_ONEOF(queue.front().event) { |
| 1678 | KJ_CASE_ONEOF(request, kj::Own<Write>) { |
| 1679 | if (request->bytes.size() == 0) { |
| 1680 | // Zero-length writes are no-ops with a pending event. If we allowed them, we'd have a hard |
| 1681 | // time distinguishing between disconnections and zero-length reads on the other end of the |
| 1682 | // TransformStream. |
| 1683 | maybeResolvePromise(js, request->promise); |
| 1684 | queue.pop_front(); |
| 1685 | |
| 1686 | // Note: we don't bother checking for an abort() here because either this write was just |
| 1687 | // queued, in which case abort() cannot have been called yet, or this write was processed |
| 1688 | // immediately after a previous write, in which case we just checked for an abort(). |
| 1689 | return writeLoop(js, ioContext); |
| 1690 | } |
| 1691 | |
| 1692 | // writeLoop() is only called with the sink in the Writable state. |
| 1693 | auto& writable = state.getUnsafe<IoOwn<Writable>>(); |
| 1694 | auto check = makeChecker(); |
| 1695 | |
| 1696 | auto amountToWrite = request->bytes.size(); |
| 1697 | |
| 1698 | auto promise = writable->sink->write(request->bytes).attach(kj::mv(request->ownBytes)); |
| 1699 | |
| 1700 | // TODO(soon): We use awaitIoLegacy() here because if the stream terminates in JavaScript in |
| 1701 | // this same isolate, then the promise may actually be waiting on JavaScript to do something, |
| 1702 | // and so should not be considered waiting on external I/O. We will need to use |
| 1703 | // registerPendingEvent() manually when reading from an external stream. Ideally, we would |
| 1704 | // refactor the implementation so that when waiting on a JavaScript stream, we strictly use |
| 1705 | // jsg::Promises and not kj::Promises, so that it doesn't look like I/O at all, and there's |
| 1706 | // no need to drop the isolate lock and take it again every time some data is read/written. |
| 1707 | // That's a larger refactor, though. |
| 1708 | return ioContext.awaitIoLegacy(js, writable->canceler.wrap(kj::mv(promise))) |
| 1709 | .then(js, |
| 1710 | ioContext.addFunctor( |
| 1711 | [this, check, maybeAbort, amountToWrite](jsg::Lock& js) -> jsg::Promise<void> { |
| 1712 | // Under some conditions, the clean up has already happened. |
| 1713 | if (queue.empty()) return js.resolvedPromise(); |
| 1714 | auto& request = check.template operator()<Write>(); |
| 1715 | maybeResolvePromise(js, request.promise); |
| 1716 | adjustWriteBufferSize(js, -amountToWrite); |
| 1717 | KJ_IF_SOME(o, observer) { |
| 1718 | o->onChunkDequeued(amountToWrite); |
| 1719 | } |
| 1720 | queue.pop_front(); |
| 1721 | maybeAbort(js); |
| 1722 | return writeLoop(js, IoContext::current()); |
| 1723 | }), |
| 1724 | ioContext.addFunctor([this, check, maybeAbort, amountToWrite]( |
| 1725 | jsg::Lock& js, jsg::Value reason) -> jsg::Promise<void> { |
| 1726 | // Under some conditions, the clean up has already happened. |
| 1727 | if (queue.empty()) return js.resolvedPromise(); |
| 1728 | auto handle = reason.getHandle(js); |
| 1729 | auto& request = check.template operator()<Write>(); |
| 1730 | auto& writable = state.getUnsafe<IoOwn<Writable>>(); |
| 1731 | adjustWriteBufferSize(js, -amountToWrite); |
| 1732 | KJ_IF_SOME(o, observer) { |
| 1733 | o->onChunkDequeued(amountToWrite); |
| 1734 | } |
| 1735 | maybeRejectPromise<void>(js, request.promise, handle); |
| 1736 | queue.pop_front(); |
| 1737 | if (!maybeAbort(js)) { |
| 1738 | auto ex = js.exceptionToKj(reason.addRef(js)); |
| 1739 | writable->abort(kj::mv(ex)); |
| 1740 | drain(js, handle); |
| 1741 | } |
| 1742 | return js.resolvedPromise(); |
| 1743 | })); |
| 1744 | } |
| 1745 | KJ_CASE_ONEOF(request, kj::Own<Pipe>) { |
| 1746 | // The destination should still be Writable, because the only way to transition to an |
| 1747 | // errored state would have been if a write request in the queue ahead of us encountered an |
| 1748 | // error. But in that case, the queue would already have been drained and we wouldn't be here. |
| 1749 | auto& writable = state.getUnsafe<IoOwn<Writable>>(); |
| 1750 | |
| 1751 | if (request->checkSignal(js)) { |
| 1752 | // If the signal is triggered, checkSignal will handle erroring the source and destination. |
| 1753 | return js.resolvedPromise(); |
| 1754 | } |
| 1755 | |
| 1756 | // The readable side should *should* still be readable here but let's double check, just |
| 1757 | // to be safe, both for closed state and errored states. |
| 1758 | if (request->source().isClosed()) { |
| 1759 | request->source().release(js); |
| 1760 | // If the source is closed, the spec requires us to close the destination unless the |
| 1761 | // preventClose option is true. |
| 1762 | if (!request->preventClose() && !isClosedOrClosing()) { |
| 1763 | doClose(js); |
| 1764 | } else { |
| 1765 | writeState.transitionTo<Unlocked>(); |
| 1766 | } |
| 1767 | return js.resolvedPromise(); |
| 1768 | } |
| 1769 | |
| 1770 | KJ_IF_SOME(errored, request->source().tryGetErrored(js)) { |
| 1771 | request->source().release(js); |
| 1772 | // If the source is errored, the spec requires us to error the destination unless the |
| 1773 | // preventAbort option is true. |
| 1774 | if (!request->preventAbort()) { |
| 1775 | auto ex = js.exceptionToKj(js.v8Ref(errored)); |
| 1776 | writable->abort(kj::mv(ex)); |
| 1777 | drain(js, errored); |
| 1778 | } else { |
| 1779 | writeState.transitionTo<Unlocked>(); |
| 1780 | } |
| 1781 | return js.resolvedPromise(); |
| 1782 | } |
| 1783 | |
| 1784 | // Up to this point, we really don't know what kind of ReadableStream source we're dealing |
| 1785 | // with. If the source is backed by a ReadableStreamSource, then the call to tryPumpTo below |
| 1786 | // will return a kj::Promise that will be resolved once the kj mechanisms for piping have |
| 1787 | // completed. From there, the only thing left to do is resolve the JavaScript pipe promise, |
| 1788 | // unlock things, and continue on. If the call to tryPumpTo returns nullptr, however, the |
| 1789 | // ReadableStream is JavaScript-backed and we need to setup a JavaScript-promise read/write |
| 1790 | // loop to pass the data into the destination. |
| 1791 | |
| 1792 | const auto handlePromise = [this, &ioContext, check = makeChecker(), |
| 1793 | preventAbort = request->preventAbort()]( |
| 1794 | jsg::Lock& js, auto promise) { |
| 1795 | return promise.then(js, ioContext.addFunctor([this, check](jsg::Lock& js) mutable { |
| 1796 | // Under some conditions, the clean up has already happened. |
| 1797 | if (queue.empty()) return js.resolvedPromise(); |
| 1798 | |
| 1799 | auto& request = check.template operator()<Pipe>(); |
| 1800 | |
| 1801 | // It's possible we got here because the source errored but preventAbort was set. |
| 1802 | // In that case, we need to treat preventAbort the same as preventClose. Be |
| 1803 | // sure to check this before calling sourceLock.close() or the error detail will |
| 1804 | // be lost. |
| 1805 | // Capture preventClose now so we can modify it locally if needed. |
| 1806 | bool preventClose = request.preventClose(); |
| 1807 | KJ_IF_SOME(errored, request.source().tryGetErrored(js)) { |
| 1808 | if (request.preventAbort()) preventClose = true; |
| 1809 | // Even through we're not going to close the destination, we still want the |
| 1810 | // pipe promise itself to be rejected in this case. |
| 1811 | maybeRejectPromise<void>(js, request.promise(), errored); |
| 1812 | } else KJ_IF_SOME(errored, state.tryGetUnsafe<StreamStates::Errored>()) { |
| 1813 | maybeRejectPromise<void>(js, request.promise(), errored.getHandle(js)); |
| 1814 | } else { |
| 1815 | maybeResolvePromise(js, request.promise()); |
| 1816 | } |
| 1817 | |
| 1818 | // Always transition the readable side to the closed state, because we read until EOF. |
| 1819 | // Note that preventClose (below) means "don't close the writable side", i.e. don't |
| 1820 | // call end(). |
| 1821 | request.source().close(js); |
| 1822 | queue.pop_front(); |
| 1823 | |
| 1824 | if (!preventClose) { |
| 1825 | // Note: unlike a real Close request, it's not possible for us to have been aborted. |
| 1826 | return close(js, true); |
| 1827 | } else { |
| 1828 | writeState.transitionTo<Unlocked>(); |
| 1829 | } |
| 1830 | return js.resolvedPromise(); |
| 1831 | }), |
| 1832 | ioContext.addFunctor( |
| 1833 | [this, check, preventAbort](jsg::Lock& js, jsg::Value reason) mutable { |
| 1834 | auto handle = reason.getHandle(js); |
| 1835 | auto& request = check.template operator()<Pipe>(); |
| 1836 | maybeRejectPromise<void>(js, request.promise(), handle); |
| 1837 | // TODO(conform): Remember all those checks we performed in ReadableStream::pipeTo()? |
| 1838 | // We're supposed to perform the same checks continually, e.g., errored writes should |
| 1839 | // cancel the readable side unless preventCancel is truthy... This would require |
| 1840 | // deeper integration with the implementation of pumpTo(). Oh well. One consequence |
| 1841 | // of this is that if there is an error on the writable side, we error the readable |
| 1842 | // side, rather than close (cancel) it, which is what the spec would have us do. |
| 1843 | // TODO(now): Warn on the console about this. |
| 1844 | request.source().error(js, handle); |
| 1845 | queue.pop_front(); |
| 1846 | if (!preventAbort) { |
| 1847 | return abort(js, handle); |
| 1848 | } |
| 1849 | doError(js, handle); |
| 1850 | return js.resolvedPromise(); |
| 1851 | })); |
| 1852 | }; |
| 1853 | |
| 1854 | KJ_IF_SOME(promise, request->source().tryPumpTo(*writable->sink, !request->preventClose())) { |
| 1855 | return handlePromise(js, |
| 1856 | ioContext.awaitIo(js, |
| 1857 | writable->canceler.wrap( |
| 1858 | AbortSignal::maybeCancelWrap(js, request->maybeSignal(), kj::mv(promise))))); |
| 1859 | } |
| 1860 | |
| 1861 | // The ReadableStream is JavaScript-backed. We can still pipe the data but it's going to be |
| 1862 | // a bit slower because we will be relying on JavaScript promises when reading the data |
| 1863 | // from the ReadableStream, then waiting on kj::Promises to write the data. We will keep |
| 1864 | // reading until either the source or destination errors or until the source signals that |
| 1865 | // it is done. |
| 1866 | return handlePromise(js, request->pipeLoop(js)); |
| 1867 | } |
| 1868 | KJ_CASE_ONEOF(request, kj::Own<Close>) { |
| 1869 | // writeLoop() is only called with the sink in the Writable state. |
| 1870 | auto& writable = state.getUnsafe<IoOwn<Writable>>(); |
| 1871 | auto check = makeChecker(); |
| 1872 | |
| 1873 | return ioContext.awaitIo(js, writable->canceler.wrap(writable->sink->end())) |
| 1874 | .then(js, ioContext.addFunctor([this, check](jsg::Lock& js) { |
| 1875 | // Under some conditions, the clean up has already happened. |
| 1876 | if (queue.empty()) return; |
| 1877 | auto& request = check.template operator()<Close>(); |
| 1878 | maybeResolvePromise(js, request.promise); |
| 1879 | queue.pop_front(); |
| 1880 | finishClose(js); |
| 1881 | }), |
| 1882 | ioContext.addFunctor([this, check](jsg::Lock& js, jsg::Value reason) { |
| 1883 | // Under some conditions, the clean up has already happened. |
| 1884 | if (queue.empty()) return; |
| 1885 | auto handle = reason.getHandle(js); |
| 1886 | auto& request = check.template operator()<Close>(); |
| 1887 | maybeRejectPromise<void>(js, request.promise, handle); |
| 1888 | queue.pop_front(); |
| 1889 | finishError(js, handle); |
| 1890 | })); |
| 1891 | } |
| 1892 | KJ_CASE_ONEOF(request, kj::Own<Flush>) { |
| 1893 | // This is not a standards-defined state for a WritableStream and is only used internally |
| 1894 | // for Socket's startTls call. |
| 1895 | // |
| 1896 | // Flushing is similar to closing the stream, the main difference is that `finishClose` |
| 1897 | // and `writable->end()` are never called. |
| 1898 | // Note: For Flush, we don't need makeChecker since we process immediately without async I/O. |
| 1899 | maybeResolvePromise(js, request->promise); |
| 1900 | queue.pop_front(); |
| 1901 | |
| 1902 | return js.resolvedPromise(); |
| 1903 | } |
| 1904 | } |
| 1905 | |
| 1906 | KJ_UNREACHABLE; |
| 1907 | } |
| 1908 | |
| 1909 | bool WritableStreamInternalController::Pipe::State::checkSignal(jsg::Lock& js) { |
| 1910 | // Returns true if the caller should bail out and stop processing. This happens in two cases: |
| 1911 | // 1. The State was aborted (e.g., by drain()) - the Pipe is being torn down |
| 1912 | // 2. The AbortSignal was triggered - we handle the abort and return true |
| 1913 | // In both cases, the caller should return a resolved promise and not continue the pipe loop. |
| 1914 | if (aborted) return true; |
| 1915 | |
| 1916 | KJ_IF_SOME(signal, maybeSignal) { |
| 1917 | if (signal->getAborted(js)) { |
| 1918 | auto reason = signal->getReason(js); |
| 1919 | |
| 1920 | // abort process might call parent.drain which will delete this, |
| 1921 | // move/copy everything we need after into temps. |
| 1922 | auto& parentRef = this->parent; |
| 1923 | auto& sourceRef = this->source; |
| 1924 | auto preventCancelCopy = this->preventCancel; |
| 1925 | auto promiseCopy = kj::mv(this->promise); |
| 1926 | |
| 1927 | if (!preventAbort) { |
| 1928 | KJ_IF_SOME(writable, parent.state.tryGetUnsafe<IoOwn<Writable>>()) { |
| 1929 | auto ex = js.exceptionToKj(reason); |
| 1930 | writable->abort(kj::mv(ex)); |
| 1931 | parentRef.drain(js, reason); |
| 1932 | } else { |
| 1933 | parent.writeState.transitionTo<Unlocked>(); |
| 1934 | } |
| 1935 | } else { |
| 1936 | parent.writeState.transitionTo<Unlocked>(); |
| 1937 | } |
| 1938 | if (!preventCancelCopy) { |
| 1939 | sourceRef.release(js, v8::Local<v8::Value>(reason)); |
| 1940 | } else { |
| 1941 | sourceRef.release(js); |
| 1942 | } |
| 1943 | maybeRejectPromise<void>(js, promiseCopy, reason); |
| 1944 | return true; |
| 1945 | } |
| 1946 | } |
| 1947 | return false; |
| 1948 | } |
| 1949 | |
| 1950 | jsg::Promise<void> WritableStreamInternalController::Pipe::State::write( |
| 1951 | v8::Local<v8::Value> handle) { |
| 1952 | auto& writable = parent.state.getUnsafe<IoOwn<Writable>>(); |
| 1953 | // TODO(soon): Once jsg::BufferSource lands and we're able to use it, this can be simplified. |
| 1954 | KJ_ASSERT(handle->IsArrayBuffer() || handle->IsArrayBufferView()); |
| 1955 | std::shared_ptr<v8::BackingStore> store; |
| 1956 | size_t byteLength = 0; |
| 1957 | size_t byteOffset = 0; |
| 1958 | if (handle->IsArrayBuffer()) { |
| 1959 | auto buffer = handle.template As<v8::ArrayBuffer>(); |
| 1960 | store = buffer->GetBackingStore(); |
| 1961 | byteLength = buffer->ByteLength(); |
| 1962 | } else { |
| 1963 | auto view = handle.template As<v8::ArrayBufferView>(); |
| 1964 | store = view->Buffer()->GetBackingStore(); |
| 1965 | byteLength = view->ByteLength(); |
| 1966 | byteOffset = view->ByteOffset(); |
| 1967 | } |
| 1968 | kj::byte* data = reinterpret_cast<kj::byte*>(store->Data()) + byteOffset; |
| 1969 | // TODO(cleanup): Have this method accept a jsg::Lock& from the caller instead of using |
| 1970 | // v8::Isolate::GetCurrent(); |
| 1971 | auto& js = jsg::Lock::current(); |
| 1972 | |
| 1973 | // For resizable ArrayBuffers or shared backing stores, we must eagerly copy |
| 1974 | // the data. A resizable ArrayBuffer's logical byte length can be changed by user |
| 1975 | // JS after write() returns but before the sink consumes the data, making the |
| 1976 | // cached byteLength stale. |
| 1977 | // But also just beacuse of V8 Sandbox requirements, we really should be copying |
| 1978 | // the data from the ArrayBuffer anyway... We incur an allocation and copy cost |
| 1979 | // here but that's to be expected. |
| 1980 | auto backing = kj::heapArray<kj::byte>(byteLength); |
| 1981 | backing.asPtr().copyFrom(kj::arrayPtr(data, byteLength)); |
| 1982 | return IoContext::current().awaitIo(js, |
| 1983 | writable->canceler.wrap(writable->sink->write(backing)).attach(kj::mv(backing)), |
| 1984 | [](jsg::Lock&) {}); |
| 1985 | } |
| 1986 | |
| 1987 | jsg::Promise<void> WritableStreamInternalController::Pipe::State::pipeLoop(jsg::Lock& js) { |
| 1988 | // This is a bit of dance. We got here because the source ReadableStream does not support |
| 1989 | // the internal, more efficient kj pipe (which means it is a JavaScript-backed ReadableStream). |
| 1990 | // We need to call read() on the source which returns a JavaScript Promise, wait on it to resolve, |
| 1991 | // then call write() which returns a kj::Promise. Before each iteration we check to see if either |
| 1992 | // the source or the destination have errored or closed and handle accordingly. At some point we |
| 1993 | // should explore if there are ways of making this more efficient. For the most part, however, |
| 1994 | // every read from the source must call into JavaScript to advance the ReadableStream. |
| 1995 | |
| 1996 | auto& ioContext = IoContext::current(); |
| 1997 | |
| 1998 | if (aborted) { |
| 1999 | return js.resolvedPromise(); |
| 2000 | } |
| 2001 | |
| 2002 | if (checkSignal(js)) { |
| 2003 | // If the signal is triggered, checkSignal will handle erroring the source and destination. |
| 2004 | return js.resolvedPromise(); |
| 2005 | } |
| 2006 | |
| 2007 | // Here we check the closed and errored states of both the source and the destination, |
| 2008 | // propagating those states to the other based on the options. This check must be |
| 2009 | // performed at the start of each iteration in the pipe loop. |
| 2010 | // |
| 2011 | // TODO(soon): These are the same checks made before we entered the loop. Try to |
| 2012 | // unify the code to reduce duplication. |
| 2013 | |
| 2014 | KJ_IF_SOME(errored, source.tryGetErrored(js)) { |
| 2015 | source.release(js); |
| 2016 | if (!preventAbort) { |
| 2017 | KJ_IF_SOME(writable, parent.state.tryGetUnsafe<IoOwn<Writable>>()) { |
| 2018 | auto ex = js.exceptionToKj(js.v8Ref(errored)); |
| 2019 | writable->abort(kj::mv(ex)); |
| 2020 | return js.rejectedPromise<void>(errored); |
| 2021 | } |
| 2022 | } |
| 2023 | |
| 2024 | // If preventAbort was true, we're going to unlock the destination now. |
| 2025 | // We are not going to propagate the error here tho. |
| 2026 | parent.writeState.transitionTo<Unlocked>(); |
| 2027 | return js.resolvedPromise(); |
| 2028 | } |
| 2029 | |
| 2030 | KJ_IF_SOME(errored, parent.state.tryGetUnsafe<StreamStates::Errored>()) { |
| 2031 | parent.writeState.transitionTo<Unlocked>(); |
| 2032 | if (!preventCancel) { |
| 2033 | auto reason = errored.getHandle(js); |
| 2034 | source.release(js, reason); |
| 2035 | return js.rejectedPromise<void>(reason); |
| 2036 | } |
| 2037 | source.release(js); |
| 2038 | return js.resolvedPromise(); |
| 2039 | } |
| 2040 | |
| 2041 | if (source.isClosed()) { |
| 2042 | source.release(js); |
| 2043 | if (!preventClose) { |
| 2044 | KJ_ASSERT(!parent.state.is<StreamStates::Errored>()); |
| 2045 | if (!parent.isClosedOrClosing()) { |
| 2046 | // We'll only be here if the sink is in the Writable state. |
| 2047 | auto& ioContext = IoContext::current(); |
| 2048 | // Capture a ref to the state to keep it alive during async operations. |
| 2049 | return ioContext |
| 2050 | .awaitIo(js, parent.state.getUnsafe<IoOwn<Writable>>()->sink->end(), [](jsg::Lock&) {}) |
| 2051 | .then(js, ioContext.addFunctor([state = kj::addRef(*this)](jsg::Lock& js) { |
| 2052 | if (state->aborted) return; |
| 2053 | state->parent.finishClose(js); |
| 2054 | }), |
| 2055 | ioContext.addFunctor([state = kj::addRef(*this)](jsg::Lock& js, jsg::Value reason) { |
| 2056 | if (state->aborted) return; |
| 2057 | state->parent.finishError(js, reason.getHandle(js)); |
| 2058 | })); |
| 2059 | } |
| 2060 | parent.writeState.transitionTo<Unlocked>(); |
| 2061 | } |
| 2062 | return js.resolvedPromise(); |
| 2063 | } |
| 2064 | |
| 2065 | if (parent.isClosedOrClosing()) { |
| 2066 | auto destClosed = js.v8TypeError("This destination writable stream is closed."_kj); |
| 2067 | parent.writeState.transitionTo<Unlocked>(); |
| 2068 | |
| 2069 | if (!preventCancel) { |
| 2070 | source.release(js, destClosed); |
| 2071 | } else { |
| 2072 | source.release(js); |
| 2073 | } |
| 2074 | |
| 2075 | return js.rejectedPromise<void>(destClosed); |
| 2076 | } |
| 2077 | |
| 2078 | return source.read(js).then(js, |
| 2079 | ioContext.addFunctor([state = kj::addRef(*this)]( |
| 2080 | jsg::Lock& js, ReadResult result) mutable -> jsg::Promise<void> { |
| 2081 | if (state->aborted || state->checkSignal(js) || result.done) { |
| 2082 | return js.resolvedPromise(); |
| 2083 | } |
| 2084 | |
| 2085 | // WritableStreamInternalControllers only support byte data. If we can't |
| 2086 | // interpret the result.value as bytes, then we error the pipe; otherwise |
| 2087 | // we sent those bytes on to the WritableStreamSink. |
| 2088 | KJ_IF_SOME(value, result.value) { |
| 2089 | auto handle = value.getHandle(js); |
| 2090 | if (handle->IsArrayBuffer() || handle->IsArrayBufferView()) { |
| 2091 | return state->write(handle).then(js, |
| 2092 | [state = kj::addRef(*state)](jsg::Lock& js) mutable -> jsg::Promise<void> { |
| 2093 | if (state->aborted) { |
| 2094 | return js.resolvedPromise(); |
| 2095 | } |
| 2096 | // The signal will be checked again at the start of the next loop iteration. |
| 2097 | return state->pipeLoop(js); |
| 2098 | }, |
| 2099 | [state = kj::addRef(*state)]( |
| 2100 | jsg::Lock& js, jsg::Value reason) mutable -> jsg::Promise<void> { |
| 2101 | if (state->aborted) { |
| 2102 | return js.resolvedPromise(); |
| 2103 | } |
| 2104 | state->parent.doError(js, reason.getHandle(js)); |
| 2105 | return state->pipeLoop(js); |
| 2106 | }); |
| 2107 | } |
| 2108 | } |
| 2109 | // Undefined and null are perfectly valid values to pass through a ReadableStream, |
| 2110 | // but we can't interpret them as bytes so if we get them here, we error the pipe. |
| 2111 | auto error = js.v8TypeError("This WritableStream only supports writing byte types."_kj); |
| 2112 | auto& writable = state->parent.state.getUnsafe<IoOwn<Writable>>(); |
| 2113 | auto ex = js.exceptionToKj(js.v8Ref(error)); |
| 2114 | writable->abort(kj::mv(ex)); |
| 2115 | // The error condition will be handled at the start of the next iteration. |
| 2116 | return state->pipeLoop(js); |
| 2117 | }), |
| 2118 | ioContext.addFunctor([state = kj::addRef(*this)]( |
| 2119 | jsg::Lock& js, jsg::Value reason) mutable -> jsg::Promise<void> { |
| 2120 | if (state->aborted) { |
| 2121 | return js.resolvedPromise(); |
| 2122 | } |
| 2123 | // The error will be processed and propagated in the next iteration. |
| 2124 | return state->pipeLoop(js); |
| 2125 | })); |
| 2126 | } |
| 2127 | |
| 2128 | void WritableStreamInternalController::drain(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2129 | doError(js, reason); |
| 2130 | while (!queue.empty()) { |
| 2131 | KJ_SWITCH_ONEOF(queue.front().event) { |
| 2132 | KJ_CASE_ONEOF(writeRequest, kj::Own<Write>) { |
| 2133 | maybeRejectPromise<void>(js, writeRequest->promise, reason); |
| 2134 | } |
| 2135 | KJ_CASE_ONEOF(pipeRequest, kj::Own<Pipe>) { |
| 2136 | if (!pipeRequest->preventCancel()) { |
| 2137 | pipeRequest->source().cancel(js, reason); |
| 2138 | } |
| 2139 | maybeRejectPromise<void>(js, pipeRequest->promise(), reason); |
| 2140 | } |
| 2141 | KJ_CASE_ONEOF(closeRequest, kj::Own<Close>) { |
| 2142 | maybeRejectPromise<void>(js, closeRequest->promise, reason); |
| 2143 | } |
| 2144 | KJ_CASE_ONEOF(flushRequest, kj::Own<Flush>) { |
| 2145 | maybeRejectPromise<void>(js, flushRequest->promise, reason); |
| 2146 | } |
| 2147 | } |
| 2148 | queue.pop_front(); |
| 2149 | } |
| 2150 | } |
| 2151 | |
| 2152 | void WritableStreamInternalController::visitForGc(jsg::GcVisitor& visitor) { |
| 2153 | for (auto& event: queue) { |
| 2154 | KJ_SWITCH_ONEOF(event.event) { |
| 2155 | KJ_CASE_ONEOF(write, kj::Own<Write>) { |
| 2156 | visitor.visit(write->promise); |
| 2157 | } |
| 2158 | KJ_CASE_ONEOF(close, kj::Own<Close>) { |
| 2159 | visitor.visit(close->promise); |
| 2160 | } |
| 2161 | KJ_CASE_ONEOF(flush, kj::Own<Flush>) { |
| 2162 | visitor.visit(flush->promise); |
| 2163 | } |
| 2164 | KJ_CASE_ONEOF(pipe, kj::Own<Pipe>) { |
| 2165 | visitor.visit(pipe->maybeSignal(), pipe->promise()); |
| 2166 | } |
| 2167 | } |
| 2168 | } |
| 2169 | KJ_IF_SOME(locked, writeState.tryGetUnsafe<WriterLocked>()) { |
| 2170 | visitor.visit(locked); |
| 2171 | } |
| 2172 | KJ_IF_SOME(pendingAbort, maybePendingAbort) { |
| 2173 | visitor.visit(*pendingAbort); |
| 2174 | } |
| 2175 | } |
| 2176 | |
| 2177 | void ReadableStreamInternalController::visitForGc(jsg::GcVisitor& visitor) { |
| 2178 | KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) { |
| 2179 | visitor.visit(locked); |
| 2180 | } |
| 2181 | } |
| 2182 | |
| 2183 | kj::Maybe<ReadableStreamController::PipeController&> ReadableStreamInternalController:: |
| 2184 | tryPipeLock() { |
| 2185 | if (isLockedToReader()) { |
| 2186 | return kj::none; |
| 2187 | } |
| 2188 | return readState.transitionTo<PipeLocked>(*this); |
| 2189 | } |
| 2190 | |
| 2191 | bool ReadableStreamInternalController::PipeLocked::isClosed() { |
| 2192 | return inner.state.is<StreamStates::Closed>(); |
| 2193 | } |
| 2194 | |
| 2195 | kj::Maybe<v8::Local<v8::Value>> ReadableStreamInternalController::PipeLocked::tryGetErrored( |
| 2196 | jsg::Lock& js) { |
| 2197 | KJ_IF_SOME(errored, inner.state.tryGetUnsafe<StreamStates::Errored>()) { |
| 2198 | return errored.getHandle(js); |
| 2199 | } |
| 2200 | return kj::none; |
| 2201 | } |
| 2202 | |
| 2203 | void ReadableStreamInternalController::PipeLocked::cancel( |
| 2204 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2205 | if (inner.state.is<Readable>()) { |
| 2206 | inner.doCancel(js, reason); |
| 2207 | } |
| 2208 | } |
| 2209 | |
| 2210 | void ReadableStreamInternalController::PipeLocked::close(jsg::Lock& js) { |
| 2211 | inner.doClose(js); |
| 2212 | } |
| 2213 | |
| 2214 | void ReadableStreamInternalController::PipeLocked::error( |
| 2215 | jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 2216 | inner.doError(js, reason); |
| 2217 | } |
| 2218 | |
| 2219 | void ReadableStreamInternalController::PipeLocked::release( |
| 2220 | jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError) { |
| 2221 | KJ_IF_SOME(error, maybeError) { |
| 2222 | cancel(js, error); |
| 2223 | } |
| 2224 | inner.readState.transitionTo<Unlocked>(); |
| 2225 | } |
| 2226 | |
| 2227 | kj::Maybe<kj::Promise<void>> ReadableStreamInternalController::PipeLocked::tryPumpTo( |
| 2228 | WritableStreamSink& sink, bool end) { |
| 2229 | // This is safe because the caller should have already checked isClosed and tryGetErrored |
| 2230 | // and handled those before calling tryPumpTo. |
| 2231 | auto& readable = KJ_ASSERT_NONNULL(inner.state.tryGetUnsafe<Readable>()); |
| 2232 | return IoContext::current().waitForDeferredProxy(readable->pumpTo(sink, end)); |
| 2233 | } |
| 2234 | |
| 2235 | jsg::Promise<ReadResult> ReadableStreamInternalController::PipeLocked::read(jsg::Lock& js) { |
| 2236 | return KJ_ASSERT_NONNULL(inner.read(js, kj::none)); |
| 2237 | } |
| 2238 | |
| 2239 | jsg::Promise<jsg::BufferSource> ReadableStreamInternalController::readAllBytes( |
| 2240 | jsg::Lock& js, uint64_t limit) { |
| 2241 | if (isLockedToReader()) { |
| 2242 | return js.rejectedPromise<jsg::BufferSource>(KJ_EXCEPTION( |
| 2243 | FAILED, "jsg.TypeError: This ReadableStream is currently locked to a reader.")); |
| 2244 | } |
| 2245 | if (isPendingClosure) { |
| 2246 | return js.rejectedPromise<jsg::BufferSource>( |
| 2247 | js.v8TypeError("This ReadableStream belongs to an object that is closing."_kj)); |
| 2248 | } |
| 2249 | KJ_SWITCH_ONEOF(state) { |
| 2250 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2251 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0); |
| 2252 | return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); |
| 2253 | } |
| 2254 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2255 | return js.rejectedPromise<jsg::BufferSource>(errored.addRef(js)); |
| 2256 | } |
| 2257 | KJ_CASE_ONEOF(readable, Readable) { |
| 2258 | auto source = KJ_ASSERT_NONNULL(removeSource(js)); |
| 2259 | auto& context = IoContext::current(); |
| 2260 | // TODO(perf): v8 sandboxing will require that backing stores are allocated within |
| 2261 | // the sandbox. This will require a change to the API of ReadableStreamSource::readAllBytes. |
| 2262 | // For now, we'll read and allocate into a proper backing store. |
| 2263 | return context.awaitIoLegacy(js, source->readAllBytes(limit).attach(kj::mv(source))) |
| 2264 | .then(js, [](jsg::Lock& js, kj::Array<kj::byte> bytes) -> jsg::BufferSource { |
| 2265 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, bytes.size()); |
| 2266 | backing.asArrayPtr().copyFrom(bytes); |
| 2267 | return jsg::BufferSource(js, kj::mv(backing)); |
| 2268 | }); |
| 2269 | } |
| 2270 | } |
| 2271 | KJ_UNREACHABLE; |
| 2272 | } |
| 2273 | |
| 2274 | jsg::Promise<kj::String> ReadableStreamInternalController::readAllText( |
| 2275 | jsg::Lock& js, uint64_t limit) { |
| 2276 | if (isLockedToReader()) { |
| 2277 | return js.rejectedPromise<kj::String>(KJ_EXCEPTION( |
| 2278 | FAILED, "jsg.TypeError: This ReadableStream is currently locked to a reader.")); |
| 2279 | } |
| 2280 | if (isPendingClosure) { |
| 2281 | return js.rejectedPromise<kj::String>( |
| 2282 | js.v8TypeError("This ReadableStream belongs to an object that is closing."_kj)); |
| 2283 | } |
| 2284 | KJ_SWITCH_ONEOF(state) { |
| 2285 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2286 | return js.resolvedPromise(kj::String()); |
| 2287 | } |
| 2288 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2289 | return js.rejectedPromise<kj::String>(errored.addRef(js)); |
| 2290 | } |
| 2291 | KJ_CASE_ONEOF(readable, Readable) { |
| 2292 | auto source = KJ_ASSERT_NONNULL(removeSource(js)); |
| 2293 | auto& context = IoContext::current(); |
| 2294 | auto option = ReadAllTextOption::NULL_TERMINATE; |
| 2295 | KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) { |
| 2296 | if (flags.getStripBomInReadAllText()) { |
| 2297 | option |= ReadAllTextOption::STRIP_BOM; |
| 2298 | } |
| 2299 | } |
| 2300 | return context.awaitIoLegacy(js, source->readAllText(limit, option).attach(kj::mv(source))); |
| 2301 | } |
| 2302 | } |
| 2303 | KJ_UNREACHABLE; |
| 2304 | } |
| 2305 | |
| 2306 | kj::Maybe<uint64_t> ReadableStreamInternalController::tryGetLength(StreamEncoding encoding) { |
| 2307 | KJ_SWITCH_ONEOF(state) { |
| 2308 | KJ_CASE_ONEOF(closed, StreamStates::Closed) { |
| 2309 | return static_cast<uint64_t>(0); |
| 2310 | } |
| 2311 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2312 | return kj::none; |
| 2313 | } |
| 2314 | KJ_CASE_ONEOF(readable, Readable) { |
| 2315 | return readable->tryGetLength(encoding); |
| 2316 | } |
| 2317 | } |
| 2318 | KJ_UNREACHABLE; |
| 2319 | } |
| 2320 | |
| 2321 | kj::Own<ReadableStreamController> ReadableStreamInternalController::detach( |
| 2322 | jsg::Lock& js, bool ignoreDetached) { |
| 2323 | return newReadableStreamInternalController( |
| 2324 | IoContext::current(), KJ_ASSERT_NONNULL(removeSource(js, ignoreDetached))); |
| 2325 | } |
| 2326 | |
| 2327 | kj::Promise<DeferredProxy<void>> ReadableStreamInternalController::pumpTo( |
| 2328 | jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) { |
| 2329 | auto source = KJ_ASSERT_NONNULL(removeSource(js)); |
| 2330 | |
| 2331 | struct Holder: public kj::Refcounted { |
| 2332 | kj::Own<WritableStreamSink> sink; |
| 2333 | kj::Own<ReadableStreamSource> source; |
| 2334 | bool done = false; |
| 2335 | |
| 2336 | Holder(kj::Own<WritableStreamSink> sink, kj::Own<ReadableStreamSource> source) |
| 2337 | : sink(kj::mv(sink)), |
| 2338 | source(kj::mv(source)) {} |
| 2339 | ~Holder() noexcept(false) { |
| 2340 | if (!done) { |
| 2341 | // It appears the pump was canceled. We should make sure this propagates back to the |
| 2342 | // source stream. This is important in particular when we're implementing the response |
| 2343 | // pump for an HTTP event (see Response::send()). Presumably it was canceled because the |
| 2344 | // client disconnected. If we don't cancel the source, then if the source is one end of |
| 2345 | // a TransformStream, the write end will just hang. Of course, this is fine if there are |
| 2346 | // no waitUntil()s running, because the whole I/O context will be canceled anyway. But if |
| 2347 | // there are waitUntil()s, then the application probably expects to get an exception from |
| 2348 | // the write() on cancellation, rather than have it hang. |
| 2349 | source->cancel(KJ_EXCEPTION(DISCONNECTED, "pump canceled")); |
| 2350 | } |
| 2351 | } |
| 2352 | }; |
| 2353 | |
| 2354 | auto holder = kj::rc<Holder>(kj::mv(sink), kj::mv(source)); |
| 2355 | return holder->source->pumpTo(*holder->sink, end) |
| 2356 | .then([holder = holder.addRef()](DeferredProxy<void> proxy) mutable -> DeferredProxy<void> { |
| 2357 | proxy.proxyTask = proxy.proxyTask.attach(holder.addRef()); |
| 2358 | holder->done = true; |
| 2359 | return kj::mv(proxy); |
| 2360 | }, [holder = holder.addRef()](kj::Exception&& ex) mutable { |
| 2361 | holder->sink->abort(ex.clone()); |
| 2362 | holder->source->cancel(ex.clone()); |
| 2363 | holder->done = true; |
| 2364 | return kj::mv(ex); |
| 2365 | }); |
| 2366 | } |
| 2367 | |
| 2368 | StreamEncoding ReadableStreamInternalController::getPreferredEncoding() { |
| 2369 | return state.tryGetUnsafe<Readable>() |
| 2370 | .map([](Readable& readable) { |
| 2371 | return readable->getPreferredEncoding(); |
| 2372 | }).orDefault(StreamEncoding::IDENTITY); |
| 2373 | } |
| 2374 | |
| 2375 | kj::Own<ReadableStreamController> newReadableStreamInternalController( |
| 2376 | IoContext& ioContext, kj::Own<ReadableStreamSource> source) { |
| 2377 | return kj::heap<ReadableStreamInternalController>(ioContext.addObject(kj::mv(source))); |
| 2378 | } |
| 2379 | |
| 2380 | kj::Own<WritableStreamController> newWritableStreamInternalController(IoContext& ioContext, |
| 2381 | kj::Own<WritableStreamSink> sink, |
| 2382 | kj::Maybe<kj::Own<ByteStreamObserver>> observer, |
| 2383 | kj::Maybe<uint64_t> maybeHighWaterMark, |
| 2384 | kj::Maybe<jsg::Promise<void>> maybeClosureWaitable) { |
| 2385 | return kj::heap<WritableStreamInternalController>( |
| 2386 | kj::mv(sink), kj::mv(observer), maybeHighWaterMark, kj::mv(maybeClosureWaitable)); |
| 2387 | } |
| 2388 | |
| 2389 | kj::StringPtr WritableStreamInternalController::jsgGetMemoryName() const { |
| 2390 | return "WritableStreamInternalController"_kjc; |
| 2391 | } |
| 2392 | |
| 2393 | size_t WritableStreamInternalController::jsgGetMemorySelfSize() const { |
| 2394 | return sizeof(WritableStreamInternalController); |
| 2395 | } |
| 2396 | void WritableStreamInternalController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 2397 | KJ_SWITCH_ONEOF(state) { |
| 2398 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 2399 | KJ_CASE_ONEOF(errored, StreamStates::Errored) { |
| 2400 | tracker.trackField("error", errored); |
| 2401 | } |
| 2402 | KJ_CASE_ONEOF(_, IoOwn<Writable>) { |
| 2403 | // Ideally we'd be able to track the size of any pending writes held in the sink's |
| 2404 | // queue but since it is behind an IoOwn and we won't be holding the IoContext here, |
| 2405 | // we can't. |
| 2406 | tracker.trackFieldWithSize("IoOwn<WritableStreamSink>", sizeof(IoOwn<WritableStreamSink>)); |
| 2407 | } |
| 2408 | } |
| 2409 | KJ_IF_SOME(writerLocked, writeState.tryGetUnsafe<WriterLocked>()) { |
| 2410 | tracker.trackField("writerLocked", writerLocked); |
| 2411 | } |
| 2412 | tracker.trackField("pendingAbort", maybePendingAbort); |
| 2413 | tracker.trackField("maybeClosureWaitable", maybeClosureWaitable); |
| 2414 | |
| 2415 | for (auto& event: queue) { |
| 2416 | tracker.trackField("event", event); |
| 2417 | } |
| 2418 | } |
| 2419 | |
| 2420 | kj::StringPtr ReadableStreamInternalController::jsgGetMemoryName() const { |
| 2421 | return "ReadableStreamInternalController"_kjc; |
| 2422 | } |
| 2423 | |
| 2424 | size_t ReadableStreamInternalController::jsgGetMemorySelfSize() const { |
| 2425 | return sizeof(ReadableStreamInternalController); |
| 2426 | } |
| 2427 | |
| 2428 | void ReadableStreamInternalController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 2429 | KJ_SWITCH_ONEOF(state) { |
| 2430 | KJ_CASE_ONEOF(closed, StreamStates::Closed) {} |
| 2431 | KJ_CASE_ONEOF(error, StreamStates::Errored) { |
| 2432 | tracker.trackField("error", error); |
| 2433 | } |
| 2434 | KJ_CASE_ONEOF(readable, Readable) { |
| 2435 | // Ideally we'd be able to track the size of any pending reads held in the source's |
| 2436 | // queue but since it is behind an IoOwn and we won't be holding the IoContext here, |
| 2437 | // we can't. |
| 2438 | tracker.trackFieldWithSize( |
| 2439 | "IoOwn<ReadableStreamSource>", sizeof(IoOwn<ReadableStreamSource>)); |
| 2440 | } |
| 2441 | } |
| 2442 | KJ_SWITCH_ONEOF(readState) { |
| 2443 | KJ_CASE_ONEOF(unlocked, Unlocked) {} |
| 2444 | KJ_CASE_ONEOF(locked, Locked) {} |
| 2445 | KJ_CASE_ONEOF(pipeLocked, PipeLocked) {} |
| 2446 | KJ_CASE_ONEOF(readerLocked, ReaderLocked) { |
| 2447 | tracker.trackField("readerLocked", readerLocked); |
| 2448 | } |
| 2449 | } |
| 2450 | } |
| 2451 | |
| 2452 | } // namespace workerd::api |