DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
IterableQueueModel.hpp
Go to the documentation of this file.
1
17
18// @author Bo Hu (bhu@fb.com)
19// @author Jordan DeLong (delong.j@fb.com)
20
21// Modification by Roland Sipos and Florian Till Groetschla
22// for DUNE-DAQ software framework
23
24#ifndef DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_ITERABLEQUEUEMODEL_HPP_
25#define DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_ITERABLEQUEUEMODEL_HPP_
26
29
31
32#include "logging/Logging.hpp"
33
34#include <folly/lang/Align.h>
35
36#include <atomic>
37#include <cassert>
38#include <cstddef>
39#include <cstdlib>
40#include <cxxabi.h>
41#include <iomanip>
42#include <iostream>
43#include <limits>
44#include <memory>
45#include <mutex>
46#include <new>
47#include <stdexcept>
48#include <thread>
49#include <type_traits>
50#include <utility>
51
52#include <xmmintrin.h>
53
54#ifdef WITH_LIBNUMA_SUPPORT
55#include <numa.h>
56#endif
57
58namespace dunedaq {
59namespace datahandlinglibs {
60
70template<class T>
72{
73 typedef T value_type;
74
75 static constexpr bool expects_order = true;
76
79
80 // Default constructor
83 , numa_aware_(false)
84 , numa_node_(0)
88 , prefill_ready_(false)
89 , prefill_done_(false)
90 , size_(2)
91 , records_(static_cast<T*>(std::malloc(sizeof(T) * 2)))
92 , readIndex_(0)
93 , writeIndex_(0)
94 {
95 }
96
97 // Constructor with alignment strategies
98 IterableQueueModel(std::size_t size, // size must be >= 2
99 bool numa_aware = false,
100 uint8_t numa_node = 0, // NOLINT (build/unsigned)
101 bool intrinsic_allocator = false,
102 std::size_t alignment_size = 0)
103 : LatencyBufferConcept<T>() // NOLINT(build/unsigned)
104 , numa_aware_(numa_aware)
105 , numa_node_(numa_node)
106 , intrinsic_allocator_(intrinsic_allocator)
107 , alignment_size_(alignment_size)
109 , prefill_ready_(false)
110 , prefill_done_(false)
111 , size_(size)
112 , readIndex_(0)
113 , writeIndex_(0)
114 {
115 assert(size >= 2);
116 allocate_memory(size, numa_aware, numa_node, intrinsic_allocator, alignment_size);
117
118 if (!records_) {
119 throw std::bad_alloc();
120 }
121#if 0
122 ptrlogger = std::thread([&](){
123 while(true) {
124 auto const currentRead = readIndex_.load(std::memory_order_relaxed);
125 auto const currentWrite = writeIndex_.load(std::memory_order_relaxed);
126 TLOG() << "BEG:" << std::hex << &records_[0] << " END:" << &records_[size] << std::dec
127 << " R:" << currentRead << " - W:" << currentWrite
128 << " OFLOW:" << overflow_ctr;
129 std::this_thread::sleep_for(std::chrono::milliseconds(100));
130 }
131 });
132#endif
133 }
134
135 // Destructor
137
138 // Free allocated memory that is different for alignment strategies and allocation policies
139 void free_memory();
140
141 // Allocate memory based on different alignment strategies and allocation policies
142 void allocate_memory(std::size_t size,
143 bool numa_aware /*= false*/,
144 uint8_t numa_node = 0, // NOLINT (build/unsigned)
145 bool intrinsic_allocator = false,
146 std::size_t alignment_size = 0);
147
148 void allocate_memory(std::size_t size) override { allocate_memory(size, false); }
149
150 // Task that fills up the LB.
151 void prefill_task();
152
153 // Issue fre-fill task
154 void force_pagefault();
155
156 // Write element into the queue
157 bool write(T&& record) override;
158
159 // Read element from a queue (move or copy the value at the front of the queue to given variable)
160 bool read(T& record) override;
161
162 // Pop element on front of queue
163 void popFront();
164
165 // Pop number of elements (X) from the front of the queue
166 void pop(std::size_t x);
167
168 // Returns true if the queue is empty
169 bool isEmpty() const;
170
171 // Returns true if write index reached read index
172 bool isFull() const;
173
174 // Returns a good-enough guess on current occupancy:
175 // * If called by consumer, then true size may be more (because producer may
176 // be adding items concurrently).
177 // * If called by producer, then true size may be less (because consumer may
178 // be removing items concurrently).
179 // * It is undefined to call this from any other thread.
180 std::size_t occupancy() const override;
181
182 // The size of the underlying buffer, not the amount of usable slots
183 std::size_t size() const { return size_; }
184
185 // Maximum number of items in the queue.
186 std::size_t capacity() const { return size_ - 1; }
187
188 // Gives a pointer to the current read index
189 const T* front() override;
190
191 // Gives a pointer to the last written element
192 const T* back() override;
193
194 // Gives a pointer to the first available slot of the queue
195 T* start_of_buffer() { return &records_[0]; }
196
197 // Gives a pointer to the last available slot of the queue
198 T* end_of_buffer() { return &records_[size_]; }
199
200 // Configures the model
201 void conf(const appmodel::LatencyBuffer* cfg) override;
202
203 // Unconfigures the model
204 void scrap(const appfwk::DAQModule::CommandData_t& /*cfg*/) override;
205
206 // Flushes the elements from the queue
207 void flush() override { pop(occupancy()); }
208
209 // Returns the current memory alignment size
210 std::size_t get_alignment_size() { return alignment_size_; }
211
212 // Iterator for elements in the queue
213 struct Iterator
214 {
215 using iterator_category = std::forward_iterator_tag;
216 using difference_type = std::ptrdiff_t;
217 using value_type = T;
218 using pointer = T*;
219 using reference = T&;
220
221 Iterator(IterableQueueModel<T>& queue, uint32_t index) // NOLINT(build/unsigned)
222 : m_queue(queue)
223 , m_index(index)
224 {
225 }
226
227 reference operator*() const { return m_queue.records_[m_index]; }
228 pointer operator->() { return &m_queue.records_[m_index]; }
229 Iterator& operator++() // NOLINT(runtime/increment_decrement) :)
230 {
231 if (good()) {
232 m_index++;
233 if (m_index == m_queue.size_) {
234 m_index = 0;
235 }
236 }
237 if (!good()) {
238 m_index = std::numeric_limits<uint32_t>::max(); // NOLINT(build/unsigned)
239 }
240 return *this;
241 }
242 Iterator operator++(int amount) // NOLINT(runtime/increment_decrement) :)
243 {
244 Iterator tmp = *this;
245 for (int i = 0; i < amount; ++i) {
246 ++(*this);
247 }
248 return tmp;
249 }
250 friend bool operator==(const Iterator& a, const Iterator& b) { return a.m_index == b.m_index; }
251 friend bool operator!=(const Iterator& a, const Iterator& b) { return a.m_index != b.m_index; }
252
253 bool good()
254 {
255 auto const currentRead = m_queue.readIndex_.load(std::memory_order_relaxed);
256 auto const currentWrite = m_queue.writeIndex_.load(std::memory_order_relaxed);
257 return (*this != m_queue.end()) &&
258 ((m_index >= currentRead && m_index < currentWrite) ||
259 (m_index >= currentRead && currentWrite < currentRead) ||
260 (currentWrite < currentRead && m_index < currentRead && m_index < currentWrite));
261 }
262
263 uint32_t get_index() { return m_index; } // NOLINT(build/unsigned)
264
265 private:
267 uint32_t m_index; // NOLINT(build/unsigned)
268 };
269
271 {
272 auto const currentRead = readIndex_.load(std::memory_order_relaxed);
273 if (currentRead == writeIndex_.load(std::memory_order_acquire)) {
274 // queue is empty
275 return end();
276 }
277 return Iterator(*this, currentRead);
278 }
279
281 {
282 return Iterator(*this, std::numeric_limits<uint32_t>::max()); // NOLINT(build/unsigned)
283 }
284
285protected:
286 virtual void generate_opmon_data() override;
287
288 // Hidden original write implementation with signature difference. Only used for pre-allocation
289 template<class... Args>
290 bool write_(Args&&... recordArgs);
291
292 // Counter for failed writes, due to the fact the queue is full
293 std::atomic<int> overflow_ctr{ 0 };
294
295 // NUMA awareness and aligned allocator usage configuration
297 uint8_t numa_node_; // NOLINT (build/unsigned)
299 std::size_t alignment_size_;
301
302 // Pre-fill and page-fault internal thread control
303 std::string prefiller_name_{ "lbpfn" };
304 std::mutex prefill_mutex_;
305 std::condition_variable prefill_cv_;
308
309 // Ptr logger for debugging
310 std::thread ptrlogger;
311
312 // Underlying buffer with padding:
313 // * hardware_destructive_interference_size is set to 128.
314 // * (Assuming cache line size of 64, so we use a cache line pair size of 128)
315 char pad0_[folly::hardware_destructive_interference_size]; // NOLINT(runtime/arrays)
316 uint32_t size_; // NOLINT(build/unsigned)
318 alignas(folly::hardware_destructive_interference_size) std::atomic<unsigned int> readIndex_; // NOLINT(build/unsigned)
319 alignas(
320 folly::hardware_destructive_interference_size) std::atomic<unsigned int> writeIndex_; // NOLINT(build/unsigned)
321 char pad1_[folly::hardware_destructive_interference_size - sizeof(writeIndex_)]; // NOLINT(runtime/arrays)
322};
323
324} // namespace datahandlinglibs
325} // namespace dunedaq
326
327// Declarations
329
330#endif // DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_ITERABLEQUEUEMODEL_HPP_
#define TLOG(...)
Definition macro.hpp:21
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
friend bool operator!=(const Iterator &a, const Iterator &b)
friend bool operator==(const Iterator &a, const Iterator &b)
Iterator(IterableQueueModel< T > &queue, uint32_t index)
char pad0_[folly::hardware_destructive_interference_size]
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.
char pad1_[folly::hardware_destructive_interference_size - sizeof(writeIndex_)]
IterableQueueModel(std::size_t size, bool numa_aware=false, uint8_t numa_node=0, bool intrinsic_allocator=false, std::size_t alignment_size=0)
bool write(T &&record) override
Move referenced object into LB.
void conf(const appmodel::LatencyBuffer *cfg) override
Configure the LB.
IterableQueueModel(const IterableQueueModel &)=delete
IterableQueueModel & operator=(const IterableQueueModel &)=delete
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.