DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
DataSubscriberModel.hpp
Go to the documentation of this file.
1
8#ifndef DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
9#define DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
10
13
16#include "iomanager/Sender.hpp"
17#include "logging/Logging.hpp"
18// #include "utilities/ReusableThread.hpp"
19
21
22// #include <folly/ProducerConsumerQueue.h>
23// #include <nlohmann/json.hpp>
24
25// #include <atomic>
26// #include <memory>
27// #include <mutex>
28#include <functional>
29
31
32template<class PayloadType>
34{
35public:
37
47
48 void init(const confmodel::DaqModule* cfg) override
49 {
50 if (cfg->get_outputs().size() != 1) {
51 throw datahandlinglibs::InitializationError(ERS_HERE, "Only 1 output supported for subscribers");
52 }
53 m_data_sender = get_iom_sender<PayloadType>(cfg->get_outputs()[0]->UID());
54
55 if (cfg->get_inputs().size() != 1) {
56 throw datahandlinglibs::InitializationError(ERS_HERE, "Only 1 input supported for subscribers");
57 }
58 m_data_receiver = get_iom_receiver<PayloadType>(cfg->get_inputs()[0]->UID());
59 }
60
61 void start()
62 {
63 m_packets = 0;
64 m_sum_packets = 0;
66 m_data_receiver->add_callback(std::bind(&DataSubscriberModel::handle_payload, this, std::placeholders::_1));
67 }
68
69 void stop() { m_data_receiver->remove_callback(); }
70
71 bool handle_payload(PayloadType& message) // NOLINT(build/unsigned)
72 {
73 ++m_packets;
75 if (!m_data_sender->try_send(std::move(message), iomanager::Sender::s_no_block)) {
77 }
78 return true;
79 }
80
81protected:
82 virtual void generate_opmon_data() override
83 {
85 info.set_num_packets(m_packets.exchange(0));
86 info.set_sum_packets(m_sum_packets);
87 info.set_num_dropped_packets(m_dropped_packets.exchange(0));
88
89 this->publish(std::move(info));
90 }
91
92private:
94 std::shared_ptr<source_t> m_data_receiver;
95
97 std::shared_ptr<sink_t> m_data_sender;
98
99 std::atomic<uint64_t> m_packets{ 0 };
100 std::atomic<uint64_t> m_sum_packets{ 0 };
101 std::atomic<uint64_t> m_dropped_packets{ 0 };
102};
103
104} // namespace dunedaq::datahandlinglibs
105
106#endif // DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
#define ERS_HERE
const std::vector< const dunedaq::confmodel::Connection * > & get_inputs() const
Get "inputs" relationship value. List of connections to/from this module.
const std::vector< const dunedaq::confmodel::Connection * > & get_outputs() const
Get "outputs" relationship value. Output connections from this module.
void init(const confmodel::DaqModule *cfg) override
dunedaq::iomanager::ReceiverConcept< PayloadType > source_t
dunedaq::iomanager::SenderConcept< PayloadType > sink_t
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