StormByte-Buffer 2.0.0
C++26 buffer module of the StormByte suite
 
Loading...
Searching...
No Matches
hopper.txx
Go to the documentation of this file.
1/*
2 * Copyright (C) 2024-2026 David C. Manuelda (StormBytePP)
3 *
4 * This file is part of StormByte-Buffer.
5 *
6 * StormByte-Buffer original source is dual-licensed:
7 *
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)
12 * any later version.
13 *
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>).
18 *
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
24 * license.
25 *
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.
29 *
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.
34 *
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>.
38 *
39 * SPDX-License-Identifier: LGPL-3.0-or-later OR LicenseRef-StormByte-Commercial
40 */
41
42#pragma once
43
44#include <StormByte/type_traits.hxx>
45
46#include <atomic>
47#include <concepts>
48#include <condition_variable>
49#include <cstddef>
50#include <memory>
51#include <mutex>
52#include <queue>
53#include <utility>
54
55namespace StormByte::Buffer {
56 /**
57 * @class Hopper<T>::Implementation
58 * @brief Internal implementation of Hopper queue details.
59 *
60 * Handles mutex-protected queue operations, atomic capacity settings,
61 * condition variable notifications, and EoF flags.
62 */
63 template<Type::MoveConstructible T>
64 class Hopper<T>::Implementation {
65 public:
66 /**
67 * @brief Constructs an unbounded Implementation instance.
68 */
69 Implementation() noexcept
70 : m_eof(false), m_wake(nullptr), m_cap(0), m_writers(1) {}
71
72 /**
73 * @brief Constructs a bounded Implementation instance.
74 * @param capacity Maximum items allowed.
75 */
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) {}
78
79 /**
80 * @brief Destructor. Marks EoF and wakes waiting producers.
81 */
82 ~Implementation() noexcept {
83 m_eof.store(true, std::memory_order_release);
84 m_space.notify_all();
85 }
86
87 /**
88 * @brief Gets capacity ceiling.
89 * @return Capacity value.
90 */
91 StormByte::Size Capacity() const noexcept {
92 return StormByte::Size{m_cap.load(std::memory_order_acquire)};
93 }
94
95 /**
96 * @brief Sets capacity ceiling.
97 * @param capacity New capacity value.
98 */
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();
102 }
103
104 /**
105 * @brief Gets item count in queue.
106 * @return Count of items.
107 */
108 StormByte::Size Size() const noexcept {
109 std::lock_guard<std::mutex> lock(m_mutex);
110 return StormByte::Size{m_items.size()};
111 }
112
113 /**
114 * @brief Checks if bounded bucket is full.
115 * @return true if capacity > 0 and size >= capacity.
116 */
117 bool Full() const noexcept {
118 const std::size_t cap = m_cap.load(std::memory_order_acquire);
119 if (cap == 0)
120 return false;
121 return Size() >= cap;
122 }
123
124 /**
125 * @brief Live writer count.
126 * @return Writers still open.
127 */
128 unsigned Writers() const noexcept {
129 return m_writers.load(std::memory_order_acquire);
130 }
131
132 /**
133 * @brief Enqueues an item, waiting if full.
134 * @param item Item to enqueue.
135 */
136 void Push(T item) noexcept {
137 if constexpr (Type::NullablePointer<T>) {
138 if (!item)
139 return;
140 }
141 {
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);
145 return cap == 0
146 || m_items.size() < cap
147 || m_eof.load(std::memory_order_acquire);
148 });
149 if (m_eof.load(std::memory_order_acquire))
150 return;
151 m_items.push(std::move(item));
152 }
153 SignalConsumer();
154 }
155
156 /**
157 * @brief Signals end of production.
158 */
159 void Eof() noexcept {
160 m_eof.store(true, std::memory_order_release);
161 SignalConsumer();
162 m_space.notify_all();
163 }
164
165 /**
166 * @brief Registers an extra writer.
167 */
168 void AddWriter() noexcept {
169 m_writers.fetch_add(1, std::memory_order_acq_rel);
170 }
171
172 /**
173 * @brief Releases one writer. Last writer force-closes.
174 */
175 void CloseWriter() noexcept {
176 unsigned prev = m_writers.load(std::memory_order_acquire);
177 while (prev > 0) {
178 if (m_writers.compare_exchange_weak(prev, prev - 1,
179 std::memory_order_acq_rel, std::memory_order_acquire)) {
180 if (prev == 1)
181 Eof();
182 return;
183 }
184 }
185 }
186
187 /**
188 * @brief Pops next item from queue without waiting.
189 * @return Next item, or default T if empty.
190 */
191 T Pop() noexcept {
192 T item{};
193 {
194 std::lock_guard<std::mutex> lock(m_mutex);
195 if (m_items.empty())
196 return T{};
197 item = std::move(m_items.front());
198 m_items.pop();
199 }
200 m_space.notify_one();
201 return item;
202 }
203
204 /**
205 * @brief Copy of front item. Does not dequeue.
206 * @return Front or default T.
207 */
208 T Front() const noexcept requires std::copy_constructible<T> {
209 std::lock_guard<std::mutex> lock(m_mutex);
210 if (m_items.empty())
211 return T{};
212 return m_items.front();
213 }
214
215 /**
216 * @brief Checks if Eof was signaled.
217 * @return true if Eof set.
218 */
219 bool EoF() const noexcept {
220 return m_eof.load(std::memory_order_acquire);
221 }
222
223 /**
224 * @brief Checks if queue is empty.
225 * @return true if empty.
226 */
227 bool Empty() const noexcept {
228 std::lock_guard<std::mutex> lock(m_mutex);
229 return m_items.empty();
230 }
231
232 /**
233 * @brief Item ready or production finished.
234 * @return true if !Empty() or EoF().
235 */
236 bool Ready() const noexcept {
237 return !Empty() || EoF();
238 }
239
240 void Notify(std::condition_variable& wake) noexcept {
241 m_wake.store(&wake, std::memory_order_release);
242 }
243
244 void Unnotify() noexcept {
245 m_wake.store(nullptr, std::memory_order_release);
246 }
247
248 private:
249 /**
250 * @brief Notifies registered consumer condition variable if set.
251 */
252 void SignalConsumer() noexcept {
253 std::condition_variable* wake = m_wake.load(std::memory_order_acquire);
254 if (wake == nullptr)
255 return;
256 wake->notify_one();
257 }
258
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.
266 };
267
268 template<Type::MoveConstructible T>
269 Hopper<T>::Hopper() noexcept
270 : m_io(std::make_unique<Implementation>()) {}
271
272 template<Type::MoveConstructible T>
273 Hopper<T>::Hopper(StormByte::Size capacity) noexcept
274 : m_io(std::make_unique<Implementation>(capacity)) {}
275
276 template<Type::MoveConstructible T>
277 Hopper<T>::~Hopper() noexcept = default;
278
279 template<Type::MoveConstructible T>
280 StormByte::Size Hopper<T>::Capacity() const noexcept {
281 return m_io->Capacity();
282 }
283
284 template<Type::MoveConstructible T>
285 void Hopper<T>::Capacity(StormByte::Size capacity) noexcept {
286 m_io->Capacity(capacity);
287 }
288
289 template<Type::MoveConstructible T>
290 StormByte::Size Hopper<T>::Size() const noexcept {
291 return m_io->Size();
292 }
293
294 template<Type::MoveConstructible T>
295 bool Hopper<T>::Full() const noexcept {
296 return m_io->Full();
297 }
298
299 template<Type::MoveConstructible T>
300 unsigned Hopper<T>::Writers() const noexcept {
301 return m_io->Writers();
302 }
303
304 template<Type::MoveConstructible T>
305 void Hopper<T>::Push(T item) noexcept {
306 m_io->Push(std::move(item));
307 }
308
309 template<Type::MoveConstructible T>
310 Hopper<T>& Hopper<T>::operator<<(T item) noexcept {
311 Push(std::move(item));
312 return *this;
313 }
314
315 template<Type::MoveConstructible T>
316 void Hopper<T>::Eof() noexcept {
317 m_io->Eof();
318 }
319
320 template<Type::MoveConstructible T>
321 void Hopper<T>::AddWriter() noexcept {
322 m_io->AddWriter();
323 }
324
325 template<Type::MoveConstructible T>
326 void Hopper<T>::CloseWriter() noexcept {
327 m_io->CloseWriter();
328 }
329
330 template<Type::MoveConstructible T>
331 T Hopper<T>::Pop() noexcept {
332 return m_io->Pop();
333 }
334
335 template<Type::MoveConstructible T>
336 T Hopper<T>::Front() const noexcept requires std::copy_constructible<T> {
337 return m_io->Front();
338 }
339
340 template<Type::MoveConstructible T>
341 Hopper<T>& Hopper<T>::operator>>(T& item) noexcept {
342 item = Pop();
343 return *this;
344 }
345
346 template<Type::MoveConstructible T>
347 bool Hopper<T>::EoF() const noexcept {
348 return m_io->EoF();
349 }
350
351 template<Type::MoveConstructible T>
352 bool Hopper<T>::Empty() const noexcept {
353 return m_io->Empty();
354 }
355
356 template<Type::MoveConstructible T>
357 bool Hopper<T>::Ready() const noexcept {
358 return m_io->Ready();
359 }
360
361 template<Type::MoveConstructible T>
362 void Hopper<T>::Notify(std::condition_variable& wake) noexcept {
363 m_io->Notify(wake);
364 }
365
366 template<Type::MoveConstructible T>
367 void Hopper<T>::Unnotify() noexcept {
368 m_io->Unnotify();
369 }
370}