ADD: added usage without a receiver queue
This commit is contained in:
@@ -23,7 +23,12 @@ namespace WHISPER {
|
||||
public:
|
||||
InternalUDPListener(std::uint16_t port, std::string address = "*");
|
||||
~InternalUDPListener();
|
||||
[[DEPRECATED]]
|
||||
void connect(std::shared_ptr<threadSafeQueue<WHISPER::Message>> receiver);
|
||||
void connect();
|
||||
|
||||
void addReceiverQueue(std::shared_ptr<threadSafeQueue<WHISPER::Message>> receiver);
|
||||
|
||||
void start();
|
||||
void stop();
|
||||
|
||||
|
||||
@@ -33,7 +33,7 @@ void InternalUDPListener::connect(std::shared_ptr<threadSafeQueue<WHISPER::Messa
|
||||
|
||||
if (receiverSocket_ == nullptr)
|
||||
{
|
||||
ctx = zmq::context_t(2);
|
||||
ctx = zmq::context_t();
|
||||
receiverSocket_ = std::make_shared<zmq::socket_t>(ctx,zmq::socket_type::dish);
|
||||
|
||||
std::string portAsString = std::to_string(port_);
|
||||
@@ -46,6 +46,30 @@ void InternalUDPListener::connect(std::shared_ptr<threadSafeQueue<WHISPER::Messa
|
||||
start();
|
||||
}
|
||||
|
||||
void InternalUDPListener::connect()
|
||||
{
|
||||
|
||||
if (receiverSocket_ == nullptr)
|
||||
{
|
||||
ctx = zmq::context_t();
|
||||
receiverSocket_ = std::make_shared<zmq::socket_t>(ctx,zmq::socket_type::dish);
|
||||
|
||||
std::string portAsString = std::to_string(port_);
|
||||
receiverSocket_->bind("udp://"+address_+":"+portAsString);
|
||||
///used to set a custom time out to the socket
|
||||
receiverSocket_->set(zmq::sockopt::rcvtimeo,100);
|
||||
}
|
||||
subscribe(WHISPER::MsgTopics::MANAGEMENT);
|
||||
|
||||
start();
|
||||
}
|
||||
|
||||
void InternalUDPListener::addReceiverQueue(std::shared_ptr<threadSafeQueue<WHISPER::Message>> receiver)
|
||||
{
|
||||
receiverQueue_ = receiver;
|
||||
}
|
||||
|
||||
|
||||
void InternalUDPListener::listen()
|
||||
{
|
||||
listening_ = true;
|
||||
@@ -59,14 +83,14 @@ void InternalUDPListener::listen()
|
||||
res = receiverSocket_->recv(msg,zmq::recv_flags::none);
|
||||
// std::string message = msg.to_string();
|
||||
|
||||
if (useHandl_ == true)
|
||||
{
|
||||
auto i = std::async(std::launch::async, MessageHandle_, msg.to_string());
|
||||
} else
|
||||
{
|
||||
WHISPER::Message receivedMessage(msg.to_string());
|
||||
receiverQueue_->addElement(receivedMessage);
|
||||
}
|
||||
if (useHandl_ == true)
|
||||
{
|
||||
auto i = std::async(std::launch::async, MessageHandle_, msg.to_string());
|
||||
} else if(receiverQueue_ != nullptr)
|
||||
{
|
||||
WHISPER::Message receivedMessage(msg.to_string());
|
||||
receiverQueue_->addElement(receivedMessage);
|
||||
}
|
||||
|
||||
}
|
||||
listening_ = false;
|
||||
|
||||
Reference in New Issue
Block a user