DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
ProcessorInternalStateBufferManager.hpp
Go to the documentation of this file.
1
8
11
12#include <array>
13#include <atomic>
14#include <cstdint>
15#include <immintrin.h>
16#include <memory>
17#include <mm_malloc.h>
18#include <string>
19#include <vector>
20
21#ifndef TPGLIBS_PROCESSORINTERNALSTATEBUFFERMANAGER_HPP_
22#define TPGLIBS_PROCESSORINTERNALSTATEBUFFERMANAGER_HPP_
23
24namespace tpglibs {
25
31template<typename T>
33{
34public:
36 using signal_t = T;
37
40
43
49
55
58
64
66 void clear()
67 {
68 for (auto& buf : m_store_buffers) {
69 _mm_free(buf.m_data);
70 }
71 }
72
73protected:
78 void allocate_buffers(size_t buffer_size);
79
84 void allocate_cast_buffers(size_t) {}; // Do nothing for the generic template
85
88
89private:
91 std::vector<std::shared_ptr<signal_t>> m_internal_state_item_ptrs;
92
95
98
100 std::atomic<ProcessorMetricArray<std::array<int16_t, 16>>*> m_cast_active_buffer{ &m_cast_store_buffers[0] };
101
103 std::atomic<ProcessorMetricArray<signal_t>*> m_write_buffer = &m_store_buffers[0];
104
106 std::atomic<ProcessorMetricArray<signal_t>*> m_read_buffer = &m_store_buffers[0];
107
109 std::atomic<uint16_t> m_write_seq{ 0 };
110
112 std::atomic<uint16_t> m_last_read_seq{ 0 };
113
116};
117
118// Template function implementations
119template<typename T>
123
124template<typename T>
129
130template<typename T>
131void
133{
134 // obtain the number of internal state items
135 auto num_items = registry->get_number_of_requested_internal_states();
136 allocate_buffers(num_items);
137 allocate_cast_buffers(num_items);
138 // Writer starts with buffer 0, reader starts with buffer 1 (they must be different!)
139 m_write_buffer.store(&m_store_buffers[0], std::memory_order_release);
140 m_read_buffer.store(&m_store_buffers[1], std::memory_order_release);
141 m_cast_active_buffer.store(&m_cast_store_buffers[0], std::memory_order_release);
142 // reset sequence counters
143 m_write_seq.store(0, std::memory_order_release);
144 m_last_read_seq.store(0, std::memory_order_release);
145 // obtain the pointers to the internal state items
147}
148
149template<typename T>
150void
152{
153 for (auto& buf : m_store_buffers) {
154 buf.m_size = buffer_size;
155 buf.m_data = static_cast<T*>(_mm_malloc(buf.m_size * sizeof(T), alignof(T)));
156 }
157}
158
159template<typename T>
160void
162{
163 // Increment seq to indicate write start (becomes odd)
164 m_write_seq.fetch_add(1, std::memory_order_release);
165
166 auto write_ptr = m_write_buffer.load(std::memory_order_acquire);
167
168 // Write to write buffer
169 for (size_t i = 0; i < m_internal_state_item_ptrs.size(); i++) {
170 // Check for nullptr before dereferencing (handles invalid state names)
171 if (m_internal_state_item_ptrs[i] != nullptr) {
172 write_ptr->m_data[i] = *m_internal_state_item_ptrs[i];
173 } else {
174 // If pointer is null, write a zeroed value
175 write_ptr->m_data[i] = T{};
176 }
177 }
178
179 // Increment seq to indicate write complete (becomes even)
180 m_write_seq.fetch_add(1, std::memory_order_release);
181}
182
183template<typename T>
186{
187 // Wait until no write is in progress
188 uint16_t current_write_seq;
189 do {
190 current_write_seq = m_write_seq.load(std::memory_order_acquire);
191 } while (current_write_seq & 1); // spin if writer is mid-write (odd seq number)
192
193 // Check if there's new data since last read
194 uint16_t last_read = m_last_read_seq.load(std::memory_order_acquire);
195
196 if (current_write_seq != last_read && current_write_seq > 0) {
197 // New data available - swap the buffers between reader and writer
198 auto current_read = m_read_buffer.load(std::memory_order_acquire);
199 auto current_write = m_write_buffer.load(std::memory_order_acquire);
200
201 // Swap: reader gets what writer just finished, writer gets what reader was using
202 m_read_buffer.store(current_write, std::memory_order_release);
203 m_write_buffer.store(current_read, std::memory_order_release);
204
205 // Update last read sequence
206 m_last_read_seq.store(current_write_seq, std::memory_order_release);
207 }
208
209 // Return data from read buffer
210 auto read_ptr = m_read_buffer.load(std::memory_order_acquire);
211 return *read_ptr;
212}
213
214// Template specializations
215// Specialization for __m256i since it has additional cast buffers.
216template<>
217inline void
219{
220 // free double buffer for temp values
221 for (auto& buf : m_store_buffers) {
222 _mm_free(buf.m_data);
223 }
224 // free double buffer for casted values
225 for (auto& buf : m_cast_store_buffers) {
226 _mm_free(buf.m_data);
227 }
228}
229
230// Specialization for __m256i -> std::array<int16_t, 16> cast
231template<>
234{
235 // First, get the raw data using switch_buffer_and_read
236 auto raw_data = switch_buffer_and_read();
237
238 auto* cast_free = m_cast_active_buffer.load(std::memory_order_acquire);
239
240 // Cast each __m256i to std::array<int16_t, 16>
241 for (size_t i = 0; i < raw_data.m_size; ++i) {
242 // Cast and save to cast buffer
243 _mm256_store_si256(reinterpret_cast<__m256i*>(cast_free->m_data[i].data()), raw_data.m_data[i]);
244 }
245
246 // Flip the active cast buffer
247
248 auto* next = (cast_free == &m_cast_store_buffers[0]) ? &m_cast_store_buffers[1] : &m_cast_store_buffers[0];
249 m_cast_active_buffer.store(next, std::memory_order_release);
250
251 return *cast_free;
252}
253
254// Specialization for __m256i
255template<>
256inline void
258{
259 for (auto& buf : m_cast_store_buffers) {
260 buf.m_size = n;
261 buf.m_data = static_cast<std::array<int16_t, 16>*>(_mm_malloc(n * sizeof(std::array<int16_t, 16>), 32));
262 }
263 m_cast_active_buffer.store(&m_cast_store_buffers[0], std::memory_order_release);
264}
265
266// Specialization for std::array<int16_t, 16> -> std::array<int16_t, 16> cast (trivial)
267template<>
270{
271 // First, get the raw data using switch_buffer_and_read
272 auto raw_data = switch_buffer_and_read();
273
274 // For std::array<int16_t, 16>, the cast is trivial - just return the data as-is
275 return raw_data;
276}
277
278} // namespace tpglibs
279
280#endif // TPGLIBS_PROCESSORINTERNALSTATEBUFFERMANAGER_HPP_
std::atomic< ProcessorMetricArray< signal_t > * > m_write_buffer
The write buffer pointer (buffer writer currently uses).
void clear()
clear all buffers and deallocate memory.
std::atomic< uint16_t > m_last_read_seq
The last sequence number that was read.
void configure_from_registry(ProcessorInternalStateNameRegistry< signal_t > *registry)
Configure and allocate correct buffer storage given the configuration string.
ProcessorMetricArray< std::array< int16_t, 16 > > m_cast_store_buffers[2]
The double buffers for storing the internal state data casted to std::array<int16_t,...
std::atomic< uint16_t > m_write_seq
The sequence number for writes (odd=writing, even=complete).
std::vector< std::shared_ptr< signal_t > > m_internal_state_item_ptrs
The vector of pointers to the internal state items.
T signal_t
Signal type to use. Generally __m256i or std::array<int16_t, 16>;.
void allocate_cast_buffers(size_t)
Allocate the correct size for the double buffer read and write buffers for the casted data.
std::atomic< ProcessorMetricArray< signal_t > * > m_read_buffer
The read buffer pointer (buffer reader currently uses).
std::atomic< ProcessorMetricArray< std::array< int16_t, 16 > > * > m_cast_active_buffer
The active buffer for the casted data.
void switch_active_buffer()
Switch the active buffer.
void allocate_buffers(size_t buffer_size)
Allocate the correct size for the double buffer read and write buffers.
ProcessorMetricArray< std::array< int16_t, 16 > > switch_buffer_and_read_casted()
Read from the inactive buffer and cast to std::array<int16_t, 16>.
ProcessorMetricArray< signal_t > switch_buffer_and_read()
Read from the inactive buffer.
ProcessorMetricArray< signal_t > m_store_buffers[2]
The double buffers for storing the internal state data.
size_t get_number_of_requested_internal_states()
Get the number of requested internal states.
std::vector< std::shared_ptr< signal_t > > get_all_requested_internal_state_item_ptrs()
Get a vector of pointers to all internal state items.
Dynamic array of processor metrics, templated on signal type.