DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
DaphneV3Interface.cpp
Go to the documentation of this file.
1
11
12#include "DaphneV3Interface.hpp"
13
14#include "logging/Logging.hpp"
15#include <fmt/format.h>
16
17#include <regex>
18#include <string>
19#include <utility>
20
21using namespace dunedaq::daphnemodules;
22using namespace daphne;
23
24DaphneV3Interface::DaphneV3Interface(std::string address, std::string routing, std::chrono::milliseconds timeout)
25 : m_context(1)
26 , m_socket(m_context, zmq::socket_type::dealer)
27 , m_timeout(timeout)
28{
29
30 m_socket.set(zmq::sockopt::routing_id, routing);
31 auto value = (int)m_timeout.count(); // NOLINT
32 TLOG() << routing << " timeout set to " << value << " ms";
33 m_socket.set(zmq::sockopt::rcvtimeo, value);
34 m_socket.set(zmq::sockopt::sndtimeo, value);
35 m_socket.set(zmq::sockopt::immediate, 1); // Don't queue messages to incomplete connections
36
37 // find out if the address has a port with a regex
38 static const std::regex ip_with_port(R"(^([^\/\s:]+)(?::(\d{1,5}))?$)");
39 std::smatch string_values;
40 if (!std::regex_match(address, string_values, ip_with_port)) {
42 }
43
44 m_connection = string_values[2].matched ? fmt::format("tcp://{}", address)
45 : fmt::format("tcp://{}:{}", address, s_default_control_port);
46 TLOG() << "Connecting to: " << m_connection;
47
48 m_socket.connect(m_connection);
49
50 auto add = string_values[1];
51 auto port = string_values[2].matched ? std::stoi(string_values[2]) : s_default_control_port;
52
53 try {
54 if (!validate_connection()) {
55 auto add = string_values[1];
56 auto port = string_values[2].matched ? std::stoi(string_values[2]) : s_default_control_port;
57 throw FailedPing(ERS_HERE, add, port);
58 }
59 } catch (const ers::Issue& e) {
60 throw FailedPing(ERS_HERE, add, port, e);
61 }
62}
63
64void
66{
67
68 const std::lock_guard<std::mutex> lock(m_access_mutex);
69 m_socket.set(zmq::sockopt::linger, 0);
70 m_socket.close();
71}
72
75{
76
77 const std::lock_guard<std::mutex> lock(m_access_mutex);
78
79 _send(std::move(message), type);
80
81 return _receive();
82}
83
84void
86{
87
89 env.set_version(2);
90 env.set_dir(DIR_REQUEST);
91
92 env.set_type(type);
93 env.set_payload(message);
94
95 // additional information
96 env.set_msg_id(m_message_counter++);
97 env.set_timestamp_ns(
98 std::chrono::duration_cast<std::chrono::nanoseconds>(std::chrono::steady_clock::now().time_since_epoch()).count());
99
100 std::string bytes = env.SerializeAsString();
101
102 if (!m_socket.send(zmq::buffer(bytes), zmq::send_flags::none)) {
104 }
105}
106
109{
110
111 zmq::message_t reply;
112 if (!m_socket.recv(reply, zmq::recv_flags::none)) {
113 // timeout or EAGAIN
114 throw FailedReceive(ERS_HERE, m_connection);
115 }
116 if (reply.size() <= 0) {
117 throw EmptyPayload(ERS_HERE, m_connection);
118 }
119
121 if (!rep.ParseFromArray(reply.data(), static_cast<int>(reply.size()))) {
122 throw FailedDecoding(ERS_HERE, rep.GetTypeName(), reply.to_string());
123 }
124
125 if (rep.version() != 2 || rep.dir() != DIR_RESPONSE) {
126 ers::warning(UnexpectedDirection(ERS_HERE, rep.version(), Direction_Name(rep.dir()), reply.to_string()));
127 }
128
129 return rep;
130}
131
132bool
134{
135 static const uint64_t good_value = 0xdeadbeef; // NOLINT
136
137 TestRegRequest req; // empty
138 auto reply = send<TestRegResponse>(
140
141 return reply.value() == good_value;
142}
#define ERS_HERE
::daphne::Direction dir() const
DaphneV3Interface(std::string address, std::string routing, std::chrono::milliseconds timeout=std::chrono::milliseconds(500))
daphne::ControlEnvelopeV2 send(std::string &&message, daphne::MessageTypeV2)
void _send(std::string &&message, daphne::MessageTypeV2)
Base class for any user define issue.
Definition Issue.hpp:76
#define TLOG(...)
Definition macro.hpp:21
const std::string & MessageTypeV2_Name(T value)
const std::string & Direction_Name(T value)
Invalid address
void warning(const Issue &issue)
Definition ers.hpp:150