DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
emulate_from_tpstream.cxx File Reference
#include "trgtools/TAEmulationWorker.hpp"
#include "trgtools/TCEmulationUnit.hpp"
#include "CLI/App.hpp"
#include "CLI/Config.hpp"
#include "CLI/Formatter.hpp"
#include <filesystem>
#include <fmt/chrono.h>
#include <fmt/core.h>
#include <fmt/format.h>
#include <optional>
#include "hdf5libs/HDF5RawDataFile.hpp"
#include "hdf5libs/HDF5SourceIDHandler.hpp"
#include "triggeralgs/TriggerCandidateFactory.hpp"
Include dependency graph for emulate_from_tpstream.cxx:

Go to the source code of this file.

Classes

struct  Options
 Struct with available cli application options. More...

Functions

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.
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.
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.
void parse_app (CLI::App &_app, Options &_opts)
 Adds options to our CLI application.
int main (int argc, char const *argv[])

Function Documentation

◆ get_available_slice_id_range()

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.

For each writer application, this function scans all of its files and finds the minimum and maximum available record IDs. It then computes the global intersection across applications:

  • global_start = max(all per-application starts)
  • global_end = min(all per-application ends)

The returned range is therefore the SliceID window for which data is expected to be available from every application.

Todo
: Rather than returning the SliceID range, should try to return a time range and have processors use that.
Parameters
_filesMap of writer application names to vectors of input HDF5 files.
_quietIf false, print per-application and global range diagnostics.
Returns
Inclusive pair {global_start, global_end} of overlapping SliceIDs.
Exceptions
std::runtime_errorIf _files is empty, if any file has no records, or if no overlapping SliceID interval exists.

Definition at line 173 of file emulate_from_tpstream.cxx.

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}

◆ main()

int main ( int argc,
char const * argv[] )

Definition at line 297 of file emulate_from_tpstream.cxx.

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);
454 tc_frag->set_type(daqdataformats::FragmentType::kTriggerCandidate);
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}
std::map< daqdataformats::SourceID, std::vector< uint64_t > > source_id_geo_id_map_t
static std::shared_ptr< AbstractFactory< TriggerCandidateMaker > > get_instance()
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.
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

◆ parse_app()

void parse_app ( CLI::App & _app,
Options & _opts )

Adds options to our CLI application.

Parameters
_appCLI application
_optsStruct with the available options

Definition at line 272 of file emulate_from_tpstream.cxx.

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}
bool latencies
do we want to measure latencies? Default: no

◆ save_fragments()

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.

Parameters
_outputfilenameName of the hdf5 file to save the output into
_sourceid_geoid_mapsourceid–geoid map required to create a HDF5 file
_fragsMap of fragments, with a vector of fragment pointers for each slice id
_quietDo we want to quiet down the cout?
Todo
Maybe in the future we will want to emulate PDS TPs.

Definition at line 33 of file emulate_from_tpstream.cxx.

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
44 hdf5libs::HDF5FileLayoutParameters layout_params;
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
77 daqdataformats::TimeSliceHeader tsh;
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}
C++ Representation of a DUNE TimeSlice, consisting of a TimeSliceHeader object and a vector of pointe...
Definition TimeSlice.hpp:27
SourceID is a generalized representation of the source of a piece of data in the DAQ....
Definition SourceID.hpp:32

◆ sort_files_per_writer()

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.

Opens each input file, validates that it is of TimeSlice type, and groups files by their application_name attribute. Within each application group, files are sorted by their file_index attribute so consecutive segments from the same writer are processed in order.

Parameters
_filesVector of input HDF5 file paths.
Returns
Map keyed by application_name, each value a vector of file handles ordered by file_index.
Exceptions
std::runtime_errorIf any input file is not of type TimeSlice.

Definition at line 118 of file emulate_from_tpstream.cxx.

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};