35 std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>> _frags,
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";
55 std::vector<hdf5libs::HDF5PathParameters> params;
56 params.push_back(params_trigger);
66 std::unique_ptr<hdf5libs::HDF5RawDataFile> output_file =
67 std::make_unique<hdf5libs::HDF5RawDataFile>(output_filename,
68 _frags.begin()->second[0]->get_run_number(),
70 "emulate_from_tpstream",
75 for (
auto& [slice_id, vec_frags] : _frags) {
79 tsh.
run_number = vec_frags[0]->get_run_number();
85 std::cout <<
"Time slice number: " << slice_id << std::endl;
88 for (std::unique_ptr<daqdataformats::Fragment>& frag_ptr : vec_frags) {
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;
100 output_file->write(ts);
120 std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>> files_sorted;
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));
129 std::string application_name = fl->get_attribute<std::string>(
"application_name");
130 files_sorted[application_name].push_back(fl);
134 for (
auto& [app_name, vec_files] : files_sorted) {
136 if (vec_files.size() <= 1) {
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");
174 const std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>>& _files,
177 if (_files.empty()) {
178 throw std::runtime_error(
"No files provided");
181 uint64_t global_start = std::numeric_limits<uint64_t>::min();
182 uint64_t global_end = std::numeric_limits<uint64_t>::max();
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();
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));
195 app_start = std::min(app_start, record_ids.begin()->first);
196 app_end = std::max(app_end, record_ids.rbegin()->first);
200 global_start = std::max(global_start, app_start);
201 global_end = std::min(global_end, app_end);
204 std::cout <<
"Application: " << appid <<
" " <<
" TimeSliceID start: " << app_start <<
" end: " << app_end
210 std::cout <<
"Global start: " << global_start <<
" Global end: " << global_end << std::endl;
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");
218 for (
auto& [appid, vec_files] : _files) {
219 for (
auto& file : vec_files) {
220 auto record_ids = file->get_all_record_ids();
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 [{}, {}].",
238 return { global_start, global_end };
297main(
int argc,
char const* argv[])
300 CLI::App app{
"Offline trigger TriggerActivity & TriggerCandidate emulatior" };
306 app.parse(argc, argv);
307 }
catch (
const CLI::ParseError& e) {
313 nlohmann::json
config = nlohmann::json::parse(config_stream);
316 std::cout <<
"Files to process:\n";
318 std::cout <<
"- " << file <<
"\n";
323 std::map<std::string, std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>>> sorted_files =
328 std::pair<uint64_t, uint64_t> processing_range = recordid_range;
331 processing_range.first = std::max(processing_range.first, opts.
slice_start.value());
336 throw std::runtime_error(
"Invalid --num-slices value: must be greater than 0");
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();
344 processing_range.second = std::min(processing_range.second, requested_end);
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 [{}, {}]",
352 recordid_range.first,
353 recordid_range.second));
357 std::cout <<
"Processing TimeSliceID range: [" << processing_range.first <<
", " << processing_range.second <<
"]"
362 std::vector<std::unique_ptr<TAEmulationWorker>> ta_emu_workers;
363 for (
auto [name, files] : sorted_files) {
364 ta_emu_workers.push_back(
369 for (
const auto& handler : ta_emu_workers) {
370 handler->start_processing();
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());
380 dest_vec.insert(dest_vec.end(), std::make_move_iterator(src_vec.begin()), std::make_move_iterator(src_vec.end()));
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];
391 dest_vec.reserve(dest_vec.size() + src_vec.size());
393 dest_vec.insert(dest_vec.end(), std::make_move_iterator(src_vec.begin()), std::make_move_iterator(src_vec.end()));
401 for (
const auto& handler : ta_emu_workers) {
403 handler->wait_to_complete_work();
406 append_tas(std::move(handler->get_tas()));
407 append_frags(std::move(handler->get_frags()));
411 sourceid_geoid_map.insert(map.begin(), map.end());
416 for (
auto& [sliceid, vec_tas] : tas) {
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);
422 n_tas += vec_tas.size();
426 std::cout <<
"Total number of TAs made: " << n_tas << std::endl;
427 std::cout <<
"Creating a TCMaker..." << std::endl;
430 std::string algo_name =
config[
"trigger_candidate_plugin"][0];
431 nlohmann::json algo_config =
config[
"trigger_candidate_config"][0];
433 std::unique_ptr<triggeralgs::TriggerCandidateMaker> tc_maker =
435 tc_maker->configure(algo_config);
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);
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());
460 frags[sliceid].push_back(std::move(tc_frag));
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()));
467 std::cout <<
"Total number of TCs made: " << tcs.size() << std::endl;
470 if (!frags.empty()) {
473 std::cout <<
"No TA/TC fragments generated. Output file will not be generated" << std::endl;