DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
dunedaq::trgtools::TAEmulationWorker Class Reference

#include <TAEmulationWorker.hpp>

Public Member Functions

 TAEmulationWorker (std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > _input_files, nlohmann::json _config, std::pair< uint64_t, uint64_t > _sliceid_range, bool _run_parallel, bool _quiet)
 Constructor, takes file input path & configuration.
 ~TAEmulationWorker ()=default
std::vector< daqdataformats::SourceID > get_valid_sourceids (daqdataformats::TimeSlice &_timeslice)
 Get the valid sourceids object from HDF5 file.
void start_processing ()
 User interaction for task processing.
void wait_to_complete_work ()
 Waits for all the tasks to complete.
std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > get_tas ()
 Retrieves all the TAs with std::move operator.
std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > get_frags ()
 Retrieves all the unique pointers to the TA fragments.
hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t get_sourceid_geoid_map ()
 Get the sourceid to geoid map object.

Private Member Functions

void process_tasks ()
 Function that processes the whole file.
void process_task (daqdataformats::SourceID _source_id, uint64_t _rec, daqdataformats::FragmentHeader _header, std::vector< trgdataformats::TriggerPrimitive > &&_tps)
 Function that processes one slice for one plane.
void worker_thread ()
 Creates & runs a worker thread.
void enqueue_task (std::function< void()> task)
 Enqueues task to process.
void wait_to_complete_tasks ()
 Waits to complete a task.

Private Attributes

std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > m_input_files
 A pointer to the input file.
nlohmann::json m_configuration
 configuration for the TA-makers
std::vector< std::string > m_input_paths
 input vector of tpstream input paths
std::map< daqdataformats::SourceID, std::unique_ptr< trgtools::TAEmulationUnit > > m_ta_emulators
 Map of SourceID : Emulator unit (TAMaker).
std::pair< uint64_t, uint64_t > m_sliceid_range
 Range of TimeSlice IDs to process.
const bool m_run_parallel
 Run the TA makers in parllel.
const bool m_quiet
 Quiet down the cout output.
std::thread m_main_thread
 The file handler thread.
std::atomic< bool > m_stop { false }
 Bool to indicate to stop the emulation.
std::mutex m_savetps_mutex
 Mutex for saving the TPs.
std::vector< std::thread > m_thread_pool
std::condition_variable m_condition
std::condition_variable m_task_complete_condition
std::queue< std::function< void()> > m_task_queue
std::mutex m_queue_mutex
std::atomic< size_t > m_active_tasks
std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > m_tas
 Output vector of TAs.
std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > m_ta_fragments
 Output vector of TA fragments.
uint16_t m_id
 Unique ID for this TAEmulationWorker.

Static Private Attributes

static uint16_t m_id_next = 0
 Global variable used to get the next ID.
static const size_t SIZE_TP = sizeof(trgdataformats::TriggerPrimitive)
 Size of the TPS.

Detailed Description

Definition at line 22 of file TAEmulationWorker.hpp.

Constructor & Destructor Documentation

◆ TAEmulationWorker()

dunedaq::trgtools::TAEmulationWorker::TAEmulationWorker ( std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > _input_files,
nlohmann::json _config,
std::pair< uint64_t, uint64_t > _sliceid_range,
bool _run_parallel,
bool _quiet )

Constructor, takes file input path & configuration.

Each TAEmulationWorker will create its own thread, so all TAEmulationWorkers are run on separate threads

Parameters
_input_filesa vector of input HDF5 shared pointers to process
_configTAMaker configuration
_sliceid_rangerange of sliceids to process
_run_parallelrun each TAMaker (one per SourceID) in parallel
_quietquiet down the cout

Definition at line 10 of file TAEmulationWorker.cpp.

