5namespace fdreadoutlibs {
7using datahandlinglibs::logging::TLVL_BOOKKEEPING;
8using datahandlinglibs::logging::TLVL_TAKE_NOTE;
10template<
class ReadoutTypeAdapter>
12 std::unique_ptr<datahandlinglibs::FrameErrorRegistry>& error_registry,
13 bool processing_enabled)
18template<
class ReadoutTypeAdapter>
42 m_t0 = std::chrono::high_resolution_clock::now();
45#ifdef TPGLIBS_ENABLE_STATE_MONITORING
47 m_state_harvester->start_collection_thread();
53template<
class ReadoutTypeAdapter>
59#ifdef TPGLIBS_ENABLE_STATE_MONITORING
60 if (m_state_harvester) {
61 m_state_harvester->stop_collection_thread();
69template<
class ReadoutTypeAdapter>
76 if (geo_id !=
nullptr) {
77 m_det_id = geo_id->get_detector_id();
84template<
class ReadoutTypeAdapter>
98template<
class ReadoutTypeAdapter>
103 const std::shared_ptr<detchannelmaps::TPCChannelMap> channel_map =
105 const std::vector<unsigned int> channel_mask_vec = proc_conf->
get_channel_mask();
107 for (
int chan = 0; chan < 64; chan++) {
110 int16_t plane = channel_map->get_plane_from_offline_channel(off_channel);
116 if (std::find(channel_mask_vec.begin(), channel_mask_vec.end(), off_channel) != channel_mask_vec.end()) {
125template<
class ReadoutTypeAdapter>
132 int plane_number = 0;
135 if (output->get_data_type() ==
"TriggerPrimitiveVector") {
138 get_iom_sender<std::vector<trigger::TriggerPrimitiveTypeAdapter>>(output->UID());
143 ers::error(datahandlinglibs::ResourceQueueError(
ERS_HERE,
"tp",
"DefaultRequestHandlerModel", excpt));
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()));
152 m_tp_generator = std::make_unique<tpglibs::TPGenerator>();
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);
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()));
167 m_tp_generator->configure(m_tpg_configs, m_channel_plane_numbers, ReadoutTypeAdapter::samples_tick_difference);
172 m_frame_limit_enabled = m_frame_count_limit > 0;
173 m_tp_limit_enabled = m_tp_count_limit > 0;
175 if (!m_frame_limit_enabled && !m_tp_limit_enabled) {
181#ifdef TPGLIBS_ENABLE_STATE_MONITORING
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;
193 if (m_tpg_metric_collect_enabled) {
194 auto processsor_references = m_tp_generator->get_all_processor_references_with_pipeline_index();
196 m_state_harvester = std::make_unique<fdreadoutlibs::TPGInternalStateHarvester>();
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);
202 <<
" channels per pipeline, " <<
static_cast<int>(pipelines) <<
" pipelines, "
203 << processsor_references.size() <<
" processor references";
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);
208 m_state_harvester->start_collection_thread();
215 m_tpg_metric_collect_enabled =
false;
218 bool warned_build_off =
false;
219 for (
const auto& name_config : m_tpg_configs) {
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)) {
229 warned_build_off =
true;
234 inherited::add_postprocess_task(
238template<
class ReadoutTypeAdapter>
248 if (proc_conf ==
nullptr) {
257template<
class ReadoutTypeAdapter>
272template<
class ReadoutTypeAdapter>
284template<
class ReadoutTypeAdapter>
317template<
class ReadoutTypeAdapter>
347 m_t0 = std::chrono::high_resolution_clock::now();
350template<
class ReadoutTypeAdapter>
364template<
class ReadoutTypeAdapter>
376 this->
publish(std::move(info));
381 auto now = std::chrono::high_resolution_clock::now();
385 double seconds = std::chrono::duration_cast<std::chrono::microseconds>(now -
m_t0).count() / 1000000.;
396 this->
publish(std::move(tp_info));
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;
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();
412 for (
int i = 0; i < top_highest_values; i++) {
416 this->
publish(std::move(tpc_info), { {
"channel", std::to_string(channel_tp_rate_vec[i].first) } });
426#ifdef TPGLIBS_ENABLE_STATE_MONITORING
428 publish_processor_metric_to_opmon();
429 publish_processor_metric_to_opmon_with_aggregation();
437#ifdef TPGLIBS_ENABLE_STATE_MONITORING
438template<
class ReadoutTypeAdapter>
442 if (!m_state_harvester) {
447 auto metrics = m_state_harvester->get_latest_results();
451 int metrics_published = 0;
454 for (
const auto& [channel, vec] : metrics) {
456 bool has_valid_metrics =
false;
458 for (
const auto& [name, val] : vec) {
459 if (name ==
"pedestal") {
461 has_valid_metrics =
true;
462 }
else if (name ==
"accum") {
464 has_valid_metrics =
true;
468 if (has_valid_metrics) {
469 this->publish(std::move(tpg_proc_info), { {
"channel", std::to_string(channel) } });
477template<
class ReadoutTypeAdapter>
482 std::tuple<float, int16_t, int16_t, float, dunedaq::trgdataformats::channel_t, dunedaq::trgdataformats::channel_t>>>
492 tuple<float, int16_t, int16_t, float, dunedaq::trgdataformats::channel_t, dunedaq::trgdataformats::channel_t>>>
497 std::map<std::string,
508 for (
const auto& [channel, vec] : metrics) {
509 if (m_channel_plane_map.empty())
512 int16_t plane = m_channel_plane_map[
channel];
514 for (
const auto& [name, val] : vec) {
515 auto& [count, mean, M2, min, max, min_channel_id, max_channel_id] = accumulators[plane][name];
519 if (count == 1 || val < min) {
523 if (count == 1 || val > max) {
535 double delta = val - mean;
536 mean += delta / count;
537 double delta2 = val - mean;
538 M2 += delta * delta2;
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;
555 stddev = std::sqrt(M2 / (count - 1));
558 all_stats[plane][metric_name] =
559 std::make_tuple(
static_cast<float>(mean), min, max, stddev, min_channel_id, max_channel_id);
566template<
class ReadoutTypeAdapter>
570 if (!m_state_harvester) {
575 auto metrics = m_state_harvester->get_latest_results();
578 auto all_stats = calculate_all_metric_summaries_across_planes(metrics);
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;
587 datahandlinglibs::opmon::TPGProcessorReducedInfo
info;
588 info.set_average(mean);
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 } });
603template<
class ReadoutTypeAdapter>
616 wfptr->daq_header.crate_id,
617 wfptr->daq_header.slot_id,
618 wfptr->daq_header.stream_id,
631 if (delta_seq_id > 0x800) {
632 delta_seq_id -= 0x1000;
633 }
else if (delta_seq_id < -0x7ff) {
634 delta_seq_id += 0x1000;
637 if (delta_seq_id == 0) {
660 TLOG() <<
"*** Data Integrity ERROR *** Sequence ID continuity is completely broken! "
661 <<
"Something is wrong with the FE source or with the configuration!";
672template<
class ReadoutTypeAdapter>
677 uint16_t tpceth_tick_difference = ReadoutTypeAdapter::expected_tick_difference;
678 uint16_t tpceth_frame_tick_difference = tpceth_tick_difference * fp->get_num_frames();
704 TLOG() <<
"*** Data Integrity ERROR *** Timestamp continuity is completely broken! "
705 <<
"Something is wrong with the FE source or with the configuration!";
717template<
class ReadoutTypeAdapter>
723 auto wfptr =
reinterpret_cast<tpcframeptr>((uint8_t*)fp);
725 std::vector<trgdataformats::TriggerPrimitive> tps = (*m_tp_generator)(wfptr);
727 uint64_t current_frame_count =
m_frame_counter.fetch_add(1, std::memory_order_relaxed) + 1;
729#ifdef TPGLIBS_ENABLE_STATE_MONITORING
731 m_state_harvester->trigger_harvest();
735 for (
const auto& tp : tps) {
749 const bool frame_limit_reached =
753 if (frame_limit_reached || tp_limit_reached) [[unlikely]] {
757 int num_new_tps = tpa_vector.size();
758 if (num_new_tps == 0) {
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;
#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 add_preprocess_task(Task &&task)
void conf(const appmodel::DataHandlerModule *conf) override
bool m_post_processing_enabled
void stop(const appfwk::DAQModule::CommandData_t &) override
virtual void generate_opmon_data() override
TaskRawDataProcessorModel(std::unique_ptr< FrameErrorRegistry > &error_registry, bool post_processing_enabled)
std::atomic< uint64_t > m_last_processed_daq_ts
void scrap(const appfwk::DAQModule::CommandData_t &) override
std::unique_ptr< FrameErrorRegistry > & m_error_registry
void set_num_tps_send_failed(::uint64_t value)
void set_num_tps_suppressed_too_long(::uint64_t value)
void set_num_tps_sent(::uint64_t value)
void set_rate_tp_hits(float value)
void set_channel_id(::uint64_t value)
void set_number_of_tps(::uint64_t value)
void set_accum(::int64_t value)
void set_pedestal(::int64_t value)
daqdataformats::SourceID m_sourceid
const ReadoutTypeAdapter * constframeptr
std::set< unsigned int > m_plane_numbers_set
uint16_t m_current_seq_id
void scrap_postprocessing()
uint32_t m_current_tp_count
void configure_postprocessing(const appmodel::DataHandlerModule *conf)
void scrap_preprocessing()
ReadoutTypeAdapter::FrameType * tpcframeptr
bool m_frame_limit_enabled
uint32_t m_metric_collect_opmon_period
void configure_channel_plane_numbers(const appmodel::TPCRawDataProcessor *proc_conf)
void stop(const appfwk::DAQModule::CommandData_t &args) override
Stop operation.
std::atomic< uint64_t > m_tps_suppressed_too_long
std::vector< std::pair< std::string, nlohmann::json > > m_tpg_configs
bool m_tpg_metric_collect_enabled
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)
uint32_t m_frame_count_limit
void scrap(const appfwk::DAQModule::CommandData_t &cfg) override
Unconfigure.
std::unique_ptr< tpglibs::TPGenerator > m_tp_generator
bool m_first_ts_missmatch
dunedaq::daqdataformats::timestamp_t m_previous_ts
std::chrono::time_point< std::chrono::high_resolution_clock > m_t0
bool m_ts_problem_reported
std::atomic< uint64_t > m_num_new_tps
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
bool m_first_seq_id_mismatch
void scrap_source_and_geo_ids()
std::map< uint, std::atomic< int > > m_tp_channel_rate_map
std::atomic< uint64_t > m_frame_counter
uint32_t m_tp_count_limit
void configure_source_and_geo_ids(const appmodel::DataHandlerModule *conf)
bool m_seq_id_problem_reported
void sequence_check(frameptr fp)
std::atomic< uint64_t > m_tps_send_failed
ReadoutTypeAdapter * frameptr
std::unordered_map< trgdataformats::channel_t, unsigned int > m_channel_plane_map
void generate_opmon_data() override
std::atomic< uint64_t > m_seq_id_error_ctr
void start(const appfwk::DAQModule::CommandData_t &args) override
Start operation.
std::atomic< uint64_t > m_ts_error_ctr
void configure_preprocessing(const appmodel::DataHandlerModule *conf)
bool m_seq_id_error_state
void conf(const appmodel::DataHandlerModule *conf) override
Set the emulator mode, if active, timestamps of processed packets are overwritten with new ones.
std::atomic< int16_t > m_seq_id_max_jump
TPCEthFrameProcessor(std::unique_ptr< datahandlinglibs::FrameErrorRegistry > &error_registry, bool processing_enabled)
dunedaq::daqdataformats::timestamp_t m_pattern_generator_current_ts
void timestamp_check(frameptr fp)
uint32_t m_frame_count_at_last_send
dunedaq::daqdataformats::timestamp_t m_pattern_generator_previous_ts
std::set< unsigned int > m_channel_mask_set
std::atomic< int16_t > m_seq_id_min_jump
dunedaq::daqdataformats::timestamp_t m_current_ts
void find_tps(constframeptr fp)
uint16_t m_previous_seq_id
static constexpr timeout_t s_no_block
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
Base class for any user define issue.
#define TLOG_DEBUG(lvl,...)
Both frame_count_limit and tp_count_limit were set FailedToSendTPVector
FrameAndTPCountersDisabled
void warning(const Issue &issue)
void error(const Issue &issue)
stats(obj, sel_links, seconds, show_udp, show_buf)
trgdataformats::TriggerPrimitive tp