DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TPGInternalStateHarvester.hpp
Go to the documentation of this file.
1
8#pragma once
9
10#include <array>
11#include <atomic>
12#include <chrono>
13#include <condition_variable>
14#include <cstdint>
15#include <memory>
16#include <mutex>
17#include <string>
18#include <thread>
19#include <unordered_map>
20#include <vector>
21
23#include "trgdataformats/Types.hpp" // or the header that defines trgdataformats::channel_t
24
25namespace dunedaq {
26namespace fdreadoutlibs {
27
29{
30public:
31 using ProcRef = std::pair<std::shared_ptr<tpglibs::AbstractProcessor<__m256i>>, int /*pipeline_id*/>;
32
33 // Destructor ensures proper thread cleanup
35
36 void set_processor_references(std::vector<ProcRef> refs);
37 const std::vector<ProcRef>& get_processor_references() const;
38
47 const std::vector<std::pair<trgdataformats::channel_t, int16_t>>& channel_plane_numbers,
48 uint8_t num_channels_per_pipeline,
49 uint8_t num_pipelines);
50
56 std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>> harvest_once();
57
58 // --- Multi-threaded interface ---
59
65
71
77
84 std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>> get_latest_results()
85 const;
86
93
94private:
95 // --- Original data structures ---
96 std::vector<ProcRef> m_processor_references;
97 // index: pipeline_id -> vector of (channel, plane) for its 16 lanes
98 std::vector<std::vector<std::pair<trgdataformats::channel_t, int16_t>>> m_channel_plane_numbers_per_pipeline;
101
102 // --- Preallocation caches (rebuilt when refs/channels change) ---
103 std::vector<std::vector<std::string>> m_metric_items_per_proc; // size == m_processor_references.size()
104 std::vector<size_t> m_expected_items_per_pipeline; // size == m_num_pipelines
105 size_t m_expected_total_channels = 0; // == m_num_channels_per_pipeline * m_num_pipelines
106
107 // --- Multi-threaded data structures ---
108 using ResultContainer = std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>>;
109
110 // Single result container with mutex protection
111 mutable std::mutex m_results_mutex;
113
114 // Thread synchronization
116 std::atomic<bool> m_thread_should_stop{ false };
117 std::atomic<bool> m_thread_running{ false };
118 std::atomic<bool> m_harvest_requested{ false };
119
120 // Synchronization primitives
121 mutable std::mutex m_config_mutex; // Protects configuration changes
122 std::mutex m_collection_mutex; // Protects collection process
123 std::condition_variable m_collection_cv; // Signals collection thread
124
125 // Thread management
127
128 // Rebuild caches after refs or channel layout changes.
130};
131
132} // namespace fdreadoutlibs
133} // namespace dunedaq
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.
std::unordered_map< trgdataformats::channel_t, std::vector< std::pair< std::string, int16_t > > > ResultContainer
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::pair< std::shared_ptr< tpglibs::AbstractProcessor< __m256i > >, int > ProcRef
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.
The DUNE-DAQ namespace.