diff --git a/CMakeLists.txt b/CMakeLists.txt index 7cfe13cd..8b2ec0b2 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1,5 +1,5 @@ cmake_minimum_required(VERSION 3.12) -project(dfmodules VERSION 6.0.1) +project(dfmodules VERSION 6.0.2) find_package(daq-cmake REQUIRED) daq_setup_environment() @@ -23,10 +23,10 @@ daq_protobuf_codegen( opmon/*.proto ) ############################################################################## daq_add_library( TriggerInhibitAgent.cpp TriggerRecordBuilderData.cpp TPBundleHandler.cpp - LINK_LIBRARIES + LINK_LIBRARIES opmonlib::opmonlib ers::ers HighFive appfwk::appfwk logging::logging stdc++fs dfmessages::dfmessages utilities::utilities trigger::trigger detdataformats::detdataformats trgdataformats::trgdataformats) - + daq_add_plugin( HDF5DataStore duneDataStore LINK_LIBRARIES dfmodules hdf5libs::hdf5libs stdc++fs) daq_add_plugin( FragmentAggregatorModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager ) diff --git a/include/dfmodules/DataStore.hpp b/include/dfmodules/DataStore.hpp index 5d10b91a..f754937a 100644 --- a/include/dfmodules/DataStore.hpp +++ b/include/dfmodules/DataStore.hpp @@ -16,13 +16,13 @@ #define DFMODULES_INCLUDE_DFMODULES_DATASTORE_HPP_ #include "appfwk/ConfigurationManager.hpp" -#include "opmonlib/MonitorableObject.hpp" #include "cetlib/BasicPluginFactory.h" #include "cetlib/compiler_macros.h" #include "daqdataformats/TimeSlice.hpp" #include "daqdataformats/TriggerRecord.hpp" #include "daqdataformats/Types.hpp" #include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG< make(const std::string& name, \ - std::shared_ptr mcfg, \ - const std::string& writer_name ) \ + std::shared_ptr mcfg, \ + const std::string& writer_name) \ { \ return std::shared_ptr(new klass(name, mcfg, writer_name)); \ } \ @@ -105,7 +105,9 @@ namespace dfmodules { /** * @brief comment */ -class DataStore : public utilities::NamedObject, public opmonlib::MonitorableObject +class DataStore + : public utilities::NamedObject + , public opmonlib::MonitorableObject { public: /** @@ -113,7 +115,8 @@ class DataStore : public utilities::NamedObject, public opmonlib::MonitorableObj * @param name Name of the DataStore instance */ explicit DataStore(const std::string& name) - : utilities::NamedObject(name), MonitorableObject() + : utilities::NamedObject(name) + , MonitorableObject() { } @@ -135,8 +138,7 @@ class DataStore : public utilities::NamedObject, public opmonlib::MonitorableObj * This allows DataStore instances to make any preparations that will be * beneficial in advance of the first data blocks being written or read. */ - virtual void prepare_for_run(daqdataformats::run_number_t run_number, - bool run_is_for_test_purposes) = 0; + virtual void prepare_for_run(daqdataformats::run_number_t run_number, bool run_is_for_test_purposes) = 0; /** * @brief Informs the DataStore that writes or reads of data blocks associated @@ -164,7 +166,7 @@ inline std::shared_ptr make_data_store(const std::string& type, const std::string& name, std::shared_ptr mcfg, - const std::string& writer_identifier) + const std::string& writer_identifier) { static cet::BasicPluginFactory bpf("duneDataStore", "make"); // NOLINT diff --git a/plugins/DFOModule.cpp b/plugins/DFOModule.cpp index 6c66e37e..0d5b9969 100644 --- a/plugins/DFOModule.cpp +++ b/plugins/DFOModule.cpp @@ -93,7 +93,6 @@ DFOModule::init(std::shared_ptr mcfg) } if (m_busy_sender == nullptr) { throw appfwk::MissingConnection(ERS_HERE, get_name(), datatype_to_string(), "output"); - } m_dfo_conf = mdal->get_configuration(); @@ -149,7 +148,8 @@ DFOModule::do_start(const CommandData_t& payload) auto sender = iom->get_sender(trb_conn); if (sender != nullptr) { bool is_ready = sender->is_ready_for_sending(std::chrono::milliseconds(100)); - TLOG_DEBUG(0) << "The TriggerDecision sender for " << trb_conn << " " << (is_ready ? "is" : "is not") << " ready."; + TLOG_DEBUG(0) << "The TriggerDecision sender for " << trb_conn << " " << (is_ready ? "is" : "is not") + << " ready."; } } iom->add_callback( @@ -196,7 +196,7 @@ DFOModule::do_stop(const CommandData_t& /*args*/) std::lock_guard guard(m_trigger_counters_mutex); m_trigger_counters.clear(); - + TLOG() << get_name() << " successfully stopped"; TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Exiting do_stop() method"; } @@ -219,18 +219,17 @@ DFOModule::receive_trigger_decision(const dfmessages::TriggerDecision& decision) << decision.trigger_number << " and run " << decision.run_number << " (current run is " << m_run_number << ")"; if (decision.run_number != m_run_number) { - ers::error(DFOModuleRunNumberMismatch( - ERS_HERE, decision.run_number, m_run_number, "MLT", decision.trigger_number)); + ers::error(DFOModuleRunNumberMismatch(ERS_HERE, decision.run_number, m_run_number, "MLT", decision.trigger_number)); return; } auto decision_received = std::chrono::steady_clock::now(); ++m_received_decisions; auto trigger_types = unpack_types(decision.trigger_type); - for ( const auto t : trigger_types ) { + for (const auto t : trigger_types) { ++get_trigger_counter(t).received; } - + std::chrono::steady_clock::time_point decision_assigned; do { @@ -255,8 +254,7 @@ DFOModule::receive_trigger_decision(const dfmessages::TriggerDecision& decision) << " to connection " << assignment->connection_name; break; } else { - ers::error( - TRBModuleAppUpdate(ERS_HERE, assignment->connection_name, "Could not send Trigger Decision")); + ers::error(TRBModuleAppUpdate(ERS_HERE, assignment->connection_name, "Could not send Trigger Decision")); m_dataflow_availability[assignment->connection_name]->set_in_error(true); } @@ -340,11 +338,11 @@ DFOModule::find_slot(const dfmessages::TriggerDecision& decision) } void -DFOModule::generate_opmon_data() +DFOModule::generate_opmon_data() { opmon::DFOInfo info; - info.set_tokens_received( m_received_tokens.exchange(0) ); + info.set_tokens_received(m_received_tokens.exchange(0)); info.set_decisions_sent(m_sent_decisions.exchange(0)); info.set_decisions_received(m_received_decisions.exchange(0)); info.set_waiting_for_decision(m_waiting_for_decision.exchange(0)); @@ -352,16 +350,16 @@ DFOModule::generate_opmon_data() info.set_forwarding_decision(m_forwarding_decision.exchange(0)); info.set_waiting_for_token(m_waiting_for_token.exchange(0)); info.set_processing_token(m_processing_token.exchange(0)); - publish( std::move(info) ); + publish(std::move(info)); - std::lock_guard guard(m_trigger_counters_mutex); - for ( auto & [type, counts] : m_trigger_counters ) { + std::lock_guard guard(m_trigger_counters_mutex); + for (auto& [type, counts] : m_trigger_counters) { opmon::TriggerInfo ti; ti.set_received(counts.received.exchange(0)); ti.set_completed(counts.completed.exchange(0)); auto name = dunedaq::trgdataformats::get_trigger_candidate_type_names()[type]; - publish( std::move(ti), {{"type", name}} ); - } + publish(std::move(ti), { { "type", name } }); + } } void @@ -382,14 +380,14 @@ DFOModule::receive_trigger_complete_token(const dfmessages::TriggerDecisionToken } TLOG_DEBUG(TLVL_TDTOKEN_RECEIVED) << get_name() << " Received TriggerDecisionToken for trigger_number " - << token.trigger_number << " and run " << token.run_number - << " (current run is " << m_run_number << ")"; + << token.trigger_number << " and run " << token.run_number << " (current run is " + << m_run_number << ")"; // add a check to see if the application data found if (token.run_number != m_run_number) { std::ostringstream oss_source; oss_source << "TRB at connection " << token.decision_destination; - ers::error(DFOModuleRunNumberMismatch( - ERS_HERE, token.run_number, m_run_number, oss_source.str(), token.trigger_number)); + ers::error( + DFOModuleRunNumberMismatch(ERS_HERE, token.run_number, m_run_number, oss_source.str(), token.trigger_number)); return; } @@ -406,7 +404,8 @@ DFOModule::receive_trigger_complete_token(const dfmessages::TriggerDecisionToken try { auto dec_ptr = app_it->second->complete_assignment(token.trigger_number, m_metadata_function); auto trigger_types = unpack_types(dec_ptr->decision.trigger_type); - for ( const auto t : trigger_types ) ++ get_trigger_counter(t).completed; + for (const auto t : trigger_types) + ++get_trigger_counter(t).completed; } catch (AssignedTriggerDecisionNotFound const& err) { ers::error(err); } diff --git a/plugins/DFOModule.hpp b/plugins/DFOModule.hpp index ad280b6c..653e941b 100644 --- a/plugins/DFOModule.hpp +++ b/plugins/DFOModule.hpp @@ -27,10 +27,10 @@ #include #include +#include #include #include #include -#include namespace dunedaq { @@ -45,9 +45,9 @@ ERS_DECLARE_ISSUE(dfmodules, ((std::string)connection_name)) ERS_DECLARE_ISSUE(dfmodules, DFOModuleRunNumberMismatch, - "DFOModule encountered run number mismatch: recvd (" - << received_run_number << ") != " << run_number << " from " << src_app << " for trigger_number " - << trig_num, + "DFOModule encountered run number mismatch: recvd (" << received_run_number << ") != " << run_number + << " from " << src_app << " for trigger_number " + << trig_num, ((uint32_t)received_run_number)((uint32_t)run_number)((std::string)src_app)( (uint32_t)trig_num)) // NOLINT(build/unsigned) ERS_DECLARE_ISSUE(dfmodules, @@ -80,11 +80,10 @@ class DFOModule : public dunedaq::appfwk::DAQModule */ explicit DFOModule(const std::string& name); - DFOModule(const DFOModule&) = delete; ///< DFOModule is not copy-constructible - DFOModule& operator=(const DFOModule&) = - delete; ///< DFOModule is not copy-assignable - DFOModule(DFOModule&&) = delete; ///< DFOModule is not move-constructible - DFOModule& operator=(DFOModule&&) = delete; ///< DFOModule is not move-assignable + DFOModule(const DFOModule&) = delete; ///< DFOModule is not copy-constructible + DFOModule& operator=(const DFOModule&) = delete; ///< DFOModule is not copy-assignable + DFOModule(DFOModule&&) = delete; ///< DFOModule is not move-constructible + DFOModule& operator=(DFOModule&&) = delete; ///< DFOModule is not move-assignable void init(std::shared_ptr mcfg) override; @@ -139,22 +138,25 @@ class DFOModule : public dunedaq::appfwk::DAQModule mutable std::mutex m_notify_trigger_mutex; // Struct for statistic - struct TriggerData { - std::atomic received{0}; - std::atomic completed{0}; + struct TriggerData + { + std::atomic received{ 0 }; + std::atomic completed{ 0 }; }; - static std::set - unpack_types( decltype(dfmessages::TriggerDecision::trigger_type) t) { + static std::set unpack_types( + decltype(dfmessages::TriggerDecision::trigger_type) t) + { std::set results; if (t == dfmessages::TypeDefaults::s_invalid_trigger_type) return results; const std::bitset<64> bits(t); - for( size_t i = 0; i < bits.size(); ++i ) { - if ( bits[i] ) results.insert((trgdataformats::TriggerCandidateData::Type)i); + for (size_t i = 0; i < bits.size(); ++i) { + if (bits[i]) + results.insert((trgdataformats::TriggerCandidateData::Type)i); } return results; } - + // Statistics std::atomic m_received_tokens{ 0 }; // NOLINT (build/unsigned) std::atomic m_sent_decisions{ 0 }; // NOLINT (build/unsigned) @@ -165,15 +167,16 @@ class DFOModule : public dunedaq::appfwk::DAQModule std::atomic m_waiting_for_token{ 0 }; // NOLINT (build/unsigned) std::atomic m_processing_token{ 0 }; // NOLINT (build/unsigned) std::map m_trigger_counters; - std::mutex m_trigger_counters_mutex; // used to safely handle the map above - TriggerData & get_trigger_counter(trgdataformats::TriggerCandidateData::Type type) { + std::mutex m_trigger_counters_mutex; // used to safely handle the map above + TriggerData& get_trigger_counter(trgdataformats::TriggerCandidateData::Type type) + { auto it = m_trigger_counters.find(type); - if (it != m_trigger_counters.end()) return it->second; - + if (it != m_trigger_counters.end()) + return it->second; + std::lock_guard guard(m_trigger_counters_mutex); return m_trigger_counters[type]; } - }; } // namespace dfmodules } // namespace dunedaq diff --git a/plugins/FragmentAggregatorModule.cpp b/plugins/FragmentAggregatorModule.cpp index f395a5b4..7329ef69 100644 --- a/plugins/FragmentAggregatorModule.cpp +++ b/plugins/FragmentAggregatorModule.cpp @@ -10,8 +10,8 @@ #include "dfmodules/CommonIssues.hpp" #include "dfmodules/opmon/FragmentAggregatorModule.pb.h" -#include "appmodel/FragmentAggregatorModule.hpp" #include "appmodel/FragmentAggregatorConf.hpp" +#include "appmodel/FragmentAggregatorModule.hpp" #include "confmodel/Connection.hpp" #include "confmodel/QueueWithSourceId.hpp" #include "daqdataformats/FragmentHeader.hpp" diff --git a/plugins/FragmentAggregatorModule.hpp b/plugins/FragmentAggregatorModule.hpp index 4ebcc98d..68cd54a5 100644 --- a/plugins/FragmentAggregatorModule.hpp +++ b/plugins/FragmentAggregatorModule.hpp @@ -1,6 +1,6 @@ /** - * @file FragmentAggregatorModule.hpp Module to dispatch data requests within an application, aggregate and send fragments - * using the IOMManager + * @file FragmentAggregatorModule.hpp Module to dispatch data requests within an application, aggregate and send + * fragments using the IOMManager * * This is part of the DUNE DAQ , copyright 2020. * Licensing/copyright details are in the COPYING file that you should have @@ -39,15 +39,12 @@ ERS_DECLARE_ISSUE(dfmodules, ///< Namespace ((daqdataformats::SourceID)src) ///< Message parameters ) - -ERS_DECLARE_ISSUE(dfmodules, ///< Namespace - AbandonedFragment, ///< Issue class name - "Fragment from " << source << " for trigger " << trigger << '-' << sequence << " of run " << run << " was dropped", - ((daqdataformats::run_number_t)run) - ((daqdataformats::trigger_number_t)trigger) - ((daqdataformats::sequence_number_t)sequence) - ((daqdataformats::SourceID)source) -) +ERS_DECLARE_ISSUE(dfmodules, ///< Namespace + AbandonedFragment, ///< Issue class name + "Fragment from " << source << " for trigger " << trigger << '-' << sequence << " of run " << run + << " was dropped", + ((daqdataformats::run_number_t)run)((daqdataformats::trigger_number_t)trigger)( + (daqdataformats::sequence_number_t)sequence)((daqdataformats::SourceID)source)) namespace dfmodules { diff --git a/plugins/HDF5DataStore.hpp b/plugins/HDF5DataStore.hpp index 1ce419e6..ee6d2b7f 100644 --- a/plugins/HDF5DataStore.hpp +++ b/plugins/HDF5DataStore.hpp @@ -193,10 +193,11 @@ class HDF5DataStore : public DataStore if (m_compression_level != 0 && m_recorded_size != 0) { // Without compression, the uncompressed raw data size is approximately the total file size, so it // serves as an approximation of what would have been written without compression - float compression_factor = (float) m_file_handle->get_uncompressed_raw_data_size() / m_file_handle->get_total_file_size(); + float compression_factor = + (float)m_file_handle->get_uncompressed_raw_data_size() / m_file_handle->get_total_file_size(); size_of_next_write = tr_size / compression_factor; } - if (! increment_file_index_if_needed(size_of_next_write)) { + if (!increment_file_index_if_needed(size_of_next_write)) { if (m_operation_mode == "one-event-per-file") { if (m_current_record_number != std::numeric_limits::max() && tr.get_header_ref().get_trigger_number() != m_current_record_number) { @@ -260,10 +261,11 @@ class HDF5DataStore : public DataStore if (m_compression_level != 0 && m_recorded_size != 0) { // Without compression, the uncompressed raw data size is approximately the total file size, so it // serves as an approximation of what would have been written without compression - float compression_factor = (float) m_file_handle->get_uncompressed_raw_data_size() / m_file_handle->get_total_file_size(); + float compression_factor = + (float)m_file_handle->get_uncompressed_raw_data_size() / m_file_handle->get_total_file_size(); size_of_next_write = ts_size / compression_factor; } - if (! increment_file_index_if_needed(size_of_next_write)) { + if (!increment_file_index_if_needed(size_of_next_write)) { if (m_operation_mode == "one-event-per-file") { if (m_current_record_number != std::numeric_limits::max() && ts.get_header().timeslice_number != m_current_record_number) { @@ -310,8 +312,7 @@ class HDF5DataStore : public DataStore * * This method may throw an exception if it finds a problem. */ - void prepare_for_run(daqdataformats::run_number_t run_number, - bool run_is_for_test_purposes) + void prepare_for_run(daqdataformats::run_number_t run_number, bool run_is_for_test_purposes) { m_run_number = run_number; m_run_is_for_test_purposes = run_is_for_test_purposes; @@ -499,8 +500,8 @@ class HDF5DataStore : public DataStore // (determined inside the HDF5RawDataFile constructor) could disagree. std::string unique_filename = file_name; if (!m_disable_unique_suffix) { - time_t now = time(0); - std::string file_creation_timestamp = boost::posix_time::to_iso_string(boost::posix_time::from_time_t(now)); + time_t now = time(0); + std::string file_creation_timestamp = boost::posix_time::to_iso_string(boost::posix_time::from_time_t(now)); // timestamp substring size_t ufn_len = unique_filename.length(); if (ufn_len > 6) { // len GT 6 gives us some confidence that we have at least x.hdf5 @@ -542,7 +543,8 @@ class HDF5DataStore : public DataStore // m_file_handle->write_attribute("data_format_version",(int)m_key_translator_ptr->get_current_version()); m_file_handle->write_attribute("operational_environment", (std::string)m_operational_environment); m_file_handle->write_attribute("offline_data_stream", (std::string)m_offline_data_stream); - m_file_handle->write_attribute("run_was_for_test_purposes", (std::string)(m_run_is_for_test_purposes ? "true" : "false")); + m_file_handle->write_attribute("run_was_for_test_purposes", + (std::string)(m_run_is_for_test_purposes ? "true" : "false")); } } else { TLOG_DEBUG(TLVL_BASIC) << get_name() << ": Pointer file to " << m_basic_name_of_open_file diff --git a/plugins/HDF5FileUtils.hpp b/plugins/HDF5FileUtils.hpp index d2ce375f..c397d793 100644 --- a/plugins/HDF5FileUtils.hpp +++ b/plugins/HDF5FileUtils.hpp @@ -22,10 +22,10 @@ #include #include -//#include "dfmodules/StorageKey.hpp" -//#include -//#include -//#include +// #include "dfmodules/StorageKey.hpp" +// #include +// #include +// #include namespace dunedaq { namespace dfmodules { diff --git a/plugins/TRBModule.cpp b/plugins/TRBModule.cpp index bd622f7b..379ebd30 100644 --- a/plugins/TRBModule.cpp +++ b/plugins/TRBModule.cpp @@ -109,7 +109,8 @@ TRBModule::init(std::shared_ptr mcfg) throw InvalidQueueFatalError(ERS_HERE, get_name(), "Fragment Input queue"); } - m_trigger_record_output = iom->get_sender>(mdal->get_trigger_record_output()->UID()); + m_trigger_record_output = + iom->get_sender>(mdal->get_trigger_record_output()->UID()); for (auto con : mdal->get_request_connections()) { for (auto source_id : con->get_source_ids()) { @@ -442,7 +443,7 @@ TRBModule::trigger_decision_callback(dfmessages::TriggerDecision& td) ++m_received_trigger_decisions; create_trigger_records_and_dispatch(td); - //check_stale_requests(); + // check_stale_requests(); auto end_time = std::chrono::steady_clock::now(); m_td_processing_us += std::chrono::duration_cast(end_time - start_time).count(); @@ -477,7 +478,9 @@ TRBModule::extract_trigger_record(const TriggerId& id) m_pending_fragment_counter -= missing_fragments; temp->get_header_ref().set_status_bit(TriggerRecordStatusBits::kIncomplete, true); - ers::error(IncompleteTriggerRecord(ERS_HERE, (m_stop_requested.load() ? "at Stop time " : ""), id, + ers::error(IncompleteTriggerRecord(ERS_HERE, + (m_stop_requested.load() ? "at Stop time " : ""), + id, temp->get_fragments_ref().size(), temp->get_header_ref().get_num_requested_components())); } @@ -667,7 +670,7 @@ TRBModule::send_trigger_record(const TriggerId& id) std::set sent_destinations; while (it != m_mon_requests.end()) { - // Only sent TR to each monitor once + // Only sent TR to each monitor once if (sent_destinations.count(it->data_destination)) { ++it; continue; @@ -683,7 +686,8 @@ TRBModule::send_trigger_record(const TriggerId& id) auto trigger_record_bytes = serialization::serialize(temp_record, serialization::SerializationType::kMsgPack); trigger_record_ptr_t record_copy = serialization::deserialize(trigger_record_bytes); - iom->get_sender(it->data_destination)->send(std::move(record_copy), m_tr_queue_timeout); + iom->get_sender(it->data_destination) + ->send(std::move(record_copy), m_tr_queue_timeout); ++m_trmon_sent_counter; wasSentSuccessfully = true; } catch (const ers::Issue& excpt) { diff --git a/plugins/TRBModule.hpp b/plugins/TRBModule.hpp index 55977ea3..abb532d8 100644 --- a/plugins/TRBModule.hpp +++ b/plugins/TRBModule.hpp @@ -9,23 +9,23 @@ #ifndef DFMODULES_PLUGINS_TRIGGERRECORDBUILDER_HPP_ #define DFMODULES_PLUGINS_TRIGGERRECORDBUILDER_HPP_ +#include "appmodel/ReadoutApplication.hpp" +#include "appmodel/SmartDaqApplication.hpp" #include "appmodel/TRBConf.hpp" #include "daqdataformats/Fragment.hpp" #include "daqdataformats/SourceID.hpp" #include "daqdataformats/TriggerRecord.hpp" #include "daqdataformats/Types.hpp" -#include "appmodel/ReadoutApplication.hpp" -#include "appmodel/SmartDaqApplication.hpp" #include "dfmessages/DataRequest.hpp" #include "dfmessages/TRMonRequest.hpp" #include "dfmessages/TriggerDecision.hpp" #include "dfmessages/Types.hpp" #include "appfwk/DAQModule.hpp" -#include "utilities/WorkerThread.hpp" -#include "iomanager/Sender.hpp" #include "iomanager/Receiver.hpp" +#include "iomanager/Sender.hpp" #include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG< mcfg) override; @@ -214,8 +216,7 @@ class TRBModule : public dunedaq::appfwk::DAQModule unsigned int create_trigger_records_and_dispatch(const dfmessages::TriggerDecision&); - bool dispatch_data_requests(dfmessages::DataRequest, - const daqdataformats::SourceID&); + bool dispatch_data_requests(dfmessages::DataRequest, const daqdataformats::SourceID&); bool send_trigger_record(const TriggerId&); // this creates a trigger record and send it @@ -233,7 +234,7 @@ class TRBModule : public dunedaq::appfwk::DAQModule void do_stop(const CommandData_t&); // Monitoring callback - void tr_requested(const dfmessages::TRMonRequest &); + void tr_requested(const dfmessages::TRMonRequest&); // Threading std::atomic m_stop_requested; @@ -253,7 +254,8 @@ class TRBModule : public dunedaq::appfwk::DAQModule // Output connections std::shared_ptr m_trigger_record_output; mutable std::mutex m_map_sourceid_connections_mutex; - std::map> m_map_sourceid_connections; ///< Mappinng between SourceID and connections + std::map> + m_map_sourceid_connections; ///< Mappinng between SourceID and connections // bookeeping using clock_type = std::chrono::steady_clock; @@ -297,7 +299,6 @@ class TRBModule : public dunedaq::appfwk::DAQModule mutable std::atomic m_td_processing_us = { 0 }; // in between calls mutable std::atomic m_fragment_processing_us = { 0 }; // in between calls - mutable std::atomic m_trmon_request_counter = { 0 }; mutable std::atomic m_trmon_sent_counter = { 0 }; diff --git a/plugins/TRMonRequestorModule.hpp b/plugins/TRMonRequestorModule.hpp old mode 100755 new mode 100644 index 236a1a51..46b06716 --- a/plugins/TRMonRequestorModule.hpp +++ b/plugins/TRMonRequestorModule.hpp @@ -11,17 +11,17 @@ #include "appfwk/DAQModule.hpp" #include "appmodel/TRMonRequestorConf.hpp" -#include "iomanager/IOManager.hpp" -#include "dfmessages/TriggerDecisionToken.hpp" #include "dfmessages/TRMonRequest.hpp" -#include "utilities/WorkerThread.hpp" +#include "dfmessages/TriggerDecisionToken.hpp" #include "dfmodules/opmon/TRMonRequestorModule.pb.h" +#include "iomanager/IOManager.hpp" +#include "utilities/WorkerThread.hpp" #include -#include #include -#include #include +#include +#include namespace dunedaq::dfmodules { class TRMonRequestorModule : public appfwk::DAQModule @@ -44,7 +44,6 @@ class TRMonRequestorModule : public appfwk::DAQModule void generate_opmon_data() override; protected: - using token_receiver_t = iomanager::ReceiverConcept; using trmon_sender_t = iomanager::SenderConcept; @@ -82,8 +81,9 @@ class TRMonRequestorModule : public appfwk::DAQModule std::shared_ptr m_token_receiver; // Monitoring - using const_metric_counter_t = std::invoke_result::type; + using const_metric_counter_t = + std::invoke_result::type; using metric_counter_t = std::remove_const::type; std::atomic m_trigger_records_requested{ 0 }; }; diff --git a/src/TPBundleHandler.cpp b/src/TPBundleHandler.cpp index bc80e3b7..ace111fa 100644 --- a/src/TPBundleHandler.cpp +++ b/src/TPBundleHandler.cpp @@ -137,14 +137,13 @@ TPBundleHandler::add_tpset(trigger::TPSet&& tpset) if (tsidx_from_begin_time <= m_slice_index_offset) { auto lk = std::lock_guard(m_accumulator_map_mutex); int64_t diff = static_cast(tsidx_from_begin_time) - static_cast(m_slice_index_offset); - if (! m_one_or_more_time_slices_have_aged_out) { - TLOG() << "Updating the slice numbers of existing accumulators by " << (1-diff); + if (!m_one_or_more_time_slices_have_aged_out) { + TLOG() << "Updating the slice numbers of existing accumulators by " << (1 - diff); for (auto& [local_tsidx, local_accum] : m_timeslice_accumulators) { local_accum.update_slice_number(1 - diff); } m_slice_index_offset -= (1 - diff); - } - else { + } else { ers::warning(TardyTPSetReceived(ERS_HERE, tpset.origin.id, tpset.start_time, diff)); return; } @@ -155,10 +154,8 @@ TPBundleHandler::add_tpset(trigger::TPSet&& tpset) { auto lk = std::lock_guard(m_accumulator_map_mutex); if (m_timeslice_accumulators.count(tsidx) == 0) { - TimeSliceAccumulator accum(tsidx * m_slice_interval, - (tsidx + 1) * m_slice_interval, - tsidx - m_slice_index_offset, - m_run_number); + TimeSliceAccumulator accum( + tsidx * m_slice_interval, (tsidx + 1) * m_slice_interval, tsidx - m_slice_index_offset, m_run_number); m_timeslice_accumulators[tsidx] = accum; } } @@ -199,7 +196,9 @@ TPBundleHandler::get_properly_aged_timeslices() m_timeslice_accumulators.erase(tsidx); } - if (list_of_timeslices.size() > 0) {m_one_or_more_time_slices_have_aged_out = true;} + if (list_of_timeslices.size() > 0) { + m_one_or_more_time_slices_have_aged_out = true; + } return list_of_timeslices; } diff --git a/src/TriggerInhibitAgent.cpp b/src/TriggerInhibitAgent.cpp index 69024fca..18288fe0 100644 --- a/src/TriggerInhibitAgent.cpp +++ b/src/TriggerInhibitAgent.cpp @@ -39,7 +39,8 @@ TriggerInhibitAgent::TriggerInhibitAgent(const std::string& parent_name, , m_trigger_inhibit_sender(our_output) , m_trigger_number_at_start_of_processing_chain(0) , m_trigger_number_at_end_of_processing_chain(0) -{} +{ +} void TriggerInhibitAgent::start_checking() diff --git a/src/TriggerRecordBuilderData.cpp b/src/TriggerRecordBuilderData.cpp index 7f44ac7d..bd13b9c0 100644 --- a/src/TriggerRecordBuilderData.cpp +++ b/src/TriggerRecordBuilderData.cpp @@ -33,7 +33,8 @@ TriggerRecordBuilderData::TriggerRecordBuilderData(std::string connection_name, , m_is_busy(false) , m_in_error(false) , m_connection_name(connection_name) -{} +{ +} TriggerRecordBuilderData::TriggerRecordBuilderData(std::string connection_name, size_t busy_threshold, @@ -104,8 +105,7 @@ TriggerRecordBuilderData::complete_assignment(daqdataformats::trigger_number_t t metadata_fun(m_metadata); ++m_complete_counter; - auto completion_time = - std::chrono::duration_cast(now - dec_ptr->assigned_time); + auto completion_time = std::chrono::duration_cast(now - dec_ptr->assigned_time); if (completion_time.count() < m_min_complete_time.load()) m_min_complete_time.store(completion_time.count()); if (completion_time.count() > m_max_complete_time.load()) @@ -113,11 +113,11 @@ TriggerRecordBuilderData::complete_assignment(daqdataformats::trigger_number_t t opmon::TRCompleteInfo i; i.set_completion_time(completion_time.count()); - i.set_tr_number( dec_ptr->decision.trigger_number ); - i.set_run_number( dec_ptr->decision.run_number ); - i.set_trigger_type( dec_ptr->decision.trigger_type ); - publish( std::move(i), {}, opmonlib::to_level(opmonlib::EntryOpMonLevel::kEventDriven) ); - + i.set_tr_number(dec_ptr->decision.trigger_number); + i.set_run_number(dec_ptr->decision.run_number); + i.set_trigger_type(dec_ptr->decision.trigger_type); + publish(std::move(i), {}, opmonlib::to_level(opmonlib::EntryOpMonLevel::kEventDriven)); + return dec_ptr; } @@ -166,10 +166,10 @@ TriggerRecordBuilderData::add_assignment(std::shared_ptr::max() ); + info.set_min_time_since_assignment(std::numeric_limits::max()); info.set_max_time_since_assignment(0); time_counter_t time = 0; @@ -187,22 +187,22 @@ TriggerRecordBuilderData::generate_opmon_data() info.set_max_time_since_assignment(us_since_assignment.count()); } lk.unlock(); - + info.set_total_time_since_assignment(time); // estimate of the capcity auto completed_trigger_records = m_complete_counter.exchange(0); - if ( completed_trigger_records > 0 ) { - m_last_average_time = 1e-6*0.5*(m_min_complete_time.exchange(0) + m_max_complete_time.exchange(0)); // in seconds + if (completed_trigger_records > 0) { + m_last_average_time = + 1e-6 * 0.5 * (m_min_complete_time.exchange(0) + m_max_complete_time.exchange(0)); // in seconds } - if ( m_last_average_time > 0. ) { + if (m_last_average_time > 0.) { // prediction rate metrics - info.set_capacity_rate( 0.5*(m_busy_threshold.load()+m_free_threshold.load())/m_last_average_time ); + info.set_capacity_rate(0.5 * (m_busy_threshold.load() + m_free_threshold.load()) / m_last_average_time); } - + publish(std::move(info)); - } std::chrono::microseconds diff --git a/src/dfmodules/TPBundleHandler.hpp b/src/dfmodules/TPBundleHandler.hpp index 81e59fe0..aa72f99a 100644 --- a/src/dfmodules/TPBundleHandler.hpp +++ b/src/dfmodules/TPBundleHandler.hpp @@ -14,10 +14,10 @@ #include "daqdataformats/TimeSlice.hpp" #include "daqdataformats/Types.hpp" -#include "trgdataformats/TriggerPrimitive.hpp" #include "ers/Issue.hpp" -#include "trigger/TPSet.hpp" #include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG< #include @@ -44,8 +44,8 @@ ERS_DECLARE_ISSUE(dfmodules, ERS_DECLARE_ISSUE(dfmodules, TardyTPSetReceived, "Received a TPSet with a timestamp that is too early compared to ones that have already " - << "been processed, sourceid=" << tpset_source_id << ", start_time=" << tpset_start_time - << ", the calculated timeslice_id is " << tsid, + << "been processed, sourceid=" << tpset_source_id << ", start_time=" << tpset_start_time + << ", the calculated timeslice_id is " << tsid, ((size_t)tpset_source_id)((daqdataformats::timestamp_t)tpset_start_time)((int64_t)tsid)) // Re-enable coverage checking LCOV_EXCL_STOP @@ -94,10 +94,7 @@ class TimeSliceAccumulator return m_update_time; } - void update_slice_number(int delta) - { - m_slice_number += delta; - } + void update_slice_number(int delta) { m_slice_number += delta; } private: daqdataformats::timestamp_t m_begin_time; diff --git a/src/dfmodules/TriggerInhibitAgent.hpp b/src/dfmodules/TriggerInhibitAgent.hpp index b9c8db70..4985e724 100644 --- a/src/dfmodules/TriggerInhibitAgent.hpp +++ b/src/dfmodules/TriggerInhibitAgent.hpp @@ -12,12 +12,12 @@ #ifndef DFMODULES_SRC_DFMODULES_TRIGGERINHIBITAGENT_HPP_ #define DFMODULES_SRC_DFMODULES_TRIGGERINHIBITAGENT_HPP_ -#include "iomanager/Sender.hpp" -#include "iomanager/Receiver.hpp" -#include "utilities/NamedObject.hpp" #include "daqdataformats/Types.hpp" #include "dfmessages/TriggerDecision.hpp" #include "dfmessages/TriggerInhibit.hpp" +#include "iomanager/Receiver.hpp" +#include "iomanager/Sender.hpp" +#include "utilities/NamedObject.hpp" #include "utilities/WorkerThread.hpp" #include @@ -43,8 +43,8 @@ class TriggerInhibitAgent : public utilities::NamedObject * @brief TriggerInhibitAgent Constructor */ explicit TriggerInhibitAgent(const std::string&, - std::shared_ptr, - std::shared_ptr); + std::shared_ptr, + std::shared_ptr); TriggerInhibitAgent(const TriggerInhibitAgent&) = delete; ///< TriggerInhibitAgent is not copy-constructible TriggerInhibitAgent& operator=(const TriggerInhibitAgent&) = delete; ///< TriggerInhibitAgent is not copy-assignable diff --git a/src/dfmodules/TriggerRecordBuilderData.hpp b/src/dfmodules/TriggerRecordBuilderData.hpp index 17cd4e2e..f8e64913 100644 --- a/src/dfmodules/TriggerRecordBuilderData.hpp +++ b/src/dfmodules/TriggerRecordBuilderData.hpp @@ -17,9 +17,9 @@ #include "dfmodules/opmon/TRBuilderData.pb.h" #include "ers/Issue.hpp" +#include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG< #include @@ -61,7 +61,8 @@ struct AssignedTriggerDecision : decision(dec) , assigned_time(std::chrono::steady_clock::now()) , connection_name(conn_name) - {} + { + } }; class TriggerRecordBuilderData : public opmonlib::MonitorableObject @@ -77,7 +78,7 @@ class TriggerRecordBuilderData : public opmonlib::MonitorableObject TriggerRecordBuilderData& operator=(TriggerRecordBuilderData&&) = delete; ~TriggerRecordBuilderData() = default; - + bool is_busy() const { return m_in_error || m_is_busy; } size_t used_slots() const { return m_assigned_trigger_decisions.size(); } @@ -118,12 +119,12 @@ class TriggerRecordBuilderData : public opmonlib::MonitorableObject // monitoring using metric_t = dunedaq::dfmodules::opmon::DFApplicationInfo; - using const_time_counter_t = std::invoke_result::type; + using const_time_counter_t = std::invoke_result::type; using time_counter_t = std::remove_const::type; std::atomic m_complete_counter{ 0 }; - std::atomic m_min_complete_time{ std::numeric_limits::max() }, m_max_complete_time{ 0 }; // in us - double m_last_average_time{0.}; + std::atomic m_min_complete_time{ std::numeric_limits::max() }, + m_max_complete_time{ 0 }; // in us + double m_last_average_time{ 0. }; }; } // namespace dfmodules } // namespace dunedaq diff --git a/unittest/DataStoreFactory_test.cxx b/unittest/DataStoreFactory_test.cxx index d197697d..78c1ba02 100644 --- a/unittest/DataStoreFactory_test.cxx +++ b/unittest/DataStoreFactory_test.cxx @@ -24,7 +24,6 @@ BOOST_AUTO_TEST_CASE(invalid_request) // we want to pass an invalid DataStore type and see if we get an exception BOOST_CHECK_THROW(make_data_store("dummy", "dummy", nullptr, "dummy_writer"), DataStoreCreationFailed); - } #if 0 diff --git a/unittest/HDF5Write_test.cxx b/unittest/HDF5Write_test.cxx index 4f791f73..fa05eee4 100644 --- a/unittest/HDF5Write_test.cxx +++ b/unittest/HDF5Write_test.cxx @@ -12,12 +12,12 @@ // We need the actual HDF5DataStore.hpp plugin header so the unit tests can access its exceptions #include "../plugins/HDF5DataStore.hpp" // NOLINT(build/include_path) -#include "appmodel/DataWriterModule.hpp" +#include "appmodel/DataStoreConf.hpp" #include "appmodel/DataWriterConf.hpp" +#include "appmodel/DataWriterModule.hpp" #include "appmodel/FilenameParams.hpp" #include "confmodel/DetectorConfig.hpp" #include "confmodel/Session.hpp" -#include "appmodel/DataStoreConf.hpp" #include "detdataformats/DetID.hpp" #define BOOST_TEST_MODULE HDF5Write_test // NOLINT @@ -208,7 +208,7 @@ BOOST_AUTO_TEST_CASE(WriteOneFile) data_store_conf_obj.set_by_val("directory_path", file_path); auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr, "dwm-01"); - + // write several events, each with several fragments for (int trigger_number = 1; trigger_number <= trigger_count; ++trigger_number) data_store_ptr->write(create_trigger_record(trigger_number, fragment_size, apa_count * link_count)); @@ -248,7 +248,7 @@ BOOST_AUTO_TEST_CASE(CheckWritingSuffix) data_store_conf_obj.set_by_val("directory_path", file_path); auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr, "dwm-01"); - + // write several events, each with several fragments for (int trigger_number = 1; trigger_number <= trigger_count; ++trigger_number) { data_store_ptr->write(create_trigger_record(trigger_number, fragment_size, apa_count * link_count)); @@ -285,8 +285,8 @@ BOOST_AUTO_TEST_CASE(NoDuplicateTimeSlices) auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr, "dwm-01"); - dunedaq::daqdataformats::TimeSlice timeslice {999, 999}; // timeslice #, run # - dunedaq::daqdataformats::TimeSlice identical_timeslice {999, 999}; + dunedaq::daqdataformats::TimeSlice timeslice{ 999, 999 }; // timeslice #, run # + dunedaq::daqdataformats::TimeSlice identical_timeslice{ 999, 999 }; data_store_ptr->write(timeslice); @@ -326,7 +326,6 @@ BOOST_AUTO_TEST_CASE(EnormousMaxFileSize) BOOST_CHECK_THROW(data_store_ptr->prepare_for_run(1, true), dunedaq::dfmodules::InsufficientDiskSpace); } - BOOST_AUTO_TEST_CASE(FileSizeLimitResultsInMultipleFiles) { std::string file_path(std::filesystem::temp_directory_path()); @@ -353,7 +352,7 @@ BOOST_AUTO_TEST_CASE(FileSizeLimitResultsInMultipleFiles) data_store_conf_obj.set_by_val("max_file_size", 3000000); // goal is 6 events per file auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr, "dwm-01"); - + // write several events, each with several fragments for (int trigger_number = 1; trigger_number <= trigger_count; ++trigger_number) data_store_ptr->write(create_trigger_record(trigger_number, fragment_size, apa_count * link_count)); @@ -396,7 +395,7 @@ BOOST_AUTO_TEST_CASE(SmallFileSizeLimitDataBlockListWrite) data_store_conf_obj.set_by_val("directory_path", file_path); data_store_conf_obj.set_by_val("max_file_size", 150000); // ~1.5 Fragment, ~0.3 TR - auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr,"dwm-01"); + auto data_store_ptr = make_data_store(data_store_conf->get_type(), data_store_conf->UID(), cfg.cfgMgr, "dwm-01"); // write several events, each with several fragments for (int trigger_number = 1; trigger_number <= trigger_count; ++trigger_number) diff --git a/unittest/TriggerRecordBuilderData_test.cxx b/unittest/TriggerRecordBuilderData_test.cxx index f381ee6b..a42b6631 100644 --- a/unittest/TriggerRecordBuilderData_test.cxx +++ b/unittest/TriggerRecordBuilderData_test.cxx @@ -7,8 +7,8 @@ * received with this code. */ -#include "opmonlib/TestOpMonManager.hpp" #include "dfmodules/TriggerRecordBuilderData.hpp" +#include "opmonlib/TestOpMonManager.hpp" #define BOOST_TEST_MODULE TriggerRecordBuilderData_test // NOLINT @@ -110,10 +110,10 @@ BOOST_AUTO_TEST_CASE(Assignments) BOOST_REQUIRE_EQUAL(trbd_p->used_slots(), 0); auto latency = - std::chrono::duration_cast(complete_time - assignment->assigned_time) - .count(); + std::chrono::duration_cast(complete_time - assignment->assigned_time).count(); - BOOST_REQUIRE_CLOSE(static_cast(trbd_p->average_latency(start_time).count()), static_cast(latency), 5); + BOOST_REQUIRE_CLOSE( + static_cast(trbd_p->average_latency(start_time).count()), static_cast(latency), 5); auto null_got_assignment = trbd_p->get_assignment(2); BOOST_REQUIRE_EQUAL(null_got_assignment, nullptr); @@ -125,7 +125,6 @@ BOOST_AUTO_TEST_CASE(Assignments) auto remnants = trbd_p->flush(); BOOST_REQUIRE_EQUAL(trbd_p->used_slots(), 0); BOOST_REQUIRE_EQUAL(remnants.size(), 1); - } BOOST_AUTO_TEST_CASE(Exceptions) @@ -185,6 +184,4 @@ BOOST_AUTO_TEST_CASE(Exceptions) trbd.add_assignment(err_assignment), NoSlotsAvailable, [](NoSlotsAvailable const&) { return true; }); } - - BOOST_AUTO_TEST_SUITE_END()