15 : m_input_files(input_files)
16 , m_sliceid_range(sliceid_range)
17 , m_run_parallel(run_parallel)
18 , m_quiet(quiet)
19 , m_id(m_id_next++)
20{
21 std::string algo_name = config["trigger_activity_plugin"][0];
22 nlohmann::json algo_config = config["trigger_activity_config"][0];
23
24 // Get the input file
25 // Extract the run number etc
26 std::vector<daqdataformats::run_number_t> run_numbers;
27 std::vector<size_t> file_indices;
28 for (const auto& input_file : input_files) {
29 if (std::find(run_numbers.begin(),
30 run_numbers.end(),
31 input_file->get_attribute<daqdataformats::run_number_t>("run_number")) == run_numbers.end()) {
32 run_numbers.push_back(input_file->get_attribute<daqdataformats::run_number_t>("run_number"));
33 }
34
35 if (std::find(file_indices.begin(),
36 file_indices.end(),
37 input_file->get_attribute<daqdataformats::run_number_t>("run_number")) == file_indices.end()) {
38 file_indices.push_back(input_file->get_attribute<size_t>("file_index"));
39 }
40 }
41
42 std::string application_name = m_input_files.front()->get_attribute<std::string>("application_name");
43
44 if (!m_quiet) {
45 fmt::print("Run Numbers: {}\nFile Indices: {}\nApp name: '{}'\n",
46 fmt::join(run_numbers, ","),
47 fmt::join(file_indices, ","),
48 application_name);
49 }
50
51 // std::set of record IDs (pair of record number & sequence number)
52 auto records = m_input_files.front()->get_all_record_ids();
53
54 // Extract the number of TAMakers to create
55 daqdataformats::TimeSlice first_timeslice = m_input_files.front()->get_timeslice(*records.begin());
56 std::vector<daqdataformats::SourceID> valid_sources = get_valid_sourceids(first_timeslice);
57 fmt::print("Number of makers to make: {}\n", valid_sources.size());
58
59 for (const daqdataformats::SourceID& sid : valid_sources) {
60 // Create TAMaker
61 std::unique_ptr<triggeralgs::TriggerActivityMaker> ta_maker =
63 ta_maker->configure(algo_config);
64
65 // Add it to the enulators
66 m_ta_emulators[sid] = std::make_unique<trgtools::TAEmulationUnit>();
67 m_ta_emulators[sid]->set_maker(ta_maker);
68
69 // Create a worker thread per emulator
70 if (m_run_parallel) {
72 }
73 }
74}
std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > m_input_files
A pointer to the input file.
std::vector< std::thread > m_thread_pool
std::map< daqdataformats::SourceID, std::unique_ptr< trgtools::TAEmulationUnit > > m_ta_emulators
Map of SourceID : Emulator unit (TAMaker).
static uint16_t m_id_next
Global variable used to get the next ID.
std::pair< uint64_t, uint64_t > m_sliceid_range
Range of TimeSlice IDs to process.
const bool m_run_parallel
Run the TA makers in parllel.
std::vector< daqdataformats::SourceID > get_valid_sourceids(daqdataformats::TimeSlice &_timeslice)
Get the valid sourceids object from HDF5 file.
uint16_t m_id
Unique ID for this TAEmulationWorker.
const bool m_quiet
Quiet down the cout output.
void worker_thread()
Creates & runs a worker thread.
static std::shared_ptr< AbstractFactory< TriggerActivityMaker > > get_instance()

◆ ~TAEmulationWorker()

dunedaq::trgtools::TAEmulationWorker::~TAEmulationWorker ( )
default

Member Function Documentation

◆ enqueue_task()

void dunedaq::trgtools::TAEmulationWorker::enqueue_task ( std::function< void()> task)
private

Enqueues task to process.

Parameters
tasktask to porcess

Definition at line 225 of file TAEmulationWorker.cpp.

226{
227 {
228 std::lock_guard<std::mutex> lock(m_queue_mutex);
229 m_task_queue.push(std::move(task));
231 }
232 m_condition.notify_one();
233}
std::queue< std::function< void()> > m_task_queue

◆ get_frags()

std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > dunedaq::trgtools::TAEmulationWorker::get_frags ( )

Retrieves all the unique pointers to the TA fragments.

Returns
A map of sourceID : TA fragment

Definition at line 316 of file TAEmulationWorker.cpp.

317{
318 return std::move(m_ta_fragments);
319}
std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > m_ta_fragments
Output vector of TA fragments.

◆ get_sourceid_geoid_map()

hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t dunedaq::trgtools::TAEmulationWorker::get_sourceid_geoid_map ( )

Get the sourceid to geoid map object.

Returns
geoid to sourceid map

Definition at line 96 of file TAEmulationWorker.cpp.

97{
98 if (!m_input_files.size()) {
99 throw "Files not set yet!";
100 }
101
102 return m_input_files.front()->get_srcid_geoid_map();
103}

◆ get_tas()

std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > dunedaq::trgtools::TAEmulationWorker::get_tas ( )

Retrieves all the TAs with std::move operator.

Returns
a map of sourceID : TA

Definition at line 310 of file TAEmulationWorker.cpp.

311{
312 return std::move(m_tas);
313}
std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > m_tas
Output vector of TAs.

◆ get_valid_sourceids()

std::vector< daqdataformats::SourceID > dunedaq::trgtools::TAEmulationWorker::get_valid_sourceids ( daqdataformats::TimeSlice & _timeslice)

