ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtClient.hpp
Go to the documentation of this file.
1
12
13#ifndef DDTCLIENT_HPP_
14#define DDTCLIENT_HPP_
15
16#include <Ddtdatatransfericd.hpp>
17#include <atomic>
18#include <condition_variable>
19#include <mal/Cii.hpp>
20#include <mal/rr/qos/ReplyTime.hpp>
21#include <mal/utility/LoadMal.hpp>
22#include <map>
23#include <mutex>
24#include <thread>
25
26#include "ddt/ddtConstants.hpp"
27#include "ddt/ddtLogger.hpp"
28
29namespace mal = ::elt::mal;
30namespace datatransfer = ::elt::ddt::datatransfer;
31
32namespace ddt {
33
57
64class DdtClient {
65 public:
69 explicit DdtClient();
70
76 DdtClient(const DdtClientConfig& config, DdtLogger* ddt_logger);
77
81 virtual ~DdtClient();
82
88 void AddUuid(std::string uuid, std::string dsi);
89
95 void UnregisterSubscriber(const std::string uuid);
96
101 bool CheckIfEmpty();
102
109 bool CheckPublisherExists(const std::string& data_stream_identifier) const;
110
122 int32_t RegisterRemoteSubscriber(const std::string& subscriber_uuid,
123 const std::string& data_stream_identifier,
124 const int32_t latency,
125 const int32_t deadline) const;
126
133 std::string GetPublishingUri(const std::string& data_stream_identifier) const;
134
142 const std::string& data_stream_identifier) const;
143
150 int32_t get_number_of_samples(
151 const std::string& data_stream_identifier) const;
152
158 bool get_compute_checksum(const std::string& data_stream_identifier) const;
159
164 std::string get_broker_uri() const;
165
166 protected:
172 void Init(const DdtClientConfig& config, DdtLogger* ddt_logger);
173
180 std::map<std::string, std::string> subscriber_map;
181
185 std::atomic<bool> connected_to_broker;
186
190 std::atomic<bool> heartbeat_active;
191
195 std::string remote_broker_uri;
196
200 std::string broker_uri;
201
205 int32_t reply_time;
206
211
216
221
226
227 private:
231 std::thread hb_thread;
232
236 std::thread reconnect_thread;
237
242 std::mutex hb_cv_mutex;
243 std::condition_variable hb_cv;
244
249 int32_t InitMalClient();
250
254 void StartHeartbeat();
255
259 void StopHeartbeat();
260
264 void HeartbeatThread();
265
270 void Reregister();
271
275 void Reconnect();
276
280 std::unique_ptr<
281 datatransfer::DataBrokerRegistrationSync,
282 std::default_delete<datatransfer::DataBrokerRegistrationSync> >
283 client;
284
288 elt::mal::rr::ListenerRegistration connection_listener;
289
293 std::mutex subscriber_mutex;
294
298 std::mutex connection_state_mutex;
299
304 const int32_t NUM_RETRIES = 10;
305
309 int32_t max_consecutive_failures;
310
315 int32_t latency;
316
320 int32_t deadline;
321};
322
323} // namespace ddt
324
325#endif /* DDTCLIENT_HPP_ */
int32_t reply_time
Definition ddtClient.hpp:205
std::atomic< bool > connected_to_broker
Definition ddtClient.hpp:185
DdtLogger * logger
Definition ddtClient.hpp:225
int32_t heartbeat_interval
Definition ddtClient.hpp:210
std::map< std::string, std::string > subscriber_map
Definition ddtClient.hpp:180
std::atomic< bool > heartbeat_active
Definition ddtClient.hpp:190
std::string remote_broker_uri
Definition ddtClient.hpp:195
bool get_compute_checksum(const std::string &data_stream_identifier) const
Definition ddtClient.cpp:371
bool CheckPublisherExists(const std::string &data_stream_identifier) const
Definition ddtClient.cpp:327
std::string GetPublishingUri(const std::string &data_stream_identifier) const
Definition ddtClient.cpp:347
int32_t heartbeat_timeout
Definition ddtClient.hpp:215
virtual ~DdtClient()
Definition ddtClient.cpp:26
void UnregisterSubscriber(const std::string uuid)
Definition ddtClient.cpp:172
std::string get_broker_uri() const
Definition ddtClient.cpp:379
void AddUuid(std::string uuid, std::string dsi)
Definition ddtClient.cpp:107
int32_t get_max_data_sample_size(const std::string &data_stream_identifier) const
Definition ddtClient.cpp:355
int32_t RegisterRemoteSubscriber(const std::string &subscriber_uuid, const std::string &data_stream_identifier, const int32_t latency, const int32_t deadline) const
Definition ddtClient.cpp:335
std::string broker_uri
Definition ddtClient.hpp:200
int32_t wait_for_connection
Definition ddtClient.hpp:220
void Init(const DdtClientConfig &config, DdtLogger *ddt_logger)
Definition ddtClient.cpp:33
int32_t get_number_of_samples(const std::string &data_stream_identifier) const
Definition ddtClient.cpp:363
bool CheckIfEmpty()
Definition ddtClient.cpp:322
Definition ddtLogger.hpp:43
Contains common used constants. This file shall contain constants that can be used by all application...
Class to wrap the usage of log4cplus as logging utility. This file provides a wrapper class for the u...
Definition ddtClient.hpp:32
Definition ddtClient.hpp:37
int32_t wait_for_connection
Definition ddtClient.hpp:47
int32_t heartbeat_timeout
Definition ddtClient.hpp:45
int32_t latency
Definition ddtClient.hpp:49
std::string broker_uri
Definition ddtClient.hpp:55
int32_t reply_time
Definition ddtClient.hpp:41
std::string remote_broker_uri
Definition ddtClient.hpp:39
int32_t deadline
Definition ddtClient.hpp:51
int32_t heartbeat_interval
Definition ddtClient.hpp:43
int32_t max_consecutive_failures
Definition ddtClient.hpp:53