DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
TPRequestHandler.cpp
Go to the documentation of this file.
4
5#include "rcif/cmd/Nljs.hpp"
6
7namespace dunedaq {
8namespace trigger {
9
10void
12{
13
14 for (auto output : conf->get_outputs()) {
15 if (output->get_data_type() == "TPSet") {
16 try {
18 } catch (const ers::Issue& excpt) {
19 throw datahandlinglibs::ResourceQueueError(ERS_HERE, "tp queue", "DefaultRequestHandlerModel", excpt);
20 }
21 }
22 }
24}
25
26void
27TPRequestHandler::scrap(const appfwk::DAQModule::CommandData_t& args)
28{
29 m_tpset_sink.reset();
31}
32
33void
34TPRequestHandler::start(const appfwk::DAQModule::CommandData_t& args)
35{
36
37 m_oldest_ts = 0;
38 m_newest_ts = 0;
40 m_end_win_ts = 0;
41 m_first_cycle = true;
42
44 rcif::cmd::StartParams start_params = args.get<rcif::cmd::StartParams>();
45 m_run_number = start_params.run;
46}
47
48void
50{
51
52 if (m_tpset_sink == nullptr)
53 return;
54
56
57 {
58 std::unique_lock<std::mutex> lock(m_cv_mutex);
59 m_cv.wait(lock, [&] { return !m_cleanup_requested; });
61 }
62 m_cv.notify_all();
63 if (m_latency_buffer->occupancy() != 0) {
64 // Prepare response
65 RequestResult rres(ResultCode::kUnknown, dr);
66 std::vector<std::pair<void*, size_t>> frag_pieces;
67
68 // Get the newest TP
69 SkipListAcc acc(inherited2::m_latency_buffer->get_skip_list());
70 auto tail = acc.last();
71 auto head = acc.first();
72 m_newest_ts = (*tail).get_timestamp();
73 m_oldest_ts = (*head).get_timestamp();
74
75 if (m_first_cycle) {
77 m_first_cycle = false;
78 }
82 auto num_tps = frag_pieces.size();
83 trigger::TPSet tpset;
86 tpset.origin = m_sourceid;
87 tpset.start_time = m_start_win_ts; // provisory timestamp, will be filled with first TP
88 tpset.end_time = m_end_win_ts; // provisory timestamp, will be filled with last TP
89 tpset.seqno = m_next_tpset_seqno++; // NOLINT(runtime/increment_decrement)
90 // reserve the space for efficiency
91 if (num_tps > 0) {
92 tpset.objects.reserve(frag_pieces.size());
93 bool first_tp = true;
94 for (auto f : frag_pieces) {
96
97 if (first_tp) {
98 tpset.start_time = tp.time_start;
99 first_tp = false;
100 }
101 tpset.end_time = tp.time_start;
102 tpset.objects.emplace_back(std::move(tp));
103 }
104 }
105 if (!m_tpset_sink->try_send(std::move(tpset), iomanager::Sender::s_no_block)) {
108 }
110
111 // remember what we sent for the next loop
113 }
114 }
115 {
116 std::lock_guard<std::mutex> lock(m_cv_mutex);
118 }
119 m_cv.notify_all();
120 return;
121}
122
123} // namespace fdreadoutlibs
124} // namespace dunedaq
#define ERS_HERE
const std::vector< const dunedaq::confmodel::Connection * > & get_outputs() const
Get "outputs" relationship value. Output connections from this module.
std::vector< std::pair< void *, size_t > > get_fragment_pieces(uint64_t start_win_ts, uint64_t end_win_ts, RequestResult &rres)
typename folly::ConcurrentSkipList< TriggerPrimitiveTypeAdapter >::Accessor SkipListAcc
virtual void start(const appfwk::DAQModule::CommandData_t &args)=0
virtual void scrap(const appfwk::DAQModule::CommandData_t &args)=0
virtual void conf(const appmodel::DataHandlerModule *conf)=0
static std::shared_ptr< IOManager > get()
Definition IOManager.hpp:40
static constexpr timeout_t s_no_block
Definition Sender.hpp:26
std::vector< T > objects
Definition Set.hpp:61
daqdataformats::run_number_t run_number
Definition Set.hpp:45
timestamp_t start_time
Definition Set.hpp:55
origin_t origin
Definition Set.hpp:48
timestamp_t end_time
Definition Set.hpp:58
void scrap(const appfwk::DAQModule::CommandData_t &args) override
std::shared_ptr< iomanager::SenderConcept< dunedaq::trigger::TPSet > > m_tpset_sink
void conf(const appmodel::DataHandlerModule *conf) override
void periodic_data_transmission() override
Periodic data transmission - relevant for trigger in particular.
void start(const appfwk::DAQModule::CommandData_t &args) override
Base class for any user define issue.
Definition Issue.hpp:76
Set< trgdataformats::TriggerPrimitive > TPSet
Definition TPSet.hpp:20
The DUNE-DAQ namespace.
void warning(const Issue &issue)
Definition ers.hpp:150
This message represents a request for data sent to a single component of the DAQ.
A single energy deposition on a TPC or PDS channel.