Skip to content
File

Blob: src/workerd/api/streams/queue.h

cpp1235 lines
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#pragma once
6 
7#include "common.h"
8 
9#include <workerd/jsg/jsg.h>
10#include <workerd/util/ring-buffer.h>
11#include <workerd/util/small-set.h>
12#include <workerd/util/state-machine.h>
13#include <workerd/util/weak-refs.h>
14 
15namespace workerd::api {
16 
17// ============================================================================
18// Queues
19//
20// There are two kinds of queues used internally by the JavaScript-backed
21// ReadableStream implementation: value queues and byte queues. Each operate
22// in generally the same way but byte queues have a number of unique complexities
23// that make them more difficult.
24//
25// A queue (of either type) has the following general characteristics:
26//
27// - Every queue has a high water mark. This is the maximum amount of data
28// that should be stored in the queue pending consumption before backpressure
29// is signaled. Additional data can always be pushed into the queue beyond
30// the high water mark, but it is not advisable to do so.
31//
32// - All data stored in the queue is in the form of entries. The
33// specific type of entry depends on the queue type. Every entry has a
34// calculated size, which is dependent on the type of entry.
35//
36// - Every queue has one or more consumers. Each consumer maintains its own
37// internal buffer of entries that it has yet to consume. Entries are
38// structured such that there is ever only one copy of any given chunk of
39// data in memory, with each entry in each consumer possessing only a reference
40// to it. Whenever data is pushed into the queue, references are pushed into
41// each of the consumers. As data is consumed from the internal buffer, the
42// entries are freed. The underlying data is freed once the last
43// reference is released.
44//
45// - Every consumer has an remaining buffer size, which is the sum of the sizes
46// of all entries remaining to be consumed in its internal buffer.
47//
48// - A queue has a total queue size, which is the remaining buffer size of the
49// consumer with the most unconsumed data.
50//
51// - A queue has a desired size, which is the amount of additional data that
52// can be pushed into the queue before backpressure is signaled. It is
53// calculated by subtracting the total queue size from the high water mark.
54//
55// - Backpressure is signaled when desired size is equal to, or less than zero.
56//
57// Each type of queue has a specific kind of consumer. These generally operate
58// in the same way but Byte Queue consumers have a number of unique details.
59//
60// - As mentioned above, every consumer maintains an internal data buffer
61// consisting of references to the data that has been pushed into
62// the queue.
63//
64// - Every consumer maintains a list of pending reads. A read is a request to
65// consume some amount of data from the internal data buffer. If there is
66// enough data in the internal buffer to immediately fulfill the read request
67// when it is received, then we do so. Otherwise, the read is moved into the
68// pending reads list and is fulfilled later once there is enough data provided
69// to the consumer to do so.
70//
71// - When data is provided to a consumer by the queue, that data is added to
72// the internal buffer only if there are no pending reads capable of
73// immediately consuming the data.
74//
75// - When data is added to the internal buffer, the remaining buffer size is
76// incremented. When data is removed from the internal buffer, the remaining
77// buffer size is decremented.
78//
79// - Whenever the remaining buffer size for a consumer is modified, the queue
80// is asked to recalculate the desired size.
81//
82// For value queues, every individual entry is queued and consumed as a whole
83// unit. It is not possible to partially consume a single value entry. The
84// size of a value entry is calculated by a JavaScript function provided by
85// user code (the "size algorithm") if one is provided. If a size algorithm
86// is not provided, the default size of a value entry is exactly 1.
87//
88// The bookkeeping for a value queue is fairly simple:
89//
90// - A single value entry is created.
91// - Clones of that single value entry are distributed to each of
92// the value queue consumers.
93// - If a consumer has a pending read, the read is fulfilled immediately
94// and the reference is never added to that consumer's internal buffer.
95// - If the consumer has no pending reads, the reference is added to the
96// consumer's internal buffer and the remaining buffer size is incremented
97// by the calculated size of the value entry.
98// - Once the value entry has been delivered to each of the consumers,
99// the total queue size is updated by setting it equal to the maximum
100// remaining buffer size among the consumers.
101// - Later, when a consumer receives a read that consumes data from the
102// internal buffer, the remaining buffer size is decremented by the calculated
103// size of the value entry, and the queue is notified to re-evaluated the
104// total queue size.
105//
106// For byte queues, the situation becomes much more complicated for two
107// specific reasons: 1) All entries are in the form of arbitrarily long
108// byte sequences that can be partially consumed, and 2) read requests
109// made to the byte queue can be "BYOB" (bring your own buffer) in which
110// the intent is to avoid being forced to copy data between buffers by
111// having the reading code allocate and provide a buffer that the stream
112// implementation will read data into. When there is only a single consumer
113// for a streams data, the BYOB model is fairly straightforward and can be
114// implemented to avoid copying entirely. However, when you have multiple
115// consumers for a byte queue, all consuming data at different rates, it is
116// not possible to avoid copying entirely. Reads that consume byte data can
117// specify a range that crosses the boundaries of the individual entries that
118// are stored within the internal buffer, further complicating the process of
119// consuming data.
120//
121// To make matters even more complicated, a stream implementation is permitted
122// to ignore the allocated buffers provided by the BYOB read request and push
123// data into the queue as if the allocated buffer were not provided at all.
124// In such cases, the BYOB read request still needs to be fulfilled with the
125// provided buffer being written into it. Unfortunately, this is not uncommon.
126// React server-side rendering, for instance, will create byte-oriented
127// ReadableStreams that support BYOB reads, but will use the controller.enqueue()
128// API to push data into the stream rather than paying any attention to the
129// BYOB buffers provided by the readers.
130//
131// The requirement to support BYOB reads makes it critical to properly sequence
132// the delivery of BYOB read requests to the stream controller implementation,
133// ensuring the proper order of bytes delivered to each consumer while respecting
134// backpressure signaling such that backpressure is always determined by the
135// consumer that is being consumed at the slowest rate.
136//
137// On top of everything else, Workers introduces the concept of a minRead,
138// that is, a minimum number of bytes that a read request should consume from
139// the queue. The read promise should not be fulfilled unless either that
140// minimum number of bytes has been provided, or the stream is closed or errored.
141 
142template <typename Self>
143class ConsumerImpl;
144 
145template <typename Self>
146class QueueImpl;
147 
148// DrainingReadResult is defined in common.h
149 
150// Provides the underlying implementation shared by ByteQueue and ValueQueue.
151template <typename Self>
152class QueueImpl final {
153 public:
154 using ConsumerImpl = ConsumerImpl<Self>;
155 using Entry = Self::Entry;
156 using State = Self::State;
157 
158 explicit QueueImpl(size_t highWaterMark)
159 : highWaterMark(highWaterMark),
160 state(QueueState::template create<Ready>()) {}
161 
162 QueueImpl(QueueImpl&&) = default;
163 QueueImpl& operator=(QueueImpl&&) = default;
164 
165 ~QueueImpl() noexcept(false) {
166 // Detach all consumers before destruction to prevent UAF.
167 // This can happen during isolate teardown when the destruction order
168 // of JS wrapper objects doesn't follow the ownership hierarchy.
169 allConsumers.forEach([&](ConsumerImpl& consumer) { consumer.detachQueue(); });
170 }
171 
172 // Closes the queue. The close is forwarded on to all consumers.
173 // If we are already closed or errored, do nothing here.
174 void close(jsg::Lock& js) {
175 if (state.isActive()) {
176#ifdef KJ_DEBUG
177 isClosingOrErroring = true;
178 KJ_DEFER(isClosingOrErroring = false);
179#endif
180 allConsumers.forEach([&](ConsumerImpl& consumer) { consumer.close(js); });
181 state.template transitionTo<Closed>();
182 }
183 }
184 
185 // The amount of data the Queue needs until it is considered full.
186 // The value can be zero or negative, in which case backpressure is
187 // signaled on the queue.
188 // If the queue is already closed or errored, return 0.
189 inline ssize_t desiredSize() const {
190 return state.isActive() ? highWaterMark - size() : 0;
191 }
192 
193 // Errors the queue. The error is forwarded on to all consumers,
194 // which will, in turn, reset their internal buffers and reject
195 // all pending consume promises.
196 // If we are already closed or errored, do nothing here.
197 void error(jsg::Lock& js, jsg::Value reason) {
198 if (state.isActive()) {
199#ifdef KJ_DEBUG
200 isClosingOrErroring = true;
201 KJ_DEFER(isClosingOrErroring = false);
202#endif
203 allConsumers.forEach([&](ConsumerImpl& consumer) { consumer.error(js, reason.addRef(js)); });
204 state.template transitionTo<Errored>(kj::mv(reason));
205 }
206 }
207 
208 // Polls all known consumers to collect their current buffer sizes
209 // so that the current queue size can be updated.
210 // If we are already closed or errored, set totalQueueSize to zero.
211 void maybeUpdateBackpressure() {
212 totalQueueSize = 0;
213 if (state.isActive()) {
214 allConsumers.forEach([&](ConsumerImpl& consumer) {
215 totalQueueSize = kj::max(totalQueueSize, consumer.size());
216 });
217 }
218 }
219 
220 // Forwards the entry to all consumers (except skipConsumer if given).
221 // For each consumer, the entry will be used to fulfill any pending consume operations.
222 // If the entry type is byteOriented and has not been fully consumed by pending consume
223 // operations, then any left over data will be pushed into the consumer's buffer.
224 // Asserts if the queue is closed or errored.
225 void push(jsg::Lock& js, kj::Rc<Entry> entry, kj::Maybe<ConsumerImpl&> skipConsumer = kj::none) {
226 state.requireActiveUnsafe("The queue is closed or errored.");
227 
228 allConsumers.forEach([&](ConsumerImpl& consumer) {
229 KJ_IF_SOME(skip, skipConsumer) {
230 if (&skip == &consumer) {
231 return;
232 }
233 }
234 consumer.push(js, entry->clone(js));
235 });
236 }
237 
238 // The current size of consumer with the most stored data.
239 size_t size() const {
240 return totalQueueSize;
241 }
242 
243 size_t getConsumerCount() const {
244 return allConsumers.size();
245 }
246 
247 bool wantsRead() const {
248 if (state.isActive()) {
249 for (const auto& weakRef: allConsumers) {
250 KJ_IF_SOME(consumer, weakRef->tryGet()) {
251 if (consumer.hasReadRequests()) return true;
252 }
253 }
254 }
255 return false;
256 }
257 
258 // Specific queue implementations may provide additional state that is attached
259 // to the Ready struct.
260 kj::Maybe<State&> getState() KJ_LIFETIMEBOUND {
261 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
262 return ready;
263 }
264 return kj::none;
265 }
266 
267 inline kj::StringPtr jsgGetMemoryName() const;
268 inline size_t jsgGetMemorySelfSize() const;
269 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
270 
271 private:
272 struct Closed {
273 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
274 };
275 struct Errored {
276 static constexpr kj::StringPtr NAME KJ_UNUSED = "errored"_kj;
277 jsg::Value reason;
278 };
279 
280 struct Ready final: public State {
281 static constexpr kj::StringPtr NAME KJ_UNUSED = "ready"_kj;
282 };
283 
284 // State machine for QueueImpl:
285 // Ready -> Closed (close() called)
286 // Ready -> Errored (error() called)
287 // Closed is terminal, Errored is implicitly terminal via ErrorState.
288 using QueueState = StateMachine<TerminalStates<Closed>,
289 ErrorState<Errored>,
290 ActiveState<Ready>,
291 Ready,
292 Closed,
293 Errored>;
294 
295 size_t highWaterMark;
296 size_t totalQueueSize = 0;
297 QueueState state;
298 // The set of consumers attached to this queue. In the typical case this
299 // will be a very small number (often just one or two), so we use SmallSet to
300 // optimize for that. This persists across state transitions so we can detach
301 // consumers even after close()/error() transitions the queue to a terminal state.
302 //
303 // We store weak references to consumers to safely handle the case where a consumer
304 // is destroyed during iteration (e.g., resolving a read request triggers JS that
305 // destroys another consumer in the same queue). When iterating, we check if the WeakRef is still valid.
306 SmallSet<kj::Rc<WeakRef<ConsumerImpl>>> allConsumers;
307 
308#ifdef KJ_DEBUG
309 // Debug flag to detect if addConsumer is called during close/error iteration.
310 // This should never happen - it would indicate a bug in the streams implementation.
311 bool isClosingOrErroring = false;
312#endif
313 
314 void addConsumer(kj::Rc<WeakRef<ConsumerImpl>> weakRef) {
315 KJ_DASSERT(
316 !isClosingOrErroring, "Cannot add a consumer while the queue is being closed or errored");
317 allConsumers.add(kj::mv(weakRef));
318 }
319 
320 void removeConsumer(ConsumerImpl& consumer) {
321 allConsumers.removeIf([&consumer](const kj::Rc<WeakRef<ConsumerImpl>>& ref) {
322 KJ_IF_SOME(c, ref->tryGet()) {
323 return &c == &consumer;
324 }
325 return false; // Already invalid, will be cleaned up later
326 });
327 maybeUpdateBackpressure();
328 }
329 
330 friend Self;
331 friend ConsumerImpl;
332};
333 
334// Provides the underlying implementation shared by ByteQueue::Consumer and ValueQueue::Consumer
335template <typename Self>
336class ConsumerImpl final {
337 public:
338 struct StateListener {
339 virtual void onConsumerClose(jsg::Lock& js) = 0;
340 virtual void onConsumerError(jsg::Lock& js, jsg::Value reason) = 0;
341 // Called when the consumer has a pending read and needs data.
342 // Returns true if the pull algorithm completed synchronously (meaning
343 // more pumping might yield additional synchronous data), false if the
344 // pull is async (promise pending) or no pull was needed.
345 virtual bool onConsumerWantsData(jsg::Lock& js) = 0;
346 };
347 
348 using QueueImpl = QueueImpl<Self>;
349 
350 // A simple utility to be allocated on any stack where consumer buffer data maybe consumed
351 // or expanded. When the stack is unwound, it ensures the backpressure is appropriately
352 // updated. Captures the weakref to the consumer as there's a chance it'll be destroyed
353 // while the scope is pending.
354 struct UpdateBackpressureScope final {
355 kj::Rc<WeakRef<ConsumerImpl<Self>>> consumer;
356 UpdateBackpressureScope(ConsumerImpl& consumer): consumer(consumer.selfRef.addRef()) {}
357 ~UpdateBackpressureScope() noexcept(false) {
358 consumer->runIfAlive([](ConsumerImpl& consumer) {
359 KJ_IF_SOME(q, consumer.queue) {
360 q.maybeUpdateBackpressure();
361 }
362 });
363 }
364 KJ_DISALLOW_COPY_AND_MOVE(UpdateBackpressureScope);
365 };
366 
367 using ReadRequest = Self::ReadRequest;
368 using Entry = Self::Entry;
369 using QueueEntry = Self::QueueEntry;
370 
371 ConsumerImpl(QueueImpl& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none)
372 : queue(queue),
373 state(ConsumerState::template create<Ready>()),
374 stateListener(stateListener) {
375 queue.addConsumer(selfRef.addRef());
376 }
377 
378 explicit ConsumerImpl(kj::Maybe<ConsumerImpl::StateListener&> stateListener)
379 : queue(kj::none),
380 state(ConsumerState::template create<Ready>()),
381 stateListener(stateListener) {}
382 
383 KJ_DISALLOW_COPY_AND_MOVE(ConsumerImpl);
384 
385 ~ConsumerImpl() noexcept(false) {
386 // queue may be none if the queue was destroyed before this consumer
387 // (e.g., during isolate teardown) or if cloned from a closed stream.
388 // We must remove ourselves before invalidating selfRef, otherwise
389 // removeConsumer won't find us (tryGet() would return none).
390 KJ_IF_SOME(q, queue) {
391 q.removeConsumer(*this);
392 }
393 // Invalidate after removal so any concurrent iteration will skip us.
394 selfRef->invalidate();
395 }
396 
397 // Called by QueueImpl destructor to detach this consumer from a queue
398 // that is about to be destroyed.
399 void detachQueue() {
400 queue = kj::none;
401 }
402 
403 void cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
404 // Already closed or errored - nothing to do.
405 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
406 for (auto& request: ready.readRequests) {
407 request->resolveAsDone(js);
408 }
409 state.template transitionTo<Closed>();
410 }
411 }
412 
413 void close(jsg::Lock& js) {
414 // If we are already closed or errored, then we do nothing here.
415 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
416 // If we are not already closing, enqueue a Close sentinel.
417 if (!isClosing()) {
418 ready.buffer.push_back(Close{});
419 }
420 
421 // Then check to see if we need to drain pending reads and
422 // update the state to Closed.
423 return maybeDrainAndSetState(js);
424 }
425 }
426 
427 inline bool empty() const {
428 return size() == 0;
429 }
430 
431 void error(jsg::Lock& js, jsg::Value reason) {
432 // If we are already closed or errored, then we do nothing here.
433 // The new error doesn't matter.
434 if (state.isActive()) {
435 maybeDrainAndSetState(js, kj::mv(reason));
436 }
437 }
438 
439 void push(jsg::Lock& js, kj::Rc<Entry> entry) {
440 // If the consumer is already closed or errored, then we do nothing here.
441 // This can happen during iteration over consumers in QueueImpl::push() when
442 // resolving a read request on one consumer triggers JavaScript code that
443 // closes or errors another consumer in the same queue.
444 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
445 // If the consumer is already closing or the entry is empty, do nothing.
446 // Also skip if queue is none (consumer cloned from closed stream).
447 if (isClosing() || entry->getSize() == 0 || queue == kj::none) {
448 return;
449 }
450 
451 UpdateBackpressureScope scope(*this);
452 Self::handlePush(js, ready, queue, kj::mv(entry));
453 }
454 }
455 
456 void read(jsg::Lock& js, ReadRequest request) {
457 if (state.template is<Closed>()) {
458 return request.resolveAsDone(js);
459 }
460 KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) {
461 return request.reject(js, errored.reason);
462 }
463 auto& ready = state.requireActiveUnsafe();
464 // Mutual exclusion with draining reads.
465 if (ready.hasPendingDrainingRead) {
466 auto error = jsg::Value(
467 js.v8Isolate, js.typeError("Cannot call read while there is a pending draining read"_kj));
468 return request.reject(js, error);
469 }
470 Self::handleRead(js, ready, *this, queue, kj::mv(request));
471 return maybeDrainAndSetState(js);
472 }
473 
474 void reset() {
475 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
476 UpdateBackpressureScope scope(*this);
477 ready.buffer.clear();
478 ready.queueTotalSize = 0;
479 }
480 }
481 
482 // The current total calculated size of the consumer's internal buffer.
483 size_t size() const {
484 return state.whenActiveOr([](const Ready& ready) { return ready.queueTotalSize; }, 0ul);
485 }
486 
487 void resolveRead(jsg::Lock& js, ReadRequest& req) {
488 auto& ready = state.requireActiveUnsafe();
489 KJ_REQUIRE(!ready.readRequests.empty());
490 KJ_REQUIRE(&req == ready.readRequests.front().get());
491 // Pop the request before resolving to ensure the request is fully owned locally.
492 auto request = kj::mv(ready.readRequests.front());
493 ready.readRequests.pop_front();
494 request->resolve(js);
495 }
496 
497 void resolveReadAsDone(jsg::Lock& js, ReadRequest& req) {
498 auto& ready = state.requireActiveUnsafe();
499 KJ_REQUIRE(!ready.readRequests.empty());
500 KJ_REQUIRE(&req == ready.readRequests.front().get());
501 // Pop the request before resolving to ensure the request is fully owned locally.
502 auto request = kj::mv(ready.readRequests.front());
503 ready.readRequests.pop_front();
504 request->resolveAsDone(js);
505 }
506 
507 void cloneTo(jsg::Lock& js, ConsumerImpl& other) {
508 if (state.template is<Closed>()) {
509 other.state.template transitionTo<Closed>();
510 return;
511 }
512 KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) {
513 other.state.template transitionTo<Errored>(errored.reason.addRef(js));
514 return;
515 }
516 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
517 // We copy the buffered state but not the readRequests.
518 auto& otherReady = KJ_REQUIRE_NONNULL(
519 other.state.tryGetActiveUnsafe(), "The new consumer should not be closed or errored.");
520 otherReady.queueTotalSize = ready.queueTotalSize;
521 for (auto& item: ready.buffer) {
522 KJ_SWITCH_ONEOF(item) {
523 KJ_CASE_ONEOF(c, Close) {
524 otherReady.buffer.push_back(Close{});
525 }
526 KJ_CASE_ONEOF(entry, QueueEntry) {
527 otherReady.buffer.push_back(entry.clone(js));
528 }
529 }
530 }
531 }
532 }
533 
534 bool hasReadRequests() const {
535 return state.whenActiveOr(
536 [](const Ready& ready) { return !ready.readRequests.empty(); }, false);
537 }
538 
539 void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) {
540 // Already closed or errored - nothing to do.
541 state.whenActive([&](Ready& ready) {
542 for (auto& request: ready.readRequests) {
543 request->resolver.reject(js, reason);
544 }
545 ready.readRequests.clear();
546 });
547 }
548 
549 void visitForGc(jsg::GcVisitor& visitor) {
550 // Technically we shouldn't really have to GC visit the stored error here but there
551 // should not be any harm in doing so.
552 KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) {
553 visitor.visit(errored.reason);
554 }
555 // There's no reason to GC visit the promise resolver or buffer in Ready state and it is
556 // potentially problematic if we do. Since the read requests are queued, if we
557 // GC visit it once, remove it from the queue, and GC happens to kick in before
558 // we access the resolver, then v8 could determine that the resolver or buffered
559 // entries are no longer reachable via tracing and free them before we can
560 // actually try to access the held resolver.
561 }
562 
563 inline kj::StringPtr jsgGetMemoryName() const;
564 inline size_t jsgGetMemorySelfSize() const;
565 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
566 
567 private:
568 // A sentinel used in the buffer to signal that close() has been called.
569 struct Close {};
570 
571 struct Closed {
572 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
573 };
574 struct Errored {
575 static constexpr kj::StringPtr NAME KJ_UNUSED = "errored"_kj;
576 jsg::Value reason;
577 };
578 struct Ready {
579 static constexpr kj::StringPtr NAME KJ_UNUSED = "ready"_kj;
580 workerd::RingBuffer<kj::OneOf<QueueEntry, Close>, 16> buffer;
581 // We use kj::Own<ReadRequest> because ByobRequest holds a reference to its associated
582 // ReadRequest. Using RingBuffer directly would invalidate those references when the buffer
583 // grows. By heap-allocating each ReadRequest, we ensure reference stability.
584 workerd::RingBuffer<kj::Own<ReadRequest>, 8> readRequests;
585 size_t queueTotalSize = 0;
586 // True if there is a pending draining read operation. Draining reads are mutually
587 // exclusive with regular reads - read() will reject if this is true, and drainingRead()
588 // will reject if there are pending readRequests.
589 bool hasPendingDrainingRead = false;
590 
591 inline kj::StringPtr jsgGetMemoryName() const;
592 inline size_t jsgGetMemorySelfSize() const;
593 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
594 };
595 
596 // State machine for ConsumerImpl:
597 // Ready -> Closed (close() called and drained)
598 // Ready -> Errored (error() called)
599 // Closed is terminal, Errored is implicitly terminal via ErrorState.
600 using ConsumerState = StateMachine<TerminalStates<Closed>,
601 ErrorState<Errored>,
602 ActiveState<Ready>,
603 Ready,
604 Closed,
605 Errored>;
606 
607 kj::Maybe<QueueImpl&> queue;
608 ConsumerState state;
609 kj::Maybe<ConsumerImpl::StateListener&> stateListener;
610 // WeakRef to this consumer, used for safe registration with QueueImpl.
611 // When this consumer is destroyed, we invalidate the WeakRef so that
612 // any iteration over allConsumers in QueueImpl will safely skip us.
613 kj::Rc<WeakRef<ConsumerImpl>> selfRef =
614 kj::rc<WeakRef<ConsumerImpl>>(kj::Badge<ConsumerImpl>{}, *this);
615 
616 bool isClosing() {
617 // Closing state is determined by whether there is a Close sentinel that has been
618 // pushed into the end of Ready state buffer.
619 return state.whenActiveOr([](Ready& ready) {
620 return !ready.buffer.empty() && ready.buffer.back().template is<Close>();
621 }, false);
622 }
623 
624 // Extract all pending read requests from the ready state into a locally-owned vector.
625 // This is used by maybeDrainAndSetState to take ownership of pending reads before
626 // performing operations that may trigger V8 GC (resolve/reject calls use wrapOpaque
627 // which does V8 allocations). Without this, GC could collect the ReadableStream that
628 // owns this ConsumerImpl (through the ownership gap: QueueImpl only holds WeakRefs),
629 // destroying the readRequests ring buffer while we're iterating it.
630 static kj::Vector<kj::Own<ReadRequest>> extractPendingReads(Ready& ready) {
631 kj::Vector<kj::Own<ReadRequest>> result(ready.readRequests.size());
632 while (!ready.readRequests.empty()) {
633 result.add(kj::mv(ready.readRequests.front()));
634 ready.readRequests.pop_front();
635 }
636 return result;
637 }
638 
639 void maybeDrainAndSetState(jsg::Lock& js, kj::Maybe<jsg::Value> maybeReason = kj::none) {
640 // If the state is already errored or closed then there is nothing to drain.
641 KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
642 UpdateBackpressureScope scope(*this);
643 KJ_IF_SOME(reason, maybeReason) {
644 // If maybeReason != nullptr, then we are draining because of an error.
645 // In that case, we want to reset/clear the buffer and reject any remaining
646 // pending read requests using the given reason.
647 
648 // We extract pending reads to local ownership before rejecting. The reject
649 // calls perform V8 allocations (wrapOpaque) which can trigger GC. If GC
650 // collects the ReadableStream that owns this ConsumerImpl (see the ownership
651 // gap: QueueImpl only holds WeakRefs to consumers, actual ownership is through
652 // ReadableStream โ†’ ReadableStreamJsController โ†’ ValueReadable โ†’ Consumer),
653 // the ConsumerImpl and its readRequests would be destroyed mid-iteration.
654 // By extracting to a local vector, the ReadRequests survive even if `this`
655 // is destroyed during the reject calls.
656 //
657 // We preserve the original ordering (reject reads, then transition state,
658 // then notify listener) and use selfRef to check if `this` is still alive
659 // after the reject calls before accessing any members.
660 auto pendingReads = extractPendingReads(ready);
661 auto weak = selfRef.addRef();
662 for (auto& request: pendingReads) {
663 request->reject(js, reason);
664 }
665 // After the reject calls, `this` may have been destroyed by GC.
666 // Use the weak ref to safely access members only if still alive.
667 weak->runIfAlive([&](ConsumerImpl& self) {
668 self.state.template transitionTo<Errored>(reason.addRef(js));
669 KJ_IF_SOME(listener, self.stateListener) {
670 listener.onConsumerError(js, kj::mv(reason));
671 // After this point, we should not assume that this consumer can
672 // be safely used at all. It's most likely the stateListener has
673 // released it.
674 }
675 });
676 } else {
677 // Otherwise, if isClosing() is true...
678 if (isClosing()) {
679 if (!empty() && !Self::handleMaybeClose(js, ready, *this, queue)) {
680 // If the queue is not empty, we'll have the implementation see
681 // if it can drain the remaining data into pending reads. If handleMaybeClose
682 // returns false, then it could not and we can't yet close. If it returns true,
683 // yay! Our queue is empty and we can continue closing down.
684 KJ_ASSERT(!empty()); // We're still not empty
685 return;
686 }
687 
688 KJ_ASSERT(empty());
689 KJ_REQUIRE(ready.buffer.size() == 1); // The close should be the only item remaining.
690 
691 // Extract pending reads and resolve them as done. Same GC safety concern
692 // as the error path above โ€” see detailed comment there.
693 auto pendingReads = extractPendingReads(ready);
694 auto weak = selfRef.addRef();
695 for (auto& request: pendingReads) {
696 request->resolveAsDone(js);
697 }
698 // After the resolve calls, `this` may have been destroyed by GC.
699 weak->runIfAlive([&](ConsumerImpl& self) {
700 self.state.template transitionTo<Closed>();
701 KJ_IF_SOME(listener, self.stateListener) {
702 listener.onConsumerClose(js);
703 // After this point, we should not assume that this consumer can
704 // be safely used at all. It's most likely the stateListener has
705 // released it.
706 }
707 });
708 }
709 }
710 }
711 }
712 
713 friend Self::Consumer;
714 friend Self;
715};
716 
717// ============================================================================
718// Value queue
719 
720class ValueQueue final {
721 public:
722 using ConsumerImpl = ConsumerImpl<ValueQueue>;
723 using QueueImpl = QueueImpl<ValueQueue>;
724 
725 struct State {
726 JSG_MEMORY_INFO(ValueQueue::State) {}
727 };
728 
729 struct ReadRequest {
730 jsg::Promise<ReadResult>::Resolver resolver;
731 
732 void resolveAsDone(jsg::Lock& js);
733 void resolve(jsg::Lock& js, jsg::Value value);
734 void reject(jsg::Lock& js, jsg::Value& value);
735 
736 JSG_MEMORY_INFO(ValueQueue::ReadRequest) {
737 tracker.trackField("resolver", resolver);
738 }
739 };
740 
741 // A value queue entry consists of an arbitrary JavaScript value and a size that is
742 // calculated by the size algorithm function provided in the stream constructor.
743 class Entry: public kj::Refcounted {
744 public:
745 explicit Entry(jsg::Value value, size_t size);
746 KJ_DISALLOW_COPY_AND_MOVE(Entry);
747 
748 jsg::Value getValue(jsg::Lock& js);
749 
750 size_t getSize() const;
751 
752 void visitForGc(jsg::GcVisitor& visitor);
753 
754 kj::Rc<Entry> clone(jsg::Lock& js);
755 
756 JSG_MEMORY_INFO(ValueQueue::Entry) {
757 tracker.trackField("value", value);
758 }
759 
760 private:
761 jsg::Value value;
762 size_t size;
763 };
764 
765 struct QueueEntry {
766 kj::Rc<Entry> entry;
767 QueueEntry clone(jsg::Lock& js);
768 
769 JSG_MEMORY_INFO(ValueQueue::QueueEntry) {
770 tracker.trackFieldWithSize("entry", entry->getSize());
771 }
772 };
773 
774 class Consumer final {
775 public:
776 Consumer(ValueQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
777 Consumer(QueueImpl& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
778 // Used when cloning a consumer whose queue has been destroyed.
779 explicit Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener);
780 Consumer(Consumer&&) = delete;
781 Consumer(Consumer&) = delete;
782 Consumer& operator=(Consumer&&) = delete;
783 Consumer& operator=(Consumer&) = delete;
784 
785 void cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
786 
787 void close(jsg::Lock& js);
788 
789 bool empty();
790 
791 void error(jsg::Lock& js, jsg::Value reason);
792 
793 void read(jsg::Lock& js, ReadRequest request);
794 
795 // Draining read for optimized pipe-to operations. Drains all currently buffered
796 // data, pumps the controller for synchronously available data, and converts
797 // all values to bytes. Values must be ArrayBuffer, ArrayBufferView, or string;
798 // other types will error the stream.
799 // Rejects if there are pending regular reads (mutual exclusion).
800 // Regular read() will reject if there is a pending draining read.
801 // The maxRead parameter is a soft limit - see ReadableStreamController::drainingRead.
802 jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead = kj::maxValue);
803 
804 void push(jsg::Lock& js, kj::Rc<Entry> entry);
805 
806 void reset();
807 
808 size_t size();
809 
810 kj::Own<Consumer> clone(
811 jsg::Lock& js, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
812 
813 bool hasReadRequests();
814 bool hasPendingDrainingRead();
815 void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason);
816 
817 void visitForGc(jsg::GcVisitor& visitor);
818 
819 inline kj::StringPtr jsgGetMemoryName() const;
820 inline size_t jsgGetMemorySelfSize() const;
821 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
822 
823 private:
824 ConsumerImpl impl;
825 
826 friend class ValueQueue;
827 };
828 
829 explicit ValueQueue(size_t highWaterMark);
830 
831 void close(jsg::Lock& js);
832 
833 ssize_t desiredSize() const;
834 
835 void error(jsg::Lock& js, jsg::Value reason);
836 
837 void maybeUpdateBackpressure();
838 
839 void push(jsg::Lock& js, kj::Rc<Entry> entry);
840 
841 size_t size() const;
842 
843 size_t getConsumerCount();
844 
845 bool wantsRead() const;
846 
847 bool hasPartiallyFulfilledRead();
848 
849 void visitForGc(jsg::GcVisitor& visitor);
850 
851 inline kj::StringPtr jsgGetMemoryName() const;
852 inline size_t jsgGetMemorySelfSize() const;
853 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
854 
855 private:
856 QueueImpl impl;
857 
858 static void handlePush(
859 jsg::Lock& js, ConsumerImpl::Ready& state, kj::Maybe<QueueImpl&> queue, kj::Rc<Entry> entry);
860 static void handleRead(jsg::Lock& js,
861 ConsumerImpl::Ready& state,
862 ConsumerImpl& consumer,
863 kj::Maybe<QueueImpl&> queue,
864 ReadRequest request);
865 static bool handleMaybeClose(jsg::Lock& js,
866 ConsumerImpl::Ready& state,
867 ConsumerImpl& consumer,
868 kj::Maybe<QueueImpl&> queue);
869 
870 friend ConsumerImpl;
871};
872 
873// ============================================================================
874// Byte queue
875 
876class ByteQueue final {
877 public:
878 using ConsumerImpl = ConsumerImpl<ByteQueue>;
879 using QueueImpl = QueueImpl<ByteQueue>;
880 
881 class ByobRequest;
882 
883 struct ReadRequest final {
884 enum class Type { DEFAULT, BYOB };
885 jsg::Promise<ReadResult>::Resolver resolver;
886 // The reference here should be cleared when the ByobRequest is invalidated,
887 // which happens either when respond(), respondWithNewView(), or invalidate()
888 // is called, or when the ByobRequest is destroyed, whichever comes first.
889 kj::Maybe<ByobRequest&> byobReadRequest;
890 
891 struct PullInto {
892 jsg::BufferSource store;
893 size_t filled = 0;
894 size_t atLeast = 1;
895 Type type = Type::DEFAULT;
896 
897 JSG_MEMORY_INFO(ByteQueue::ReadRequest::PullInto) {
898 tracker.trackField("store", store);
899 }
900 } pullInto;
901 
902 ReadRequest(jsg::Promise<ReadResult>::Resolver resolver, PullInto pullInto);
903 ReadRequest(ReadRequest&&) = default;
904 ReadRequest& operator=(ReadRequest&&) = default;
905 ~ReadRequest() noexcept(false);
906 void resolveAsDone(jsg::Lock& js);
907 void resolve(jsg::Lock& js);
908 void reject(jsg::Lock& js, jsg::Value& value);
909 
910 kj::Own<ByobRequest> makeByobReadRequest(ConsumerImpl& consumer, QueueImpl& queue);
911 
912 JSG_MEMORY_INFO(ByteQueue::ReadRequest) {
913 tracker.trackField("resolver", resolver);
914 tracker.trackField("pullInto", pullInto);
915 }
916 };
917 
918 // The ByobRequest is essentially a handle to the ByteQueue::ReadRequest that can be given to a
919 // ReadableStreamBYOBRequest object to fulfill the request using the BYOB API pattern.
920 //
921 // When isInvalidated() is false, respond() or respondWithNewView() can be called to fulfill
922 // the BYOB read request. Once either of those are called, or once invalidate() is called,
923 // the ByobRequest is no longer usable and should be discarded.
924 class ByobRequest final {
925 public:
926 ByobRequest(ReadRequest& request, ConsumerImpl& consumer, QueueImpl& queue)
927 : request(request),
928 consumer(consumer),
929 queue(queue) {}
930 
931 KJ_DISALLOW_COPY_AND_MOVE(ByobRequest);
932 
933 ~ByobRequest() noexcept(false);
934 
935 inline ReadRequest& getRequest() {
936 return KJ_ASSERT_NONNULL(request);
937 }
938 
939 bool respond(jsg::Lock& js, size_t amount);
940 
941 bool respondWithNewView(jsg::Lock& js, jsg::BufferSource view);
942 
943 // Disconnects this ByobRequest instance from the associated ByteQueue::ReadRequest.
944 // The term "invalidate" is adopted from the streams spec for handling BYOB requests.
945 void invalidate();
946 
947 inline bool isInvalidated() const {
948 return request == kj::none;
949 }
950 
951 bool isPartiallyFulfilled();
952 
953 size_t getAtLeast() const;
954 
955 v8::Local<v8::Uint8Array> getView(jsg::Lock& js);
956 
957 // Returns the byte length of the original underlying ArrayBuffer.
958 size_t getOriginalBufferByteLength(jsg::Lock& js) const;
959 
960 // Returns the byte offset of the original view plus bytes filled.
961 size_t getOriginalByteOffsetPlusBytesFilled() const;
962 
963 JSG_MEMORY_INFO(ByteQueue::ByobRequest) {}
964 
965 private:
966 kj::Maybe<ReadRequest&> request;
967 ConsumerImpl& consumer;
968 QueueImpl& queue;
969 };
970 
971 struct State {
972 // We use a ring buffer for pending BYOB read requests. Since we store kj::Own<ByobRequest>,
973 // the actual ByobRequest objects are heap-allocated and won't be invalidated by buffer growth.
974 workerd::RingBuffer<kj::Own<ByobRequest>, 8> pendingByobReadRequests;
975 
976 JSG_MEMORY_INFO(ByteQueue::State) {
977 for (auto& request: pendingByobReadRequests) {
978 tracker.trackField("pendingByobReadRequest", request);
979 }
980 }
981 };
982 
983 // A byte queue entry consists of a jsg::BufferSource containing a non-zero-length
984 // sequence of bytes. The size is determined by the number of bytes in the entry.
985 class Entry: public kj::Refcounted {
986 public:
987 explicit Entry(jsg::BufferSource store);
988 
989 kj::ArrayPtr<kj::byte> toArrayPtr();
990 
991 size_t getSize() const;
992 
993 void visitForGc(jsg::GcVisitor& visitor);
994 
995 kj::Rc<Entry> clone(jsg::Lock& js);
996 
997 JSG_MEMORY_INFO(ByteQueue::Entry) {
998 tracker.trackField("store", store);
999 }
1000 
1001 private:
1002 jsg::BufferSource store;
1003 };
1004 
1005 struct QueueEntry {
1006 kj::Rc<Entry> entry;
1007 size_t offset;
1008 
1009 QueueEntry clone(jsg::Lock& js);
1010 
1011 JSG_MEMORY_INFO(ByteQueue::QueueEntry) {
1012 tracker.trackFieldWithSize("entry", entry->getSize());
1013 }
1014 };
1015 
1016 class Consumer {
1017 public:
1018 Consumer(ByteQueue& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
1019 Consumer(QueueImpl& queue, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
1020 // Used when cloning a consumer whose queue has been destroyed.
1021 explicit Consumer(kj::Maybe<ConsumerImpl::StateListener&> stateListener);
1022 Consumer(Consumer&&) = delete;
1023 Consumer(Consumer&) = delete;
1024 Consumer& operator=(Consumer&&) = delete;
1025 Consumer& operator=(Consumer&) = delete;
1026 
1027 void cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
1028 
1029 void close(jsg::Lock& js);
1030 
1031 bool empty() const;
1032 
1033 void error(jsg::Lock& js, jsg::Value reason);
1034 
1035 void read(jsg::Lock& js, ReadRequest request);
1036 
1037 // Draining read for optimized pipe-to operations. Drains all currently buffered
1038 // data and pumps the controller for synchronously available data.
1039 // Returns bytes directly without conversion (data is already bytes).
1040 // Rejects if there are pending regular reads (mutual exclusion).
1041 // Regular read() will reject if there is a pending draining read.
1042 // The maxRead parameter is a soft limit - see ReadableStreamController::drainingRead.
1043 jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead = kj::maxValue);
1044 
1045 void push(jsg::Lock& js, kj::Rc<Entry> entry);
1046 
1047 void reset();
1048 
1049 size_t size() const;
1050 
1051 kj::Own<Consumer> clone(
1052 jsg::Lock& js, kj::Maybe<ConsumerImpl::StateListener&> stateListener = kj::none);
1053 bool hasReadRequests();
1054 bool hasPendingDrainingRead();
1055 void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason);
1056 
1057 void visitForGc(jsg::GcVisitor& visitor);
1058 
1059 inline kj::StringPtr jsgGetMemoryName() const;
1060 inline size_t jsgGetMemorySelfSize() const;
1061 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
1062 
1063 private:
1064 ConsumerImpl impl;
1065 };
1066 
1067 explicit ByteQueue(size_t highWaterMark);
1068 
1069 void close(jsg::Lock& js);
1070 
1071 ssize_t desiredSize() const;
1072 
1073 void error(jsg::Lock& js, jsg::Value reason);
1074 
1075 void maybeUpdateBackpressure();
1076 
1077 void push(jsg::Lock& js, kj::Rc<Entry> entry);
1078 
1079 size_t size() const;
1080 
1081 size_t getConsumerCount();
1082 
1083 bool wantsRead() const;
1084 
1085 bool hasPartiallyFulfilledRead();
1086 
1087 // nextPendingByobReadRequest will be used to support the ReadableStreamBYOBRequest interface
1088 // that is part of ReadableByteStreamController. When user code calls the `controller.byobRequest`
1089 // API on a ReadableByteStreamController, they are going to get an instance of a
1090 // ReadableStreamBYOBRequest object. That object will own the `kj::Own<ByobReadRequest>` that
1091 // is returned here. User code could end up doing something silly like holding a reference to
1092 // that byobRequest long after it has been invalidated. We heap-allocate these just to allow
1093 // their lifespan to be attached to the ReadableStreamBYOBRequest object but internally they
1094 // will be disconnected as appropriate.
1095 kj::Maybe<kj::Own<ByobRequest>> nextPendingByobReadRequest();
1096 
1097 void visitForGc(jsg::GcVisitor& visitor);
1098 
1099 inline kj::StringPtr jsgGetMemoryName() const;
1100 inline size_t jsgGetMemorySelfSize() const;
1101 inline void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
1102 
1103 private:
1104 QueueImpl impl;
1105 
1106 static void handlePush(
1107 jsg::Lock& js, ConsumerImpl::Ready& state, kj::Maybe<QueueImpl&> queue, kj::Rc<Entry> entry);
1108 static void handleRead(jsg::Lock& js,
1109 ConsumerImpl::Ready& state,
1110 ConsumerImpl& consumer,
1111 kj::Maybe<QueueImpl&> queue,
1112 ReadRequest request);
1113 static bool handleMaybeClose(jsg::Lock& js,
1114 ConsumerImpl::Ready& state,
1115 ConsumerImpl& consumer,
1116 kj::Maybe<QueueImpl&> queue);
1117 
1118 friend ConsumerImpl;
1119 friend class Consumer;
1120};
1121 
1122template <typename Self>
1123kj::StringPtr QueueImpl<Self>::jsgGetMemoryName() const {
1124 return "QueueImpl"_kjc;
1125}
1126 
1127template <typename Self>
1128size_t QueueImpl<Self>::jsgGetMemorySelfSize() const {
1129 return sizeof(QueueImpl<Self>);
1130}
1131 
1132template <typename Self>
1133void QueueImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1134 KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) {
1135 tracker.trackField("error", errored.reason);
1136 }
1137}
1138 
1139template <typename Self>
1140kj::StringPtr ConsumerImpl<Self>::jsgGetMemoryName() const {
1141 return "ConsumerImpl"_kjc;
1142}
1143 
1144template <typename Self>
1145size_t ConsumerImpl<Self>::jsgGetMemorySelfSize() const {
1146 return sizeof(ConsumerImpl<Self>);
1147}
1148 
1149template <typename Self>
1150void ConsumerImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1151 KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) {
1152 tracker.trackField("error", errored.reason);
1153 } else KJ_IF_SOME(ready, state.tryGetActiveUnsafe()) {
1154 tracker.trackField("inner", ready);
1155 }
1156}
1157 
1158template <typename Self>
1159kj::StringPtr ConsumerImpl<Self>::Ready::jsgGetMemoryName() const {
1160 return "ConsumerImpl::Ready"_kjc;
1161}
1162 
1163template <typename Self>
1164size_t ConsumerImpl<Self>::Ready::jsgGetMemorySelfSize() const {
1165 return sizeof(Ready);
1166}
1167 
1168template <typename Self>
1169void ConsumerImpl<Self>::Ready::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1170 for (auto& entry: buffer) {
1171 KJ_SWITCH_ONEOF(entry) {
1172 KJ_CASE_ONEOF(c, Close) {
1173 tracker.trackFieldWithSize("pendingClose", sizeof(Close));
1174 }
1175 KJ_CASE_ONEOF(e, QueueEntry) {
1176 tracker.trackField("entry", e);
1177 }
1178 }
1179 }
1180 
1181 for (auto& request: readRequests) {
1182 tracker.trackField("pendingRead", *request);
1183 }
1184}
1185 
1186kj::StringPtr ValueQueue::Consumer::jsgGetMemoryName() const {
1187 return "ValueQueue::Consumer"_kjc;
1188}
1189 
1190size_t ValueQueue::Consumer::jsgGetMemorySelfSize() const {
1191 return sizeof(ValueQueue::Consumer);
1192}
1193 
1194void ValueQueue::Consumer::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1195 tracker.trackField("impl", impl);
1196}
1197 
1198kj::StringPtr ValueQueue::jsgGetMemoryName() const {
1199 return "ValueQueue"_kjc;
1200}
1201 
1202size_t ValueQueue::jsgGetMemorySelfSize() const {
1203 return sizeof(ValueQueue);
1204}
1205 
1206void ValueQueue::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1207 tracker.trackField("impl", impl);
1208}
1209 
1210kj::StringPtr ByteQueue::Consumer::jsgGetMemoryName() const {
1211 return "ByteQueue::Consumer"_kjc;
1212}
1213 
1214size_t ByteQueue::Consumer::jsgGetMemorySelfSize() const {
1215 return sizeof(ByteQueue::Consumer);
1216}
1217 
1218void ByteQueue::Consumer::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1219 tracker.trackField("impl", impl);
1220}
1221 
1222kj::StringPtr ByteQueue::jsgGetMemoryName() const {
1223 return "ByteQueue"_kjc;
1224}
1225 
1226size_t ByteQueue::jsgGetMemorySelfSize() const {
1227 return sizeof(ByteQueue);
1228}
1229 
1230void ByteQueue::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
1231 tracker.trackField("impl", impl);
1232}
1233 
1234} // namespace workerd::api