File
Blob: src/workerd/api/streams/queue.h
| 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 | |
| 15 | namespace 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 | |
| 142 | template <typename Self> |
| 143 | class ConsumerImpl; |
| 144 | |
| 145 | template <typename Self> |
| 146 | class QueueImpl; |
| 147 | |
| 148 | // DrainingReadResult is defined in common.h |
| 149 | |
| 150 | // Provides the underlying implementation shared by ByteQueue and ValueQueue. |
| 151 | template <typename Self> |
| 152 | class 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 |
| 335 | template <typename Self> |
| 336 | class 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 | |
| 720 | class 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 | |
| 876 | class 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 | |
| 1122 | template <typename Self> |
| 1123 | kj::StringPtr QueueImpl<Self>::jsgGetMemoryName() const { |
| 1124 | return "QueueImpl"_kjc; |
| 1125 | } |
| 1126 | |
| 1127 | template <typename Self> |
| 1128 | size_t QueueImpl<Self>::jsgGetMemorySelfSize() const { |
| 1129 | return sizeof(QueueImpl<Self>); |
| 1130 | } |
| 1131 | |
| 1132 | template <typename Self> |
| 1133 | void QueueImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 1134 | KJ_IF_SOME(errored, state.tryGetErrorUnsafe()) { |
| 1135 | tracker.trackField("error", errored.reason); |
| 1136 | } |
| 1137 | } |
| 1138 | |
| 1139 | template <typename Self> |
| 1140 | kj::StringPtr ConsumerImpl<Self>::jsgGetMemoryName() const { |
| 1141 | return "ConsumerImpl"_kjc; |
| 1142 | } |
| 1143 | |
| 1144 | template <typename Self> |
| 1145 | size_t ConsumerImpl<Self>::jsgGetMemorySelfSize() const { |
| 1146 | return sizeof(ConsumerImpl<Self>); |
| 1147 | } |
| 1148 | |
| 1149 | template <typename Self> |
| 1150 | void 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 | |
| 1158 | template <typename Self> |
| 1159 | kj::StringPtr ConsumerImpl<Self>::Ready::jsgGetMemoryName() const { |
| 1160 | return "ConsumerImpl::Ready"_kjc; |
| 1161 | } |
| 1162 | |
| 1163 | template <typename Self> |
| 1164 | size_t ConsumerImpl<Self>::Ready::jsgGetMemorySelfSize() const { |
| 1165 | return sizeof(Ready); |
| 1166 | } |
| 1167 | |
| 1168 | template <typename Self> |
| 1169 | void 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 | |
| 1186 | kj::StringPtr ValueQueue::Consumer::jsgGetMemoryName() const { |
| 1187 | return "ValueQueue::Consumer"_kjc; |
| 1188 | } |
| 1189 | |
| 1190 | size_t ValueQueue::Consumer::jsgGetMemorySelfSize() const { |
| 1191 | return sizeof(ValueQueue::Consumer); |
| 1192 | } |
| 1193 | |
| 1194 | void ValueQueue::Consumer::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 1195 | tracker.trackField("impl", impl); |
| 1196 | } |
| 1197 | |
| 1198 | kj::StringPtr ValueQueue::jsgGetMemoryName() const { |
| 1199 | return "ValueQueue"_kjc; |
| 1200 | } |
| 1201 | |
| 1202 | size_t ValueQueue::jsgGetMemorySelfSize() const { |
| 1203 | return sizeof(ValueQueue); |
| 1204 | } |
| 1205 | |
| 1206 | void ValueQueue::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 1207 | tracker.trackField("impl", impl); |
| 1208 | } |
| 1209 | |
| 1210 | kj::StringPtr ByteQueue::Consumer::jsgGetMemoryName() const { |
| 1211 | return "ByteQueue::Consumer"_kjc; |
| 1212 | } |
| 1213 | |
| 1214 | size_t ByteQueue::Consumer::jsgGetMemorySelfSize() const { |
| 1215 | return sizeof(ByteQueue::Consumer); |
| 1216 | } |
| 1217 | |
| 1218 | void ByteQueue::Consumer::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 1219 | tracker.trackField("impl", impl); |
| 1220 | } |
| 1221 | |
| 1222 | kj::StringPtr ByteQueue::jsgGetMemoryName() const { |
| 1223 | return "ByteQueue"_kjc; |
| 1224 | } |
| 1225 | |
| 1226 | size_t ByteQueue::jsgGetMemorySelfSize() const { |
| 1227 | return sizeof(ByteQueue); |
| 1228 | } |
| 1229 | |
| 1230 | void ByteQueue::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 1231 | tracker.trackField("impl", impl); |
| 1232 | } |
| 1233 | |
| 1234 | } // namespace workerd::api |