File
Blob: src/workerd/api/streams/readable-source-adapter.c++
| 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 | |
| 9 | namespace workerd::api::streams { |
| 10 | |
| 11 | namespace { |
| 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. |
| 18 | jsg::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. |
| 36 | struct 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 | |
| 126 | ReadableStreamSourceJsAdapter::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 | |
| 132 | ReadableStreamSourceJsAdapter::~ReadableStreamSourceJsAdapter() noexcept(false) { |
| 133 | selfRef->invalidate(); |
| 134 | } |
| 135 | |
| 136 | void 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 | |
| 143 | void ReadableStreamSourceJsAdapter::cancel(jsg::Lock& js, const jsg::JsValue& reason) { |
| 144 | cancel(js.exceptionToKj(reason)); |
| 145 | } |
| 146 | |
| 147 | void 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 | |
| 155 | bool ReadableStreamSourceJsAdapter::isClosed() { |
| 156 | return state.is<Closed>(); |
| 157 | } |
| 158 | |
| 159 | kj::Maybe<const kj::Exception&> ReadableStreamSourceJsAdapter::isCanceled() { |
| 160 | return state.tryGetErrorUnsafe(); |
| 161 | } |
| 162 | |
| 163 | jsg::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. |
| 282 | jsg::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 | |
| 320 | jsg::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 | |
| 380 | jsg::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 | |
| 445 | kj::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 | |
| 452 | kj::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 | |
| 479 | struct 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. |
| 551 | struct 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 | |
| 576 | namespace { |
| 577 | constexpr size_t kMinRemainingForAdditionalRead = 512; |
| 578 | |
| 579 | jsg::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 | |
| 586 | using JsByteSource = kj::OneOf<jsg::JsRef<jsg::JsString>, |
| 587 | jsg::JsRef<jsg::JsArrayBuffer>, |
| 588 | jsg::JsRef<jsg::JsArrayBufferView>>; |
| 589 | |
| 590 | kj::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. |
| 603 | kj::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 | |
| 685 | ReadableSourceKjAdapter::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 | |
| 692 | ReadableSourceKjAdapter::Active::~Active() noexcept(false) { |
| 693 | cancel(KJ_EXCEPTION(DISCONNECTED, "ReadableSourceKjAdapter is canceled.")); |
| 694 | } |
| 695 | |
| 696 | void 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 | |
| 717 | ReadableSourceKjAdapter::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 | |
| 724 | ReadableSourceKjAdapter::~ReadableSourceKjAdapter() noexcept(false) { |
| 725 | selfRef->invalidate(); |
| 726 | } |
| 727 | |
| 728 | jsg::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. |
| 836 | kj::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 | |
| 977 | kj::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 | |
| 1032 | kj::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 | |
| 1052 | void 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 | |
| 1059 | kj::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 | |
| 1073 | void ReadableSourceKjAdapter::throwIfCancelingOrCanceled(Active& active) { |
| 1074 | KJ_IF_SOME(exception, checkCancelingOrCanceled(active)) { |
| 1075 | kj::throwFatalException(kj::mv(exception)); |
| 1076 | } |
| 1077 | } |
| 1078 | |
| 1079 | kj::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 | |
| 1181 | kj::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 | |
| 1233 | ReadableSource::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 | |
| 1247 | kj::Promise<kj::Array<const kj::byte>> ReadableSourceKjAdapter::readAllBytes(size_t limit) { |
| 1248 | co_return co_await readAllImpl<kj::byte>(limit); |
| 1249 | } |
| 1250 | |
| 1251 | kj::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 | |
| 1256 | template <typename T> |
| 1257 | kj::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 | |
| 1313 | template <typename T> |
| 1314 | jsg::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 |