StormByte-Buffer 2.0.0
C++26 buffer module of the StormByte suite
 
Loading...
Searching...
No Matches
sink.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 <condition_variable>
48#include <cstddef>
49#include <functional>
50#include <map>
51#include <memory>
52#include <mutex>
53#include <set>
54#include <utility>
55#include <vector>
56
57namespace StormByte::Buffer {
58 /**
59 * @class Sink<T>::Implementation
60 * @brief Internal implementation class for Sink.
61 *
62 * Manages the map of key-to-Hopper buckets, thread synchronization,
63 * wiring condition variables, and pop selection algorithms.
64 */
65 template<Type::MoveConstructible T>
66 class Sink<T>::Implementation {
67 public:
68 /**
69 * @brief Constructs the Sink Implementation instance.
70 */
71 Implementation() noexcept
72 : m_rr(0), m_consumer(nullptr), m_closed(false), m_drain(false) {}
73
74 /**
75 * @brief Destructor. Marks Sink closed and wakes any waiting threads on m_wired.
76 */
77 ~Implementation() noexcept {
78 m_closed.store(true, std::memory_order_release);
79 m_wired.notify_all();
80 }
81
82 /**
83 * @brief Enqueues an item into the hopper for key.
84 * @param key Bucket key identifier.
85 * @param item Item to push.
86 */
87 void Push(int key, T item) noexcept {
88 if constexpr (Type::NullablePointer<T>) {
89 if (!item)
90 return;
91 }
92 std::shared_ptr<Hopper<T>> hopper;
93 {
94 std::unique_lock<std::mutex> lock(m_mutex);
95 m_wired.wait(lock, [this, key] {
96 return m_closed.load(std::memory_order_acquire)
97 || m_drain.load(std::memory_order_acquire)
98 || m_buckets.find(key) != m_buckets.end();
99 });
100 if (m_buckets.find(key) == m_buckets.end())
101 return;
102 hopper = m_buckets[key];
103 }
104 if (!hopper)
105 return;
106 hopper->Push(std::move(item));
107 }
108
109 /**
110 * @brief Closes this Sink and returns hoppers it writes, for last-writer Eof.
111 * @param cv Set to the consumer condition variable, if any.
112 * @return Writer hoppers to CloseWriter (empty if already closed).
113 */
114 std::vector<std::shared_ptr<Hopper<T>>> Close(std::condition_variable*& cv) noexcept {
115 std::vector<std::shared_ptr<Hopper<T>>> writers;
116 {
117 std::lock_guard<std::mutex> lock(m_mutex);
118 const bool already = m_closed.exchange(true, std::memory_order_acq_rel);
119 if (!already) {
120 for (auto& hopper : m_order) {
121 if (m_writers.contains(hopper))
122 writers.push_back(hopper);
123 }
124 }
125 m_wired.notify_all();
126 cv = m_consumer.load(std::memory_order_acquire);
127 }
128 return writers;
129 }
130
131 /**
132 * @brief Shares all existing hoppers with consumer.
133 * @param consumer Consumer Sink implementation reference.
134 */
135 void Bind(Implementation& consumer) {
136 std::scoped_lock lock(m_mutex, consumer.m_mutex);
137 const bool closed = m_closed.load(std::memory_order_acquire)
138 || consumer.m_closed.load(std::memory_order_acquire);
139 std::condition_variable* cv = consumer.m_consumer.load(std::memory_order_acquire);
140 for (auto& [key, hopper] : m_buckets) {
141 if (closed)
142 hopper->Eof();
143 if (cv != nullptr)
144 hopper->Notify(*cv);
145 consumer.m_buckets[key] = hopper;
146 }
147 if (closed)
148 consumer.m_closed.store(true, std::memory_order_release);
149 consumer.RebuildOrder();
150 consumer.m_wired.notify_all();
151 m_wired.notify_all();
152 }
153
154 /**
155 * @brief Creates or shares the hopper for key with consumer.
156 * @param key Bucket key.
157 * @param consumer Consumer Sink implementation reference.
158 * @return Hopper when this Sink already held it (extra writer); empty otherwise.
159 */
160 std::shared_ptr<Hopper<T>> Bind(int key, Implementation& consumer) {
161 std::scoped_lock lock(m_mutex, consumer.m_mutex);
162 const bool closed = m_closed.load(std::memory_order_acquire)
163 || consumer.m_closed.load(std::memory_order_acquire);
164 std::condition_variable* cv = consumer.m_consumer.load(std::memory_order_acquire);
165 const bool existed = m_buckets.contains(key);
166 auto hopper = Ensure(key);
167 const bool already_writer = existed && consumer.m_writers.contains(hopper);
168 if (existed)
169 consumer.m_writers.insert(hopper);
170 if (closed)
171 hopper->Eof();
172 if (cv != nullptr)
173 hopper->Notify(*cv);
174 consumer.m_buckets[key] = hopper;
175 if (closed)
176 consumer.m_closed.store(true, std::memory_order_release);
177 consumer.RebuildOrder();
178 consumer.m_wired.notify_all();
179 m_wired.notify_all();
180 return (existed && !already_writer) ? hopper : nullptr;
181 }
182
183 /**
184 * @brief Sets Drain mode.
185 */
186 void Drain() noexcept {
187 m_drain.store(true, std::memory_order_release);
188 m_wired.notify_all();
189 }
190
191 /**
192 * @brief Checks if Drain was set.
193 * @return true if draining.
194 */
195 bool Draining() const noexcept {
196 return m_drain.load(std::memory_order_acquire);
197 }
198
199 /**
200 * @brief Registers condition variable for consumer notifications.
201 * @param consumer Condition variable reference.
202 */
203 void Notify(std::condition_variable& consumer) noexcept {
204 m_consumer.store(&consumer, std::memory_order_release);
205 const auto hoppers = Order();
206 for (auto& hopper : hoppers)
207 hopper->Notify(consumer);
208 }
209
210 /**
211 * @brief Drops the consumer condition variable on this Sink and its hoppers.
212 */
213 void Unnotify() noexcept {
214 m_consumer.store(nullptr, std::memory_order_release);
215 const auto hoppers = Order();
216 for (auto& hopper : hoppers) {
217 if (hopper)
218 hopper->Unnotify();
219 }
220 }
221
222 /**
223 * @brief Snapshot of wired keys in map order.
224 * @return Keys, empty if none.
225 */
226 std::vector<int> Keys() const noexcept {
227 std::lock_guard<std::mutex> lock(m_mutex);
228 std::vector<int> keys;
229 keys.reserve(m_buckets.size());
230 for (const auto& [key, hopper] : m_buckets)
231 keys.push_back(key);
232 return keys;
233 }
234
235 /**
236 * @brief Number of wired hoppers.
237 * @return Bucket count.
238 */
239 StormByte::Size Buckets() const noexcept {
240 std::lock_guard<std::mutex> lock(m_mutex);
241 return StormByte::Size{m_buckets.size()};
242 }
243
244 /**
245 * @brief Whether key is wired.
246 * @param key Bucket key.
247 * @return true if present.
248 */
249 bool Contains(int key) const noexcept {
250 return static_cast<bool>(Bucket(key));
251 }
252
253 /**
254 * @brief Gets capacity of key hopper.
255 * @param key Bucket key.
256 * @return Capacity value.
257 */
258 StormByte::Size Capacity(int key) const noexcept {
259 const auto hopper = Bucket(key);
260 if (!hopper)
261 return StormByte::Size{0};
262 return hopper->Capacity();
263 }
264
265 /**
266 * @brief Sets capacity of key hopper.
267 * @param key Bucket key.
268 * @param capacity New capacity.
269 */
270 void Capacity(int key, StormByte::Size capacity) noexcept {
271 std::lock_guard<std::mutex> lock(m_mutex);
272 auto found = m_buckets.find(key);
273 if (found != m_buckets.end() && found->second)
274 found->second->Capacity(capacity);
275 }
276
277 /**
278 * @brief Gets pending item count of key hopper.
279 * @param key Bucket key.
280 * @return Item count.
281 */
282 StormByte::Size Size(int key) const noexcept {
283 const auto hopper = Bucket(key);
284 if (!hopper)
285 return StormByte::Size{0};
286 return hopper->Size();
287 }
288
289 /**
290 * @brief Checks if key hopper is full.
291 * @param key Bucket key.
292 * @return true if full.
293 */
294 bool Full(int key) const noexcept {
295 const auto hopper = Bucket(key);
296 if (!hopper)
297 return false;
298 return hopper->Full();
299 }
300
301 /**
302 * @brief Whether key hopper has no items.
303 * @param key Bucket key.
304 * @return true if missing or empty.
305 */
306 bool Empty(int key) const noexcept {
307 const auto hopper = Bucket(key);
308 if (!hopper)
309 return true;
310 return hopper->Empty();
311 }
312
313 /**
314 * @brief Whether producers marked Eof on key hopper.
315 * @param key Bucket key.
316 * @return Hopper EoF, or false if missing.
317 */
318 bool EoF(int key) const noexcept {
319 const auto hopper = Bucket(key);
320 if (!hopper)
321 return false;
322 return hopper->EoF();
323 }
324
325 /**
326 * @brief Whether key hopper has an item or is finished.
327 * @param key Bucket key.
328 * @return false if missing.
329 */
330 bool Ready(int key) const noexcept {
331 const auto hopper = Bucket(key);
332 if (!hopper)
333 return false;
334 return !hopper->Empty() || hopper->EoF();
335 }
336
337 /**
338 * @brief Copy of front item of key hopper. Does not dequeue.
339 * @param key Bucket key.
340 * @return Front or default T.
341 */
342 T Front(int key) const noexcept requires Type::CopyConstructible<T> {
343 const auto hopper = Bucket(key);
344 if (!hopper)
345 return T{};
346 return hopper->Front();
347 }
348
349 /**
350 * @brief Pops item using default selection.
351 * @return Popped item or default T.
352 */
353 T Pop() noexcept {
354 return Pop(typename Sink<T>::Select{});
355 }
356
357 /**
358 * @brief Pops item using specified selection function.
359 * @param select Bucket index chooser.
360 * @return Popped item or default T.
361 */
362 T Pop(const typename Sink<T>::Select& select) noexcept {
363 std::vector<std::shared_ptr<Hopper<T>>> hoppers;
364 {
365 std::unique_lock<std::mutex> lock(m_mutex);
366 m_wired.wait(lock, [this] {
367 return m_closed.load(std::memory_order_acquire) || !m_order.empty();
368 });
369 hoppers = m_order;
370 }
371 if (hoppers.empty())
372 return T{};
373 if (hoppers.size() == 1)
374 return hoppers.front()->Pop();
375
376 const std::size_t count = hoppers.size();
377 std::size_t start = 0;
378 if (select)
379 start = static_cast<std::size_t>(select(StormByte::Size{count})) % count;
380 else
381 start = m_rr.fetch_add(1, std::memory_order_relaxed) % count;
382
383 for (std::size_t offset = 0; offset < count; ++offset) {
384 auto& hopper = hoppers[(start + offset) % count];
385 if (!hopper->Empty())
386 return hopper->Pop();
387 }
388 return T{};
389 }
390
391 /**
392 * @brief Pops from one key only.
393 * @param key Bucket key.
394 * @return Item or default T.
395 */
396 T Pop(int key) noexcept {
397 std::shared_ptr<Hopper<T>> hopper;
398 {
399 std::unique_lock<std::mutex> lock(m_mutex);
400 m_wired.wait(lock, [this, key] {
401 return m_closed.load(std::memory_order_acquire)
402 || m_buckets.contains(key);
403 });
404 auto found = m_buckets.find(key);
405 if (found == m_buckets.end() || !found->second)
406 return T{};
407 hopper = found->second;
408 }
409 return hopper->Pop();
410 }
411
412 /**
413 * @brief Checks if Sink is finished.
414 * @return true if closed and all hoppers drained.
415 */
416 bool EoF() const noexcept {
417 const auto hoppers = Order();
418 if (hoppers.empty())
419 return m_closed.load(std::memory_order_acquire);
420 for (const auto& hopper : hoppers) {
421 if (!hopper->EoF())
422 return false;
423 if (!hopper->Empty())
424 return false;
425 }
426 return true;
427 }
428
429 /**
430 * @brief Checks if Pop can return immediately.
431 * @return true if item is ready or EoF reached.
432 */
433 bool Ready() const noexcept {
434 const auto hoppers = Order();
435 if (hoppers.empty())
436 return m_closed.load(std::memory_order_acquire);
437 bool drained = true;
438 for (const auto& hopper : hoppers) {
439 if (!hopper->Empty())
440 return true;
441 if (!hopper->EoF())
442 drained = false;
443 }
444 return drained;
445 }
446
447 private:
448 /**
449 * @brief Ensures hopper for key exists. Caller holds m_mutex.
450 * @param key Bucket key.
451 * @return Shared hopper instance.
452 */
453 std::shared_ptr<Hopper<T>> Ensure(int key) {
454 auto found = m_buckets.find(key);
455 if (found != m_buckets.end())
456 return found->second;
457 auto hopper = std::make_shared<Hopper<T>>();
458 std::condition_variable* cv = m_consumer.load(std::memory_order_acquire);
459 if (cv != nullptr)
460 hopper->Notify(*cv);
461 m_buckets.emplace(key, hopper);
462 m_writers.insert(hopper);
463 RebuildOrder();
464 return hopper;
465 }
466
467 /**
468 * @brief Rebuilds order vector from m_buckets. Caller holds m_mutex.
469 */
470 void RebuildOrder() {
471 m_order.clear();
472 m_order.reserve(m_buckets.size());
473 for (auto& [key, hopper] : m_buckets)
474 m_order.push_back(hopper);
475 }
476
477 /**
478 * @brief Returns snapshot of current hoppers in order.
479 * @return Vector of hoppers.
480 */
481 std::vector<std::shared_ptr<Hopper<T>>> Order() const {
482 std::lock_guard<std::mutex> lock(m_mutex);
483 return m_order;
484 }
485
486 /**
487 * @brief Retrieves hopper for key.
488 * @param key Bucket key.
489 * @return Hopper pointer or nullptr.
490 */
491 std::shared_ptr<Hopper<T>> Bucket(int key) const {
492 std::lock_guard<std::mutex> lock(m_mutex);
493 auto found = m_buckets.find(key);
494 if (found == m_buckets.end())
495 return nullptr;
496 return found->second;
497 }
498
499 mutable std::mutex m_mutex; ///< Guards bucket map and order vector.
500 std::condition_variable m_wired; ///< Waits for bucket binding or closure.
501 std::map<int, std::shared_ptr<Hopper<T>>> m_buckets;///< Map of integer keys to Hopper buckets.
502 std::set<std::shared_ptr<Hopper<T>>> m_writers; ///< Hoppers this Sink writes (CloseWriter on Eof).
503 std::vector<std::shared_ptr<Hopper<T>>> m_order; ///< Order vector of hoppers for Pop.
504 std::atomic<std::size_t> m_rr; ///< Round-robin counter.
505 std::atomic<std::condition_variable*> m_consumer; ///< Registered consumer condition variable.
506 std::atomic<bool> m_closed; ///< Closed flag.
507 std::atomic<bool> m_drain; ///< Drain mode flag.
508 };
509
510 template<Type::MoveConstructible T>
511 Sink<T>::Sink() noexcept
512 : m_io(std::make_unique<Implementation>()) {}
513
514 template<Type::MoveConstructible T>
515 Sink<T>::~Sink() noexcept = default;
516
517 template<Type::MoveConstructible T>
518 void Sink<T>::Push(int key, T item) noexcept {
519 m_io->Push(key, std::move(item));
520 }
521
522 template<Type::MoveConstructible T>
523 void Sink<T>::Eof() noexcept {
524 std::condition_variable* cv = nullptr;
525 const auto writers = m_io->Close(cv);
526 for (const auto& hopper : writers)
527 hopper->CloseWriter();
528 if (cv)
529 cv->notify_all();
530 }
531
532 template<Type::MoveConstructible T>
533 Sink<T>::Lane::Lane(Sink& from, int key) noexcept
534 : m_from(&from), m_key(key) {}
535
536 template<Type::MoveConstructible T>
537 typename Sink<T>::Lane Sink<T>::To(int key) noexcept {
538 return Lane(*this, key);
539 }
540
541 template<Type::MoveConstructible T>
542 Sink<T>& Sink<T>::Lane::operator>>(Sink& dest) noexcept {
543 if (auto extra = m_from->m_io->Bind(m_key, *dest.m_io))
544 extra->AddWriter();
545 return dest;
546 }
547
548 template<Type::MoveConstructible T>
549 Sink<T>& Sink<T>::operator>>(Sink& dest) noexcept {
550 m_io->Bind(*dest.m_io);
551 return dest;
552 }
553
554 template<Type::MoveConstructible T>
555 Sink<T>& Sink<T>::operator<<(Sink& src) noexcept {
556 src >> *this;
557 return *this;
558 }
559
560 template<Type::MoveConstructible T>
561 Sink<T>& Sink<T>::operator<<(Lane lane) noexcept {
562 lane >> *this;
563 return *this;
564 }
565
566 template<Type::MoveConstructible T>
567 void Sink<T>::Drain() noexcept {
568 m_io->Drain();
569 }
570
571 template<Type::MoveConstructible T>
572 bool Sink<T>::Draining() const noexcept {
573 return m_io->Draining();
574 }
575
576 template<Type::MoveConstructible T>
577 void Sink<T>::Notify(std::condition_variable& consumer) noexcept {
578 m_io->Notify(consumer);
579 }
580
581 template<Type::MoveConstructible T>
582 void Sink<T>::Unnotify() noexcept {
583 m_io->Unnotify();
584 }
585
586 template<Type::MoveConstructible T>
587 std::vector<int> Sink<T>::Keys() const noexcept {
588 return m_io->Keys();
589 }
590
591 template<Type::MoveConstructible T>
592 StormByte::Size Sink<T>::Buckets() const noexcept {
593 return m_io->Buckets();
594 }
595
596 template<Type::MoveConstructible T>
597 bool Sink<T>::Contains(int key) const noexcept {
598 return m_io->Contains(key);
599 }
600
601 template<Type::MoveConstructible T>
602 StormByte::Size Sink<T>::Capacity(int key) const noexcept {
603 return m_io->Capacity(key);
604 }
605
606 template<Type::MoveConstructible T>
607 void Sink<T>::Capacity(int key, StormByte::Size capacity) noexcept {
608 m_io->Capacity(key, capacity);
609 }
610
611 template<Type::MoveConstructible T>
612 StormByte::Size Sink<T>::Size(int key) const noexcept {
613 return m_io->Size(key);
614 }
615
616 template<Type::MoveConstructible T>
617 bool Sink<T>::Full(int key) const noexcept {
618 return m_io->Full(key);
619 }
620
621 template<Type::MoveConstructible T>
622 bool Sink<T>::Empty(int key) const noexcept {
623 return m_io->Empty(key);
624 }
625
626 template<Type::MoveConstructible T>
627 bool Sink<T>::EoF(int key) const noexcept {
628 return m_io->EoF(key);
629 }
630
631 template<Type::MoveConstructible T>
632 bool Sink<T>::Ready(int key) const noexcept {
633 return m_io->Ready(key);
634 }
635
636 template<Type::MoveConstructible T>
637 T Sink<T>::Front(int key) const noexcept requires Type::CopyConstructible<T> {
638 return m_io->Front(key);
639 }
640
641 template<Type::MoveConstructible T>
642 T Sink<T>::Pop() noexcept {
643 return m_io->Pop();
644 }
645
646 template<Type::MoveConstructible T>
647 T Sink<T>::Pop(const Select& select) noexcept {
648 return m_io->Pop(select);
649 }
650
651 template<Type::MoveConstructible T>
652 T Sink<T>::Pop(int key) noexcept {
653 return m_io->Pop(key);
654 }
655
656 template<Type::MoveConstructible T>
657 bool Sink<T>::EoF() const noexcept {
658 return m_io->EoF();
659 }
660
661 template<Type::MoveConstructible T>
662 bool Sink<T>::Ready() const noexcept {
663 return m_io->Ready();
664 }
665}