DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
emulate_from_tpstream.cxx
Go to the documentation of this file.
3
4#include "CLI/App.hpp"
5#include "CLI/Config.hpp"
6#include "CLI/Formatter.hpp"
7
8#include <filesystem>
9#include <fmt/chrono.h>
10#include <fmt/core.h>
11#include <fmt/format.h>
12#include <optional>
13
16
18
19//-----------------------------------------------------------------------------
20
21using namespace dunedaq;
22using namespace trgtools;
23
32void
33save_fragments(const std::string& _outputfilename,
35 std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>> _frags,
36 bool _quiet)
37{
38 std::string output_filename = _outputfilename;
39 if (output_filename.size() < 5 || output_filename.compare(output_filename.size() - 5, 5, ".hdf5") != 0) {
40 output_filename += ".hdf5";
41 }
42
43 // Create layout parameter object required for HDF5 creation
45
46 // Create HDF5 parameter path required for the layout
47 hdf5libs::HDF5PathParameters params_trigger;
48 params_trigger.detector_group_type = "Detector_Readout";
50 params_trigger.detector_group_name = "TPC";
51 params_trigger.element_name_prefix = "Link";
52 params_trigger.digits_for_element_number = 5;
53
54 // Fill the HDF5 layout
55 std::vector<hdf5libs::HDF5PathParameters> params;
56 params.push_back(params_trigger);
57 layout_params.record_name_prefix = "TimeSlice";
58 layout_params.digits_for_record_number = 6;
59 layout_params.digits_for_sequence_number = 0;
60 layout_params.record_header_dataset_name = "TimeSliceHeader";
61 layout_params.raw_data_group_name = "RawData";
62 layout_params.view_group_name = "Views";
63 layout_params.path_params_list = { params };
64
65 // Create pointer to a new output HDF5 file
66 std::unique_ptr<hdf5libs::HDF5RawDataFile> output_file =
67 std::make_unique<hdf5libs::HDF5RawDataFile>(output_filename,
68 _frags.begin()->second[0]->get_run_number(),
69 0,
70 "emulate_from_tpstream",
71 layout_params,
72 _sourceid_geoid_map);
73
74 // Iterate over the time slices & save all the fragments
75 for (auto& [slice_id, vec_frags] : _frags) {
76 // Create a new timeslice header
78 tsh.timeslice_number = slice_id;
79 tsh.run_number = vec_frags[0]->get_run_number();
81
82 // Create a new timeslice
84 if (!_quiet) {
85 std::cout << "Time slice number: " << slice_id << std::endl;
86 }
87 // Add the fragments to the timeslices
88 for (std::unique_ptr<daqdataformats::Fragment>& frag_ptr : vec_frags) {
89 if (!_quiet) {
90 std::cout << " Writing elementid: " << frag_ptr->get_element_id()
91 << " trigger number: " << frag_ptr->get_trigger_number()
92 << " trigger_timestamp: " << frag_ptr->get_trigger_timestamp()
93 << " window_begin: " << frag_ptr->get_window_begin()
94 << " sequence_no: " << frag_ptr->get_sequence_number() << std::endl;
95 }
96 ts.add_fragment(std::move(frag_ptr));
97 }
98
99 // Write the timeslice to output file
100 output_file->write(ts);
101 }
102}
103
117std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>>
118sort_files_per_writer(const std::vector<std::string>& _files)
119{
120 std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>> files_sorted;
121
122 // Put files into per-writer app groups
123 for (const std::string& file : _files) {
124 std::shared_ptr<hdf5libs::HDF5RawDataFile> fl = std::make_shared<hdf5libs::HDF5RawDataFile>(file);
125 if (!fl->is_timeslice_type()) {
126 throw std::runtime_error(fmt::format("ERROR: input file '{}' not of type 'TimeSlice'", file));
127 }
128
129 std::string application_name = fl->get_attribute<std::string>("application_name");
130 files_sorted[application_name].push_back(fl);
131 }
132
133 // Sort files for each writer application individually
134 for (auto& [app_name, vec_files] : files_sorted) {
135 // Don't sort if we have 0 or 1 files in the application...
136 if (vec_files.size() <= 1) {
137 continue;
138 }
139
140 // Sort w.r.t. file index attribute
141 std::sort(
142 vec_files.begin(),
143 vec_files.end(),
144 [](const std::shared_ptr<hdf5libs::HDF5RawDataFile>& a, const std::shared_ptr<hdf5libs::HDF5RawDataFile>& b) {
145 return a->get_attribute<size_t>("file_index") < b->get_attribute<size_t>("file_index");
146 });
147 }
148
149 return files_sorted;
150};
151
172std::pair<uint64_t, uint64_t>
174 const std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>>& _files,
175 bool _quiet)
176{
177 if (_files.empty()) {
178 throw std::runtime_error("No files provided");
179 }
180
181 uint64_t global_start = std::numeric_limits<uint64_t>::min();
182 uint64_t global_end = std::numeric_limits<uint64_t>::max();
183
184 // Get the min & max record id per application, and the global
185 for (auto& [appid, vec_files] : _files) {
186 uint64_t app_start = std::numeric_limits<uint64_t>::max();
187 uint64_t app_end = std::numeric_limits<uint64_t>::min();
188
189 // Find min / max record id for this application
190 for (auto& file : vec_files) {
191 auto record_ids = file->get_all_record_ids();
192 if (record_ids.empty()) {
193 throw std::runtime_error(fmt::format("File from application {} contains no records.", appid));
194 }
195 app_start = std::min(app_start, record_ids.begin()->first);
196 app_end = std::max(app_end, record_ids.rbegin()->first);
197 }
198
199 // Update the global min / max record id
200 global_start = std::max(global_start, app_start);
201 global_end = std::min(global_end, app_end);
202
203 if (!_quiet) {
204 std::cout << "Application: " << appid << " " << " TimeSliceID start: " << app_start << " end: " << app_end
205 << std::endl;
206 }
207 }
208
209 if (!_quiet) {
210 std::cout << "Global start: " << global_start << " Global end: " << global_end << std::endl;
211 }
212 if (global_start > global_end) {
213 throw std::runtime_error("One of the provided files' id range did not overlap with the rest. Please select files "
214 "with overlapping TimeSlice IDs");
215 }
216
217 // Extra validation / error handling
218 for (auto& [appid, vec_files] : _files) {
219 for (auto& file : vec_files) {
220 auto record_ids = file->get_all_record_ids();
221
222 uint64_t file_start = record_ids.begin()->first;
223 uint64_t file_end = record_ids.rbegin()->first;
224 if ((file_start > global_end || file_end < global_start)) {
225 uint64_t file_index = file->get_attribute<size_t>("file_index");
226 throw std::runtime_error(fmt::format("File from TPStreamWrite application '{}' (index '{}') has record id "
227 "range [{}, {}], which does not overlap with global range [{}, {}].",
228 appid,
229 file_index,
230 file_start,
231 file_end,
232 global_start,
233 global_end));
234 }
235 }
236 }
237
238 return { global_start, global_end };
239}
240
245{
247 std::vector<std::string> input_files;
249 std::string output_filename;
251 std::string config_name;
253 bool quiet = false;
256 bool latencies = false;
258 bool run_parallel = false;
260 std::optional<uint64_t> slice_start = std::nullopt;
262 std::optional<uint64_t> num_slices = std::nullopt;
263};
264
271void
272parse_app(CLI::App& _app, Options& _opts)
273{
274 _app.add_option("-i,--input-files", _opts.input_files, "List of input files (required)")
275 ->required()
276 ->check(CLI::ExistingFile); // Validate that each file exists
277
278 _app.add_option("-o,--output-file", _opts.output_filename, "Output file (required)")
279 ->required(); // make the argument required
280
281 _app
282 .add_option("-j,--json-config", _opts.config_name, "Trigger Activity and Candidate config JSON to use (required)")
283 ->required()
284 ->check(CLI::ExistingFile);
285
286 _app.add_flag("--parallel", _opts.run_parallel, "Run the TAMakers in parallel");
287
288 _app.add_flag("--quiet", _opts.quiet, "Quiet outputs.");
289
290 _app.add_flag("--latencies", _opts.latencies, "Saves latencies per TP into csv");
291
292 _app.add_option("-s,--slice-start", _opts.slice_start, "Inclusive lower bound for TimeSlice ID to process");
293 _app.add_option("-n,--num-slices", _opts.num_slices, "Number of TimeSlices to process");
294}
295
296int
297main(int argc, char const* argv[])
298{
299 // Do all the CLI processing first
300 CLI::App app{ "Offline trigger TriggerActivity & TriggerCandidate emulatior" };
301 Options opts{};
302
303 parse_app(app, opts);
304
305 try {
306 app.parse(argc, argv);
307 } catch (const CLI::ParseError& e) {
308 return app.exit(e);
309 }
310
311 // Get the configuration file
312 std::ifstream config_stream(opts.config_name);
313 nlohmann::json config = nlohmann::json::parse(config_stream);
314
315 if (!opts.quiet) {
316 std::cout << "Files to process:\n";
317 for (const std::string& file : opts.input_files) {
318 std::cout << "- " << file << "\n";
319 }
320 }
321
322 // Sort the files into a map writer_id::vector<HDF5>
323 std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>> sorted_files =
325
326 // Get the available record_id range
327 std::pair<uint64_t, uint64_t> recordid_range = get_available_slice_id_range(sorted_files, opts.quiet);
328 std::pair<uint64_t, uint64_t> processing_range = recordid_range;
329
330 if (opts.slice_start.has_value()) {
331 processing_range.first = std::max(processing_range.first, opts.slice_start.value());
332 }
333
334 if (opts.num_slices.has_value()) {
335 if (opts.num_slices.value() == 0) {
336 throw std::runtime_error("Invalid --num-slices value: must be greater than 0");
337 }
338
339 uint64_t requested_end = processing_range.first + (opts.num_slices.value() - 1);
340 if (requested_end < processing_range.first) {
341 requested_end = std::numeric_limits<uint64_t>::max();
342 }
343
344 processing_range.second = std::min(processing_range.second, requested_end);
345 }
346
347 if (processing_range.first > processing_range.second) {
348 throw std::runtime_error(fmt::format(
349 "Requested timeslice subset (slice-start={}, num-slices={}) does not overlap available range [{}, {}]",
350 opts.slice_start.value_or(0),
351 opts.num_slices.value_or(0),
352 recordid_range.first,
353 recordid_range.second));
354 }
355
356 if (!opts.quiet) {
357 std::cout << "Processing TimeSliceID range: [" << processing_range.first << ", " << processing_range.second << "]"
358 << std::endl;
359 }
360
361 // Create the file handlers
362 std::vector<std::unique_ptr<TAEmulationWorker>> ta_emu_workers;
363 for (auto [name, files] : sorted_files) {
364 ta_emu_workers.push_back(
365 std::make_unique<TAEmulationWorker>(files, config, processing_range, opts.run_parallel, opts.quiet));
366 }
367
368 // Start each file handler
369 for (const auto& handler : ta_emu_workers) {
370 handler->start_processing();
371 }
372
373 // Output map of TA vectors & function that appends TAs to that vector
374 std::map<uint64_t, std::vector<triggeralgs::TriggerActivity>> tas;
375 auto append_tas = [&tas](std::map<uint64_t, std::vector<triggeralgs::TriggerActivity>>&& _tas) {
376 for (auto& [sliceid, src_vec] : _tas) {
377 auto& dest_vec = tas[sliceid];
378 dest_vec.reserve(dest_vec.size() + src_vec.size());
379
380 dest_vec.insert(dest_vec.end(), std::make_move_iterator(src_vec.begin()), std::make_move_iterator(src_vec.end()));
381 src_vec.clear();
382 }
383 };
384
385 // Output map of Fragment vectors & function that appends TAs to that vector
386 std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>> frags;
387 auto append_frags = [&frags](std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>>&& _frags) {
388 for (auto& [sliceid, src_vec] : _frags) {
389 auto& dest_vec = frags[sliceid];
390
391 dest_vec.reserve(dest_vec.size() + src_vec.size());
392
393 dest_vec.insert(dest_vec.end(), std::make_move_iterator(src_vec.begin()), std::make_move_iterator(src_vec.end()));
394 src_vec.clear();
395 }
396 };
397
398 // Iterate over the handlers, wait for them to complete their job & append
399 // their TAs to our vector when ready.
401 for (const auto& handler : ta_emu_workers) {
402 // Wait for all TAs to be made
403 handler->wait_to_complete_work();
404
405 // Append output TA/fragments
406 append_tas(std::move(handler->get_tas()));
407 append_frags(std::move(handler->get_frags()));
408
409 // Get and merge the source ID map
410 hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t map = handler->get_sourceid_geoid_map();
411 sourceid_geoid_map.insert(map.begin(), map.end());
412 }
413
414 // Sort the TAs in each slice before pushing them into the TCMaker
415 size_t n_tas = 0;
416 for (auto& [sliceid, vec_tas] : tas) {
417 std::sort(
418 vec_tas.begin(), vec_tas.end(), [](const triggeralgs::TriggerActivity& a, const triggeralgs::TriggerActivity& b) {
419 return std::tie(a.time_start, a.channel_start, a.time_end) <
420 std::tie(b.time_start, b.channel_start, b.time_end);
421 });
422 n_tas += vec_tas.size();
423 }
424
425 if (!opts.quiet) {
426 std::cout << "Total number of TAs made: " << n_tas << std::endl;
427 std::cout << "Creating a TCMaker..." << std::endl;
428 }
429 // Create the TC emulator
430 std::string algo_name = config["trigger_candidate_plugin"][0];
431 nlohmann::json algo_config = config["trigger_candidate_config"][0];
432
433 std::unique_ptr<triggeralgs::TriggerCandidateMaker> tc_maker =
435 tc_maker->configure(algo_config);
436
437 trgtools::TCEmulationUnit tc_emulator;
438 tc_emulator.set_maker(tc_maker);
439
440 // Emulate the TriggerCandidates
441 std::vector<triggeralgs::TriggerCandidate> tcs;
442 for (auto& [sliceid, vec_tas] : tas) {
443 std::unique_ptr<daqdataformats::Fragment> tc_frag = tc_emulator.emulate_vector(vec_tas);
444 if (!tc_frag) {
445 continue;
446 }
447
448 // Manipulate the fragment header
449 daqdataformats::FragmentHeader frag_hdr = tc_frag->get_header();
450 frag_hdr.element_id =
451 daqdataformats::SourceID{ daqdataformats::SourceID::Subsystem::kTrigger, tc_frag->get_element_id().id + 10000 };
452
453 tc_frag->set_header_fields(frag_hdr);
455 tc_frag->set_trigger_number(sliceid);
456 tc_frag->set_window_begin(frags[sliceid][0]->get_window_begin());
457 tc_frag->set_window_end(frags[sliceid][0]->get_window_end());
458
459 // Push the fragment to our list
460 frags[sliceid].push_back(std::move(tc_frag));
461
462 std::vector<triggeralgs::TriggerCandidate> tmp_tcs = tc_emulator.get_last_output_buffer();
463 tcs.reserve(tcs.size() + tmp_tcs.size());
464 tcs.insert(tcs.end(), std::make_move_iterator(tmp_tcs.begin()), std::make_move_iterator(tmp_tcs.end()));
465 }
466 if (!opts.quiet) {
467 std::cout << "Total number of TCs made: " << tcs.size() << std::endl;
468 }
469
470 if (!frags.empty()) {
471 save_fragments(opts.output_filename, sourceid_geoid_map, std::move(frags), opts.quiet);
472 } else {
473 std::cout << "No TA/TC fragments generated. Output file will not be generated" << std::endl;
474 }
475
476 return 0;
477}
C++ Representation of a DUNE TimeSlice, consisting of a TimeSliceHeader object and a vector of pointe...
Definition TimeSlice.hpp:27
void add_fragment(std::unique_ptr< Fragment > &&fragment)
Add a Fragment pointer to the Fragments vector.
Definition TimeSlice.hpp:67
std::map< daqdataformats::SourceID, std::vector< uint64_t > > source_id_geo_id_map_t
std::unique_ptr< daqdataformats::Fragment > emulate_vector(const std::vector< input_t > &inputs)
std::vector< output_t > get_last_output_buffer()
void set_maker(std::unique_ptr< maker_t > &maker)
static std::shared_ptr< AbstractFactory< TriggerCandidateMaker > > get_instance()
int main(int argc, char *argv[])
std::pair< uint64_t, uint64_t > get_available_slice_id_range(const std::map< std::string, std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > > &_files, bool _quiet)
Compute the common SliceID interval shared by all writer groups.
std::map< std::string, std::vector< std::shared_ptr< hdf5libs::HDF5RawDataFile > > > sort_files_per_writer(const std::vector< std::string > &_files)
Group and order input TimeSlice files per TPStream writer application.
void parse_app(CLI::App &_app, Options &_opts)
Adds options to our CLI application.
void save_fragments(const std::string &_outputfilename, const hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t &_sourceid_geoid_map, std::map< uint64_t, std::vector< std::unique_ptr< daqdataformats::Fragment > > > _frags, bool _quiet)
Saves fragments in timeslice format into output HDF5 file.
The DUNE-DAQ namespace.
Struct with available cli application options.
std::string output_filename
output filename
std::optional< uint64_t > num_slices
optional number of TimeSlices to process
std::string config_name
the configuration filename
std::vector< std::string > input_files
vector of input filenames
bool quiet
do we want to quiet down the cout? Default: no
bool run_parallel
runs each TAMaker on a separate thread
std::optional< uint64_t > slice_start
optional lower bound (inclusive) for TimeSlice IDs to process
bool latencies
do we want to measure latencies? Default: no
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
Data fields associated with a TimeSliceHeader.
timeslice_number_t timeslice_number
Slice number of this TimeSlice within the stream.