File
Blob: src/workerd/util/batch-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 <kj/debug.h> |
| 8 | #include <kj/vector.h> |
| 9 | |
| 10 | #include <utility> |
| 11 | |
| 12 | namespace workerd { |
| 13 | |
| 14 | using kj::uint; |
| 15 | |
| 16 | // A double-buffered batch queue which enforces an upper bound on buffer growth. |
| 17 | // |
| 18 | // Objects of this type have two buffers -- the push buffer and the pop buffer -- and support |
| 19 | // `push()` and `pop()` operations. `push()` adds elements to the push buffer. `pop()` swaps the |
| 20 | // push and the pop buffers and returns a RAII object which provides a view onto the pop buffer. |
| 21 | // When the RAII object is destroyed, it resets the size and capacity of the pop buffer. |
| 22 | // |
| 23 | // This class is useful when the cost of context switching between producers and consumers is |
| 24 | // high and/or when you must be able to gracefully handle bursts of pushes, such as when |
| 25 | // transferring objects between threads. Note that this class implements no cross-thread |
| 26 | // synchronization itself, but it can become an effective multiple-producer, single-consumer queue |
| 27 | // when wrapped as a `kj::MutexGuarded<BatchQueue<T>>`. |
| 28 | template <typename T> |
| 29 | class BatchQueue { |
| 30 | public: |
| 31 | // `initialCapacity` is the number of elements of type T for which we should allocate space in the |
| 32 | // initial buffers, and any reconstructed buffers. Buffers will be reconstructed if they are |
| 33 | // observed to grow beyond `maxCapacity` after a completed pop operation. |
| 34 | explicit BatchQueue(uint initialCapacity, uint maxCapacity) |
| 35 | : pushBuffer(initialCapacity), |
| 36 | popBuffer(initialCapacity), |
| 37 | initialCapacity(initialCapacity), |
| 38 | maxCapacity(maxCapacity) {} |
| 39 | |
| 40 | // This is the return type of `pop()` (in fact, `pop()` is the only way to construct a non-empty |
| 41 | // Batch). Default-constructible, moveable, and non-copyable. |
| 42 | // |
| 43 | // A Batch can be converted to an ArrayPtr<T>. When a Batch is destroyed, it clears the pop |
| 44 | // buffer and resets the pop buffer capacity to `initialCapacity` if necessary. |
| 45 | class Batch { |
| 46 | public: |
| 47 | Batch() = default; |
| 48 | Batch(Batch&&) = default; |
| 49 | Batch& operator=(Batch&&) = default; |
| 50 | ~Batch() noexcept(false); |
| 51 | KJ_DISALLOW_COPY(Batch); |
| 52 | |
| 53 | operator kj::ArrayPtr<T>() { |
| 54 | return batchQueue.map([](auto& bq) -> kj::ArrayPtr<T> { |
| 55 | return bq.popBuffer; |
| 56 | }).orDefault(nullptr); |
| 57 | } |
| 58 | |
| 59 | kj::ArrayPtr<T> asArrayPtr() { |
| 60 | return *this; |
| 61 | } |
| 62 | |
| 63 | private: |
| 64 | explicit Batch(BatchQueue& batchQueue): batchQueue(batchQueue) {} |
| 65 | friend BatchQueue; |
| 66 | |
| 67 | kj::Maybe<BatchQueue<T>&> batchQueue; |
| 68 | // It's a Maybe so we can support move operations. |
| 69 | }; |
| 70 | |
| 71 | // If a batch is available, swap the buffers and return a Batch object backed by the pop buffer. |
| 72 | // The caller should destroy the Batch object as soon as they are done with it. Destruction will |
| 73 | // clear the pop buffer and, if necessary, reconstruct it to stay under `maxCapacity`. |
| 74 | // |
| 75 | // Throws if `pop()` is called again before the previous Batch object was destroyed. Note |
| 76 | // that this exception is only reliable if the previous `pop()` returned a non-empty Batch. |
| 77 | // |
| 78 | // `pop()` accesses both buffers, so it must be synchronized with `push()` operations across |
| 79 | // threads. Batch objects and `push()` access different buffers, so they require no explicit |
| 80 | // cross-thread synchronization with each other. |
| 81 | Batch pop() { |
| 82 | KJ_REQUIRE(popBuffer.empty(), "pop()'s previous result not yet destroyed."); |
| 83 | |
| 84 | Batch batch; |
| 85 | |
| 86 | if (!pushBuffer.empty()) { |
| 87 | std::swap(pushBuffer, popBuffer); |
| 88 | batch = Batch(*this); |
| 89 | } |
| 90 | |
| 91 | return batch; |
| 92 | } |
| 93 | |
| 94 | // Add an item to the current batch. |
| 95 | template <typename U> |
| 96 | void push(U&& value) { |
| 97 | pushBuffer.add(kj::fwd<U>(value)); |
| 98 | } |
| 99 | |
| 100 | auto empty() const { |
| 101 | return pushBuffer.empty(); |
| 102 | } |
| 103 | auto size() const { |
| 104 | return pushBuffer.size(); |
| 105 | } |
| 106 | |
| 107 | private: |
| 108 | kj::Vector<T> pushBuffer; |
| 109 | kj::Vector<T> popBuffer; |
| 110 | uint initialCapacity; |
| 111 | uint maxCapacity; |
| 112 | }; |
| 113 | |
| 114 | // ======================================================================================= |
| 115 | // Inline implementation details |
| 116 | |
| 117 | template <typename T> |
| 118 | BatchQueue<T>::Batch::~Batch() noexcept(false) { |
| 119 | KJ_IF_SOME(bq, batchQueue) { |
| 120 | bq.popBuffer.clear(); |
| 121 | if (auto capacity = bq.popBuffer.capacity(); capacity > bq.maxCapacity) { |
| 122 | // Reset the queue to avoid letting it grow unbounded. |
| 123 | bq.popBuffer = kj::Vector<T>(bq.initialCapacity); |
| 124 | } |
| 125 | } |
| 126 | } |
| 127 | |
| 128 | } // namespace workerd |