Skip to content
File

Blob: src/workerd/api/streams/readable-source-adapter.c++

57.6 KB
1#include "readable-source-adapter.h"
2 
3#include "writable-sink.h"
4 
5#include <workerd/util/checked-queue.h>
6 
7#include <bit>
8 
9namespace workerd::api::streams {
10 
11namespace {
12// Per the ReadableStream spec, when a read(buf) is performed on a BYOB reader,
13// if the stream is already closed, we still need to return the allocated buffer
14// back to the caller, but it must be in a zero-length view. This utility function
15// does that. It takes the original allocation and wraps it into a new ArrayBuffer
16// instance that is wrapped by a zero-length view of the same type as the original
17// TypedArray we were given.
18jsg::BufferSource transferToEmptyBuffer(jsg::Lock& js, jsg::BufferSource buffer) {
19 KJ_DASSERT(!buffer.isDetached() && buffer.canDetach(js));
20 auto backing = buffer.detach(js);
21 backing.limit(0);
22 auto buf = jsg::BufferSource(js, kj::mv(backing));
23 KJ_DASSERT(buf.size() == 0);
24 return kj::mv(buf);
25}
26} // namespace
27 
28// The Active state maintains a queue of tasks, such as read or close operations. Each task
29// contains a promise-returning function object and a fulfiller. When the first task is
30// enqueued, the active state begins processing the queue asynchronously. Each function
31// is invoked in order, its promise awaited, and the result passed to the fulfiller. The
32// fulfiller notifies the code which enqueued the task that the task has completed. In
33// this way, read and close operations are safely executed in serial, even if one operation
34// is called before the previous completes. This mechanism satisfies KJ's restriction on
35// concurrent operations on streams.
36struct ReadableStreamSourceJsAdapter::Active {
37 struct Task {
38 kj::Function<kj::Promise<size_t>()> task;
39 kj::Own<kj::PromiseFulfiller<size_t>> fulfiller;
40 Task(kj::Function<kj::Promise<size_t>()> task, kj::Own<kj::PromiseFulfiller<size_t>> fulfiller)
41 : task(kj::mv(task)),
42 fulfiller(kj::mv(fulfiller)) {}
43 KJ_DISALLOW_COPY_AND_MOVE(Task);
44 };
45 using TaskQueue = workerd::util::Queue<kj::Own<Task>>;
46 
47 kj::Own<ReadableSource> source;
48 kj::Canceler canceler;
49 TaskQueue queue;
50 bool canceled = false;
51 bool running = false;
52 bool closePending = false;
53 kj::Maybe<kj::Exception> pendingCancel;
54 
55 Active(kj::Own<ReadableSource> source): source(kj::mv(source)) {}
56 KJ_DISALLOW_COPY_AND_MOVE(Active);
57 ~Active() noexcept(false) {
58 // When the Active is dropped, we cancel any remaining pending reads and
59 // abort the sink.
60 cancel(KJ_EXCEPTION(DISCONNECTED, "Writable stream is canceled or closed."));
61 
62 // Check invariants for safety.
63 // 1. Our canceler should be empty because we canceled it.
64 KJ_DASSERT(canceler.isEmpty());
65 // 2. The write queue should be empty.
66 KJ_DASSERT(queue.empty());
67 }
68 
69 // Explicitly cancel all in-flight and pending tasks in the queue.
70 // This is a non-op if cancel has already been called.
71 void cancel(kj::Exception&& exception) {
72 if (canceled) return;
73 canceled = true;
74 // 1. Cancel our in-flight "runLoop", if any.
75 pendingCancel = exception.clone();
76 canceler.cancel(exception.clone());
77 // 2. Drop our queue of pending tasks.
78 queue.drainTo(
79 [&exception](kj::Own<Task>&& task) { task->fulfiller->reject(exception.clone()); });
80 // 3. Cancel and drop the source itself. We're done with it.
81 if (exception.getType() != kj::Exception::Type::DISCONNECTED) {
82 source->cancel(kj::mv(exception));
83 }
84 auto dropped KJ_UNUSED = kj::mv(source);
85 }
86 
87 kj::Promise<size_t> enqueue(kj::Function<kj::Promise<size_t>()> task) {
88 KJ_DASSERT(!canceled, "cannot enqueue tasks on a canceled queue");
89 auto paf = kj::newPromiseAndFulfiller<size_t>();
90 queue.push(kj::heap<Task>(kj::mv(task), kj::mv(paf.fulfiller)));
91 if (!running) {
92 IoContext::current().addTask(canceler.wrap(run()));
93 }
94 return kj::mv(paf.promise);
95 }
96 
97 kj::Promise<void> run() {
98 KJ_DEFER(running = false);
99 running = true;
100 while (!queue.empty() && !canceled) {
101 auto task = KJ_ASSERT_NONNULL(queue.pop());
102 KJ_DEFER({
103 if (task->fulfiller->isWaiting()) {
104 KJ_IF_SOME(pending, pendingCancel) {
105 task->fulfiller->reject(kj::mv(pending));
106 } else {
107 task->fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "Task was canceled."));
108 }
109 }
110 });
111 bool taskFailed = false;
112 try {
113 task->fulfiller->fulfill(co_await task->task());
114 } catch (...) {
115 auto ex = kj::getCaughtExceptionAsKj();
116 task->fulfiller->reject(kj::mv(ex));
117 taskFailed = true;
118 }
119 // If the task failed, we exit the loop. We're going to abort the
120 // entire remaining queue anyway so there's no point in continuing.
121 if (taskFailed) co_return;
122 }
123 }
124};
125 
126ReadableStreamSourceJsAdapter::ReadableStreamSourceJsAdapter(
127 jsg::Lock& js, IoContext& ioContext, kj::Own<ReadableSource> source)
128 : state(State::create<Open>(ioContext.addObject(kj::heap<Active>(kj::mv(source))))),
129 selfRef(kj::rc<WeakRef<ReadableStreamSourceJsAdapter>>(
130 kj::Badge<ReadableStreamSourceJsAdapter>{}, *this)) {}
131 
132ReadableStreamSourceJsAdapter::~ReadableStreamSourceJsAdapter() noexcept(false) {
133 selfRef->invalidate();
134}
135 
136void ReadableStreamSourceJsAdapter::cancel(kj::Exception exception) {
137 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
138 open.active->cancel(exception.clone());
139 }
140 state.forceTransitionTo<kj::Exception>(kj::mv(exception));
141}
142 
143void ReadableStreamSourceJsAdapter::cancel(jsg::Lock& js, const jsg::JsValue& reason) {
144 cancel(js.exceptionToKj(reason));
145}
146 
147void ReadableStreamSourceJsAdapter::shutdown(jsg::Lock& js) {
148 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
149 open.active->cancel(KJ_EXCEPTION(DISCONNECTED, "Stream was shut down."));
150 state.transitionTo<Closed>();
151 }
152 // If we are are already closed or canceled, this is a no-op.
153}
154 
155bool ReadableStreamSourceJsAdapter::isClosed() {
156 return state.is<Closed>();
157}
158 
159kj::Maybe<const kj::Exception&> ReadableStreamSourceJsAdapter::isCanceled() {
160 return state.tryGetErrorUnsafe();
161}
162 
163jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> ReadableStreamSourceJsAdapter::read(
164 jsg::Lock& js, ReadOptions options) {
165 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
166 // Really should not have been called if errored but just in case,
167 // return a rejected promise.
168 return js.rejectedPromise<ReadResult>(js.exceptionToJs(exception.clone()));
169 }
170 
171 if (state.is<Closed>()) {
172 // We are already in a closed state. This is a no-op, just return
173 // an empty buffer.
174 return js.resolvedPromise(ReadResult{
175 .buffer = transferToEmptyBuffer(js, kj::mv(options.buffer)),
176 .done = true,
177 });
178 }
179 
180 auto& open = state.requireActiveUnsafe();
181 // Deference the IoOwn once to get the active state.
182 Active& active = *open.active;
183 
184 // If close is pending, we cannot accept any more reads.
185 // Treat them as if the stream is closed.
186 if (active.closePending) {
187 return js.resolvedPromise(ReadResult{
188 .buffer = transferToEmptyBuffer(js, kj::mv(options.buffer)),
189 .done = true,
190 });
191 }
192 
193 // Ok, we are in a readable state, there are no pending closes.
194 // Let's enqueue our read request.
195 auto& ioContext = IoContext::current();
196 
197 auto buffer = kj::mv(options.buffer);
198 auto elementSize = buffer.getElementSize();
199 
200 // The buffer size should always be a multiple of the element size and should
201 // always be at least as large as minBytes. This should be handled for us by
202 // the jsg::BufferSource, but just to be safe, we will double-check with a
203 // debug assert here.
204 KJ_DASSERT(buffer.size() % elementSize == 0);
205 
206 auto minBytes = kj::min(options.minBytes.orDefault(elementSize), buffer.size());
207 // We want to be sure that minBytes is a multiple of the element size
208 // of the buffer, otherwise we might never be able to satisfy the request
209 // correcty. If the caller provided a minBytes, and it is not a multiple
210 // of the element size, we will round it up to the next multiple.
211 if (elementSize > 1) {
212 minBytes = minBytes + (elementSize - (minBytes % elementSize)) % elementSize;
213 }
214 
215 // Note: We do not enforce that the source must provide at least minBytes
216 // if available here as that is part of the contract of the source itself.
217 // We will simply pass minBytes along to the source and it is up to the
218 // source to honor it. We do, however, enforce that the source must
219 // never return more than the size of the buffer we provided.
220 
221 // We only pass a kj::ArrayPtr to the buffer into the read call, keeping
222 // the actual buffer instance alive by attaching it to the JS promise
223 // chain that follows the read in order to keep it alive.
224 auto promise = active.enqueue(kj::coCapture(
225 [&active, buffer = buffer.asArrayPtr(), minBytes]() mutable -> kj::Promise<size_t> {
226 // TODO(soon): The underlying kj streams API now supports passing the
227 // kj::ArrayPtr directly to the read call, but ReadableStreamSource has
228 // not yet been updated to do so. When it is, we can update this read to
229 // pass `buffer` directly rather than passing the begin() and size().
230 co_return co_await active.source->read(buffer, minBytes);
231 }));
232 return ioContext
233 .awaitIo(js, kj::mv(promise),
234 [buffer = kj::mv(buffer), self = selfRef.addRef()](jsg::Lock& js,
235 size_t bytesRead) mutable -> jsg::Promise<ReadableStreamSourceJsAdapter::ReadResult> {
236 // If the bytesRead is 0, that indicates the stream is closed. We will
237 // move the stream to a closed state and return the empty buffer.
238 if (bytesRead == 0) {
239 self->runIfAlive([](ReadableStreamSourceJsAdapter& self) {
240 KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) {
241 open.active->closePending = true;
242 }
243 });
244 return js.resolvedPromise(ReadResult{
245 .buffer = transferToEmptyBuffer(js, kj::mv(buffer)),
246 .done = true,
247 });
248 }
249 KJ_DASSERT(bytesRead <= buffer.size());
250 
251 // If bytesRead is not a multiple of the element size, that indicates
252 // that the source either read less than minBytes (and ended), or is
253 // simply unable to satisfy the element size requirement. We cannot
254 // provide a partial element to the caller, so reject the read.
255 if (bytesRead % buffer.getElementSize() != 0) {
256 return js.rejectedPromise<ReadResult>(
257 js.typeError(kj::str("The underlying stream failed to provide a multiple of the "
258 "target element size ",
259 buffer.getElementSize())));
260 }
261 
262 auto backing = buffer.detach(js);
263 backing.limit(bytesRead);
264 return js.resolvedPromise(ReadResult{
265 .buffer = jsg::BufferSource(js, kj::mv(backing)),
266 .done = false,
267 });
268 })
269 .catch_(js,
270 [self = selfRef.addRef()](
271 jsg::Lock& js, jsg::Value exception) -> ReadableStreamSourceJsAdapter::ReadResult {
272 // If an error occurred while reading, we need to transition the adapter
273 // to the canceled state, but only if the adapter is still alive.
274 auto error = jsg::JsValue(exception.getHandle(js));
275 self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); });
276 js.throwException(kj::mv(exception));
277 });
278}
279 
280// Transitions the adapter into the closing state. Once the read queue
281// is empty, we will close the source and transition to the closed state.
282jsg::Promise<void> ReadableStreamSourceJsAdapter::close(jsg::Lock& js) {
283 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
284 // Really should not have been called if errored but just in case,
285 // return a rejected promise.
286 return js.rejectedPromise<void>(js.exceptionToJs(exception.clone()));
287 }
288 
289 if (state.is<Closed>()) {
290 // We are already in a closed state. This is a no-op. This really
291 // should not have been called if closed but just in case, return
292 // a resolved promise.
293 return js.resolvedPromise();
294 }
295 
296 auto& open = state.requireActiveUnsafe();
297 auto& ioContext = IoContext::current();
298 auto& active = *open.active;
299 
300 if (active.closePending) {
301 return js.rejectedPromise<void>(js.typeError("Close already pending, cannot close again."));
302 }
303 
304 active.closePending = true;
305 auto promise = active.enqueue([]() -> kj::Promise<size_t> { co_return 0; });
306 
307 return ioContext
308 .awaitIo(js, kj::mv(promise), [self = selfRef.addRef()](jsg::Lock&, size_t) {
309 self->runIfAlive(
310 [](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo<Closed>(); });
311 }).catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value&& exception) {
312 // Likewise, while nothing should be waiting on the ready promise, we
313 // should still reject it just in case.
314 auto error = jsg::JsValue(exception.getHandle(js));
315 self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); });
316 js.throwException(kj::mv(exception));
317 });
318}
319 
320jsg::Promise<jsg::JsRef<jsg::JsString>> ReadableStreamSourceJsAdapter::readAllText(
321 jsg::Lock& js, uint64_t limit) {
322 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
323 // Really should not have been called if errored but just in case,
324 // return a rejected promise.
325 return js.rejectedPromise<jsg::JsRef<jsg::JsString>>(js.exceptionToJs(exception.clone()));
326 }
327 
328 if (state.is<Closed>()) {
329 // We are already in a closed state. This is a no-op. This really
330 // should not have been called if closed but just in case, return
331 // a resolved promise.
332 return js.resolvedPromise(jsg::JsRef(js, js.str()));
333 }
334 
335 auto& open = state.requireActiveUnsafe();
336 auto& ioContext = IoContext::current();
337 auto& active = *open.active;
338 
339 if (active.closePending) {
340 return js.rejectedPromise<jsg::JsRef<jsg::JsString>>(
341 js.typeError("Close already pending, cannot read."));
342 }
343 active.closePending = true;
344 
345 struct Holder {
346 kj::Maybe<kj::String> result;
347 };
348 auto holder = kj::heap<Holder>();
349 
350 auto promise = active.enqueue([&active, &holder = *holder, limit]() -> kj::Promise<size_t> {
351 auto str = co_await active.source->readAllText(limit);
352 size_t amount = str.size();
353 holder.result = kj::mv(str);
354 co_return amount;
355 });
356 
357 return ioContext
358 .awaitIo(js, kj::mv(promise),
359 [self = selfRef.addRef(), holder = kj::mv(holder)](jsg::Lock& js, size_t amount) {
360 self->runIfAlive(
361 [&](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo<Closed>(); });
362 KJ_IF_SOME(result, holder->result) {
363 KJ_DASSERT(result.size() == amount);
364 return jsg::JsRef(js, js.str(result));
365 } else {
366 return jsg::JsRef(js, js.str());
367 }
368 })
369 .catch_(js,
370 [self = selfRef.addRef()](
371 jsg::Lock& js, jsg::Value&& exception) -> jsg::JsRef<jsg::JsString> {
372 // Likewise, while nothing should be waiting on the ready promise, we
373 // should still reject it just in case.
374 auto error = jsg::JsValue(exception.getHandle(js));
375 self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); });
376 js.throwException(kj::mv(exception));
377 });
378}
379 
380jsg::Promise<jsg::BufferSource> ReadableStreamSourceJsAdapter::readAllBytes(
381 jsg::Lock& js, uint64_t limit) {
382 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
383 // Really should not have been called if errored but just in case,
384 // return a rejected promise.
385 return js.rejectedPromise<jsg::BufferSource>(js.exceptionToJs(exception.clone()));
386 }
387 
388 if (state.is<Closed>()) {
389 // We are already in a closed state. This is a no-op. This really
390 // should not have been called if closed but just in case, return
391 // a resolved promise.
392 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0);
393 return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing)));
394 }
395 
396 auto& open = state.requireActiveUnsafe();
397 auto& ioContext = IoContext::current();
398 auto& active = *open.active;
399 
400 if (active.closePending) {
401 return js.rejectedPromise<jsg::BufferSource>(
402 js.typeError("Close already pending, cannot read."));
403 }
404 active.closePending = true;
405 
406 struct Holder {
407 kj::Maybe<kj::Array<const kj::byte>> result;
408 };
409 auto holder = kj::heap<Holder>();
410 
411 auto promise = active.enqueue([&active, &holder = *holder, limit]() -> kj::Promise<size_t> {
412 auto str = co_await active.source->readAllBytes(limit);
413 size_t amount = str.size();
414 holder.result = kj::mv(str);
415 co_return amount;
416 });
417 
418 return ioContext
419 .awaitIo(js, kj::mv(promise),
420 [self = selfRef.addRef(), holder = kj::mv(holder)](jsg::Lock& js, size_t amount) {
421 self->runIfAlive(
422 [&](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo<Closed>(); });
423 KJ_IF_SOME(result, holder->result) {
424 KJ_DASSERT(result.size() == amount);
425 // We have to copy the data into the backing store because of the
426 // v8 sandboxing rules.
427 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, amount);
428 backing.asArrayPtr().copyFrom(result);
429 return jsg::BufferSource(js, kj::mv(backing));
430 } else {
431 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0);
432 return jsg::BufferSource(js, kj::mv(backing));
433 }
434 })
435 .catch_(js,
436 [self = selfRef.addRef()](jsg::Lock& js, jsg::Value&& exception) -> jsg::BufferSource {
437 // Likewise, while nothing should be waiting on the ready promise, we
438 // should still reject it just in case.
439 auto error = jsg::JsValue(exception.getHandle(js));
440 self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); });
441 js.throwException(kj::mv(exception));
442 });
443}
444 
445kj::Maybe<uint64_t> ReadableStreamSourceJsAdapter::tryGetLength(StreamEncoding encoding) {
446 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
447 return open.active->source->tryGetLength(encoding);
448 }
449 return kj::none;
450}
451 
452kj::Maybe<ReadableStreamSourceJsAdapter::Tee> ReadableStreamSourceJsAdapter::tryTee(
453 jsg::Lock& js, uint64_t limit) {
454 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
455 js.throwException(js.exceptionToJs(exception.clone()));
456 }
457 
458 if (state.is<Closed>()) {
459 // We are already closed, cannot tee.
460 return kj::none;
461 }
462 
463 auto& open = state.requireActiveUnsafe();
464 auto& active = *open.active;
465 // If we are closing, or have pending tasks, we cannot tee.
466 JSG_REQUIRE(!active.closePending && !active.running && active.queue.empty(), Error,
467 "Cannot tee a stream that is closing or has pending reads.");
468 auto tee = active.source->tee(limit);
469 auto& ioContext = IoContext::current();
470 state.transitionTo<Closed>();
471 return Tee{
472 .branch1 = kj::heap<ReadableStreamSourceJsAdapter>(js, ioContext, kj::mv(tee.branch1)),
473 .branch2 = kj::heap<ReadableStreamSourceJsAdapter>(js, ioContext, kj::mv(tee.branch2)),
474 };
475}
476 
477// ===============================================================================================
478 
479struct ReadableSourceKjAdapter::Active {
480 IoContext& ioContext;
481 jsg::Ref<ReadableStream> stream;
482 jsg::Ref<ReadableStreamDefaultReader> reader;
483 kj::Canceler canceler;
484 
485 struct Idle {
486 static constexpr kj::StringPtr NAME KJ_UNUSED = "idle"_kj;
487 };
488 struct Readable {
489 static constexpr kj::StringPtr NAME KJ_UNUSED = "readable"_kj;
490 // Previously read but unconsumed bytes. We keep these around for the next read call.
491 kj::Array<const kj::byte> data;
492 kj::ArrayPtr<const kj::byte> view;
493 
494 Readable(kj::Array<const kj::byte>&& data): data(kj::mv(data)), view(this->data) {}
495 };
496 struct Reading {
497 static constexpr kj::StringPtr NAME KJ_UNUSED = "reading"_kj;
498 // The contract for ReadableStreamSource is that there can be only one read() in-flight
499 // against the underlying stream at a time.
500 };
501 struct Done {
502 static constexpr kj::StringPtr NAME KJ_UNUSED = "done"_kj;
503 // If a read returns fewer than the requested minBytes, that indicates the stream is done. We
504 // make note of that here to prevent any further reads. We cannot transition to the closed
505 // state in the promise chain of the read because the adapter will cancel the read promise
506 // itself once Active is destroyed, and that would be a bad thing.
507 };
508 struct Canceling {
509 static constexpr kj::StringPtr NAME KJ_UNUSED = "canceling"_kj;
510 kj::Exception exception;
511 };
512 struct Canceled {
513 static constexpr kj::StringPtr NAME KJ_UNUSED = "canceled"_kj;
514 kj::Exception exception;
515 };
516 
517 // Inner state machine for tracking read operation state:
518 // Idle -> Reading (start read)
519 // Reading -> Idle (read complete, no leftover)
520 // Reading -> Readable (read complete, has leftover)
521 // Reading -> Done (read returned less than minBytes)
522 // Any -> Canceling (error during read)
523 // Any -> Canceled (explicit cancel)
524 // Done, Canceling, and Canceled are terminal states.
525 using InnerState = StateMachine<TerminalStates<Done, Canceling, Canceled>,
526 Idle,
527 Readable,
528 Reading,
529 Done,
530 Canceling,
531 Canceled>;
532 InnerState state;
533 
534 Active(jsg::Lock& js, IoContext& ioContext, jsg::Ref<ReadableStream> stream);
535 KJ_DISALLOW_COPY_AND_MOVE(Active);
536 ~Active() noexcept(false);
537 
538 void cancel(kj::Exception reason);
539};
540 
541// The ReadContext struct holds all the state needed to perform a read,
542// including the JS objects that need to be kept alive during the
543// read operation, the buffer we are reading into, and the total
544// number of bytes read so far. This must be kept alive until the
545// read is fully complete and returned back to the adapter when
546// the read is complete.
547//
548// Ownership of the ReadContext is passed into the isolate lock and
549// held by JS promise continuations, so it must not contain any
550// kj I/O objects or references without an IoOwn wrapper.
551struct ReadableSourceKjAdapter::ReadContext {
552 jsg::Ref<ReadableStream> stream;
553 jsg::Ref<ReadableStreamDefaultReader> reader;
554 kj::ArrayPtr<kj::byte> buffer;
555 // Only set to back the buffer if we need to keep it alive.
556 kj::Maybe<kj::Array<kj::byte>> backingBuffer;
557 size_t totalRead = 0;
558 size_t minBytes = 0;
559 kj::Maybe<Active::Readable> maybeLeftOver;
560 // We keep a weak reference to the adapter itself so we can track
561 // whether it is still alive while we are in a JS promise chain.
562 // If the adapter is gone, or transitions to a closed or canceled
563 // state we will abandon the read. If the ref is not set, then we
564 // are in a pump operation and do not need to check for liveness.
565 kj::Maybe<kj::Rc<WeakRef<ReadableSourceKjAdapter>>> adapterRef;
566 
567 void reset() {
568 // Resetting is only allowed if we have the backing buffer.
569 buffer = KJ_ASSERT_NONNULL(backingBuffer);
570 totalRead = 0;
571 minBytes = 0;
572 maybeLeftOver = kj::none;
573 }
574};
575 
576namespace {
577constexpr size_t kMinRemainingForAdditionalRead = 512;
578 
579jsg::Ref<ReadableStreamDefaultReader> initReader(jsg::Lock& js, jsg::Ref<ReadableStream>& stream) {
580 JSG_REQUIRE(!stream->isLocked(), TypeError, "ReadableStream is locked.");
581 JSG_REQUIRE(!stream->isDisturbed(), TypeError, "ReadableStream is disturbed.");
582 auto reader = stream->getReader(js, kj::none);
583 return kj::mv(KJ_ASSERT_NONNULL(reader.tryGet<jsg::Ref<ReadableStreamDefaultReader>>()));
584}
585 
586using JsByteSource = kj::OneOf<jsg::JsRef<jsg::JsString>,
587 jsg::JsRef<jsg::JsArrayBuffer>,
588 jsg::JsRef<jsg::JsArrayBufferView>>;
589 
590kj::Maybe<JsByteSource> tryExtractJsByteSource(jsg::Lock& js, const jsg::JsValue& jsval) {
591 KJ_IF_SOME(abView, jsval.tryCast<jsg::JsArrayBuffer>()) {
592 return kj::Maybe(jsg::JsRef(js, abView));
593 } else KJ_IF_SOME(ab, jsval.tryCast<jsg::JsArrayBufferView>()) {
594 return kj::Maybe(jsg::JsRef(js, ab));
595 } else KJ_IF_SOME(str, jsval.tryCast<jsg::JsString>()) {
596 return kj::Maybe(jsg::JsRef(js, str));
597 }
598 return kj::none;
599}
600 
601// Copies as much data from source into the context as possible, returning
602// the number of bytes copied.
603kj::Maybe<kj::Array<const kj::byte>> copyFromSource(
604 jsg::Lock& js, ReadableSourceKjAdapter::ReadContext& context, const JsByteSource& source) {
605 KJ_SWITCH_ONEOF(source) {
606 KJ_CASE_ONEOF(str, jsg::JsRef<jsg::JsString>) {
607 auto view = str.getHandle(js);
608 size_t len = view.length(js);
609 size_t toCopy = kj::min(len, context.buffer.size());
610 
611 if (toCopy == 0) {
612 return kj::none;
613 }
614 
615 if (toCopy < len) {
616 // We are going to have left-over data. Unfortunately in this case
617 // we have to copy the data twice... once into a kj::String and
618 // again into our buffer. This is because the V8 string UTF-8
619 // write API does not support partial writes with an offset.
620 auto data = view.toUSVString(js);
621 context.buffer.first(toCopy).copyFrom(data.asBytes().first(toCopy));
622 context.totalRead += toCopy;
623 context.buffer = context.buffer.slice(toCopy);
624 KJ_DASSERT(context.buffer.size() == 0);
625 return kj::Maybe(data.asBytes().slice(toCopy).attach(kj::mv(data)));
626 }
627 
628 // We can copy everything in one go. Yay! This is great because we
629 // can avoid a double copy here.
630 auto ret KJ_UNUSED = view.writeInto(js, context.buffer.asChars().first(toCopy),
631 jsg::JsString::WriteFlags::REPLACE_INVALID_UTF8);
632 KJ_DASSERT(ret.written == toCopy);
633 context.totalRead += toCopy;
634 context.buffer = context.buffer.slice(toCopy);
635 return kj::none;
636 }
637 KJ_CASE_ONEOF(ab, jsg::JsRef<jsg::JsArrayBuffer>) {
638 auto src = ab.getHandle(js).asArrayPtr();
639 size_t toCopy = kj::min(src.size(), context.buffer.size());
640 if (toCopy == 0) {
641 return kj::none;
642 }
643 
644 context.buffer.first(toCopy).copyFrom(src.first(toCopy));
645 context.totalRead += toCopy;
646 context.buffer = context.buffer.slice(toCopy);
647 
648 if (toCopy < src.size()) {
649 KJ_DASSERT(context.buffer.size() == 0);
650 // TODO(mpk): For now, we have to copy the left-over data into a new array.
651 // Why? I'm happy you asked! Because the src is backed by a
652 // v8::BackingStore protected by the v8 sandboxing rules and we
653 // don't yet have the memory protection key logic in place to safely
654 // share that memory outside of the v8 heap. For now, copy. Later
655 // we can revisit this to hopefully avoid the additinal copy.
656 return kj::Maybe(kj::heapArray(src.slice(toCopy)));
657 }
658 
659 return kj::none;
660 }
661 KJ_CASE_ONEOF(view, jsg::JsRef<jsg::JsArrayBufferView>) {
662 auto src = view.getHandle(js).asArrayPtr();
663 size_t toCopy = kj::min(src.size(), context.buffer.size());
664 if (toCopy == 0) {
665 // Copy nothing. Return 0.
666 return kj::none;
667 }
668 
669 context.buffer.first(toCopy).copyFrom(src.first(toCopy));
670 context.totalRead += toCopy;
671 context.buffer = context.buffer.slice(toCopy);
672 
673 if (toCopy < src.size()) {
674 KJ_DASSERT(context.buffer.size() == 0);
675 return kj::Maybe(kj::heapArray(src.slice(toCopy)));
676 }
677 
678 return kj::none;
679 }
680 }
681 KJ_UNREACHABLE;
682}
683} // namespace
684 
685ReadableSourceKjAdapter::Active::Active(
686 jsg::Lock& js, IoContext& ioContext, jsg::Ref<ReadableStream> stream)
687 : ioContext(ioContext),
688 stream(kj::mv(stream)),
689 reader(initReader(js, this->stream)),
690 state(InnerState::create<Idle>()) {}
691 
692ReadableSourceKjAdapter::Active::~Active() noexcept(false) {
693 cancel(KJ_EXCEPTION(DISCONNECTED, "ReadableSourceKjAdapter is canceled."));
694}
695 
696void ReadableSourceKjAdapter::Active::cancel(kj::Exception reason) {
697 if (state.is<Canceled>()) {
698 return;
699 }
700 bool wasDone = state.is<Done>();
701 state.forceTransitionTo<Canceled>(reason.clone());
702 canceler.cancel(reason.clone());
703 if (!wasDone) {
704 // If the previous read indicated that it was the last read, then
705 // the reader will have already been dropped. We do not need to
706 // cancel it here.
707 ioContext.addTask(ioContext.run([readable = kj::mv(stream), reader = kj::mv(reader),
708 exception = kj::mv(reason)](jsg::Lock& js) mutable {
709 auto& ioContext = IoContext::current();
710 auto error = js.exceptionToJsValue(kj::mv(exception));
711 auto promise = reader->cancel(js, error.getHandle(js));
712 return ioContext.awaitJs(js, kj::mv(promise));
713 }));
714 }
715}
716 
717ReadableSourceKjAdapter::ReadableSourceKjAdapter(
718 jsg::Lock& js, IoContext& ioContext, jsg::Ref<ReadableStream> stream, Options options)
719 : state(KjState::create<KjOpen>(kj::heap<Active>(js, ioContext, kj::mv(stream)))),
720 options(options),
721 selfRef(
722 kj::rc<WeakRef<ReadableSourceKjAdapter>>(kj::Badge<ReadableSourceKjAdapter>{}, *this)) {}
723 
724ReadableSourceKjAdapter::~ReadableSourceKjAdapter() noexcept(false) {
725 selfRef->invalidate();
726}
727 
728jsg::Promise<kj::Own<ReadableSourceKjAdapter::ReadContext>> ReadableSourceKjAdapter::readInternal(
729 jsg::Lock& js, kj::Own<ReadContext> context, MinReadPolicy minReadPolicy) {
730 auto& ioContext = IoContext::current();
731 // Pay close attention to the lambda captures here. There are no raw references
732 // captured! The adapter itself may be destroyed or closed while we are in the
733 // promise chain below, so we have to be careful to only hold weak references
734 // and pass ownership of the context along the promise chain.
735 //
736 // The other important thing here is to remember that everything in this function
737 // is running within the isolate lock. The idea is to keep the entire read of the
738 // underlying stream entirely within the lock so that we don't have to bounce
739 // in and out of the isolate lock multiple times. We only return to the kj world
740 // once the entire read is complete.
741 //
742 // Note the uses of addFunctor below. This is important because it ensures
743 // that the promise continuations are run within the correct IoContext.
744 return context->reader->read(js).then(js,
745 ioContext.addFunctor([context = kj::mv(context), minReadPolicy](jsg::Lock& js,
746 ReadResult result) mutable -> jsg::Promise<kj::Own<ReadContext>> {
747 if (result.done || result.value == kj::none) {
748 // Stream is ended.
749 return js.resolvedPromise(kj::mv(context));
750 }
751 
752 auto& value = KJ_ASSERT_NONNULL(result.value);
753 
754 // Ok, we have some data. Let's make sure it is bytes.
755 // We accept either an ArrayBuffer, ArrayBufferView, or string.
756 auto jsval = jsg::JsValue(value.getHandle(js));
757 KJ_IF_SOME(result, tryExtractJsByteSource(js, jsval)) {
758 // Process the resulting data.
759 KJ_IF_SOME(leftOver, copyFromSource(js, *context, result)) {
760 KJ_ASSERT(context->buffer.size() == 0);
761 if (leftOver.size() > 0) {
762 context->maybeLeftOver = Active::Readable(kj::mv(leftOver));
763 } else {
764 context->maybeLeftOver = kj::none;
765 }
766 return js.resolvedPromise(kj::mv(context));
767 }
768 
769 // At this point, we should have no left over data.
770 KJ_DASSERT(context->maybeLeftOver == kj::none);
771 
772 // If the buffer is exactly full (the chunk filled it perfectly), we're done.
773 if (context->buffer.size() == 0) {
774 return js.resolvedPromise(kj::mv(context));
775 }
776 
777 // We might continue reading only if the adapter is still alive and
778 // in an active state...
779 bool continueReading = true;
780 KJ_IF_SOME(adapterRef, context->adapterRef) {
781 continueReading = adapterRef->isValid();
782 adapterRef->runIfAlive(
783 [&](ReadableSourceKjAdapter& adapter) { continueReading = adapter.state.isActive(); });
784 }
785 
786 // If we have satisfied the minimum read requirement and either
787 // (a) the minReadPolicy is IMMEDIATE or (b) there are fewer
788 // than 512 bytes left in the buffer, we will just return what we
789 // have. The idea here is that while we could just return what we have
790 // and let the caller call read again, that would be inefficient if
791 // the caller has a large buffer and is trying to read a lot of data.
792 // Instead of returning early with a minimally filled buffer, let's
793 // try to fill it up a bit more before returning. The 512 byte limit
794 // is somewhat arbitrary. The risk, of course, is that the next read
795 // will return too much data to fit into the buffer, which will then
796 // have to be stashed away as left over data. There's also a risk that
797 // the stream is slow and we end up with more latency waiting for
798 // the next chunk of data to arrive. In practice, this seems unlikely
799 // to be a problem. The IMMEDIATE policy is useful in the latter case,
800 // when the caller wants to get whatever data is available as soon
801 // as possible, even if it is just a small amount. The downside of the
802 // IMMEDIATE policy is that it can lead to a lot of small reads that
803 // are expensive because they have to grab the isolate lock each time.
804 bool minReadSatisfied = context->totalRead >= context->minBytes &&
805 (minReadPolicy == MinReadPolicy::IMMEDIATE ||
806 context->buffer.size() < kMinRemainingForAdditionalRead);
807 
808 if (!continueReading || minReadSatisfied) {
809 return js.resolvedPromise(kj::mv(context));
810 }
811 
812 // We still have not satisfied the minimum read requirement or we are
813 // trying to fill up a larger buffer. We will need to read more. Let's
814 // call readInternal again to get the next chunk of data. Keep in mind
815 // that this is not a true recursive call because readInternal returns
816 // a jsg::Promise. We're just chaining the promises together here.
817 return readInternal(js, kj::mv(context), minReadPolicy);
818 }
819 
820 // Oooo, invalid type. We cannot handle this and must treat this as a fatal error.
821 // We will cancel the stream and return an error.
822 auto error = js.typeError("ReadableStream provided a non-bytes value. Only ArrayBuffer, "
823 "ArrayBufferView, or string are supported.");
824 context->reader->cancel(js, error);
825 return js.rejectedPromise<kj::Own<ReadContext>>(error);
826 }),
827 ioContext.addFunctor([](jsg::Lock& js, jsg::Value exception) {
828 // In this case, the reader should already be in an errored state
829 // since it it the read that failed. Just propagate the error.
830 return js.rejectedPromise<kj::Own<ReadContext>>(kj::mv(exception));
831 }));
832}
833 
834// We separate out the actual read implementation so that it can be used by
835// both read and the pumpToImpl implementation.
836kj::Promise<size_t> ReadableSourceKjAdapter::readImpl(
837 Active& active, kj::ArrayPtr<kj::byte> dest, size_t minBytes) {
838 
839 KJ_IF_SOME(readable, active.state.tryGetUnsafe<Active::Readable>()) {
840 // We have some data left over from a previous read. Use that first.
841 
842 // If we have enough left over to fully satisfy this read,
843 // Use it, then update our left over view.
844 if (readable.view.size() >= dest.size()) {
845 dest.copyFrom(readable.view.first(dest.size()));
846 readable.view = readable.view.slice(dest.size());
847 if (readable.view.size() == 0) {
848 // We used up all our left over data. We can transition to the idle state.
849 active.state.transitionTo<Active::Idle>();
850 }
851 // Otherwise we still have some left over data. That
852 // is ok, we will keep it around for the next read.
853 // We intentionally do not transition to the idle state
854 // here because we want to keep the left over data for
855 // the next read.
856 return dest.size();
857 }
858 
859 // Otherwise, consume what we do have left over.
860 auto size = readable.view.size();
861 dest.first(size).copyFrom(readable.view);
862 dest = dest.slice(size);
863 
864 active.state.transitionTo<Active::Idle>();
865 
866 // Did we at least satisfy the minimum bytes?
867 if (size >= minBytes) {
868 // Awesome, we are technically done with this read.
869 // While we might actually have more room in our buffer, and the
870 // minReadyPolicy might be OPPORTUNISTIC, we will not try to
871 // read more from the stream right now so that we can avoid having
872 // to grab the isolate lock for this read. Instead, let's return
873 // what we have and let the caller call read again if/when they want.
874 // This risks leaving a fair amount of unused space in the buffer
875 // and requiring more read calls but it avoids the overhead of
876 // an additional isolate lock grab when we know we can at least
877 // provide some data right now.
878 return size;
879 }
880 }
881 
882 // If we got here, we still have not satisfied the minimum bytes,
883 // so we will continue on to read more from the stream. But, we
884 // also should not have any more data left over. Let's verify.
885 KJ_ASSERT(active.state.is<Active::Idle>());
886 active.state.transitionTo<Active::Reading>();
887 
888 // Our read context holds all the state needed to perform the read.
889 // Ownership of the context is passed into the read operation and
890 // returned back to us when the read is complete.
891 auto context = kj::heap<ReadContext>({
892 .stream = active.stream.addRef(),
893 .reader = active.reader.addRef(),
894 .buffer = dest,
895 .totalRead = 0,
896 .minBytes = minBytes,
897 .adapterRef = selfRef.addRef(),
898 });
899 
900 return active.canceler
901 .wrap(
902 // Warning: Do *not* capture "active" in this lambda! It may be destroyed
903 // while we are in the promise chain. Instead, we capture a weak
904 // reference to the adapter itself and check that we are still alive
905 // and active before trying to update any state.
906 active.ioContext.run([context = kj::mv(context), self = selfRef.addRef(),
907 minReadPolicy = options.minReadPolicy](
908 jsg::Lock& js) mutable -> kj::Promise<size_t> {
909 auto& ioContext = IoContext::current();
910 
911 // Perform the actual read.
912 return ioContext.awaitJs(js, readInternal(js, kj::mv(context), minReadPolicy))
913 .then([self = kj::mv(self)](kj::Own<ReadContext> context) mutable -> kj::Promise<size_t> {
914 // By the time we get here, it is possible that the adapter has been
915 // destroyed. If that's the case, it's okay, that's what our weak ref
916 // is here for. We will only try to update our state if we are still
917 // alive and active.
918 
919 self->runIfAlive([&](ReadableSourceKjAdapter& self) {
920 // Ok, we're still alive! Yay! But, let's check to make sure we didn't
921 // change state while we were reading.
922 KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) {
923 auto& active = *open.active;
924 // Ok, we're still active. Let's see if we have any left over data
925 // that we need to stash away for the next read.
926 KJ_IF_SOME(leftOver, context->maybeLeftOver) {
927 // We have some left over data. Stash it away for the next read.
928 active.state.transitionTo<Active::Readable>(kj::mv(leftOver));
929 // In this branch, we must have filled the entire destination
930 // buffer and satisfied the minimum read requirement or else
931 // we wouldn't have any left over data. Let's just assert that
932 // invariant just in case.
933 KJ_DASSERT(context->totalRead >= context->minBytes);
934 } else if (context->totalRead < context->minBytes) {
935 // We returned fewer than the minimum bytes requested. This is our
936 // signal that we're done.
937 active.state.transitionTo<Active::Done>();
938 // We cannot change the state to Closed here because we are still
939 // inside the kj::Promise chain wrapped by the canceler. If we
940 // change the state to Closed, the Active would be destroyed, causing
941 // this promise chain to be canceled.
942 auto droppedReader KJ_UNUSED = kj::mv(active.reader);
943 auto droppedStream KJ_UNUSED = kj::mv(active.stream);
944 // In this branch, we should not have any left over data.
945 // Let's assert that invariant just in case.
946 KJ_DASSERT(context->maybeLeftOver == kj::none);
947 } else {
948 // Our read is complete. Return to the idle state and we're done.
949 active.state.transitionTo<Active::Idle>();
950 
951 // In this branch, we must have satisfied the minimum read
952 // requirement. Let's just assert that invariant just in case.
953 KJ_DASSERT(context->totalRead >= context->minBytes);
954 // We should not have any left over data.
955 KJ_DASSERT(context->maybeLeftOver == kj::none);
956 }
957 } else {
958 // We were closed or canceled while we were reading. Doh!
959 // That's ok, there's nothing more we can or need to do
960 // here. Just fall-through to the return below.
961 }
962 });
963 return context->totalRead;
964 });
965 })).catch_([self = selfRef.addRef()](kj::Exception exception) -> kj::Promise<size_t> {
966 self->runIfAlive([&](ReadableSourceKjAdapter& self) {
967 KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) {
968 open.active->state.forceTransitionTo<Active::Canceling>(Active::Canceling{
969 .exception = exception.clone(),
970 });
971 }
972 });
973 return kj::mv(exception);
974 });
975}
976 
977kj::Promise<size_t> ReadableSourceKjAdapter::read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes) {
978 
979 if (buffer.size() == 0) {
980 // Nothing to read. This is a no-op.
981 return static_cast<size_t>(0);
982 }
983 
984 // Clamp the minBytes to [1, buffer.size()].
985 minBytes = kj::min(buffer.size(), kj::max(minBytes, 1UL));
986 KJ_DASSERT(minBytes >= 1 && minBytes <= buffer.size(),
987 "minBytes must be less than or equal to the buffer size.");
988 
989 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
990 return exception.clone();
991 }
992 
993 if (state.is<KjClosed>()) {
994 return static_cast<size_t>(0);
995 }
996 
997 auto& open = state.requireActiveUnsafe();
998 auto& active = *open.active;
999 KJ_SWITCH_ONEOF(active.state) {
1000 KJ_CASE_ONEOF(_, Active::Reading) {
1001 KJ_FAIL_REQUIRE("Cannot have multiple concurrent reads.");
1002 }
1003 KJ_CASE_ONEOF(_, Active::Done) {
1004 // The previous read indicated that it was the last read by returning
1005 // less than the minimum bytes requested. We have to treat this as
1006 // the stream being closed.
1007 state.transitionTo<KjClosed>();
1008 return static_cast<size_t>(0);
1009 }
1010 KJ_CASE_ONEOF(canceling, Active::Canceling) {
1011 // The stream is being canceled. Propagate the exception and complete
1012 // the state transition.
1013 return KJ_ASSERT_NONNULL(checkCancelingOrCanceled(active));
1014 }
1015 KJ_CASE_ONEOF(canceled, Active::Canceled) {
1016 // The stream was canceled. Propagate the exception and complete
1017 // the state transition.
1018 return KJ_ASSERT_NONNULL(checkCancelingOrCanceled(active));
1019 }
1020 KJ_CASE_ONEOF(r, Active::Readable) {
1021 // There is some data left over from a previous read.
1022 return readImpl(active, buffer, minBytes);
1023 }
1024 KJ_CASE_ONEOF(_, Active::Idle) {
1025 // There are no pending reads and no left over data.
1026 return readImpl(active, buffer, minBytes);
1027 }
1028 }
1029 KJ_UNREACHABLE;
1030}
1031 
1032kj::Maybe<size_t> ReadableSourceKjAdapter::tryGetLength(StreamEncoding encoding) {
1033 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
1034 auto& active = *open.active;
1035 if (active.state.is<Active::Done>() || active.state.is<Active::Canceled>()) {
1036 // If the previous read indicated that it was the last, then
1037 // let's just transition to the closed state now and return kj::none.
1038 state.transitionTo<KjClosed>();
1039 return kj::none;
1040 }
1041 if (checkCancelingOrCanceled(active) != kj::none) {
1042 return kj::none;
1043 }
1044 return active.stream->tryGetLength(encoding).map(
1045 [](uint64_t len) { return static_cast<size_t>(len); });
1046 }
1047 
1048 // The stream is either closed or errored.
1049 return kj::none;
1050}
1051 
1052void ReadableSourceKjAdapter::cancel(kj::Exception reason) {
1053 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
1054 open.active->cancel(reason.clone());
1055 }
1056 state.forceTransitionTo<kj::Exception>(kj::mv(reason));
1057}
1058 
1059kj::Maybe<kj::Exception> ReadableSourceKjAdapter::checkCancelingOrCanceled(Active& active) {
1060 KJ_IF_SOME(canceling, active.state.tryGetUnsafe<Active::Canceling>()) {
1061 auto exception = kj::mv(canceling.exception);
1062 state.forceTransitionTo<kj::Exception>(exception.clone());
1063 return kj::mv(exception);
1064 }
1065 KJ_IF_SOME(canceled, active.state.tryGetUnsafe<Active::Canceled>()) {
1066 auto exception = kj::mv(canceled.exception);
1067 state.forceTransitionTo<kj::Exception>(exception.clone());
1068 return kj::mv(exception);
1069 }
1070 return kj::none;
1071}
1072 
1073void ReadableSourceKjAdapter::throwIfCancelingOrCanceled(Active& active) {
1074 KJ_IF_SOME(exception, checkCancelingOrCanceled(active)) {
1075 kj::throwFatalException(kj::mv(exception));
1076 }
1077}
1078 
1079kj::Promise<void> ReadableSourceKjAdapter::pumpToImpl(
1080 kj::Own<Active> active, WritableSink& output, EndAfterPump end) {
1081 // This implementation uses DrainingReader to efficiently pull all synchronously
1082 // available data from the underlying JS stream in each iteration. This minimizes
1083 // the number of isolate lock acquisitions by getting all available data at once
1084 // rather than reading into fixed-size buffers.
1085 
1086 KJ_DASSERT(active->state.is<Active::Idle>() || active->state.is<Active::Readable>(),
1087 "pumpToImpl called when stream is not in an active state.");
1088 
1089 bool writeFailed = false;
1090 
1091 // First, if the active state is in the Readable state, we need to drain the
1092 // left over data before starting the main read loop.
1093 // This is unlikely to occur in the typical case, but we need to handle it
1094 // nonetheless.
1095 KJ_IF_SOME(readable, active->state.tryGetUnsafe<Active::Readable>()) {
1096 co_await output.write(readable.view);
1097 active->state.transitionTo<Active::Idle>();
1098 }
1099 
1100 // We hold the DrainingReader during the pump. The pointer remains valid because
1101 // the reader is created and owned during the pump loop lifetime.
1102 kj::Maybe<kj::Own<DrainingReader>> maybeReader;
1103 
1104 // Initialize the pump by releasing the default reader and creating a DrainingReader.
1105 // This requires the isolate lock.
1106 co_await active->ioContext.run(
1107 [&active, &maybeReader](jsg::Lock& js) mutable -> kj::Promise<void> {
1108 // Release the existing reader's lock so we can create a DrainingReader.
1109 active->reader->releaseLock(js);
1110 
1111 // Create the DrainingReader for the stream.
1112 maybeReader = KJ_ASSERT_NONNULL(DrainingReader::create(js, *active->stream),
1113 "Failed to create DrainingReader - stream should not be locked");
1114 return kj::READY_NOW;
1115 });
1116 
1117 auto& reader = KJ_ASSERT_NONNULL(maybeReader);
1118 kj::Maybe<kj::Exception> pendingException;
1119 
1120 try {
1121 while (true) {
1122 // Perform a draining read to get all synchronously available data.
1123 // Pass raw pointer to reader into the lambda - safe because we own it
1124 // and keep it alive for the duration of the pump.
1125 // The draining reader grabs all data currently available in the stream's
1126 // queue, then tries to read more data up to a limit as long as the data
1127 // can be provided synchronously. The idea is to drain off as much data
1128 // from the stream as possible each time we are holding the isolate lock
1129 // to minimize the number of times we need to re-enter the lock.
1130 DrainingReader* readerPtr = reader.get();
1131 DrainingReadResult result =
1132 co_await active->ioContext.run([readerPtr](jsg::Lock& js) mutable {
1133 auto& ioContext = IoContext::current();
1134 // Use a 256KB limit to allow periodic yielding to the event loop,
1135 // preventing a fast producer from monopolizing the thread. This limit
1136 // only affects subsequent pump iterations after the initial buffer drain.
1137 constexpr size_t kMaxReadPerCycle = 256 * 1024;
1138 return ioContext.awaitJs(js, readerPtr->read(js, kMaxReadPerCycle));
1139 });
1140 
1141 // Write all the chunks we received using vectored write for efficiency.
1142 if (result.chunks.size() > 0) {
1143 KJ_ON_SCOPE_FAILURE(writeFailed = true);
1144 // Convert Array<Array<byte>> to ArrayPtr<ArrayPtr<const byte>> for vectored write.
1145 auto pieces =
1146 KJ_MAP(chunk, result.chunks) -> kj::ArrayPtr<const kj::byte> { return chunk.asPtr(); };
1147 co_await output.write(pieces);
1148 }
1149 
1150 // If the stream is done, end the output if needed and exit.
1151 if (result.done) {
1152 KJ_ON_SCOPE_FAILURE(writeFailed = true);
1153 if (end) {
1154 co_await output.end();
1155 }
1156 co_return;
1157 }
1158 }
1159 } catch (...) {
1160 auto exception = kj::getCaughtExceptionAsKj();
1161 if (!writeFailed) {
1162 // If we got an error and it wasn't the write that failed, abort the output.
1163 output.abort(exception.clone());
1164 }
1165 // Store the exception to handle after the catch block.
1166 pendingException = kj::mv(exception);
1167 }
1168 
1169 // If there was an error, cancel the reader and propagate the exception.
1170 KJ_IF_SOME(exception, pendingException) {
1171 DrainingReader* readerPtr = reader.get();
1172 co_await active->ioContext.run([readerPtr, ex = exception.clone()](jsg::Lock& js) mutable {
1173 auto& ioContext = IoContext::current();
1174 auto error = js.exceptionToJsValue(kj::mv(ex));
1175 return ioContext.awaitJs(js, readerPtr->cancel(js, error.getHandle(js)));
1176 });
1177 kj::throwFatalException(kj::mv(exception));
1178 }
1179}
1180 
1181kj::Promise<DeferredProxy<void>> ReadableSourceKjAdapter::pumpTo(
1182 WritableSink& output, EndAfterPump end) {
1183 // The pumpTo operation continually reads from the stream and writes
1184 // to the output until the stream is closed or an error occurs. Once
1185 // the pump starts, the adapter transitions to the closed state and
1186 // ownership of the underlying stream is transferred to the pump
1187 // operation.
1188 
1189 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
1190 return kj::Promise<DeferredProxy<void>>(DeferredProxy<void>{exception.clone()});
1191 }
1192 
1193 if (state.is<KjClosed>()) {
1194 // Already closed, nothing to do.
1195 return newNoopDeferredProxy();
1196 }
1197 
1198 auto& open = state.requireActiveUnsafe();
1199 auto& active = *open.active;
1200 // Per the contract for ReadableStreamSource::pumpTo, the pump operation
1201 // will take over ownership of the underlying stream until it is complete,
1202 // leaving the adapter itself in a closed state once the pump starts.
1203 // Dropping the returned promise will cancel the pump operation.
1204 // We do, however, need to first make sure that our active state is
1205 // not already pending a read or terminal state change.
1206 KJ_REQUIRE(!active.state.is<Active::Reading>(), "Cannot have multiple concurrent reads.");
1207 
1208 if (active.state.is<Active::Done>()) {
1209 // The previous read indicated that it was the last read by returning
1210 // less than the minimum bytes requested, or the stream was fully
1211 // canceled. We have to treat this as the stream being closed.
1212 state.transitionTo<KjClosed>();
1213 return newNoopDeferredProxy();
1214 }
1215 
1216 KJ_IF_SOME(exception, checkCancelingOrCanceled(active)) {
1217 return kj::Promise<DeferredProxy<void>>(kj::mv(exception));
1218 }
1219 
1220 // The active state should be Readable of Idle here. Let's verify.
1221 KJ_DASSERT(active.state.is<Active::Readable>() || active.state.is<Active::Idle>());
1222 
1223 // The Active state will be transferred into the pumpImpl operation.
1224 auto activeState = kj::mv(open.active);
1225 state.transitionTo<KjClosed>(); // transition to closed immediately
1226 
1227 // Because pumpToImpl is wrapping a JavaScript stream, it is not eligible
1228 // for deferred proxying. We will return a noopDeferredProxy that wraps the
1229 // promise from pumpToImpl();
1230 return addNoopDeferredProxy(pumpToImpl(kj::mv(activeState), output, end));
1231}
1232 
1233ReadableSource::Tee ReadableSourceKjAdapter::tee(size_t) {
1234 KJ_UNIMPLEMENTED("Teeing a ReadableSourceKjAdapter is not supported.");
1235 // Explanation: Teeing a ReadableStream must be done under the isolate lock,
1236 // as does creating a new ReadableSourceKjAdapter. However, when tee()
1237 // is called we are not guaranteed to be under the isolate lock, nor can
1238 // we acquire the lock here because this is a synchronous operation and
1239 // acquiring the isolate lock requires waiting for a promise to resolve.
1240 //
1241 // Teeing here is unlikely to be necessary. If you do need a tee, it's
1242 // necessary to tee the underlying ReadableStream directly and create
1243 // two separate ReadableSourceKjAdapters, one for each branch of
1244 // that tee while the lock is held.
1245}
1246 
1247kj::Promise<kj::Array<const kj::byte>> ReadableSourceKjAdapter::readAllBytes(size_t limit) {
1248 co_return co_await readAllImpl<kj::byte>(limit);
1249}
1250 
1251kj::Promise<kj::String> ReadableSourceKjAdapter::readAllText(size_t limit) {
1252 auto array = co_await readAllImpl<char>(limit);
1253 co_return kj::String(kj::mv(array));
1254}
1255 
1256template <typename T>
1257kj::Promise<kj::Array<T>> ReadableSourceKjAdapter::readAllImpl(size_t limit) {
1258 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
1259 kj::throwFatalException(exception.clone());
1260 }
1261 
1262 if (state.is<KjClosed>()) {
1263 co_return kj::Array<T>();
1264 }
1265 
1266 auto& open = state.requireActiveUnsafe();
1267 auto& active = *open.active;
1268 KJ_REQUIRE(!active.state.is<Active::Reading>(), "Cannot have multiple concurrent reads.");
1269 
1270 if (active.state.is<Active::Done>()) {
1271 // The previous read indicated that it was the last read by returning
1272 // less than the minimum bytes requested. We have to treat this as
1273 // the stream being closed.
1274 state.transitionTo<KjClosed>();
1275 co_return kj::Array<T>();
1276 }
1277 
1278 throwIfCancelingOrCanceled(active);
1279 
1280 // Our readAll operation will accumulate data into a buffer up to the
1281 // specified limit. If the limit is exceeded, the returned promise will
1282 // be rejected. Once the readAll operation starts, the adapter is moved
1283 // into a closed state and ownership of the underlying stream is transferred
1284 // to the readAll promise.
1285 auto activeState = kj::mv(open.active);
1286 state.transitionTo<KjClosed>(); // transition to closed immediately
1287 
1288 KJ_DASSERT(activeState->state.is<Active::Readable>() || activeState->state.is<Active::Idle>());
1289 
1290 // We do not use the canceler here. The adapter is closed and can be safely dropped.
1291 // This promise, however, will keep the stream alive until the read is completed.
1292 // If the returned promise is dropped, the readAll operation will be canceled.
1293 CancelationToken cancelationToken;
1294 co_return co_await IoContext::current().run(
1295 [limit, active = kj::mv(activeState), cancelationToken = cancelationToken.getWeakRef()](
1296 jsg::Lock& js) mutable -> kj::Promise<kj::Array<T>> {
1297 kj::Vector<T> accumulated;
1298 // If we know the length of the stream ahead of time, and it is within the limit,
1299 // we can reserve that much space in the accumulator to avoid multiple allocations.
1300 KJ_IF_SOME(length, active->stream->tryGetLength(StreamEncoding::IDENTITY)) {
1301 if (length <= limit) {
1302 accumulated.reserve(length); // Pre-allocate
1303 }
1304 }
1305 
1306 auto& ioContext = IoContext::current();
1307 return ioContext.awaitJs(js,
1308 readAllReadImpl(js, ioContext.addObject(kj::mv(active)), kj::mv(accumulated), limit,
1309 kj::mv(cancelationToken)));
1310 });
1311}
1312 
1313template <typename T>
1314jsg::Promise<kj::Array<T>> ReadableSourceKjAdapter::readAllReadImpl(jsg::Lock& js,
1315 IoOwn<Active> active,
1316 kj::Vector<T> accumulated,
1317 size_t limit,
1318 kj::Rc<WeakRef<CancelationToken>> cancelationToken) {
1319 
1320 // Check for cancelation. The cancelation token is a weak ref. If the promise
1321 // that represents the readAll operation is dropped, the token will be invalidated.
1322 // Since there is no way to directly cancel a JavaScript promise, this is the best
1323 // we can do to interrupt the loop.
1324 if (!cancelationToken->isValid()) {
1325 return js.rejectedPromise<kj::Array<T>>(js.error("readAll operation was canceled."));
1326 }
1327 
1328 // First, drain any leftover data if the active state is in Readable mode.
1329 KJ_IF_SOME(readable, active->state.tryGetUnsafe<Active::Readable>()) {
1330 auto leftover = readable.view.asBytes();
1331 if (leftover.size() > limit) {
1332 auto error = js.rangeError("Memory limit would be exceeded before EOF.");
1333 return active->reader->cancel(js, error).then(
1334 js, [ex = jsg::JsRef(js, error)](jsg::Lock& js) {
1335 return js.rejectedPromise<kj::Array<T>>(ex.getHandle(js));
1336 });
1337 }
1338 if constexpr (kj::isSameType<T, char>()) {
1339 accumulated.addAll(leftover.asChars());
1340 } else {
1341 accumulated.addAll(leftover);
1342 }
1343 active->state.transitionTo<Active::Idle>();
1344 }
1345 
1346 return active->reader->read(js).then(js,
1347 [active = kj::mv(active), accumulated = kj::mv(accumulated), limit,
1348 cancelationToken = kj::mv(cancelationToken)](
1349 jsg::Lock& js, ReadResult result) mutable -> jsg::Promise<kj::Array<T>> {
1350 // Check for cancelation.
1351 if (!cancelationToken->isValid()) {
1352 return js.rejectedPromise<kj::Array<T>>(js.error("readAll operation was canceled."));
1353 }
1354 
1355 if (result.done || result.value == kj::none) {
1356 // Stream ended. Return accumulated data.
1357 // If we're reading text, add NUL terminator.
1358 if constexpr (kj::isSameType<T, char>()) {
1359 accumulated.add('\0');
1360 }
1361 return js.resolvedPromise(accumulated.releaseAsArray());
1362 }
1363 
1364 auto& value = KJ_ASSERT_NONNULL(result.value);
1365 auto jsval = jsg::JsValue(value.getHandle(js));
1366 
1367 kj::ArrayPtr<const kj::byte> bytes;
1368 kj::Maybe<kj::String> maybeOwnedString;
1369 
1370 KJ_IF_SOME(str, jsval.tryCast<jsg::JsString>()) {
1371 auto data = str.toUSVString(js);
1372 bytes = data.asBytes();
1373 maybeOwnedString = kj::mv(data);
1374 } else KJ_IF_SOME(ab, jsval.tryCast<jsg::JsArrayBuffer>()) {
1375 bytes = ab.asArrayPtr();
1376 } else KJ_IF_SOME(view, jsval.tryCast<jsg::JsArrayBufferView>()) {
1377 bytes = view.asArrayPtr();
1378 } else {
1379 auto error = js.typeError("ReadableStream provided a non-bytes value. Only ArrayBuffer, "
1380 "ArrayBufferView, or string are supported.");
1381 return active->reader->cancel(js, error).then(
1382 js, [err = jsg::JsRef(js, error)](jsg::Lock& js) {
1383 return js.rejectedPromise<kj::Array<T>>(err.getHandle(js));
1384 });
1385 }
1386 
1387 if (accumulated.size() + bytes.size() > limit) {
1388 auto error = js.rangeError("Memory limit would be exceeded before EOF.");
1389 return active->reader->cancel(js, error).then(
1390 js, [err = jsg::JsRef(js, error)](jsg::Lock& js) {
1391 return js.rejectedPromise<kj::Array<T>>(err.getHandle(js));
1392 });
1393 }
1394 
1395 // Accumulate the bytes.
1396 if constexpr (kj::isSameType<T, char>()) {
1397 accumulated.addAll(bytes.asChars());
1398 } else {
1399 accumulated.addAll(bytes);
1400 }
1401 
1402 // Continue reading.
1403 return readAllReadImpl(
1404 js, kj::mv(active), kj::mv(accumulated), limit, kj::mv(cancelationToken));
1405 });
1406}
1407 
1408} // namespace workerd::api::streams