DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TAProcessor.cpp
Go to the documentation of this file.
1
8#include "trigger/TAProcessor.hpp" // NOLINT(build/include)
9
10// #include "appfwk/DAQModuleHelper.hpp"
11#include "iomanager/Sender.hpp"
12#include "logging/Logging.hpp"
13
18
19// #include "detchannelmaps/TPCChannelMap.hpp"
20
21#include "trigger/TAWrapper.hpp"
23
26
29
32
33// THIS SHOULDN'T BE HERE!!!!! But it is necessary.....
35
36namespace dunedaq {
37namespace trigger {
38
39TAProcessor::TAProcessor(std::unique_ptr<datahandlinglibs::FrameErrorRegistry>& error_registry,
40 bool post_processing_enabled)
41 : datahandlinglibs::TaskRawDataProcessorModel<TAWrapper>(error_registry, post_processing_enabled)
42{
43}
44
46
47void
48TAProcessor::start(const appfwk::DAQModule::CommandData_t& args)
49{
50 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Entering start() method";
51
52 // Reset stats
53 m_ta_received_count.store(0);
54 m_tc_made_count.store(0);
55 m_tc_sent_count.store(0);
57
58 m_running_flag.store(true);
59
60 inherited::start(args);
61
62 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Exiting start() method";
63}
64
65void
66TAProcessor::stop(const appfwk::DAQModule::CommandData_t& args)
67{
68 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Entering stop() method";
69
70 inherited::stop(args);
71 m_running_flag.store(false);
73
74 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Exiting stop() method";
75}
76
77void
79{
80 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Entering conf() method";
81
82 for (auto output : conf->get_outputs()) {
83 try {
84 if (output->get_data_type() == "TriggerCandidate") {
85 m_tc_sink = get_iom_sender<triggeralgs::TriggerCandidate>(output->UID());
86 }
87 } catch (const ers::Issue& excpt) {
88 ers::error(datahandlinglibs::ResourceQueueError(ERS_HERE, "tc", "DefaultRequestHandlerModel", excpt));
89 }
90 }
91
94 std::vector<const appmodel::TCAlgorithm*> tc_algorithms;
96 auto proc_conf = dp->cast<appmodel::TADataProcessor>();
97 if (proc_conf != nullptr && m_post_processing_enabled) {
98 tc_algorithms = proc_conf->get_algorithms();
99 }
100
101 for (auto algo : tc_algorithms) {
102 TLOG() << "Selected TC algorithm: " << algo->UID();
103 std::shared_ptr<triggeralgs::TriggerCandidateMaker> maker = make_tc_maker(algo->class_name());
104 nlohmann::json algo_json = algo->to_json(true);
105 maker->configure(algo_json[algo->UID()]);
106 inherited::add_postprocess_task(std::bind(&TAProcessor::find_tc, this, std::placeholders::_1, maker));
107 m_tcms.push_back(maker);
108 }
109 m_latency_monitoring.store(dp->get_latency_monitoring());
111
112 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Exiting conf() method";
113}
114
115void
116TAProcessor::scrap(const appfwk::DAQModule::CommandData_t& args)
117{
118 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Entering scrap() method";
119 m_tcms.clear();
120 m_tc_sink.reset();
121 inherited::scrap(args);
122 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TAProcessor: Exiting scrap() method";
123}
124
125void
127{
129
130 info.set_ta_received_count(m_ta_received_count.load());
131 info.set_tc_made_count(m_tc_made_count.load());
132 info.set_tc_sent_count(m_tc_sent_count.load());
133 info.set_tc_failed_sent_count(m_tc_failed_sent_count.load());
134
135 this->publish(std::move(info));
136
137 if (m_latency_monitoring.load() && m_running_flag.load()) {
138 opmon::TriggerLatency lat_info;
139
142
143 this->publish(std::move(lat_info));
144 }
145}
146
150void
151TAProcessor::find_tc(const TAWrapper* ta, std::shared_ptr<triggeralgs::TriggerCandidateMaker> tca)
152{
153 // time_activity gave 0 :/
154 if (m_latency_monitoring.load())
157 std::vector<triggeralgs::TriggerCandidate> tcs;
158 tca->operator()(ta->activity, tcs);
159 for (auto tc : tcs) {
161 if (m_latency_monitoring.load())
162 m_latency_instance.update_latency_out(tc.time_candidate);
163 if (!m_tc_sink->try_send(std::move(tc), iomanager::Sender::s_no_block)) {
164 ers::warning(TCDropped(ERS_HERE, tc.time_start, m_sourceid.id));
166 } else {
168 }
169 }
171 return;
172}
173
174void
176{
177 TLOG() << "TAProcessor opmon counters summary:";
178 TLOG() << "------------------------------";
179 TLOG() << "TAs received: \t\t" << m_ta_received_count;
180 TLOG() << "TCs made: \t\t" << m_tc_made_count;
181 TLOG() << "TCs sent: \t\t" << m_tc_sent_count;
182 TLOG() << "TCs failed to send: \t" << m_tc_failed_sent_count;
183 TLOG();
184}
185
186} // namespace fdreadoutlibs
187} // 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.
uint32_t get_source_id() const
Get "source_id" attribute value.
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
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
void update_latency_out(uint64_t latency)
Definition Latency.hpp:46
latency get_latency_in() const
Definition Latency.hpp:49
latency get_latency_out() const
Definition Latency.hpp:52
void update_latency_in(uint64_t latency)
Definition Latency.hpp:43
std::atomic< metric_counter_type > m_ta_received_count
void scrap(const appfwk::DAQModule::CommandData_t &args) override
Unconfigure.
std::atomic< bool > m_running_flag
void generate_opmon_data() override
dunedaq::trigger::Latency m_latency_instance
std::atomic< metric_counter_type > m_tc_failed_sent_count
void find_tc(const TAWrapper *ta, std::shared_ptr< triggeralgs::TriggerCandidateMaker > tcm)
std::vector< std::shared_ptr< triggeralgs::TriggerCandidateMaker > > m_tcms
void stop(const appfwk::DAQModule::CommandData_t &args) override
Stop operation.
daqdataformats::SourceID m_sourceid
TAProcessor(std::unique_ptr< datahandlinglibs::FrameErrorRegistry > &error_registry, bool post_processing_enabled)
std::atomic< metric_counter_type > m_tc_made_count
std::atomic< metric_counter_type > m_tc_sent_count
void start(const appfwk::DAQModule::CommandData_t &args) override
Start operation.
std::atomic< bool > m_latency_monitoring
std::shared_ptr< iomanager::SenderConcept< triggeralgs::TriggerCandidate > > m_tc_sink
void conf(const appmodel::DataHandlerModule *conf) override
Set the emulator mode, if active, timestamps of processed packets are overwritten with new ones.
Base class for any user define issue.
Definition Issue.hpp:76
#define TLVL_ENTER_EXIT_METHODS
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
std::unique_ptr< triggeralgs::TriggerCandidateMaker > make_tc_maker(std::string const &plugin_name)
Load a TriggerCandidateMaker plugin and return a unique_ptr to the contained class.
The DUNE-DAQ namespace.
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101
Subsystem subsystem
The general subsystem of the source of the data.
Definition SourceID.hpp:56
ID_t id
Unique identifier of the source of the data.
Definition SourceID.hpp:59
static const constexpr daqdataformats::SourceID::Subsystem subsystem
Definition TAWrapper.hpp:76
triggeralgs::TriggerActivity activity
Definition TAWrapper.hpp:25