You cannot select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
208 lines
6.3 KiB
C++
208 lines
6.3 KiB
C++
// Copyright 2016 Proyectos y Sistemas de Mantenimiento SL (eProsima).
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
/*!
|
|
* @file Publisher.cxx
|
|
* This file contains the implementation of the publisher functions.
|
|
*
|
|
* This file was generated by the tool fastddsgen.
|
|
*/
|
|
|
|
#include "Publisher.hpp"
|
|
|
|
#include <condition_variable>
|
|
#include <csignal>
|
|
#include <stdexcept>
|
|
#include <thread>
|
|
|
|
#include <fastdds/dds/domain/DomainParticipantFactory.hpp>
|
|
#include <fastdds/dds/log/Log.hpp>
|
|
#include <fastdds/dds/publisher/DataWriter.hpp>
|
|
#include <fastdds/dds/publisher/Publisher.hpp>
|
|
#include <fastdds/dds/publisher/qos/DataWriterQos.hpp>
|
|
#include <fastdds/dds/publisher/qos/PublisherQos.hpp>
|
|
|
|
#include "SystemPubSubTypes.hpp"
|
|
|
|
#include "msg.hpp"
|
|
#include "MsgHandler.hpp"
|
|
|
|
using namespace eprosima::fastdds::dds;
|
|
|
|
PublisherApp::PublisherApp(
|
|
const int& domain_id)
|
|
: factory_(nullptr)
|
|
, participant_(nullptr)
|
|
, publisher_(nullptr)
|
|
, topic_(nullptr)
|
|
, writer_(nullptr)
|
|
, type_(new PrintRspPubSubType())
|
|
, matched_(0)
|
|
, samples_sent_(0)
|
|
, stop_(false)
|
|
{
|
|
//
|
|
|
|
// Create the participant
|
|
DomainParticipantQos pqos = PARTICIPANT_QOS_DEFAULT;
|
|
pqos.name("Print_pub_participant");
|
|
pqos.wire_protocol().builtin.discovery_config.leaseDuration = Duration_t(60, 0);
|
|
pqos.wire_protocol().builtin.discovery_config.leaseDuration_announcementperiod = Duration_t(30, 0);
|
|
factory_ = DomainParticipantFactory::get_shared_instance();
|
|
participant_ = factory_->create_participant(domain_id, pqos, nullptr, StatusMask::none());
|
|
if (participant_ == nullptr)
|
|
{
|
|
throw std::runtime_error("PrintRsp Participant initialization failed");
|
|
}
|
|
|
|
// Register the type
|
|
type_.register_type(participant_);
|
|
|
|
// Create the publisher
|
|
PublisherQos pub_qos = PUBLISHER_QOS_DEFAULT;
|
|
participant_->get_default_publisher_qos(pub_qos);
|
|
publisher_ = participant_->create_publisher(pub_qos, nullptr, StatusMask::none());
|
|
if (publisher_ == nullptr)
|
|
{
|
|
throw std::runtime_error("PrintRsp Publisher initialization failed");
|
|
}
|
|
|
|
// Create the topic
|
|
TopicQos topic_qos = TOPIC_QOS_DEFAULT;
|
|
participant_->get_default_topic_qos(topic_qos);
|
|
topic_ = participant_->create_topic("PrintRspTopic", type_.get_type_name(), topic_qos);
|
|
if (topic_ == nullptr)
|
|
{
|
|
throw std::runtime_error("PrintRsp Topic initialization failed");
|
|
}
|
|
|
|
// Create the data writer
|
|
DataWriterQos writer_qos = DATAWRITER_QOS_DEFAULT;
|
|
publisher_->get_default_datawriter_qos(writer_qos);
|
|
writer_qos.reliability().kind = ReliabilityQosPolicyKind::RELIABLE_RELIABILITY_QOS;
|
|
writer_qos.durability().kind = DurabilityQosPolicyKind::TRANSIENT_LOCAL_DURABILITY_QOS;
|
|
writer_qos.history().kind = HistoryQosPolicyKind::KEEP_LAST_HISTORY_QOS;
|
|
writer_ = publisher_->create_datawriter(topic_, writer_qos, this, StatusMask::all());
|
|
if (writer_ == nullptr)
|
|
{
|
|
throw std::runtime_error("PrintRsp DataWriter initialization failed");
|
|
}
|
|
}
|
|
|
|
PublisherApp::~PublisherApp()
|
|
{
|
|
if (nullptr != participant_)
|
|
{
|
|
// Delete DDS entities contained within the DomainParticipant
|
|
participant_->delete_contained_entities();
|
|
|
|
// Delete DomainParticipant
|
|
factory_->delete_participant(participant_);
|
|
}
|
|
}
|
|
|
|
void PublisherApp::on_publication_matched(
|
|
DataWriter* writer,
|
|
const PublicationMatchedStatus& info)
|
|
{
|
|
if (info.current_count_change == 1)
|
|
{
|
|
{
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
matched_ = info.current_count;
|
|
}
|
|
std::cout << writer->get_topic()->get_name() << " Publisher matched." << std::endl;
|
|
cv_.notify_one();
|
|
}
|
|
else if (info.current_count_change == -1)
|
|
{
|
|
{
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
matched_ = info.current_count;
|
|
}
|
|
std::cout << writer->get_topic()->get_name() << " Publisher unmatched." << std::endl;
|
|
}
|
|
else
|
|
{
|
|
std::cout << info.current_count_change
|
|
<< " is not a valid value for PublicationMatchedStatus current count change" << std::endl;
|
|
}
|
|
}
|
|
|
|
void PublisherApp::run(std::shared_ptr<MsgHandler> handler)
|
|
{
|
|
while (!is_stopped())
|
|
{
|
|
std::unique_lock<std::mutex> lock(MsgData::queue_cv_mtx_, std::try_to_lock);
|
|
if (lock.owns_lock())
|
|
{
|
|
if(!MsgData::PrintReq_queue_.empty())
|
|
{
|
|
PrintReq req = std::move(MsgData::PrintReq_queue_.front());
|
|
MsgData::PrintReq_queue_.pop();
|
|
lock.unlock();
|
|
|
|
std::cout << "index: " << req.index() << std::endl;
|
|
handler->HandleDdsMsg(req.msg());
|
|
|
|
PrintRsp rsp;
|
|
rsp.index() = req.index();
|
|
rsp.msg() = "success";
|
|
writer_->write(&rsp);
|
|
}
|
|
else
|
|
{
|
|
lock.unlock();
|
|
}
|
|
}
|
|
// Wait for period or stop event
|
|
std::unique_lock<std::mutex> period_lock(mutex_);
|
|
cv_.wait_for(period_lock, std::chrono::milliseconds(period_ms_), [this]()
|
|
{
|
|
return is_stopped();
|
|
});
|
|
}
|
|
}
|
|
|
|
|
|
bool PublisherApp::publish()
|
|
{
|
|
bool ret = false;
|
|
// Wait for the data endpoints discovery
|
|
std::unique_lock<std::mutex> matched_lock(mutex_);
|
|
cv_.wait(matched_lock, [&]()
|
|
{
|
|
// at least one has been discovered
|
|
return ((matched_ > 0) || is_stopped());
|
|
});
|
|
|
|
if (!is_stopped())
|
|
{
|
|
/* Initialize your structure here */
|
|
WeighRsp sample_;
|
|
ret = (RETCODE_OK == writer_->write(&sample_));
|
|
}
|
|
return ret;
|
|
}
|
|
|
|
bool PublisherApp::is_stopped()
|
|
{
|
|
return stop_.load();
|
|
}
|
|
|
|
void PublisherApp::stop()
|
|
{
|
|
stop_.store(true);
|
|
cv_.notify_one();
|
|
} |