ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtDataTransferLib.hpp
Go to the documentation of this file.
1
9
10#ifndef DDTDATATRANSFERLIB_HPP_
11#define DDTDATATRANSFERLIB_HPP_
12
13#include <Ddtdatatransfericd.hpp>
14#include <boost/property_tree/ini_parser.hpp>
15#include <boost/property_tree/ptree.hpp>
16#include <boost/signals2/signal.hpp>
17#include <condition_variable>
18#include <functional>
19#include <iostream>
20#include <mal/Cii.hpp>
21#include <mal/rr/qos/ReplyTime.hpp>
22#include <mal/utility/LoadMal.hpp>
23#include <mutex>
24#include <thread>
25
26#include "ddt/ddtErrorCodes.hpp"
27#include "ddt/ddtLogger.hpp"
29
30namespace mal = ::elt::mal;
31namespace datatransfer = ::elt::ddt::datatransfer;
32
33namespace ddt {
34
39 public:
53
57 using ConnectionStateSignalT = boost::signals2::signal<void(ConnectionState)>;
58
63 explicit DdtDataTransferLib(DdtLogger* ddt_logger);
64
70 explicit DdtDataTransferLib(log4cplus::Logger const& log4cplus_logger);
71
75 virtual ~DdtDataTransferLib();
76
84 void SetQoS(const int ddt_latency, const int ddt_deadline);
85
93 const std::string VerifyPathInBrokerUri(std::string broker_uri);
94
100 int InitMAL(const std::string broker_uri);
101
106 std::unique_ptr<datatransfer::DataBrokerRegistrationSync,
107 std::default_delete<datatransfer::DataBrokerRegistrationSync>>
109
118 virtual int RegisterPublisher(const std::string uri, const std::string dsi,
119 const bool compute_crc) {
120 return 0;
121 };
122
127 virtual int UnregisterPublisher() { return 0; };
128
132 virtual void PublishData() {
133 // intentionally-blank override
134 }
135
145 virtual int RegisterSubscriber(const std::string uri, const std::string dsi,
146 const std::string remote_uri,
147 const int32_t interval = 10) {
148 return 0;
149 };
150
155 virtual int UnregisterSubscriber() { return 0; };
156
161 virtual DataSample* ReadData() { return nullptr; };
162
163
171 boost::signals2::connection ConnectToConnectionState(
172 const std::function<void(ConnectionState)>& callback);
173
174 protected:
181 void StartHeartbeat(const int32_t interval, const std::string id);
182
186 void StopHeartbeat();
187
191 virtual void StopThreads();
192
196 void Reconnect();
197
201 virtual void Reregister() {}
202
209 void CheckHeartbeatTimeout(int32_t& new_reply_time);
210
217
223 const std::string GetConfigFilePath();
224
228 int latency = 0;
229
233 int deadline = 0;
234
238 int32_t reply_time = 0;
239
244
248 std::thread hb_thread;
249
255 std::thread reconnect_thread;
256
260 std::atomic<bool> reconnect_pending{false};
261
267
272 std::mutex hb_cv_mutex;
273 std::condition_variable hb_cv;
274
279 std::atomic<bool> heartbeat_active{false};
280
284 std::atomic<bool> shutdown_in_progress{false};
285
289 std::unique_ptr<
290 datatransfer::DataBrokerRegistrationSync,
291 std::default_delete<datatransfer::DataBrokerRegistrationSync> >
293
297 std::atomic<bool> connected_to_broker{false};
298
303
308
312 std::string broker_uri;
313
317 elt::mal::rr::ListenerRegistration connection_listener;
318
322 DdtLogger* logger = nullptr;
323
328
332 const int32_t REPLY_TIME_DEFAULT = 6;
333
337 const int32_t REPLY_TIME_MIN = 2;
338
343
347 const int32_t WAIT_FOR_CONNECTION_DEFAULT = 10;
348
352 const int32_t WAIT_FOR_CONNECTION_MIN = 2;
353
358
363
368
369 private:
373 void Init(DdtLogger* ddt_logger);
374
378 const int32_t NUM_RECONNECT_RETRIES = 10;
379
383 void HeartbeatThread();
384
385 std::string identifier;
386
387 const std::string BROKER_PATH{"/broker/Broker1"};
388 const int LATENCY_DEFAULT = 10000;
389 const int DEADLINE_DEFAULT = 10;
390 const int32_t HEARTBEAT_INTERVAL_DEFAULT = 1;
391};
392
393} // namespace ddt
394
395#endif /* DDTDATATRANSFERLIB_HPP_ */
std::thread reconnect_thread
Definition ddtDataTransferLib.hpp:255
const int32_t WAIT_FOR_CONNECTION_MIN
Definition ddtDataTransferLib.hpp:352
const std::string VerifyPathInBrokerUri(std::string broker_uri)
Definition ddtDataTransferLib.cpp:331
DdtLogger * logger
Definition ddtDataTransferLib.hpp:322
virtual DataSample * ReadData()
Definition ddtDataTransferLib.hpp:161
void SetQoS(const int ddt_latency, const int ddt_deadline)
Definition ddtDataTransferLib.cpp:93
DdtDataTransferLib(DdtLogger *ddt_logger)
Definition ddtDataTransferLib.cpp:16
void SpawnReconnectionThread()
Definition ddtDataTransferLib.cpp:61
virtual int UnregisterSubscriber()
Definition ddtDataTransferLib.hpp:155
std::atomic< bool > reconnect_pending
Definition ddtDataTransferLib.hpp:260
void StopHeartbeat()
Definition ddtDataTransferLib.cpp:118
std::condition_variable hb_cv
Definition ddtDataTransferLib.hpp:273
boost::signals2::signal< void(ConnectionState)> ConnectionStateSignalT
Definition ddtDataTransferLib.hpp:57
boost::signals2::connection ConnectToConnectionState(const std::function< void(ConnectionState)> &callback)
Definition ddtDataTransferLib.cpp:56
const std::string GetConfigFilePath()
Definition ddtDataTransferLib.cpp:364
virtual int RegisterSubscriber(const std::string uri, const std::string dsi, const std::string remote_uri, const int32_t interval=10)
Definition ddtDataTransferLib.hpp:145
int32_t wait_for_connection
Definition ddtDataTransferLib.hpp:342
const int32_t REPLY_TIME_DEFAULT
Definition ddtDataTransferLib.hpp:332
std::string broker_uri
Definition ddtDataTransferLib.hpp:312
void StartHeartbeat(const int32_t interval, const std::string id)
Definition ddtDataTransferLib.cpp:98
const int32_t REPLY_TIME_MIN
Definition ddtDataTransferLib.hpp:337
std::unique_ptr< datatransfer::DataBrokerRegistrationSync, std::default_delete< datatransfer::DataBrokerRegistrationSync > > client
Definition ddtDataTransferLib.hpp:292
int deadline
Definition ddtDataTransferLib.hpp:233
const int32_t MAX_CONSECUTIVE_FAILURES_MIN
Definition ddtDataTransferLib.hpp:367
std::atomic< bool > shutdown_in_progress
Definition ddtDataTransferLib.hpp:284
std::unique_ptr< datatransfer::DataBrokerRegistrationSync, std::default_delete< datatransfer::DataBrokerRegistrationSync > > GetBrokerClient()
Definition ddtDataTransferLib.cpp:360
int32_t heartbeat_interval
Definition ddtDataTransferLib.hpp:243
int32_t max_consecutive_failures
Definition ddtDataTransferLib.hpp:357
DdtLogger * my_logger
Definition ddtDataTransferLib.hpp:327
std::mutex hb_cv_mutex
Definition ddtDataTransferLib.hpp:272
std::atomic< bool > heartbeat_active
Definition ddtDataTransferLib.hpp:279
int latency
Definition ddtDataTransferLib.hpp:228
virtual void PublishData()
Definition ddtDataTransferLib.hpp:132
ConnectionState
Definition ddtDataTransferLib.hpp:43
@ Connected
Connection is up.
Definition ddtDataTransferLib.hpp:47
@ PublisherDisconnected
Broker is reachable but the publisher of data stream is not.
Definition ddtDataTransferLib.hpp:51
@ Reconnecting
Reconnection in progress.
Definition ddtDataTransferLib.hpp:49
@ Disconnected
Connection is down.
Definition ddtDataTransferLib.hpp:45
virtual void StopThreads()
Definition ddtDataTransferLib.cpp:35
void CheckHeartbeatTimeout(int32_t &new_reply_time)
Definition ddtDataTransferLib.cpp:235
void NotifyConnectionState(ConnectionState state)
Definition ddtDataTransferLib.cpp:75
std::atomic< bool > connected_to_broker
Definition ddtDataTransferLib.hpp:297
ConnectionStateSignalT connection_state_signal
Definition ddtDataTransferLib.hpp:302
virtual int UnregisterPublisher()
Definition ddtDataTransferLib.hpp:127
elt::mal::rr::ListenerRegistration connection_listener
Definition ddtDataTransferLib.hpp:317
virtual void Reregister()
Definition ddtDataTransferLib.hpp:201
const int32_t MAX_CONSECUTIVE_FAILURES_DEFAULT
Definition ddtDataTransferLib.hpp:362
std::thread hb_thread
Definition ddtDataTransferLib.hpp:248
const int32_t WAIT_FOR_CONNECTION_DEFAULT
Definition ddtDataTransferLib.hpp:347
virtual ~DdtDataTransferLib()
Definition ddtDataTransferLib.cpp:26
std::atomic< ConnectionState > last_notified_state
Definition ddtDataTransferLib.hpp:307
void Reconnect()
Definition ddtDataTransferLib.cpp:188
int InitMAL(const std::string broker_uri)
Definition ddtDataTransferLib.cpp:270
int32_t reply_time
Definition ddtDataTransferLib.hpp:238
virtual int RegisterPublisher(const std::string uri, const std::string dsi, const bool compute_crc)
Definition ddtDataTransferLib.hpp:118
Definition ddtLogger.hpp:43
Contains common used error codes. This file shall contain error codes that can be used by all applica...
Class to wrap the usage of log4cplus as logging utility. This file provides a wrapper class for the u...
Accessor for a shared memory. This class provides the functionalities to access created shared memori...
Definition ddtClient.hpp:32
Definition ddtMemoryAccessor.hpp:175