ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtDataConsumer.hpp
Go to the documentation of this file.
1
12
13#ifndef DDTDATACONSUMER_HPP_
14#define DDTDATACONSUMER_HPP_
15
17#include "ddt/ddtStatistics.hpp"
18
19namespace mal = ::elt::mal;
20namespace datatransfer = ::elt::ddt::datatransfer;
21
25typedef boost::signals2::signal<void(datatransfer::NotificationType,
26 const std::string&)>
28
32typedef signal_n::slot_type slot_n;
33
34namespace ddt {
35
44 public:
54 DdtDataConsumer(const std::string& data_stream_identifier,
55 const int32_t latency, const int32_t deadline,
56 DdtLogger* ddt_logger);
57
68 DdtDataConsumer(const std::string& data_stream_identifier,
69 const std::string& subscription_uri, const int32_t latency,
70 const int32_t deadline, DdtLogger* ddt_logger);
71
75 ~DdtDataConsumer() override;
76
80 void StartSubscription();
81
85 void StopSubscription();
86
100
106 void AddUuid(std::string uuid, SubscriberType type);
107
112 void RemoveUuid(const std::string uuid);
113
118 void Notify(const NotificationType type) override;
119
125
132
138 std::map<std::string, SubscriberType> get_subscribers();
139
144 std::string get_remote_broker_uri() const;
145
150 int32_t get_notification_port() const;
151
157
161 void ResetStatistics();
162
167 std::string get_publishing_uri() const;
168
173 void set_remote_broker_uri(const std::string& remote_uri);
174
179 void set_number_of_samples(const int32_t num_samples);
180
186 void set_memory_accessor(DdtMemoryAccessor* mem_accessor);
187
192 void set_notification_port(const int32_t noti_port);
193
198 void set_publishing_uri(const std::string pub_uri);
199
204 void set_originating_broker(const std::string orig_broker);
205
210
211 protected:
217 void Init(const std::string& ds_id, DdtLogger* ddt_logger);
218
223
228
233
238
242 std::string publishing_uri;
243
247 std::chrono::system_clock::time_point last_received;
248
252 uint64_t total_samples = 0;
253
257 uint64_t total_bytes = 0;
258
262 uint64_t total_latency = 0;
263
264 private:
273 void CreateSubscriber(const std::string& subscription_uri,
274 const int32_t latency, const int32_t deadline);
275
279 void Subscribe();
280
288 void CreateNotifier(const int32_t latency, const int32_t deadline);
289
294 void ReceiveDataEvent(
295 const mal::ps::DataEvent<datatransfer::DataPacket>& event);
296
297 std::unique_ptr<
298 mal::ps::Subscriber<datatransfer::DataPacket>,
299 std::default_delete<mal::ps::Subscriber<datatransfer::DataPacket> > >
300 data_subscriber;
301 mal::ps::DataEventFilter<datatransfer::DataPacket> filter;
302 std::shared_ptr<datatransfer::DataPacket> ddt_key_sample;
303 std::shared_ptr<datatransfer::DataPacket> ddt_data_packet;
304
305 std::string remote_broker_uri;
306 std::string originating_broker;
307 std::mutex subscriber_mutex;
308 std::mutex statistics_mutex;
309
314 std::map<std::string, SubscriberType> subscriber_map;
315
316 std::promise<void> exit_signal;
317 std::future<void> future_object;
318
323 std::unique_ptr<mal::ps::InstancePublisher<datatransfer::NotificationSample>,
324 std::default_delete<mal::ps::InstancePublisher<
325 datatransfer::NotificationSample> > >
326 notifier;
327 std::shared_ptr<datatransfer::NotificationSample> ddt_notification_sample;
328};
329
330} // namespace ddt
331
332#endif /* DDTDATACONSUMER_HPP_ */
void Notify(const NotificationType type) override
Definition ddtDataConsumer.cpp:85
uint64_t total_latency
Definition ddtDataConsumer.hpp:262
void AddUuid(std::string uuid, SubscriberType type)
Definition ddtDataConsumer.cpp:220
void set_originating_broker(const std::string orig_broker)
Definition ddtDataConsumer.cpp:320
signal_n notification_signal
Definition ddtDataConsumer.hpp:209
int32_t number_of_samples
Definition ddtDataConsumer.hpp:232
void set_notification_port(const int32_t noti_port)
Definition ddtDataConsumer.cpp:312
void set_memory_accessor(DdtMemoryAccessor *mem_accessor)
Definition ddtDataConsumer.cpp:308
int32_t get_number_of_remote_subscribers()
Definition ddtDataConsumer.cpp:242
DdtDataConsumer(const std::string &data_stream_identifier, const int32_t latency, const int32_t deadline, DdtLogger *ddt_logger)
Definition ddtDataConsumer.cpp:17
void ResetStatistics()
Definition ddtDataConsumer.cpp:289
~DdtDataConsumer() override
void StartSubscription()
Definition ddtDataConsumer.cpp:324
std::string get_publishing_uri() const
Definition ddtDataConsumer.cpp:296
int32_t notification_port
Definition ddtDataConsumer.hpp:237
std::string data_stream_identifier
Definition ddtDataConsumer.hpp:227
std::map< std::string, SubscriberType > get_subscribers()
Definition ddtDataConsumer.cpp:262
int32_t get_number_of_subscribers()
Definition ddtDataConsumer.cpp:237
uint64_t total_samples
Definition ddtDataConsumer.hpp:252
void set_publishing_uri(const std::string pub_uri)
Definition ddtDataConsumer.cpp:316
int32_t get_notification_port() const
Definition ddtDataConsumer.cpp:271
void RemoveUuid(const std::string uuid)
Definition ddtDataConsumer.cpp:226
SubscriberType
Definition ddtDataConsumer.hpp:90
@ REMOTE
Definition ddtDataConsumer.hpp:98
@ LOCAL
Definition ddtDataConsumer.hpp:94
void set_remote_broker_uri(const std::string &remote_uri)
Definition ddtDataConsumer.cpp:300
DdtStatistics get_statistics()
Definition ddtDataConsumer.cpp:275
uint64_t total_bytes
Definition ddtDataConsumer.hpp:257
std::string get_remote_broker_uri() const
Definition ddtDataConsumer.cpp:267
DdtMemoryAccessor * memory_accessor
Definition ddtDataConsumer.hpp:222
std::string publishing_uri
Definition ddtDataConsumer.hpp:242
void StopSubscription()
Definition ddtDataConsumer.cpp:333
void Init(const std::string &ds_id, DdtLogger *ddt_logger)
Definition ddtDataConsumer.cpp:37
std::chrono::system_clock::time_point last_received
Definition ddtDataConsumer.hpp:247
void set_number_of_samples(const int32_t num_samples)
Definition ddtDataConsumer.cpp:304
Definition ddtLogger.hpp:43
Definition ddtMemoryAccessor.hpp:265
NotificationType
Definition ddtProducerConsumerBase.hpp:50
DdtProducerConsumerBase(DdtLogger *ddt_logger)
Definition ddtProducerConsumerBase.cpp:16
boost::signals2::signal< void(datatransfer::NotificationType, const std::string &)> signal_n
Definition ddtDataConsumer.hpp:27
signal_n::slot_type slot_n
Definition ddtDataConsumer.hpp:32
Base class for DdtDataProducer and DdtDataConsumer. This class serves as a base class for DdtDataProd...
Statistics for the monitoring API. This struct contains the raw values for the monitoring API.
Definition ddtClient.hpp:32
Definition ddtStatistics.hpp:18