Get the valid sourceids object from HDF5 file.

Parameters
_timeslicetimeslice to load the sourceIDs from
Returns
vector of SoureIDs from this file

Definition at line 77 of file TAEmulationWorker.cpp.

78{
79 const auto& fragments = _timeslice.get_fragments_ref();
80
81 std::vector<daqdataformats::SourceID> ret;
82 for (const auto& fragment : fragments) {
83 if (fragment->get_fragment_type() != daqdataformats::FragmentType::kTriggerPrimitive) {
84 continue;
85 }
86
87 daqdataformats::SourceID sourceid = fragment->get_element_id();
88
89 ret.push_back(sourceid);
90 }
91
92 return ret;
93}
@ kTriggerPrimitive
Trigger format TPs produced by trigger code.
Cannot add TPSet with sourceid

◆ process_task()

void dunedaq::trgtools::TAEmulationWorker::process_task ( daqdataformats::SourceID _source_id,
uint64_t _rec,
daqdataformats::FragmentHeader _header,
std::vector< trgdataformats::TriggerPrimitive > && _tps )
private

Function that processes one slice for one plane.

Parameters
_source_idsourceID to process
_recrecord ID to process
_headerfragment header
_tpsvectors of trigger primitives to process (with move operator)

Definition at line 267 of file TAEmulationWorker.cpp.

271{
272 // Get te last fragment
273 std::unique_ptr<daqdataformats::Fragment> frag = m_ta_emulators[_source_id]->emulate_vector(_tps);
274
275 // Don't do anything if no fragments found
276 if (!frag) {
277 return;
278 }
279
280 // Get all the TriggerActivities from the TA Emulator buffer
281 std::vector<triggeralgs::TriggerActivity> ta_buffer = m_ta_emulators[_source_id]->get_last_output_buffer();
282
283 // Don't continue if no TAs found
284 size_t n_tas = ta_buffer.size();
285 if (!n_tas) {
286 return;
287 }
288
289 if (!m_quiet && n_tas) {
290 fmt::print(" Found {} TAs!\n", n_tas);
291 }
292
293 // Set the fragment header & push into our output (with locking!)
294 {
295 if (m_run_parallel) {
296 std::lock_guard<std::mutex> lock(m_savetps_mutex);
297 }
298 m_tas[_rec].reserve(m_tas[_rec].size() + ta_buffer.size());
299 m_tas[_rec].insert(
300 m_tas[_rec].end(), std::make_move_iterator(ta_buffer.begin()), std::make_move_iterator(ta_buffer.end()));
301
302 frag->set_header_fields(_header);
304
305 m_ta_fragments[_rec].push_back(std::move(frag));
306 }
307}
std::mutex m_savetps_mutex
Mutex for saving the TPs.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size

◆ process_tasks()

void dunedaq::trgtools::TAEmulationWorker::process_tasks ( )
private

Function that processes the whole file.

Definition at line 138 of file TAEmulationWorker.cpp.

