14 if (!std::is_trivially_destructible<T>::value) {
17 while (readIndex != endIndex) {
19 if (++readIndex ==
size_) {
28#ifdef WITH_LIBNUMA_SUPPORT
42 bool intrinsic_allocator,
43 std::size_t alignment_size)
50#ifdef WITH_LIBNUMA_SUPPORT
51 numa_set_preferred((
unsigned)numa_node);
52#ifdef WITH_LIBNUMA_BIND_POLICY
53 numa_set_bind_policy(WITH_LIBNUMA_BIND_POLICY);
55#ifdef WITH_LIBNUMA_STRICT_POLICY
56 numa_set_strict(WITH_LIBNUMA_STRICT_POLICY);
58 records_ =
static_cast<T*
>(numa_alloc_onnode(
sizeof(T) *
size, numa_node));
61 "NUMA allocation was requested but program was built without USE_LIBNUMA");
63 }
else if (intrinsic_allocator && alignment_size > 0) {
64 records_ =
static_cast<T*
>(_mm_malloc(
sizeof(T) *
size, alignment_size));
65 }
else if (!intrinsic_allocator && alignment_size > 0) {
66 records_ =
static_cast<T*
>(std::aligned_alloc(alignment_size,
sizeof(T) *
size));
67 }
else if (!numa_aware && !intrinsic_allocator && alignment_size == 0) {
69 records_ =
static_cast<T*
>(std::malloc(
sizeof(T) *
size));
92 for (
size_t i = 0; i <
size_ - 1; ++i) {
94 write_(std::move(element));
116 auto handle = prefill_thread.native_handle();
117 pthread_setname_np(handle, tname);
119#ifdef WITH_LIBNUMA_SUPPORT
120 cpu_set_t affinitymask;
121 CPU_ZERO(&affinitymask);
122 struct bitmask* nodecpumask = numa_allocate_cpumask();
125 ret = numa_node_to_cpus(
numa_node_, nodecpumask);
128 for (
int i = 0; i < numa_num_configured_cpus(); ++i) {
129 if (numa_bitmask_isbitset(nodecpumask, i)) {
130 CPU_SET(i, &affinitymask);
133 ret = pthread_setaffinity_np(handle,
sizeof(cpu_set_t), &affinitymask);
135 numa_free_cpumask(nodecpumask);
150 prefill_thread.join();
158 auto const currentWrite =
writeIndex_.load(std::memory_order_relaxed);
159 auto nextRecord = currentWrite + 1;
160 if (nextRecord ==
size_) {
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);
180 auto const currentRead =
readIndex_.load(std::memory_order_relaxed);
181 if (currentRead ==
writeIndex_.load(std::memory_order_acquire)) {
186 auto nextRecord = currentRead + 1;
187 if (nextRecord ==
size_) {
190 record = std::move(
records_[currentRead]);
192 readIndex_.store(nextRecord, std::memory_order_release);
201 auto const currentRead =
readIndex_.load(std::memory_order_relaxed);
202 assert(currentRead !=
writeIndex_.load(std::memory_order_acquire));
204 auto nextRecord = currentRead + 1;
205 if (nextRecord ==
size_) {
210 readIndex_.store(nextRecord, std::memory_order_release);
218 for (std::size_t i = 0; i < x; i++) {
236 auto nextRecord =
writeIndex_.load(std::memory_order_acquire) + 1;
237 if (nextRecord ==
size_) {
240 if (nextRecord !=
readIndex_.load(std::memory_order_acquire)) {
257 int ret =
static_cast<int>(
writeIndex_.load(std::memory_order_acquire)) -
258 static_cast<int>(
readIndex_.load(std::memory_order_acquire));
260 ret +=
static_cast<int>(
size_);
262 return static_cast<std::size_t
>(ret);
270 auto const currentRead =
readIndex_.load(std::memory_order_relaxed);
271 if (currentRead ==
writeIndex_.load(std::memory_order_acquire)) {
282 auto const currentWrite =
writeIndex_.load(std::memory_order_relaxed);
283 if (currentWrite ==
readIndex_.load(std::memory_order_acquire)) {
286 int currentLast = currentWrite;
287 if (currentLast == 0) {
288 currentLast =
size_ - 1;
312 throw std::bad_alloc();
334 records_ =
static_cast<T*
>(std::malloc(
sizeof(T) * 2));
341template<
class... Args>
346 auto const currentWrite =
writeIndex_.load(std::memory_order_relaxed);
347 auto nextRecord = currentWrite + 1;
348 if (nextRecord ==
size_) {
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);
376 info.set_num_buffer_elements(this->
occupancy());
377 this->
publish(std::move(info));
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
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::atomic< int > overflow_ctr
std::size_t occupancy() const override
Occupancy of LB.
virtual void generate_opmon_data() override
std::mutex prefill_mutex_
std::atomic< unsigned int > readIndex_
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)
std::condition_variable prefill_cv_
bool write_(Args &&... recordArgs)
const T * back() override
Get pointer to the back of the LB.
const T * front() override
Get pointer to the front of the LB.
std::size_t alignment_size_
bool invalid_configuration_requested_
std::string prefiller_name_
bool write(T &&record) override
Move referenced object into LB.
void conf(const appmodel::LatencyBuffer *cfg) override
Configure the LB.
bool intrinsic_allocator_
void pop(std::size_t x)
Pop specified amount of elements from LB.
std::atomic< unsigned int > writeIndex_
void scrap(const appfwk::DAQModule::CommandData_t &) override
Unconfigure the LB.
bool read(T &record) override
Move object from LB to referenced.