DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
fdreadoutlibs
include
fdreadoutlibs
tpg
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
22
#include "
tpglibs/AbstractProcessor.hpp
"
23
#include "
trgdataformats/Types.hpp
"
// or the header that defines trgdataformats::channel_t
24
25
namespace
dunedaq
{
26
namespace
fdreadoutlibs
{
27
28
class
TPGInternalStateHarvester
29
{
30
public
:
31
using
ProcRef
= std::pair<std::shared_ptr<tpglibs::AbstractProcessor<__m256i>>,
int
/*pipeline_id*/
>;
32
33
// Destructor ensures proper thread cleanup
34
~TPGInternalStateHarvester
();
35
36
void
set_processor_references
(std::vector<ProcRef> refs);
37
const
std::vector<ProcRef>&
get_processor_references
()
const
;
38
46
void
update_channel_plane_numbers
(
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
64
void
start_collection_thread
();
65
70
void
stop_collection_thread
();
71
76
void
trigger_harvest
();
77
84
std::unordered_map<trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>>
get_latest_results
()
85
const
;
86
92
bool
is_collection_thread_running
()
const
;
93
94
private
:
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
;
99
uint8_t
m_num_channels_per_pipeline
;
100
uint8_t
m_num_pipelines
;
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
;
112
ResultContainer
m_latest_results
;
113
114
// Thread synchronization
115
std::thread
m_collection_thread
;
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
126
void
collection_thread_worker_
();
127
128
// Rebuild caches after refs or channel layout changes.
129
void
rebuild_prealloc_caches_
();
130
};
131
132
}
// namespace fdreadoutlibs
133
}
// namespace dunedaq
AbstractProcessor.hpp
dunedaq::fdreadoutlibs::TPGInternalStateHarvester
Definition
TPGInternalStateHarvester.hpp:29
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_expected_total_channels
size_t m_expected_total_channels
Definition
TPGInternalStateHarvester.hpp:105
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_collection_cv
std::condition_variable m_collection_cv
Definition
TPGInternalStateHarvester.hpp:123
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::stop_collection_thread
void stop_collection_thread()
Stop the background collection thread Blocks until thread is fully stopped.
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::is_collection_thread_running
bool is_collection_thread_running() const
Check if collection thread is running.
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::get_processor_references
const std::vector< ProcRef > & get_processor_references() const
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::start_collection_thread
void start_collection_thread()
Start the background collection thread Must be called before using async collection features.
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::ResultContainer
std::unordered_map< trgdataformats::channel_t, std::vector< std::pair< std::string, int16_t > > > ResultContainer
Definition
TPGInternalStateHarvester.hpp:108
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_thread_running
std::atomic< bool > m_thread_running
Definition
TPGInternalStateHarvester.hpp:117
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_num_channels_per_pipeline
uint8_t m_num_channels_per_pipeline
Definition
TPGInternalStateHarvester.hpp:99
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::trigger_harvest
void trigger_harvest()
Signal the collection thread to perform one harvest cycle Non-blocking - returns immediately.
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_harvest_requested
std::atomic< bool > m_harvest_requested
Definition
TPGInternalStateHarvester.hpp:118
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_num_pipelines
uint8_t m_num_pipelines
Definition
TPGInternalStateHarvester.hpp:100
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_config_mutex
std::mutex m_config_mutex
Definition
TPGInternalStateHarvester.hpp:121
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::harvest_once
std::unordered_map< trgdataformats::channel_t, std::vector< std::pair< std::string, int16_t > > > harvest_once()
Harvest once, outputs channel -> [(metric_name, value)...].
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::get_latest_results
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...
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::collection_thread_worker_
void collection_thread_worker_()
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_expected_items_per_pipeline
std::vector< size_t > m_expected_items_per_pipeline
Definition
TPGInternalStateHarvester.hpp:104
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_processor_references
std::vector< ProcRef > m_processor_references
Definition
TPGInternalStateHarvester.hpp:96
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_thread_should_stop
std::atomic< bool > m_thread_should_stop
Definition
TPGInternalStateHarvester.hpp:116
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::set_processor_references
void set_processor_references(std::vector< ProcRef > refs)
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_latest_results
ResultContainer m_latest_results
Definition
TPGInternalStateHarvester.hpp:112
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::~TPGInternalStateHarvester
~TPGInternalStateHarvester()
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_channel_plane_numbers_per_pipeline
std::vector< std::vector< std::pair< trgdataformats::channel_t, int16_t > > > m_channel_plane_numbers_per_pipeline
Definition
TPGInternalStateHarvester.hpp:98
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_collection_thread
std::thread m_collection_thread
Definition
TPGInternalStateHarvester.hpp:115
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::ProcRef
std::pair< std::shared_ptr< tpglibs::AbstractProcessor< __m256i > >, int > ProcRef
Definition
TPGInternalStateHarvester.hpp:31
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_results_mutex
std::mutex m_results_mutex
Definition
TPGInternalStateHarvester.hpp:111
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_collection_mutex
std::mutex m_collection_mutex
Definition
TPGInternalStateHarvester.hpp:122
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::m_metric_items_per_proc
std::vector< std::vector< std::string > > m_metric_items_per_proc
Definition
TPGInternalStateHarvester.hpp:103
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::update_channel_plane_numbers
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.
dunedaq::fdreadoutlibs::TPGInternalStateHarvester::rebuild_prealloc_caches_
void rebuild_prealloc_caches_()
dunedaq::fdreadoutlibs
Definition
CRTBernFrameProcessor.hpp:16
dunedaq
The DUNE-DAQ namespace.
Definition
cib_utilities.cpp:5
Types.hpp
Generated on
for DUNE-DAQ by
1.18.0