DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TPCEthFrameProcessor.hxx
Go to the documentation of this file.
2DUNE_DAQ_TYPESTRING(std::vector<dunedaq::trigger::TriggerPrimitiveTypeAdapter>, "TriggerPrimitiveVector")
3
4namespace dunedaq {
5namespace fdreadoutlibs {
6
7using datahandlinglibs::logging::TLVL_BOOKKEEPING;
8using datahandlinglibs::logging::TLVL_TAKE_NOTE;
9
10template<class ReadoutTypeAdapter>
12 std::unique_ptr<datahandlinglibs::FrameErrorRegistry>& error_registry,
13 bool processing_enabled)
14 : datahandlinglibs::TaskRawDataProcessorModel<ReadoutTypeAdapter>(error_registry, processing_enabled)
15{
16}
17
18template<class ReadoutTypeAdapter>
19void
20TPCEthFrameProcessor<ReadoutTypeAdapter>::start(const appfwk::DAQModule::CommandData_t& args)
21{
22 // Reset software TPG resources
23 if (this->m_post_processing_enabled) {
26 }
27
28 // Reset timestamp check
29 m_previous_ts = 0;
30 m_current_ts = 0;
33 m_ts_error_state = false;
35
40
41 // Reset stats
42 m_t0 = std::chrono::high_resolution_clock::now();
43 m_num_new_tps.exchange(0);
44
45#ifdef TPGLIBS_ENABLE_STATE_MONITORING
46 if (m_state_harvester && m_tpg_metric_collect_enabled) {
47 m_state_harvester->start_collection_thread();
48 }
49#endif
50 inherited::start(args);
51}
52
53template<class ReadoutTypeAdapter>
54void
55TPCEthFrameProcessor<ReadoutTypeAdapter>::stop(const appfwk::DAQModule::CommandData_t& args)
56{
57 inherited::stop(args);
58 if (this->m_post_processing_enabled) {
59#ifdef TPGLIBS_ENABLE_STATE_MONITORING
60 if (m_state_harvester) {
61 m_state_harvester->stop_collection_thread();
62 }
63#endif
64 // Clears the pipelines and resets with the given configs.
65 m_tp_generator->configure(m_tpg_configs, m_channel_plane_numbers, ReadoutTypeAdapter::samples_tick_difference);
66 }
67}
68
69template<class ReadoutTypeAdapter>
70void
74 m_sourceid.subsystem = ReadoutTypeAdapter::subsystem;
75 auto geo_id = conf->get_geo_id();
76 if (geo_id != nullptr) {
77 m_det_id = geo_id->get_detector_id();
78 m_crate_id = geo_id->get_crate_id();
79 m_slot_id = geo_id->get_slot_id();
80 m_stream_id = geo_id->get_stream_id();
81 }
82}
84template<class ReadoutTypeAdapter>
85void
91 std::bind(&TPCEthFrameProcessor<ReadoutTypeAdapter>::sequence_check, this, std::placeholders::_1));
92 }
95 std::bind(&TPCEthFrameProcessor<ReadoutTypeAdapter>::timestamp_check, this, std::placeholders::_1));
96}
98template<class ReadoutTypeAdapter>
99void
101 const appmodel::TPCRawDataProcessor* proc_conf)
102{
103 const std::shared_ptr<detchannelmaps::TPCChannelMap> channel_map =
104 dunedaq::detchannelmaps::make_tpc_map(proc_conf->get_channel_map());
105 const std::vector<unsigned int> channel_mask_vec = proc_conf->get_channel_mask();
106
107 for (int chan = 0; chan < 64; chan++) {
108 trgdataformats::channel_t off_channel = channel_map->get_offline_channel_from_det_crate_slot_stream_chan(
110 int16_t plane = channel_map->get_plane_from_offline_channel(off_channel);
111 m_channel_plane_numbers.push_back(std::make_pair(off_channel, plane));
112
113 // This processor only needs to handle some (maybe 0) of the masked channels.
114 // Only get those relevant channels for the later check.
115 // Only get the planes for the channels that are not masked.
116 if (std::find(channel_mask_vec.begin(), channel_mask_vec.end(), off_channel) != channel_mask_vec.end()) {
117 m_channel_mask_set.insert(off_channel);
118 } else {
119 m_channel_plane_map[off_channel] = plane;
120 m_plane_numbers_set.insert(plane);
121 }
122 }
123}
124
125template<class ReadoutTypeAdapter>
126void
128 const appmodel::TPCRawDataProcessor* proc_conf)
129{
130 // Setting TP sinks.
131 // Configurations currently have the sinks iterate in order, but there may be more sinks than planes.
132 int plane_number = 0;
133 for (auto output : conf->get_outputs()) {
134 try {
135 if (output->get_data_type() == "TriggerPrimitiveVector") {
136 if (m_plane_numbers_set.contains(plane_number)) {
137 m_plane_to_tp_sink_map[plane_number] =
138 get_iom_sender<std::vector<trigger::TriggerPrimitiveTypeAdapter>>(output->UID());
139 }
140 plane_number++;
141 }
142 } catch (const ers::Issue& excpt) {
143 ers::error(datahandlinglibs::ResourceQueueError(ERS_HERE, "tp", "DefaultRequestHandlerModel", excpt));
144 }
145 }
146
147 // We do need a coverage for all planes.
148 if (m_plane_numbers_set.size() > m_plane_to_tp_sink_map.size()) {
149 ers::error(DetectorPlaneToTPSinkMismatch(ERS_HERE, m_plane_numbers_set.size(), m_plane_to_tp_sink_map.size()));
150 }
151
152 m_tp_generator = std::make_unique<tpglibs::TPGenerator>();
153
154 // Set the minimum TP samples over threshold.
155 auto conf_sot_minima = proc_conf->get_sot_minima();
156 std::vector<uint16_t> sot_minima{ conf_sot_minima->get_sot_minimum_plane0(),
157 conf_sot_minima->get_sot_minimum_plane1(),
158 conf_sot_minima->get_sot_minimum_plane2() };
159 m_tp_generator->set_sot_minima(sot_minima);
160
161 std::vector<const appmodel::ProcessingStep*> processing_steps = proc_conf->get_processing_steps();
162 for (auto step : processing_steps) {
163 m_tpg_configs.push_back(std::make_pair(step->class_name(), step->to_json(false).back()));
164 }
165
166 // Let the TPG generator configure
167 m_tp_generator->configure(m_tpg_configs, m_channel_plane_numbers, ReadoutTypeAdapter::samples_tick_difference);
168
169 // Set the limits on when to send TPs and check that we can actually send on these limits.
170 m_frame_count_limit = proc_conf->get_frame_count_limit();
171 m_tp_count_limit = proc_conf->get_tp_count_limit();
172 m_frame_limit_enabled = m_frame_count_limit > 0;
173 m_tp_limit_enabled = m_tp_count_limit > 0;
174
175 if (!m_frame_limit_enabled && !m_tp_limit_enabled) {
177 }
178
179 m_metric_collect_opmon_period = proc_conf->get_metric_collect_opmon_period();
180
181#ifdef TPGLIBS_ENABLE_STATE_MONITORING
182 // Monitoring enabled at build time — set up harvester
183 // Still read toggle_state for backwards compat (honor it for now)
184 m_tpg_metric_collect_enabled = false;
185 for (const auto& name_config : m_tpg_configs) {
186 if (name_config.second.contains("metric_collect_toggle_state") &&
187 name_config.second["metric_collect_toggle_state"] == true) {
188 m_tpg_metric_collect_enabled = true;
189 break;
190 }
191 }
192
193 if (m_tpg_metric_collect_enabled) {
194 auto processsor_references = m_tp_generator->get_all_processor_references_with_pipeline_index();
195
196 m_state_harvester = std::make_unique<fdreadoutlibs::TPGInternalStateHarvester>();
197
198 const uint8_t channels_per_pipeline = 16;
199 const uint8_t pipelines = static_cast<uint8_t>(m_channel_plane_numbers.size() / channels_per_pipeline);
200
201 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Configuring state harvester with " << static_cast<int>(channels_per_pipeline)
202 << " channels per pipeline, " << static_cast<int>(pipelines) << " pipelines, "
203 << processsor_references.size() << " processor references";
204
205 m_state_harvester->update_channel_plane_numbers(m_channel_plane_numbers, channels_per_pipeline, pipelines);
206 m_state_harvester->set_processor_references(processsor_references);
207
208 m_state_harvester->start_collection_thread();
209
210 TLOG_DEBUG(TLVL_BOOKKEEPING) << "State harvester configured and started successfully";
211 }
212
213#else
214 // Monitoring disabled at build time
215 m_tpg_metric_collect_enabled = false;
216
217 // Warn if per-processor monitoring params are configured but will have no effect
218 bool warned_build_off = false;
219 for (const auto& name_config : m_tpg_configs) {
220 // Check if any of the monitoring configs are not the default values: someone is requesting them.
221 const bool toggle_state_set = name_config.second.value("metric_collect_toggle_state", false) == true;
222 const bool time_sample_period_set =
223 name_config.second.value("metric_collect_time_sample_period", uint64_t{ 256 }) != 256;
224 const bool requested_internal_states_set =
225 !name_config.second.value("requested_internal_states", std::string{}).empty();
226 if (!warned_build_off && (toggle_state_set || time_sample_period_set || requested_internal_states_set)) {
227 // warn only once.
228 ers::warning(TPGStateMonitoringConfigIgnored(ERS_HERE));
229 warned_build_off = true;
230 }
231 }
232#endif
233
234 inherited::add_postprocess_task(
235 std::bind(&TPCEthFrameProcessor<ReadoutTypeAdapter>::find_tps, this, std::placeholders::_1));
236}
237
238template<class ReadoutTypeAdapter>
239void
241{
243 if (dp == nullptr) {
244 return;
245 }
246
248 if (proc_conf == nullptr) {
249 return;
250 }
251
252 // Need TPCRawDataProcessor configurations to configure the following.
254 configure_find_tps(conf, proc_conf);
255}
256
257template<class ReadoutTypeAdapter>
258void
271
272template<class ReadoutTypeAdapter>
273void
283
284template<class ReadoutTypeAdapter>
285void
287{
288 m_emulator_mode = false;
289 m_first_frame = true;
290
291 // Timestamps.
292 m_previous_ts = 0;
293 m_current_ts = 0;
294
297
299 m_ts_problem_reported = false;
300 m_ts_error_state = false;
301 m_ts_error_ctr = 0;
302
303 // Sequence ID.
306
309 m_seq_id_error_state = false;
313
314 // The preprocessing tasks scrap is handled by inherited::scrap().
315}
316
317template<class ReadoutTypeAdapter>
318void
320{
321 // Channel-plane variables
322 m_channel_mask_set.clear();
323 m_plane_numbers_set.clear();
325 m_channel_plane_map.clear();
326
327 // TP variables
328 m_tp_generator->reset();
329 m_tpg_configs.clear();
332
333 m_frame_limit_enabled = false;
334 m_tp_limit_enabled = false;
338
339 // OpMon variables
342 m_tp_channel_rate_map.clear();
343 m_num_new_tps.exchange(0);
344 m_tps_suppressed_too_long.exchange(0);
345 m_tps_send_failed.exchange(0);
346 m_frame_counter.exchange(0);
347 m_t0 = std::chrono::high_resolution_clock::now();
348}
349
350template<class ReadoutTypeAdapter>
351void
352TPCEthFrameProcessor<ReadoutTypeAdapter>::scrap(const appfwk::DAQModule::CommandData_t& cfg)
353{
356
357 if (this->m_post_processing_enabled) {
359 }
360
361 inherited::scrap(cfg);
362}
363
364template<class ReadoutTypeAdapter>
365void
367{
369
370 info.set_num_seq_id_errors(m_seq_id_error_ctr.load());
371 info.set_min_seq_id_jump(m_seq_id_min_jump.exchange(0));
372 info.set_max_seq_id_jump(m_seq_id_max_jump.exchange(0));
373
374 info.set_num_ts_errors(m_ts_error_ctr.load());
375
376 this->publish(std::move(info));
377
378 this->m_error_registry->log_registered_errors();
379
380 if (this->m_post_processing_enabled) {
381 auto now = std::chrono::high_resolution_clock::now();
382 int num_new_tps = m_num_new_tps.exchange(0);
383 int num_new_tps_suppressed_too_long = m_tps_suppressed_too_long.exchange(0);
384 int num_new_tps_send_failed = m_tps_send_failed.exchange(0);
385 double seconds = std::chrono::duration_cast<std::chrono::microseconds>(now - m_t0).count() / 1000000.;
386 TLOG_DEBUG(TLVL_BOOKKEEPING) << "TP rate: " << std::to_string(num_new_tps / seconds / 1000.) << " [kHz]";
387 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Total new TPs: " << num_new_tps;
388
390 tp_info.set_rate_tp_hits(num_new_tps / seconds / 1000.);
391
392 tp_info.set_num_tps_sent(num_new_tps);
393 tp_info.set_num_tps_suppressed_too_long(num_new_tps_suppressed_too_long);
394 tp_info.set_num_tps_send_failed(num_new_tps_send_failed);
395
396 this->publish(std::move(tp_info));
397 // Find the channels with the top TP rates
398 // Create a vector of pairs to store the map elements
399 std::vector<std::pair<uint, int>> channel_tp_rate_vec(m_tp_channel_rate_map.begin(), m_tp_channel_rate_map.end());
400 // Sort the vector in descending order of the value of the pairs
401 sort(channel_tp_rate_vec.begin(), channel_tp_rate_vec.end(), [](std::pair<uint, int>& a, std::pair<uint, int>& b) {
402 return a.second > b.second;
403 });
404 // Add the metrics to opmon
405 // For convenience we are selecting only the top 10 elements
406 if (channel_tp_rate_vec.size() != 0) {
407 int top_highest_values = 10;
408 if (channel_tp_rate_vec.size() < 10) {
409 top_highest_values = channel_tp_rate_vec.size();
410 }
411 // datahandlinglibs::opmon::TPChannelsInfo channels_info;
412 for (int i = 0; i < top_highest_values; i++) {
414 tpc_info.set_number_of_tps(channel_tp_rate_vec[i].second);
415 tpc_info.set_channel_id(channel_tp_rate_vec[i].first);
416 this->publish(std::move(tpc_info), { { "channel", std::to_string(channel_tp_rate_vec[i].first) } });
417 }
418 }
419
420 // Reset the counter in the channel rate map
421 for (auto& el : m_tp_channel_rate_map) {
422 el.second = 0;
423 }
424 m_t0 = now;
425
426#ifdef TPGLIBS_ENABLE_STATE_MONITORING
427 if (m_tpg_metric_collect_enabled && m_state_harvester) {
428 publish_processor_metric_to_opmon();
429 publish_processor_metric_to_opmon_with_aggregation();
430 }
431#endif
432 }
433
435}
436
437#ifdef TPGLIBS_ENABLE_STATE_MONITORING
438template<class ReadoutTypeAdapter>
439void
441{
442 if (!m_state_harvester) {
443 return;
444 }
445
446 // Get latest results from background collection thread
447 auto metrics = m_state_harvester->get_latest_results();
448
449 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Publishing processor metrics for " << metrics.size() << " channels";
450
451 int metrics_published = 0;
452
453 // Publish per-channel metrics
454 for (const auto& [channel, vec] : metrics) {
456 bool has_valid_metrics = false;
457
458 for (const auto& [name, val] : vec) {
459 if (name == "pedestal") {
460 tpg_proc_info.set_pedestal(val);
461 has_valid_metrics = true;
462 } else if (name == "accum") {
463 tpg_proc_info.set_accum(val);
464 has_valid_metrics = true;
465 }
466 }
467
468 if (has_valid_metrics) {
469 this->publish(std::move(tpg_proc_info), { { "channel", std::to_string(channel) } });
470 metrics_published++;
471 }
472 }
473
474 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Published " << metrics_published << " channel metrics";
475}
476
477template<class ReadoutTypeAdapter>
478std::map<
479 int16_t,
480 std::map<
481 std::string,
482 std::tuple<float, int16_t, int16_t, float, dunedaq::trgdataformats::channel_t, dunedaq::trgdataformats::channel_t>>>
484 const std::unordered_map<dunedaq::trgdataformats::channel_t, std::vector<std::pair<std::string, int16_t>>>& metrics)
485{
486 // Structure to hold all statistics: plane -> metric -> (mean, min, max, stddev, min_channel_id, max_channel_id)
487 std::map<
488 int16_t,
489 std::map<
490 std::string,
491 std::
492 tuple<float, int16_t, int16_t, float, dunedaq::trgdataformats::channel_t, dunedaq::trgdataformats::channel_t>>>
493 all_stats;
494
495 // Structure to accumulate statistics: plane -> metric -> (count, mean, M2, min, max, min_channel_id, max_channel_id)
496 std::map<int16_t,
497 std::map<std::string,
498 std::tuple<size_t,
499 double,
500 double,
501 int16_t,
502 int16_t,
505 accumulators;
506
507 // Single pass through all metrics to collect data using Welford's online algorithm for variance
508 for (const auto& [channel, vec] : metrics) {
509 if (m_channel_plane_map.empty())
510 continue;
511
512 int16_t plane = m_channel_plane_map[channel];
513
514 for (const auto& [name, val] : vec) {
515 auto& [count, mean, M2, min, max, min_channel_id, max_channel_id] = accumulators[plane][name];
516
517 count++;
518
519 if (count == 1 || val < min) {
520 min = val;
521 min_channel_id = channel;
522 }
523 if (count == 1 || val > max) {
524 max = val;
525 max_channel_id = channel;
526 }
527
528 // Welford's online algorithm for variance calculation
529 if (count == 1) {
530 // First value: initialize mean and M2
531 mean = val;
532 M2 = 0.0;
533 } else {
534 // Update mean and M2 using Welford's algorithm
535 double delta = val - mean;
536 mean += delta / count;
537 double delta2 = val - mean;
538 M2 += delta * delta2;
539 }
540 }
541 }
542
543 // Calculate final statistics from accumulated data
544 for (const auto& [plane, metric_map] : accumulators) {
545 for (const auto& [metric_name, acc_data] : metric_map) {
546 const auto& [count, mean, M2, min, max, min_channel_id, max_channel_id] = acc_data;
547
548 if (count == 0)
549 continue;
550
551 float stddev = 0.0f;
552
553 // Calculate standard deviation using accumulated M2
554 if (count > 1) {
555 stddev = std::sqrt(M2 / (count - 1));
556 }
557
558 all_stats[plane][metric_name] =
559 std::make_tuple(static_cast<float>(mean), min, max, stddev, min_channel_id, max_channel_id);
560 }
561 }
562
563 return all_stats;
564}
565
566template<class ReadoutTypeAdapter>
567void
569{
570 if (!m_state_harvester) {
571 return;
572 }
573
574 // Get latest results from background collection thread
575 auto metrics = m_state_harvester->get_latest_results();
576
577 // Use optimized single-pass calculation for all metrics across all planes
578 auto all_stats = calculate_all_metric_summaries_across_planes(metrics);
579
580 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Publishing aggregated metrics for " << all_stats.size() << " planes";
581
582 // Publish all calculated statistics
583 for (const auto& [plane, metric_map] : all_stats) {
584 for (const auto& [metric_name, stats] : metric_map) {
585 const auto& [mean, min, max, stddev, min_channel_id, max_channel_id] = stats;
586
587 datahandlinglibs::opmon::TPGProcessorReducedInfo info;
588 info.set_average(mean);
589 info.set_max(max);
590 info.set_min(min);
591 info.set_standard_dev(stddev);
592 info.set_max_channel_id(max_channel_id);
593 info.set_min_channel_id(min_channel_id);
594 this->publish(std::move(info), { { "plane", std::to_string(plane) }, { "metric", metric_name } });
595 }
596 }
597}
598#endif // TPGLIBS_ENABLE_STATE_MONITORING
599
603template<class ReadoutTypeAdapter>
604void
606{
607 // Acquire timestamp
608 auto wfptr = reinterpret_cast<tpcframeptr>(fp); // NOLINT
609 m_current_seq_id = wfptr->daq_header.seq_id;
610
611 // Check that the system is properly configured from the first frame.
612 if (m_first_frame) [[unlikely]] {
613 if (wfptr->daq_header.crate_id != m_crate_id || wfptr->daq_header.slot_id != m_slot_id ||
614 wfptr->daq_header.stream_id != m_stream_id) {
615 ers::error(LinkMisconfiguration(ERS_HERE,
616 wfptr->daq_header.crate_id,
617 wfptr->daq_header.slot_id,
618 wfptr->daq_header.stream_id,
620 m_slot_id,
621 m_stream_id));
622 }
623
624 m_first_frame = false;
625 }
626
627 // Check sequence id
628 // Calculate the next sequence id (12 bits)
629 uint16_t expected_seq_id = (m_previous_seq_id + fp->get_num_frames()) & 0xfff;
630 int16_t delta_seq_id = m_current_seq_id - expected_seq_id;
631 if (delta_seq_id > 0x800) {
632 delta_seq_id -= 0x1000;
633 } else if (delta_seq_id < -0x7ff) {
634 delta_seq_id += 0x1000;
635 }
636
637 if (delta_seq_id == 0) {
638 m_seq_id_error_state = false;
639 } else {
640 // uint16_t delta_seq_id = (m_current_seq_id-expected_seq_id);
642 m_seq_id_max_jump = std::max(delta_seq_id, m_seq_id_max_jump.load());
643 m_seq_id_min_jump = std::min(delta_seq_id, m_seq_id_min_jump.load());
644
645 if (m_first_seq_id_mismatch) { // log once
646 TLOG_DEBUG(TLVL_BOOKKEEPING) << "First sequence id MISMATCH! -> | previous: " << std::to_string(m_previous_seq_id)
647 << " current: " + std::to_string(m_current_seq_id);
649 } else {
651 this->m_error_registry->add_error(
652 "Sequence ID jump", datahandlinglibs::FrameErrorRegistry::ErrorInterval(expected_seq_id, m_current_seq_id));
654 }
655 }
656 }
657
658 if (m_seq_id_error_ctr > 1000) {
660 TLOG() << "*** Data Integrity ERROR *** Sequence ID continuity is completely broken! "
661 << "Something is wrong with the FE source or with the configuration!";
663 }
664 }
665
667}
668
672template<class ReadoutTypeAdapter>
673void
675{
676
677 uint16_t tpceth_tick_difference = ReadoutTypeAdapter::expected_tick_difference;
678 uint16_t tpceth_frame_tick_difference = tpceth_tick_difference * fp->get_num_frames();
679
680 auto wfptr = reinterpret_cast<tpcframeptr>(fp); // NOLINT
681 m_current_ts = wfptr->get_timestamp();
682
683 // Check timestamp
684 if (m_previous_ts > 0 && m_current_ts - m_previous_ts != tpceth_frame_tick_difference) [[unlikely]] {
686 if (m_first_ts_missmatch) { // log once
687 TLOG_DEBUG(TLVL_BOOKKEEPING) << "First timestamp MISMATCH! -> | previous: " << std::to_string(m_previous_ts)
688 << " current: " + std::to_string(m_current_ts);
689 m_first_ts_missmatch = false;
690 } else {
691 if (!m_ts_error_state) {
692 this->m_error_registry->add_error("Timestamp jump",
694 m_previous_ts + tpceth_frame_tick_difference, m_current_ts));
695 m_ts_error_state = true;
696 }
697 }
698 } else {
699 m_ts_error_state = false;
700 }
701
702 if (m_ts_error_ctr > 1000) {
704 TLOG() << "*** Data Integrity ERROR *** Timestamp continuity is completely broken! "
705 << "Something is wrong with the FE source or with the configuration!";
707 }
708 }
709
712}
713
717template<class ReadoutTypeAdapter>
718void
720{
721 if (!fp)
722 return;
723 auto wfptr = reinterpret_cast<tpcframeptr>((uint8_t*)fp); // NOLINT
724
725 std::vector<trgdataformats::TriggerPrimitive> tps = (*m_tp_generator)(wfptr);
726
727 uint64_t current_frame_count = m_frame_counter.fetch_add(1, std::memory_order_relaxed) + 1;
728
729#ifdef TPGLIBS_ENABLE_STATE_MONITORING
730 if (m_tpg_metric_collect_enabled && m_state_harvester && current_frame_count % m_metric_collect_opmon_period == 0) {
731 m_state_harvester->trigger_harvest();
732 }
733#endif
734
735 for (const auto& tp : tps) {
736 // If this TP is on a masked channel, skip it.
737 if (std::binary_search(m_channel_mask_set.begin(), m_channel_mask_set.end(), uint32_t(tp.channel)))
738 continue;
739 // Need to move into a type adapter.
741 tpa.tp = tp;
742
743 tpa.tp.detid = m_det_id; // Last missing piece.
744 m_plane_to_tpa_vector_map[m_channel_plane_map[uint32_t(tp.channel)]].push_back(tpa);
745 m_tp_channel_rate_map[uint32_t(tp.channel)]++;
747 }
748
749 const bool frame_limit_reached =
751 const bool tp_limit_reached = m_tp_limit_enabled && (m_current_tp_count >= m_tp_count_limit);
752
753 if (frame_limit_reached || tp_limit_reached) [[unlikely]] {
755 m_frame_count_at_last_send = current_frame_count;
756 for (auto& [plane_num, tpa_vector] : m_plane_to_tpa_vector_map) {
757 int num_new_tps = tpa_vector.size();
758 if (num_new_tps == 0) {
759 continue;
760 }
761 const auto ts_begin = tpa_vector.front().tp.time_start;
762 const auto channel_begin = tpa_vector.front().tp.channel;
763 const auto ts_end = tpa_vector.back().tp.time_start;
764 const auto channel_end = tpa_vector.back().tp.channel;
765 if (!m_plane_to_tp_sink_map[plane_num]->try_send(std::move(tpa_vector), iomanager::Sender::s_no_block)) {
766 ers::warning(FailedToSendTPVector(ERS_HERE, ts_begin, channel_begin, ts_end, channel_end));
768 } else {
769 m_num_new_tps += num_new_tps;
770 }
771 }
772 }
773 return;
774}
775
776} // namespace fdreadoutlibs
777} // namespace dunedaq
#define ERS_HERE
#define DUNE_DAQ_TYPESTRING(Type, typestring)
Declare the datatype_to_string method for the given type.
const dunedaq::appmodel::DataProcessor * get_data_processor() const
Get "data_processor" relationship value.
const dunedaq::appmodel::DataHandlerConf * get_module_configuration() const
Get "module_configuration" relationship value.
bool get_emulation_mode() const
Get "emulation_mode" attribute value.
uint32_t get_source_id() const
Get "source_id" attribute value.
const dunedaq::confmodel::GeoId * get_geo_id() const
Get "geo_id" relationship value.
const std::vector< uint32_t > & get_channel_mask() const
Get "channel_mask" attribute value. List of channels to be masked from TP generation.
const std::string & get_channel_map() const
Get "channel_map" attribute value.
uint32_t get_frame_count_limit() const
Get "frame_count_limit" attribute value. When this number of frames is reached the TPs are sent to th...
const dunedaq::appmodel::SamplesOverThresholdMinima * get_sot_minima() const
Get "sot_minima" relationship value. TP samples over threshold minimum requirement by plane.
const std::vector< const dunedaq::appmodel::ProcessingStep * > & get_processing_steps() const
Get "processing_steps" relationship value.
uint32_t get_metric_collect_opmon_period() const
Get "metric_collect_opmon_period" attribute value. The rate at which processor metric is polled from ...
uint32_t get_tp_count_limit() const
Get "tp_count_limit" attribute value. When this number of TPs is reached, the TPs are sent to the sin...
const TARGET * cast() const noexcept
Casts object to different class.
const std::vector< const dunedaq::confmodel::Connection * > & get_outputs() const
Get "outputs" relationship value. Output connections from this module.
void start(const appfwk::DAQModule::CommandData_t &) override
void conf(const appmodel::DataHandlerModule *conf) override
void stop(const appfwk::DAQModule::CommandData_t &) override
TaskRawDataProcessorModel(std::unique_ptr< FrameErrorRegistry > &error_registry, bool post_processing_enabled)
void scrap(const appfwk::DAQModule::CommandData_t &) override
void configure_postprocessing(const appmodel::DataHandlerModule *conf)
void configure_channel_plane_numbers(const appmodel::TPCRawDataProcessor *proc_conf)
void stop(const appfwk::DAQModule::CommandData_t &args) override
Stop operation.
std::vector< std::pair< std::string, nlohmann::json > > m_tpg_configs
std::unordered_map< unsigned int, std::vector< trigger::TriggerPrimitiveTypeAdapter > > m_plane_to_tpa_vector_map
void configure_find_tps(const appmodel::DataHandlerModule *conf, const appmodel::TPCRawDataProcessor *proc_conf)
void scrap(const appfwk::DAQModule::CommandData_t &cfg) override
Unconfigure.
std::unique_ptr< tpglibs::TPGenerator > m_tp_generator
dunedaq::daqdataformats::timestamp_t m_previous_ts
std::chrono::time_point< std::chrono::high_resolution_clock > m_t0
std::vector< std::pair< trgdataformats::channel_t, int16_t > > m_channel_plane_numbers
std::unordered_map< unsigned int, std::shared_ptr< iomanager::SenderConcept< std::vector< trigger::TriggerPrimitiveTypeAdapter > > > > m_plane_to_tp_sink_map
std::map< uint, std::atomic< int > > m_tp_channel_rate_map
void configure_source_and_geo_ids(const appmodel::DataHandlerModule *conf)
std::unordered_map< trgdataformats::channel_t, unsigned int > m_channel_plane_map
void start(const appfwk::DAQModule::CommandData_t &args) override
Start operation.
void configure_preprocessing(const appmodel::DataHandlerModule *conf)
void conf(const appmodel::DataHandlerModule *conf) override
Set the emulator mode, if active, timestamps of processed packets are overwritten with new ones.
TPCEthFrameProcessor(std::unique_ptr< datahandlinglibs::FrameErrorRegistry > &error_registry, bool processing_enabled)
dunedaq::daqdataformats::timestamp_t m_pattern_generator_current_ts
dunedaq::daqdataformats::timestamp_t m_pattern_generator_previous_ts
dunedaq::daqdataformats::timestamp_t m_current_ts
static constexpr timeout_t s_no_block
Definition Sender.hpp:26
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
Base class for any user define issue.
Definition Issue.hpp:76
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
The DUNE-DAQ namespace.
Both frame_count_limit and tp_count_limit were set FailedToSendTPVector
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101
stats(obj, sel_links, seconds, show_udp, show_buf)
SourceID is a generalized representation of the source of a piece of data in the DAQ....
Definition SourceID.hpp:32