Skip to content
File

Blob: src/workerd/api/streams/internal.c++

97.9 KB
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 
19namespace workerd::api {
20 
21namespace {
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 
36kj::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++.
54class 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 
255kj::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()`.
270class 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 
286class 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 
375kj::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 
386kj::Maybe<uint64_t> ReadableStreamSource::tryGetLength(StreamEncoding encoding) {
387 return kj::none;
388}
389 
390kj::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 
404kj::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 
419void ReadableStreamSource::cancel(kj::Exception reason) {}
420 
421kj::Maybe<ReadableStreamSource::Tee> ReadableStreamSource::tryTee(uint64_t limit) {
422 return kj::none;
423}
424 
425kj::Maybe<kj::Promise<DeferredProxy<void>>> WritableStreamSink::tryPumpFrom(
426 ReadableStreamSource& input, bool end) {
427 return kj::none;
428}
429 
430// =======================================================================================
431 
432ReadableStreamInternalController::~ReadableStreamInternalController() noexcept(false) {
433 if (readState.is<ReaderLocked>()) {
434 readState.transitionTo<Unlocked>();
435 }
436}
437 
438jsg::Ref<ReadableStream> ReadableStreamInternalController::addRef() {
439 return KJ_ASSERT_NONNULL(owner).addRef();
440}
441 
442kj::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 
689kj::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 
786jsg::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 
807jsg::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 
820void 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 
834void 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 
846void 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 
858ReadableStreamController::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 
910kj::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 
947bool 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 
975void 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 
1000void WritableStreamInternalController::Writable::abort(kj::Exception&& ex) {
1001 canceler.cancel(ex.clone());
1002 sink->abort(kj::mv(ex));
1003}
1004 
1005WritableStreamInternalController::~WritableStreamInternalController() noexcept(false) {
1006 if (writeState.is<WriterLocked>()) {
1007 writeState.transitionTo<Unlocked>();
1008 }
1009}
1010 
1011jsg::Ref<WritableStream> WritableStreamInternalController::addRef() {
1012 return KJ_ASSERT_NONNULL(owner).addRef();
1013}
1014 
1015jsg::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 
1101void 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 
1110void 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 
1127void WritableStreamInternalController::setHighWaterMark(uint64_t highWaterMark) {
1128 maybeHighWaterMark = highWaterMark;
1129}
1130 
1131jsg::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 
1165jsg::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 
1187jsg::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 
1222jsg::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 
1230jsg::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 
1281kj::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 
1399kj::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 
1424void 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 
1448kj::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 
1467bool 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 
1499void 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 
1520bool 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 
1527bool WritableStreamInternalController::isPiping() {
1528 return state.is<IoOwn<Writable>>() && !queue.empty() && queue.back().event.is<kj::Own<Pipe>>();
1529}
1530 
1531bool WritableStreamInternalController::isErrored() {
1532 return state.is<StreamStates::Errored>();
1533}
1534 
1535void 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 
1550void 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 
1565void 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 
1572jsg::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 
1584void WritableStreamInternalController::finishClose(jsg::Lock& js) {
1585 KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) {
1586 pendingAbort->complete(js);
1587 }
1588 
1589 doClose(js);
1590}
1591 
1592void 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 
1603jsg::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 
1909bool 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 
1950jsg::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 
1987jsg::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 
2128void 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 
2152void 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 
2177void ReadableStreamInternalController::visitForGc(jsg::GcVisitor& visitor) {
2178 KJ_IF_SOME(locked, readState.tryGetUnsafe<ReaderLocked>()) {
2179 visitor.visit(locked);
2180 }
2181}
2182 
2183kj::Maybe<ReadableStreamController::PipeController&> ReadableStreamInternalController::
2184 tryPipeLock() {
2185 if (isLockedToReader()) {
2186 return kj::none;
2187 }
2188 return readState.transitionTo<PipeLocked>(*this);
2189}
2190 
2191bool ReadableStreamInternalController::PipeLocked::isClosed() {
2192 return inner.state.is<StreamStates::Closed>();
2193}
2194 
2195kj::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 
2203void 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 
2210void ReadableStreamInternalController::PipeLocked::close(jsg::Lock& js) {
2211 inner.doClose(js);
2212}
2213 
2214void ReadableStreamInternalController::PipeLocked::error(
2215 jsg::Lock& js, v8::Local<v8::Value> reason) {
2216 inner.doError(js, reason);
2217}
2218 
2219void 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 
2227kj::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 
2235jsg::Promise<ReadResult> ReadableStreamInternalController::PipeLocked::read(jsg::Lock& js) {
2236 return KJ_ASSERT_NONNULL(inner.read(js, kj::none));
2237}
2238 
2239jsg::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 
2274jsg::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 
2306kj::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 
2321kj::Own<ReadableStreamController> ReadableStreamInternalController::detach(
2322 jsg::Lock& js, bool ignoreDetached) {
2323 return newReadableStreamInternalController(
2324 IoContext::current(), KJ_ASSERT_NONNULL(removeSource(js, ignoreDetached)));
2325}
2326 
2327kj::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 
2368StreamEncoding ReadableStreamInternalController::getPreferredEncoding() {
2369 return state.tryGetUnsafe<Readable>()
2370 .map([](Readable& readable) {
2371 return readable->getPreferredEncoding();
2372 }).orDefault(StreamEncoding::IDENTITY);
2373}
2374 
2375kj::Own<ReadableStreamController> newReadableStreamInternalController(
2376 IoContext& ioContext, kj::Own<ReadableStreamSource> source) {
2377 return kj::heap<ReadableStreamInternalController>(ioContext.addObject(kj::mv(source)));
2378}
2379 
2380kj::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 
2389kj::StringPtr WritableStreamInternalController::jsgGetMemoryName() const {
2390 return "WritableStreamInternalController"_kjc;
2391}
2392 
2393size_t WritableStreamInternalController::jsgGetMemorySelfSize() const {
2394 return sizeof(WritableStreamInternalController);
2395}
2396void 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 
2420kj::StringPtr ReadableStreamInternalController::jsgGetMemoryName() const {
2421 return "ReadableStreamInternalController"_kjc;
2422}
2423 
2424size_t ReadableStreamInternalController::jsgGetMemorySelfSize() const {
2425 return sizeof(ReadableStreamInternalController);
2426}
2427 
2428void 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