DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
DataHandlingModel.hpp
Go to the documentation of this file.
1
9#ifndef DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_READOUTMODEL_HPP_
10#define DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_READOUTMODEL_HPP_
11
20
22
25#include "iomanager/Sender.hpp"
26
27#include "logging/Logging.hpp"
28
31
34
38
41
45
48
49#include <folly/coro/Baton.h>
50#include <folly/coro/Task.h>
51#include <folly/futures/Future.h>
52
53#include <algorithm>
54#include <functional>
55#include <memory>
56#include <string>
57#include <utility>
58#include <vector>
59
64
65namespace dunedaq {
66namespace datahandlinglibs {
67
68template<class ReadoutType,
69 class RequestHandlerType,
70 class LatencyBufferType,
71 class RawDataProcessorType,
72 class InputDataType = ReadoutType>
74{
75public:
76 // Using shorter typenames
77 using RDT = ReadoutType;
78 using RHT = RequestHandlerType;
79 using LBT = LatencyBufferType;
80 using RPT = RawDataProcessorType;
81 using IDT = InputDataType;
82
83 // Using timestamp typenames
84 using timestamp_t = std::uint64_t; // NOLINT(build/unsigned)
85 static inline constexpr timestamp_t ns = 1;
86 static inline constexpr timestamp_t us = 1000 * ns;
87 static inline constexpr timestamp_t ms = 1000 * us;
88 static inline constexpr timestamp_t s = 1000 * ms;
89
90 // Explicit constructor with run marker pass-through
91 explicit DataHandlingModel(std::atomic<bool>& run_marker)
93 , m_fake_trigger(false)
98 , m_raw_data_receiver(nullptr)
100 , m_latency_buffer_impl(nullptr)
101 , m_raw_processor_impl(nullptr)
102 {
103 }
104
105 virtual ~DataHandlingModel() = default;
106
107 // Initializes the readoutmodel and its internals
108 void init(const appmodel::DataHandlerModule* modconf);
109
110 // Configures the readoutmodel and its internals
111 void conf(const appfwk::DAQModule::CommandData_t& args);
112
113 // Unconfigures readoutmodel's internals
114 void scrap(const appfwk::DAQModule::CommandData_t& args)
115 {
116 m_request_handler_impl->scrap(args);
117 m_latency_buffer_impl->scrap(args);
118 m_raw_processor_impl->scrap(args);
119 }
120
121 // Starts readoutmodel's internals
122 void start(const appfwk::DAQModule::CommandData_t& args);
123
124 // Stops readoutmodel's internals
125 void stop(const appfwk::DAQModule::CommandData_t& args);
126
127 // Record function: invokes request handler's record implementation
128 void record(const appfwk::DAQModule::CommandData_t& args) override { m_request_handler_impl->record(args); }
129
130 // Opmon get_info call implementation
131 // void get_info(opmonlib::InfoCollector& ci, int level);
132
133 // Consume callback
134 std::function<void(IDT&&)> m_consume_callback;
135
136protected:
138 {
139 public:
140 PostprocessScheduleAlgorithm(LatencyBufferType& latency_buffer_impl,
141 RawDataProcessorType& raw_processor_impl,
142 uint64_t processing_delay_ticks, // NOLINT(build/unsigned)
143 uint64_t post_processing_delay_min_wait, // NOLINT(build/unsigned)
144 uint64_t post_processing_delay_max_wait) // NOLINT(build/unsigned)
145 : m_latency_buffer_impl{ latency_buffer_impl }
146 , m_raw_processor_impl{ raw_processor_impl }
147 , m_processing_delay_ticks{ processing_delay_ticks }
148 , m_post_processing_delay_min_wait{ post_processing_delay_min_wait }
149 , m_post_processing_delay_max_wait{ post_processing_delay_max_wait }
150 , m_first_cycle{ true }
152 , m_last_post_proc_time{ std::chrono::system_clock::now() }
154 , m_max_wait_in_ticks{ post_processing_delay_max_wait * 62500 } // FIXME: hardcoded clock frequency
155 {
156 }
157
158 // High-level interface
159 // Schedule deferred post-processing and notify timeout expiration to the processor
160 int run(bool timeout)
161 {
162 int processed = this->do_run(timeout);
163
164 if (timeout) {
166 m_raw_processor_impl.invoke_postprocess_schedule_timeout_policy(timeout_accumulated);
167 }
168
169 return processed;
170 }
171
172 // Deferral of the post processing, to allow elements being reordered in the LB
173 // Basically, find data older than a certain timestamp and process all data since the last post-processed element up
174 // to that value
175 int do_run(bool timeout)
176 {
177 if (m_latency_buffer_impl.occupancy() == 0) {
178 TLOG_DEBUG(TLVL_WORK_STEPS) << "Nothing to postprocess (empty buffer)";
179 return 0;
180 }
181
182 if (m_first_cycle) {
183 auto head = m_latency_buffer_impl.front();
184 m_processed_up_to.set_timestamp(head->get_timestamp());
185 m_first_cycle = false;
186 TLOG() << "***** First pass post processing *****";
187 }
188
189 // Get the LB boundaries
190 auto tail = m_latency_buffer_impl.back();
191 auto newest_ts = tail->get_timestamp();
192
193 timestamp_t end_win_ts = 0;
194 std::chrono::time_point<std::chrono::system_clock> now{ std::chrono::system_clock::now() };
195
196 if (timeout) {
197 // Return if the last processed timestamp is greater than the newest timestamp
198 // This condition occurs after a timeout
199 if (m_processed_up_to.get_timestamp() >= newest_ts + 1) {
200 TLOG_DEBUG(TLVL_WORK_STEPS) << "Nothing to postprocess (at or past cap)";
201 return 0;
202 }
203
206
207 end_win_ts = newest_ts - m_processing_delay_ticks + timeout_accumulated;
208 end_win_ts = std::min(end_win_ts, newest_ts + 1); // Cap to prevent end_win_ts from becoming unnecessarily large
209 } else {
211
212 if (m_processed_up_to.get_timestamp() >= newest_ts + 1) {
213 TLOG_DEBUG(TLVL_WORK_STEPS) << "Nothing to postprocess (data arrived too late, will be ignored)";
214 return 0;
215 }
216
217 auto milliseconds = std::chrono::duration_cast<std::chrono::milliseconds>(now - m_last_post_proc_time);
218
219 if (milliseconds.count() > m_post_processing_delay_min_wait) {
220 if (newest_ts - m_processed_up_to.get_timestamp() > m_processing_delay_ticks) {
221 end_win_ts = newest_ts - m_processing_delay_ticks;
222 } else {
223 TLOG_DEBUG(TLVL_WORK_STEPS) << "Not ready to postprocess (m_processing_delay_ticks is greater)";
224 return 0;
225 }
226 } else {
227 TLOG_DEBUG(TLVL_WORK_STEPS) << "Not ready to postprocess (too fast)";
228 return 0;
229 }
230 }
231
232 auto start_iter = m_latency_buffer_impl.lower_bound(m_processed_up_to, false);
233 m_processed_up_to.set_timestamp(end_win_ts);
234 auto end_iter = m_latency_buffer_impl.lower_bound(m_processed_up_to, false);
235
236 // This likely happens when RDT uses a composite key
237 // The current algorithm does not support composite keys
238 // Our search item `m_processed_up_to` will have its other keys set to their defaults
239 // E.g., for TriggerPrimitive, channel = INVALID_TP_CHANNEL
240 // Even if an entry with the same ts exists in the buffer, its channel will be a valid (smaller) value,
241 // so `lower_bound` will not be able to find it
242 // We should verify that this is the only scenario in which we end up here
243 if (!start_iter.good()) {
244 TLOG_DEBUG(TLVL_WORK_STEPS) << "Nothing to postprocess (!start_iter.good())";
245 return 0;
246 }
247
248 if (start_iter == end_iter) {
249 TLOG_DEBUG(TLVL_WORK_STEPS) << "Nothing to postprocess (start_iter == end_iter)";
250 return 0;
251 }
252
253 int processed = 0;
254 for (auto it = start_iter; it != end_iter; ++it) {
255 // Just to be completely safe
256 // We should understand why we end up here
257 if (!it.good()) {
258 TLOG_DEBUG(TLVL_WORK_STEPS) << "Invalid iterator in postprocessing loop";
259 break;
260 }
261 m_raw_processor_impl.postprocess_item(&(*it));
262 ++processed;
263 }
264
266
267 return processed;
268 }
269
270 private:
271 LatencyBufferType& m_latency_buffer_impl;
272 RawDataProcessorType& m_raw_processor_impl;
273 const uint64_t m_processing_delay_ticks; // NOLINT(build/unsigned)
274 const uint64_t m_post_processing_delay_min_wait; // NOLINT(build/unsigned)
275 const uint64_t m_post_processing_delay_max_wait; // NOLINT(build/unsigned)
280 std::chrono::time_point<std::chrono::system_clock> m_last_post_proc_time;
281 };
282
283 // Perform processing operations on payload
284 void process_item(RDT&& payload);
285
286 // Transform payload if needed, then perform processing
287 void transform_and_process(IDT&& payload);
288
289 // Raw data consume callback
290 void consume_callback(IDT&& payload);
291
292 // Raw data consumer's work function
294
295 // Timesync thread's work function
297
298 // Postprocess scheduler thread's work function
300
301 // Postprocess schedule coroutine
302 folly::coro::Task<void> postprocess_schedule();
303
304 // Dispatch data request
306
307 // Transform input data type to readout
308 virtual std::vector<RDT> transform_payload(IDT& original) const { return { reinterpret_cast<RDT&>(original) }; }
309
310 // Actions postprocess scheduler takes if no data arrives in a configured time
312 {
313 return; // No-op for this class
314 }
315
316 // Operational monitoring
317 virtual void generate_opmon_data() override;
318
319 // Constructor params
320 std::atomic<bool>& m_run_marker;
321
322 // CONFIGURATION
323 // appfwk::app::ModInit m_queue_config;
329 uint64_t m_processing_delay_ticks; // NOLINT(build/unsigned)
330 uint64_t m_post_processing_delay_min_wait; // NOLINT(build/unsigned)
331 uint64_t m_post_processing_delay_max_wait; // NOLINT(build/unsigned)
332
333 // STATS
335 using num_payload_t = std::remove_const<std::invoke_result<decltype(&metric_t::num_payloads), metric_t>::type>::type;
336 using sum_payload_t = std::remove_const<std::invoke_result<decltype(&metric_t::sum_payloads), metric_t>::type>::type;
337 using num_request_t = std::remove_const<std::invoke_result<decltype(&metric_t::num_requests), metric_t>::type>::type;
338 using sum_request_t = std::remove_const<std::invoke_result<decltype(&metric_t::sum_requests), metric_t>::type>::type;
340 std::remove_const<std::invoke_result<decltype(&metric_t::num_data_input_timeouts), metric_t>::type>::type;
342 std::remove_const<std::invoke_result<decltype(&metric_t::num_lb_insert_failures), metric_t>::type>::type;
343 using num_post_processing_delay_max_waits_t = std::remove_const<
344 std::invoke_result<decltype(&metric_t::num_post_processing_delay_max_waits), metric_t>::type>::type;
345
346 std::atomic<num_payload_t> m_num_payloads{ 0 };
347 std::atomic<sum_payload_t> m_sum_payloads{ 0 };
348 std::atomic<num_request_t> m_num_requests{ 0 };
349 std::atomic<sum_request_t> m_sum_requests{ 0 };
350 std::atomic<rawq_timeout_count_t> m_rawq_timeout_count{ 0 };
351 std::atomic<num_lb_insert_failures_t> m_num_lb_insert_failures{ 0 };
352 std::atomic<num_post_processing_delay_max_waits_t> m_num_post_processing_delay_max_waits{ 0 };
353 std::atomic<int> m_stats_packet_count{ 0 };
354
355 // CONSUMER
357
358 // RAW RECEIVER
359 std::chrono::milliseconds m_raw_receiver_timeout_ms;
360 std::chrono::microseconds m_raw_receiver_sleep_us;
362 std::shared_ptr<raw_receiver_ct> m_raw_data_receiver;
365
366 // REQUEST RECEIVERS
368 std::shared_ptr<request_receiver_ct> m_data_request_receiver;
369
370 // FRAGMENT SENDER
371 // std::chrono::milliseconds m_fragment_sender_timeout_ms;
372 // using fragment_sender_ct = iomanager::SenderConcept<std::pair<std::unique_ptr<daqdataformats::Fragment>,
373 // std::string>>; std::shared_ptr<fragment_sender_ct> m_fragment_sender;
374
375 // TIME-SYNC
377 std::shared_ptr<timesync_sender_ct> m_timesync_sender;
380
381 // POSTPROCESS SCHEDULER
383 folly::coro::Baton m_baton;
384 std::unique_ptr<folly::Timekeeper> m_timekeeper;
385
386 // LATENCY BUFFER
387 std::shared_ptr<LatencyBufferType> m_latency_buffer_impl;
388
389 // RAW PROCESSING
390 std::shared_ptr<RawDataProcessorType> m_raw_processor_impl;
391
392 // REQUEST HANDLER
393 std::shared_ptr<RequestHandlerType> m_request_handler_impl;
395
396 // ERROR REGISTRY
397 std::unique_ptr<FrameErrorRegistry> m_error_registry;
398
399 // RUN START T0
400 std::chrono::time_point<std::chrono::high_resolution_clock> m_t0;
401};
402
403} // namespace datahandlinglibs
404} // namespace dunedaq
405
406// Declarations
408
409#endif // DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_MODELS_READOUTMODEL_HPP_
PostprocessScheduleAlgorithm(LatencyBufferType &latency_buffer_impl, RawDataProcessorType &raw_processor_impl, uint64_t processing_delay_ticks, uint64_t post_processing_delay_min_wait, uint64_t post_processing_delay_max_wait)
std::chrono::time_point< std::chrono::system_clock > m_last_post_proc_time
std::remove_const< std::invoke_result< decltype(&metric_t::num_lb_insert_failures), metric_t >::type >::type num_lb_insert_failures_t
DataHandlingModel(std::atomic< bool > &run_marker)
void record(const appfwk::DAQModule::CommandData_t &args) override
std::remove_const< std::invoke_result< decltype(&metric_t::num_data_input_timeouts), metric_t >::type >::type rawq_timeout_count_t
iomanager::ReceiverConcept< dfmessages::DataRequest > request_receiver_ct
void init(const appmodel::DataHandlerModule *modconf)
Forward calls from the appfwk.
std::remove_const< std::invoke_result< decltype(&metric_t::num_requests), metric_t >::type >::type num_request_t
std::remove_const< std::invoke_result< decltype(&metric_t::num_payloads), metric_t >::type >::type num_payload_t
iomanager::ReceiverConcept< InputDataType > raw_receiver_ct
std::remove_const< std::invoke_result< decltype(&metric_t::num_post_processing_delay_max_waits), metric_t >::type >::type num_post_processing_delay_max_waits_t
void stop(const appfwk::DAQModule::CommandData_t &args)
virtual std::vector< RDT > transform_payload(IDT &original) const
void start(const appfwk::DAQModule::CommandData_t &args)
std::remove_const< std::invoke_result< decltype(&metric_t::sum_payloads), metric_t >::type >::type sum_payload_t
void scrap(const appfwk::DAQModule::CommandData_t &args)
std::remove_const< std::invoke_result< decltype(&metric_t::sum_requests), metric_t >::type >::type sum_request_t
dunedaq::datahandlinglibs::opmon::DataHandlerInfo metric_t
void conf(const appfwk::DAQModule::CommandData_t &args)
void dispatch_requests(dfmessages::DataRequest &data_request)
void run_consume()
Function that will be run in its own thread to read the raw packets from the connection and add them ...
void run_timesync()
Function that will be run in its own thread and sends periodic timesync messages by pushing them to t...
iomanager::SenderConcept< dfmessages::TimeSync > timesync_sender_ct
std::atomic< bool > run_marker
Global atomic for process lifetime.
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
ReadoutType
Which type of readout to use for TriggerDecision and DataRequest.
Definition Types.hpp:57
The DUNE-DAQ namespace.
SourceID is a generalized representation of the source of a piece of data in the DAQ....
Definition SourceID.hpp:32
This message represents a request for data sent to a single component of the DAQ.