2 * Copyright (C) 2024-2026 David C. Manuelda (StormBytePP)
4 * This file is part of StormByte-Buffer.
6 * StormByte-Buffer original source is dual-licensed:
8 * 1. GNU Lesser General Public License v3.0 (or later)
9 * You may redistribute and/or modify this file under the terms of the
10 * GNU Lesser General Public License as published by the Free Software
11 * Foundation, either version 3 of the License, or (at your option)
14 * 2. Commercial license
15 * Alternatively, this file may be used under the terms of a commercial
16 * license agreement with the copyright holder
17 * (David C. Manuelda <StormByte@gmail.com>).
19 * Both licenses apply only to original StormByte-Buffer source in this
20 * repository. They do not cover other StormByte modules or any third-party
21 * material shipped with this repository (including everything under
22 * thirdparty/, and in particular the bundled StormByte-Logger tree and
23 * the rest of the StormByte suite it vendors), which remains under its own
26 * Neither license grants any patent rights. Any patent licenses required
27 * to use this software or third-party components must be obtained separately
28 * from the patent holders.
30 * StormByte-Buffer is distributed in the hope that it will be useful,
31 * but WITHOUT ANY WARRANTY; without even the implied warranty of
32 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
33 * GNU Lesser General Public License for more details.
35 * You should have received a copy of the GNU Lesser General Public License
36 * version 3 along with StormByte-Buffer. If not, see
37 * <https://www.gnu.org/licenses/lgpl-3.0.html>.
39 * SPDX-License-Identifier: LGPL-3.0-or-later OR LicenseRef-StormByte-Commercial
44#include <StormByte/type_traits.hxx>
48#include <condition_variable>
55namespace StormByte::Buffer {
57 * @class Hopper<T>::Implementation
58 * @brief Internal implementation of Hopper queue details.
60 * Handles mutex-protected queue operations, atomic capacity settings,
61 * condition variable notifications, and EoF flags.
63 template<Type::MoveConstructible T>
64 class Hopper<T>::Implementation {
67 * @brief Constructs an unbounded Implementation instance.
69 Implementation() noexcept
70 : m_eof(false), m_wake(nullptr), m_cap(0), m_writers(1) {}
73 * @brief Constructs a bounded Implementation instance.
74 * @param capacity Maximum items allowed.
76 explicit Implementation(StormByte::Size capacity) noexcept
77 : m_eof(false), m_wake(nullptr), m_cap(static_cast<std::size_t>(capacity)), m_writers(1) {}
80 * @brief Destructor. Marks EoF and wakes waiting producers.
82 ~Implementation() noexcept {
83 m_eof.store(true, std::memory_order_release);
88 * @brief Gets capacity ceiling.
89 * @return Capacity value.
91 StormByte::Size Capacity() const noexcept {
92 return StormByte::Size{m_cap.load(std::memory_order_acquire)};
96 * @brief Sets capacity ceiling.
97 * @param capacity New capacity value.
99 void Capacity(StormByte::Size capacity) noexcept {
100 m_cap.store(static_cast<std::size_t>(capacity), std::memory_order_release);
101 m_space.notify_all();
105 * @brief Gets item count in queue.
106 * @return Count of items.
108 StormByte::Size Size() const noexcept {
109 std::lock_guard<std::mutex> lock(m_mutex);
110 return StormByte::Size{m_items.size()};
114 * @brief Checks if bounded bucket is full.
115 * @return true if capacity > 0 and size >= capacity.
117 bool Full() const noexcept {
118 const std::size_t cap = m_cap.load(std::memory_order_acquire);
121 return Size() >= cap;
125 * @brief Live writer count.
126 * @return Writers still open.
128 unsigned Writers() const noexcept {
129 return m_writers.load(std::memory_order_acquire);
133 * @brief Enqueues an item, waiting if full.
134 * @param item Item to enqueue.
136 void Push(T item) noexcept {
137 if constexpr (Type::NullablePointer<T>) {
142 std::unique_lock<std::mutex> lock(m_mutex);
143 m_space.wait(lock, [this]() {
144 const std::size_t cap = m_cap.load(std::memory_order_acquire);
146 || m_items.size() < cap
147 || m_eof.load(std::memory_order_acquire);
149 if (m_eof.load(std::memory_order_acquire))
151 m_items.push(std::move(item));
157 * @brief Signals end of production.
159 void Eof() noexcept {
160 m_eof.store(true, std::memory_order_release);
162 m_space.notify_all();
166 * @brief Registers an extra writer.
168 void AddWriter() noexcept {
169 m_writers.fetch_add(1, std::memory_order_acq_rel);
173 * @brief Releases one writer. Last writer force-closes.
175 void CloseWriter() noexcept {
176 unsigned prev = m_writers.load(std::memory_order_acquire);
178 if (m_writers.compare_exchange_weak(prev, prev - 1,
179 std::memory_order_acq_rel, std::memory_order_acquire)) {
188 * @brief Pops next item from queue without waiting.
189 * @return Next item, or default T if empty.
194 std::lock_guard<std::mutex> lock(m_mutex);
197 item = std::move(m_items.front());
200 m_space.notify_one();
205 * @brief Copy of front item. Does not dequeue.
206 * @return Front or default T.
208 T Front() const noexcept requires std::copy_constructible<T> {
209 std::lock_guard<std::mutex> lock(m_mutex);
212 return m_items.front();
216 * @brief Checks if Eof was signaled.
217 * @return true if Eof set.
219 bool EoF() const noexcept {
220 return m_eof.load(std::memory_order_acquire);
224 * @brief Checks if queue is empty.
225 * @return true if empty.
227 bool Empty() const noexcept {
228 std::lock_guard<std::mutex> lock(m_mutex);
229 return m_items.empty();
233 * @brief Item ready or production finished.
234 * @return true if !Empty() or EoF().
236 bool Ready() const noexcept {
237 return !Empty() || EoF();
240 void Notify(std::condition_variable& wake) noexcept {
241 m_wake.store(&wake, std::memory_order_release);
244 void Unnotify() noexcept {
245 m_wake.store(nullptr, std::memory_order_release);
250 * @brief Notifies registered consumer condition variable if set.
252 void SignalConsumer() noexcept {
253 std::condition_variable* wake = m_wake.load(std::memory_order_acquire);
259 mutable std::mutex m_mutex; ///< Guards queue access.
260 std::condition_variable m_space; ///< Producer wait condition when full.
261 std::queue<T> m_items; ///< Queue of stored items.
262 std::atomic<bool> m_eof; ///< End of production flag.
263 std::atomic<std::condition_variable*> m_wake; ///< Consumer condition variable.
264 std::atomic<std::size_t> m_cap; ///< Capacity ceiling (0 = unbounded).
265 std::atomic<unsigned> m_writers; ///< Live writers; last CloseWriter Eofs.
268 template<Type::MoveConstructible T>
269 Hopper<T>::Hopper() noexcept
270 : m_io(std::make_unique<Implementation>()) {}
272 template<Type::MoveConstructible T>
273 Hopper<T>::Hopper(StormByte::Size capacity) noexcept
274 : m_io(std::make_unique<Implementation>(capacity)) {}
276 template<Type::MoveConstructible T>
277 Hopper<T>::~Hopper() noexcept = default;
279 template<Type::MoveConstructible T>
280 StormByte::Size Hopper<T>::Capacity() const noexcept {
281 return m_io->Capacity();
284 template<Type::MoveConstructible T>
285 void Hopper<T>::Capacity(StormByte::Size capacity) noexcept {
286 m_io->Capacity(capacity);
289 template<Type::MoveConstructible T>
290 StormByte::Size Hopper<T>::Size() const noexcept {
294 template<Type::MoveConstructible T>
295 bool Hopper<T>::Full() const noexcept {
299 template<Type::MoveConstructible T>
300 unsigned Hopper<T>::Writers() const noexcept {
301 return m_io->Writers();
304 template<Type::MoveConstructible T>
305 void Hopper<T>::Push(T item) noexcept {
306 m_io->Push(std::move(item));
309 template<Type::MoveConstructible T>
310 Hopper<T>& Hopper<T>::operator<<(T item) noexcept {
311 Push(std::move(item));
315 template<Type::MoveConstructible T>
316 void Hopper<T>::Eof() noexcept {
320 template<Type::MoveConstructible T>
321 void Hopper<T>::AddWriter() noexcept {
325 template<Type::MoveConstructible T>
326 void Hopper<T>::CloseWriter() noexcept {
330 template<Type::MoveConstructible T>
331 T Hopper<T>::Pop() noexcept {
335 template<Type::MoveConstructible T>
336 T Hopper<T>::Front() const noexcept requires std::copy_constructible<T> {
337 return m_io->Front();
340 template<Type::MoveConstructible T>
341 Hopper<T>& Hopper<T>::operator>>(T& item) noexcept {
346 template<Type::MoveConstructible T>
347 bool Hopper<T>::EoF() const noexcept {
351 template<Type::MoveConstructible T>
352 bool Hopper<T>::Empty() const noexcept {
353 return m_io->Empty();
356 template<Type::MoveConstructible T>
357 bool Hopper<T>::Ready() const noexcept {
358 return m_io->Ready();
361 template<Type::MoveConstructible T>
362 void Hopper<T>::Notify(std::condition_variable& wake) noexcept {
366 template<Type::MoveConstructible T>
367 void Hopper<T>::Unnotify() noexcept {