File
Blob: src/workerd/api/streams/queue.c++
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #include "queue.h" |
| 6 | |
| 7 | #include <workerd/io/features.h> |
| 8 | #include <workerd/jsg/jsg.h> |
| 9 | |
| 10 | #include <kj/common.h> |
| 11 | |
| 12 | #include <algorithm> |
| 13 | |
| 14 | namespace workerd::api { |
| 15 | |
| 16 | // ====================================================================================== |
| 17 | // ValueQueue |
| 18 | #pragma region ValueQueue |
| 19 | |
| 20 | #pragma region ValueQueue::ReadRequest |
| 21 | |
| 22 | void ValueQueue::ReadRequest::resolveAsDone(jsg::Lock& js) { |
| 23 | resolver.resolve(js, ReadResult{.done = true}); |
| 24 | } |
| 25 | |
| 26 | void ValueQueue::ReadRequest::resolve(jsg::Lock& js, jsg::Value value) { |
| 27 | resolver.resolve(js, ReadResult{.value = kj::mv(value), .done = false}); |
| 28 | } |
| 29 | |
| 30 | void ValueQueue::ReadRequest::reject(jsg::Lock& js, jsg::Value& value) { |
| 31 | resolver.reject(js, value.getHandle(js)); |
| 32 | } |
| 33 | |
| 34 | #pragma endregion ValueQueue::ReadRequest |
| 35 | |
| 36 | #pragma region ValueQueue::Entry |
| 37 | |
| 38 | ValueQueue::Entry::Entry(jsg::Value value, size_t size): value(kj::mv(value)), size(size) {} |
| 39 | |
| 40 | jsg::Value ValueQueue::Entry::getValue(jsg::Lock& js) { |
| 41 | return value.addRef(js); |
| 42 | } |
| 43 | |
| 44 | size_t ValueQueue::Entry::getSize() const { |
| 45 | return size; |
| 46 | } |
| 47 | |
| 48 | void ValueQueue::Entry::visitForGc(jsg::GcVisitor& visitor) { |
| 49 | visitor.visit(value); |
| 50 | } |
| 51 | |
| 52 | #pragma endregion ValueQueue::Entry |
| 53 | |
| 54 | #pragma region ValueQueue::QueueEntry |
| 55 | |
| 56 | kj::Rc<ValueQueue::Entry> ValueQueue::Entry::clone(jsg::Lock& js) { |
| 57 | return addRefToThis(); |
| 58 | } |
| 59 | |
| 60 | ValueQueue::QueueEntry ValueQueue::QueueEntry::clone(jsg::Lock& js) { |
| 61 | return QueueEntry{.entry = entry->clone(js)}; |
| 62 | } |
| 63 | |
| 64 | #pragma endregion ValueQueue::QueueEntry |
| 65 | |
| 66 | #pragma region ValueQueue::Consumer |
| 67 | |
| 68 | ValueQueue::Consumer::Consumer( |
| 69 | ValueQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 70 | : impl(queue.impl, stateListener) {} |
| 71 | |
| 72 | ValueQueue::Consumer::Consumer( |
| 73 | QueueImpl& impl, kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 74 | : impl(impl, stateListener) {} |
| 75 | |
| 76 | ValueQueue::Consumer::Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 77 | : impl(stateListener) {} |
| 78 | |
| 79 | void ValueQueue::Consumer::cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 80 | impl.cancel(js, maybeReason); |
| 81 | } |
| 82 | |
| 83 | void ValueQueue::Consumer::close(jsg::Lock& js) { |
| 84 | impl.close(js); |
| 85 | }; |
| 86 | |
| 87 | bool ValueQueue::Consumer::empty() { |
| 88 | return impl.empty(); |
| 89 | } |
| 90 | |
| 91 | void ValueQueue::Consumer::error(jsg::Lock& js, jsg::Value reason) { |
| 92 | impl.error(js, kj::mv(reason)); |
| 93 | }; |
| 94 | |
| 95 | void ValueQueue::Consumer::read(jsg::Lock& js, ReadRequest request) { |
| 96 | impl.read(js, kj::mv(request)); |
| 97 | } |
| 98 | |
| 99 | void ValueQueue::Consumer::push(jsg::Lock& js, kj::Rc<Entry> entry) { |
| 100 | impl.push(js, kj::mv(entry)); |
| 101 | } |
| 102 | |
| 103 | void ValueQueue::Consumer::reset() { |
| 104 | impl.reset(); |
| 105 | }; |
| 106 | |
| 107 | size_t ValueQueue::Consumer::size() { |
| 108 | return impl.size(); |
| 109 | } |
| 110 | |
| 111 | kj::Own<ValueQueue::Consumer> ValueQueue::Consumer::clone( |
| 112 | jsg::Lock& js, kj::Maybe<ConsumerImpl::StateListener&> stateListener) { |
| 113 | // If the queue was destroyed (e.g., stream was closed), we can still clone |
| 114 | // the consumer - the cloneTo() will copy the closed/errored state. |
| 115 | kj::Own<Consumer> consumer; |
| 116 | KJ_IF_SOME(q, impl.queue) { |
| 117 | consumer = kj::heap<Consumer>(q, stateListener); |
| 118 | } else { |
| 119 | consumer = kj::heap<Consumer>(stateListener); |
| 120 | } |
| 121 | impl.cloneTo(js, consumer->impl); |
| 122 | return kj::mv(consumer); |
| 123 | } |
| 124 | |
| 125 | bool ValueQueue::Consumer::hasReadRequests() { |
| 126 | return impl.hasReadRequests(); |
| 127 | } |
| 128 | |
| 129 | bool ValueQueue::Consumer::hasPendingDrainingRead() { |
| 130 | return impl.state.whenActiveOr( |
| 131 | [](const ConsumerImpl::Ready& ready) { return ready.hasPendingDrainingRead; }, false); |
| 132 | } |
| 133 | |
| 134 | namespace { |
| 135 | // Helper to convert a JS value to bytes. Returns kj::none if the value cannot be converted. |
| 136 | kj::Maybe<kj::Array<kj::byte>> valueToBytes(jsg::Lock& js, jsg::Value& value) { |
| 137 | auto jsval = jsg::JsValue(value.getHandle(js)); |
| 138 | |
| 139 | // Try ArrayBuffer first. |
| 140 | KJ_IF_SOME(ab, jsval.tryCast<jsg::JsArrayBuffer>()) { |
| 141 | auto src = ab.asArrayPtr(); |
| 142 | return kj::heapArray(src); |
| 143 | } |
| 144 | |
| 145 | // Try ArrayBufferView. |
| 146 | KJ_IF_SOME(abView, jsval.tryCast<jsg::JsArrayBufferView>()) { |
| 147 | auto src = abView.asArrayPtr(); |
| 148 | return kj::heapArray(src); |
| 149 | } |
| 150 | |
| 151 | // Try string - convert to UTF-8. |
| 152 | KJ_IF_SOME(str, jsval.tryCast<jsg::JsString>()) { |
| 153 | auto data = str.toUSVString(js); |
| 154 | return kj::heapArray(data.asBytes()); |
| 155 | } |
| 156 | |
| 157 | // Unsupported type. |
| 158 | return kj::none; |
| 159 | } |
| 160 | } // namespace |
| 161 | |
| 162 | jsg::Promise<DrainingReadResult> ValueQueue::Consumer::drainingRead(jsg::Lock& js, size_t maxRead) { |
| 163 | // If there are pending regular reads, reject - mutual exclusion. |
| 164 | if (hasReadRequests()) { |
| 165 | return js.rejectedPromise<DrainingReadResult>( |
| 166 | js.typeError("Cannot call drainingRead while there are pending reads"_kj)); |
| 167 | } |
| 168 | |
| 169 | // Check if already closed or errored. |
| 170 | if (impl.state.template is<ConsumerImpl::Closed>()) { |
| 171 | return js.resolvedPromise(DrainingReadResult{.chunks = nullptr, .done = true}); |
| 172 | } |
| 173 | KJ_IF_SOME(errored, impl.state.tryGetErrorUnsafe()) { |
| 174 | return js.rejectedPromise<DrainingReadResult>(errored.reason.getHandle(js)); |
| 175 | } |
| 176 | |
| 177 | auto& ready = impl.state.requireActiveUnsafe(); |
| 178 | ConsumerImpl::UpdateBackpressureScope scope(impl); |
| 179 | |
| 180 | // Mark that we're doing a draining read. This allows onConsumerWantsData() |
| 181 | // to use forcePull() which bypasses backpressure checks. The flag is cleared |
| 182 | // either synchronously (for immediate returns) or in promise callbacks (for async). |
| 183 | ready.hasPendingDrainingRead = true; |
| 184 | |
| 185 | // Collect all buffered data, converting values to bytes. |
| 186 | kj::Vector<kj::Array<kj::byte>> chunks; |
| 187 | bool isClosing = false; |
| 188 | size_t totalRead = 0; |
| 189 | |
| 190 | // Drains buffered data, converting values to bytes. Returns a rejected promise if a value |
| 191 | // cannot be converted; otherwise returns kj::none to indicate success. Stops draining |
| 192 | // when totalRead reaches or exceeds maxRead (after finishing the current item). |
| 193 | static const auto drainBuffer = |
| 194 | [](jsg::Lock& js, ConsumerImpl& impl, ConsumerImpl::Ready& ready, |
| 195 | kj::Vector<kj::Array<kj::byte>>& chunks, size_t& totalRead, bool& isClosing, |
| 196 | size_t maxRead) -> kj::Maybe<jsg::Promise<DrainingReadResult>> { |
| 197 | while (!ready.buffer.empty() && !isClosing && totalRead < maxRead) { |
| 198 | auto& item = ready.buffer.front(); |
| 199 | KJ_SWITCH_ONEOF(item) { |
| 200 | KJ_CASE_ONEOF(close, ConsumerImpl::Close) { |
| 201 | isClosing = true; |
| 202 | break; |
| 203 | } |
| 204 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 205 | auto value = entry.entry->getValue(js); |
| 206 | KJ_IF_SOME(bytes, valueToBytes(js, value)) { |
| 207 | totalRead += bytes.size(); |
| 208 | chunks.add(kj::mv(bytes)); |
| 209 | ready.queueTotalSize -= entry.entry->getSize(); |
| 210 | ready.buffer.pop_front(); |
| 211 | } else { |
| 212 | auto error = js.typeError( |
| 213 | "Draining read encountered a value that cannot be converted to bytes"_kj); |
| 214 | impl.error(js, jsg::Value(js.v8Isolate, error)); |
| 215 | return js.rejectedPromise<DrainingReadResult>(error); |
| 216 | } |
| 217 | } |
| 218 | } |
| 219 | } |
| 220 | return kj::none; |
| 221 | }; |
| 222 | |
| 223 | // Drain the buffer up to maxRead bytes, then pump for more if under the limit. |
| 224 | KJ_IF_SOME(errorPromise, drainBuffer(js, impl, ready, chunks, totalRead, isClosing, maxRead)) { |
| 225 | return kj::mv(errorPromise); |
| 226 | } |
| 227 | |
| 228 | // Pump the controller for more synchronously available data. |
| 229 | // maxRead is checked here: we only proceed with pumping if we haven't exceeded it. |
| 230 | KJ_IF_SOME(listener, impl.stateListener) { |
| 231 | while (!isClosing && totalRead < maxRead) { |
| 232 | size_t prevChunkCount = chunks.size(); |
| 233 | bool pullCompletedSync = listener.onConsumerWantsData(js); |
| 234 | |
| 235 | // The pull callback may have closed or errored the consumer, which |
| 236 | // destroys the Ready state (and its RingBuffer). We must not touch |
| 237 | // `ready` after that. |
| 238 | if (!impl.state.isActive()) break; |
| 239 | |
| 240 | // Drain buffered data that was added by the pull, respecting maxRead. |
| 241 | KJ_IF_SOME(errorPromise, |
| 242 | drainBuffer(js, impl, ready, chunks, totalRead, isClosing, maxRead)) { |
| 243 | return kj::mv(errorPromise); |
| 244 | } |
| 245 | |
| 246 | // If pull is async or no new data was added, stop pumping. |
| 247 | if (!pullCompletedSync || chunks.size() == prevChunkCount) { |
| 248 | break; |
| 249 | } |
| 250 | } |
| 251 | } |
| 252 | |
| 253 | // If the consumer was closed or errored during pumping, the `ready` |
| 254 | // reference is dangling. Return what we have or the appropriate error. |
| 255 | if (!impl.state.isActive()) { |
| 256 | KJ_IF_SOME(errored, impl.state.tryGetErrorUnsafe()) { |
| 257 | return js.rejectedPromise<DrainingReadResult>(errored.reason.getHandle(js)); |
| 258 | } |
| 259 | // Closed — all data was already drained. Return collected chunks. |
| 260 | return js.resolvedPromise(DrainingReadResult{ |
| 261 | .chunks = chunks.releaseAsArray(), |
| 262 | .done = true, |
| 263 | }); |
| 264 | } |
| 265 | |
| 266 | // If the controller was canceled during pumping (e.g., pull callback called |
| 267 | // controller.cancel()), the QueueImpl is destroyed and the consumer's queue |
| 268 | // reference is detached (set to kj::none). The consumer state is still Active |
| 269 | // because cancel on the controller doesn't notify consumers — it only closes |
| 270 | // the controller's own state. No more data will ever arrive. Drain remaining |
| 271 | // buffer data respecting maxRead, and signal done only when the buffer is empty. |
| 272 | if (impl.queue == kj::none) { |
| 273 | // Drain remaining buffer up to maxRead. If there's still more, the caller |
| 274 | // will loop back and we'll drain the rest on subsequent calls. |
| 275 | KJ_IF_SOME(errorPromise, drainBuffer(js, impl, ready, chunks, totalRead, isClosing, maxRead)) { |
| 276 | return kj::mv(errorPromise); |
| 277 | } |
| 278 | ready.hasPendingDrainingRead = false; |
| 279 | bool done = ready.buffer.empty() || isClosing; |
| 280 | // If isClosing, finalize the consumer so onConsumerClose fires promptly. |
| 281 | // maybeDrainAndSetState may transition consumer to Closed, making `ready` dangling. |
| 282 | if (isClosing) { |
| 283 | impl.maybeDrainAndSetState(js); |
| 284 | } |
| 285 | return js.resolvedPromise(DrainingReadResult{ |
| 286 | .chunks = chunks.releaseAsArray(), |
| 287 | .done = done, |
| 288 | }); |
| 289 | } |
| 290 | |
| 291 | // If we collected data, return it immediately. |
| 292 | if (!chunks.empty() || isClosing) { |
| 293 | ready.hasPendingDrainingRead = false; |
| 294 | // If isClosing, finalize the consumer so onConsumerClose fires promptly. |
| 295 | // maybeDrainAndSetState may transition consumer to Closed, making `ready` dangling. |
| 296 | if (isClosing) { |
| 297 | impl.maybeDrainAndSetState(js); |
| 298 | } |
| 299 | return js.resolvedPromise(DrainingReadResult{ |
| 300 | .chunks = chunks.releaseAsArray(), |
| 301 | .done = isClosing, |
| 302 | }); |
| 303 | } |
| 304 | |
| 305 | // No data available - need to wait. Queue a pending draining read. |
| 306 | // We create a ReadResult promise and transform it to DrainingReadResult. |
| 307 | // The flag remains set (was set at the start) and will be cleared by the promise callbacks. |
| 308 | auto prp = js.newPromiseAndResolver<ReadResult>(); |
| 309 | |
| 310 | ReadRequest request{.resolver = kj::mv(prp.resolver)}; |
| 311 | ready.readRequests.push_back(kj::heap<ReadRequest>(kj::mv(request))); |
| 312 | |
| 313 | KJ_IF_SOME(listener, impl.stateListener) { |
| 314 | listener.onConsumerWantsData(js); |
| 315 | } |
| 316 | |
| 317 | // Transform the ReadResult promise to DrainingReadResult. |
| 318 | return prp.promise.then( |
| 319 | js, [this](jsg::Lock& js, ReadResult result) mutable -> DrainingReadResult { |
| 320 | KJ_IF_SOME(ready, impl.state.tryGetActiveUnsafe()) { |
| 321 | ready.hasPendingDrainingRead = false; |
| 322 | } |
| 323 | |
| 324 | if (result.done) { |
| 325 | return DrainingReadResult{.chunks = nullptr, .done = true}; |
| 326 | } |
| 327 | |
| 328 | // Convert the value to bytes. |
| 329 | kj::Vector<kj::Array<kj::byte>> chunks; |
| 330 | KJ_IF_SOME(val, result.value) { |
| 331 | KJ_IF_SOME(bytes, valueToBytes(js, val)) { |
| 332 | chunks.add(kj::mv(bytes)); |
| 333 | } |
| 334 | // If valueToBytes returned kj::none, we just return empty chunks. |
| 335 | // The error case should have been caught earlier. |
| 336 | } |
| 337 | |
| 338 | return DrainingReadResult{ |
| 339 | .chunks = chunks.releaseAsArray(), |
| 340 | .done = false, |
| 341 | }; |
| 342 | }, [this](jsg::Lock& js, jsg::Value exception) mutable -> DrainingReadResult { |
| 343 | KJ_IF_SOME(ready, impl.state.tryGetActiveUnsafe()) { |
| 344 | ready.hasPendingDrainingRead = false; |
| 345 | } |
| 346 | js.throwException(kj::mv(exception)); |
| 347 | }); |
| 348 | } |
| 349 | |
| 350 | void ValueQueue::Consumer::cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { |
| 351 | impl.cancelPendingReads(js, reason); |
| 352 | } |
| 353 | |
| 354 | void ValueQueue::Consumer::visitForGc(jsg::GcVisitor& visitor) { |
| 355 | visitor.visit(impl); |
| 356 | } |
| 357 | |
| 358 | #pragma endregion ValueQueue::Consumer |
| 359 | |
| 360 | ValueQueue::ValueQueue(size_t highWaterMark): impl(highWaterMark) {} |
| 361 | |
| 362 | void ValueQueue::close(jsg::Lock& js) { |
| 363 | impl.close(js); |
| 364 | } |
| 365 | |
| 366 | ssize_t ValueQueue::desiredSize() const { |
| 367 | return impl.desiredSize(); |
| 368 | } |
| 369 | |
| 370 | void ValueQueue::error(jsg::Lock& js, jsg::Value reason) { |
| 371 | impl.error(js, kj::mv(reason)); |
| 372 | } |
| 373 | |
| 374 | void ValueQueue::maybeUpdateBackpressure() { |
| 375 | impl.maybeUpdateBackpressure(); |
| 376 | } |
| 377 | |
| 378 | void ValueQueue::push(jsg::Lock& js, kj::Rc<Entry> entry) { |
| 379 | impl.push(js, kj::mv(entry)); |
| 380 | } |
| 381 | |
| 382 | size_t ValueQueue::size() const { |
| 383 | return impl.size(); |
| 384 | } |
| 385 | |
| 386 | void ValueQueue::handlePush( |
| 387 | jsg::Lock& js, ConsumerImpl::Ready& state, kj::Maybe<QueueImpl&> queue, kj::Rc<Entry> entry) { |
| 388 | // If there are no pending reads, just add the entry to the buffer and return, adjusting |
| 389 | // the size of the queue in the process. |
| 390 | if (state.readRequests.empty()) { |
| 391 | state.queueTotalSize += entry->getSize(); |
| 392 | state.buffer.push_back(QueueEntry{.entry = kj::mv(entry)}); |
| 393 | return; |
| 394 | } |
| 395 | |
| 396 | // Otherwise, pop the next pending read and resolve it. There should be nothing in the queue. |
| 397 | KJ_REQUIRE(state.buffer.empty() && state.queueTotalSize == 0); |
| 398 | auto request = kj::mv(state.readRequests.front()); |
| 399 | state.readRequests.pop_front(); |
| 400 | request->resolve(js, entry->getValue(js)); |
| 401 | } |
| 402 | |
| 403 | void ValueQueue::handleRead(jsg::Lock& js, |
| 404 | ConsumerImpl::Ready& state, |
| 405 | ConsumerImpl& consumer, |
| 406 | kj::Maybe<QueueImpl&> queue, |
| 407 | ReadRequest request) { |
| 408 | // If there are no pending read requests and there is data in the buffer, |
| 409 | // we will try to fulfill the read request immediately. |
| 410 | if (state.queueTotalSize > 0 && state.buffer.empty()) { |
| 411 | // Is our queue accounting correct? |
| 412 | LOG_WARNING_ONCE("ValueQueue::handleRead encountered a queueTotalSize > 0 " |
| 413 | "with an empty buffer. This should not happen.", |
| 414 | state.queueTotalSize); |
| 415 | } |
| 416 | if (state.readRequests.empty() && !state.buffer.empty()) { |
| 417 | auto& entry = state.buffer.front(); |
| 418 | |
| 419 | KJ_SWITCH_ONEOF(entry) { |
| 420 | KJ_CASE_ONEOF(c, ConsumerImpl::Close) { |
| 421 | // This case shouldn't actually happen. The queueTotalSize should be zero if the |
| 422 | // only item remaining in the queue is the close sentinel because we decrement the |
| 423 | // queueTotalSize every time we remove an item. If we get here, something is wrong. |
| 424 | // We'll handle it by resolving the read request and keep going but let's emit a log |
| 425 | // warning so we can investigate. |
| 426 | // Note that we do not want to remove the close sentinel here so that the next call to |
| 427 | // maybeDrainAndSetState will see it and handle the transition to the closed state. |
| 428 | KJ_LOG(ERROR, |
| 429 | "ValueQueue::handleRead encountered a close sentinel in the queue " |
| 430 | "with queueTotalSize > 0. This should not happen.", |
| 431 | state.queueTotalSize); |
| 432 | request.resolveAsDone(js); |
| 433 | return; |
| 434 | } |
| 435 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 436 | auto freed = kj::mv(entry); |
| 437 | state.buffer.pop_front(); |
| 438 | request.resolve(js, freed.entry->getValue(js)); |
| 439 | state.queueTotalSize -= freed.entry->getSize(); |
| 440 | return; |
| 441 | } |
| 442 | } |
| 443 | KJ_UNREACHABLE; |
| 444 | } else if (state.queueTotalSize == 0 && consumer.isClosing()) { |
| 445 | // Otherwise, if state.queueTotalSize is zero and isClosing() is true there won't be any |
| 446 | // more data coming. Just resolve the read as done and move on. |
| 447 | request.resolveAsDone(js); |
| 448 | } else { |
| 449 | // Otherwise, push the read request into the pending readRequests. It will be |
| 450 | // resolved either as soon as there is data available or the consumer closes |
| 451 | // or errors. |
| 452 | state.readRequests.push_back(kj::heap<ReadRequest>(kj::mv(request))); |
| 453 | KJ_IF_SOME(listener, consumer.stateListener) { |
| 454 | listener.onConsumerWantsData(js); |
| 455 | } |
| 456 | } |
| 457 | } |
| 458 | |
| 459 | bool ValueQueue::handleMaybeClose(jsg::Lock& js, |
| 460 | ConsumerImpl::Ready& state, |
| 461 | ConsumerImpl& consumer, |
| 462 | kj::Maybe<QueueImpl&> queue) { |
| 463 | // If the value queue is not yet empty we have to keep waiting for more reads to consume it. |
| 464 | // Return false to indicate that we cannot close yet. |
| 465 | return false; |
| 466 | } |
| 467 | |
| 468 | size_t ValueQueue::getConsumerCount() { |
| 469 | return impl.getConsumerCount(); |
| 470 | } |
| 471 | |
| 472 | bool ValueQueue::wantsRead() const { |
| 473 | return impl.wantsRead(); |
| 474 | } |
| 475 | |
| 476 | bool ValueQueue::hasPartiallyFulfilledRead() { |
| 477 | // A ValueQueue can never have a partially fulfilled read. |
| 478 | return false; |
| 479 | } |
| 480 | |
| 481 | void ValueQueue::visitForGc(jsg::GcVisitor& visitor) {} |
| 482 | |
| 483 | #pragma endregion ValueQueue |
| 484 | |
| 485 | // ====================================================================================== |
| 486 | // ByteQueue |
| 487 | #pragma region ByteQueue |
| 488 | |
| 489 | #pragma region ByteQueue::ReadRequest |
| 490 | |
| 491 | namespace { |
| 492 | void maybeInvalidateByobRequest(kj::Maybe<ByteQueue::ByobRequest&>& req) { |
| 493 | KJ_IF_SOME(byobRequest, req) { |
| 494 | byobRequest.invalidate(); |
| 495 | // The call to byobRequest->invalidate() should have cleared the reference. |
| 496 | KJ_ASSERT(req == kj::none); |
| 497 | } |
| 498 | } |
| 499 | } // namespace |
| 500 | |
| 501 | ByteQueue::ReadRequest::ReadRequest( |
| 502 | jsg::Promise<ReadResult>::Resolver resolver, ByteQueue::ReadRequest::PullInto pullInto) |
| 503 | : resolver(kj::mv(resolver)), |
| 504 | pullInto(kj::mv(pullInto)) {} |
| 505 | |
| 506 | ByteQueue::ReadRequest::~ReadRequest() noexcept(false) { |
| 507 | maybeInvalidateByobRequest(byobReadRequest); |
| 508 | } |
| 509 | |
| 510 | void ByteQueue::ReadRequest::resolveAsDone(jsg::Lock& js) { |
| 511 | if (pullInto.filled > 0) { |
| 512 | // There's been at least some data written, we need to respond but not |
| 513 | // set done to true since that's what the streams spec requires. |
| 514 | pullInto.store.trim(js, pullInto.store.size() - pullInto.filled); |
| 515 | resolver.resolve( |
| 516 | js, ReadResult{.value = js.v8Ref(pullInto.store.getHandle(js)), .done = false}); |
| 517 | } else { |
| 518 | // Otherwise, we set the length to zero |
| 519 | pullInto.store.trim(js, pullInto.store.size()); |
| 520 | KJ_ASSERT(pullInto.store.size() == 0); |
| 521 | resolver.resolve(js, ReadResult{.value = js.v8Ref(pullInto.store.getHandle(js)), .done = true}); |
| 522 | } |
| 523 | maybeInvalidateByobRequest(byobReadRequest); |
| 524 | } |
| 525 | |
| 526 | void ByteQueue::ReadRequest::resolve(jsg::Lock& js) { |
| 527 | pullInto.store.trim(js, pullInto.store.size() - pullInto.filled); |
| 528 | resolver.resolve(js, ReadResult{.value = js.v8Ref(pullInto.store.getHandle(js)), .done = false}); |
| 529 | maybeInvalidateByobRequest(byobReadRequest); |
| 530 | } |
| 531 | |
| 532 | void ByteQueue::ReadRequest::reject(jsg::Lock& js, jsg::Value& value) { |
| 533 | resolver.reject(js, value.getHandle(js)); |
| 534 | maybeInvalidateByobRequest(byobReadRequest); |
| 535 | } |
| 536 | |
| 537 | kj::Own<ByteQueue::ByobRequest> ByteQueue::ReadRequest::makeByobReadRequest( |
| 538 | ConsumerImpl& consumer, QueueImpl& queue) { |
| 539 | auto req = kj::heap<ByobRequest>(*this, consumer, queue); |
| 540 | byobReadRequest = *req; |
| 541 | return kj::mv(req); |
| 542 | } |
| 543 | |
| 544 | #pragma endregion ByteQueue::ReadRequest |
| 545 | |
| 546 | #pragma region ByteQueue::Entry |
| 547 | |
| 548 | ByteQueue::Entry::Entry(jsg::BufferSource store): store(kj::mv(store)) {} |
| 549 | |
| 550 | kj::ArrayPtr<kj::byte> ByteQueue::Entry::toArrayPtr() { |
| 551 | return store.asArrayPtr(); |
| 552 | } |
| 553 | |
| 554 | size_t ByteQueue::Entry::getSize() const { |
| 555 | return store.size(); |
| 556 | } |
| 557 | |
| 558 | kj::Rc<ByteQueue::Entry> ByteQueue::Entry::clone(jsg::Lock& js) { |
| 559 | return addRefToThis(); |
| 560 | } |
| 561 | |
| 562 | void ByteQueue::Entry::visitForGc(jsg::GcVisitor& visitor) {} |
| 563 | |
| 564 | #pragma endregion ByteQueue::Entry |
| 565 | |
| 566 | #pragma region ByteQueue::QueueEntry |
| 567 | |
| 568 | ByteQueue::QueueEntry ByteQueue::QueueEntry::clone(jsg::Lock& js) { |
| 569 | return QueueEntry{ |
| 570 | .entry = entry->clone(js), |
| 571 | .offset = offset, |
| 572 | }; |
| 573 | } |
| 574 | |
| 575 | #pragma endregion ByteQueue::QueueEntry |
| 576 | |
| 577 | #pragma region ByteQueue::Consumer |
| 578 | |
| 579 | ByteQueue::Consumer::Consumer( |
| 580 | ByteQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 581 | : impl(queue.impl, stateListener) {} |
| 582 | |
| 583 | ByteQueue::Consumer::Consumer( |
| 584 | QueueImpl& impl, kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 585 | : impl(impl, stateListener) {} |
| 586 | |
| 587 | ByteQueue::Consumer::Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener) |
| 588 | : impl(stateListener) {} |
| 589 | |
| 590 | void ByteQueue::Consumer::cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) { |
| 591 | impl.cancel(js, maybeReason); |
| 592 | } |
| 593 | |
| 594 | void ByteQueue::Consumer::close(jsg::Lock& js) { |
| 595 | impl.close(js); |
| 596 | } |
| 597 | |
| 598 | bool ByteQueue::Consumer::empty() const { |
| 599 | return impl.empty(); |
| 600 | } |
| 601 | |
| 602 | void ByteQueue::Consumer::error(jsg::Lock& js, jsg::Value reason) { |
| 603 | impl.error(js, kj::mv(reason)); |
| 604 | } |
| 605 | |
| 606 | void ByteQueue::Consumer::read(jsg::Lock& js, ReadRequest request) { |
| 607 | impl.read(js, kj::mv(request)); |
| 608 | } |
| 609 | |
| 610 | void ByteQueue::Consumer::push(jsg::Lock& js, kj::Rc<Entry> entry) { |
| 611 | impl.push(js, kj::mv(entry)); |
| 612 | } |
| 613 | |
| 614 | void ByteQueue::Consumer::reset() { |
| 615 | impl.reset(); |
| 616 | } |
| 617 | |
| 618 | size_t ByteQueue::Consumer::size() const { |
| 619 | return impl.size(); |
| 620 | } |
| 621 | |
| 622 | kj::Own<ByteQueue::Consumer> ByteQueue::Consumer::clone( |
| 623 | jsg::Lock& js, kj::Maybe<ConsumerImpl::StateListener&> stateListener) { |
| 624 | // If the queue was destroyed (e.g., stream was closed), we can still clone |
| 625 | // the consumer - the cloneTo() will copy the closed/errored state. |
| 626 | kj::Own<Consumer> consumer; |
| 627 | KJ_IF_SOME(q, impl.queue) { |
| 628 | consumer = kj::heap<Consumer>(q, stateListener); |
| 629 | } else { |
| 630 | consumer = kj::heap<Consumer>(stateListener); |
| 631 | } |
| 632 | impl.cloneTo(js, consumer->impl); |
| 633 | return kj::mv(consumer); |
| 634 | } |
| 635 | |
| 636 | bool ByteQueue::Consumer::hasReadRequests() { |
| 637 | return impl.hasReadRequests(); |
| 638 | } |
| 639 | |
| 640 | bool ByteQueue::Consumer::hasPendingDrainingRead() { |
| 641 | return impl.state.whenActiveOr( |
| 642 | [](const ConsumerImpl::Ready& ready) { return ready.hasPendingDrainingRead; }, false); |
| 643 | } |
| 644 | |
| 645 | jsg::Promise<DrainingReadResult> ByteQueue::Consumer::drainingRead(jsg::Lock& js, size_t maxRead) { |
| 646 | // If there are pending regular reads, reject - mutual exclusion. |
| 647 | if (hasReadRequests()) { |
| 648 | return js.rejectedPromise<DrainingReadResult>( |
| 649 | js.typeError("Cannot call drainingRead while there are pending reads"_kj)); |
| 650 | } |
| 651 | |
| 652 | // Check if already closed or errored. |
| 653 | if (impl.state.template is<ConsumerImpl::Closed>()) { |
| 654 | return js.resolvedPromise(DrainingReadResult{.chunks = nullptr, .done = true}); |
| 655 | } |
| 656 | KJ_IF_SOME(errored, impl.state.tryGetErrorUnsafe()) { |
| 657 | return js.rejectedPromise<DrainingReadResult>(errored.reason.getHandle(js)); |
| 658 | } |
| 659 | |
| 660 | auto& ready = impl.state.requireActiveUnsafe(); |
| 661 | ConsumerImpl::UpdateBackpressureScope scope(impl); |
| 662 | |
| 663 | // Mark that we're doing a draining read. This allows onConsumerWantsData() |
| 664 | // to use forcePull() which bypasses backpressure checks. The flag is cleared |
| 665 | // either synchronously (for immediate returns) or in promise callbacks (for async). |
| 666 | ready.hasPendingDrainingRead = true; |
| 667 | |
| 668 | // Collect all buffered data (already bytes for ByteQueue). |
| 669 | kj::Vector<kj::Array<kj::byte>> chunks; |
| 670 | bool isClosing = false; |
| 671 | size_t totalRead = 0; |
| 672 | |
| 673 | // Drains buffered byte data into chunks. Stops draining when totalRead reaches |
| 674 | // or exceeds maxRead (after finishing the current item). |
| 675 | static const auto drainBuffer = [](ConsumerImpl::Ready& ready, |
| 676 | kj::Vector<kj::Array<kj::byte>>& chunks, size_t& totalRead, |
| 677 | bool& isClosing, size_t maxRead) { |
| 678 | while (!ready.buffer.empty() && !isClosing && totalRead < maxRead) { |
| 679 | auto& item = ready.buffer.front(); |
| 680 | KJ_SWITCH_ONEOF(item) { |
| 681 | KJ_CASE_ONEOF(close, ConsumerImpl::Close) { |
| 682 | isClosing = true; |
| 683 | break; |
| 684 | } |
| 685 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 686 | auto ptr = entry.entry->toArrayPtr(); |
| 687 | auto offset = entry.offset; |
| 688 | auto size = ptr.size() - offset; |
| 689 | totalRead += size; |
| 690 | chunks.add(kj::heapArray(ptr.slice(offset, offset + size))); |
| 691 | ready.queueTotalSize -= size; |
| 692 | ready.buffer.pop_front(); |
| 693 | } |
| 694 | } |
| 695 | } |
| 696 | }; |
| 697 | |
| 698 | // Drain the buffer up to maxRead bytes, then pump for more if under the limit. |
| 699 | drainBuffer(ready, chunks, totalRead, isClosing, maxRead); |
| 700 | |
| 701 | // Pump the controller for more synchronously available data. |
| 702 | // maxRead is checked here: we only proceed with pumping if we haven't exceeded it. |
| 703 | KJ_IF_SOME(listener, impl.stateListener) { |
| 704 | while (!isClosing && totalRead < maxRead) { |
| 705 | size_t prevChunkCount = chunks.size(); |
| 706 | bool pullCompletedSync = listener.onConsumerWantsData(js); |
| 707 | |
| 708 | // The pull callback may have closed or errored the consumer, which |
| 709 | // destroys the Ready state (and its RingBuffer). We must not touch |
| 710 | // `ready` after that. |
| 711 | if (!impl.state.isActive()) break; |
| 712 | |
| 713 | // Drain buffered data that was added by the pull, respecting maxRead. |
| 714 | drainBuffer(ready, chunks, totalRead, isClosing, maxRead); |
| 715 | |
| 716 | // If pull is async or no new data was added, stop pumping. |
| 717 | if (!pullCompletedSync || chunks.size() == prevChunkCount) { |
| 718 | break; |
| 719 | } |
| 720 | } |
| 721 | } |
| 722 | |
| 723 | // If the consumer was closed or errored during pumping, the `ready` |
| 724 | // reference is dangling. Return what we have or the appropriate error. |
| 725 | if (!impl.state.isActive()) { |
| 726 | KJ_IF_SOME(errored, impl.state.tryGetErrorUnsafe()) { |
| 727 | return js.rejectedPromise<DrainingReadResult>(errored.reason.getHandle(js)); |
| 728 | } |
| 729 | return js.resolvedPromise(DrainingReadResult{ |
| 730 | .chunks = chunks.releaseAsArray(), |
| 731 | .done = true, |
| 732 | }); |
| 733 | } |
| 734 | |
| 735 | // If the controller was canceled during pumping (e.g., pull callback called |
| 736 | // controller.cancel()), the QueueImpl is destroyed and the consumer's queue |
| 737 | // reference is detached (set to kj::none). The consumer state is still Active |
| 738 | // because cancel on the controller doesn't notify consumers — it only closes |
| 739 | // the controller's own state. No more data will ever arrive. Drain remaining |
| 740 | // buffer data respecting maxRead, and signal done only when the buffer is empty. |
| 741 | if (impl.queue == kj::none) { |
| 742 | // Drain remaining buffer up to maxRead. If there's still more, the caller |
| 743 | // will loop back and we'll drain the rest on subsequent calls. |
| 744 | drainBuffer(ready, chunks, totalRead, isClosing, maxRead); |
| 745 | ready.hasPendingDrainingRead = false; |
| 746 | bool done = ready.buffer.empty() || isClosing; |
| 747 | // If isClosing, finalize the consumer so onConsumerClose fires promptly. |
| 748 | // maybeDrainAndSetState may transition consumer to Closed, making `ready` dangling. |
| 749 | if (isClosing) { |
| 750 | impl.maybeDrainAndSetState(js); |
| 751 | } |
| 752 | return js.resolvedPromise(DrainingReadResult{ |
| 753 | .chunks = chunks.releaseAsArray(), |
| 754 | .done = done, |
| 755 | }); |
| 756 | } |
| 757 | |
| 758 | // If we collected data, return it immediately. |
| 759 | if (!chunks.empty() || isClosing) { |
| 760 | ready.hasPendingDrainingRead = false; |
| 761 | // If isClosing, finalize the consumer so onConsumerClose fires promptly. |
| 762 | // maybeDrainAndSetState may transition consumer to Closed, making `ready` dangling. |
| 763 | if (isClosing) { |
| 764 | impl.maybeDrainAndSetState(js); |
| 765 | } |
| 766 | return js.resolvedPromise(DrainingReadResult{ |
| 767 | .chunks = chunks.releaseAsArray(), |
| 768 | .done = isClosing, |
| 769 | }); |
| 770 | } |
| 771 | |
| 772 | // No data available - need to wait. Create a default read request. |
| 773 | // We allocate a buffer for the read - the data will be copied into it. |
| 774 | // The flag remains set (was set at the start) and will be cleared by the promise callbacks. |
| 775 | constexpr size_t kDefaultReadSize = 16384; // 16KB default buffer |
| 776 | KJ_IF_SOME(store, jsg::BufferSource::tryAllocUnsafe(js, kDefaultReadSize)) { |
| 777 | auto prp = js.newPromiseAndResolver<ReadResult>(); |
| 778 | |
| 779 | ReadRequest::PullInto pullInto{ |
| 780 | .store = kj::mv(store), |
| 781 | .filled = 0, |
| 782 | .atLeast = 1, |
| 783 | .type = ReadRequest::Type::DEFAULT, |
| 784 | }; |
| 785 | ReadRequest request(kj::mv(prp.resolver), kj::mv(pullInto)); |
| 786 | ready.readRequests.push_back(kj::heap<ReadRequest>(kj::mv(request))); |
| 787 | |
| 788 | KJ_IF_SOME(listener, impl.stateListener) { |
| 789 | listener.onConsumerWantsData(js); |
| 790 | } |
| 791 | |
| 792 | // Transform the ReadResult promise to DrainingReadResult. |
| 793 | return prp.promise.then( |
| 794 | js, [this](jsg::Lock& js, ReadResult result) mutable -> DrainingReadResult { |
| 795 | KJ_IF_SOME(ready, impl.state.tryGetActiveUnsafe()) { |
| 796 | ready.hasPendingDrainingRead = false; |
| 797 | } |
| 798 | |
| 799 | if (result.done) { |
| 800 | return DrainingReadResult{.chunks = nullptr, .done = true}; |
| 801 | } |
| 802 | |
| 803 | kj::Vector<kj::Array<kj::byte>> chunks; |
| 804 | KJ_IF_SOME(val, result.value) { |
| 805 | auto jsval = jsg::JsValue(val.getHandle(js)); |
| 806 | KJ_IF_SOME(ab, jsval.tryCast<jsg::JsArrayBuffer>()) { |
| 807 | chunks.add(kj::heapArray(ab.asArrayPtr())); |
| 808 | } else KJ_IF_SOME(abView, jsval.tryCast<jsg::JsArrayBufferView>()) { |
| 809 | chunks.add(kj::heapArray(abView.asArrayPtr())); |
| 810 | } |
| 811 | } |
| 812 | |
| 813 | return DrainingReadResult{ |
| 814 | .chunks = chunks.releaseAsArray(), |
| 815 | .done = false, |
| 816 | }; |
| 817 | }, [this](jsg::Lock& js, jsg::Value exception) mutable -> DrainingReadResult { |
| 818 | KJ_IF_SOME(ready, impl.state.tryGetActiveUnsafe()) { |
| 819 | ready.hasPendingDrainingRead = false; |
| 820 | } |
| 821 | js.throwException(kj::mv(exception)); |
| 822 | }); |
| 823 | } else { |
| 824 | return js.rejectedPromise<DrainingReadResult>( |
| 825 | js.error("Failed to allocate buffer for draining read"_kj)); |
| 826 | } |
| 827 | } |
| 828 | |
| 829 | void ByteQueue::Consumer::cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { |
| 830 | impl.cancelPendingReads(js, reason); |
| 831 | } |
| 832 | |
| 833 | void ByteQueue::Consumer::visitForGc(jsg::GcVisitor& visitor) { |
| 834 | visitor.visit(impl); |
| 835 | } |
| 836 | |
| 837 | #pragma endregion ByteQueue::Consumer |
| 838 | |
| 839 | #pragma region ByteQueue::ByobRequest |
| 840 | |
| 841 | ByteQueue::ByobRequest::~ByobRequest() noexcept(false) { |
| 842 | invalidate(); |
| 843 | } |
| 844 | |
| 845 | void ByteQueue::ByobRequest::invalidate() { |
| 846 | KJ_IF_SOME(req, request) { |
| 847 | req.byobReadRequest = kj::none; |
| 848 | request = kj::none; |
| 849 | } |
| 850 | } |
| 851 | |
| 852 | bool ByteQueue::ByobRequest::isPartiallyFulfilled() { |
| 853 | return !isInvalidated() && getRequest().pullInto.filled > 0 && |
| 854 | getRequest().pullInto.store.getElementSize() > 1; |
| 855 | } |
| 856 | |
| 857 | bool ByteQueue::ByobRequest::respond(jsg::Lock& js, size_t amount) { |
| 858 | // So what happens here? The read request has been fulfilled directly by writing |
| 859 | // into the storage buffer of the request. Unfortunately, this will only resolve |
| 860 | // the data for the one consumer from which the request was received. We have to |
| 861 | // copy the data into a refcounted ByteQueue::Entry that is pushed into the other |
| 862 | // known consumers. |
| 863 | |
| 864 | // First, we check to make sure that the request hasn't been invalidated already. |
| 865 | // Here, invalidated is a fancy word for the promise having been resolved or |
| 866 | // rejected already. |
| 867 | auto& req = KJ_REQUIRE_NONNULL(request, "the pending byob read request was already invalidated"); |
| 868 | |
| 869 | // The amount cannot be more than the total space in the request store. |
| 870 | JSG_REQUIRE(req.pullInto.filled + amount <= req.pullInto.store.size(), RangeError, |
| 871 | kj::str("Too many bytes [", amount, "] in response to a BYOB read request.")); |
| 872 | |
| 873 | auto sourcePtr = req.pullInto.store.asArrayPtr(); |
| 874 | |
| 875 | if (queue.getConsumerCount() > 1) { |
| 876 | // Allocate the entry into which we will be copying the provided data for the |
| 877 | // other consumers of the queue. |
| 878 | KJ_IF_SOME(store, jsg::BufferSource::tryAllocUnsafe(js, amount)) { |
| 879 | auto entry = kj::rc<Entry>(kj::mv(store)); |
| 880 | |
| 881 | auto start = sourcePtr.slice(req.pullInto.filled); |
| 882 | |
| 883 | // Safely copy the data over into the entry. |
| 884 | entry->toArrayPtr().first(amount).copyFrom(start.first(amount)); |
| 885 | |
| 886 | // Push the entry into the other consumers. |
| 887 | queue.push(js, kj::mv(entry), consumer); |
| 888 | } else { |
| 889 | js.throwException(js.error("Failed to allocate memory for the byob read response."_kj)); |
| 890 | } |
| 891 | } |
| 892 | |
| 893 | // For this consumer, if the number of bytes provided in the response does not |
| 894 | // align with the element size of the read into buffer, we need to shave off |
| 895 | // those extra bytes and push them into the consumers queue so they can be picked |
| 896 | // up by the next read. |
| 897 | req.pullInto.filled += amount; |
| 898 | |
| 899 | if (amount < req.pullInto.atLeast) { |
| 900 | // The response has not yet met the minimal requirement of this byob read. |
| 901 | // In this case, we do not want to resolve the read yet, and we do not |
| 902 | // want the byob request to be invalidated. We don't need to worry about |
| 903 | // unaligned bytes yet. We're just going to return false to tell the caller |
| 904 | // not to invalidate and to update the view over this store. |
| 905 | |
| 906 | // We do want to decrease the atLeast by the amount of bytes we received. |
| 907 | req.pullInto.atLeast -= amount; |
| 908 | return false; |
| 909 | } |
| 910 | |
| 911 | // There is no need to adjust the pullInto.atLeast here because we are resolving |
| 912 | // the read immediately. |
| 913 | |
| 914 | auto unaligned = req.pullInto.filled % req.pullInto.store.getElementSize(); |
| 915 | // It is possible that the request was partially filled already. |
| 916 | req.pullInto.filled -= unaligned; |
| 917 | |
| 918 | // Fulfill this request! |
| 919 | consumer.resolveRead(js, req); |
| 920 | |
| 921 | if (unaligned > 0) { |
| 922 | auto start = sourcePtr.slice(amount - unaligned); |
| 923 | |
| 924 | KJ_IF_SOME(store, jsg::BufferSource::tryAllocUnsafe(js, unaligned)) { |
| 925 | auto excess = kj::rc<Entry>(kj::mv(store)); |
| 926 | excess->toArrayPtr().first(unaligned).copyFrom(start.first(unaligned)); |
| 927 | consumer.push(js, kj::mv(excess)); |
| 928 | } else { |
| 929 | js.throwException(js.error("Failed to allocate memory for the byob read response."_kj)); |
| 930 | } |
| 931 | } |
| 932 | |
| 933 | return true; |
| 934 | } |
| 935 | |
| 936 | bool ByteQueue::ByobRequest::respondWithNewView(jsg::Lock& js, jsg::BufferSource view) { |
| 937 | // The idea here is that rather than filling the view that the controller was given, |
| 938 | // it chose to create its own view and fill that, likely over the same ArrayBuffer. |
| 939 | // What we do here is perform some basic validations on what we were given, and if |
| 940 | // those pass, we'll replace the backing store held in the req.pullInto with the one |
| 941 | // given, then continue on issuing the respond as normal. |
| 942 | auto& req = KJ_REQUIRE_NONNULL(request, "the pending byob read request was already invalidated"); |
| 943 | auto amount = view.size(); |
| 944 | |
| 945 | JSG_REQUIRE(view.canDetach(js), TypeError, "Unable to use non-detachable ArrayBuffer."); |
| 946 | JSG_REQUIRE(req.pullInto.store.getOffset() + req.pullInto.filled == view.getOffset(), RangeError, |
| 947 | "The given view has an invalid byte offset."); |
| 948 | JSG_REQUIRE(req.pullInto.store.size() == view.underlyingArrayBufferSize(js), RangeError, |
| 949 | "The underlying ArrayBuffer is not the correct length."); |
| 950 | JSG_REQUIRE(req.pullInto.filled + amount <= req.pullInto.store.size(), RangeError, |
| 951 | "The view is not the correct length."); |
| 952 | |
| 953 | req.pullInto.store = jsg::BufferSource(js, view.detach(js)); |
| 954 | return respond(js, amount); |
| 955 | } |
| 956 | |
| 957 | size_t ByteQueue::ByobRequest::getAtLeast() const { |
| 958 | KJ_IF_SOME(req, request) { |
| 959 | return req.pullInto.atLeast; |
| 960 | } |
| 961 | return 0; |
| 962 | } |
| 963 | |
| 964 | v8::Local<v8::Uint8Array> ByteQueue::ByobRequest::getView(jsg::Lock& js) { |
| 965 | KJ_IF_SOME(req, request) { |
| 966 | return req.pullInto.store |
| 967 | .getTypedViewSlice<v8::Uint8Array>(js, req.pullInto.filled, req.pullInto.store.size()) |
| 968 | .getHandle(js) |
| 969 | .As<v8::Uint8Array>(); |
| 970 | } |
| 971 | return v8::Local<v8::Uint8Array>(); |
| 972 | } |
| 973 | |
| 974 | size_t ByteQueue::ByobRequest::getOriginalBufferByteLength(jsg::Lock& js) const { |
| 975 | KJ_IF_SOME(req, request) { |
| 976 | KJ_IF_SOME(size, req.pullInto.store.underlyingArrayBufferSize(js)) { |
| 977 | return size; |
| 978 | } |
| 979 | } |
| 980 | return 0; |
| 981 | } |
| 982 | |
| 983 | size_t ByteQueue::ByobRequest::getOriginalByteOffsetPlusBytesFilled() const { |
| 984 | KJ_IF_SOME(req, request) { |
| 985 | return req.pullInto.store.getOffset() + req.pullInto.filled; |
| 986 | } |
| 987 | return 0; |
| 988 | } |
| 989 | |
| 990 | #pragma endregion ByteQueue::ByobRequest |
| 991 | |
| 992 | ByteQueue::ByteQueue(size_t highWaterMark): impl(highWaterMark) {} |
| 993 | |
| 994 | void ByteQueue::close(jsg::Lock& js) { |
| 995 | // Note: We intentionally do NOT invalidate pending byob requests here. |
| 996 | // According to the spec, the byobRequest should remain accessible after close |
| 997 | // so that respondWithNewView() can be called on it (which should throw |
| 998 | // appropriate errors for invalid views). The byob request will be invalidated |
| 999 | // when respond() or respondWithNewView() is called. |
| 1000 | if (!FeatureFlags::get(js).getPedanticWpt()) { |
| 1001 | KJ_IF_SOME(ready, impl.state.tryGetUnsafe<ByteQueue::QueueImpl::Ready>()) { |
| 1002 | while (!ready.pendingByobReadRequests.empty()) { |
| 1003 | ready.pendingByobReadRequests.front()->invalidate(); |
| 1004 | ready.pendingByobReadRequests.pop_front(); |
| 1005 | } |
| 1006 | } |
| 1007 | } |
| 1008 | impl.close(js); |
| 1009 | } |
| 1010 | |
| 1011 | ssize_t ByteQueue::desiredSize() const { |
| 1012 | return impl.desiredSize(); |
| 1013 | } |
| 1014 | |
| 1015 | void ByteQueue::error(jsg::Lock& js, jsg::Value reason) { |
| 1016 | impl.error(js, kj::mv(reason)); |
| 1017 | } |
| 1018 | |
| 1019 | void ByteQueue::maybeUpdateBackpressure() { |
| 1020 | KJ_IF_SOME(state, impl.getState()) { |
| 1021 | // Invalidated byob read requests will accumulate if we do not take |
| 1022 | // care of them from time to time. Since maybeUpdateBackpressure |
| 1023 | // is going to be called regularly while the queue is actively in use, |
| 1024 | // this is as good a place to clean them out as any. |
| 1025 | // |
| 1026 | // We iterate through the ring buffer and remove invalidated items from the front. |
| 1027 | // Since items are typically invalidated in order, this should be efficient for |
| 1028 | // the common case. |
| 1029 | while (!state.pendingByobReadRequests.empty() && |
| 1030 | state.pendingByobReadRequests.front()->isInvalidated()) { |
| 1031 | state.pendingByobReadRequests.pop_front(); |
| 1032 | } |
| 1033 | } |
| 1034 | impl.maybeUpdateBackpressure(); |
| 1035 | } |
| 1036 | |
| 1037 | void ByteQueue::push(jsg::Lock& js, kj::Rc<Entry> entry) { |
| 1038 | impl.push(js, kj::mv(entry)); |
| 1039 | } |
| 1040 | |
| 1041 | size_t ByteQueue::size() const { |
| 1042 | return impl.size(); |
| 1043 | } |
| 1044 | |
| 1045 | void ByteQueue::handlePush(jsg::Lock& js, |
| 1046 | ConsumerImpl::Ready& state, |
| 1047 | kj::Maybe<QueueImpl&> queue, |
| 1048 | kj::Rc<Entry> newEntry) { |
| 1049 | const auto bufferData = [&](size_t offset) { |
| 1050 | state.queueTotalSize += newEntry->getSize() - offset; |
| 1051 | state.buffer.emplace_back(QueueEntry{ |
| 1052 | .entry = kj::mv(newEntry), |
| 1053 | .offset = offset, |
| 1054 | }); |
| 1055 | }; |
| 1056 | |
| 1057 | // If there are no pending reads add the entry to the buffer. |
| 1058 | if (state.readRequests.empty()) { |
| 1059 | return bufferData(0); |
| 1060 | } |
| 1061 | |
| 1062 | // Otherwise, check the the pending reads in the buffer. If the amount |
| 1063 | // of data in the queue + the amount of data provided by this entry |
| 1064 | // are >= the pending reads atLeast, then we will fulfill the pending |
| 1065 | // read, and keep fulfilling pending reads as long as they are available. |
| 1066 | // Once we are out of pending reads, we will buffer the remaining data. |
| 1067 | auto entrySize = newEntry->getSize(); |
| 1068 | auto amountAvailable = state.queueTotalSize + entrySize; |
| 1069 | size_t entryOffset = 0; |
| 1070 | |
| 1071 | while (!state.readRequests.empty() && amountAvailable > 0) { |
| 1072 | auto& pending = *state.readRequests.front(); |
| 1073 | |
| 1074 | // If the amountAvailable is less than the pending read request's atLeast, |
| 1075 | // then we're just going to buffer the data and bailout without fulfilling |
| 1076 | // the read. We will take care of fulfilling the read later once there |
| 1077 | // is enough data. |
| 1078 | |
| 1079 | if (amountAvailable < pending.pullInto.atLeast) { |
| 1080 | return bufferData(0); |
| 1081 | } |
| 1082 | |
| 1083 | // There might be at least some data in the buffer. If there is, it should |
| 1084 | // not be more than the current pending.pullInfo.atLeast or something went |
| 1085 | // wrong somewhere else. |
| 1086 | KJ_REQUIRE(state.queueTotalSize < pending.pullInto.atLeast); |
| 1087 | |
| 1088 | // First, we copy any data in the buffer out to the pending.pullInto. This |
| 1089 | // should completely consume the current buffer. |
| 1090 | while (!state.buffer.empty()) { |
| 1091 | auto& next = state.buffer.front(); |
| 1092 | KJ_SWITCH_ONEOF(next) { |
| 1093 | KJ_CASE_ONEOF(c, ConsumerImpl::Close) { |
| 1094 | // This should have been caught by the isClosing() check above. |
| 1095 | KJ_FAIL_ASSERT("The consumer is closed."); |
| 1096 | } |
| 1097 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 1098 | auto sourcePtr = entry.entry->toArrayPtr(); |
| 1099 | auto sourceSize = sourcePtr.size() - entry.offset; |
| 1100 | |
| 1101 | auto destPtr = pending.pullInto.store.asArrayPtr().slice(pending.pullInto.filled); |
| 1102 | auto destAmount = pending.pullInto.store.size() - pending.pullInto.filled; |
| 1103 | |
| 1104 | // sourceSize is the amount of data remaining in the current entry to copy. |
| 1105 | // destAmount is the amount of space remaining to be filled in the pending read. |
| 1106 | // Because destAmount should be greater than or equal to atLeast, and because we |
| 1107 | // already checked that the queueTotalSize is less than atLeast, it should not be |
| 1108 | // possible for sourceSize to be zero nor greater than or equal to destAmount, |
| 1109 | // so let's verify. |
| 1110 | KJ_REQUIRE(sourceSize > 0 && sourceSize < destAmount); |
| 1111 | |
| 1112 | // Safely copy sourceSize bytes from sourcePtr to destPtr |
| 1113 | destPtr.first(sourceSize).copyFrom(sourcePtr.slice(entry.offset)); |
| 1114 | |
| 1115 | // We have completely consumed the data in this entry and can safely free |
| 1116 | // our reference to it now. Yay! |
| 1117 | auto released = kj::mv(next); |
| 1118 | state.buffer.pop_front(); |
| 1119 | |
| 1120 | pending.pullInto.filled += sourceSize; |
| 1121 | |
| 1122 | // There is no reason to adjust the pullInto.atLeast here because we |
| 1123 | // will be immediately resolving the read in the next step. |
| 1124 | |
| 1125 | state.queueTotalSize -= sourceSize; |
| 1126 | amountAvailable -= sourceSize; |
| 1127 | } |
| 1128 | } |
| 1129 | } |
| 1130 | |
| 1131 | // At this point, there shouldn't be any data remaining in the buffer. |
| 1132 | KJ_REQUIRE(state.queueTotalSize == 0); |
| 1133 | |
| 1134 | // And there should be data remaining in the pending pullInto destination. |
| 1135 | KJ_REQUIRE(pending.pullInto.filled < pending.pullInto.store.size()); |
| 1136 | |
| 1137 | // And the amountAvailable should be equal to the current push size. |
| 1138 | KJ_REQUIRE(amountAvailable == entrySize - entryOffset); |
| 1139 | |
| 1140 | // Now, we determine how much of the current entry we can copy into the |
| 1141 | // destination pullInto by taking the lesser of amountAvailable and |
| 1142 | // destination pullInto size - filled (which gives us the amount of space |
| 1143 | // remaining in the destination). |
| 1144 | auto amountToCopy = |
| 1145 | kj::min(amountAvailable, pending.pullInto.store.size() - pending.pullInto.filled); |
| 1146 | |
| 1147 | // The amountToCopy should not be more than the entry size minus the entryOffset |
| 1148 | // (which is the amount of data remaining to be consumed in the current entry). |
| 1149 | KJ_REQUIRE(amountToCopy <= entrySize - entryOffset); |
| 1150 | |
| 1151 | // The amountToCopy plus pending.pullInto.filled should be more than or equal to atLeast |
| 1152 | // and less than or equal pending.pullInto.store.size(). |
| 1153 | KJ_REQUIRE(amountToCopy + pending.pullInto.filled >= pending.pullInto.atLeast && |
| 1154 | amountToCopy + pending.pullInto.filled <= pending.pullInto.store.size()); |
| 1155 | |
| 1156 | // Awesome, so now we safely copy amountToCopy bytes from the current entry into |
| 1157 | // the remaining space in pending.pullInto.store, being careful to account for |
| 1158 | // the entryOffset and pending.pullInto.filled offsets to determine the range |
| 1159 | // where we start copying. |
| 1160 | auto entryPtr = newEntry->toArrayPtr(); |
| 1161 | auto destPtr = pending.pullInto.store.asArrayPtr().slice(pending.pullInto.filled); |
| 1162 | destPtr.first(amountToCopy).copyFrom(entryPtr.slice(entryOffset).first(amountToCopy)); |
| 1163 | |
| 1164 | // Yay! this pending read has been fulfilled. There might be more tho. Let's adjust |
| 1165 | // the amountAvailable and continue trying to consume data. |
| 1166 | amountAvailable -= amountToCopy; |
| 1167 | entryOffset += amountToCopy; |
| 1168 | pending.pullInto.filled += amountToCopy; |
| 1169 | |
| 1170 | // We do not need to adjust the pullInto.atLeast here since we are immediately |
| 1171 | // fulfilling the read at this point. |
| 1172 | |
| 1173 | auto request = kj::mv(state.readRequests.front()); |
| 1174 | state.readRequests.pop_front(); |
| 1175 | request->resolve(js); |
| 1176 | } |
| 1177 | |
| 1178 | // If the entry was consumed completely by the pending read, then we're done! |
| 1179 | // We don't have to buffer any data and shouldn't have any data in the buffer! |
| 1180 | // Since we possibly consumed data from the buffer, however, let's make sure |
| 1181 | // we tell the queue to update backpressure signaling. |
| 1182 | if (entryOffset == entrySize) { |
| 1183 | KJ_REQUIRE(state.queueTotalSize == 0); |
| 1184 | return; |
| 1185 | } |
| 1186 | |
| 1187 | // Otherwise, we need to buffer the remaining data, being careful to set the offset |
| 1188 | // for the data that we have already consumed. |
| 1189 | bufferData(entryOffset); |
| 1190 | } |
| 1191 | |
| 1192 | void ByteQueue::handleRead(jsg::Lock& js, |
| 1193 | ConsumerImpl::Ready& state, |
| 1194 | ConsumerImpl& consumer, |
| 1195 | kj::Maybe<QueueImpl&> queue, |
| 1196 | ReadRequest request) { |
| 1197 | const auto pendingRead = [&]() { |
| 1198 | bool isByob = request.pullInto.type == ReadRequest::Type::BYOB; |
| 1199 | state.readRequests.push_back(kj::heap<ReadRequest>(kj::mv(request))); |
| 1200 | if (isByob) { |
| 1201 | // Because ReadRequest is movable, and because the ByobRequest captures |
| 1202 | // a reference to the ReadRequest, we wait until after it is added to |
| 1203 | // state.readRequests to create the associated ByobRequest. |
| 1204 | // If the queue is none, the consumer was cloned from a closed stream |
| 1205 | // and we can't create a ByobRequest. If the queue state is none, |
| 1206 | // the queue has already been closed. |
| 1207 | KJ_IF_SOME(q, queue) { |
| 1208 | KJ_IF_SOME(queueState, q.getState()) { |
| 1209 | queueState.pendingByobReadRequests.push_back( |
| 1210 | state.readRequests.back()->makeByobReadRequest(consumer, q)); |
| 1211 | } |
| 1212 | } |
| 1213 | } |
| 1214 | KJ_IF_SOME(listener, consumer.stateListener) { |
| 1215 | listener.onConsumerWantsData(js); |
| 1216 | } |
| 1217 | }; |
| 1218 | |
| 1219 | const auto consume = [&](size_t amountToConsume) { |
| 1220 | while (amountToConsume > 0) { |
| 1221 | KJ_REQUIRE(!state.buffer.empty()); |
| 1222 | // There must be at least one item in the buffer. |
| 1223 | auto& item = state.buffer.front(); |
| 1224 | |
| 1225 | KJ_SWITCH_ONEOF(item) { |
| 1226 | KJ_CASE_ONEOF(c, ConsumerImpl::Close) { |
| 1227 | // We reached the end of the buffer! All data has been consumed. |
| 1228 | return true; |
| 1229 | } |
| 1230 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 1231 | // The amount to copy is the lesser of the current entry size minus |
| 1232 | // offset and the data remaining in the destination to fill. |
| 1233 | auto entrySize = entry.entry->getSize(); |
| 1234 | auto amountToCopy = kj::min( |
| 1235 | entrySize - entry.offset, request.pullInto.store.size() - request.pullInto.filled); |
| 1236 | auto elementSize = request.pullInto.store.getElementSize(); |
| 1237 | if (amountToCopy > elementSize) { |
| 1238 | amountToCopy -= amountToCopy % elementSize; |
| 1239 | } |
| 1240 | if (amountToConsume > elementSize) { |
| 1241 | amountToConsume -= amountToConsume % elementSize; |
| 1242 | } |
| 1243 | |
| 1244 | // Once we have the amount, we safely copy amountToCopy bytes from the |
| 1245 | // entry into the destination request, accounting properly for the offsets. |
| 1246 | auto sourcePtr = entry.entry->toArrayPtr().slice(entry.offset); |
| 1247 | auto destPtr = request.pullInto.store.asArrayPtr().slice(request.pullInto.filled); |
| 1248 | |
| 1249 | destPtr.first(amountToCopy).copyFrom(sourcePtr.first(amountToCopy)); |
| 1250 | |
| 1251 | request.pullInto.filled += amountToCopy; |
| 1252 | |
| 1253 | // If pullInto.atLeast is greater than amountToCopy, let's adjust |
| 1254 | // atLeast down by the number of bytes we've consumed, indicating |
| 1255 | // a smaller minimum read requirement. |
| 1256 | if (request.pullInto.atLeast > amountToCopy) { |
| 1257 | request.pullInto.atLeast -= amountToCopy; |
| 1258 | } else if (request.pullInto.atLeast == amountToCopy) { |
| 1259 | request.pullInto.atLeast = 1; |
| 1260 | } |
| 1261 | entry.offset += amountToCopy; |
| 1262 | amountToConsume -= amountToCopy; |
| 1263 | state.queueTotalSize -= amountToCopy; |
| 1264 | |
| 1265 | // If the entry.offset is equal to the size of the entry, then we've consumed the |
| 1266 | // entire thing and can free it and continue iterating. The amountToConsume might |
| 1267 | // be >= 0, we will check it at the start of the next iteration. |
| 1268 | if (entry.offset == entrySize) { |
| 1269 | auto released = kj::mv(item); |
| 1270 | state.buffer.pop_front(); |
| 1271 | continue; |
| 1272 | } |
| 1273 | |
| 1274 | // Otherwise, it is OK that there is data remaining but the amountToConsume |
| 1275 | // should be 0. Specifically, we either consume the entire entry and there |
| 1276 | // is data left over to consume, or we did not consume the entire entry |
| 1277 | // but read all that we can. |
| 1278 | KJ_REQUIRE(amountToConsume == 0); |
| 1279 | } |
| 1280 | } |
| 1281 | } |
| 1282 | return false; |
| 1283 | }; |
| 1284 | |
| 1285 | // If there are no pending read requests and there is data in the buffer, |
| 1286 | // we will try to fulfill the read request immediately. |
| 1287 | if (state.readRequests.empty() && state.queueTotalSize > 0) { |
| 1288 | // If the available size is less than the read requests atLeast, then |
| 1289 | // push the read request into the pending so we can wait for more data... |
| 1290 | |
| 1291 | if (state.queueTotalSize < request.pullInto.atLeast) { |
| 1292 | // If there is anything in the consumers queue at this point, We need to |
| 1293 | // copy those bytes into the byob buffer and advance the filled counter |
| 1294 | // forward that number of bytes. |
| 1295 | if (state.queueTotalSize > 0 && consume(state.queueTotalSize)) { |
| 1296 | return request.resolveAsDone(js); |
| 1297 | } |
| 1298 | return pendingRead(); |
| 1299 | } |
| 1300 | |
| 1301 | // Awesome, ok, it looks like we have enough data in the queue for us |
| 1302 | // to minimally fill this read request! The amount to copy is the lesser |
| 1303 | // of the queue total size and the maximum amount of space in the request |
| 1304 | // pull into. |
| 1305 | if (consume(kj::min(state.queueTotalSize, request.pullInto.store.size()))) { |
| 1306 | |
| 1307 | // If consume returns true, the consumer hit the end and we need to |
| 1308 | // just resolve the request as done and return. |
| 1309 | return request.resolveAsDone(js); |
| 1310 | } |
| 1311 | |
| 1312 | // Now, we can resolve the read promise. Since we consumed data from the |
| 1313 | // buffer, we also want to make sure to notify the queue so it can update |
| 1314 | // backpressure signaling. |
| 1315 | request.resolve(js); |
| 1316 | } else if (state.queueTotalSize == 0 && consumer.isClosing()) { |
| 1317 | // Otherwise, if size() is zero and isClosing() is true, we should have already |
| 1318 | // drained but let's take care of that now. Specifically, in this case there's |
| 1319 | // no data in the queue and close() has already been called, so there won't be |
| 1320 | // any more data coming. |
| 1321 | request.resolveAsDone(js); |
| 1322 | } else { |
| 1323 | // Otherwise, push the read request into the pending readRequests. It will be |
| 1324 | // resolved either as soon as there is data available or the consumer closes |
| 1325 | // or errors. |
| 1326 | return pendingRead(); |
| 1327 | } |
| 1328 | } |
| 1329 | |
| 1330 | bool ByteQueue::handleMaybeClose(jsg::Lock& js, |
| 1331 | ConsumerImpl::Ready& state, |
| 1332 | ConsumerImpl& consumer, |
| 1333 | kj::Maybe<QueueImpl&> queue) { |
| 1334 | // This is called when we know that we are closing and we still have data in |
| 1335 | // the queue. We want to see if we can drain as much of it into pending reads |
| 1336 | // as possible. If we're able to drain all of it, then yay! We can go ahead and |
| 1337 | // close. Otherwise we stay open and wait for more reads to consume the rest. |
| 1338 | |
| 1339 | // We should only be here if there is data remaining in the queue. |
| 1340 | KJ_ASSERT(state.queueTotalSize > 0); |
| 1341 | |
| 1342 | // We should also only be here if the consumer is closing. |
| 1343 | KJ_ASSERT(consumer.isClosing()); |
| 1344 | |
| 1345 | const auto consume = [&] { |
| 1346 | // Consume will copy as much of the remaining data in the buffer as possible |
| 1347 | // to the next pending read. If the remaining data can fit into the remaining |
| 1348 | // space in the read, awesome, we've consumed everything and we will return |
| 1349 | // true. If the remaining data cannot fit into the remaining space in the read, |
| 1350 | // then we'll return false to indicate that there's more data to consume. In |
| 1351 | // either case, the pending read is popped off the pending queue and resolved. |
| 1352 | |
| 1353 | KJ_ASSERT(!state.readRequests.empty()); |
| 1354 | auto& pending = *state.readRequests.front(); |
| 1355 | |
| 1356 | while (!state.buffer.empty()) { |
| 1357 | auto& next = state.buffer.front(); |
| 1358 | KJ_SWITCH_ONEOF(next) { |
| 1359 | KJ_CASE_ONEOF(c, ConsumerImpl::Close) { |
| 1360 | // We've reached the end! queueTotalSize should be zero. We need to |
| 1361 | // resolve and pop the current read and return true to indicate that |
| 1362 | // we're all done. |
| 1363 | // |
| 1364 | // Technically, we really shouldn't get here but the case is covered |
| 1365 | // just in case. |
| 1366 | KJ_ASSERT(state.queueTotalSize == 0); |
| 1367 | auto request = kj::mv(state.readRequests.front()); |
| 1368 | state.readRequests.pop_front(); |
| 1369 | request->resolve(js); |
| 1370 | return true; |
| 1371 | } |
| 1372 | KJ_CASE_ONEOF(entry, QueueEntry) { |
| 1373 | auto sourcePtr = entry.entry->toArrayPtr(); |
| 1374 | auto sourceSize = sourcePtr.size() - entry.offset; |
| 1375 | |
| 1376 | auto destPtr = pending.pullInto.store.asArrayPtr().slice(pending.pullInto.filled); |
| 1377 | auto destAmount = pending.pullInto.store.size() - pending.pullInto.filled; |
| 1378 | |
| 1379 | // There should be space available to copy into and data to copy from, or |
| 1380 | // something else went wrong. |
| 1381 | KJ_ASSERT(destAmount > 0); |
| 1382 | KJ_ASSERT(sourceSize > 0); |
| 1383 | |
| 1384 | // sourceSize is the amount of data remaining in the current entry to copy. |
| 1385 | // destAmount is the amount of space remaining to be filled in the pending read. |
| 1386 | auto amountToCopy = kj::min(sourceSize, destAmount); |
| 1387 | |
| 1388 | auto sourceStart = sourcePtr.slice(entry.offset); |
| 1389 | |
| 1390 | // It shouldn't be possible for sourceEnd to extend past the sourcePtr.end() |
| 1391 | // but let's make sure just to be safe. |
| 1392 | KJ_ASSERT(amountToCopy <= sourceStart.size()); |
| 1393 | |
| 1394 | // Safely copy amountToCopy bytes from the source into the destination. |
| 1395 | destPtr.first(amountToCopy).copyFrom(sourceStart.first(amountToCopy)); |
| 1396 | pending.pullInto.filled += amountToCopy; |
| 1397 | |
| 1398 | // We do not need to adjust down the atLeast here because, no matter what, |
| 1399 | // the read is going to be resolved either here or in the next iteration. |
| 1400 | |
| 1401 | state.queueTotalSize -= amountToCopy; |
| 1402 | entry.offset += amountToCopy; |
| 1403 | |
| 1404 | KJ_ASSERT(entry.offset <= sourcePtr.size()); |
| 1405 | |
| 1406 | if (amountToCopy == sourcePtr.size()) { |
| 1407 | // If amountToCopy is equal to sourcePtr.size(), we've consumed the entire entry |
| 1408 | // and we can free it. |
| 1409 | auto released = kj::mv(next); |
| 1410 | state.buffer.pop_front(); |
| 1411 | |
| 1412 | if (amountToCopy == destAmount) { |
| 1413 | // If the amountToCopy is equal to destAmount, then we've completely filled |
| 1414 | // this read request with the data remaining. Resolve the read request. If |
| 1415 | // state.queueTotalSize happens to be zero, we can safely indicate that we |
| 1416 | // have read the remaining data as this may have been the last actual value |
| 1417 | // entry in the buffer. |
| 1418 | auto request = kj::mv(state.readRequests.front()); |
| 1419 | state.readRequests.pop_front(); |
| 1420 | request->resolve(js); |
| 1421 | |
| 1422 | if (state.queueTotalSize == 0) { |
| 1423 | // If the queueTotalSize is zero at this point, the next item in the queue |
| 1424 | // must be a close and we can return true. All of the data has been consumed. |
| 1425 | KJ_ASSERT(state.buffer.front().is<ConsumerImpl::Close>()); |
| 1426 | return true; |
| 1427 | } |
| 1428 | |
| 1429 | // Otherwise, there's still data to consume, return false here to move on |
| 1430 | // to the next pending read (if any). |
| 1431 | return false; |
| 1432 | } |
| 1433 | |
| 1434 | // We know that amountToCopy cannot be greater than destAmount because |
| 1435 | // of the kj::min above. |
| 1436 | |
| 1437 | // Continuing here means that our pending read still has space to fill |
| 1438 | // and we might still have value entries to fill it. We'll iterate around |
| 1439 | // and see where we get. |
| 1440 | continue; |
| 1441 | } |
| 1442 | |
| 1443 | // This read did not consume everything in this entry but doesn't have |
| 1444 | // any more space to fill. We will resolve this read and return false |
| 1445 | // to indicate that the outer loop should continue with the next read |
| 1446 | // request if there is one. |
| 1447 | |
| 1448 | // At this point, it should be impossible for state.queueTotalSize to |
| 1449 | // be zero because there is still data remaining to be consumed in this |
| 1450 | // buffer. |
| 1451 | KJ_ASSERT(state.queueTotalSize > 0); |
| 1452 | |
| 1453 | auto request = kj::mv(state.readRequests.front()); |
| 1454 | state.readRequests.pop_front(); |
| 1455 | request->resolve(js); |
| 1456 | return false; |
| 1457 | } |
| 1458 | } |
| 1459 | } |
| 1460 | |
| 1461 | return state.queueTotalSize == 0; |
| 1462 | }; |
| 1463 | |
| 1464 | // We can only consume here if there are pending reads! |
| 1465 | while (!state.readRequests.empty()) { |
| 1466 | // We ignore the read request atLeast here since we are closing. Our goal is to |
| 1467 | // consume as much of the data as possible. |
| 1468 | |
| 1469 | if (consume()) { |
| 1470 | // If consume returns true, we reached the end and have no more data to |
| 1471 | // consume. That's a good thing! It means we can go ahead and close down. |
| 1472 | return true; |
| 1473 | } |
| 1474 | |
| 1475 | // If consume() returns false, there is still data left to consume in the queue. |
| 1476 | // We will loop around and try again so long as there are still read requests |
| 1477 | // pending. |
| 1478 | } |
| 1479 | |
| 1480 | // At this point, we shouldn't have any read requests and there should be data |
| 1481 | // left in the queue. We have to keep waiting for more reads to consume the |
| 1482 | // remaining data. |
| 1483 | KJ_ASSERT(state.queueTotalSize > 0); |
| 1484 | KJ_ASSERT(state.readRequests.empty()); |
| 1485 | |
| 1486 | return false; |
| 1487 | } |
| 1488 | |
| 1489 | kj::Maybe<kj::Own<ByteQueue::ByobRequest>> ByteQueue::nextPendingByobReadRequest() { |
| 1490 | KJ_IF_SOME(state, impl.getState()) { |
| 1491 | while (!state.pendingByobReadRequests.empty()) { |
| 1492 | auto request = kj::mv(state.pendingByobReadRequests.front()); |
| 1493 | state.pendingByobReadRequests.pop_front(); |
| 1494 | if (!request->isInvalidated()) { |
| 1495 | return kj::mv(request); |
| 1496 | } |
| 1497 | } |
| 1498 | } |
| 1499 | return kj::none; |
| 1500 | } |
| 1501 | |
| 1502 | bool ByteQueue::hasPartiallyFulfilledRead() { |
| 1503 | KJ_IF_SOME(state, impl.getState()) { |
| 1504 | if (!state.pendingByobReadRequests.empty()) { |
| 1505 | auto& pending = state.pendingByobReadRequests.front(); |
| 1506 | if (pending->isPartiallyFulfilled()) { |
| 1507 | return true; |
| 1508 | } |
| 1509 | } |
| 1510 | } |
| 1511 | return false; |
| 1512 | } |
| 1513 | |
| 1514 | bool ByteQueue::wantsRead() const { |
| 1515 | return impl.wantsRead(); |
| 1516 | } |
| 1517 | |
| 1518 | size_t ByteQueue::getConsumerCount() { |
| 1519 | return impl.getConsumerCount(); |
| 1520 | } |
| 1521 | |
| 1522 | void ByteQueue::visitForGc(jsg::GcVisitor& visitor) {} |
| 1523 | |
| 1524 | #pragma endregion ByteQueue |
| 1525 | |
| 1526 | } // namespace workerd::api |