Skip to content
File

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

56.4 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "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 
14namespace workerd::api {
15 
16// ======================================================================================
17// ValueQueue
18#pragma region ValueQueue
19 
20#pragma region ValueQueue::ReadRequest
21 
22void ValueQueue::ReadRequest::resolveAsDone(jsg::Lock& js) {
23 resolver.resolve(js, ReadResult{.done = true});
24}
25 
26void ValueQueue::ReadRequest::resolve(jsg::Lock& js, jsg::Value value) {
27 resolver.resolve(js, ReadResult{.value = kj::mv(value), .done = false});
28}
29 
30void 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 
38ValueQueue::Entry::Entry(jsg::Value value, size_t size): value(kj::mv(value)), size(size) {}
39 
40jsg::Value ValueQueue::Entry::getValue(jsg::Lock& js) {
41 return value.addRef(js);
42}
43 
44size_t ValueQueue::Entry::getSize() const {
45 return size;
46}
47 
48void ValueQueue::Entry::visitForGc(jsg::GcVisitor& visitor) {
49 visitor.visit(value);
50}
51 
52#pragma endregion ValueQueue::Entry
53 
54#pragma region ValueQueue::QueueEntry
55 
56kj::Rc<ValueQueue::Entry> ValueQueue::Entry::clone(jsg::Lock& js) {
57 return addRefToThis();
58}
59 
60ValueQueue::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 
68ValueQueue::Consumer::Consumer(
69 ValueQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener)
70 : impl(queue.impl, stateListener) {}
71 
72ValueQueue::Consumer::Consumer(
73 QueueImpl& impl, kj::Maybe<ConsumerImpl::StateListener&> stateListener)
74 : impl(impl, stateListener) {}
75 
76ValueQueue::Consumer::Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener)
77 : impl(stateListener) {}
78 
79void ValueQueue::Consumer::cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
80 impl.cancel(js, maybeReason);
81}
82 
83void ValueQueue::Consumer::close(jsg::Lock& js) {
84 impl.close(js);
85};
86 
87bool ValueQueue::Consumer::empty() {
88 return impl.empty();
89}
90 
91void ValueQueue::Consumer::error(jsg::Lock& js, jsg::Value reason) {
92 impl.error(js, kj::mv(reason));
93};
94 
95void ValueQueue::Consumer::read(jsg::Lock& js, ReadRequest request) {
96 impl.read(js, kj::mv(request));
97}
98 
99void ValueQueue::Consumer::push(jsg::Lock& js, kj::Rc<Entry> entry) {
100 impl.push(js, kj::mv(entry));
101}
102 
103void ValueQueue::Consumer::reset() {
104 impl.reset();
105};
106 
107size_t ValueQueue::Consumer::size() {
108 return impl.size();
109}
110 
111kj::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 
125bool ValueQueue::Consumer::hasReadRequests() {
126 return impl.hasReadRequests();
127}
128 
129bool ValueQueue::Consumer::hasPendingDrainingRead() {
130 return impl.state.whenActiveOr(
131 [](const ConsumerImpl::Ready& ready) { return ready.hasPendingDrainingRead; }, false);
132}
133 
134namespace {
135// Helper to convert a JS value to bytes. Returns kj::none if the value cannot be converted.
136kj::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 
162jsg::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 
350void ValueQueue::Consumer::cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) {
351 impl.cancelPendingReads(js, reason);
352}
353 
354void ValueQueue::Consumer::visitForGc(jsg::GcVisitor& visitor) {
355 visitor.visit(impl);
356}
357 
358#pragma endregion ValueQueue::Consumer
359 
360ValueQueue::ValueQueue(size_t highWaterMark): impl(highWaterMark) {}
361 
362void ValueQueue::close(jsg::Lock& js) {
363 impl.close(js);
364}
365 
366ssize_t ValueQueue::desiredSize() const {
367 return impl.desiredSize();
368}
369 
370void ValueQueue::error(jsg::Lock& js, jsg::Value reason) {
371 impl.error(js, kj::mv(reason));
372}
373 
374void ValueQueue::maybeUpdateBackpressure() {
375 impl.maybeUpdateBackpressure();
376}
377 
378void ValueQueue::push(jsg::Lock& js, kj::Rc<Entry> entry) {
379 impl.push(js, kj::mv(entry));
380}
381 
382size_t ValueQueue::size() const {
383 return impl.size();
384}
385 
386void 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 
403void 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 
459bool 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 
468size_t ValueQueue::getConsumerCount() {
469 return impl.getConsumerCount();
470}
471 
472bool ValueQueue::wantsRead() const {
473 return impl.wantsRead();
474}
475 
476bool ValueQueue::hasPartiallyFulfilledRead() {
477 // A ValueQueue can never have a partially fulfilled read.
478 return false;
479}
480 
481void 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 
491namespace {
492void 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 
501ByteQueue::ReadRequest::ReadRequest(
502 jsg::Promise<ReadResult>::Resolver resolver, ByteQueue::ReadRequest::PullInto pullInto)
503 : resolver(kj::mv(resolver)),
504 pullInto(kj::mv(pullInto)) {}
505 
506ByteQueue::ReadRequest::~ReadRequest() noexcept(false) {
507 maybeInvalidateByobRequest(byobReadRequest);
508}
509 
510void 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 
526void 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 
532void ByteQueue::ReadRequest::reject(jsg::Lock& js, jsg::Value& value) {
533 resolver.reject(js, value.getHandle(js));
534 maybeInvalidateByobRequest(byobReadRequest);
535}
536 
537kj::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 
548ByteQueue::Entry::Entry(jsg::BufferSource store): store(kj::mv(store)) {}
549 
550kj::ArrayPtr<kj::byte> ByteQueue::Entry::toArrayPtr() {
551 return store.asArrayPtr();
552}
553 
554size_t ByteQueue::Entry::getSize() const {
555 return store.size();
556}
557 
558kj::Rc<ByteQueue::Entry> ByteQueue::Entry::clone(jsg::Lock& js) {
559 return addRefToThis();
560}
561 
562void ByteQueue::Entry::visitForGc(jsg::GcVisitor& visitor) {}
563 
564#pragma endregion ByteQueue::Entry
565 
566#pragma region ByteQueue::QueueEntry
567 
568ByteQueue::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 
579ByteQueue::Consumer::Consumer(
580 ByteQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener)
581 : impl(queue.impl, stateListener) {}
582 
583ByteQueue::Consumer::Consumer(
584 QueueImpl& impl, kj::Maybe<ConsumerImpl::StateListener&> stateListener)
585 : impl(impl, stateListener) {}
586 
587ByteQueue::Consumer::Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener)
588 : impl(stateListener) {}
589 
590void ByteQueue::Consumer::cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
591 impl.cancel(js, maybeReason);
592}
593 
594void ByteQueue::Consumer::close(jsg::Lock& js) {
595 impl.close(js);
596}
597 
598bool ByteQueue::Consumer::empty() const {
599 return impl.empty();
600}
601 
602void ByteQueue::Consumer::error(jsg::Lock& js, jsg::Value reason) {
603 impl.error(js, kj::mv(reason));
604}
605 
606void ByteQueue::Consumer::read(jsg::Lock& js, ReadRequest request) {
607 impl.read(js, kj::mv(request));
608}
609 
610void ByteQueue::Consumer::push(jsg::Lock& js, kj::Rc<Entry> entry) {
611 impl.push(js, kj::mv(entry));
612}
613 
614void ByteQueue::Consumer::reset() {
615 impl.reset();
616}
617 
618size_t ByteQueue::Consumer::size() const {
619 return impl.size();
620}
621 
622kj::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 
636bool ByteQueue::Consumer::hasReadRequests() {
637 return impl.hasReadRequests();
638}
639 
640bool ByteQueue::Consumer::hasPendingDrainingRead() {
641 return impl.state.whenActiveOr(
642 [](const ConsumerImpl::Ready& ready) { return ready.hasPendingDrainingRead; }, false);
643}
644 
645jsg::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 
829void ByteQueue::Consumer::cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) {
830 impl.cancelPendingReads(js, reason);
831}
832 
833void ByteQueue::Consumer::visitForGc(jsg::GcVisitor& visitor) {
834 visitor.visit(impl);
835}
836 
837#pragma endregion ByteQueue::Consumer
838 
839#pragma region ByteQueue::ByobRequest
840 
841ByteQueue::ByobRequest::~ByobRequest() noexcept(false) {
842 invalidate();
843}
844 
845void ByteQueue::ByobRequest::invalidate() {
846 KJ_IF_SOME(req, request) {
847 req.byobReadRequest = kj::none;
848 request = kj::none;
849 }
850}
851 
852bool ByteQueue::ByobRequest::isPartiallyFulfilled() {
853 return !isInvalidated() && getRequest().pullInto.filled > 0 &&
854 getRequest().pullInto.store.getElementSize() > 1;
855}
856 
857bool 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 
936bool 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 
957size_t ByteQueue::ByobRequest::getAtLeast() const {
958 KJ_IF_SOME(req, request) {
959 return req.pullInto.atLeast;
960 }
961 return 0;
962}
963 
964v8::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 
974size_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 
983size_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 
992ByteQueue::ByteQueue(size_t highWaterMark): impl(highWaterMark) {}
993 
994void 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 
1011ssize_t ByteQueue::desiredSize() const {
1012 return impl.desiredSize();
1013}
1014 
1015void ByteQueue::error(jsg::Lock& js, jsg::Value reason) {
1016 impl.error(js, kj::mv(reason));
1017}
1018 
1019void 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 
1037void ByteQueue::push(jsg::Lock& js, kj::Rc<Entry> entry) {
1038 impl.push(js, kj::mv(entry));
1039}
1040 
1041size_t ByteQueue::size() const {
1042 return impl.size();
1043}
1044 
1045void 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 
1192void 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 
1330bool 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 
1489kj::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 
1502bool 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 
1514bool ByteQueue::wantsRead() const {
1515 return impl.wantsRead();
1516}
1517 
1518size_t ByteQueue::getConsumerCount() {
1519 return impl.getConsumerCount();
1520}
1521 
1522void ByteQueue::visitForGc(jsg::GcVisitor& visitor) {}
1523 
1524#pragma endregion ByteQueue
1525 
1526} // namespace workerd::api