From 8b6676f4b8cf4018424673bc0690a2c4e65c7e3f Mon Sep 17 00:00:00 2001 From: Eric Flumerfelt Date: Mon, 13 Jul 2026 07:56:12 -0500 Subject: [PATCH 1/2] Remove unused TokenManager code --- CMakeLists.txt | 5 +- plugins/MLTModule.hpp | 1 - plugins/RandomTCMakerModule.hpp | 2 - src/TokenManager.cpp | 104 -------------------------------- src/trigger/TokenManager.hpp | 102 ------------------------------- unittest/TokenManager_test.cxx | 93 ---------------------------- 6 files changed, 2 insertions(+), 305 deletions(-) delete mode 100644 src/TokenManager.cpp delete mode 100644 src/trigger/TokenManager.hpp delete mode 100644 unittest/TokenManager_test.cxx diff --git a/CMakeLists.txt b/CMakeLists.txt index bd2d02f7..386a95f2 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -26,7 +26,7 @@ find_package(CLI11 REQUIRED) daq_protobuf_codegen( opmon/*.proto ) ############################################################################## -# Main library +# Main library set(TRIGGER_DEPENDENCIES appmodel::appmodel appfwk::appfwk @@ -54,7 +54,7 @@ daq_add_library( ## merge a chance to succeed. #daq_codegen( # txbuffer.jsonnet DEP_PKGS datahandlinglibs TEMPLATES Structs.hpp.j2 Nljs.hpp.j2 ) - + ############################################################################## # Plugins @@ -80,7 +80,6 @@ daq_add_application( generate_tpset_from_hdf5 generate_tpset_from_hdf5.cxx TEST ############################################################################## # Unit Tests -#daq_add_unit_test(TokenManager_test LINK_LIBRARIES trigger) #daq_add_unit_test(TxSet_test LINK_LIBRARIES trigger) #daq_add_unit_test(BufferManager_test LINK_LIBRARIES trigger) #daq_add_unit_test(TriggerObjectOverlay_test LINK_LIBRARIES trigger) diff --git a/plugins/MLTModule.hpp b/plugins/MLTModule.hpp index 677d0e4e..d4c18206 100644 --- a/plugins/MLTModule.hpp +++ b/plugins/MLTModule.hpp @@ -17,7 +17,6 @@ #include "trigger/Issues.hpp" #include "trigger/Latency.hpp" #include "trigger/LivetimeCounter.hpp" -#include "trigger/TokenManager.hpp" #include "trigger/opmon/latency_info.pb.h" #include "trigger/opmon/moduleleveltrigger_info.pb.h" diff --git a/plugins/RandomTCMakerModule.hpp b/plugins/RandomTCMakerModule.hpp index 63ebaa39..44b102a1 100644 --- a/plugins/RandomTCMakerModule.hpp +++ b/plugins/RandomTCMakerModule.hpp @@ -7,8 +7,6 @@ #ifndef TRIGGER_PLUGINS_RANDOMTRIGGERCANDIDATEMAKER_HPP_ #define TRIGGER_PLUGINS_RANDOMTRIGGERCANDIDATEMAKER_HPP_ -#include "trigger/TokenManager.hpp" - #include "appfwk/ConfigurationManager.hpp" #include "appfwk/DAQModule.hpp" #include "confmodel/Connection.hpp" diff --git a/src/TokenManager.cpp b/src/TokenManager.cpp deleted file mode 100644 index 729d0ad9..00000000 --- a/src/TokenManager.cpp +++ /dev/null @@ -1,104 +0,0 @@ -/** - * @file TokenManager.cpp - * - * This is part of the DUNE DAQ Application Framework, copyright 2020. - * Licensing/copyright details are in the COPYING file that you should have - * received with this code. - */ - -#include "trigger/TokenManager.hpp" -#include "trigger/LivetimeCounter.hpp" - -#include "iomanager/IOManager.hpp" - -#include -#include - -namespace dunedaq::trigger { - -TokenManager::TokenManager(const std::string& connection_name, - int initial_tokens, - daqdataformats::run_number_t run_number, - std::shared_ptr livetime_counter) - : m_connection_name(connection_name) - , m_n_tokens(initial_tokens) - , m_run_number(run_number) - , m_livetime_counter(livetime_counter) - , m_token_receiver(nullptr) -{ - m_open_trigger_time = std::chrono::steady_clock::now(); - - m_token_receiver = get_iom_receiver(m_connection_name); - m_token_receiver->add_callback(std::bind(&TokenManager::receive_token, this, std::placeholders::_1)); -} - -TokenManager::~TokenManager() -{ - m_token_receiver->remove_callback(); - - if (!m_open_trigger_decisions.empty()) { - - auto now = std::chrono::steady_clock::now(); - if (std::chrono::duration_cast(now - m_open_trigger_time) > - std::chrono::milliseconds(3000)) { - std::ostringstream o; - o << "Open Trigger Decisions: ["; - { // Scope for lock_guard - bool first = true; - std::lock_guard lk(m_open_trigger_decisions_mutex); - for (auto& td : m_open_trigger_decisions) { - if (!first) - o << ", "; - o << td; - first = false; - } - o << "]"; - } - TLOG_DEBUG(0) << o.str(); - } - } -} - -int -TokenManager::get_n_tokens() const -{ - return m_n_tokens.load(); -} - -void -TokenManager::trigger_sent(dfmessages::trigger_number_t trigger_number) -{ - std::lock_guard lk(m_open_trigger_decisions_mutex); - m_open_trigger_decisions.insert(trigger_number); - m_n_tokens--; - if (m_n_tokens.load() == 0) { - m_livetime_counter->set_state(LivetimeCounter::State::kDead); - } -} - -void -TokenManager::receive_token(dfmessages::TriggerDecisionToken& token) -{ - TLOG_DEBUG(1) << "Received token with run number " << token.run_number << ", current run number " << m_run_number; - if (token.run_number == m_run_number) { - if (m_n_tokens.load() == 0) { - m_livetime_counter->set_state(LivetimeCounter::State::kLive); - } - m_n_tokens++; - TLOG_DEBUG(1) << "There are now " << m_n_tokens.load() << " tokens available"; - - if (token.trigger_number != dfmessages::TypeDefaults::s_invalid_trigger_number) { - if (m_open_trigger_decisions.count(token.trigger_number)) { - std::lock_guard lk(m_open_trigger_decisions_mutex); - m_open_trigger_decisions.erase(token.trigger_number); - TLOG_DEBUG(1) << "Token indicates that trigger decision " << token.trigger_number - << " has been completed. There are now " << m_open_trigger_decisions.size() - << " triggers in flight"; - } else { - // ERS warning: received token for trigger number I don't recognize - } - } - } -} - -} // namespace dunedaq::trigger diff --git a/src/trigger/TokenManager.hpp b/src/trigger/TokenManager.hpp deleted file mode 100644 index 5078bdd3..00000000 --- a/src/trigger/TokenManager.hpp +++ /dev/null @@ -1,102 +0,0 @@ -/** - * @file TokenManager.hpp - * - * This is part of the DUNE DAQ Application Framework, copyright 2020. - * Licensing/copyright details are in the COPYING file that you should have - * received with this code. - */ - -#ifndef TRIGGER_SRC_TRIGGER_TOKENMANAGER_HPP_ -#define TRIGGER_SRC_TRIGGER_TOKENMANAGER_HPP_ - -#include "LivetimeCounter.hpp" - -#include "dfmessages/TriggerDecisionToken.hpp" -#include "dfmessages/Types.hpp" -#include "iomanager/Receiver.hpp" - -#include -#include -#include -#include -#include -#include - -namespace dunedaq { -namespace trigger { - -/** - * @brief TokenManager keeps track of the number of in-flight trigger decisions. - * - * TokenManager implements a credit-based system for trigger - * inhibits. It is constructed with an initial number of tokens and a - * queue of @a dfmessages::TriggerDecisionToken. When a trigger - * decision is sent, the number of tokens is decremented, and when a - * TriggerDecisionToken is received on the queue, the number of tokens - * is incremented. When the count of available tokens reaches zero, no - * further TriggerDecisions may be issued. - */ -class TokenManager -{ -public: - TokenManager(const std::string& connection_name, - int initial_tokens, - daqdataformats::run_number_t run_number, - std::shared_ptr livetime_counter); - - virtual ~TokenManager(); - - TokenManager(TokenManager const&) = delete; - TokenManager(TokenManager&&) = default; - TokenManager& operator=(TokenManager const&) = delete; - TokenManager& operator=(TokenManager&&) = default; - - /** - * Get the number of available tokens - */ - int get_n_tokens() const; - - /** - * Are tokens currently available, to allow sending of new trigger decisions? - */ - bool triggers_allowed() const { return get_n_tokens() > 0; } - - /** - * Notify TokenManager that a trigger decision has been sent. This - * decreases the number of available tokens by one. - * - * Note: you should call this function *before* pushing the corresponding TriggerDecision - * to its output queue. If you do these steps in the other order, the TriggerComplete message - * may be returned before TokenManager is aware of the corresponding trigger decision - */ - void trigger_sent(dfmessages::trigger_number_t); - -private: - // The main thread - void receive_token(dfmessages::TriggerDecisionToken& token); - - std::string m_connection_name; - - // Are we running? - std::atomic m_running_flag; - // How many tokens are currently available? - std::atomic m_n_tokens; - - // The currently-in-flight trigger decisions, and a mutex to guard it - std::set m_open_trigger_decisions; - std::mutex m_open_trigger_decisions_mutex; - - daqdataformats::run_number_t m_run_number; - std::shared_ptr m_livetime_counter; - - // open strigger report time - std::chrono::time_point m_open_trigger_time; - - // the IOManager receiver instance - std::shared_ptr> m_token_receiver; -}; - -} // namespace trigger -} // namespace dunedaq - -#endif // TRIGGER_SRC_TRIGGER_TOKENMANAGER_HPP_ diff --git a/unittest/TokenManager_test.cxx b/unittest/TokenManager_test.cxx deleted file mode 100644 index f85f1376..00000000 --- a/unittest/TokenManager_test.cxx +++ /dev/null @@ -1,93 +0,0 @@ -/** - * @file TokenManager_test.cxx TokenManager class Unit Tests - * - * This is part of the DUNE DAQ Application Framework, copyright 2020. - * Licensing/copyright details are in the COPYING file that you should have - * received with this code. - */ - -#include "trigger/LivetimeCounter.hpp" -#include "trigger/TokenManager.hpp" - -#include "iomanager/IOManager.hpp" -#include "logging/Logging.hpp" - -/** - * @brief Name of this test module - */ -#define BOOST_TEST_MODULE TokenManager_test // NOLINT - -#include "boost/test/unit_test.hpp" - -#include -#include -#include -#include - -using namespace dunedaq; - -BOOST_AUTO_TEST_SUITE(BOOST_TEST_MODULE) - -/** - * @brief Initializes the IOManager - */ -struct IOManagerTestFixture -{ - IOManagerTestFixture() - { - setenv("DUNEDAQ_PARTITION", "TokenManager_t", 0); - - dunedaq::iomanager::Connections_t connections; - dunedaq::iomanager::ConnectionId cid{ "foo", "TriggerDecisionToken" }; - connections.emplace_back( - dunedaq::iomanager::Connection{ cid, "inproc://foo", dunedaq::iomanager::ConnectionType::kSendRecv }); - get_iomanager()->configure({}, connections, false, 0ms); // Not using ConfigClient - } - ~IOManagerTestFixture() - { - get_iomanager()->reset(); - } - - IOManagerTestFixture(IOManagerTestFixture const&) = default; - IOManagerTestFixture(IOManagerTestFixture&&) = default; - IOManagerTestFixture& operator=(IOManagerTestFixture const&) = default; - IOManagerTestFixture& operator=(IOManagerTestFixture&&) = default; -}; - -BOOST_TEST_GLOBAL_FIXTURE(IOManagerTestFixture); - -BOOST_AUTO_TEST_CASE(Basics) -{ - using namespace std::chrono_literals; - - int initial_tokens = 10; - daqdataformats::run_number_t run_number = 1; - auto livetime_counter = std::make_shared(trigger::LivetimeCounter::State::kPaused); - trigger::TokenManager tm("foo", initial_tokens, run_number, livetime_counter); - - BOOST_CHECK_EQUAL(tm.get_n_tokens(), initial_tokens); - BOOST_CHECK_EQUAL(tm.triggers_allowed(), true); - - for (int i = 0; i < initial_tokens - 1; ++i) { - tm.trigger_sent(i); - BOOST_CHECK_EQUAL(tm.get_n_tokens(), initial_tokens - i - 1); - BOOST_CHECK_EQUAL(tm.triggers_allowed(), true); - } - - tm.trigger_sent(initial_tokens); - BOOST_CHECK_EQUAL(tm.get_n_tokens(), 0); - BOOST_CHECK_EQUAL(tm.triggers_allowed(), false); - - // Send a token and check that triggers become allowed again - dfmessages::TriggerDecisionToken token; - token.run_number = run_number; - token.trigger_number = 1; - get_iom_sender("foo")->send(std::move(token), std::chrono::milliseconds(10)); - - // Give TokenManager a little time to pop the token off the queue - std::this_thread::sleep_for(100ms); - BOOST_CHECK_EQUAL(tm.get_n_tokens(), 1); - BOOST_CHECK_EQUAL(tm.triggers_allowed(), true); -} - -BOOST_AUTO_TEST_SUITE_END() From 484af925b8678a2d60b676f73fc7e81ebcf5461b Mon Sep 17 00:00:00 2001 From: Eric Flumerfelt Date: Tue, 4 Aug 2026 14:19:26 -0500 Subject: [PATCH 2/2] Allow sending TriggerDecisions to multiple DFOs. Enable auto-discovery via TriggerInhibit messages --- plugins/MLTModule.cpp | 17 +++++++++++++---- plugins/MLTModule.hpp | 2 +- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/plugins/MLTModule.cpp b/plugins/MLTModule.cpp index 9e2a6151..d3b3eed2 100644 --- a/plugins/MLTModule.cpp +++ b/plugins/MLTModule.cpp @@ -84,7 +84,7 @@ MLTModule::do_configure(const CommandData_t& /*obj*/) // Get the inputs for (auto con : m_mtrg->get_inputs()) { if (con->get_data_type() == datatype_to_string()) { - m_decision_input = get_iom_receiver(con->UID()); + m_decision_input = get_iom_receiver(con->UID()); } else if (con->get_data_type() == datatype_to_string()) { m_inhibit_input = get_iom_receiver(con->UID()); } @@ -93,7 +93,7 @@ MLTModule::do_configure(const CommandData_t& /*obj*/) // Get the outputs for (auto con : m_mtrg->get_outputs()) { if (con->get_data_type() == datatype_to_string()) - m_decision_output = get_iom_sender(con->UID()); + m_decision_outputs[con->UID()] = get_iom_sender(con->UID()); } hdf5libs::HDF5SourceIDHandler::source_id_geo_id_map_t geoidmap = @@ -151,7 +151,7 @@ MLTModule::do_scrap(const CommandData_t& /*obj*/) TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Entering scrap() method"; m_decision_input.reset(); - m_decision_output.reset(); + m_decision_outputs.clear(); m_inhibit_input.reset(); m_srcid_detid_map.clear(); @@ -365,7 +365,11 @@ MLTModule::trigger_decisions_callback(dfmessages::TriggerDecision& decision) } try { - m_decision_output->send(std::move(decision), std::chrono::milliseconds(1)); + for (auto& [conn_name, sender] : m_decision_outputs) { + auto decision_copy = dfmessages::TriggerDecision(decision); + sender->send(std::move(decision_copy), std::chrono::milliseconds(1)); + } + // m_decision_output->send(std::move(decision), std::chrono::milliseconds(1)); m_td_sent_count++; for (const auto t : trigger_types) { @@ -416,6 +420,11 @@ MLTModule::dfo_busy_callback(dfmessages::TriggerInhibit& inhibit) m_dfo_is_busy = inhibit.busy; LivetimeCounter::State state = (inhibit.busy) ? LivetimeCounter::State::kDead : LivetimeCounter::State::kLive; m_livetime_counter->set_state(state); + + if (!m_decision_outputs.count(inhibit.decision_destination)) { + m_decision_outputs[inhibit.decision_destination] = + get_iom_sender(inhibit.decision_destination); + } } } diff --git a/plugins/MLTModule.hpp b/plugins/MLTModule.hpp index d4c18206..17d039c5 100644 --- a/plugins/MLTModule.hpp +++ b/plugins/MLTModule.hpp @@ -95,7 +95,7 @@ class MLTModule : public dunedaq::appfwk::DAQModule // Queue sources and sinks std::shared_ptr> m_decision_input; - std::shared_ptr> m_decision_output; + std::map>> m_decision_outputs; std::shared_ptr> m_inhibit_input; /* TD requests