DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
ElinkModel.hpp
Go to the documentation of this file.
1
8#ifndef FLXLIBS_SRC_ELINKMODEL_HPP_
9#define FLXLIBS_SRC_ELINKMODEL_HPP_
10
11#include "ElinkConcept.hpp"
12
14
15#include "packetformat/block_format.hpp"
16
17// #include "appfwk/DAQModuleHelper.hpp"
19#include "iomanager/Sender.hpp"
20#include "logging/Logging.hpp"
22
24
25#include <folly/ProducerConsumerQueue.h>
26#include <nlohmann/json.hpp>
27
28#include <atomic>
29#include <memory>
30#include <mutex>
31#include <string>
32
33namespace dunedaq::flxlibs {
34
35template<class TargetPayloadType>
37{
38public:
41 using data_t = nlohmann::json;
42
48 : ElinkConcept()
49 , m_run_marker{ false }
51 {
52 }
54
55 std::shared_ptr<err_sink_t>& get_error_sink() { return m_error_sink_queue; }
56
57 void init(const size_t block_queue_capacity)
58 {
59 m_block_addr_queue = std::make_unique<folly::ProducerConsumerQueue<uint64_t>>(block_queue_capacity); // NOLINT
60 }
61
62 void conf(size_t block_size, bool is_32b_trailers)
63 {
64 if (m_configured) {
65 TLOG_DEBUG(5) << "ElinkModel is already configured!";
66 } else {
68 // if (inconsistency)
69 // ers::fatal(ElinkConfigurationInconsistency(ERS_HERE, m_num_links));
70
71 m_parser->configure(block_size, is_32b_trailers); // unsigned bsize, bool trailer_is_32bit
72 m_configured = true;
73 }
74 }
75
76 void start()
77 {
78 m_t0 = std::chrono::high_resolution_clock::now();
79 if (!m_run_marker.load()) {
80 set_running(true);
82 TLOG() << "Started ElinkModel of link " << inherited::m_link_id << "...";
83 } else {
84 TLOG_DEBUG(5) << "ElinkModel of link " << inherited::m_link_id << " is already running!";
85 }
86 }
87
88 void stop()
89 {
90 if (m_run_marker.load()) {
91 set_running(false);
93 std::this_thread::sleep_for(std::chrono::milliseconds(10));
94 }
95 TLOG_DEBUG(5) << "Stopped ElinkModel of link " << m_link_id << "!";
96 } else {
97 TLOG_DEBUG(5) << "ElinkModel of link " << m_link_id << " is already stopped!";
98 }
99 }
100
101 void set_running(bool should_run)
102 {
103 bool was_running = m_run_marker.exchange(should_run);
104 TLOG_DEBUG(5) << "Active state was toggled from " << was_running << " to " << should_run;
105 }
106
107 bool queue_in_block_address(uint64_t block_addr) // NOLINT(build/unsigned)
108 {
109 if (m_block_addr_queue->write(block_addr)) { // ok write
110 return true;
111 } else { // failed write
112 return false;
113 }
114 }
115
116 void acquire_callback() override
117 {
119 TLOG_DEBUG(5) << "SourceModel callback is already acquired!";
120 } else {
121 // Getting DataMoveCBRegistry
123 m_sink_callback = dmcbr->get_callback<TargetPayloadType>(inherited::m_sink_conf);
125 }
126 }
127
128 // Callbacks
130 using sink_cb_t = std::shared_ptr<std::function<void(TargetPayloadType&&)>>;
132
133protected:
134 void generate_opmon_data() override
135 {
136
138 auto now = std::chrono::high_resolution_clock::now();
139 auto& stats = m_parser_impl.get_stats();
140
141 double seconds = std::chrono::duration_cast<std::chrono::microseconds>(now - m_t0).count() / 1000000.;
142
143 info.set_num_short_chunks_processed(stats.short_ctr.exchange(0));
144 info.set_num_chunks_processed(stats.chunk_ctr.exchange(0));
145 info.set_num_subchunks_processed(stats.subchunk_ctr.exchange(0));
146 info.set_num_blocks_processed(stats.block_ctr.exchange(0));
147
148 info.set_rate_blocks_processed(info.num_blocks_processed() / seconds / 1000.);
149 info.set_rate_chunks_processed(info.num_chunks_processed() / seconds / 1000.);
150
151 info.set_num_short_chunks_processed_with_error(stats.error_short_ctr.exchange(0));
152 info.set_num_chunks_processed_with_error(stats.error_chunk_ctr.exchange(0));
153 info.set_num_subchunks_processed_with_error(stats.error_subchunk_ctr.exchange(0));
154 info.set_num_blocks_processed_with_error(stats.error_block_ctr.exchange(0));
155 info.set_num_subchunk_crc_errors(stats.subchunk_crc_error_ctr.exchange(0));
156 info.set_num_subchunk_trunc_errors(stats.subchunk_trunc_error_ctr.exchange(0));
157 info.set_num_subchunk_errors(stats.subchunk_error_ctr.exchange(0));
158
159 TLOG_DEBUG(2) << inherited::m_elink_str // Move to TLVL_TAKE_NOTE from readout
160 << " Parser stats ->"
161 << " Blocks: " << info.num_blocks_processed() << " Block rate: " << info.rate_blocks_processed()
162 << " [kHz]"
163 << " Chunks: " << info.num_chunks_processed() << " Chunk rate: " << info.rate_chunks_processed()
164 << " [kHz]"
165 << " Shorts: " << info.num_short_chunks_processed() << " Subchunks:" << info.num_subchunks_processed()
166 << " Error Chunks: " << info.num_chunks_processed_with_error()
167 << " Error Shorts: " << info.num_short_chunks_processed_with_error()
168 << " Error Subchunks: " << info.num_subchunks_processed_with_error()
169 << " Error Block: " << info.num_blocks_processed_with_error();
170
171 m_t0 = now;
172
173 publish(std::move(info),
174 { { "card", std::to_string(m_card_id) },
175 { "logical_unit", std::to_string(m_logical_unit) },
176 { "link", std::to_string(m_link_id) },
177 { "tag", std::to_string(m_link_tag) } });
178 }
179
180private:
181 // Types
182 using UniqueBlockAddrQueue = std::unique_ptr<folly::ProducerConsumerQueue<uint64_t>>; // NOLINT(build/unsigned)
183
184 // Internals
185 std::atomic<bool> m_run_marker;
186 bool m_configured{ false };
187
188 // Sink
189 bool m_sink_is_set{ false };
190 std::shared_ptr<err_sink_t> m_error_sink_queue;
191
192 // blocks to process
194
195 // Processor
196 inline static const std::string m_parser_thread_name = "elinkp";
199 {
200 while (m_run_marker.load()) {
201 uint64_t block_addr; // NOLINT
202 if (m_block_addr_queue->read(block_addr)) { // read success
203 const auto* block = const_cast<felix::packetformat::block*>(
204 felix::packetformat::block_from_bytes(reinterpret_cast<const char*>(block_addr)) // NOLINT
205 );
206 m_parser->process(block);
207 } else { // couldn't read from queue
208 std::this_thread::sleep_for(std::chrono::milliseconds(10));
209 }
210 }
211 }
212};
213
214} // namespace dunedaq::flxlibs
215
216#endif // FLXLIBS_SRC_ELINKMODEL_HPP_
static std::shared_ptr< DataMoveCallbackRegistry > get()
std::unique_ptr< felix::packetformat::BlockParser< DefaultParserImpl > > m_parser
std::chrono::time_point< std::chrono::high_resolution_clock > m_t0
const appmodel::DataMoveCallbackConf * m_sink_conf
utilities::ReusableThread m_parser_thread
UniqueBlockAddrQueue m_block_addr_queue
static const std::string m_parser_thread_name
void conf(size_t block_size, bool is_32b_trailers)
void set_running(bool should_run)
bool queue_in_block_address(uint64_t block_addr)
void acquire_callback() override
ElinkModel()
ElinkModel Constructor.
std::atomic< bool > m_run_marker
iomanager::SenderConcept< felix::packetformat::chunk > err_sink_t
std::shared_ptr< std::function< void(TargetPayloadType &&)> > sink_cb_t
std::shared_ptr< err_sink_t > m_error_sink_queue
void generate_opmon_data() override
std::unique_ptr< folly::ProducerConsumerQueue< uint64_t > > UniqueBlockAddrQueue
void init(const size_t block_queue_capacity)
std::shared_ptr< err_sink_t > & get_error_sink()
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
bool set_work(Function &&f, Args &&... args)
void set_name(const std::string &name, int tid)
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21