DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TAEmulationWorker.cpp
Go to the documentation of this file.
1#ifndef TRGTOOLS_TAEMULATIONWORKER_CPP_
2#define TRGTOOLS_TAEMULATIONWORKER_CPP_
3
5
6namespace dunedaq::trgtools {
7
9
10TAEmulationWorker::TAEmulationWorker(std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>> input_files,
11 nlohmann::json config,
12 std::pair<uint64_t, uint64_t> sliceid_range,
13 bool run_parallel,
14 bool quiet)
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}
75
76std::vector<daqdataformats::SourceID>
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}
94
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}
104
105void
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}
136
137void
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
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)
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}
217
218void
223
224void
225TAEmulationWorker::enqueue_task(std::function<void()> task)
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}
234
235void
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}
241
242void
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}
265
266void
268 uint64_t _rec,
270 std::vector<trgdataformats::TriggerPrimitive>&& _tps)
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}
308
309std::map<uint64_t, std::vector<triggeralgs::TriggerActivity>>
311{
312 return std::move(m_tas);
313}
314
315std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>>
317{
318 return std::move(m_ta_fragments);
319}
320
321}; // namespace dunedaq::trgtools
322
323#endif // TRGTOOLS_TAEMULATIONWORKER_CXX_
C++ Representation of a DUNE TimeSlice, consisting of a TimeSliceHeader object and a vector of pointe...
Definition TimeSlice.hpp:27
const std::vector< std::unique_ptr< Fragment > > & get_fragments_ref() const
Get a handle to the Fragments.
Definition TimeSlice.hpp:55
std::map< daqdataformats::SourceID, std::vector< uint64_t > > source_id_geo_id_map_t
std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > get_frags()
Retrieves all the unique pointers to the TA fragments.
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).
void enqueue_task(std::function< void()> task)
Enqueues task to process.
static uint16_t m_id_next
Global variable used to get the next ID.
void start_processing()
User interaction for task processing.
std::pair< uint64_t, uint64_t > m_sliceid_range
Range of TimeSlice IDs to process.
std::queue< std::function< void()> > m_task_queue
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.
std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > m_tas
Output vector of TAs.
std::condition_variable m_task_complete_condition
std::mutex m_savetps_mutex
Mutex for saving the TPs.
const bool m_run_parallel
Run the TA makers in parllel.
void wait_to_complete_work()
Waits for all the tasks to complete.
static const size_t SIZE_TP
Size of the TPS.
std::vector< daqdataformats::SourceID > get_valid_sourceids(daqdataformats::TimeSlice &_timeslice)
Get the valid sourceids object from HDF5 file.
std::map< uint64_t, std::vector< triggeralgs::TriggerActivity > > get_tas()
Retrieves all the TAs with std::move operator.
void wait_to_complete_tasks()
Waits to complete a task.
uint16_t m_id
Unique ID for this TAEmulationWorker.
std::thread m_main_thread
The file handler thread.
std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > m_ta_fragments
Output vector of TA fragments.
void process_tasks()
Function that processes the whole file.
std::atomic< bool > m_stop
Bool to indicate to stop the emulation.
const bool m_quiet
Quiet down the cout output.
void worker_thread()
Creates & runs a worker thread.
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.
hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t get_sourceid_geoid_map()
Get the sourceid to geoid map object.
static std::shared_ptr< AbstractFactory< TriggerActivityMaker > > get_instance()
@ kTriggerPrimitive
Trigger format TPs produced by trigger code.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
Cannot add TPSet with sourceid
The header for a DUNE Fragment.
SourceID element_id
Component that generated the data in this Fragment.
SourceID is a generalized representation of the source of a piece of data in the DAQ....
Definition SourceID.hpp:32
A single energy deposition on a TPC or PDS channel.