DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
process_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
19
20using namespace dunedaq;
21
23{
24private:
25 /* data */
26 std::unique_ptr<hdf5libs::HDF5RawDataFile> m_input_file;
27 std::unique_ptr<hdf5libs::HDF5RawDataFile> m_output_file;
28
29 void open_files(std::string input_path, std::string output_path);
30 void close_files();
31
33
34 // Can modify?
36
37public:
38 TimeSliceProcessor(std::string input_path, std::string output_path);
40
41 void set_processor(std::function<void(daqdataformats::TimeSlice&)> processor);
42 void loop(uint64_t num_records = 0, uint64_t offset = 0, bool quiet = false);
43};
44
45//-----------------------------------------------------------------------------
46TimeSliceProcessor::TimeSliceProcessor(std::string input_path, std::string output_path)
47{
48 this->open_files(input_path, output_path);
49}
50
51//-----------------------------------------------------------------------------
56
57//-----------------------------------------------------------------------------
58void
59TimeSliceProcessor::open_files(std::string input_path, std::string output_path)
60{
61 // Open input file
62 m_input_file = std::make_unique<hdf5libs::HDF5RawDataFile>(input_path);
63
64 if (!m_input_file->is_timeslice_type()) {
65 fmt::print("ERROR: input file '{}' not of type 'TimeSlice'\n", input_path);
66 throw std::runtime_error(fmt::format("ERROR: input file '{}' not of type 'TimeSlice'", input_path));
67 }
68
69 auto run_number = m_input_file->get_attribute<daqdataformats::run_number_t>("run_number");
70 auto file_index = m_input_file->get_attribute<size_t>("file_index");
71 auto application_name = m_input_file->get_attribute<std::string>("application_name");
72
73 fmt::print("Run Number: {}\nFile Index: {}\nApp name: '{}'\n", run_number, file_index, application_name);
74
75 if (!output_path.empty()) {
76 // Open output file
77 m_output_file = std::make_unique<hdf5libs::HDF5RawDataFile>(
78 output_path,
79 m_input_file->get_attribute<daqdataformats::run_number_t>("run_number"),
80 m_input_file->get_attribute<size_t>("file_index"),
81 m_input_file->get_attribute<std::string>("application_name"),
82 m_input_file->get_file_layout().get_file_layout_params(),
83 m_input_file->get_srcid_geoid_map());
84 }
85}
86
87//-----------------------------------------------------------------------------
88void
90{
91 // Do something?
92}
93
94//-----------------------------------------------------------------------------
95void
97{
98 m_processor = processor;
99}
100
101//-----------------------------------------------------------------------------
102void
108
109//-----------------------------------------------------------------------------
110void
111TimeSliceProcessor::loop(uint64_t num_records, uint64_t offset, bool quiet)
112{
113
114 // Replace with a record selection?
115 auto records = m_input_file->get_all_record_ids();
116
117 if (!num_records) {
118 num_records = (records.size() - offset);
119 }
120
121 uint64_t first_rec = offset, last_rec = offset + num_records;
122
123 uint64_t i_rec(0);
124 for (const auto& rid : records) {
125
126 if (i_rec < first_rec || i_rec >= last_rec) {
127 ++i_rec;
128 continue;
129 }
130
131 if (!quiet)
132 fmt::print("\n-- Processing TSL {}:{}\n\n", rid.first, rid.second);
133 auto tsl = m_input_file->get_timeslice(rid);
134 // Or filter on a selection here using a lambda?
135
136 // if (!quiet)
137 // fmt::print("TSL number {}\n", tsl.get_header().timeslice_number);
138
139 // Add a process method
140 this->process(tsl);
141
142 if (m_output_file)
143 m_output_file->write(tsl);
144
145 ++i_rec;
146 if (!quiet)
147 fmt::print("\n-- Finished TSL {}:{}\n\n", rid.first, rid.second);
148 }
149}
150
151//-----------------------------------------------------------------------------
152int
153main(int argc, char const* argv[])
154{
155
156 CLI::App app{ "tapipe" };
157 // argv = app.ensure_utf8(argv);
158
159 std::string input_file_path;
160 app.add_option("-i", input_file_path, "Input TPStream file path")->required();
161 std::string output_file_path;
162 app.add_option("-o", output_file_path, "Output TPStream file path");
163 std::string channel_map_name = "VDColdboxTPCChannelMap";
164 app.add_option("-m", channel_map_name, "Detector Channel Map");
165 std::string config_name;
166 app.add_option("-j", config_name, "Trigger Activity and Candidate config JSON to use.")->required();
167 uint64_t skip_rec(0);
168 app.add_option("-s", skip_rec, "Skip records");
169 uint64_t num_rec(0);
170 app.add_option("-n", num_rec, "Process records");
171
172 bool quiet = false;
173 app.add_flag("--quiet", quiet, "Quiet outputs.");
174
175 bool latencies = false;
176 app.add_flag("--latencies", latencies, "Saves latencies per TP into csv");
177 CLI11_PARSE(app, argc, argv);
178
179 if (!quiet)
180 fmt::print("TPStream file: {}\n", input_file_path);
181
182 TimeSliceProcessor rp(input_file_path, output_file_path);
183
184 // TP source id (subsystem)
185 auto tp_subsystem_requirement = daqdataformats::SourceID::Subsystem::kTrigger;
186
187 auto channel_map = dunedaq::detchannelmaps::make_tpc_map(channel_map_name);
188
189 // Read configuration
190 std::ifstream config_stream(config_name);
191 nlohmann::json config = nlohmann::json::parse(config_stream);
192
193 // Only use the first plugin for now.
194 nlohmann::json ta_algo = config["trigger_activity_plugin"][0];
195 nlohmann::json ta_config = config["trigger_activity_config"][0];
196
197 nlohmann::json tc_algo = config["trigger_candidate_plugin"][0];
198 nlohmann::json tc_config = config["trigger_candidate_config"][0];
199
200 // Finally create a TA maker
201 std::unique_ptr<triggeralgs::TriggerActivityMaker> ta_maker =
203 ta_maker->configure(ta_config);
204 std::unique_ptr<trgtools::TAEmulationUnit> ta_emulator = std::make_unique<trgtools::TAEmulationUnit>();
205 ta_emulator->set_maker(ta_maker);
206 // TODO: Use a better file naming scheme for CSV.
207 if (latencies) {
208 std::filesystem::path output_path(output_file_path);
209 ta_emulator->set_timing_file(
210 (output_path.parent_path() / ("ta_timings_" + output_path.stem().string() + ".csv")).string());
211 ta_emulator->write_csv_header("TP Time Start,TP ADC Integral,Time Diffs,Is Last TP In TA");
212 }
213
214 // Finally create a TA maker
215 std::unique_ptr<triggeralgs::TriggerCandidateMaker> tc_maker =
217 tc_maker->configure(tc_config);
218 std::unique_ptr<trgtools::TCEmulationUnit> tc_emulator = std::make_unique<trgtools::TCEmulationUnit>();
219 tc_emulator->set_maker(tc_maker);
220 // TODO: Use a better file naming scheme for CSV.
221 if (latencies) {
222 std::filesystem::path output_path(output_file_path);
223 tc_emulator->set_timing_file(
224 (output_path.parent_path() / ("tc_timings_" + output_path.stem().string() + ".csv")).string());
225 tc_emulator->write_csv_header("Time Diffs");
226 }
227
228 // Generic filter hook
229 std::function<bool(const trgdataformats::TriggerPrimitive&)> tp_filter;
230
231 auto z_plane_filter = [&](const trgdataformats::TriggerPrimitive& tp) -> bool {
232 return (channel_map->get_plane_from_offline_channel(tp.channel) != 2);
233 };
234
235 tp_filter = z_plane_filter;
236
237 rp.set_processor([&](daqdataformats::TimeSlice& tsl) -> void {
238 const std::vector<std::unique_ptr<daqdataformats::Fragment>>& frags = tsl.get_fragments_ref();
239 const size_t num_frags = frags.size();
240 if (!quiet)
241 fmt::print("The number of fragments: {}\n", num_frags);
242
243 uint64_t average_ta_time = 0;
244 uint64_t average_tc_time = 0;
245
246 size_t num_tas = 0;
247 size_t num_tcs = 0;
248
249 // Need a static for-loop: adding fragments to tsl will mutate frags even though it's const.
250 for (size_t i = 0; i < num_frags; i++) {
251 const auto& frag = frags[i];
252
253 // The fragment has to be for the trigger (not e.g. for retreival from readout)
254 if (frag->get_element_id().subsystem != tp_subsystem_requirement) {
255 if (!quiet)
256 fmt::print(" Warning, got non kTrigger SourceID {}\n", frag->get_element_id().to_string());
257 continue;
258 }
259
260 // The fragment has to be TriggerPrimitive
261 if (frag->get_fragment_type() != daqdataformats::FragmentType::kTriggerPrimitive) {
262 if (!quiet)
263 fmt::print(" Error: FragmentType is: {}!\n", fragment_type_to_string(frag->get_fragment_type()));
264 continue;
265 }
266
267 // This bit should be outside the loop
268 if (!quiet)
269 fmt::print(" Fragment id: {} [{}]\n",
270 frag->get_element_id().to_string(),
271 daqdataformats::fragment_type_to_string(frag->get_fragment_type()));
272
273 // Pull tps out
274 size_t n_tps = frag->get_data_size() / sizeof(trgdataformats::TriggerPrimitive);
275 if (!quiet) {
276 fmt::print(" TP fragment size: {}\n", frag->get_data_size());
277 fmt::print(" Num TPs: {}\n", n_tps);
278 }
279
280 // Create a TP buffer
281 std::vector<trgdataformats::TriggerPrimitive> tp_buffer;
282 // Prepare the TP buffer, checking for time ordering
283 tp_buffer.reserve(tp_buffer.size() + n_tps);
284
285 // Populate the TP buffer
286 trgdataformats::TriggerPrimitive* tp_array = static_cast<trgdataformats::TriggerPrimitive*>(frag->get_data());
287 uint64_t last_ts = 0;
288 for (size_t i(0); i < n_tps; ++i) {
289 auto& tp = tp_array[i];
290 if (tp.time_start <= last_ts && !quiet) {
291 fmt::print(" ERROR: {} {} ", +tp.time_start, last_ts);
292 }
293 tp_buffer.push_back(tp);
294 }
295
296 // Print some useful info
297 uint64_t d_ts = tp_array[n_tps - 1].time_start - tp_array[0].time_start;
298 if (!quiet)
299 fmt::print(" TS gap: {} {} ms\n", d_ts, d_ts * 16.0 / 1'000'000);
300
301 //
302 // TA Processing
303 //
304
305 const auto ta_start = std::chrono::steady_clock::now();
306 std::unique_ptr<daqdataformats::Fragment> ta_frag = ta_emulator->emulate_vector(tp_buffer);
307 const auto ta_end = std::chrono::steady_clock::now();
308
309 if (ta_frag == nullptr) // Buffer was empty.
310 continue;
311 num_tas += ta_emulator->get_last_output_buffer().size();
312
313 // TA time calculation.
314 const uint64_t ta_diff = std::chrono::nanoseconds(ta_end - ta_start).count();
315 average_ta_time += ta_diff;
316 if (!quiet) {
317 fmt::print("\tTA Time Process: {} ns.\n", ta_diff);
318 }
319
320 daqdataformats::FragmentHeader frag_hdr = frag->get_header();
321
322 // Customise the source id (add 1000 to id)
323 frag_hdr.element_id =
325
326 ta_frag->set_header_fields(frag_hdr);
328
329 tsl.add_fragment(std::move(ta_frag));
330 //
331 // TA Processing Ends
332 //
333
334 //
335 // TC Processing
336 //
337
338 std::vector<triggeralgs::TriggerActivity> ta_buffer = ta_emulator->get_last_output_buffer();
339 const auto tc_start = std::chrono::steady_clock::now();
340 std::unique_ptr<daqdataformats::Fragment> tc_frag = tc_emulator->emulate_vector(ta_buffer);
341 const auto tc_end = std::chrono::steady_clock::now();
342
343 if (tc_frag == nullptr) // Buffer was empty.
344 continue;
345 num_tcs += tc_emulator->get_last_output_buffer().size();
346
347 // TC time calculation.
348 const uint64_t tc_diff = std::chrono::nanoseconds(tc_end - tc_start).count();
349 average_tc_time += tc_diff;
350 if (!quiet) {
351 fmt::print("\tTC Time Process: {} ns.\n", tc_diff);
352 }
353
354 // Shares the same frag_hdr.
355 tc_frag->set_header_fields(frag_hdr);
357
358 tsl.add_fragment(std::move(tc_frag));
359
360 } // Fragment for loop
361
362 if (num_tas == 0)
363 average_ta_time = 0;
364 else
365 average_ta_time /= num_tas;
366 if (num_tcs == 0)
367 average_tc_time = 0;
368 else
369 average_tc_time /= num_tcs;
370 if (!quiet) {
371 fmt::print("\t\tAverage TA Time Process ({} TAs): {} ns.\n", num_tas, average_ta_time);
372 fmt::print("\t\tAverage TC Time Process ({} TCs): {} ns.\n", num_tcs, average_tc_time);
373 }
374 });
375
376 rp.loop(num_rec, skip_rec);
377
378 /* code */
379 return 0;
380}
void set_processor(std::function< void(daqdataformats::TimeSlice &)> processor)
void process(daqdataformats::TimeSlice &tls)
std::unique_ptr< hdf5libs::HDF5RawDataFile > m_input_file
TimeSliceProcessor(std::string input_path, std::string output_path)
std::unique_ptr< hdf5libs::HDF5RawDataFile > m_output_file
void open_files(std::string input_path, std::string output_path)
std::function< void(daqdataformats::TimeSlice &)> m_processor
void loop(uint64_t num_records=0, uint64_t offset=0, bool quiet=false)
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
const std::vector< std::unique_ptr< Fragment > > & get_fragments_ref() const
Get a handle to the Fragments.
Definition TimeSlice.hpp:55
static std::shared_ptr< AbstractFactory< TriggerActivityMaker > > get_instance()
double offset
int main(int argc, char *argv[])
@ kTriggerPrimitive
Trigger format TPs produced by trigger code.
std::string fragment_type_to_string(const FragmentType &type)
The DUNE-DAQ namespace.
Unknown serialization type<< t,((char) t)) template< typename T > inline std::string datatype_to_string() { return "Unknown";} namespace serialization { template< typename T > struct is_serializable :std::false_type {};enum SerializationType { kMsgPack };inline SerializationType from_string(const std::string s) { if(s=="msgpack") return kMsgPack;throw UnknownSerializationTypeString(ERS_HERE, s);} constexpr uint8_t serialization_type_byte(SerializationType stype) { switch(stype) { case kMsgPack:return 'M';default:throw UnknownSerializationTypeEnum(ERS_HERE);} } constexpr SerializationType DEFAULT_SERIALIZATION_TYPE=kMsgPack;template< class T > std::vector< uint8_t > serialize(const T &obj, SerializationType stype=DEFAULT_SERIALIZATION_TYPE) { switch(stype) { case kMsgPack:{ msgpack::sbuffer buf;msgpack::pack(buf, obj);std::vector< uint8_t > ret(buf.size()+1);ret[0]=serialization_type_byte(stype);std::copy(buf.data(), buf.data()+buf.size(), ret.begin()+1);return ret;} default:throw UnknownSerializationTypeEnum(ERS_HERE);} } template< class T, typename CharType=unsigned char > T deserialize(const std::vector< CharType > &v) { switch(v[0]) { case serialization_type_byte(kMsgPack):{ try { msgpack::object_handle oh=msgpack::unpack(const_cast< char * >(reinterpret_cast< const char * >(v.data()+1)), v.size() - 1,[](msgpack::type::object_type, std::size_t, void *) -> bool
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.