123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232 |
- #ifdef USE_FASTRTPS
- #include <fastrtps/participant/Participant.h>
- #include <fastrtps/attributes/ParticipantAttributes.h>
- #include <fastrtps/subscriber/Subscriber.h>
- #include <fastrtps/attributes/SubscriberAttributes.h>
- #include <fastrtps/Domain.h>
- #include <fastcdr/FastBuffer.h>
- #include <fastcdr/FastCdr.h>
- #include <fastcdr/Cdr.h>
- #include <fastrtps/participant/Participant.h>
- #include <fastrtps/attributes/ParticipantAttributes.h>
- #include <fastrtps/attributes/PublisherAttributes.h>
- #include <fastrtps/publisher/Publisher.h>
- #include <fastrtps/Domain.h>
- #include <fastdds/rtps/transport/shared_mem/SharedMemTransportDescriptor.h>
- #include <fastdds/rtps/transport/UDPv4TransportDescriptor.h>
- #include <fastdds/rtps/transport/TCPv4TransportDescriptor.h>
- #include "TopicsSubscriber.h"
- using namespace eprosima::fastrtps;
- using namespace eprosima::fastrtps::rtps;
- using namespace eprosima::fastdds::rtps;
- #include <QDateTime>
- TopicsSubscriber::TopicsSubscriber() : mp_participant(nullptr), mp_subscriber(nullptr) {}
- TopicsSubscriber::~TopicsSubscriber() { Domain::removeParticipant(mp_participant);}
- bool TopicsSubscriber::init(const char * strtopic,const char * strpubip,const unsigned short nPort,int ntype)
- {
- strncpy(mstrtopic,strtopic,255);
-
- ParticipantAttributes PParam;
- PParam.rtps.setName("Participant_subscriber");
- PParam.rtps.builtin.discovery_config.discoveryProtocol = DiscoveryProtocol_t::SIMPLE;
- PParam.rtps.builtin.discovery_config.use_SIMPLE_EndpointDiscoveryProtocol = true;
- PParam.rtps.builtin.discovery_config.m_simpleEDP.use_PublicationReaderANDSubscriptionWriter = true;
- PParam.rtps.builtin.discovery_config.m_simpleEDP.use_PublicationWriterANDSubscriptionReader = true;
- PParam.rtps.builtin.discovery_config.leaseDuration = c_TimeInfinite;
-
- PParam.rtps.useBuiltinTransports = false;
- if(ntype == 0)
- {
- auto sm_transport = std::make_shared<SharedMemTransportDescriptor>();
- sm_transport->segment_size(100*1024*1024);
- PParam.rtps.userTransports.push_back(sm_transport);
-
- auto udp_transport = std::make_shared<UDPv4TransportDescriptor>();
-
- PParam.rtps.userTransports.push_back(udp_transport);
- }
- else
- {
- auto tcp2_transport = std::make_shared<TCPv4TransportDescriptor>();
-
- Locator_t initial_peer_locator;
- initial_peer_locator.kind = LOCATOR_KIND_TCPv4;
- IPLocator::setIPv4(initial_peer_locator, strpubip);
- initial_peer_locator.port = nPort;
- PParam.rtps.builtin.initialPeersList.push_back(initial_peer_locator);
-
-
-
- tcp2_transport->wait_for_tcp_negotiation = false;
-
- PParam.rtps.userTransports.push_back(tcp2_transport);
- }
- mp_participant = Domain::createParticipant(PParam);
- if(mp_participant == nullptr)
- {
- return false;
- }
-
- Domain::registerType(mp_participant, static_cast<TopicDataType*>(&myType));
-
- SubscriberAttributes Rparam;
- Rparam.topic.topicKind = NO_KEY;
- Rparam.topic.topicDataType = myType.getName();
- Rparam.topic.topicName = strtopic;
- Rparam.topic.historyQos.depth = 30;
- Rparam.topic.resourceLimitsQos.max_samples = 50;
- Rparam.topic.resourceLimitsQos.allocated_samples = 20;
- Rparam.qos.m_reliability.kind = RELIABLE_RELIABILITY_QOS;
- mp_subscriber = Domain::createSubscriber(mp_participant,Rparam, static_cast<SubscriberListener*>(&m_listener));
- if(mp_subscriber == nullptr)
- {
- return false;
- }
- return true;
- }
- void TopicsSubscriber::SubListener::onSubscriptionMatched(Subscriber* sub,MatchingInfo& info)
- {
- (void)sub;
- if (info.status == MATCHED_MATCHING)
- {
- n_matched++;
- std::cout << "Subscriber matched" << std::endl;
- }
- else
- {
- n_matched--;
- std::cout << "Subscriber unmatched" << std::endl;
- }
- }
- void TopicsSubscriber::SubListener::setReceivedTopicFunction(ModuleFun xFun)
- {
- mFun = xFun;
- mbSetFun = true;
- }
- void TopicsSubscriber::SubListener::onNewDataMessage(Subscriber* sub)
- {
-
- TopicSample::Message st;
- static int ncount = 0;
- static int nmaxlatancy = 0;
- std::cout<<"new msg"<<std::endl;
- std::cout<<"count is "<<sub->getUnreadCount()<<std::endl;
- if(sub->takeNextData(&st, &m_info))
- {
- if(m_info.sampleKind == ALIVE)
- {
-
- ++n_msg;
- ncount++;
- std::cout << "Sample received, count=" << st.counter() <<" total: "<<ncount<<std::endl;
- qint64 timex = QDateTime::currentMSecsSinceEpoch();
- int nlatancy = (timex - st.sendtime());
- if(nlatancy>nmaxlatancy)nmaxlatancy = nlatancy;
- std::cout<<" latency is "<<nlatancy<<" max: "<<nmaxlatancy<<std::endl;
- std::cout<<" size is "<<st.xdata().size()<<std::endl;
- QDateTime dt = QDateTime::fromMSecsSinceEpoch(st.sendtime());
- if(mbSetFun) mFun((char *)(st.xdata().data()),st.xdata().size(),st.counter(),&dt,st.msgname().data());
- }
- }
- }
- void TopicsSubscriber::setReceivedTopicFunction(ModuleFun xFun)
- {
- m_listener.setReceivedTopicFunction(xFun);
- }
- void TopicsSubscriber::run()
- {
- std::cout << "Waiting for Data, press Enter to stop the Subscriber. "<<std::endl;
- std::cin.ignore();
- std::cout << "Shutting down the Subscriber." << std::endl;
- }
- #endif
|