70std::vector<const confmodel::ExcludableEntity*>
80 TLOG_DEBUG(6) <<
"Generating modules for application " << this->
UID();
91 if (reader_conf == 0) {
92 throw(BadConf(
ERS_HERE,
"No DataReaderModule configuration given"));
94 std::string reader_class = reader_conf->get_template_for();
99 auto dlh_class = dlh_conf->get_template_for();
103 throw(BadConf(
ERS_HERE,
"TP generation is enabled but there is no TP data handler configuration"));
106 std::string tph_class =
"";
108 tph_class = tph_conf->get_template_for();
114 const QueueDescriptor* dlh_reqinput_qdesc =
nullptr;
115 const QueueDescriptor* tp_input_qdesc =
nullptr;
117 const QueueDescriptor* fa_output_qdesc =
nullptr;
120 auto destination_class = rule->get_destination_class();
121 auto data_type = rule->get_descriptor()->get_data_type();
124 if (destination_class ==
"DataHandlerModule" || destination_class == dlh_class || destination_class == tph_class) {
125 if (data_type ==
"DataRequest") {
126 dlh_reqinput_qdesc = rule->get_descriptor();
127 }
else if ((data_type ==
"TriggerPrimitive" || data_type ==
"TriggerPrimitiveVector") &&
129 tp_input_qdesc = rule->get_descriptor();
131 }
else if (destination_class ==
"FragmentAggregatorModule") {
132 fa_output_qdesc = rule->get_descriptor();
136 if (dlh_reqinput_qdesc ==
nullptr) {
137 throw(BadConf(
ERS_HERE,
"No data link handler request input queue descriptor given"));
143 const NetworkConnectionDescriptor* fa_net_desc =
nullptr;
144 const NetworkConnectionDescriptor* tp_net_desc =
nullptr;
145 const NetworkConnectionDescriptor* ta_net_desc =
nullptr;
146 const NetworkConnectionDescriptor* ts_net_desc =
nullptr;
148 auto endpoint_class = rule->get_endpoint_class();
149 auto data_type = rule->get_descriptor()->get_data_type();
151 if (endpoint_class ==
"FragmentAggregatorModule") {
152 fa_net_desc = rule->get_descriptor();
153 }
else if (data_type ==
"TPSet") {
154 tp_net_desc = rule->get_descriptor();
155 }
else if (data_type ==
"TriggerActivity") {
156 ta_net_desc = rule->get_descriptor();
157 }
else if (data_type ==
"TimeSync") {
158 ts_net_desc = rule->get_descriptor();
162 if (fa_net_desc ==
nullptr) {
163 throw(BadConf(
ERS_HERE,
"No Fragment Aggregator network descriptor given"));
165 if (ts_net_desc ==
nullptr && dlh_conf->get_generate_timesync()) {
166 throw(BadConf(
ERS_HERE,
"No Time Sync network descriptor given but time sync generation is enabled"));
174 if (raw_data_callback_desc ==
nullptr) {
175 throw(BadConf(
ERS_HERE,
"No Raw Data Callback descriptor given"));
180 if (fa_output_qdesc ==
nullptr) {
181 throw(BadConf(
ERS_HERE,
"No fragment output queue descriptor given"));
183 std::vector<const confmodel::Connection*> req_queues;
184 conffwk::ConfigObject frag_queue_obj = obj_fac.create_queue_obj(fa_output_qdesc);
190 std::vector<const confmodel::DaqModule*> modules;
196 std::vector<const confmodel::DetectorStream*> all_enabled_det_streams;
197 std::map<uint32_t, const appmodel::DataMoveCallbackConf*> callback_confs_by_sid;
200 uint16_t conn_idx = 0;
203 if (helper->is_excluded(d2d_conn)) {
204 TLOG_DEBUG(7) <<
"Ignoring excluded DetectorToDaqConnection " << d2d_conn->UID();
208 TLOG_DEBUG(6) <<
"Processing DetectorToDaqConnection " << d2d_conn->UID();
212 if (d2d_conn->senders().empty()) {
213 throw(BadConf(
ERS_HERE,
"DetectorToDaqConnection does not contain senders"));
215 if (d2d_conn->receiver() ==
nullptr) {
216 throw(BadConf(
ERS_HERE,
"DetectorToDaqConnection does not contain a receiver"));
220 auto det_senders = d2d_conn->senders();
221 auto det_receiver = d2d_conn->receiver();
223 std::vector<const confmodel::DetectorStream*> enabled_det_streams;
225 for (
auto stream : d2d_conn->streams()) {
228 if (helper->is_excluded(stream)) {
229 TLOG_DEBUG(7) <<
"Ignoring excluded DetectorStream " << stream->UID();
233 all_enabled_det_streams.push_back(stream);
234 enabled_det_streams.push_back(stream);
240 if (reader_class ==
"DPDKReaderModule") {
241 if (!d2d_conn->castable(
"NetworkDetectorToDaqConnection")) {
243 fmt::format(
"{} requires NetworkDetectorToDaqConnection, found {} of class {}",
246 d2d_conn->class_name())));
248 if (!det_receiver->cast<appmodel::DPDKReceiver>()) {
250 fmt::format(
"{} requires NWDetDataReceiver, found {} of class {}",
253 det_receiver->class_name())));
255 }
else if (reader_class ==
"SocketReaderModule") {
256 if (!d2d_conn->castable(
"SocketDetectorToDaqConnection")) {
258 fmt::format(
"{} requires SocketDetectorToDaqConnection, found {} of class {}",
261 d2d_conn->class_name())));
263 if (!det_receiver->cast<appmodel::SocketReceiver>()) {
265 fmt::format(
"{} requires SocketReceiver, found {} of class {}",
268 det_receiver->class_name())));
270 }
else if (reader_class ==
"FelixReaderModule") {
271 if (!d2d_conn->castable(
"FelixDetectorToDaqConnection")) {
273 fmt::format(
"{} requires FelixDetectorToDaqConnection, found {} of class {}",
276 d2d_conn->class_name())));
278 if (!det_receiver->cast<appmodel::FelixDataReceiver>()) {
280 fmt::format(
"FelixReaderModule requires FelixDataReceiver, found {} of class {}",
282 det_receiver->class_name())));
298 std::string reader_uid(fmt::format(
"datareader-{}-{}", this->
UID(), std::to_string(conn_idx++)));
300 "creating OKS configuration object for Data reader class {} with id {}", reader_class, reader_uid);
301 auto reader_obj = obj_fac.create(reader_class, reader_uid);
304 reader_obj.set_obj(
"configuration", &reader_conf->config_object());
305 reader_obj.set_objs(
"connections", { &d2d_conn->config_object() });
308 std::vector<const conffwk::ConfigObject*> raw_data_callback_objs;
311 for (
auto ds : enabled_det_streams) {
312 conffwk::ConfigObject callback_obj = obj_fac.create_callback_sid_obj(raw_data_callback_desc, ds->get_source_id());
313 const auto* callback_conf = obj_fac.get_dal<DataMoveCallbackConf>(callback_obj.UID());
314 raw_data_callback_objs.push_back(&callback_conf->config_object());
315 callback_confs_by_sid[ds->get_source_id()] = callback_conf;
318 reader_obj.set_objs(
"raw_data_callbacks", raw_data_callback_objs);
320 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(reader_obj.UID()));
327 std::vector<const confmodel::Connection*> tp_queues;
329 if (tp_input_qdesc ==
nullptr) {
330 throw(BadConf(
ERS_HERE,
"TP generation is enabled but no TP input queue descriptor given"));
332 if (tp_net_desc ==
nullptr) {
333 throw(BadConf(
ERS_HERE,
"TP generation is enabled but no TPSet network descriptor given"));
335 if (ta_net_desc ==
nullptr) {
336 throw(BadConf(
ERS_HERE,
"TP generation is enabled but no TriggerActivity network descriptor given"));
339 auto tph_conf_obj = tph_conf->config_object();
342 for (
auto sid : tpsrc_ids) {
343 conffwk::ConfigObject tp_queue_obj;
344 conffwk::ConfigObject tpreq_queue_obj;
345 std::string tp_uid(
"tphandler-" + std::to_string(sid->get_sid()));
346 auto tph_obj = obj_fac.create(tph_class, tp_uid);
347 tph_obj.set_by_val<uint32_t>(
"source_id", sid->get_sid());
348 tph_obj.set_by_val<uint32_t>(
"detector_id", 1);
350 tph_obj.set_obj(
"module_configuration", &tph_conf_obj);
353 tp_queue_obj = obj_fac.create_queue_sid_obj(tp_input_qdesc, sid->get_sid());
354 tp_queue_obj.set_by_val<uint32_t>(
"recv_timeout_ms", 50);
355 tp_queue_obj.set_by_val<uint32_t>(
"send_timeout_ms", 1);
357 tp_queues.push_back(obj_fac.get_dal<confmodel::Connection>(tp_queue_obj.UID()));
359 tpreq_queue_obj = obj_fac.create_queue_sid_obj(dlh_reqinput_qdesc, sid->get_sid());
360 req_queues.push_back(obj_fac.get_dal<confmodel::Connection>(tpreq_queue_obj.UID()));
363 conffwk::ConfigObject tp_net_obj = obj_fac.create_net_obj(tp_net_desc, tp_uid);
366 conffwk::ConfigObject ta_net_obj = obj_fac.create_net_obj(ta_net_desc, tp_uid);
369 tph_obj.
set_objs(
"inputs", { &tp_queue_obj, &tpreq_queue_obj });
370 tph_obj.set_objs(
"outputs", { &tp_net_obj, &ta_net_obj, &frag_queue_obj });
372 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(tph_obj.UID()));
377 std::vector<const conffwk::ConfigObject*> tp_queue_objs;
378 for (
auto q : tp_queues) {
379 tp_queue_objs.push_back(&q->config_object());
387 auto emulation_mode = reader_conf->get_emulation_mode();
388 for (
auto ds : all_enabled_det_streams) {
390 uint32_t sid = ds->get_source_id();
392 "Processing stream {}, id {}, det id {}", ds->UID(), ds->get_source_id(), ds->get_geo_id()->get_detector_id());
393 std::string uid(fmt::format(
"DLH-{}", sid));
395 "creating OKS configuration object for Data Link Handler class {}, if {}", dlh_class, sid);
396 auto dlh_obj = obj_fac.create(dlh_class, uid);
397 dlh_obj.set_by_val<uint32_t>(
"source_id", sid);
398 dlh_obj.set_by_val<uint32_t>(
"detector_id", ds->get_geo_id()->get_detector_id());
400 dlh_obj.set_by_val<
bool>(
"emulation_mode", emulation_mode);
401 dlh_obj.set_obj(
"geo_id", &ds->get_geo_id()->config_object());
402 dlh_obj.set_obj(
"module_configuration", &dlh_conf->config_object());
403 dlh_obj.set_obj(
"raw_data_callback", &callback_confs_by_sid[sid]->
config_object());
405 std::vector<const conffwk::ConfigObject*> dlh_ins, dlh_outs;
408 conffwk::ConfigObject req_queue_obj = obj_fac.create_queue_sid_obj(dlh_reqinput_qdesc, ds);
411 req_queues.push_back(obj_fac.get_dal<confmodel::Connection>(req_queue_obj.UID()));
412 dlh_ins.push_back(&req_queue_obj);
413 dlh_outs.push_back(&frag_queue_obj);
416 if (dlh_conf->get_generate_timesync()) {
418 conffwk::ConfigObject ts_net_obj = obj_fac.create_net_obj(ts_net_desc, std::to_string(sid));
419 dlh_outs.push_back(&ts_net_obj);
422 for (
auto tpq : tp_queue_objs) {
423 dlh_outs.push_back(tpq);
425 dlh_obj.set_objs(
"inputs", dlh_ins);
426 dlh_obj.set_objs(
"outputs", dlh_outs);
428 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(dlh_obj.UID()));
433 if (aggregator_conf == 0) {
434 throw(BadConf(
ERS_HERE,
"No FragmentAggregatorModule configuration given"));
436 std::string faUid(
"fragmentaggregator-" +
UID());
438 TLOG_DEBUG(7) <<
"creating OKS configuration object for Fragment Aggregator class ";
439 auto frag_aggr = obj_fac.create(
"FragmentAggregatorModule", faUid);
440 conffwk::ConfigObject fa_net_obj = obj_fac.create_net_obj(fa_net_desc);
444 std::vector<conffwk::ConfigObject> fragOutObjs;
445 for (
auto [uid,
descriptor] : helper->get_netdescriptors(
"Fragment",
"DFApplication")) {
446 std::string dreqNetUid(
descriptor->get_uid_base() + uid);
447 auto frag_conn = obj_fac.create(
"NetworkConnection", dreqNetUid);
449 frag_conn.set_by_val<std::string>(
"data_type",
descriptor->get_data_type());
450 frag_conn.set_by_val<std::string>(
"connection_type",
descriptor->get_connection_type());
452 frag_conn.set_by_val<
int>(
"capacity", all_enabled_det_streams.size() * 2);
454 auto serviceObj =
descriptor->get_associated_service()->config_object();
455 frag_conn.set_obj(
"associated_service", &serviceObj);
456 fragOutObjs.push_back(frag_conn);
460 std::vector<const conffwk::ConfigObject*> fa_output_objs;
461 for (
auto& fNet : fragOutObjs) {
462 fa_output_objs.push_back(&fNet);
465 for (
auto& q : req_queues) {
466 fa_output_objs.push_back(&q->config_object());
469 frag_aggr.set_obj(
"configuration", &aggregator_conf->config_object());
470 frag_aggr.set_objs(
"inputs", { &fa_net_obj, &frag_queue_obj });
471 frag_aggr.set_objs(
"outputs", fa_output_objs);
472 modules.push_back(obj_fac.get_dal<confmodel::DaqModule>(frag_aggr.UID()));
474 obj_fac.update_modules(modules);
set_objs(self, name, value)
const dunedaq::appmodel::DataHandlerConf * get_tp_handler() const
Get "tp_handler" relationship value.
virtual std::vector< const ExcludableEntity * > contained_excludable_entities() const override
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.
void generate_modules(std::shared_ptr< appmodel::ConfigurationHelper >) const override
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.
const ConfigObject & config_object() const
const std::string & UID() const noexcept
std::vector< const dunedaq::confmodel::ExcludableEntity * > to_resources(const std::vector< T * > &vector_of_children)
#define TLOG_DEBUG(lvl,...)
CIB Buffer std::string descriptor Message from std::string descriptor CIB process std::string descriptor descriptor