139{
140 // Iterate over the input files
141 for (auto& input_file : m_input_files) {
142 // std::set of record IDs (pair of record number & sequence number)
143 auto records = input_file->get_all_record_ids();
144
145 for (const auto& record : records) {
146 if (record.first < m_sliceid_range.first || record.first > m_sliceid_range.second) {
147 if (!m_quiet)
148 fmt::print(" Will not process RecordID {} because it's outside of our range!", record.first);
149 continue;
150 }
151
152 // Get all the fragments
153 daqdataformats::TimeSlice timeslice = input_file->get_timeslice(record);
154 const auto& fragments = timeslice.get_fragments_ref();
155
156 // Iterate over the fragments & process each fragment
157 for (const auto& fragment : fragments) {
158 daqdataformats::SourceID sid = fragment->get_element_id();
159
160 if (!m_ta_emulators.contains(sid)) {
161 continue;
162 }
163
164 // Pull tps out
165 size_t n_tps = fragment->get_data_size() / SIZE_TP;
166 if (!m_quiet) {
167 fmt::print(" TP fragment size: {}\n", fragment->get_data_size());
168 fmt::print(" Num TPs: {}\n", n_tps);
169 }
170
171 // Create a TP buffer
172 std::vector<trgdataformats::TriggerPrimitive> tp_buffer;
173 // Prepare the TP buffer, checking for time ordering
174 tp_buffer.reserve(n_tps);
175
176 // Populate the TP buffer
177 trgdataformats::TriggerPrimitive* tp_array =
178 static_cast<trgdataformats::TriggerPrimitive*>(fragment->get_data());
179 uint64_t last_ts = 0;
180 for (size_t tpid(0); tpid < n_tps; ++tpid) {
181 auto& tp = tp_array[tpid];
182 if (tp.time_start <= last_ts && !m_quiet) {
183 fmt::print(" ERROR: {} {} ", +tp.time_start, last_ts);
184 }
185 tp_buffer.push_back(tp);
186 }
187
188 daqdataformats::FragmentHeader frag_hdr = fragment->get_header();
189
190 // Customise the source id (add 1000 to id)
191 frag_hdr.element_id = daqdataformats::SourceID{ daqdataformats::SourceID::Subsystem::kTrigger,
192 fragment->get_element_id().id + 1000 };
193
194 // Either enqueue the task if using parallel processing, or execute the task now
195 if (m_run_parallel) {
196 enqueue_task([this, sid, record, frag_hdr, tp_buffer = std::move(tp_buffer)]() mutable {
197 this->process_task(sid, record.first, frag_hdr, std::move(tp_buffer));
198 });
199 } else {
200 this->process_task(sid, record.first, frag_hdr, std::move(tp_buffer));
201 }
202 }
203 // If running in parallel, wait to process entire slice before we move to
204 // the next one
205 if (m_run_parallel) {
207 }
208 }
209
210 size_t total = 0;
211 for (auto& [key, vec_tas] : m_tas) {
212 total += vec_tas.size();
213 }
214 std::cout << "We have a total of " << total << " TAs!" << std::endl;
215 }
216}
void enqueue_task(std::function< void()> task)
Enqueues task to process.
static const size_t SIZE_TP
Size of the TPS.
void wait_to_complete_tasks()
Waits to complete a task.
void process_task(daqdataformats::SourceID _source_id, uint64_t _rec, daqdataformats::FragmentHeader _header, std::vector< trgdataformats::TriggerPrimitive > &&_tps)
Function that processes one slice for one plane.

◆ start_processing()

void dunedaq::trgtools::TAEmulationWorker::start_processing ( )

User interaction for task processing.

Definition at line 219 of file TAEmulationWorker.cpp.

220{
222}
std::thread m_main_thread
The file handler thread.
void process_tasks()
Function that processes the whole file.

◆ wait_to_complete_tasks()

void dunedaq::trgtools::TAEmulationWorker::wait_to_complete_tasks ( )
private

Waits to complete a task.

Definition at line 236 of file TAEmulationWorker.cpp.

237{
238 std::unique_lock<std::mutex> lock(m_queue_mutex);
239 m_task_complete_condition.wait(lock, [this]() { return m_active_tasks == 0; });
240}
std::condition_variable m_task_complete_condition

◆ wait_to_complete_work()

void dunedaq::trgtools::TAEmulationWorker::wait_to_complete_work ( )

Waits for all the tasks to complete.

Definition at line 243 of file TAEmulationWorker.cpp.

244{
245 // Wait for the main threads to join
246 m_main_thread.join();
247 fmt::print("TAEmulationWorker_{} work completed\n", m_id);
248
249 // Wait for the tasks to complete
250 if (m_run_parallel) {
252
253 {
254 std::lock_guard<std::mutex> lock(m_queue_mutex);
255 m_stop = true;
256 }
257 m_condition.notify_all();
258
259 fmt::print("m_stop issued\n");
260 for (std::thread& thread : m_thread_pool) {
261 thread.join();
262 }
263 }
264}
std::atomic< bool > m_stop
Bool to indicate to stop the emulation.

◆ worker_thread()

void dunedaq::trgtools::TAEmulationWorker::worker_thread ( )
private

Creates & runs a worker thread.

Get a task from the queue (with locking)

Definition at line 106 of file TAEmulationWorker.cpp.

107{
108 while (true) {
110 std::function<void()> task;
111 {
112 std::unique_lock<std::mutex> lock(m_queue_mutex);
113 m_condition.wait(lock, [this]() { return m_stop || !m_task_queue.empty(); });
114
115 if (m_stop && m_task_queue.empty()) {
116 return;
117 }
118
119 task = std::move(m_task_queue.front());
120 m_task_queue.pop();
121 }
122
123 // Run & complete a task
124 task();
125
126 // Notify that task was completed
127 {
128 std::lock_guard<std::mutex> lock(m_queue_mutex);
130 if (m_active_tasks == 0) {
131 m_task_complete_condition.notify_all();
132 }
133 }
134 }
135}

Member Data Documentation

◆ m_active_tasks

