10#ifndef DDTDATATRANSFERLIB_HPP_
11#define DDTDATATRANSFERLIB_HPP_
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>
21#include <mal/rr/qos/ReplyTime.hpp>
22#include <mal/utility/LoadMal.hpp>
30namespace mal = ::elt::mal;
31namespace datatransfer = ::elt::ddt::datatransfer;
84 void SetQoS(
const int ddt_latency,
const int ddt_deadline);
106 std::unique_ptr<datatransfer::DataBrokerRegistrationSync,
107 std::default_delete<datatransfer::DataBrokerRegistrationSync>>
119 const bool compute_crc) {
146 const std::string remote_uri,
147 const int32_t interval = 10) {
181 void StartHeartbeat(
const int32_t interval,
const std::string
id);
290 datatransfer::DataBrokerRegistrationSync,
291 std::default_delete<datatransfer::DataBrokerRegistrationSync> >
378 const int32_t NUM_RECONNECT_RETRIES = 10;
383 void HeartbeatThread();
385 std::string identifier;
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;
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