DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TAEmulationWorker.hpp
Go to the documentation of this file.
1#ifndef TRGTOOLS_TAEMULATIONWORKER_HPP_
2#define TRGTOOLS_TAEMULATIONWORKER_HPP_
3
5
6#include "CLI/App.hpp"
7#include "CLI/Config.hpp"
8#include "CLI/Formatter.hpp"
9
10#include <filesystem>
11#include <fmt/chrono.h>
12#include <fmt/core.h>
13#include <fmt/format.h>
14
19
20namespace dunedaq::trgtools {
21
23{
24public:
37 TAEmulationWorker(std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>> _input_files,
38 nlohmann::json _config,
39 std::pair<uint64_t, uint64_t> _sliceid_range,
40 bool _run_parallel,
41 bool _quiet);
42
43 ~TAEmulationWorker() = default;
44
51 std::vector<daqdataformats::SourceID> get_valid_sourceids(daqdataformats::TimeSlice& _timeslice);
52
54 void start_processing();
55
58
64 std::map<uint64_t, std::vector<triggeralgs::TriggerActivity>> get_tas();
65
71 std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>> get_frags();
72
79
80private:
82 void process_tasks();
83
93 uint64_t _rec,
95 std::vector<trgdataformats::TriggerPrimitive>&& _tps);
96
98 void worker_thread();
99
105 void enqueue_task(std::function<void()> task);
106
109
110private:
112 std::vector<std::shared_ptr<hdf5libs::HDF5RawDataFile>> m_input_files;
113
115 nlohmann::json m_configuration;
116
118 std::vector<std::string> m_input_paths;
119
121 std::map<daqdataformats::SourceID, std::unique_ptr<trgtools::TAEmulationUnit>> m_ta_emulators;
122
124 std::pair<uint64_t, uint64_t> m_sliceid_range;
125
127 const bool m_run_parallel;
128
130 const bool m_quiet;
131
132 /*
133 * Threading objects for the main file handler
134 */
135
137 std::thread m_main_thread;
138
140 std::atomic<bool> m_stop{ false };
141
142 /*
143 * Optional threading objects for the tasks
144 * i.e. one thread per TAMaker
145 */
146
148 std::mutex m_savetps_mutex;
149 std::vector<std::thread> m_thread_pool;
150 std::condition_variable m_condition;
151 std::condition_variable m_task_complete_condition;
152 std::queue<std::function<void()>> m_task_queue;
153 std::mutex m_queue_mutex;
154 std::atomic<size_t> m_active_tasks;
155
157 std::map<uint64_t, std::vector<triggeralgs::TriggerActivity>> m_tas;
158
160 std::map<uint64_t, std::vector<std::unique_ptr<daqdataformats::Fragment>>> m_ta_fragments;
161
163 uint16_t m_id;
165 static uint16_t m_id_next;
167 static const size_t SIZE_TP = sizeof(trgdataformats::TriggerPrimitive);
168};
169
170};
171
172#endif
C++ Representation of a DUNE TimeSlice, consisting of a TimeSliceHeader object and a vector of pointe...
Definition TimeSlice.hpp:27
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::vector< std::string > m_input_paths
input vector of tpstream input paths
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.
nlohmann::json m_configuration
configuration for the TA-makers
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.
The header for a DUNE 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.