DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
NP02ReadoutApplication.cpp
Go to the documentation of this file.
1
10
18#include "confmodel/Session.hpp"
19
23
26
29
33#include "confmodel/GeoId.hpp"
35#include "confmodel/Service.hpp"
36
43
55
57
58#include "logging/Logging.hpp"
59#include <fmt/core.h>
60
61#include <string>
62#include <vector>
63
64// using namespace dunedaq;
65// using namespace dunedaq::appmodel;
66
67namespace dunedaq {
68namespace appmodel {
69
70//-----------------------------------------------------------------------------
71void
72NP02ReadoutApplication::generate_modules(std::shared_ptr<appmodel::ConfigurationHelper> helper) const
73{
74
75 TLOG_DEBUG(6) << "Generating modules for application " << this->UID();
76
77 ConfigObjectFactory obj_fac(this);
78
79 //
80 // Extract basic configuration objects
81 //
82
83 // Data reader
84 auto reader_conf = get_data_reader();
85 if (reader_conf == 0) {
86 throw(BadConf(ERS_HERE, "No DataReaderModule configuration given"));
87 }
88 std::string reader_class = reader_conf->get_template_for();
89
90 // Link handler
91 auto dlh_conf = get_link_handler();
92 // What is template for?
93 auto dlh_class = dlh_conf->get_template_for();
94
95 auto tph_conf = get_tp_handler();
96 if (tph_conf == nullptr && get_tp_generation_enabled()) {
97 throw(BadConf(ERS_HERE, "TP generation is enabled but there is no TP data handler configuration"));
98 }
99
100 std::string tph_class = "";
101 if (tph_conf != nullptr && get_tp_generation_enabled()) {
102 tph_class = tph_conf->get_template_for();
103 }
104
105 //
106 // Process the queue rules looking for inputs to our DL/TP handler modules
107 //
108 const QueueDescriptor* dlh_reqinput_qdesc = nullptr;
109 const QueueDescriptor* tp_input_qdesc = nullptr;
110 // const QueueDescriptor* tpReqInputQDesc = nullptr;
111 const QueueDescriptor* fa_output_qdesc = nullptr;
112
113 for (auto rule : get_queue_rules()) {
114 auto destination_class = rule->get_destination_class();
115 auto data_type = rule->get_descriptor()->get_data_type();
116 // Why datahander here?
117 if (destination_class == "DataHandlerModule" || destination_class == dlh_class || destination_class == tph_class) {
118 if (data_type == "DataRequest") {
119 dlh_reqinput_qdesc = rule->get_descriptor();
120 } else if ((data_type == "TriggerPrimitive" || data_type == "TriggerPrimitiveVector") &&
122 tp_input_qdesc = rule->get_descriptor();
123 }
124 } else if (destination_class == "FragmentAggregatorModule") {
125 fa_output_qdesc = rule->get_descriptor();
126 }
127 }
128
129 //
130 // Process the network rules looking for the Fragment Aggregator and TP handler data reuest inputs
131 //
132 const NetworkConnectionDescriptor* fa_net_desc = nullptr;
133 const NetworkConnectionDescriptor* tp_net_desc = nullptr;
134 const NetworkConnectionDescriptor* ta_net_desc = nullptr;
135 const NetworkConnectionDescriptor* ts_net_desc = nullptr;
136 for (auto rule : get_network_rules()) {
137 auto endpoint_class = rule->get_endpoint_class();
138 auto data_type = rule->get_descriptor()->get_data_type();
139
140 if (endpoint_class == "FragmentAggregatorModule") {
141 fa_net_desc = rule->get_descriptor();
142 } else if (data_type == "TPSet") {
143 tp_net_desc = rule->get_descriptor();
144 } else if (data_type == "TriggerActivity") {
145 ta_net_desc = rule->get_descriptor();
146 } else if (data_type == "TimeSync") {
147 ts_net_desc = rule->get_descriptor();
148 }
149 }
150
151 // Create here the Queue on which all data fragments are forwarded to the fragment aggregator
152 // and a container for the queues of data request to TP handler and DLH
153 if (fa_output_qdesc == nullptr) {
154 throw(BadConf(ERS_HERE, "No fragment output queue descriptor given"));
155 }
156 std::vector<const confmodel::Connection*> req_queues;
157 conffwk::ConfigObject frag_queue_obj = obj_fac.create_queue_obj(fa_output_qdesc);
158
159 //
160 // Get the callback descriptor
161 //
162 const DataMoveCallbackDescriptor* raw_data_callback_desc = get_callback_desc();
163
164 if (raw_data_callback_desc == nullptr) {
165 throw(BadConf(ERS_HERE, "No Raw Data Callback descriptor given"));
166 }
167
168 //
169 // Scan Detector 2 DAQ connections to extract sender, receiver and stream information
170 //
171
172 std::vector<const confmodel::DaqModule*> modules;
173
174 // Loop over the detector to daq connections and generate one data reader per connection
175 // and the cooresponding datalink handlers
176
177 // Collect all streams
178 std::vector<std::pair<int16_t, const confmodel::DetectorStream*>> all_enabled_det_streams;
179 std::map<uint32_t, const appmodel::DataMoveCallbackConf*> callback_confs_by_sid;
180
181 std::vector<const conffwk::ConfigObject*> d2d_conn_objs;
182 uint16_t conn_idx = 0;
183
184 std::set<int16_t> numas;
185 for (auto d2d_conn : get_detector_connections()) {
186 uint16_t receiver_numa = 0;
187
188 // Are we sure?
189 if (helper->is_excluded(d2d_conn)) {
190 TLOG_DEBUG(7) << "Ignoring excluded DetectorToDaqConnection " << d2d_conn->UID();
191 continue;
192 }
193
194 d2d_conn_objs.push_back(&d2d_conn->config_object());
195
196 TLOG_DEBUG(6) << "Processing DetectorToDaqConnection " << d2d_conn->UID();
197 // get the readout groups and the interfaces and streams therein; 1 reaout group corresponds to 1 data reader module
198
199 if (d2d_conn->senders().empty()) {
200 throw(BadConf(ERS_HERE, "DetectorToDaqConnection does not contain sebders or receivers"));
201 }
202 if (d2d_conn->receiver() == nullptr) {
203 throw(BadConf(ERS_HERE, "DetectorToDaqConnection does not contain a receiver"));
204 }
205
206 // Loop over detector 2 daq connections to find senders and receivers
207 auto det_senders = d2d_conn->senders();
208 auto det_receiver = d2d_conn->receiver();
209
210 // Here I want to resolve the type of connection (network, felix, or?)
211 // Rules of engagement: if the receiver interface is network or felix, the receivers should be castable to the
212 // counterpart
213 bool requires_dpdk = (reader_class == "DPDKReaderModule" || reader_class == "FDFakeReaderModule");
214
215 if (reader_class == "DPDKReaderModule" || reader_class == "SocketReaderModule" ||
216 reader_class == "FDFakeReaderModule") {
217 if ((requires_dpdk &&
218 !det_receiver
219 ->cast<appmodel::DPDKReceiver>()) || // SSB: Note here, we are intrinsically locking FakeCard readout to
220 // only emulate DPDK data reception. Given NP02ReadoutApplication is
221 // intended for TDE readout at NP02, assuming this is OK.
222 (reader_class == "SocketReaderModule" && !det_receiver->cast<appmodel::SocketReceiver>())) {
223 std::string required_class = requires_dpdk ? "DPDKReceiver" : "SocketReceiver";
224 throw(BadConf(ERS_HERE,
225 fmt::format("{} requires {}, found {} of class {}",
226 reader_class,
227 required_class,
228 det_receiver->UID(),
229 det_receiver->class_name())));
230 }
231
232 // SSB: Note that here you need to include FDFakeCardReader as well, because emulated readout needs some way to
233 // map NUMA to streams Since we require a receiver in the NetworkDetector2DAQConnections this would still work if
234 // the receiver type is a DPDKReceiver
235 if (reader_class == "DPDKReaderModule" || reader_class == "FDFakeReaderModule") {
236 auto dpdk_reciever = det_receiver->cast<appmodel::DPDKReceiver>();
237 receiver_numa = (int16_t)dpdk_reciever->get_uses()->get_numa_id();
238 TLOG_DEBUG(7) << "receiver numa: " << receiver_numa;
239 }
240
241 bool all_nw_senders = true;
242 for (auto s : det_senders) {
243 all_nw_senders &= (s->cast<appmodel::NWDetDataSender>() != nullptr);
244 }
245
246 // Ensure that all senders are compatible with receiver
247 if (!all_nw_senders) {
248 throw(BadConf(ERS_HERE, "Non-network DetDataSener found with NWreceiver"));
249 }
250 }
251
252 std::vector<const confmodel::DetectorStream*> enabled_det_streams;
253 // Loop over senders
254 for (auto stream : d2d_conn->streams()) {
255
256 // Are we sure?
257 if (helper->is_excluded(stream)) {
258 TLOG_DEBUG(7) << "Ignoring excluded DetectorStream " << stream->UID();
259 continue;
260 }
261
262 // loop over streams
263 all_enabled_det_streams.push_back(std::make_pair(receiver_numa, stream));
264 enabled_det_streams.push_back(stream);
265 numas.insert(receiver_numa);
266 }
267 }
268
269 //-----------------------------------------------------------------
270 //
271 // Create DataReaderModule object
272 //
273
274 //
275 // Instantiate DataReaderModule of type DPDKReaderModule
276 //
277
278 // Create the Data reader object
279
280 std::string reader_uid(fmt::format("datareader-{}-{}", this->UID(), std::to_string(conn_idx++)));
281 TLOG_DEBUG(6) << fmt::format(
282 "creating OKS configuration object for Data reader class {} with id {}", reader_class, reader_uid);
283 auto reader_obj = obj_fac.create(reader_class, reader_uid);
284
285 // Populate configuration and interfaces (leave output queues for later)
286 reader_obj.set_obj("configuration", &reader_conf->config_object());
287 reader_obj.set_objs("connections", d2d_conn_objs);
288
289 // Create the raw data callbacks
290 std::vector<const conffwk::ConfigObject*> raw_data_callback_objs;
291
292 // Create data queues
293 for (auto& [numa, ds] : all_enabled_det_streams) {
294 conffwk::ConfigObject callback_obj = obj_fac.create_callback_sid_obj(raw_data_callback_desc, ds->get_source_id());
295 const auto* callback_conf = obj_fac.get_dal<DataMoveCallbackConf>(callback_obj.UID());
296 raw_data_callback_objs.push_back(&callback_conf->config_object());
297 callback_confs_by_sid[ds->get_source_id()] = callback_conf;
298 }
299
300 reader_obj.set_objs("raw_data_callbacks", raw_data_callback_objs);
301
302 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(reader_obj.UID()));
303
304 //-----------------------------------------------------------------
305 //
306 // Prepare the tp handlers and related queues
307 //
308 std::vector<std::pair<uint32_t, const confmodel::Connection*>> tp_queues;
309
311
312 // Create TP handler object
313 auto tph_conf_obj = tph_conf->config_object();
314 auto tpsrc_ids = get_tp_source_ids();
315
316 if ((tpsrc_ids.size() % 3) > 0) {
317 throw(
318 BadConf(ERS_HERE,
319 fmt::format("number of TP source IDs must be a multiple of 3, current amount: {}", tpsrc_ids.size())));
320 }
321
322 for (auto sid : tpsrc_ids) {
323 conffwk::ConfigObject tp_queue_obj;
324 conffwk::ConfigObject tpreq_queue_obj;
325 std::string tp_uid("tphandler-" + std::to_string(sid->get_sid()));
326 auto tph_obj = obj_fac.create(tph_class, tp_uid);
327 tph_obj.set_by_val<uint32_t>("source_id", sid->get_sid());
328 tph_obj.set_by_val<uint32_t>("detector_id", 1); // 1 == kDAQ
329 tph_obj.set_by_val<bool>("post_processing_enabled", get_ta_generation_enabled());
330 tph_obj.set_obj("module_configuration", &tph_conf_obj);
331
332 // Create the TPs aggregator queue (from RawData Handlers to TP handlers)
333 tp_queue_obj = obj_fac.create_queue_sid_obj(tp_input_qdesc, sid->get_sid());
334 tp_queue_obj.set_by_val<uint32_t>("recv_timeout_ms", 50);
335 tp_queue_obj.set_by_val<uint32_t>("send_timeout_ms", 1);
336
337 tp_queues.push_back(std::make_pair(sid->get_sid(), obj_fac.get_dal<confmodel::Connection>(tp_queue_obj.UID())));
338 // Create tp data requests queue from Fragment Aggregator
339 tpreq_queue_obj = obj_fac.create_queue_sid_obj(dlh_reqinput_qdesc, sid->get_sid());
340 req_queues.push_back(obj_fac.get_dal<confmodel::Connection>(tpreq_queue_obj.UID()));
341
342 // Create the tp(set) publishing service
343 conffwk::ConfigObject tp_net_obj = obj_fac.create_net_obj(tp_net_desc, tp_uid);
344
345 // Create the ta(set) publishing service
346 conffwk::ConfigObject ta_net_obj = obj_fac.create_net_obj(ta_net_desc, tp_uid);
347
348 // Register queues with tp handler
349 tph_obj.set_objs("inputs", { &tp_queue_obj, &tpreq_queue_obj });
350 tph_obj.set_objs("outputs", { &tp_net_obj, &ta_net_obj, &frag_queue_obj });
351 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(tph_obj.UID()));
352 }
353 }
354
355 // Add output queueus of tps
356 std::vector<std::pair<uint32_t, const conffwk::ConfigObject*>> tp_queue_objs;
357 for (auto q : tp_queues) {
358 tp_queue_objs.push_back(std::make_pair(q.first, &q.second->config_object()));
359 }
360
361 //-----------------------------------------------------------------
362 //
363 // Create datalink handlers
364 //
365 // Recover the emulation flag
366
367 auto lb_conf = dlh_conf->get_latency_buffer();
368
369 std::map<int16_t, conffwk::ConfigObject> numa_dhlconf_map;
370 for (int16_t numa : numas) {
371 auto lb_confobj_numa = obj_fac.create(lb_conf->class_name(), fmt::format("{}-numa{}", lb_conf->UID(), numa));
372 lb_confobj_numa.set_by_val<uint32_t>("size", lb_conf->get_size());
373 lb_confobj_numa.set_by_val<bool>("numa_aware", lb_conf->get_numa_aware());
374 lb_confobj_numa.set_by_val<int16_t>("numa_node", numa);
375 lb_confobj_numa.set_by_val<bool>("intrinsic_allocator", lb_conf->get_intrinsic_allocator());
376 lb_confobj_numa.set_by_val<uint32_t>("alignment_size", lb_conf->get_alignment_size());
377 lb_confobj_numa.set_by_val<bool>("preallocation", lb_conf->get_preallocation());
378
379 auto dhl_confobj_numa = obj_fac.create(dlh_conf->class_name(), fmt::format("{}-numa{}", dlh_conf->UID(), numa));
380 dhl_confobj_numa.set_by_val<std::string>("template_for", dlh_conf->get_template_for());
381 dhl_confobj_numa.set_by_val<std::string>("input_data_type", dlh_conf->get_input_data_type());
382 dhl_confobj_numa.set_by_val<bool>("generate_timesync", dlh_conf->get_generate_timesync());
383 dhl_confobj_numa.set_by_val<uint64_t>("post_processing_delay_ticks", dlh_conf->get_post_processing_delay_ticks());
384 dhl_confobj_numa.set_by_val<std::string>("input_data_type", dlh_conf->get_input_data_type());
385 dhl_confobj_numa.set_obj("request_handler", &dlh_conf->get_request_handler()->config_object());
386 dhl_confobj_numa.set_obj("latency_buffer", &lb_confobj_numa);
387 dhl_confobj_numa.set_obj("data_processor", &dlh_conf->get_data_processor()->config_object());
388
389 numa_dhlconf_map[numa] = dhl_confobj_numa;
390 }
391
392 auto emulation_mode = reader_conf->get_emulation_mode();
393 for (auto& [numa, ds] : all_enabled_det_streams) {
394 uint32_t sid = ds->get_source_id();
395 TLOG_DEBUG(6) << fmt::format(
396 "Processing stream {}, id {}, det id {}", ds->UID(), ds->get_source_id(), ds->get_geo_id()->get_detector_id());
397 std::string uid(fmt::format("DLH-{}", sid));
398 TLOG_DEBUG(6) << fmt::format(
399 "creating OKS configuration object for Data Link Handler class {}, if {}", dlh_class, sid);
400 auto dlh_obj = obj_fac.create(dlh_class, uid);
401 dlh_obj.set_by_val<uint32_t>("source_id", sid);
402 dlh_obj.set_by_val<uint32_t>("detector_id", ds->get_geo_id()->get_detector_id());
403 dlh_obj.set_by_val<bool>("post_processing_enabled", get_tp_generation_enabled());
404 dlh_obj.set_by_val<bool>("emulation_mode", emulation_mode);
405 dlh_obj.set_obj("geo_id", &ds->get_geo_id()->config_object());
406 dlh_obj.set_obj("module_configuration", &numa_dhlconf_map[numa]);
407 dlh_obj.set_obj("raw_data_callback", &callback_confs_by_sid[sid]->config_object());
408
409 std::vector<const conffwk::ConfigObject*> dlh_ins, dlh_outs;
410
411 // Create request queue
412 conffwk::ConfigObject req_queue_obj = obj_fac.create_queue_sid_obj(dlh_reqinput_qdesc, ds);
413
414 // Add the requessts queue dal pointer to the outputs of the FragmentAggregatorModule
415 req_queues.push_back(obj_fac.get_dal<confmodel::Connection>(req_queue_obj.UID()));
416 dlh_ins.push_back(&req_queue_obj);
417 dlh_outs.push_back(&frag_queue_obj);
418
419 // Time Sync network connection
420 if (dlh_conf->get_generate_timesync()) {
421 // Add timestamp endpoint
422 conffwk::ConfigObject ts_net_obj = obj_fac.create_net_obj(ts_net_desc, std::to_string(sid));
423 dlh_outs.push_back(&ts_net_obj);
424 }
425
426 // here, we want to select which tp queues to add to the output, to separate mutiple detector elements
427 for (auto tpq : tp_queue_objs) {
428 if ((sid / 100) == (tpq.first / 10)) {
429 dlh_outs.push_back(tpq.second);
430 }
431 }
432 dlh_obj.set_objs("inputs", dlh_ins);
433 dlh_obj.set_objs("outputs", dlh_outs);
434
435 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(dlh_obj.UID()));
436 }
437
438 // Finally create Fragment Aggregator
439 auto aggregator_conf = get_fragment_aggregator();
440 if (aggregator_conf == 0) {
441 throw(BadConf(ERS_HERE, "No FragmentAggregatorModule configuration given"));
442 }
443 std::string faUid("fragmentaggregator-" + UID());
444 // conffwk::ConfigObject frag_aggr;
445 TLOG_DEBUG(7) << "creating OKS configuration object for Fragment Aggregator class ";
446 auto frag_aggr = obj_fac.create("FragmentAggregatorModule", faUid);
447 conffwk::ConfigObject fa_net_obj = obj_fac.create_net_obj(fa_net_desc);
448
449 // Process special Network rules!
450 // Looking for Fragment rules from DFAppplications in current Session
451 std::vector<conffwk::ConfigObject> fragOutObjs;
452 for (auto [uid, descriptor] : helper->get_netdescriptors("Fragment", "DFApplication")) {
453 std::string dreqNetUid(descriptor->get_uid_base() + uid);
454 auto frag_conn = obj_fac.create("NetworkConnection", dreqNetUid);
455
456 frag_conn.set_by_val<std::string>("data_type", descriptor->get_data_type());
457 frag_conn.set_by_val<std::string>("connection_type", descriptor->get_connection_type());
458 // Override capacity, set to 2x expected number of Fragments
459 frag_conn.set_by_val<int>("capacity", all_enabled_det_streams.size() * 2);
460
461 auto serviceObj = descriptor->get_associated_service()->config_object();
462 frag_conn.set_obj("associated_service", &serviceObj);
463 fragOutObjs.push_back(frag_conn);
464 }
465
466 // Add output queueus of data requests and Fragments
467 std::vector<const conffwk::ConfigObject*> fa_output_objs;
468 for (auto& fNet : fragOutObjs) {
469 fa_output_objs.push_back(&fNet);
470 }
471
472 for (auto& q : req_queues) {
473 fa_output_objs.push_back(&q->config_object());
474 }
475
476 frag_aggr.set_obj("configuration", &aggregator_conf->config_object());
477 frag_aggr.set_objs("inputs", { &fa_net_obj, &frag_queue_obj });
478 frag_aggr.set_objs("outputs", fa_output_objs);
479
480 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(frag_aggr.UID()));
481
482 obj_fac.update_modules(modules);
483} // NOLINT
484
485} // namespace appmodel
486} // namespace dunedaq
#define ERS_HERE
void update_modules(const std::vector< const confmodel::DaqModule * > &modules)
conffwk::ConfigObject create_queue_sid_obj(const QueueDescriptor *qdesc, uint32_t src_id) const
conffwk::ConfigObject create_callback_sid_obj(const DataMoveCallbackDescriptor *cdesc, uint32_t src_id) const
conffwk::ConfigObject create_net_obj(const NetworkConnectionDescriptor *ndesc, std::string uid) const
Helper function that gets a network connection config.
const T * get_dal(std::string uid) const
conffwk::ConfigObject create_queue_obj(const QueueDescriptor *qdesc, std::string uid="") const
conffwk::ConfigObject create(const std::string &class_name, const std::string &id) const
const dunedaq::confmodel::NetworkDevice * get_uses() const
Get "uses" relationship value.
void generate_modules(std::shared_ptr< appmodel::ConfigurationHelper >) const override
const dunedaq::appmodel::DataHandlerConf * get_tp_handler() const
Get "tp_handler" relationship value.
const std::vector< const dunedaq::confmodel::DetectorToDaqConnection * > & get_detector_connections() const
Get "detector_connections" relationship value. The list of detector channels to be read out by this r...
bool get_ta_generation_enabled() const
Get "ta_generation_enabled" attribute value.
const dunedaq::appmodel::FragmentAggregatorConf * get_fragment_aggregator() const
Get "fragment_aggregator" relationship value.
const std::vector< const dunedaq::appmodel::SourceIDConf * > & get_tp_source_ids() const
Get "tp_source_ids" relationship value.
const dunedaq::appmodel::DataHandlerConf * get_link_handler() const
Get "link_handler" relationship value.
const dunedaq::appmodel::DataMoveCallbackDescriptor * get_callback_desc() const
Get "callback_desc" relationship value.
const dunedaq::appmodel::DataReaderConf * get_data_reader() const
Get "data_reader" relationship value.
bool get_tp_generation_enabled() const
Get "tp_generation_enabled" attribute value.
const std::vector< const dunedaq::appmodel::NetworkConnectionRule * > & get_network_rules() const
Get "network_rules" relationship value.
const std::vector< const dunedaq::appmodel::QueueConnectionRule * > & get_queue_rules() const
Get "queue_rules" relationship value.
void set_by_val(const std::string &name, T value)
Set attribute value.
void set_objs(const std::string &name, const std::vector< const ConfigObject * > &o, bool skip_non_null_check=false)
Set relationship multi-value.
const std::string & UID() const noexcept
Return object identity.
void set_obj(const std::string &name, const ConfigObject *o, bool skip_non_null_check=false)
Set relationship single-value.
const TARGET * cast() const noexcept
Casts object to different class.
const ConfigObject & config_object() const
const std::string & UID() const noexcept
uint8_t get_numa_id() const
Get "numa_id" attribute value.
conffwk entry point
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
The DUNE-DAQ namespace.
CIB Buffer std::string descriptor Message from std::string descriptor CIB process std::string descriptor descriptor