DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
datahandlinglibs
include
datahandlinglibs
models
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
27
#include "
datahandlinglibs/DataHandlingIssues.hpp
"
28
#include "
datahandlinglibs/concepts/LatencyBufferConcept.hpp
"
29
30
#include "
datahandlinglibs/opmon/datahandling_info.pb.h
"
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
58
namespace
dunedaq
{
59
namespace
datahandlinglibs
{
60
70
template
<
class
T>
71
struct
IterableQueueModel
:
public
LatencyBufferConcept
<T>
72
{
73
typedef
T
value_type
;
74
75
static
constexpr
bool
expects_order
=
true
;
76
77
IterableQueueModel
(
const
IterableQueueModel
&) =
delete
;
78
IterableQueueModel
&
operator=
(
const
IterableQueueModel
&) =
delete
;
79
80
// Default constructor
81
IterableQueueModel
()
82
:
LatencyBufferConcept
<T>()
83
,
numa_aware_
(false)
84
,
numa_node_
(0)
85
,
intrinsic_allocator_
(false)
86
,
alignment_size_
(0)
87
,
invalid_configuration_requested_
(false)
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)
108
,
invalid_configuration_requested_
(false)
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
136
~IterableQueueModel
() {
free_memory
(); }
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
:
266
IterableQueueModel<T>
&
m_queue
;
267
uint32_t
m_index
;
// NOLINT(build/unsigned)
268
};
269
270
Iterator
begin
()
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
280
Iterator
end
()
281
{
282
return
Iterator
(*
this
, std::numeric_limits<uint32_t>::max());
// NOLINT(build/unsigned)
283
}
284
285
protected
:
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
296
bool
numa_aware_
;
297
uint8_t
numa_node_
;
// NOLINT (build/unsigned)
298
bool
intrinsic_allocator_
;
299
std::size_t
alignment_size_
;
300
bool
invalid_configuration_requested_
;
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_
;
306
bool
prefill_ready_
;
307
bool
prefill_done_
;
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)
317
T*
records_
;
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
328
#include "
detail/IterableQueueModel.hxx
"
329
330
#endif
// DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_ITERABLEQUEUEMODEL_HPP_
DataHandlingIssues.hpp
IterableQueueModel.hxx
LatencyBufferConcept.hpp
dunedaq::appmodel::LatencyBuffer
Definition
LatencyBuffer.hpp:20
dunedaq::datahandlinglibs::LatencyBufferConcept::LatencyBufferConcept
LatencyBufferConcept()
Definition
LatencyBufferConcept.hpp:34
datahandling_info.pb.h
Logging.hpp
TLOG
#define TLOG(...)
Definition
macro.hpp:21
dunedaq::datahandlinglibs
Definition
DataHandlingConcept.hpp:16
dunedaq
The DUNE-DAQ namespace.
Definition
cib_utilities.cpp:5
dunedaq::size
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
Definition
FelixIssues.hpp:30
std
Definition
SchemaUtils.hpp:118
dunedaq::datahandlinglibs::IterableQueueModel::Iterator
Definition
IterableQueueModel.hpp:214
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::reference
T & reference
Definition
IterableQueueModel.hpp:219
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::iterator_category
std::forward_iterator_tag iterator_category
Definition
IterableQueueModel.hpp:215
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator!=
friend bool operator!=(const Iterator &a, const Iterator &b)
Definition
IterableQueueModel.hpp:251
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::good
bool good()
Definition
IterableQueueModel.hpp:253
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::get_index
uint32_t get_index()
Definition
IterableQueueModel.hpp:263
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator*
reference operator*() const
Definition
IterableQueueModel.hpp:227
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator++
Iterator operator++(int amount)
Definition
IterableQueueModel.hpp:242
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::m_queue
IterableQueueModel< T > & m_queue
Definition
IterableQueueModel.hpp:266
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator++
Iterator & operator++()
Definition
IterableQueueModel.hpp:229
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::value_type
T value_type
Definition
IterableQueueModel.hpp:217
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::pointer
T * pointer
Definition
IterableQueueModel.hpp:218
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator->
pointer operator->()
Definition
IterableQueueModel.hpp:228
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::difference_type
std::ptrdiff_t difference_type
Definition
IterableQueueModel.hpp:216
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::m_index
uint32_t m_index
Definition
IterableQueueModel.hpp:267
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::operator==
friend bool operator==(const Iterator &a, const Iterator &b)
Definition
IterableQueueModel.hpp:250
dunedaq::datahandlinglibs::IterableQueueModel::Iterator::Iterator
Iterator(IterableQueueModel< T > &queue, uint32_t index)
Definition
IterableQueueModel.hpp:221
dunedaq::datahandlinglibs::IterableQueueModel::pad0_
char pad0_[folly::hardware_destructive_interference_size]
Definition
IterableQueueModel.hpp:315
dunedaq::datahandlinglibs::IterableQueueModel::capacity
std::size_t capacity() const
Definition
IterableQueueModel.hpp:186
dunedaq::datahandlinglibs::IterableQueueModel::prefill_ready_
bool prefill_ready_
Definition
IterableQueueModel.hpp:306
dunedaq::datahandlinglibs::IterableQueueModel::get_alignment_size
std::size_t get_alignment_size()
Definition
IterableQueueModel.hpp:210
dunedaq::datahandlinglibs::IterableQueueModel::isFull
bool isFull() const
Definition
IterableQueueModel.hxx:234
dunedaq::datahandlinglibs::IterableQueueModel::size_
uint32_t size_
Definition
IterableQueueModel.hpp:316
dunedaq::datahandlinglibs::IterableQueueModel::IterableQueueModel
IterableQueueModel()
Definition
IterableQueueModel.hpp:81
dunedaq::datahandlinglibs::IterableQueueModel::flush
void flush() override
Flush all elements from the latency buffer.
Definition
IterableQueueModel.hpp:207
dunedaq::datahandlinglibs::IterableQueueModel::overflow_ctr
std::atomic< int > overflow_ctr
Definition
IterableQueueModel.hpp:293
dunedaq::datahandlinglibs::IterableQueueModel::value_type
T value_type
Definition
IterableQueueModel.hpp:73
dunedaq::datahandlinglibs::IterableQueueModel::end
Iterator end()
Definition
IterableQueueModel.hpp:280
dunedaq::datahandlinglibs::IterableQueueModel::occupancy
std::size_t occupancy() const override
Occupancy of LB.
Definition
IterableQueueModel.hxx:255
dunedaq::datahandlinglibs::IterableQueueModel::numa_node_
uint8_t numa_node_
Definition
IterableQueueModel.hpp:297
dunedaq::datahandlinglibs::IterableQueueModel::generate_opmon_data
virtual void generate_opmon_data() override
Definition
IterableQueueModel.hxx:373
dunedaq::datahandlinglibs::IterableQueueModel::expects_order
static constexpr bool expects_order
Definition
IterableQueueModel.hpp:75
dunedaq::datahandlinglibs::IterableQueueModel::prefill_mutex_
std::mutex prefill_mutex_
Definition
IterableQueueModel.hpp:304
dunedaq::datahandlinglibs::IterableQueueModel::isEmpty
bool isEmpty() const
Definition
IterableQueueModel.hxx:226
dunedaq::datahandlinglibs::IterableQueueModel::readIndex_
std::atomic< unsigned int > readIndex_
Definition
IterableQueueModel.hpp:318
dunedaq::datahandlinglibs::IterableQueueModel::allocate_memory
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)
Definition
IterableQueueModel.hxx:39
dunedaq::datahandlinglibs::IterableQueueModel::end_of_buffer
T * end_of_buffer()
Definition
IterableQueueModel.hpp:198
dunedaq::datahandlinglibs::IterableQueueModel::prefill_done_
bool prefill_done_
Definition
IterableQueueModel.hpp:307
dunedaq::datahandlinglibs::IterableQueueModel::prefill_cv_
std::condition_variable prefill_cv_
Definition
IterableQueueModel.hpp:305
dunedaq::datahandlinglibs::IterableQueueModel::write_
bool write_(Args &&... recordArgs)
Definition
IterableQueueModel.hxx:343
dunedaq::datahandlinglibs::IterableQueueModel::back
const T * back() override
Get pointer to the back of the LB.
Definition
IterableQueueModel.hxx:280
dunedaq::datahandlinglibs::IterableQueueModel::~IterableQueueModel
~IterableQueueModel()
Definition
IterableQueueModel.hpp:136
dunedaq::datahandlinglibs::IterableQueueModel::popFront
void popFront()
Definition
IterableQueueModel.hxx:199
dunedaq::datahandlinglibs::IterableQueueModel::begin
Iterator begin()
Definition
IterableQueueModel.hpp:270
dunedaq::datahandlinglibs::IterableQueueModel::front
const T * front() override
Get pointer to the front of the LB.
Definition
IterableQueueModel.hxx:268
dunedaq::datahandlinglibs::IterableQueueModel::pad1_
char pad1_[folly::hardware_destructive_interference_size - sizeof(writeIndex_)]
Definition
IterableQueueModel.hpp:321
dunedaq::datahandlinglibs::IterableQueueModel::alignment_size_
std::size_t alignment_size_
Definition
IterableQueueModel.hpp:299
dunedaq::datahandlinglibs::IterableQueueModel::invalid_configuration_requested_
bool invalid_configuration_requested_
Definition
IterableQueueModel.hpp:300
dunedaq::datahandlinglibs::IterableQueueModel::IterableQueueModel
IterableQueueModel(std::size_t size, bool numa_aware=false, uint8_t numa_node=0, bool intrinsic_allocator=false, std::size_t alignment_size=0)
Definition
IterableQueueModel.hpp:98
dunedaq::datahandlinglibs::IterableQueueModel::prefiller_name_
std::string prefiller_name_
Definition
IterableQueueModel.hpp:303
dunedaq::datahandlinglibs::IterableQueueModel::start_of_buffer
T * start_of_buffer()
Definition
IterableQueueModel.hpp:195
dunedaq::datahandlinglibs::IterableQueueModel::records_
T * records_
Definition
IterableQueueModel.hpp:317
dunedaq::datahandlinglibs::IterableQueueModel::write
bool write(T &&record) override
Move referenced object into LB.
Definition
IterableQueueModel.hxx:156
dunedaq::datahandlinglibs::IterableQueueModel::prefill_task
void prefill_task()
Definition
IterableQueueModel.hxx:85
dunedaq::datahandlinglibs::IterableQueueModel::conf
void conf(const appmodel::LatencyBuffer *cfg) override
Configure the LB.
Definition
IterableQueueModel.hxx:298
dunedaq::datahandlinglibs::IterableQueueModel::IterableQueueModel
IterableQueueModel(const IterableQueueModel &)=delete
dunedaq::datahandlinglibs::IterableQueueModel::operator=
IterableQueueModel & operator=(const IterableQueueModel &)=delete
dunedaq::datahandlinglibs::IterableQueueModel::intrinsic_allocator_
bool intrinsic_allocator_
Definition
IterableQueueModel.hpp:298
dunedaq::datahandlinglibs::IterableQueueModel::size
std::size_t size() const
Definition
IterableQueueModel.hpp:183
dunedaq::datahandlinglibs::IterableQueueModel::pop
void pop(std::size_t x)
Pop specified amount of elements from LB.
Definition
IterableQueueModel.hxx:216
dunedaq::datahandlinglibs::IterableQueueModel::writeIndex_
std::atomic< unsigned int > writeIndex_
Definition
IterableQueueModel.hpp:320
dunedaq::datahandlinglibs::IterableQueueModel::scrap
void scrap(const appfwk::DAQModule::CommandData_t &) override
Unconfigure the LB.
Definition
IterableQueueModel.hxx:323
dunedaq::datahandlinglibs::IterableQueueModel::read
bool read(T &record) override
Move object from LB to referenced.
Definition
IterableQueueModel.hxx:178
dunedaq::datahandlinglibs::IterableQueueModel::allocate_memory
void allocate_memory(std::size_t size) override
Definition
IterableQueueModel.hpp:148
dunedaq::datahandlinglibs::IterableQueueModel::free_memory
void free_memory()
Definition
IterableQueueModel.hxx:9
dunedaq::datahandlinglibs::IterableQueueModel::numa_aware_
bool numa_aware_
Definition
IterableQueueModel.hpp:296
dunedaq::datahandlinglibs::IterableQueueModel::force_pagefault
void force_pagefault()
Definition
IterableQueueModel.hxx:108
dunedaq::datahandlinglibs::IterableQueueModel::ptrlogger
std::thread ptrlogger
Definition
IterableQueueModel.hpp:310
Generated on
for DUNE-DAQ by
1.18.0