ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtDataSubscriber.hpp
Go to the documentation of this file.
1
11
12#ifndef DDTDATASUBSCRIBER_HPP_
13#define DDTDATASUBSCRIBER_HPP_
14
15#define BOOST_BIND_GLOBAL_PLACEHOLDERS
16
17#include <boost/bind.hpp>
18#include <boost/uuid/uuid.hpp>
19#include <boost/uuid/uuid_generators.hpp>
20#include <boost/uuid/uuid_io.hpp>
22
24
25namespace ddt {
26
33 public:
39
45 explicit DdtDataSubscriber(log4cplus::Logger const &log4cplus_logger);
46
50 ~DdtDataSubscriber() override;
51
52 int RegisterSubscriber(const std::string uri, const std::string dsi,
53 const std::string remote_uri,
54 const int32_t interval = 10) override;
55
56 int UnregisterSubscriber() override;
57
58 DataSample *ReadData() override;
59
64
69
75
81 boost::signals2::connection connect(const SignalT::slot_type &event_listener);
82
83 protected:
87 void LoadDefaults();
88
92 void StopThreads() override;
93
97 void ReadIni();
98
103
107 const int32_t MAX_AGE_DATA_SAMPLE_DEFAULT = 10000;
108
109 private:
113 void Init();
114
118 void PrintConfigValues();
119
123 bool CheckPath();
124
128 void InitializeNotificationSubscriber(
129 const std::string data_stream_identifier,
130 const int32_t notification_port);
131
136 void Subscribe();
137
142 void Reregister() override;
143
150 void NotificationEvent(
151 const mal::ps::DataEvent<datatransfer::NotificationSample> &event);
152
156 int CheckPublisher();
157
161 int CreateAccessor();
162
166 void ReopenShm();
167
168 void LogSubscriberParameter();
169
170 void PrintErrorMessage(const int result);
171
172 DdtStatisticsClient *statistics_client = nullptr;
173 DdtMemoryAccessor *memory_accessor = nullptr;
174 std::string shm_id;
175 std::string data_stream_identifier;
176 std::string subscriber_uuid;
177 std::string broker_uri;
178 std::string remote_broker_uri;
179 int32_t reading_interval = 0;
180 std::atomic<bool> event_active{false};
181
182 std::unique_ptr<mal::ps::Subscriber<datatransfer::NotificationSample>,
183 std::default_delete<
184 mal::ps::Subscriber<datatransfer::NotificationSample> > >
185 notification_subscriber;
186 std::shared_ptr<datatransfer::NotificationSample> ddt_key_notification;
187 std::shared_ptr<datatransfer::NotificationSample> ddt_notification;
188 mal::ps::DataEventFilter<datatransfer::NotificationSample> filter;
189
190 std::promise<void> exit_signal;
191 std::future<void> future_object;
192
196 std::thread subscribe_thread;
197
201 std::thread reregister_thread;
202
203 const int32_t NUM_RETRIES = 10;
204 const int32_t MAX_AGE_DATA_SAMPLE_MIN = 2000;
205};
206
207} // namespace ddt
208
209#endif /* DDTDATASUBSCRIBER_HPP_ */
void StopThreads() override
Definition ddtDataSubscriber.cpp:42
DataSample * ReadData() override
Definition ddtDataSubscriber.cpp:466
void ReadIni()
Definition ddtDataSubscriber.cpp:72
void StartNotificationSubscription()
Definition ddtDataSubscriber.cpp:527
void StopNotificationSubscription()
Definition ddtDataSubscriber.cpp:543
boost::signals2::connection connect(const SignalT::slot_type &event_listener)
Definition ddtDataSubscriber.cpp:403
const int32_t MAX_AGE_DATA_SAMPLE_DEFAULT
Definition ddtDataSubscriber.hpp:107
~DdtDataSubscriber() override
Definition ddtDataSubscriber.cpp:33
DdtStatistics get_statistics()
Definition ddtDataSubscriber.cpp:446
int RegisterSubscriber(const std::string uri, const std::string dsi, const std::string remote_uri, const int32_t interval=10) override
Definition ddtDataSubscriber.cpp:175
int32_t max_age_data_sample
Definition ddtDataSubscriber.hpp:102
int UnregisterSubscriber() override
Definition ddtDataSubscriber.cpp:408
void LoadDefaults()
Definition ddtDataSubscriber.cpp:64
DdtDataSubscriber(DdtLogger *logger)
Definition ddtDataSubscriber.cpp:17
DdtLogger * logger
Definition ddtDataTransferLib.hpp:322
DdtDataTransferLib(DdtLogger *ddt_logger)
Definition ddtDataTransferLib.cpp:16
Definition ddtLogger.hpp:43
Definition ddtMemoryAccessor.hpp:265
Definition ddtStatisticsClient.hpp:26
Base class for DdtDataPublishers and DdtDataSubscribers. This is the base class for DdtDataPublishers...
Class for providing statistics. This class provides an API to query the statistics from a broker.
Definition ddtClient.hpp:32
Definition ddtMemoryAccessor.hpp:175
Definition ddtStatistics.hpp:18