DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TriggerGenericMaker.hpp
Go to the documentation of this file.
8
9#ifndef TRIGGER_SRC_TRIGGER_TRIGGERGENERICMAKER_HPP_
10#define TRIGGER_SRC_TRIGGER_TRIGGERGENERICMAKER_HPP_
11
12#include "trigger/Issues.hpp"
13#include "trigger/Set.hpp"
14#include "trigger/TimeSliceInputBuffer.hpp"
15#include "trigger/TimeSliceOutputBuffer.hpp"
16#include "trigger/triggergenericmakerinfo/InfoNljs.hpp"
17
18#include "appfwk/DAQModule.hpp"
20#include "dfmessages/Types.hpp"
23#include "iomanager/Sender.hpp"
24#include "logging/Logging.hpp"
27
28#include <algorithm>
29#include <memory>
30#include <string>
31#include <vector>
32
33namespace dunedaq::trigger {
34
35// Forward declare the class encapsulating partial specifications of do_work
36template<class IN, class OUT, class MAKER>
38
39// This template class reads IN items from queues, passes them to MAKER objects,
40// and writes the resulting OUT objects to another queue. The behavior of
41// passing IN objects to the MAKER and creating OUT objects from the MAKER is
42// encapsulated by TriggerGenericWorker<IN,OUT,MAKER> templates, defined later
43// in this file
44template<class IN, class OUT, class MAKER>
45class TriggerGenericMaker : public dunedaq::appfwk::DAQModule
46{
47 friend class TriggerGenericWorker<IN, OUT, MAKER>;
48
49public:
50 explicit TriggerGenericMaker(const std::string& name)
51 : DAQModule(name)
52 , m_thread(std::bind(&TriggerGenericMaker::do_work, this, std::placeholders::_1))
54 , m_sent_count(0)
55 , m_run_number(0)
56 , m_input_queue(nullptr)
57 , m_output_queue(nullptr)
58 , m_queue_timeout(100)
59 , m_algorithm_name("[uninitialized]")
60 , m_sourceid(dunedaq::daqdataformats::SourceID::s_invalid_id)
61 , m_buffer_time(0)
62 , m_window_time(625000)
63 , worker(*this) // should be last; may use other members
64 {
65 register_command("start", &TriggerGenericMaker::do_start);
66 register_command("stop", &TriggerGenericMaker::do_stop);
67 register_command("conf", &TriggerGenericMaker::do_configure);
68 register_command("scrap", &TriggerGenericMaker::do_scrap);
69 }
70
72
77
78 // void init(const CommandData_t& obj) override
79 //{
80 // // TODO: Reimplement as OKS
81 // //m_input_queue = get_iom_receiver<IN>(appfwk::connection_uid(obj, "input"));
82 // //m_output_queue = get_iom_sender<OUT>(appfwk::connection_uid(obj, "output"));
83 // }
84
85 void get_info(opmonlib::InfoCollector& ci, int /*level*/) override
86 {
87 triggergenericmakerinfo::Info i;
88
89 i.received_count = m_received_count.load();
90 i.sent_count = m_sent_count.load();
91 if (m_maker) {
92 i.data_vs_system_ms = m_maker->m_data_vs_system_time;
93 } else
94 i.data_vs_system_ms = 0;
95
96 ci.add(i);
97 }
98
99protected:
100 void set_algorithm_name(const std::string& name) { m_algorithm_name = name; }
101
102 // Only applies to makers that output Set<B>
103 void set_sourceid(uint32_t element_id) // NOLINT(build/unsigned)
104 {
105 m_sourceid = element_id;
106 }
107
108 // Only applies to makers that output Set<B>
110 {
111 m_window_time = window_time;
112 m_buffer_time = buffer_time;
113 }
114
115private:
117
118 using metric_counter_type = decltype(triggergenericmakerinfo::Info::received_count);
119 std::atomic<metric_counter_type> m_received_count;
120 std::atomic<metric_counter_type> m_sent_count;
122
124 std::shared_ptr<source_t> m_input_queue;
125
127 std::shared_ptr<sink_t> m_output_queue;
128
129 std::chrono::milliseconds m_queue_timeout;
130
131 std::string m_algorithm_name;
132
133 uint32_t m_sourceid; // NOLINT(build/unsigned)
134
137
138 std::unique_ptr<MAKER> m_maker;
139 nlohmann::json m_maker_conf;
140
142
143 // This should return a unique_ptr to the MAKER created from conf command arguments.
144 // Should also call set_algorithm_name and set_geoid/set_windowing (if desired)
145 virtual std::unique_ptr<MAKER> make_maker(const nlohmann::json& obj) = 0;
146
147 void do_start(const CommandData_t& startobj)
148 {
150 m_sent_count = 0;
153 m_thread.start_working_thread(get_name());
154 m_run_number = startobj.value<dunedaq::daqdataformats::run_number_t>("run", 0);
155 }
156
157 void do_stop(const CommandData_t& /*obj*/) { m_thread.stop_working_thread(); }
158
159 void do_configure(const CommandData_t& obj)
160 {
161 // P. Rodrigues 2022-07-13
162 // We stash the config here and don't actually create the maker
163 // algorithm until start time, so that the algorithm doesn't
164 // persist between runs and hold onto its state from the previous
165 // run
167
168 // worker should be notified that configuration potentially changed
170 }
171
172 void do_scrap(const CommandData_t& obj)
173 {
174 m_input_queue.reset();
175 m_output_queue.reset();
176 m_maker.reset();
177 m_maker_conf.clear();
178 }
179
180 void do_work(std::atomic<bool>& m_running_flag)
181 {
182 // Loop until a stop is received
183 while (m_running_flag.load()) {
184 // While there are items in the input queue, continue draining even if
185 // the running_flag is false, but stop _immediately_ when input is empty
186 IN in;
187 while (receive(in)) {
188 if (m_running_flag.load()) {
189 worker.process(in);
190 }
191 }
192 }
193 // P. Rodrigues 2022-06-01. The argument here is whether to drop
194 // buffered outputs. We choose 'true' because some significant
195 // time can pass between the last input sent by readout and when
196 // we receive a stop. (This happens because stop is sent serially
197 // to readout units before trigger, and each RU takes ~1s to
198 // stop). So by the time we receive a stop command, our buffered
199 // outputs are stale and will cause tardy warnings from the zipper
200 // downstream
201 worker.drain(true);
202 TLOG() << get_name() << ": Exiting do_work() method for run " << m_run_number << ", received " << m_received_count
203 << " inputs (" << worker.get_low_level_input_count() << " sub-inputs) and successfully sent " << m_sent_count
204 << " outputs. ";
205 worker.reset();
206 }
207
208 bool receive(IN& in)
209 {
210 try {
211 in = m_input_queue->receive(m_queue_timeout);
212 } catch (const dunedaq::iomanager::TimeoutExpired& excpt) {
213 // it is perfectly reasonable that there might be no data in the queue
214 // some fraction of the times that we check, so we just continue on and try again
215 return false;
216 }
218 return true;
219 }
220
221 bool send(OUT&& out)
222 {
223 try {
224 m_output_queue->send(std::move(out), m_queue_timeout);
225 } catch (const dunedaq::iomanager::TimeoutExpired& excpt) {
226 ers::warning(excpt);
227 return false;
228 }
229 ++m_sent_count;
230 return true;
231 }
232};
233
234// To handle the different unpacking schemes implied by different templates,
235// do_work is broken out into its own template class that is a friend of
236// TriggerGenericMaker. C++ still does not support partial specification of a
237// single method in a template class, so this approach is the least redundant
238// way to achieve that functionality
239
240// The base template assumes the MAKER has an operator() with the signature
241// operator()(IN, std::vector<OUT>)
242template<class IN, class OUT, class MAKER>
244{
245public:
251
253
254 void reconfigure() {}
255
257
258 void process(IN& in)
259 {
261 std::vector<OUT> out_vec; // one input -> many outputs
262 try {
263 m_parent.m_maker->operator()(in, out_vec);
264 } catch (...) { // NOLINT TODO Benjamin Land <BenLand100@github.com> May 28-2021 can we restrict the possible
265 // exceptions triggeralgs might raise?
266 ers::fatal(AlgorithmFatalError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
267 return;
268 }
269
270 while (out_vec.size()) {
271 if (!m_parent.send(std::move(out_vec.back()))) {
272 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
273 // out_vec.back() is dropped
274 }
275 out_vec.pop_back();
276 }
277 }
278
279 void drain(bool) {}
280
283};
284
285// Partial specialization for IN = Set<A>, OUT = Set<B> and assumes the MAKER has:
286// operator()(A, std::vector<B>)
287template<class A, class B, class MAKER>
288class TriggerGenericWorker<Set<A>, Set<B>, MAKER>
289{
290public: // NOLINT
292 : m_parent(parent)
293 , m_in_buffer(parent.get_name(), parent.m_algorithm_name)
294 , m_out_buffer(parent.get_name(), parent.m_algorithm_name, parent.m_buffer_time)
296 {
297 }
298
300
301 TimeSliceInputBuffer<A> m_in_buffer;
302 TimeSliceOutputBuffer<B> m_out_buffer;
303
305
307 {
308 m_out_buffer.set_window_time(m_parent.m_window_time);
309 m_out_buffer.set_buffer_time(m_parent.m_buffer_time);
310 }
311
312 void reset()
313 {
315 m_out_buffer.reset();
317 }
318
319 void process_slice(const std::vector<A>& time_slice, std::vector<B>& out_vec)
320 {
321 // time_slice is a full slice (all Set<A> combined), time ordered, vector of A
322 // call operator for each of the objects in the vector
323 for (const A& x : time_slice) {
324 try {
325 m_parent.m_maker->operator()(x, out_vec);
326 } catch (...) { // NOLINT TODO Benjamin Land <BenLand100@github.com> May 28-2021 can we restrict the possible
327 // exceptions triggeralgs might raise?
328 ers::fatal(AlgorithmFatalError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
329 return;
330 }
331 }
332 }
333
334 void process(Set<A>& in)
335 {
336 std::vector<B> elems; // Bs to buffer for the next window
337 switch (in.type) {
340 ers::warning(OutOfOrderSets(ERS_HERE, m_parent.get_name(), m_prev_start_time, in.start_time));
341 }
343 std::vector<A> time_slice;
345 if (!m_in_buffer.buffer(in, time_slice, start_time, end_time)) {
346 return; // no complete time slice yet (`in` was part of buffered slice)
347 }
348 m_low_level_input_count += time_slice.size();
349 process_slice(time_slice, elems);
350 } break;
352 // PAR 2022-04-27 We've got a heartbeat for time T, so we know
353 // we won't receive any more inputs for times t < T. Therefore
354 // we can flush all items in the input buffer, which have
355 // times t < T, because the input is time-ordered. We put the
356 // heartbeat in the output buffer, which will handle it
357 // appropriately
358
359 std::vector<A> time_slice;
361 if (m_in_buffer.flush(time_slice, start_time, end_time)) {
362 if (end_time > in.start_time) {
363 // This should never happen, but we check here so we at least get some output if it did
364 ers::fatal(OutOfOrderSets(ERS_HERE, m_parent.get_name(), end_time, in.start_time));
365 }
366 m_low_level_input_count += time_slice.size();
367 process_slice(time_slice, elems);
368 }
369
370 Set<B> heartbeat;
371 heartbeat.type = Set<B>::Type::kHeartbeat;
372 heartbeat.start_time = in.start_time;
373 heartbeat.end_time = in.end_time;
375
376 TLOG_DEBUG(4) << "Buffering heartbeat with start time " << heartbeat.start_time;
377 m_out_buffer.buffer_heartbeat(heartbeat);
378
379 // flush the maker
380 try {
381 // TODO Benjamin Land <BenLand100@github.com> July-14-2021 flushed events go into the buffer... until a window
382 // is ready?
383 m_parent.m_maker->flush(in.end_time, elems);
384 } catch (...) { // NOLINT TODO Benjamin Land <BenLand100@github.com> May-28-2021 can we restrict the possible
385 // exceptions triggeralgs might raise?
386 ers::fatal(AlgorithmFatalError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
387 return;
388 }
389 } break;
391 ers::error(UnknownSetError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
392 break;
393 }
394
395 // add new elements to output buffer
396 if (elems.size() > 0) {
397 m_out_buffer.buffer(elems);
398 }
399
400 size_t n_output_windows = 0;
401 // emit completed windows
402 while (m_out_buffer.ready()) {
403 ++n_output_windows;
404 Set<B> out;
405 m_out_buffer.flush(out);
406 out.seqno = m_parent.m_sent_count;
408
409 if (out.type == Set<B>::Type::kHeartbeat) {
410 TLOG_DEBUG(4) << "Sending heartbeat with start time " << out.start_time;
411 if (!m_parent.send(std::move(out))) {
412 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
413 // out is dropped
414 }
415 }
416 // Only form and send Set<B> if it has a nonzero number of objects
417 else if (out.type == Set<B>::Type::kPayload && out.objects.size() != 0) {
418 TLOG_DEBUG(4) << "Output set window ready with start time " << out.start_time << " end time " << out.end_time
419 << " and " << out.objects.size() << " members";
420 if (!m_parent.send(std::move(out))) {
421 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
422 // out is dropped
423 }
424 }
425 }
426 TLOG_DEBUG(4) << "process() done. Advanced output buffer by " << n_output_windows << " output windows";
427 }
428
429 void drain(bool drop)
430 {
431 // First, send anything in the input buffer to the algorithm, and add any
432 // results to output buffer
433 std::vector<A> time_slice;
435 if (m_in_buffer.flush(time_slice, start_time, end_time)) {
436 std::vector<B> elems;
437 m_low_level_input_count += time_slice.size();
438 process_slice(time_slice, elems);
439 if (elems.size() > 0) {
440 m_out_buffer.buffer(elems);
441 }
442 }
443 // Second, drain the output buffer onto the queue. These may not be "fully
444 // formed" windows, but at this point we're getting no more data anyway.
445 while (!m_out_buffer.empty()) {
446 Set<B> out;
447 m_out_buffer.flush(out);
448 out.seqno = m_parent.m_sent_count;
450
451 if (out.type == Set<B>::Type::kHeartbeat) {
452 if (!drop) {
453 if (!m_parent.send(std::move(out))) {
454 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
455 // out is dropped
456 }
457 }
458 }
459 // Only form and send Set<B> if it has a nonzero number of objects
460 else if (out.type == Set<B>::Type::kPayload && out.objects.size() != 0) {
461 TLOG_DEBUG(1) << "Output set window ready with start time " << out.start_time << " end time " << out.end_time
462 << " and " << out.objects.size() << " members";
463 if (!drop) {
464 if (!m_parent.send(std::move(out))) {
465 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
466 // out is dropped
467 }
468 }
469 }
470 }
471 }
472
475};
476
477// Partial specialization for IN = Set<A> and assumes the the MAKER has:
478// operator()(A, std::vector<OUT>)
479template<class A, class OUT, class MAKER>
480class TriggerGenericWorker<Set<A>, OUT, MAKER>
481{
482public: // NOLINT
483 explicit TriggerGenericWorker(TriggerGenericMaker<Set<A>, OUT, MAKER>& parent)
484 : m_parent(parent)
485 , m_in_buffer(parent.get_name(), parent.m_algorithm_name)
487 {
488 }
489
491
492 TimeSliceInputBuffer<A> m_in_buffer;
493
494 void reconfigure() {}
495
497
498 void process_slice(const std::vector<A>& time_slice, std::vector<OUT>& out_vec)
499 {
500 // time_slice is a full slice (all Set<A> combined), time ordered, vector of A
501 // call operator for each of the objects in the vector
502 for (const A& x : time_slice) {
503 try {
504 m_parent.m_maker->operator()(x, out_vec);
505 } catch (...) { // NOLINT TODO Benjamin Land <BenLand100@github.com> May 28-2021 can we restrict the possible
506 // exceptions triggeralgs might raise?
507 ers::fatal(AlgorithmFatalError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
508 return;
509 }
510 }
511 }
512
513 void process(Set<A>& in)
514 {
515 std::vector<OUT> out_vec; // either a whole time slice, heartbeat flushed, or empty
516 switch (in.type) {
518 std::vector<A> time_slice;
520 if (!m_in_buffer.buffer(in, time_slice, start_time, end_time)) {
521 return; // no complete time slice yet (`in` was part of buffered slice)
522 }
523 m_low_level_input_count += time_slice.size();
524 process_slice(time_slice, out_vec);
525 } break;
527 // TODO BJL May-28-2021 should anything happen with the heartbeat when OUT is not a Set<T>?
528 //
529 // PAR 2022-01-21 We've got a heartbeat for time T, so we know
530 // we won't receive any more inputs for times t < T. Therefore
531 // we can flush all items in the input buffer, which have
532 // times t < T, because the input is time-ordered
533 try {
534 std::vector<A> time_slice;
536 if (m_in_buffer.flush(time_slice, start_time, end_time)) {
537 if (end_time > in.start_time) {
538 // This should never happen, but we check here so we at least get some output if it did
539 ers::fatal(OutOfOrderSets(ERS_HERE, m_parent.get_name(), end_time, in.start_time));
540 }
541 m_low_level_input_count += time_slice.size();
542 process_slice(time_slice, out_vec);
543 }
544 m_parent.m_maker->flush(in.end_time, out_vec);
545 } catch (...) { // NOLINT TODO Benjamin Land <BenLand100@github.com> May 28-2021 can we restrict the possible
546 // exceptions triggeralgs might raise?
547 ers::fatal(AlgorithmFatalError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
548 return;
549 }
550 break;
552 ers::error(UnknownSetError(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
553 break;
554 }
555
556 while (out_vec.size()) {
557 if (!m_parent.send(std::move(out_vec.back()))) {
558 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
559 // out.back() is dropped
560 }
561 out_vec.pop_back();
562 }
563 }
564
565 void drain(bool drop)
566 {
567 // Send anything in the input buffer to the algorithm, and put any results
568 // on the output queue
569 std::vector<A> time_slice;
571 if (m_in_buffer.flush(time_slice, start_time, end_time)) {
572 std::vector<OUT> out_vec;
573 m_low_level_input_count += time_slice.size();
574 process_slice(time_slice, out_vec);
575 while (out_vec.size()) {
576 if (!drop) {
577 if (!m_parent.send(std::move(out_vec.back()))) {
578 ers::error(AlgorithmFailedToSend(ERS_HERE, m_parent.get_name(), m_parent.m_algorithm_name));
579 // out.back() is dropped
580 }
581 }
582 out_vec.pop_back();
583 }
584 }
585 }
586
589};
590
591} // namespace dunedaq::trigger
592
593#endif // TRIGGER_SRC_TRIGGER_TRIGGERGENERICMAKER_HPP_
#define ERS_HERE
A set of TPs or TAs in a given time window, defined by its start and end times.
Definition Set.hpp:26
timestamp_t start_time
Definition Set.hpp:55
origin_t origin
Definition Set.hpp:48
timestamp_t end_time
Definition Set.hpp:58
void do_configure(const CommandData_t &obj)
dunedaq::iomanager::SenderConcept< OUT > sink_t
void do_work(std::atomic< bool > &m_running_flag)
TriggerGenericMaker(TriggerGenericMaker &&)=delete
void do_scrap(const CommandData_t &obj)
dunedaq::utilities::WorkerThread m_thread
std::atomic< metric_counter_type > m_received_count
virtual std::unique_ptr< MAKER > make_maker(const nlohmann::json &obj)=0
TriggerGenericMaker & operator=(const TriggerGenericMaker &)=delete
void do_start(const CommandData_t &startobj)
void set_windowing(daqdataformats::timestamp_t window_time, daqdataformats::timestamp_t buffer_time)
void get_info(opmonlib::InfoCollector &ci, int) override
decltype(triggergenericmakerinfo::Info::received_count) metric_counter_type
std::atomic< metric_counter_type > m_sent_count
TriggerGenericMaker & operator=(TriggerGenericMaker &&)=delete
void set_algorithm_name(const std::string &name)
dunedaq::iomanager::ReceiverConcept< IN > source_t
TriggerGenericMaker(const TriggerGenericMaker &)=delete
TriggerGenericWorker< IN, OUT, MAKER > worker
void process_slice(const std::vector< A > &time_slice, std::vector< OUT > &out_vec)
TriggerGenericWorker(TriggerGenericMaker< Set< A >, OUT, MAKER > &parent)
TriggerGenericWorker(TriggerGenericMaker< Set< A >, Set< B >, MAKER > &parent)
void process_slice(const std::vector< A > &time_slice, std::vector< B > &out_vec)
TriggerGenericMaker< Set< A >, Set< B >, MAKER > & m_parent
TriggerGenericWorker(TriggerGenericMaker< IN, OUT, MAKER > &parent)
TriggerGenericMaker< IN, OUT, MAKER > & m_parent
WorkerThread contains a thread which runs the do_work() function.
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 TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
uint64_t timestamp_t
Type used to represent DUNE timing system timestamps.
Definition Types.hpp:26
daqdataformats::run_number_t run_number_t
Copy daqdataformats::run_number_t.
Definition Types.hpp:34
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed out
msgpack::object obj
Cannot add TPSet with start_time
void warning(const Issue &issue)
Definition ers.hpp:150
void fatal(const Issue &issue)
Definition ers.hpp:111
void error(const Issue &issue)
Definition ers.hpp:101
SourceID is a generalized representation of the source of a piece of data in the DAQ....
Definition SourceID.hpp:32