DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TCProcessor.cpp
Go to the documentation of this file.
1
8#include "trigger/TCProcessor.hpp" // NOLINT(build/include)
9
10#include "iomanager/Sender.hpp"
11#include "logging/Logging.hpp"
12
17#include "trigger/TCWrapper.hpp"
19
22
25
26// THIS SHOULDN'T BE HERE!!!!! But it is necessary.....
28
29namespace dunedaq {
30namespace trigger {
31
32TCProcessor::TCProcessor(std::unique_ptr<datahandlinglibs::FrameErrorRegistry>& error_registry,
33 bool post_processing_enabled)
34 : datahandlinglibs::TaskRawDataProcessorModel<TCWrapper>(error_registry, post_processing_enabled)
35{
36}
37
39
40void
41TCProcessor::start(const appfwk::DAQModule::CommandData_t& args)
42{
43 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Entering start() method";
44
45 m_running_flag.store(true);
47 pthread_setname_np(m_send_trigger_decisions_thread.native_handle(), "mlt-dec"); // TODO: originally mlt-trig-dec
48
49 // Reset stats
50 m_tds_created_count.store(0);
51 m_tds_sent_count.store(0);
52 m_tds_dropped_count.store(0);
54 m_tds_cleared_count.store(0);
55 // per TC
56 m_tc_received_count.store(0);
58 m_tds_sent_tc_count.store(0);
62 m_tc_ignored_count.store(0);
63 inherited::start(args);
64
65 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Exiting start() method";
66}
67
68void
69TCProcessor::stop(const appfwk::DAQModule::CommandData_t& args)
70{
71 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Entering stop() method";
72
73 inherited::stop(args);
74 m_running_flag.store(false);
75
76 // Make sure condition_variable knows we flipped running flag
77 {
78 std::lock_guard<std::mutex> lock(m_td_vector_mutex);
79 m_cv.notify_all();
80 }
81
82 // Wait for the TD-sending thread to stop
84
85 // Drop all TDs in vectors at run stage change. Have to do this
86 // after joining m_send_trigger_decisions_thread so we don't
87 // concurrently access the vectors
89
91
92 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Exiting stop() method";
93}
94
95void
97{
98 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Entering conf() method";
99
100 auto mtrg = cfg->cast<appmodel::TriggerDataHandlerModule>();
101 if (mtrg == nullptr) {
102 throw(InvalidConfiguration(ERS_HERE, "Provided null TriggerDataHandlerModule configuration!"));
103 }
104 for (auto output : mtrg->get_outputs()) {
105 try {
106 if (output->get_data_type() == "TriggerDecision") {
107 m_td_sink = get_iom_sender<dfmessages::TriggerDecision>(output->UID());
108 }
109 } catch (const ers::Issue& excpt) {
110 ers::error(datahandlinglibs::ResourceQueueError(ERS_HERE, "td", "DefaultRequestHandlerModel", excpt));
111 }
112 }
113
114 auto dp = mtrg->get_module_configuration()->get_data_processor();
115 auto proc_conf = dp->cast<appmodel::TCDataProcessor>();
116
117 // Add all Source IDs to mandatoy links for now...
118 for (auto const& link : mtrg->get_mandatory_source_ids()) {
119 m_mandatory_links.push_back(
120 dfmessages::SourceID{ daqdataformats::SourceID::string_to_subsystem(link->get_subsystem()), link->get_sid() });
121 }
122 for (auto const& link : mtrg->get_enabled_source_ids()) {
123 m_mandatory_links.push_back(
124 dfmessages::SourceID{ daqdataformats::SourceID::string_to_subsystem(link->get_subsystem()), link->get_sid() });
125 }
126
127 // TODO: Group links!
128 // m_group_links_data = conf->get_groups_links();
132 TLOG_DEBUG(3) << "Total group links: " << m_total_group_links;
133
134 m_tc_merging = proc_conf->get_merge_overlapping_tcs();
135 m_ignore_tc_pileup = proc_conf->get_ignore_overlapping_tcs();
136 m_buffer_timeout = proc_conf->get_buffer_timeout();
137 m_send_timed_out_tds = (m_ignore_tc_pileup) ? false : proc_conf->get_td_out_of_timeout();
138 m_td_readout_limit = proc_conf->get_td_readout_limit();
139 m_ignored_tc_types = proc_conf->get_ignore_tc();
141
142 // Trigger bitwords
143 std::vector<const appmodel::TriggerBitword*> bitwords = proc_conf->get_trigger_bitwords();
144 m_use_bitwords = !bitwords.empty();
145 if (m_use_bitwords) {
146 set_trigger_bitwords(bitwords);
148 }
149 TLOG_DEBUG(3) << "Use bitwords: " << m_use_bitwords;
150 TLOG_DEBUG(3) << "Allow merging: " << m_tc_merging;
151 TLOG_DEBUG(3) << "Ignore pileup: " << m_ignore_tc_pileup;
152 TLOG_DEBUG(3) << "Buffer timeout: " << m_buffer_timeout;
153 TLOG_DEBUG(3) << "Should send timed out TDs: " << m_send_timed_out_tds;
154 TLOG_DEBUG(3) << "TD readout limit: " << m_td_readout_limit;
155
156 // ROI map
157 m_roi_conf_data = proc_conf->get_roi_group_conf();
159 if (m_use_roi_readout) {
162 }
163 TLOG_DEBUG(3) << "Use ROI readout?: " << m_use_roi_readout;
164
165 // Custom readout map
166 m_readout_window_map_data = proc_conf->get_tc_readout_map();
168 if (m_use_readout_map) {
171 }
172 TLOG_DEBUG(3) << "Use readout map: " << m_use_readout_map;
173
174 // Ignoring TC types
175 TLOG_DEBUG(3) << "Ignoring TC types: " << m_ignoring_tc_types;
177 TLOG_DEBUG(3) << "TC types to ignore: ";
178 for (std::vector<unsigned int>::iterator it = m_ignored_tc_types.begin(); it != m_ignored_tc_types.end();) {
179 TLOG_DEBUG(3) << *it;
180 ++it;
181 }
182 }
183 m_latency_monitoring.store(dp->get_latency_monitoring());
184 inherited::add_postprocess_task(std::bind(&TCProcessor::make_td, this, std::placeholders::_1));
185
186 inherited::conf(mtrg);
187
188 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Exiting conf() method";
189}
190
191void
192TCProcessor::scrap(const appfwk::DAQModule::CommandData_t& args)
193{
194 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Entering scrap() method";
195
196 m_mandatory_links.clear();
197 m_group_links.clear();
198 m_roi_conf.clear();
199 m_roi_conf_data.clear();
200 m_roi_conf_ids.clear();
201 m_roi_conf_probs.clear();
202 m_roi_conf_probs_c.clear();
203 m_pending_tds.clear();
205 m_readout_window_map.clear();
206 m_ignored_tc_types.clear();
207
208 m_td_sink.reset();
209
210 m_group_links_data.clear();
211
212 inherited::scrap(args);
213
214 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "TCProcessor: Exiting scrap() method";
215}
216
217void
219{
221
222 info.set_tds_created_count(m_tds_created_count.load());
223 info.set_tds_sent_count(m_tds_sent_count.load());
224 info.set_tds_dropped_count(m_tds_dropped_count.load());
225 info.set_tds_failed_bitword_count(m_tds_failed_bitword_count.load());
226 info.set_tds_cleared_count(m_tds_cleared_count.load());
227 info.set_tc_received_count(m_tc_received_count.load());
228 info.set_tc_ignored_count(m_tc_ignored_count.load());
229 info.set_tds_created_tc_count(m_tds_created_tc_count.load());
230 info.set_tds_sent_tc_count(m_tds_sent_tc_count.load());
231 info.set_tds_dropped_tc_count(m_tds_dropped_tc_count.load());
232 info.set_tds_failed_bitword_tc_count(m_tds_failed_bitword_tc_count.load());
233 info.set_tds_cleared_tc_count(m_tds_cleared_tc_count.load());
234
235 this->publish(std::move(info));
236
237 if (m_latency_monitoring.load() && m_running_flag.load()) {
238 opmon::TriggerLatency lat_info;
239
242
243 this->publish(std::move(lat_info));
244 }
245}
246
250void
252{
253
254 auto tc = tcw->candidate;
255 if (m_latency_monitoring.load())
258
259 if ((m_use_readout_map) && (m_readout_window_map.count(tc.type))) {
260 TLOG_DEBUG(3) << "Got TC of type " << static_cast<int>(tc.type) << ", timestamp " << tc.time_candidate
261 << ", start/end " << tc.time_start << "/" << tc.time_end << ", readout start/end "
262 << tc.time_candidate - m_readout_window_map[tc.type].first << "/"
263 << tc.time_candidate + m_readout_window_map[tc.type].second;
264 } else {
265 TLOG_DEBUG(3) << "Got TC of type " << static_cast<int>(tc.type) << ", timestamp " << tc.time_candidate
266 << ", start/end " << tc.time_start << "/" << tc.time_end;
267 }
268
269 // Option to ignore TC types (if given by config)
270 if (m_ignoring_tc_types == true && check_trigger_type_ignore(static_cast<unsigned int>(tc.type)) == true) {
271 TLOG_DEBUG(3) << " Ignore TC type: " << static_cast<unsigned int>(tc.type);
273
274 /*FIXME: comment out this block: if a TC is to be ignored it shall just be ignored!
275 if (m_tc_merging) {
276 // Still need to check for overlap with existing TD, if overlaps, include in the TD, but don't extend
277 // readout
278 std::lock_guard<std::mutex> lock(m_td_vector_mutex);
279 add_tc_ignored(*tc);
280 }
281 */
282 } else {
283 std::lock_guard<std::mutex> lock(m_td_vector_mutex);
284 add_tc(tc);
285 m_cv.notify_one();
286 TLOG_DEBUG(10) << "pending tds size: " << m_pending_tds.size();
287 }
288 m_last_processed_daq_ts = tc.time_start;
289 return;
290}
291
294{
296 TLOG_DEBUG(5) << "earliest TC index: " << m_earliest_tc_index;
297
298 if (pending_td.contributing_tcs.size() > 1) {
299 TLOG_DEBUG(5) << "!!! TD created from " << pending_td.contributing_tcs.size() << " TCs !!!";
300 }
301
303 decision.trigger_number = 0; // filled by MLT
304 decision.run_number = 0; // filled by MLT
305 decision.trigger_timestamp = pending_td.contributing_tcs[m_earliest_tc_index].time_candidate;
307
308 TDBitset td_bitword = get_TD_bitword(pending_td);
309 TLOG_DEBUG(5) << "[MLT] TD has bitword: " << td_bitword << " "
310 << static_cast<dfmessages::trigger_type_t>(td_bitword.to_ulong());
311 decision.trigger_type = static_cast<dfmessages::trigger_type_t>(td_bitword.to_ulong()); // m_trigger_type;
312
313 // decision.trigger_type = 1; // m_trigger_type;
314
315 TLOG_DEBUG(3) << ", TC detid: " << pending_td.contributing_tcs[m_earliest_tc_index].detid
316 << ", TC type: " << static_cast<int>(pending_td.contributing_tcs[m_earliest_tc_index].type)
317 << ", TC cont number: " << pending_td.contributing_tcs.size()
318 << ", DECISION trigger type: " << decision.trigger_type
319 << ", DECISION timestamp: " << decision.trigger_timestamp
320 << ", request window begin: " << pending_td.readout_start
321 << ", request window end: " << pending_td.readout_end;
322
323 std::vector<dfmessages::ComponentRequest> requests =
325 add_requests_to_decision(decision, requests);
326
327 if (!m_use_roi_readout) {
328 for (const auto& [key, value] : m_group_links) {
329 std::vector<dfmessages::ComponentRequest> group_requests =
330 create_all_decision_requests(value, pending_td.readout_start, pending_td.readout_end);
331 add_requests_to_decision(decision, group_requests);
332 }
333 } else { // using ROI readout
335 }
336
338 m_tds_created_tc_count += pending_td.contributing_tcs.size();
339
340 return decision;
341}
342
343void
345{
346 // A unique lock that can be locked and unlocked
347 std::unique_lock<std::mutex> lock(m_td_vector_mutex);
348
349 while (m_running_flag) {
350 // TODO: think about better implementation (notify?, something event driven)
351 m_cv.wait_for(lock, std::chrono::microseconds(100));
352 auto ready_tds = get_ready_tds(m_pending_tds);
353 TLOG_DEBUG(10) << "ready tds: " << ready_tds.size() << ", updated pending tds: " << m_pending_tds.size();
354
355 for (std::vector<PendingTD>::iterator it = ready_tds.begin(); it != ready_tds.end();) {
356 call_tc_decision(*it);
357 ++it;
358 }
359 }
360}
361
362void
364{
365
366 if (m_use_bitwords) {
367 // Check trigger bitwords
368 TDBitset td_bitword = get_TD_bitword(pending_td);
369 if (!check_trigger_bitwords(td_bitword)) {
370 // Don't process further if the bitword check failed
373 return;
374 }
375 }
376
377 dfmessages::TriggerDecision decision = create_decision(pending_td);
378 auto tn = decision.trigger_number;
379 auto td_ts = decision.trigger_timestamp;
380
381 if (m_latency_monitoring.load())
382 m_latency_instance.update_latency_out(pending_td.contributing_tcs.front().time_start);
383 if (!m_td_sink->try_send(std::move(decision), iomanager::Sender::s_no_block)) {
384 ers::warning(TDDropped(ERS_HERE, tn, td_ts));
386 m_tds_dropped_tc_count += pending_td.contributing_tcs.size();
387 } else {
389 m_tds_sent_tc_count += pending_td.contributing_tcs.size();
390 }
391}
392
393void
395{
396 bool tc_dealt = false;
397 int64_t tc_wallclock_arrived =
398 std::chrono::duration_cast<std::chrono::milliseconds>(std::chrono::steady_clock::now().time_since_epoch()).count();
399
401
402 for (std::vector<PendingTD>::iterator it = m_pending_tds.begin(); it != m_pending_tds.end();) {
403 // Don't deal with TC here if there's no overlap
404 if (!check_overlap(tc, *it)) {
405 it++;
406 continue;
407 }
408
409 // If overlap and ignoring, we drop the TC and flag it as dealt with.
410 if (m_ignore_tc_pileup) {
412 tc_dealt = true;
413 TLOG_DEBUG(3) << "TC overlapping with a previous TD, dropping!";
414 break;
415 }
416
417 // If we're here, TC merging must be on, in which case we're actually
418 // going to merge the TC into the TD.
419 it->contributing_tcs.push_back(tc);
420 if ((m_use_readout_map) && (m_readout_window_map.count(tc.type))) {
421 TLOG_DEBUG(3) << "TC with start/end times " << tc.time_candidate - m_readout_window_map[tc.type].first << "/"
422 << tc.time_candidate + m_readout_window_map[tc.type].second
423 << " overlaps with pending TD with start/end times " << it->readout_start << "/"
424 << it->readout_end;
425 it->readout_start = ((tc.time_candidate - m_readout_window_map[tc.type].first) >= it->readout_start)
426 ? it->readout_start
427 : (tc.time_candidate - m_readout_window_map[tc.type].first);
428 it->readout_end = ((tc.time_candidate + m_readout_window_map[tc.type].second) >= it->readout_end)
429 ? (tc.time_candidate + m_readout_window_map[tc.type].second)
430 : it->readout_end;
431 } else {
432 TLOG_DEBUG(3) << "TC with start/end times " << tc.time_start << "/" << tc.time_end
433 << " overlaps with pending TD with start/end times " << it->readout_start << "/"
434 << it->readout_end;
435 it->readout_start = (tc.time_start >= it->readout_start) ? it->readout_start : tc.time_start;
436 it->readout_end = (tc.time_end >= it->readout_end) ? tc.time_end : it->readout_end;
437 }
438 it->walltime_expiration = tc_wallclock_arrived + m_buffer_timeout;
439 tc_dealt = true;
440 break;
441 }
442 }
443
444 // Don't do anything else if we've already dealt with the TC
445 if (tc_dealt) {
446 return;
447 }
448
449 // Create a new TD out of the TC
450 PendingTD td_candidate;
451 td_candidate.contributing_tcs.push_back(tc);
452 if ((m_use_readout_map) && (m_readout_window_map.count(tc.type))) {
453 td_candidate.readout_start = tc.time_candidate - m_readout_window_map[tc.type].first;
454 td_candidate.readout_end = tc.time_candidate + m_readout_window_map[tc.type].second;
455 } else {
456 td_candidate.readout_start = tc.time_start;
457 td_candidate.readout_end = tc.time_end;
458 }
459 td_candidate.walltime_expiration = tc_wallclock_arrived + m_buffer_timeout;
460 m_pending_tds.push_back(td_candidate);
461}
462
463void
465{
466 for (std::vector<PendingTD>::iterator it = m_pending_tds.begin(); it != m_pending_tds.end();) {
467 if (check_overlap(tc, *it)) {
468 if ((m_use_readout_map) && (m_readout_window_map.count(tc.type))) {
469 TLOG_DEBUG(3) << "!Ignored! TC with start/end times " << tc.time_candidate - m_readout_window_map[tc.type].first
470 << "/" << tc.time_candidate + m_readout_window_map[tc.type].second
471 << " overlaps with pending TD with start/end times " << it->readout_start << "/"
472 << it->readout_end;
473 } else {
474 TLOG_DEBUG(3) << "!Ignored! TC with start/end times " << tc.time_start << "/" << tc.time_end
475 << " overlaps with pending TD with start/end times " << it->readout_start << "/"
476 << it->readout_end;
477 }
478 it->contributing_tcs.push_back(tc);
479 break;
480 }
481 ++it;
482 }
483 return;
484}
485
486bool
488{
489 if ((m_use_readout_map) && (m_readout_window_map.count(tc.type))) {
490 return !(((tc.time_candidate + m_readout_window_map[tc.type].second) < pending_td.readout_start) ||
491 ((tc.time_candidate - m_readout_window_map[tc.type].first > pending_td.readout_end)));
492 } else {
493 return !((tc.time_end < pending_td.readout_start) || (tc.time_start > pending_td.readout_end));
494 }
495}
496
497std::vector<TCProcessor::PendingTD>
498TCProcessor::get_ready_tds(std::vector<PendingTD>& pending_tds)
499{
500 std::vector<PendingTD> return_tds;
501 for (std::vector<PendingTD>::iterator it = pending_tds.begin(); it != pending_tds.end();) {
502 auto timestamp_now =
503 std::chrono::duration_cast<std::chrono::milliseconds>(std::chrono::steady_clock::now().time_since_epoch())
504 .count();
505 if (timestamp_now >= it->walltime_expiration) {
506 return_tds.push_back(*it);
507 it = pending_tds.erase(it);
508 } else if (check_td_readout_length(*it)) { // Also pass on TDs with (too) long readout window
509 return_tds.push_back(*it);
510 it = pending_tds.erase(it);
511 } else {
512 ++it;
513 }
514 }
515 return return_tds;
516}
517
518int
520{
521 int earliest_tc_index = -1;
522 triggeralgs::timestamp_t earliest_tc_time;
523 for (int i = 0; i < static_cast<int>(pending_td.contributing_tcs.size()); i++) {
524 if (earliest_tc_index == -1) {
525 earliest_tc_time = pending_td.contributing_tcs[i].time_candidate;
526 earliest_tc_index = i;
527 } else {
528 if (pending_td.contributing_tcs[i].time_candidate < earliest_tc_time) {
529 earliest_tc_time = pending_td.contributing_tcs[i].time_candidate;
530 earliest_tc_index = i;
531 }
532 }
533 }
534 return earliest_tc_index;
535}
536
537bool
539{
540 bool td_too_long = false;
541 if (static_cast<int64_t>(pending_td.readout_end - pending_td.readout_start) >= m_td_readout_limit) {
542 td_too_long = true;
543 TLOG_DEBUG(3) << "Too long readout window: " << (pending_td.readout_end - pending_td.readout_start)
544 << ", sending immediate TD!";
545 }
546 return td_too_long;
547}
548
549void
551{
552 std::lock_guard<std::mutex> lock(m_td_vector_mutex);
554 // Use std::accumulate to sum up the sizes of all contributing_tcs vectors
555 size_t tds_cleared_tc_count =
556 std::accumulate(m_pending_tds.begin(), m_pending_tds.end(), 0, [](size_t sum, const PendingTD& ptd) {
557 return sum + ptd.contributing_tcs.size();
558 });
559 m_tds_cleared_tc_count += tds_cleared_tc_count;
560 m_pending_tds.clear();
561}
562bool
564{
565 bool ignore = false;
566 for (std::vector<unsigned int>::iterator it = m_ignored_tc_types.begin(); it != m_ignored_tc_types.end();) {
567 if (tc_type == *it) {
568 ignore = true;
569 break;
570 }
571 ++it;
572 }
573 return ignore;
574}
575
576void
578{
579 TLOG_DEBUG(3) << "Configured trigger words:";
580 for (const auto& bitword : m_trigger_bitwords) {
581 TLOG_DEBUG(3) << bitword;
582 }
583}
584
585bool
587{
588 bool trigger_check = false;
589 for (const auto& bitword : m_trigger_bitwords) {
590 TLOG_DEBUG(15) << "TD word: " << td_bitword << ", bitword: " << bitword;
591 trigger_check = ((td_bitword & bitword) == bitword);
592 TLOG_DEBUG(15) << "&: " << (td_bitword & bitword);
593 TLOG_DEBUG(15) << "trigger?: " << trigger_check;
594 if (trigger_check == true)
595 break;
596 }
597 return trigger_check;
598}
599
600void
601TCProcessor::set_trigger_bitwords(const std::vector<const appmodel::TriggerBitword*>& _bitwords)
602{
603 for (const appmodel::TriggerBitword* bitword : _bitwords) {
604 TDBitset temp_bitword;
605
606 for (const std::string& tctype_str : bitword->get_bitword()) {
608
609 if (tc_type == TCType::kUnknown) {
610 throw(InvalidConfiguration(ERS_HERE, "Provided an unknown/non-existent TC type as a trigger bitword!"));
611 }
612
613 temp_bitword.set(static_cast<uint64_t>(tc_type));
614 }
615
616 m_trigger_bitwords.push_back(temp_bitword);
617 }
618}
619
620void
621TCProcessor::parse_readout_map(const std::vector<const appmodel::TCReadoutMap*>& data)
622{
623 for (auto readout_type : data) {
624 TCType tc_type =
625 static_cast<TCType>(dunedaq::trgdataformats::string_to_trigger_candidate_type(readout_type->get_tc_type_name()));
626
627 // Throw error if unknown TC type
628 if (tc_type == TCType::kUnknown) {
629 throw(InvalidConfiguration(ERS_HERE, "Provided an unknown TC type in the TCReadoutMap for the TCProcessor"));
630 }
631
632 m_readout_window_map[tc_type] = { readout_type->get_time_before(), readout_type->get_time_after() };
633 }
634 return;
635}
636void
637TCProcessor::print_readout_map(std::map<TCType, std::pair<triggeralgs::timestamp_t, triggeralgs::timestamp_t>> map)
638{
639 TLOG_DEBUG(3) << "MLT TD Readout map:";
640 for (auto const& [key, val] : map) {
641 TLOG_DEBUG(3) << "Type: " << static_cast<int>(key) << ", before: " << val.first << ", after: " << val.second;
642 }
643 return;
644}
645
646void
647TCProcessor::parse_group_links(const nlohmann::json& data)
648{
649 for (auto group : data) {
650 const nlohmann::json& temp_links_data = group["links"];
651 std::vector<dfmessages::SourceID> temp_links;
652 for (auto link : temp_links_data) {
653 temp_links.push_back(
654 dfmessages::SourceID{ daqdataformats::SourceID::string_to_subsystem(link["subsystem"]), link["element"] });
655 }
656 m_group_links.insert({ group["group"], temp_links });
657 }
658 return;
659}
660
661void
663{
664 TLOG_DEBUG(3) << "MLT Group Links:";
665 for (auto const& [key, val] : m_group_links) {
666 TLOG_DEBUG(3) << "Group: " << key;
667 for (auto const& link : val) {
668 TLOG_DEBUG(3) << link;
669 }
670 }
671 TLOG_DEBUG(3) << " ";
672 return;
673}
678{
680 request.component = link;
681 request.window_begin = start;
682 request.window_end = end;
683
684 TLOG_DEBUG(10) << "setting request start: " << request.window_begin;
685 TLOG_DEBUG(10) << "setting request end: " << request.window_end;
686
687 return request;
688}
689
690std::vector<dfmessages::ComponentRequest>
691TCProcessor::create_all_decision_requests(std::vector<dfmessages::SourceID> links,
694{
695 std::vector<dfmessages::ComponentRequest> requests;
696 for (auto link : links) {
697 requests.push_back(create_request_for_link(link, start, end));
698 }
699 return requests;
700}
701
702void
704 std::vector<dfmessages::ComponentRequest> requests)
705{
706 for (auto request : requests) {
707 decision.components.push_back(request);
708 }
709}
710
711void
712TCProcessor::parse_roi_conf(const std::vector<const appmodel::ROIGroupConf*>& data)
713{
714 int counter = 0;
715 float run_sum = 0;
716 for (auto group : data) {
717 roi_group temp_roi_group;
718 temp_roi_group.n_links = group->get_number_of_link_groups();
719 temp_roi_group.prob = group->get_probability();
720 temp_roi_group.time_window = group->get_time_window();
721 temp_roi_group.mode = group->get_groups_selection_mode();
722 m_roi_conf.insert({ counter, temp_roi_group });
723 m_roi_conf_ids.push_back(counter);
724 m_roi_conf_probs.push_back(group->get_probability());
725 run_sum += static_cast<float>(group->get_probability());
726 m_roi_conf_probs_c.push_back(run_sum);
727 counter++;
728 }
729 return;
730}
731
732void
733TCProcessor::print_roi_conf(std::map<int, roi_group> roi_conf)
734{
735 TLOG_DEBUG(3) << "ROI CONF";
736 for (const auto& [key, value] : roi_conf) {
737 TLOG_DEBUG(3) << "ID: " << key;
738 TLOG_DEBUG(3) << "n links: " << value.n_links;
739 TLOG_DEBUG(3) << "prob: " << value.prob;
740 TLOG_DEBUG(3) << "time: " << value.time_window;
741 TLOG_DEBUG(3) << "mode: " << value.mode;
742 }
743 TLOG_DEBUG(3) << " ";
744 return;
745}
746
747float
749{
750 float rnd = (double)rand() / RAND_MAX;
751 return rnd * (limit);
752}
753
754int
756{
757 float rnd_num = get_random_num_float(m_roi_conf_probs_c.back());
758 for (int i = 0; i < static_cast<int>(m_roi_conf_probs_c.size()); i++) {
759 if (rnd_num < m_roi_conf_probs_c[i]) {
760 return i;
761 }
762 }
763 return -1;
764}
765
766int
768{
770 int rnd = rand() % range;
771 return rnd;
772}
773void
775{
776 // Get configuration at random (weighted)
777 int group_pick = pick_roi_group_conf();
778 if (group_pick != -1) {
779 roi_group this_group = m_roi_conf[m_roi_conf_ids[group_pick]];
780 std::vector<dfmessages::SourceID> links;
781
782 // If mode is random, pick groups to request at random
783 if (this_group.mode == "kRandom") {
784 TLOG_DEBUG(10) << "RAND";
785 std::set<int> groups;
786 while (static_cast<int>(groups.size()) < this_group.n_links) {
787 groups.insert(get_random_num_int());
788 }
789 for (auto r_id : groups) {
790 links.insert(links.end(), m_group_links[r_id].begin(), m_group_links[r_id].end());
791 }
792 // Otherwise, read sequntially by IDs, starting at 0
793 } else {
794 TLOG_DEBUG(10) << "SEQ";
795 int r_id = 0;
796 while (r_id < this_group.n_links) {
797 links.insert(links.end(), m_group_links[r_id].begin(), m_group_links[r_id].end());
798 r_id++;
799 }
800 }
801
802 TLOG_DEBUG(10) << "TD timestamp: " << decision.trigger_timestamp;
803 TLOG_DEBUG(10) << "group window: " << this_group.time_window;
804
805 // Once the components are prepared, create requests and append them to decision
806 std::vector<dfmessages::ComponentRequest> requests = create_all_decision_requests(
807 links, decision.trigger_timestamp - this_group.time_window, decision.trigger_timestamp + this_group.time_window);
808 add_requests_to_decision(decision, requests);
809 links.clear();
810 }
811 return;
812}
813
816{
817 // get only unique types
818 std::vector<int> tc_types;
819 for (auto tc : ready_td.contributing_tcs) {
820 tc_types.push_back(static_cast<int>(tc.type));
821 }
822 tc_types.erase(std::unique(tc_types.begin(), tc_types.end()), tc_types.end());
823
824 // form TD bitword
825 TDBitset td_bitword;
826 for (auto tc_type : tc_types) {
827 td_bitword.set(tc_type);
828 }
829 return td_bitword;
830}
831
832void
834{
835 TLOG() << "TCProcessor opmon counters summary:";
836 TLOG() << "------------------------------";
837 TLOG() << "TDs created: \t\t\t" << m_tds_created_count << " \t(" << m_tds_created_tc_count << " TCs)";
838 TLOG() << "TDs sent: \t\t\t" << m_tds_sent_count << " \t(" << m_tds_sent_tc_count << " TCs)";
839 TLOG() << "TDs dropped: \t\t\t" << m_tds_dropped_count << " \t(" << m_tds_dropped_tc_count << " TCs)";
840 TLOG() << "TDs failed bitword check: \t" << m_tds_failed_bitword_count << " \t(" << m_tds_failed_bitword_tc_count
841 << " TCs)";
842 TLOG() << "TDs cleared: \t\t\t" << m_tds_cleared_count << " \t(" << m_tds_cleared_tc_count << " TCs)";
843 TLOG() << "------------------------------";
844 TLOG() << "TCs received: \t" << m_tc_received_count;
845 TLOG() << "TCs ignored: \t" << m_tc_ignored_count;
846 TLOG();
847}
848
849} // namespace fdreadoutlibs
850} // 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.
const TARGET * cast() const noexcept
Casts object to different class.
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::vector< const appmodel::TCReadoutMap * > m_readout_window_map_data
std::atomic< metric_counter_type > m_tc_ignored_count
std::atomic< metric_counter_type > m_tds_dropped_count
bool check_trigger_type_ignore(unsigned int tc_type)
void call_tc_decision(const PendingTD &pending_td)
TCProcessor(std::unique_ptr< datahandlinglibs::FrameErrorRegistry > &error_registry, bool post_processing_enabled)
dunedaq::trigger::Latency m_latency_instance
std::vector< unsigned int > m_ignored_tc_types
std::map< int, std::vector< daqdataformats::SourceID > > m_group_links
void print_roi_conf(std::map< int, roi_group > roi_conf)
TDBitset get_TD_bitword(const PendingTD &ready_td) const
void parse_roi_conf(const std::vector< const appmodel::ROIGroupConf * > &data)
std::thread m_send_trigger_decisions_thread
std::atomic< metric_counter_type > m_tds_failed_bitword_tc_count
void add_requests_to_decision(dfmessages::TriggerDecision &decision, std::vector< dfmessages::ComponentRequest > requests)
std::shared_ptr< iomanager::SenderConcept< dfmessages::TriggerDecision > > m_td_sink
dfmessages::TriggerDecision create_decision(const PendingTD &pending_td)
void start(const appfwk::DAQModule::CommandData_t &args) override
Start operation.
std::vector< dfmessages::ComponentRequest > create_all_decision_requests(std::vector< daqdataformats::SourceID > links, triggeralgs::timestamp_t start, triggeralgs::timestamp_t end)
std::atomic< metric_counter_type > m_tds_cleared_count
void conf(const appmodel::DataHandlerModule *conf) override
Set the emulator mode, if active, timestamps of processed packets are overwritten with new ones.
std::atomic< bool > m_tc_merging
bool check_trigger_bitwords(const TDBitset &td_bitword) const
void roi_readout_make_requests(dfmessages::TriggerDecision &decision)
std::map< int, roi_group > m_roi_conf
std::atomic< metric_counter_type > m_tds_sent_tc_count
void scrap(const appfwk::DAQModule::CommandData_t &args) override
Unconfigure.
void print_readout_map(std::map< TCType, std::pair< triggeralgs::timestamp_t, triggeralgs::timestamp_t > > map)
void make_td(const TCWrapper *tc)
void generate_opmon_data() override
std::atomic< metric_counter_type > m_tds_cleared_tc_count
void set_trigger_bitwords(const std::vector< const appmodel::TriggerBitword * > &_bitwords)
bool check_td_readout_length(const PendingTD &)
std::atomic< metric_counter_type > m_tds_dropped_tc_count
std::atomic< metric_counter_type > m_tds_failed_bitword_count
std::vector< int > m_roi_conf_ids
std::vector< PendingTD > get_ready_tds(std::vector< PendingTD > &pending_tds)
std::vector< float > m_roi_conf_probs_c
std::vector< PendingTD > m_pending_tds
std::atomic< metric_counter_type > m_tc_received_count
void stop(const appfwk::DAQModule::CommandData_t &args) override
Stop operation.
std::atomic< metric_counter_type > m_tds_created_tc_count
bool m_ignore_tc_pileup
Ignore TCs that overlap with already made TD.
dfmessages::ComponentRequest create_request_for_link(daqdataformats::SourceID link, triggeralgs::timestamp_t start, triggeralgs::timestamp_t end)
std::vector< const appmodel::ROIGroupConf * > m_roi_conf_data
std::atomic< bool > m_running_flag
void parse_readout_map(const std::vector< const appmodel::TCReadoutMap * > &data)
std::condition_variable m_cv
void add_tc(const triggeralgs::TriggerCandidate tc)
std::atomic< metric_counter_type > m_tds_sent_count
std::vector< float > m_roi_conf_probs
std::atomic< bool > m_latency_monitoring
triggeralgs::TriggerCandidate::Type TCType
bool check_overlap(const triggeralgs::TriggerCandidate &tc, const PendingTD &pending_td)
void parse_group_links(const nlohmann::json &data)
std::atomic< metric_counter_type > m_tds_created_count
float get_random_num_float(float limit)
std::vector< TDBitset > m_trigger_bitwords
void add_tc_ignored(const triggeralgs::TriggerCandidate tc)
std::atomic< bool > m_send_timed_out_tds
std::vector< daqdataformats::SourceID > m_mandatory_links
int get_earliest_tc_index(const PendingTD &pending_td)
std::map< TCType, std::pair< triggeralgs::timestamp_t, triggeralgs::timestamp_t > > m_readout_window_map
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
@ kLocalized
Local readout, send Fragments to dataflow.
Definition Types.hpp:59
daqdataformats::ComponentRequest ComponentRequest
Copy daqdataformats::ComponentRequest.
Definition Types.hpp:33
daqdataformats::trigger_type_t trigger_type_t
Copy daqdataforamts::trigger_type_t.
Definition Types.hpp:51
daqdataformats::SourceID SourceID
Copy daqdataformats::SourceID.
Definition Types.hpp:32
TriggerCandidateData::Type string_to_trigger_candidate_type(const std::string &name)
The DUNE-DAQ namespace.
DAC value out of range
Message.
Definition DACNode.hpp:32
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101
dunedaq::trgdataformats::timestamp_t timestamp_t
Definition Types.hpp:16
timestamp_t window_end
End of the data collection window.
SourceID component
The ID of the Requested Component.
timestamp_t window_begin
Start of the data collection window.
static Subsystem string_to_subsystem(const std::string &typestring)
Definition SourceID.hxx:81
A message containing information about a Trigger from Data Selection (or a TriggerDecisionEmulator).
std::vector< ComponentRequest > components
The DAQ components which should be read out to create the TriggerRecord.
ReadoutType readout_type
The type of readout to use (i.e. where to route data).
run_number_t run_number
The current run number.
trigger_number_t trigger_number
The trigger number assigned to this TriggerDecision.
timestamp_t trigger_timestamp
The DAQ timestamp.
trigger_type_t trigger_type
The type of the trigger.
std::vector< triggeralgs::TriggerCandidate > contributing_tcs
triggeralgs::TriggerCandidate candidate
Definition TCWrapper.hpp:23
Factory couldn t std::string alg_name InvalidConfiguration
Definition Issues.hpp:33