DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TriggerInhibitAgent.cpp
Go to the documentation of this file.
1
8
11
12#include "logging/Logging.hpp"
13
14#include <memory>
15#include <string>
16#include <utility>
17
21#define TRACE_NAME "TriggerInhibitAgent" // NOLINT
22enum
23{
26};
27
28namespace dunedaq {
29namespace dfmodules {
30
31TriggerInhibitAgent::TriggerInhibitAgent(const std::string& parent_name,
32 std::shared_ptr<trigdecreceiver_t> our_input,
33 std::shared_ptr<triginhsender_t> our_output)
34 : NamedObject(parent_name + "::TriggerInhibitAgent")
35 , m_thread(std::bind(&TriggerInhibitAgent::do_work, this, std::placeholders::_1))
36 , m_queue_timeout(100)
39 , m_trigger_inhibit_sender(our_output)
42{
43}
44
45void
47{
48 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Entering start_checking() method";
50 TLOG() << get_name() << " successfully started";
51 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Exiting start_checking() method";
52}
53
54void
56{
57 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Entering stop_checking() method";
59 TLOG() << get_name() << " successfully stopped";
60 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Exiting stop_checking() method";
61}
62
63void
64TriggerInhibitAgent::do_work(std::atomic<bool>& running_flag)
65{
66 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Entering do_work() method";
67
68 // configuration (hard-coded, for now; will be input from calling code later)
69 int fake_busy_interval_sec = 0;
70 std::chrono::seconds chrono_fake_busy_interval(fake_busy_interval_sec);
71 int fake_busy_duration_sec = 0;
72 std::chrono::seconds chrono_fake_busy_duration(fake_busy_duration_sec);
73 int min_interval_between_inhibit_messages_msec = 0;
74 std::chrono::milliseconds chrono_min_interval_between_inhibit_messages(min_interval_between_inhibit_messages_msec);
75
76 // initialization
77 enum LocalState
78 {
79 no_update,
80 free_state,
81 busy_state
82 };
83 std::chrono::steady_clock::time_point current_time = std::chrono::steady_clock::now();
84 // std::chrono::steady_clock::time_point start_time_of_latest_fake_busy = current_time - chrono_fake_busy_duration;
85 std::chrono::steady_clock::time_point last_sent_time = current_time;
86 LocalState requested_state = no_update;
87 LocalState current_state = free_state;
88 int32_t received_message_count = 0;
89 int32_t sent_message_count = 0;
90
91 // work loop
92 while (running_flag.load()) {
93
94 // check if a TriggerDecision message has arrived, and save the trigger
95 // number contained within it, if one has arrived
96 try {
98 ++received_message_count;
99 TLOG_DEBUG(TLVL_WORK_STEPS) << get_name() << ": Popped the TriggerDecision for trigger number "
100 << trig_dec.trigger_number << " off the input queue";
102 } catch (const iomanager::TimeoutExpired& excpt) {
103 // it is perfectly reasonable that there will be no data in the queue some
104 // fraction of the times that we check, so we just continue on and try again later
105 }
106
107 // to-do: add some logic to fake inhibits
108
109 // check if A) we are supposed to be checking the trigger_number difference, and
110 // B) if so, whether an Inhibit should be asserted or cleared
111 uint32_t threshold = m_threshold_for_inhibit.load(); // NOLINT
112 if (threshold > 0) {
115 if (temp_trig_num_at_start >= temp_trig_num_at_end &&
116 (temp_trig_num_at_start - temp_trig_num_at_end) >= threshold) {
117 if (current_state == free_state) {
118 requested_state = busy_state;
119 }
120 } else {
121 if (current_state == busy_state) {
122 requested_state = free_state;
123 }
124 }
125 }
126
127 // to-do: add some logic to periodically send a message even if nothing has changed
128
129 // send an Inhibit messages, if needed (either Busy or Free state)
130 if (requested_state != no_update && requested_state != current_state) {
131 if ((std::chrono::steady_clock::now() - last_sent_time) >= chrono_min_interval_between_inhibit_messages) {
132 dfmessages::TriggerInhibit inhibit_message;
133 if (requested_state == busy_state) {
134 inhibit_message.busy = true;
135 } else {
136 inhibit_message.busy = false;
137 }
138
139 TLOG_DEBUG(TLVL_WORK_STEPS) << get_name() << ": Pushing a TriggerInhibit message with busy state set to "
140 << inhibit_message.busy << " onto the output queue";
141 try {
142 m_trigger_inhibit_sender->send(std::move(inhibit_message), m_queue_timeout);
143 ++sent_message_count;
144#if 0
145 // temporary logging
146 std::ostringstream oss_sent;
147 oss_sent << ": Successfully pushed a TriggerInhibit message with busy state set to " << inhibit_message.busy
148 << " onto the output queue";
149 TLOG() << ProgressUpdate(ERS_HERE, get_name(), oss_sent.str());
150#endif
151 // if we successfully pushed the message to the Sink, then we assume that the
152 // receiver will get it, and we update our internal state accordingly
153 current_state = requested_state;
154 requested_state = no_update;
155 last_sent_time = std::chrono::steady_clock::now();
156 } catch (const iomanager::TimeoutExpired& excpt) {
157 // It is not ideal if we fail to send the inhibit message out, but rather than
158 // retrying some unknown number of times, we simply output a TRACE message and
159 // go on. This has the benefit of being responsive with pulling TriggerDecision
160 // messages off the input queue, and maybe our Busy/Free state will have changed
161 // by the time that the receiver is ready to receive more messages.
163 << ": TIMEOUT pushing a TriggerInhibit message onto the output queue";
164 }
165 }
166 }
167 }
168
169 std::ostringstream oss_summ;
170 oss_summ << ": Exiting the do_work() method, received " << received_message_count
171 << " TriggerDecision messages and sent " << sent_message_count
172 << " TriggerInhibit messages of all types (both Busy and Free).";
173 TLOG() << ProgressUpdate(ERS_HERE, get_name(), oss_summ.str());
174 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << get_name() << ": Exiting do_work() method";
175}
176
177} // namespace dfmodules
178} // namespace dunedaq
#define ERS_HERE
std::shared_ptr< triginhsender_t > m_trigger_inhibit_sender
std::shared_ptr< trigdecreceiver_t > m_trigger_decision_receiver
std::atomic< daqdataformats::trigger_number_t > m_trigger_number_at_start_of_processing_chain
dunedaq::utilities::WorkerThread m_thread
std::atomic< daqdataformats::trigger_number_t > m_trigger_number_at_end_of_processing_chain
TriggerInhibitAgent(const std::string &, std::shared_ptr< trigdecreceiver_t >, std::shared_ptr< triginhsender_t >)
TriggerInhibitAgent Constructor.
NamedObject(const std::string &name)
NamedObject Constructor.
const std::string & get_name() const final
Get the name of this NamedObejct.
void stop_working_thread()
Stop the working thread.
void start_working_thread(const std::string &name="noname")
Start the working thread (which executes the do_work() function).
#define TLVL_ENTER_EXIT_METHODS
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
uint64_t trigger_number_t
Definition Types.hpp:18
An ERS Issue for DataStore creation failure.
Definition DataStore.hpp:91
The DUNE-DAQ namespace.
A message containing information about a Trigger from Data Selection (or a TriggerDecisionEmulator).
trigger_number_t trigger_number
The trigger number assigned to this TriggerDecision.
Represents a message indicating whether TriggerDecisions should be inhibited.
bool busy
Whether the system is busy.