ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtConnectionManager.hpp
Go to the documentation of this file.
1
10
11#ifndef DDTCONNECTIONMANAGER_HPP_
12#define DDTCONNECTIONMANAGER_HPP_
13
14#define BOOST_BIND_GLOBAL_PLACEHOLDERS
15
16#include <boost/bind/bind.hpp>
17#include <boost/property_tree/ini_parser.hpp>
18#include <boost/property_tree/ptree.hpp>
19#include <boost/signals2/signal.hpp>
20#include <boost/bind/bind.hpp>
21#include <chrono>
22#include <mal/rr/ServerAmi.hpp>
23#include <mal/rr/ServerContextProvider.hpp>
24#include <mal/rr/qos/ReplyTime.hpp>
25#include <memory>
26
27#include "ddt/ddtClient.hpp"
30#include "ddt/ddtLogger.hpp"
32
33namespace mal = ::elt::mal;
34namespace datatransfer = ::elt::ddt::datatransfer;
35
36namespace ddt {
37
43 : public virtual datatransfer::DataBrokerRegistration {
44 public:
51 explicit DdtConnectionManager(DdtLogger* ddt_logger,
52 const std::string uri_string,
53 const std::string config);
54
62 explicit DdtConnectionManager(DdtMemoryManager* const mmgr,
63 DdtLogger* ddt_logger,
64 const std::string uri_string,
65 const std::string config);
66
70 ~DdtConnectionManager() override;
71
89 int32_t RegisterPublisher(const std::string& data_stream_identifier,
90 const int32_t latency, int32_t deadline,
91 const int32_t max_data_sample_size,
92 const int32_t number_of_samples,
93 const bool compute_checksum,
94 const std::string& publishing_uri) override;
95
102 int32_t UnregisterPublisher(
103 const std::string& data_stream_identifier) override;
104
109 int32_t UnregisterPublishers();
110
120 void PublishData(const std::string& data_stream_identifier) override;
121
142 int32_t RegisterSubscriber(const std::string& subscriber_uuid,
143 const std::string& data_stream_identifier,
144 const std::string& remote_broker_uri,
145 const int32_t latency,
146 const int32_t deadline) override;
147
161 int32_t RegisterRemoteSubscriber(const std::string& remote_broker,
162 const std::string& subscriber_uuid,
163 const std::string& data_stream_identifier,
164 const int32_t latency,
165 const int32_t deadline) override;
166
178 int32_t UnregisterSubscriber(const std::string& data_stream_identifier,
179 const std::string& subscriber_uuid) override;
180
185 int32_t UnregisterSubscribers();
186
193 const std::string& data_stream_identifier) override;
194
200 int32_t get_number_of_samples(
201 const std::string& data_stream_identifier) override;
202
208 std::string get_publishing_uri(
209 const std::string& data_stream_identifier) override;
210
216 int32_t get_notification_port(
217 const std::string& data_stream_identifier) override;
218
224 std::vector<std::string> get_statistics(
225 const std::string& data_stream_identifier) override;
226
231 int32_t get_heartbeat_interval() override;
232
237 int32_t get_heartbeat_timeout() override;
238
244 bool get_compute_checksum(const std::string& data_stream_identifier) override;
245
251 std::string get_shm_id(const std::string& data_stream_identifier) override;
252
258 std::string get_shm_full_path(const std::string& data_stream_identifier) override;
259
264 std::string get_broker_uri() override;
265
270 void UpdateHeartbeat(const std::string& identifier) override;
271
281 const std::string& data_stream_identifier) override;
282
289 bool CheckPublisherExists(const std::string& data_stream_identifier) override;
290
299 const std::string& remote_broker_uri,
300 const std::string& data_stream_identifier) override;
301
308 void UpdateStatistics(const std::string& data_stream_identifier,
309 const int32_t datavec_size,
310 const uint64_t source_timestamp) override;
311
317 int32_t GetMaxPossibleBufferSize(int32_t max_data_sample_size) override;
318
323 std::vector<std::string> GetRegisteredStreams() override;
324
329 std::vector<std::string> GetConnectedBrokers() override;
330
331 protected:
335 void LoadDefaults();
336
342 std::string GetConfigPath() const;
343
347 void ReadIni();
348
355 std::string CreateSubscriptionUri(const std::string publishing_uri,
356 const std::string remote_broker_uri) const;
357
365 bool CheckStreamIdInUse(const std::string& data_stream_identifier);
366
373 void CreateStatistics(const std::string& data_stream_identifier,
374 const int32_t number_of_samples);
375
380 void ResetStatistics(const std::string& data_stream_identifier);
381
387 const int32_t SHM_TIMEOUT_DEFAULT = 10;
388
393 const int32_t WAITING_TIME_DEFAULT = 1000;
394
399 const int32_t REPLY_TIME_DEFAULT = 6;
400
405 const int32_t HEARTBEAT_INTERVAL_DEFAULT = 1;
406
412 const int32_t HEARTBEAT_TIMEOUT_DEFAULT = 10;
413
417 const int32_t WAIT_FOR_CONNECTION_DEFAULT = 10;
418
423 const int32_t LATENCY_DEFAULT = 10000;
424
428 const int32_t DEADLINE_DEFAULT = 10;
429
434
438 int32_t shm_timeout;
439
444
448 int32_t reply_time;
449
454
459
464
468 int32_t latency;
469
473 int32_t deadline;
474
479
484 std::set<std::string> registered_publishers;
485
491 std::map<std::string, DdtStatistics> statistics_map;
492
493 private:
501 void Init(DdtMemoryManager* mmgr, DdtLogger* ddt_logger,
502 const std::string uri_string, const std::string config);
503
507 void PrintConfigValues();
508
512 void StartHeartbeat();
513
517 void StopHeartbeat();
518
524 void HeartbeatThread();
525
531 void NotifySubscribers(const std::string& data_stream_identifier,
532 std::unique_lock<std::mutex>& producer_lock);
533
541 int32_t CheckPublisherUsingStream(const std::string data_stream_identifier,
542 const std::string remote_broker_uri);
543
559 void CreateDataConsumer(std::unique_lock<std::mutex>& consumer_lock,
560 const std::string subscriber_uuid,
561 const std::string data_stream_identifier,
562 const std::string remote_broker_uri,
563 const std::string originating_broker,
564 const int32_t latency, const int32_t deadline,
565 const std::string subscription_uri,
566 const int32_t number_of_samples);
567
573 void SearchAndUnregSubscriber(const std::string identifier);
574
584 int32_t RegisterLocalSubscriber(const std::string subscriber_uuid,
585 const std::string data_stream_identifier,
586 const int32_t latency,
587 const int32_t deadline);
588
599 int32_t RegisterLocalSubscriberRemote(
600 const std::string subscriber_uuid,
601 const std::string data_stream_identifier,
602 const std::string remote_broker_uri, const int32_t latency,
603 const int32_t deadline);
604
611 void ProcessNotificationEvent(const datatransfer::NotificationType type,
612 const std::string& data_stream_identifier);
613
623 void FreeShmThread(const std::string& data_stream_identifier,
624 const int shm_timeout);
625
632 bool CheckPublisherReregistration(const std::string& data_stream_identifier);
633
643 bool CheckSharedMemoryRecreation(const std::string& data_stream_identifier,
644 const int32_t max_data_sample_size,
645 const int32_t number_of_samples);
646
651 void Publish(const std::string& data_stream_identifier);
652
657 void PubRegNotification(const std::string& data_stream_identifier);
658
663 void PubUnregNotification(const std::string& data_stream_identifier);
664
670 std::map<std::string, DdtDataProducer*> producer_map;
671
677 std::map<std::string, DdtDataConsumer*> consumer_map;
678
685 std::map<std::string,
686 std::chrono::time_point<std::chrono::high_resolution_clock>>
687 client_map;
688
694 std::map<std::string, std::string> connected_brokers_map;
695
699 DdtMemoryManager* memory_manager;
700
706 std::map<std::string, DdtClient*> ddt_clients;
707
711 std::atomic<bool> stop_threads;
712
716 std::atomic<int> thread_counter;
717
721 std::mutex producer_mutex;
722
726 std::mutex consumer_mutex;
727
731 std::mutex client_mutex;
732
736 std::mutex statistics_mutex;
737
741 std::mutex ddt_clients_mutex;
742
746 std::mutex registered_publishers_mutex;
747
751 std::mutex connected_brokers_mutex;
752
756 DdtLogger* logger;
757
761 std::promise<void> exit_signal;
762
766 std::future<void> future_object;
767
771 std::atomic<bool> heartbeat_active;
772
776 std::string broker_uri;
777
781 std::string config_file;
782
786 boost::signals2::connection connection;
787
791 const int32_t SHM_TIMEOUT_MIN = 2;
792
796 const int32_t WAITING_TIME_MIN = 1000;
797
801 const int32_t REPLY_TIME_MIN = 2;
802
806 const int32_t HEARTBEAT_INTERVAL_MIN = 0;
807
811 const int32_t HEARTBEAT_TIMEOUT_MIN = 3;
812
816 const int32_t WAIT_FOR_CONNECTION_MIN = 2;
817
821 const int32_t LATENCY_MIN = 1000;
822
826 const int32_t DEADLINE_MIN = 1;
827
831 const int32_t MAX_CONSECUTIVE_FAILURES_MIN = 1;
832
836 const int32_t NUM_RETRIES = 30;
837};
838
839} // namespace ddt
840
841#endif /* DDTCONNECTIONMANAGER_HPP_ */
int32_t get_heartbeat_timeout() override
Definition ddtConnectionManager.cpp:1543
int32_t waiting_time
Definition ddtConnectionManager.hpp:443
std::string get_shm_full_path(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1559
std::string GetConfigPath() const
Definition ddtConnectionManager.cpp:81
int32_t RegisterRemoteSubscriber(const std::string &remote_broker, const std::string &subscriber_uuid, const std::string &data_stream_identifier, const int32_t latency, const int32_t deadline) override
Definition ddtConnectionManager.cpp:1196
const int32_t REPLY_TIME_DEFAULT
Definition ddtConnectionManager.hpp:399
std::string CreateSubscriptionUri(const std::string publishing_uri, const std::string remote_broker_uri) const
Definition ddtConnectionManager.cpp:1566
int32_t get_heartbeat_interval() override
Definition ddtConnectionManager.cpp:1539
int32_t UnregisterSubscribers()
Definition ddtConnectionManager.cpp:1334
void UpdateHeartbeat(const std::string &identifier) override
Definition ddtConnectionManager.cpp:407
bool get_compute_checksum(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1547
const int32_t WAITING_TIME_DEFAULT
Definition ddtConnectionManager.hpp:393
int32_t deadline
Definition ddtConnectionManager.hpp:473
std::string get_broker_uri() override
Definition ddtConnectionManager.cpp:1564
std::map< std::string, DdtStatistics > statistics_map
Definition ddtConnectionManager.hpp:491
const int32_t LATENCY_DEFAULT
Definition ddtConnectionManager.hpp:423
int32_t UnregisterPublisher(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:578
int32_t RegisterPublisher(const std::string &data_stream_identifier, const int32_t latency, int32_t deadline, const int32_t max_data_sample_size, const int32_t number_of_samples, const bool compute_checksum, const std::string &publishing_uri) override
Definition ddtConnectionManager.cpp:416
const int32_t HEARTBEAT_TIMEOUT_DEFAULT
Definition ddtConnectionManager.hpp:412
std::vector< std::string > GetRegisteredStreams() override
Definition ddtConnectionManager.cpp:1699
void UpdateStatistics(const std::string &data_stream_identifier, const int32_t datavec_size, const uint64_t source_timestamp) override
Definition ddtConnectionManager.cpp:1632
std::string get_publishing_uri(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1462
bool CheckRemotePublisherExists(const std::string &remote_broker_uri, const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1413
void CreateStatistics(const std::string &data_stream_identifier, const int32_t number_of_samples)
Definition ddtConnectionManager.cpp:1654
void PublishData(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:706
bool CheckPubRegistrationPermitted(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1365
std::vector< std::string > get_statistics(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1486
void ReadIni()
Definition ddtConnectionManager.cpp:98
const int32_t WAIT_FOR_CONNECTION_DEFAULT
Definition ddtConnectionManager.hpp:417
~DdtConnectionManager() override
Definition ddtConnectionManager.cpp:33
void LoadDefaults()
Definition ddtConnectionManager.cpp:69
int32_t GetMaxPossibleBufferSize(int32_t max_data_sample_size) override
Definition ddtConnectionManager.cpp:1694
const int32_t HEARTBEAT_INTERVAL_DEFAULT
Definition ddtConnectionManager.hpp:405
bool CheckPublisherExists(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1400
const int32_t SHM_TIMEOUT_DEFAULT
Definition ddtConnectionManager.hpp:387
int32_t max_consecutive_failures
Definition ddtConnectionManager.hpp:478
int32_t UnregisterSubscriber(const std::string &data_stream_identifier, const std::string &subscriber_uuid) override
Definition ddtConnectionManager.cpp:1250
int32_t shm_timeout
Definition ddtConnectionManager.hpp:438
const int32_t MAX_CONSECUTIVE_FAILURES_DEFAULT
Definition ddtConnectionManager.hpp:433
const int32_t DEADLINE_DEFAULT
Definition ddtConnectionManager.hpp:428
int32_t heartbeat_interval
Definition ddtConnectionManager.hpp:453
int32_t get_notification_port(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1474
int32_t heartbeat_timeout
Definition ddtConnectionManager.hpp:458
int32_t UnregisterPublishers()
Definition ddtConnectionManager.cpp:682
int32_t reply_time
Definition ddtConnectionManager.hpp:448
std::vector< std::string > GetConnectedBrokers() override
Definition ddtConnectionManager.cpp:1724
int32_t latency
Definition ddtConnectionManager.hpp:468
std::string get_shm_id(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1554
bool CheckStreamIdInUse(const std::string &data_stream_identifier)
Definition ddtConnectionManager.cpp:1384
int32_t get_number_of_samples(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1457
int32_t RegisterSubscriber(const std::string &subscriber_uuid, const std::string &data_stream_identifier, const std::string &remote_broker_uri, const int32_t latency, const int32_t deadline) override
Definition ddtConnectionManager.cpp:796
std::set< std::string > registered_publishers
Definition ddtConnectionManager.hpp:484
DdtConnectionManager(DdtLogger *ddt_logger, const std::string uri_string, const std::string config)
Definition ddtConnectionManager.cpp:15
int32_t wait_for_connection
Definition ddtConnectionManager.hpp:463
void ResetStatistics(const std::string &data_stream_identifier)
Definition ddtConnectionManager.cpp:1680
int32_t get_max_data_sample_size(const std::string &data_stream_identifier) override
Definition ddtConnectionManager.cpp:1452
Definition ddtLogger.hpp:43
Definition ddtMemoryManager.hpp:63
Client class for the connection to remote brokers. This class creates MAL clients to connect to remot...
Data Consumer. This class provides the functionality to subscribe to a data stream,...
Data Producer. This class provides the functionality to publish data over network and enables sending...
Class to wrap the usage of log4cplus as logging utility. This file provides a wrapper class for the u...
Manager for shared memories. This class manages the handling of shared memories.
Definition ddtClient.hpp:32