Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
mpsc_queue_example.cc
Go to the documentation of this file.
1
2/*
3 Aleph_w
4
5 Data structures & Algorithms
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
40#include <print_rule.H>
41#include <tpl_mpsc_queue.H>
42
43#include <atomic>
44#include <iostream>
45#include <thread>
46#include <vector>
47
48using namespace Aleph;
49
50namespace
51{
52struct Job
53{
54 int producer;
55 int seq;
56};
57
59{
60 std::cout << "[1] Single-threaded push/try_pop walkthrough\n";
61 print_rule();
62
64 q.push(1);
65 q.push(2);
66 q.emplace(3);
67
68 int out = 0;
69 while (q.try_pop(out))
70 std::cout << "popped " << out << "\n";
71 std::cout << "queue empty: " << std::boolalpha << q.is_empty() << "\n\n";
72}
73
75{
76 std::cout << "[2] Four producers feeding one consumer\n";
77 print_rule();
78
79 constexpr int producer_count = 4;
80 constexpr int items_per_producer = 5000;
81
83 std::atomic<bool> start{false};
84 std::vector<std::thread> producers;
85
86 for (int p = 0; p < producer_count; ++p)
87 producers.emplace_back([&, p]
88 {
89 while (not start.load(std::memory_order_acquire))
90 ; // spin until every producer thread has been created
91 for (int i = 0; i < items_per_producer; ++i)
92 q.push(Job{p, i});
93 });
94
95 std::vector<int> per_producer_count(producer_count, 0);
96 const int total = producer_count * items_per_producer;
97 int consumed = 0;
98
99 start.store(true, std::memory_order_release);
100
101 // The consumer runs on the main thread, matching MpscQueue's
102 // single-consumer contract.
103 while (consumed < total)
104 {
105 Job job{};
106 if (not q.try_pop(job))
107 continue; // transient: a producer's push is mid-flight, or done
108 ++per_producer_count[static_cast<size_t>(job.producer)];
109 ++consumed;
110 }
111
112 for (auto & t : producers)
113 t.join();
114
115 std::cout << "Consumed " << consumed << " jobs from " << producer_count
116 << " producers:\n";
117 for (int p = 0; p < producer_count; ++p)
118 std::cout << " producer " << p << ": " << per_producer_count[static_cast<size_t>(p)]
119 << " jobs (expected " << items_per_producer << ")\n";
120 std::cout << "queue empty after drain: " << std::boolalpha << q.is_empty()
121 << "\n\n";
122}
123} // namespace
124
125int main()
126{
127 std::cout << "\n=== Aleph::MpscQueue: multi-producer/single-consumer queue ===\n\n";
128
131
132 std::cout << "Done.\n";
133 return 0;
134}
size_t size_t int32_t * out
Definition ca-c-api.h:120
Unbounded lock-free multi-producer/single-consumer queue.
void emplace(Args &&... args)
Construct a new element in place at the back of the queue.
void push(const T &value)
Push a copy of value onto the queue.
bool is_empty() const noexcept
Advisory check for whether the queue currently has no elements.
bool try_pop(T &out)
Attempt to pop the front element into out.
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
int main()
Main namespace for Aleph-w library functions.
Definition ah-arena.H:89
void print_rule()
Prints a horizontal rule for example output separation.
Definition print_rule.H:39
std::ostream & join(const C &c, const std::string &sep, std::ostream &out)
Join elements of an Aleph-style container into a stream.
Unbounded lock-free multi-producer/single-consumer queue (Aleph::MpscQueue).