DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
Receiver.hpp
Go to the documentation of this file.
1
22
23#ifndef IPM_INCLUDE_IPM_RECEIVER_HPP_
24#define IPM_INCLUDE_IPM_RECEIVER_HPP_
25
26#include "cetlib/BasicPluginFactory.h"
27#include "cetlib/compiler_macros.h"
28#include "ers/Issue.hpp"
29#include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG<<issue wont work.
31
32#include <atomic>
33#include <memory>
34#include <string>
35#include <vector>
36
37namespace dunedaq {
38// Disable coverage collection LCOV_EXCL_START
39ERS_DECLARE_ISSUE(ipm, KnownStateForbidsReceive, "Receiver not in a state to receive data", )
41 UnexpectedNumberOfBytes,
42 connection_name << ": Expected " << bytes1 << " bytes in message but received " << bytes2,
43 ((std::string)connection_name)((int)bytes1)((int)bytes2)) // NOLINT
46 connection_name << ": Unable to receive within timeout period (timeout period was " << timeout
47 << " milliseconds)",
48 ((std::string)connection_name)((int)timeout)) // NOLINT
49// Reenable coverage collection LCOV_EXCL_STOP
50} // namespace dunedaq
51
52#ifndef EXTERN_C_FUNC_DECLARE_START
53// NOLINTNEXTLINE(build/define_used)
54#define EXTERN_C_FUNC_DECLARE_START \
55 extern "C" \
56 {
57#endif
58
63// NOLINTNEXTLINE
64#define DEFINE_DUNE_IPM_RECEIVER(klass) \
65 EXTERN_C_FUNC_DECLARE_START \
66 std::shared_ptr<dunedaq::ipm::Receiver> make() \
67 { \
68 return std::shared_ptr<dunedaq::ipm::Receiver>(new klass()); \
69 } \
70 }
71
72namespace dunedaq::ipm {
73
75{
76
77public:
79 {
80 std::string connection_name{ "" };
81 std::string connection_string{ "" };
82 std::vector<std::string> connection_strings{};
83 };
84 using duration_t = std::chrono::milliseconds;
85 static constexpr duration_t s_block = duration_t::max();
86 static constexpr duration_t s_no_block = duration_t::zero();
87
88 using message_size_t = int;
89 static constexpr message_size_t s_any_size =
90 0; // Since "I want 0 bytes" is pointless, "0" denotes "I don't care about the size"
91
92 Receiver() = default;
93 virtual ~Receiver() = default;
94
95 virtual std::string connect_for_receives(const ConnectionInfo& connection_info) = 0;
96
97 virtual bool can_receive() const noexcept = 0;
98
99 // receive() will perform some universally-desirable checks before calling user-implemented receive_:
100 // -Throws KnownStateForbidsReceive if can_receive() == false
101 // -Throws UnexpectedNumberOfBytes if the "nbytes" argument isn't anysize, and the
102 // received bytes inside the function aren't the same number as nbytes
103
104 struct Response
105 {
106 std::string metadata{ "" };
107 std::vector<char> data{};
108 };
109
110 Response receive(const duration_t& timeout, message_size_t num_bytes = s_any_size, bool no_tmoexcept_mode = false);
111 virtual bool data_pending() = 0;
112
113 virtual void register_callback(std::function<void(Response&)>) = 0;
114 virtual void unregister_callback() = 0;
115
116 Receiver(const Receiver&) = delete;
117 Receiver& operator=(const Receiver&) = delete;
118
119 Receiver(Receiver&&) = delete;
121
122protected:
124 void generate_opmon_data() override;
125
126 virtual Response receive_(const duration_t& timeout, bool no_tmoexcept_mode) = 0;
127
128private:
129 mutable std::atomic<size_t> m_bytes = { 0 };
130 mutable std::atomic<size_t> m_messages = { 0 };
131};
132
133inline std::shared_ptr<Receiver>
134make_ipm_receiver(std::string const& plugin_name)
135{
136 static cet::BasicPluginFactory bpf("duneIPM", "make");
137 return bpf.makePlugin<std::shared_ptr<Receiver>>(plugin_name);
138}
139
140} // namespace dunedaq::ipm
141
142#endif // IPM_INCLUDE_IPM_RECEIVER_HPP_
ConnectionInfo m_connection_info
Definition Receiver.hpp:123
Receiver(const Receiver &)=delete
virtual void unregister_callback()=0
virtual ~Receiver()=default
static constexpr duration_t s_block
Definition Receiver.hpp:85
std::chrono::milliseconds duration_t
Definition Receiver.hpp:84
static constexpr message_size_t s_any_size
Definition Receiver.hpp:89
Receiver(Receiver &&)=delete
virtual bool data_pending()=0
Receiver & operator=(const Receiver &)=delete
std::atomic< size_t > m_bytes
Definition Receiver.hpp:129
static constexpr duration_t s_no_block
Definition Receiver.hpp:86
void generate_opmon_data() override
Definition Receiver.cpp:39
Receiver & operator=(Receiver &&)=delete
virtual Response receive_(const duration_t &timeout, bool no_tmoexcept_mode)=0
virtual bool can_receive() const noexcept=0
std::atomic< size_t > m_messages
Definition Receiver.hpp:130
virtual void register_callback(std::function< void(Response &)>)=0
virtual std::string connect_for_receives(const ConnectionInfo &connection_info)=0
An ERS Error indicating that an exception was thrown from ZMQ while performing an operation.
std::shared_ptr< Receiver > make_ipm_receiver(std::string const &plugin_name)
Definition Receiver.hpp:134
The DUNE-DAQ namespace.
ReceiveTimeoutExpired
Definition Receiver.hpp:45
ERS_DECLARE_ISSUE(cibmodules, CIBCommunicationError, " CIB Hardware Communication Error: "<< descriptor,((std::string) descriptor)) ERS_DECLARE_ISSUE(cibmodules
std::vector< std::string > connection_strings
Definition Receiver.hpp:82