std::atomic<size_t> dunedaq::trgtools::TAEmulationWorker::m_active_tasks
private

Definition at line 154 of file TAEmulationWorker.hpp.

◆ m_condition

std::condition_variable dunedaq::trgtools::TAEmulationWorker::m_condition
private

Definition at line 150 of file TAEmulationWorker.hpp.

◆ m_configuration

nlohmann::json dunedaq::trgtools::TAEmulationWorker::m_configuration
private

configuration for the TA-makers

Definition at line 115 of file TAEmulationWorker.hpp.

◆ m_id

uint16_t dunedaq::trgtools::TAEmulationWorker::m_id
private

Unique ID for this TAEmulationWorker.

Definition at line 163 of file TAEmulationWorker.hpp.

◆ m_id_next

uint16_t dunedaq::trgtools::TAEmulationWorker::m_id_next = 0
staticprivate

Global variable used to get the next ID.

Definition at line 165 of file TAEmulationWorker.hpp.

◆ m_input_files

std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile> > dunedaq::trgtools::TAEmulationWorker::m_input_files
private

A pointer to the input file.

Definition at line 112 of file TAEmulationWorker.hpp.

◆ m_input_paths

std::vector<std::string> dunedaq::trgtools::TAEmulationWorker::m_input_paths
private

input vector of tpstream input paths

Definition at line 118 of file TAEmulationWorker.hpp.

◆ m_main_thread

std::thread dunedaq::trgtools::TAEmulationWorker::m_main_thread
private

The file handler thread.

Definition at line 137 of file TAEmulationWorker.hpp.

◆ m_queue_mutex

std::mutex dunedaq::trgtools::TAEmulationWorker::m_queue_mutex
private

Definition at line 153 of file TAEmulationWorker.hpp.

◆ m_quiet

const bool dunedaq::trgtools::TAEmulationWorker::m_quiet
private

Quiet down the cout output.

Definition at line 130 of file TAEmulationWorker.hpp.

◆ m_run_parallel

const bool dunedaq::trgtools::TAEmulationWorker::m_run_parallel
private

Run the TA makers in parllel.

Definition at line 127 of file TAEmulationWorker.hpp.

◆ m_savetps_mutex

std::mutex dunedaq::trgtools::TAEmulationWorker::m_savetps_mutex
private

Mutex for saving the TPs.

Definition at line 148 of file TAEmulationWorker.hpp.

◆ m_sliceid_range

std::pair<uint64_t, uint64_t> dunedaq::trgtools::TAEmulationWorker::m_sliceid_range
private

Range of TimeSlice IDs to process.

Definition at line 124 of file TAEmulationWorker.hpp.

◆ m_stop

std::atomic<bool> dunedaq::trgtools::TAEmulationWorker::m_stop { false }
private

Bool to indicate to stop the emulation.

Definition at line 140 of file TAEmulationWorker.hpp.

◆ m_ta_emulators

std::map<daqdataformats::SourceID, std::unique_ptr<trgtools::TAEmulationUnit> > dunedaq::trgtools::TAEmulationWorker::m_ta_emulators
private

Map of SourceID : Emulator unit (TAMaker).

Definition at line 121 of file TAEmulationWorker.hpp.

◆ m_ta_fragments

std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment> > > dunedaq::trgtools::TAEmulationWorker::m_ta_fragments
private

Output vector of TA fragments.

Definition at line 160 of file TAEmulationWorker.hpp.

◆ m_tas

std::map<uint64_t, std::vector<triggeralgs::TriggerActivity> > dunedaq::trgtools::TAEmulationWorker::m_tas
private

Output vector of TAs.

Definition at line 157 of file TAEmulationWorker.hpp.

◆ m_task_complete_condition

std::condition_variable dunedaq::trgtools::TAEmulationWorker::m_task_complete_condition
private

Definition at line 151 of file TAEmulationWorker.hpp.

◆ m_task_queue

std::queue<std::function<void()> > dunedaq::trgtools::TAEmulationWorker::m_task_queue
private

Definition at line 152 of file TAEmulationWorker.hpp.

◆ m_thread_pool

std::vector<std::thread> dunedaq::trgtools::TAEmulationWorker::m_thread_pool
private

Definition at line 149 of file TAEmulationWorker.hpp.

◆ SIZE_TP

const size_t dunedaq::trgtools::TAEmulationWorker::SIZE_TP = sizeof(trgdataformats::TriggerPrimitive)
staticprivate

Size of the TPS.

Definition at line 167 of file TAEmulationWorker.hpp.


The documentation for this class was generated from the following files: