Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
@@ -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()
Expand All @@ -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 )
Expand Down
18 changes: 10 additions & 8 deletions include/dfmodules/DataStore.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<<issue wont work.
#include "opmonlib/MonitorableObject.hpp"
#include "utilities/NamedObject.hpp"

#include "nlohmann/json.hpp"
Expand All @@ -47,8 +47,8 @@
#define DEFINE_DUNE_DATA_STORE(klass) \
EXTERN_C_FUNC_DECLARE_START \
std::shared_ptr<dunedaq::dfmodules::DataStore> make(const std::string& name, \
std::shared_ptr<dunedaq::appfwk::ConfigurationManager> mcfg, \
const std::string& writer_name ) \
std::shared_ptr<dunedaq::appfwk::ConfigurationManager> mcfg, \
const std::string& writer_name) \
{ \
return std::shared_ptr<dunedaq::dfmodules::DataStore>(new klass(name, mcfg, writer_name)); \
} \
Expand Down Expand Up @@ -105,15 +105,18 @@ namespace dfmodules {
/**
* @brief comment
*/
class DataStore : public utilities::NamedObject, public opmonlib::MonitorableObject
class DataStore
: public utilities::NamedObject
, public opmonlib::MonitorableObject
{
public:
/**
* @brief DataStore Constructor
* @param name Name of the DataStore instance
*/
explicit DataStore(const std::string& name)
: utilities::NamedObject(name), MonitorableObject()
: utilities::NamedObject(name)
, MonitorableObject()
{
}

Expand All @@ -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
Expand Down Expand Up @@ -164,7 +166,7 @@ inline std::shared_ptr<DataStore>
make_data_store(const std::string& type,
const std::string& name,
std::shared_ptr<dunedaq::appfwk::ConfigurationManager> mcfg,
const std::string& writer_identifier)
const std::string& writer_identifier)
{
static cet::BasicPluginFactory bpf("duneDataStore", "make"); // NOLINT

Expand Down
41 changes: 20 additions & 21 deletions plugins/DFOModule.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,6 @@ DFOModule::init(std::shared_ptr<appfwk::ConfigurationManager> mcfg)
}
if (m_busy_sender == nullptr) {
throw appfwk::MissingConnection(ERS_HERE, get_name(), datatype_to_string<dfmessages::TriggerInhibit>(), "output");

}

m_dfo_conf = mdal->get_configuration();
Expand Down Expand Up @@ -149,7 +148,8 @@ DFOModule::do_start(const CommandData_t& payload)
auto sender = iom->get_sender<dfmessages::TriggerDecision>(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<dfmessages::TriggerDecisionToken>(
Expand Down Expand Up @@ -196,7 +196,7 @@ DFOModule::do_stop(const CommandData_t& /*args*/)

std::lock_guard<std::mutex> 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";
}
Expand All @@ -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 {

Expand All @@ -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);
}

Expand Down Expand Up @@ -340,28 +338,28 @@ 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));
info.set_deciding_destination(m_deciding_destination.exchange(0));
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<std::mutex> guard(m_trigger_counters_mutex);
for ( auto & [type, counts] : m_trigger_counters ) {
std::lock_guard<std::mutex> 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
Expand All @@ -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;
}

Expand All @@ -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);
}
Expand Down
47 changes: 25 additions & 22 deletions plugins/DFOModule.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,10 +27,10 @@

#include <map>
#include <memory>
#include <mutex>
#include <string>
#include <utility>
#include <vector>
#include <mutex>

namespace dunedaq {

Expand All @@ -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,
Expand Down Expand Up @@ -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<appfwk::ConfigurationManager> mcfg) override;

Expand Down Expand Up @@ -139,22 +138,25 @@ class DFOModule : public dunedaq::appfwk::DAQModule
mutable std::mutex m_notify_trigger_mutex;

// Struct for statistic
struct TriggerData {
std::atomic<uint64_t> received{0};
std::atomic<uint64_t> completed{0};
struct TriggerData
{
std::atomic<uint64_t> received{ 0 };
std::atomic<uint64_t> completed{ 0 };
};
static std::set<trgdataformats::TriggerCandidateData::Type>
unpack_types( decltype(dfmessages::TriggerDecision::trigger_type) t) {
static std::set<trgdataformats::TriggerCandidateData::Type> unpack_types(
decltype(dfmessages::TriggerDecision::trigger_type) t)
{
std::set<trgdataformats::TriggerCandidateData::Type> 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<uint64_t> m_received_tokens{ 0 }; // NOLINT (build/unsigned)
std::atomic<uint64_t> m_sent_decisions{ 0 }; // NOLINT (build/unsigned)
Expand All @@ -165,15 +167,16 @@ class DFOModule : public dunedaq::appfwk::DAQModule
std::atomic<uint64_t> m_waiting_for_token{ 0 }; // NOLINT (build/unsigned)
std::atomic<uint64_t> m_processing_token{ 0 }; // NOLINT (build/unsigned)
std::map<dunedaq::trgdataformats::TriggerCandidateData::Type, TriggerData> 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<std::mutex> guard(m_trigger_counters_mutex);
return m_trigger_counters[type];
}

};
} // namespace dfmodules
} // namespace dunedaq
Expand Down
2 changes: 1 addition & 1 deletion plugins/FragmentAggregatorModule.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
19 changes: 8 additions & 11 deletions plugins/FragmentAggregatorModule.hpp
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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 {

Expand Down
Loading
Loading