DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TPGInternalStateHarvester.cpp
Go to the documentation of this file.
1
8#ifdef TPGLIBS_ENABLE_STATE_MONITORING
11#include "logging/Logging.hpp"
12#include <algorithm>
13#include <cassert>
14#include <iostream>
15
17
18namespace dunedaq {
19namespace fdreadoutlibs {
20
22{
23 // Ensure thread is properly stopped
25}
26
27void
29{
30 m_processor_references = std::move(refs);
31 rebuild_prealloc_caches_(); // if num_pipelines is not set, this will gracefully get empty per-pipeline statistics,
32 // which will be filled later by update
33}
34
35const std::vector<TPGInternalStateHarvester::ProcRef>&
37{
39}
40
41void
43 const std::vector<std::pair<trgdataformats::channel_t, int16_t>>& channel_plane_numbers,
44 uint8_t num_channels_per_pipeline,
45 uint8_t num_pipelines)
46{
48 m_channel_plane_numbers_per_pipeline.resize(num_pipelines);
49
50 // Cut in order: each pipeline has num_channels_per_pipeline channels
51 for (uint8_t p = 0; p < num_pipelines; ++p) {
52 auto begin = channel_plane_numbers.begin() + p * num_channels_per_pipeline;
53 auto end = begin + num_channels_per_pipeline;
54 m_channel_plane_numbers_per_pipeline[p] = std::vector<std::pair<trgdataformats::channel_t, int16_t>>(begin, end);
55 }
56 m_num_channels_per_pipeline = num_channels_per_pipeline;
57 m_num_pipelines = num_pipelines;
59
60 rebuild_prealloc_caches_(); // now pipelines number is known, can accurately calculate the expected number of metric
61 // items per pipeline
62}
63
64void
66{
67 // pre-clean
70
71 // first clear/reset the expected number of metric items per pipeline
74
75 // cache the metric names for each processor and accumulate to the corresponding pipeline
76 for (const auto& [proc, pipeline_id] : m_processor_references) {
77 if (proc) {
78 // Use the processor's interface method which delegates to the registry
79 auto names = proc->get_requested_internal_state_names();
80 if (static_cast<size_t>(pipeline_id) < m_expected_items_per_pipeline.size()) {
81 m_expected_items_per_pipeline[pipeline_id] += names.size();
82 }
83 m_metric_items_per_proc.emplace_back(std::move(names));
84 } else {
85 m_metric_items_per_proc.emplace_back(); // empty
86 }
87 }
88}
89
90std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>>
92{
93 std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>> out;
94
95 // 1. pre-estimate the capacity of the buckets (usually 64)
98 } else {
99 out.reserve(static_cast<size_t>(m_num_channels_per_pipeline) * m_num_pipelines);
100 }
101
102 // defensive: if the cache is not complete for the current pipelines, rebuild it
105 }
106
107 for (size_t i = 0; i < m_processor_references.size(); ++i) {
108 const auto& [proc, pipeline_id] = m_processor_references[i];
109 if (!proc) {
110 continue;
111 }
112
113 // 1) get the cached metric names; if empty, fall back to pulling directly from processor
114 std::vector<std::string> metric_names_cached;
115 if (m_metric_items_per_proc.size() > i) {
116 metric_names_cached = m_metric_items_per_proc[i];
117 } else {
118 // Fallback: use processor's interface method
119 metric_names_cached = proc->get_requested_internal_state_names();
120 }
121
122 // 2) get the current snapshot
123 const auto arr = proc->read_internal_states_as_integer_array();
124
125 // basic consistency: the number of items should be the same as the number of snapshot items
126 if (metric_names_cached.size() != arr.m_size) {
127 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Processor " << i << " size mismatch: metric_names=" << metric_names_cached.size()
128 << " vs array=" << arr.m_size;
129 continue;
130 }
131
132 // the current pipeline's lane -> (channel, plane)
133 assert(static_cast<size_t>(pipeline_id) < m_channel_plane_numbers_per_pipeline.size());
134 const auto& chan_plane_vec = m_channel_plane_numbers_per_pipeline[pipeline_id];
135
136 // 3) first reserve the capacity for all channels in the current pipeline in out
137 // so that subsequent emplace_back will not trigger allocation
138 std::vector<decltype(out.begin())> iters_for_lanes;
139 iters_for_lanes.resize(chan_plane_vec.size());
140
141 const size_t expected_items_here =
142 (pipeline_id < m_expected_items_per_pipeline.size())
143 ? m_expected_items_per_pipeline[pipeline_id]
144 : metric_names_cached.size(); // fallback, use the number of items of the current processor
145
146 for (size_t lane = 0; lane < chan_plane_vec.size(); ++lane) {
147 const auto ch = chan_plane_vec[lane].first; // offline channel id
148 auto [it, inserted] = out.try_emplace(ch, std::vector<std::pair<std::string, int16_t>>{});
149 if (inserted) {
150 // only reserve the capacity for the first time the channel is encountered
151 it->second.reserve(expected_items_here);
152 }
153 iters_for_lanes[lane] = it;
154 }
155
156 // 4) map the 16-lane array of each metric to the corresponding channel
157 for (size_t item = 0; item < arr.m_size; ++item) {
158 const auto& lanes = arr.m_data[item]; // std::array<int16_t, 16>
159 const auto& name = metric_names_cached[item];
160 const size_t L = lanes.size(); // e.g. 16
161
162 const size_t up_to = std::min(L, chan_plane_vec.size());
163 for (size_t lane = 0; lane < up_to; ++lane) {
164 // directly use the iterator, avoid the repeated lookup of the map
165 iters_for_lanes[lane]->second.emplace_back(name, lanes[lane]);
166 }
167 }
168 }
169
170 return out;
171}
172
173// --- Multi-threaded implementation ---
174
175void
177{
178 std::lock_guard<std::mutex> lock(m_config_mutex);
179
180 if (m_thread_running.load()) {
181 return; // Already running
182 }
183
184 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Starting internal state collection thread";
185
186 // Initialize result container
187 m_latest_results.clear();
188
189 // Reset thread control flags
190 m_thread_should_stop.store(false);
191 m_harvest_requested.store(false);
192
193 // Start the collection thread
195 m_thread_running.store(true);
196}
197
198void
200{
201 {
202 std::lock_guard<std::mutex> config_lock(m_config_mutex);
203
204 if (!m_thread_running.load()) {
205 return; // Already stopped
206 }
207
208 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Stopping internal state collection thread";
209
210 // Signal thread to stop
211 m_thread_should_stop.store(true);
212 }
213
214 // Notify the collection thread using the correct mutex
215 {
216 std::lock_guard<std::mutex> collection_lock(m_collection_mutex);
217 m_collection_cv.notify_all();
218 }
219
220 // Wait for thread to finish
221 if (m_collection_thread.joinable()) {
222 m_collection_thread.join();
223 }
224
225 m_thread_running.store(false);
226
227 // Clear results
228 {
229 std::lock_guard<std::mutex> lock(m_results_mutex);
230 m_latest_results.clear();
231 }
232}
233
234void
236{
237 if (!m_thread_running.load()) {
238 return; // Thread not running
239 }
240
241 m_harvest_requested.store(true);
242
243 // Notify the collection thread using the correct mutex
244 {
245 std::lock_guard<std::mutex> lock(m_collection_mutex);
246 m_collection_cv.notify_all();
247 }
248}
249
250std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>>
252{
253 std::lock_guard<std::mutex> lock(m_results_mutex);
254 return m_latest_results; // Return copy under lock (blocking read is acceptable)
255}
256
257bool
259{
260 return m_thread_running.load();
261}
262
263void
265{
266 while (!m_thread_should_stop.load()) {
267 std::unique_lock<std::mutex> lock(m_collection_mutex);
268
269 // Wait for harvest request or stop signal
270 m_collection_cv.wait(lock, [this] { return m_harvest_requested.load() || m_thread_should_stop.load(); });
271
272 if (m_thread_should_stop.load()) {
273 break;
274 }
275
276 // Reset the harvest request flag
277 m_harvest_requested.store(false);
278
279 // Release lock during collection to allow concurrent reads
280 lock.unlock();
281
282 // Perform the actual harvest (expensive operation in background)
283 auto new_results = harvest_once();
284
285 // Update results with simple mutex protection
286 {
287 std::lock_guard<std::mutex> results_lock(m_results_mutex);
288 m_latest_results = std::move(new_results);
289 }
290 }
291}
292
293} // namespace fdreadoutlibs
294} // namespace dunedaq
295#endif // TPGLIBS_ENABLE_STATE_MONITORING
void stop_collection_thread()
Stop the background collection thread Blocks until thread is fully stopped.
bool is_collection_thread_running() const
Check if collection thread is running.
const std::vector< ProcRef > & get_processor_references() const
void start_collection_thread()
Start the background collection thread Must be called before using async collection features.
void trigger_harvest()
Signal the collection thread to perform one harvest cycle Non-blocking - returns immediately.
std::unordered_map< trgdataformats::channel_t, std::vector< std::pair< std::string, int16_t > > > harvest_once()
Harvest once, outputs channel -> [(metric_name, value)...].
std::unordered_map< trgdataformats::channel_t, std::vector< std::pair< std::string, int16_t > > > get_latest_results() const
Get the latest collected results (thread-safe, non-blocking) Returns a copy of the most recent harves...
void set_processor_references(std::vector< ProcRef > refs)
std::vector< std::vector< std::pair< trgdataformats::channel_t, int16_t > > > m_channel_plane_numbers_per_pipeline
std::vector< std::vector< std::string > > m_metric_items_per_proc
void update_channel_plane_numbers(const std::vector< std::pair< trgdataformats::channel_t, int16_t > > &channel_plane_numbers, uint8_t num_channels_per_pipeline, uint8_t num_pipelines)
Cuts a full list of (channel, plane) into per-pipeline lists of 16 lanes each.
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed out