DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
IterableQueueModel.hxx
Go to the documentation of this file.
1// Declarations for IterableQueueModel
2
3namespace dunedaq {
4namespace datahandlinglibs {
5
6// Free allocated memory that is different for alignment strategies and allocation policies
7template<class T>
8void
10{
11 // We need to destruct anything that may still exist in our queue.
12 // (No real synchronization needed at destructor time: only one
13 // thread can be doing this.)
14 if (!std::is_trivially_destructible<T>::value) {
15 std::size_t readIndex = readIndex_;
16 std::size_t endIndex = writeIndex_;
17 while (readIndex != endIndex) {
18 records_[readIndex].~T();
19 if (++readIndex == size_) { // NOLINT(runtime/increment_decrement)
20 readIndex = 0;
21 }
22 }
23 }
24 // Different allocators require custom free functions
26 _mm_free(records_);
27 } else if (numa_aware_) {
28#ifdef WITH_LIBNUMA_SUPPORT
29 numa_free(records_, sizeof(T) * size_);
30#endif
31 } else {
32 std::free(records_);
33 }
34}
35
36// Allocate memory based on different alignment strategies and allocation policies
37template<class T>
38void
40 bool numa_aware,
41 uint8_t numa_node, // NOLINT (build/unsigned)
42 bool intrinsic_allocator,
43 std::size_t alignment_size)
44{
45 assert(size >= 2);
46 // TODO: check for valid alignment sizes! | July-21-2021 | Roland Sipos | rsipos@cern.ch
47
48 if (numa_aware &&
49 numa_node < 8) { // numa allocator from libnuma; we get "numa_node >= 0" for free, given its datatype
50#ifdef WITH_LIBNUMA_SUPPORT
51 numa_set_preferred((unsigned)numa_node); // https://linux.die.net/man/3/numa_set_preferred
52#ifdef WITH_LIBNUMA_BIND_POLICY
53 numa_set_bind_policy(WITH_LIBNUMA_BIND_POLICY); // https://linux.die.net/man/3/numa_set_bind_policy
54#endif
55#ifdef WITH_LIBNUMA_STRICT_POLICY
56 numa_set_strict(WITH_LIBNUMA_STRICT_POLICY); // https://linux.die.net/man/3/numa_set_strict
57#endif
58 records_ = static_cast<T*>(numa_alloc_onnode(sizeof(T) * size, numa_node));
59#else
61 "NUMA allocation was requested but program was built without USE_LIBNUMA");
62#endif
63 } else if (intrinsic_allocator && alignment_size > 0) { // _mm allocator
64 records_ = static_cast<T*>(_mm_malloc(sizeof(T) * size, alignment_size));
65 } else if (!intrinsic_allocator && alignment_size > 0) { // std aligned allocator
66 records_ = static_cast<T*>(std::aligned_alloc(alignment_size, sizeof(T) * size));
67 } else if (!numa_aware && !intrinsic_allocator && alignment_size == 0) {
68 // Standard allocator
69 records_ = static_cast<T*>(std::malloc(sizeof(T) * size));
70
71 } else {
72 // Let it fail, as expected combination might be invalid
73 // records_ = static_cast<T*>(std::malloc(sizeof(T) * size_);
74 }
75
76 size_ = size;
77 numa_aware_ = numa_aware;
78 numa_node_ = numa_node;
79 intrinsic_allocator_ = intrinsic_allocator;
80 alignment_size_ = alignment_size;
81}
82
83template<class T>
84void
86{
87 // Wait until LB issues ready
88 std::unique_lock lk(prefill_mutex_);
89 prefill_cv_.wait(lk, [this] { return prefill_ready_; });
90
91 // After wait, we are ready to force page-fault
92 for (size_t i = 0; i < size_ - 1; ++i) {
93 T element = T();
94 write_(std::move(element));
95 }
96 flush();
97
98 // Preallocation done
99 prefill_done_ = true;
100
101 // Manual unlock is done before notify: avoid waking up the waiting thread only to block again.
102 lk.unlock();
103 prefill_cv_.notify_one();
104}
105
106template<class T>
107void
109{
110 // Local prefiller thread
111 std::thread prefill_thread(&IterableQueueModel<T>::prefill_task, this);
112
113 // Tweak prefiller thread
114 char tname[16];
115 snprintf(tname, 16, "%s-%d", prefiller_name_.c_str(), numa_node_);
116 auto handle = prefill_thread.native_handle();
117 pthread_setname_np(handle, tname);
118
119#ifdef WITH_LIBNUMA_SUPPORT
120 cpu_set_t affinitymask;
121 CPU_ZERO(&affinitymask);
122 struct bitmask* nodecpumask = numa_allocate_cpumask();
123 int ret = 0;
124 // Get NODE CPU mask
125 ret = numa_node_to_cpus(numa_node_, nodecpumask);
126 assert(ret == 0);
127 // Apply corresponding NODE CPUs to affinity mask
128 for (int i = 0; i < numa_num_configured_cpus(); ++i) {
129 if (numa_bitmask_isbitset(nodecpumask, i)) {
130 CPU_SET(i, &affinitymask);
131 }
132 }
133 ret = pthread_setaffinity_np(handle, sizeof(cpu_set_t), &affinitymask);
134 assert(ret == 0);
135 numa_free_cpumask(nodecpumask);
136#endif
137
138 // Trigger prefiller thread
139 {
140 std::lock_guard lk(prefill_mutex_);
141 prefill_ready_ = true;
142 }
143 prefill_cv_.notify_one();
144 // Wait for prefiller thread to finish
145 {
146 std::unique_lock lk(prefill_mutex_);
147 prefill_cv_.wait(lk, [this] { return prefill_done_; });
148 }
149 // Join with prefiller thread
150 prefill_thread.join();
151}
152
153// Write element into the queue
154template<class T>
155bool
157{
158 auto const currentWrite = writeIndex_.load(std::memory_order_relaxed);
159 auto nextRecord = currentWrite + 1;
160 if (nextRecord == size_) {
161 nextRecord = 0;
162 }
163
164 if (nextRecord != readIndex_.load(std::memory_order_acquire)) {
165 new (&records_[currentWrite]) T(std::move(record));
166 writeIndex_.store(nextRecord, std::memory_order_release);
167 return true;
168 }
169
170 // queue is full
171 ++overflow_ctr;
172 return false;
173}
174
175// Read element from a queue (move or copy the value at the front of the queue to given variable)
176template<class T>
177bool
179{
180 auto const currentRead = readIndex_.load(std::memory_order_relaxed);
181 if (currentRead == writeIndex_.load(std::memory_order_acquire)) {
182 // queue is empty
183 return false;
184 }
185
186 auto nextRecord = currentRead + 1;
187 if (nextRecord == size_) {
188 nextRecord = 0;
189 }
190 record = std::move(records_[currentRead]);
191 records_[currentRead].~T();
192 readIndex_.store(nextRecord, std::memory_order_release);
193 return true;
194}
195
196// Pop element on front of queue
197template<class T>
198void
200{
201 auto const currentRead = readIndex_.load(std::memory_order_relaxed);
202 assert(currentRead != writeIndex_.load(std::memory_order_acquire));
203
204 auto nextRecord = currentRead + 1;
205 if (nextRecord == size_) {
206 nextRecord = 0;
207 }
208
209 records_[currentRead].~T();
210 readIndex_.store(nextRecord, std::memory_order_release);
211}
212
213// Pop number of elements (X) from the front of the queue
214template<class T>
215void
217{
218 for (std::size_t i = 0; i < x; i++) {
219 popFront();
220 }
221}
222
223// Returns true if the queue is empty
224template<class T>
225bool
227{
228 return readIndex_.load(std::memory_order_acquire) == writeIndex_.load(std::memory_order_acquire);
229}
230
231// Returns true if write index reached read index
232template<class T>
233bool
235{
236 auto nextRecord = writeIndex_.load(std::memory_order_acquire) + 1;
237 if (nextRecord == size_) {
238 nextRecord = 0;
239 }
240 if (nextRecord != readIndex_.load(std::memory_order_acquire)) {
241 return false;
242 }
243 // queue is full
244 return true;
245}
246
247// Returns a good-enough guess on current occupancy:
248// * If called by consumer, then true size may be more (because producer may
249// be adding items concurrently).
250// * If called by producer, then true size may be less (because consumer may
251// be removing items concurrently).
252// * It is undefined to call this from any other thread.
253template<class T>
254std::size_t
256{
257 int ret = static_cast<int>(writeIndex_.load(std::memory_order_acquire)) -
258 static_cast<int>(readIndex_.load(std::memory_order_acquire));
259 if (ret < 0) {
260 ret += static_cast<int>(size_);
261 }
262 return static_cast<std::size_t>(ret);
263}
264
265// Gives a pointer to the current read index
266template<class T>
267const T*
269{
270 auto const currentRead = readIndex_.load(std::memory_order_relaxed);
271 if (currentRead == writeIndex_.load(std::memory_order_acquire)) {
272 return nullptr;
273 }
274 return &records_[currentRead];
275}
276
277// Gives a pointer to the current write index
278template<class T>
279const T*
281{
282 auto const currentWrite = writeIndex_.load(std::memory_order_relaxed);
283 if (currentWrite == readIndex_.load(std::memory_order_acquire)) {
284 return nullptr;
285 }
286 int currentLast = currentWrite;
287 if (currentLast == 0) {
288 currentLast = size_ - 1;
289 } else {
290 currentLast--;
291 }
292 return &records_[currentLast];
293}
294
295// Configures the model
296template<class T>
297void
299{
300 assert(cfg->get_size() >= 2);
301 free_memory();
302
304 cfg->get_numa_aware(),
305 cfg->get_numa_node(),
307 cfg->get_alignment_size());
308 readIndex_ = 0;
309 writeIndex_ = 0;
310
311 if (!records_) {
312 throw std::bad_alloc();
313 }
314
315 if (cfg->get_preallocation()) {
317 }
318}
319
320// Unconfigures the model
321template<class T>
322void
323IterableQueueModel<T>::scrap(const appfwk::DAQModule::CommandData_t& /*cfg*/)
324{
325 free_memory();
326 numa_aware_ = false;
327 numa_node_ = 0;
328 intrinsic_allocator_ = false;
329 alignment_size_ = 0;
331 prefill_ready_ = false;
332 prefill_done_ = false;
333 size_ = 2;
334 records_ = static_cast<T*>(std::malloc(sizeof(T) * 2));
335 readIndex_ = 0;
336 writeIndex_ = 0;
337}
338
339// Hidden original write implementation with signature difference. Only used for pre-allocation
340template<class T>
341template<class... Args>
342bool
343IterableQueueModel<T>::write_(Args&&... recordArgs)
344{
345 // const std::lock_guard<std::mutex> lock(m_mutex);
346 auto const currentWrite = writeIndex_.load(std::memory_order_relaxed);
347 auto nextRecord = currentWrite + 1;
348 if (nextRecord == size_) {
349 nextRecord = 0;
350 }
351 // if (nextRecord == readIndex_.load(std::memory_order_acquire)) {
352 // std::cout << "SPSC WARNING -> Queue is full! WRITE PASSES READ!!! \n";
353 //}
354 // new (&records_[currentWrite]) T(std::forward<Args>(recordArgs)...);
355 // writeIndex_.store(nextRecord, std::memory_order_release);
356 // return true;
357
358 // ORIGINAL:
359 if (nextRecord != readIndex_.load(std::memory_order_acquire)) {
360 new (&records_[currentWrite]) T(std::forward<Args>(recordArgs)...);
361 writeIndex_.store(nextRecord, std::memory_order_release);
362 return true;
363 }
364 // queue is full
365
366 ++overflow_ctr;
367
368 return false;
369}
370
371template<class T>
372void
374{
376 info.set_num_buffer_elements(this->occupancy());
377 this->publish(std::move(info));
378}
379
380} // namespace datahandlinglibs
381} // namespace dunedaq
#define ERS_HERE
uint32_t get_alignment_size() const
Get "alignment_size" attribute value.
bool get_preallocation() const
Get "preallocation" attribute value.
int16_t get_numa_node() const
Get "numa_node" attribute value.
uint32_t get_size() const
Get "size" attribute value.
bool get_intrinsic_allocator() const
Get "intrinsic_allocator" attribute value.
bool get_numa_aware() const
Get "numa_aware" attribute value.
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
The DUNE-DAQ namespace.
SourceID[" << sourceid << "] Command daqdataformats::SourceID Readout Initialization std::string initerror Configuration std::string conferror GenericConfigurationError
void flush() override
Flush all elements from the latency buffer.
std::size_t occupancy() const override
Occupancy of LB.
void allocate_memory(std::size_t size, bool numa_aware, uint8_t numa_node=0, bool intrinsic_allocator=false, std::size_t alignment_size=0)
const T * back() override
Get pointer to the back of the LB.
const T * front() override
Get pointer to the front of the LB.
bool write(T &&record) override
Move referenced object into LB.
void conf(const appmodel::LatencyBuffer *cfg) override
Configure the LB.
void pop(std::size_t x)
Pop specified amount of elements from LB.
void scrap(const appfwk::DAQModule::CommandData_t &) override
Unconfigure the LB.
bool read(T &record) override
Move object from LB to referenced.