DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
NetworkSenderModel.hxx
Go to the documentation of this file.
4
5#include "ipm/Sender.hpp"
6#include "logging/Logging.hpp"
8
9#include <memory>
10#include <string>
11#include <typeinfo>
12#include <utility>
13
14using namespace std::chrono_literals; // NOLINT
15
16namespace dunedaq::iomanager {
17
18template<typename Datatype>
20 : SenderConcept<Datatype>(conn_id)
21{
22 TLOG() << "NetworkSenderModel created with DT! Addr: " << static_cast<void*>(this) << ", uid=" << conn_id.uid
23 << ", data_type=" << conn_id.data_type;
24 get_sender(std::chrono::milliseconds(1000));
25 if (m_network_sender_ptr == nullptr) {
26 TLOG() << "Initial connection attempt failed for uid=" << conn_id.uid << ", data_type=" << conn_id.data_type;
27 }
28}
29
30template<typename Datatype>
37
38template<typename Datatype>
39inline void
40NetworkSenderModel<Datatype>::send(Datatype&& data, Sender::timeout_t timeout) // NOLINT
41{
42 try {
43 write_network<Datatype>(data, timeout);
44 } catch (ipm::SendTimeoutExpired& ex) {
45 throw TimeoutExpired(ERS_HERE, this->id().uid, "send", timeout.count(), ex);
46 }
47}
48
49template<typename Datatype>
50inline bool
52{
53 return try_write_network<Datatype>(data, timeout);
54}
55
56template<typename Datatype>
57inline void
58NetworkSenderModel<Datatype>::send_with_topic(Datatype&& data, Sender::timeout_t timeout, std::string topic) // NOLINT
59{
60 try {
61 write_network_with_topic<Datatype>(data, timeout, topic);
62 } catch (ipm::SendTimeoutExpired& ex) {
63 throw TimeoutExpired(ERS_HERE, this->id().uid, "send", timeout.count(), ex);
64 }
65}
66
67template<typename Datatype>
68inline bool
70{
71 get_sender(timeout);
72 return (m_network_sender_ptr != nullptr);
73}
74
75template<typename Datatype>
76inline void
78{
79 auto start = std::chrono::steady_clock::now();
80 while (m_network_sender_ptr == nullptr &&
81 std::chrono::duration_cast<Sender::timeout_t>(std::chrono::steady_clock::now() - start) <= timeout) {
82 // get network resources
83 try {
85
86 if (NetworkManager::get().is_pubsub_connection(this->id())) {
87 TLOG() << "Setting topic to " << this->id().data_type;
88 m_topic = this->id().data_type;
89 }
90 } catch (ers::Issue const& ex) {
91 m_network_sender_ptr = nullptr;
92 std::this_thread::sleep_for(std::chrono::milliseconds(1));
93 }
94 }
95}
96
97template<typename Datatype>
98template<typename MessageType>
99inline typename std::enable_if<dunedaq::serialization::is_serializable<MessageType>::value, void>::type
101{
102 std::lock_guard<std::mutex> lk(m_send_mutex);
103 get_sender(timeout);
104 if (m_network_sender_ptr == nullptr) {
105 throw TimeoutExpired(
106 ERS_HERE, this->id().uid, "send", timeout.count(), ConnectionInstanceNotFound(ERS_HERE, this->id().uid));
107 }
108
109 auto serialized = dunedaq::serialization::serialize(message);
110 // TLOG() << "Serialized message for network sending: " << serialized.size() << ", topic=" <<
111 // m_topic << ", this="
112 // << (void*)this;
113
114 try {
115 m_network_sender_ptr->send(serialized.data(), serialized.size(), extend_first_timeout(timeout), m_topic);
116 } catch (ipm::SendTimeoutExpired const& ex) {
117 TLOG() << "Timeout detected, removing sender to re-acquire connection";
119 m_network_sender_ptr = nullptr;
120 throw;
121 }
122}
123
124template<typename Datatype>
125template<typename MessageType>
126inline typename std::enable_if<!dunedaq::serialization::is_serializable<MessageType>::value, void>::type
128{
129 throw NetworkMessageNotSerializable(ERS_HERE, typeid(MessageType).name()); // NOLINT(runtime/rtti)
130}
131
132template<typename Datatype>
133template<typename MessageType>
134inline typename std::enable_if<dunedaq::serialization::is_serializable<MessageType>::value, bool>::type
136{
137 std::lock_guard<std::mutex> lk(m_send_mutex);
138 get_sender(timeout);
139 if (m_network_sender_ptr == nullptr) {
140 TLOG_DEBUG(5) << ConnectionInstanceNotFound(ERS_HERE, this->id().uid);
141 return false;
142 }
143
144 auto serialized = dunedaq::serialization::serialize(message);
145 // TLOG() << "Serialized message for network sending: " << serialized.size() << ", topic=" <<
146 // m_topic <<
147 // ", this=" << (void*)this;
148
149 auto res =
150 m_network_sender_ptr->send(serialized.data(), serialized.size(), extend_first_timeout(timeout), m_topic, true);
151 if (!res) {
152 TLOG() << "Timeout detected, removing sender to re-acquire connection";
154 m_network_sender_ptr = nullptr;
155 }
156 return res;
157}
158
159template<typename Datatype>
160template<typename MessageType>
161inline typename std::enable_if<!dunedaq::serialization::is_serializable<MessageType>::value, bool>::type
163{
164 ers::error(NetworkMessageNotSerializable(ERS_HERE, typeid(MessageType).name())); // NOLINT(runtime/rtti)
165 return false;
166}
167
168template<typename Datatype>
169template<typename MessageType>
170inline typename std::enable_if<dunedaq::serialization::is_serializable<MessageType>::value, void>::type
172 Sender::timeout_t const& timeout,
173 std::string topic)
174{
175 std::lock_guard<std::mutex> lk(m_send_mutex);
176 get_sender(timeout);
177 if (m_network_sender_ptr == nullptr) {
178 throw TimeoutExpired(
179 ERS_HERE, this->id().uid, "send", timeout.count(), ConnectionInstanceNotFound(ERS_HERE, this->id().uid));
180 }
181
182 auto serialized = dunedaq::serialization::serialize(message);
183 // TLOG() << "Serialized message for network sending: " << serialized.size() << ", topic=" <<
184 // m_topic << ", this="
185 // << (void*)this;
186
187 try {
188 m_network_sender_ptr->send(serialized.data(), serialized.size(), timeout, topic);
189 } catch (TimeoutExpired const& ex) {
190 m_network_sender_ptr = nullptr;
191 throw;
192 }
193}
194
195template<typename Datatype>
196template<typename MessageType>
197inline typename std::enable_if<!dunedaq::serialization::is_serializable<MessageType>::value, void>::type
199{
200 throw NetworkMessageNotSerializable(ERS_HERE, typeid(MessageType).name()); // NOLINT(runtime/rtti)
201}
202
203template<typename Datatype>
206{
207 if (m_first) {
208 m_first = false;
209 if (timeout > 1000ms) {
210 return timeout;
211 }
212 return 1000ms;
213 }
214
215 return timeout;
216}
217
218} // namespace dunedaq::iomanager
#define ERS_HERE
static NetworkManager & get()
std::shared_ptr< ipm::Sender > get_sender(ConnectionId const &conn_id)
void remove_sender(ConnectionId const &conn_id)
void get_sender(Sender::timeout_t const &timeout)
std::enable_if< serialization::is_serializable< MessageType >::value, void >::type write_network(MessageType &message, Sender::timeout_t const &timeout)
bool try_send(Datatype &&data, Sender::timeout_t timeout) override
void send(Datatype &&data, Sender::timeout_t timeout) override
NetworkSenderModel(ConnectionId const &conn_id)
std::shared_ptr< ipm::Sender > m_network_sender_ptr
bool is_ready_for_sending(Sender::timeout_t timeout) override
std::enable_if< serialization::is_serializable< MessageType >::value, void >::type write_network_with_topic(MessageType &message, Sender::timeout_t const &timeout, std::string topic)
void send_with_topic(Datatype &&data, Sender::timeout_t timeout, std::string topic) override
std::enable_if< serialization::is_serializable< MessageType >::value, bool >::type try_write_network(MessageType &message, Sender::timeout_t const &timeout)
Sender::timeout_t extend_first_timeout(Sender::timeout_t timeout)
SenderConcept(ConnectionId const &conn_id)
Definition Sender.hpp:46
std::chrono::milliseconds timeout_t
Definition Sender.hpp:24
ConnectionId id() const
Definition Sender.hpp:35
Base class for any user define issue.
Definition Issue.hpp:76
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
void error(const Issue &issue)
Definition ers.hpp:101