Skip to content

Commit a0a4fee

Browse files
Create InitialConnection for TCP initial peers 2.13 & 2.10 & 2.6 (#4947)
* Refs #20650: Add test Signed-off-by: cferreiragonz <carlosferreira@eprosima.com> * Refs #20650: Create initial connect for initial peers Signed-off-by: cferreiragonz <carlosferreira@eprosima.com> --------- Signed-off-by: cferreiragonz <carlosferreira@eprosima.com>
1 parent 98ef5c7 commit a0a4fee

2 files changed

Lines changed: 73 additions & 3 deletions

File tree

src/cpp/rtps/builtin/discovery/participant/PDPSimple.cpp

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -379,7 +379,17 @@ bool PDPSimple::create_dcps_participant_endpoints()
379379

380380
WriterAttributes watt = create_builtin_writer_attributes();
381381
watt.endpoint.reliabilityKind = BEST_EFFORT;
382-
watt.endpoint.remoteLocatorList = m_discovery.initialPeersList;
382+
if (!m_discovery.initialPeersList.empty())
383+
{
384+
if (mp_RTPSParticipant->has_tcp_transports())
385+
{
386+
mp_RTPSParticipant->create_tcp_connections(m_discovery.initialPeersList);
387+
}
388+
else
389+
{
390+
watt.endpoint.remoteLocatorList = m_discovery.initialPeersList;
391+
}
392+
}
383393

384394
if (pattr.throughputController.bytesPerPeriod != UINT32_MAX && pattr.throughputController.periodMillisecs != 0)
385395
{

test/blackbox/common/BlackboxTestsTransportTCP.cpp

Lines changed: 62 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -674,7 +674,7 @@ TEST(TransportTCP, Client_reconnection)
674674
delete requester;
675675
}
676676

677-
// Test copy constructor and copy assignment for TCPv4
677+
// Test zero listening port for TCPv4
678678
TEST_P(TransportTCP, TCPv4_autofill_port)
679679
{
680680
PubSubReader<HelloWorldPubSubType> p1(TEST_TOPIC_NAME);
@@ -704,7 +704,7 @@ TEST_P(TransportTCP, TCPv4_autofill_port)
704704
EXPECT_TRUE(IPLocator::getPhysicalPort(p2_locators.begin()[0]) == port);
705705
}
706706

707-
// Test copy constructor and copy assignment for TCPv6
707+
// Test zero listening port for TCPv6
708708
TEST_P(TransportTCP, TCPv6_autofill_port)
709709
{
710710
PubSubReader<HelloWorldPubSubType> p1(TEST_TOPIC_NAME);
@@ -1238,6 +1238,66 @@ TEST_P(TransportTCP, large_message_large_data_send_receive)
12381238
reader.block_for_all();
12391239
}
12401240

1241+
// Test CreateInitialConnection for TCP
1242+
TEST_P(TransportTCP, TCP_initial_peers_connection)
1243+
{
1244+
PubSubWriter<HelloWorldPubSubType> p1(TEST_TOPIC_NAME);
1245+
PubSubReader<HelloWorldPubSubType> p2(TEST_TOPIC_NAME);
1246+
PubSubReader<HelloWorldPubSubType> p3(TEST_TOPIC_NAME);
1247+
1248+
// Add TCP Transport with listening port
1249+
auto p1_transport = std::make_shared<eprosima::fastdds::rtps::TCPv4TransportDescriptor>();
1250+
p1_transport->add_listener_port(global_port);
1251+
auto p2_transport = std::make_shared<eprosima::fastdds::rtps::TCPv4TransportDescriptor>();
1252+
p2_transport->add_listener_port(global_port + 1);
1253+
auto p3_transport = std::make_shared<eprosima::fastdds::rtps::TCPv4TransportDescriptor>();
1254+
p3_transport->add_listener_port(global_port - 1);
1255+
1256+
// Add initial peer to client
1257+
Locator_t initialPeerLocator;
1258+
initialPeerLocator.kind = LOCATOR_KIND_TCPv4;
1259+
IPLocator::setIPv4(initialPeerLocator, 127, 0, 0, 1);
1260+
initialPeerLocator.port = global_port;
1261+
LocatorList_t initial_peer_list;
1262+
initial_peer_list.push_back(initialPeerLocator);
1263+
1264+
// Setup participants
1265+
p1.disable_builtin_transport()
1266+
.add_user_transport_to_pparams(p1_transport);
1267+
1268+
p2.disable_builtin_transport()
1269+
.initial_peers(initial_peer_list)
1270+
.add_user_transport_to_pparams(p2_transport);
1271+
1272+
p3.disable_builtin_transport()
1273+
.initial_peers(initial_peer_list)
1274+
.add_user_transport_to_pparams(p3_transport);
1275+
1276+
// Init participants
1277+
p1.init();
1278+
p2.init();
1279+
p3.init();
1280+
ASSERT_TRUE(p1.isInitialized());
1281+
ASSERT_TRUE(p2.isInitialized());
1282+
ASSERT_TRUE(p3.isInitialized());
1283+
1284+
// Wait for discovery
1285+
p1.wait_discovery(2, std::chrono::seconds(0));
1286+
p2.wait_discovery(std::chrono::seconds(0), 1);
1287+
p3.wait_discovery(std::chrono::seconds(0), 1);
1288+
1289+
// Send and receive data
1290+
auto data = default_helloworld_data_generator();
1291+
p2.startReception(data);
1292+
p3.startReception(data);
1293+
1294+
p1.send(data);
1295+
EXPECT_TRUE(data.empty());
1296+
1297+
p2.block_for_all();
1298+
p3.block_for_all();
1299+
}
1300+
12411301
#ifdef INSTANTIATE_TEST_SUITE_P
12421302
#define GTEST_INSTANTIATE_TEST_MACRO(x, y, z, w) INSTANTIATE_TEST_SUITE_P(x, y, z, w)
12431303
#else

0 commit comments

Comments
 (0)