Skip to content
4 changes: 3 additions & 1 deletion include/fastdds/dds/core/policy/ParameterTypes.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -184,7 +184,9 @@ enum ParameterId_t : uint16_t
PID_RTPS_RELIABLE_READER = 0x8201,
PID_READER_RESOURCE_LIMITS = 0x8202,
/* Participant specific */
PID_WIREPROTOCOL_CONFIG = 0x8300
PID_WIREPROTOCOL_CONFIG = 0x8300,
/* RPC specific */
PID_RPC_MORE_REPLIES = 0x8400
};

/*!
Expand Down
3 changes: 3 additions & 0 deletions include/fastdds/dds/subscriber/SampleInfo.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,9 @@ struct SampleInfo
//!Related Sample Identity (Extension for RPC)
rtps::SampleIdentity related_sample_identity;

//!Flag to indicate if there are more replies (Extension for RPC)
bool has_more_replies = false;

};

} // namespace dds
Expand Down
55 changes: 2 additions & 53 deletions include/fastdds/rtps/common/SampleIdentity.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,57 +34,6 @@ class FASTDDS_EXPORTED_API SampleIdentity
{
public:

/*!
* @brief Default constructor. Constructs an unknown SampleIdentity.
*/
SampleIdentity()
: writer_guid_(GUID_t::unknown())
, sequence_number_(SequenceNumber_t::unknown())
{
}

/*!
* @brief Copy constructor.
*/
SampleIdentity(
const SampleIdentity& sample_id)
: writer_guid_(sample_id.writer_guid_)
, sequence_number_(sample_id.sequence_number_)
{
}

/*!
* @brief Move constructor.
*/
SampleIdentity(
SampleIdentity&& sample_id)
: writer_guid_(std::move(sample_id.writer_guid_))
, sequence_number_(std::move(sample_id.sequence_number_))
{
}

/*!
* @brief Assignment operator.
*/
SampleIdentity& operator =(
const SampleIdentity& sample_id)
{
writer_guid_ = sample_id.writer_guid_;
sequence_number_ = sample_id.sequence_number_;
return *this;
}

/*!
* @brief Move constructor.
*/
SampleIdentity& operator =(
SampleIdentity&& sample_id)
{
writer_guid_ = std::move(sample_id.writer_guid_);
sequence_number_ = std::move(sample_id.sequence_number_);
return *this;
}

/*!
* @brief
*/
Expand Down Expand Up @@ -171,9 +120,9 @@ class FASTDDS_EXPORTED_API SampleIdentity

private:

GUID_t writer_guid_;
GUID_t writer_guid_ = GUID_t::unknown();

SequenceNumber_t sequence_number_;
SequenceNumber_t sequence_number_ = SequenceNumber_t::unknown();

friend std::istream& operator >>(
std::istream& input,
Expand Down
14 changes: 14 additions & 0 deletions include/fastdds/rtps/common/WriteParams.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,18 @@ class FASTDDS_EXPORTED_API WriteParams
return *this;
}

bool has_more_replies() const
{
return has_more_replies_;
}

WriteParams& has_more_replies(
bool more_replies)
{
has_more_replies_ = more_replies;
return *this;
}

static WriteParams WRITE_PARAM_DEFAULT;

/**
Expand Down Expand Up @@ -259,6 +271,8 @@ class FASTDDS_EXPORTED_API WriteParams
Time_t source_timestamp_{ -1, TIME_T_INFINITE_NANOSECONDS };
/// User write data
UserWriteDataPtr user_write_data_{nullptr};
/// Flag to indicate if there are more replies
bool has_more_replies_ = false;
};

} // namespace rtps
Expand Down
13 changes: 13 additions & 0 deletions src/cpp/fastdds/core/policy/ParameterList.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,19 @@ bool ParameterList::updateCacheChangeFromInlineQos(
break;
}

case PID_RPC_MORE_REPLIES:
{
// Ignore custom PID when coming from other vendors
if (rtps::c_VendorId_eProsima != change.vendor_id)
{
return true;
}

change.write_params.has_more_replies(true);

break;
}

case PID_STATUS_INFO:
{
ParameterStatusInfo_t p(pid, plength);
Expand Down
19 changes: 18 additions & 1 deletion src/cpp/fastdds/core/policy/ParameterSerializer.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -155,7 +155,9 @@ class ParameterSerializer<Parameter_t>
rtps::CDRMessage_t* cdr_message,
const rtps::SampleIdentity& sample_id)
{
if (cdr_message->pos + 28 > cdr_message->max_size)
uint32_t required_size = 24 + 4; // 24 for the sample identity and 4 for the PID and length

if (cdr_message->pos + required_size > cdr_message->max_size)
{
return false;
}
Expand All @@ -168,6 +170,7 @@ class ParameterSerializer<Parameter_t>
sample_id.writer_guid().entityId.value, rtps::EntityId_t::size);
rtps::CDRMessage::addInt32(cdr_message, sample_id.sequence_number().high);
rtps::CDRMessage::addUInt32(cdr_message, sample_id.sequence_number().low);

return true;
}

Expand Down Expand Up @@ -198,6 +201,19 @@ class ParameterSerializer<Parameter_t>
return true;
}

static inline bool add_parameter_more_replies(
rtps::CDRMessage_t* cdr_message)
{
if (cdr_message->pos + 4 > cdr_message->max_size)
{
return false;
}

rtps::CDRMessage::addUInt16(cdr_message, dds::PID_RPC_MORE_REPLIES);
rtps::CDRMessage::addUInt16(cdr_message, 0);
return true;
}

static inline uint32_t cdr_serialized_size(
const fastcdr::string_255& str)
{
Expand Down Expand Up @@ -846,6 +862,7 @@ inline bool ParameterSerializer<ParameterSampleIdentity_t>::read_content_from_cd
parameter.sample_id.writer_guid().entityId.value, rtps::EntityId_t::size);
valid &= rtps::CDRMessage::readInt32(cdr_message, &parameter.sample_id.sequence_number().high);
valid &= rtps::CDRMessage::readUInt32(cdr_message, &parameter.sample_id.sequence_number().low);

return valid;
}

Expand Down
1 change: 1 addition & 0 deletions src/cpp/fastdds/rpc/ReplierImpl.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ ReturnCode_t ReplierImpl::send_reply(

rtps::WriteParams wparams;
wparams.related_sample_identity(info.related_sample_identity);
wparams.has_more_replies(info.has_more_replies);

return replier_writer_->write(data, wparams);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ struct ReadTakeCommand
info.sample_identity.writer_guid(item->writerGUID);
info.sample_identity.sequence_number(item->sequenceNumber);
info.related_sample_identity = item->write_params.sample_identity();
info.has_more_replies = item->write_params.has_more_replies();

info.valid_data = true;

Expand Down
4 changes: 4 additions & 0 deletions src/cpp/rtps/history/WriterHistory.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -475,6 +475,10 @@ void WriterHistory::set_fragments(
{
inline_qos_size += (2 * fastdds::dds::ParameterSerializer<Parameter_t>::PARAMETER_SAMPLE_IDENTITY_SIZE);
}
if (change->write_params.has_more_replies())
{
inline_qos_size += 4u;
}
if (ChangeKind_t::ALIVE != change->kind && TopicKind_t::WITH_KEY == mp_writer->getAttributes().topicKind)
{
inline_qos_size += fastdds::dds::ParameterSerializer<Parameter_t>::PARAMETER_KEY_SIZE;
Expand Down
8 changes: 7 additions & 1 deletion src/cpp/rtps/messages/submessages/DataMsg.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,8 @@ struct DataMsgUtils
(nullptr != inlineQos) ||
((WITH_KEY == topicKind) &&
(!change->writerGUID.is_builtin() || expectsInlineQos || change->kind != ALIVE)) ||
(change->write_params.related_sample_identity() != SampleIdentity::unknown());
(change->write_params.related_sample_identity() != SampleIdentity::unknown()) ||
(change->write_params.has_more_replies());

dataFlag = ALIVE == change->kind &&
change->serializedPayload.length > 0 && nullptr != change->serializedPayload.data;
Expand Down Expand Up @@ -134,6 +135,11 @@ struct DataMsgUtils
change->write_params.related_sample_identity());
}

if (change->write_params.has_more_replies())
{
fastdds::dds::ParameterSerializer<fastdds::dds::Parameter_t>::add_parameter_more_replies(msg);
}

if (WITH_KEY == topicKind && (!change->writerGUID.is_builtin() || expectsInlineQos || ALIVE != change->kind))
{
fastdds::dds::ParameterSerializer<fastdds::dds::Parameter_t>::add_parameter_key(msg,
Expand Down
121 changes: 121 additions & 0 deletions test/blackbox/common/DDSBlackboxTestsBasic.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -512,6 +512,127 @@ TEST(DDSBasic, PidRelatedSampleIdentity)
ASSERT_EQ(related_sample_identity_, info.related_sample_identity);
}

/**
* Read a parameterList from a CDRMessage.
* Search for PID_RPC_MORE_REPLIES in msg.
* @param [in] msg Reference to the message.
* @param [out] exists_pid_rpc_more_replies True if the parameter is inside msg.
* @return true if parsing was correct, false otherwise.
*/
bool check_rpc_has_more_replies_field(
fastdds::rtps::CDRMessage_t& msg,
bool& exists_pid_rpc_more_replies)
{
uint32_t qos_size = 0;

uint32_t original_pos = msg.pos;
bool is_sentinel = false;
while (!is_sentinel)
{
msg.pos = original_pos + qos_size;

ParameterId_t pid{PID_SENTINEL};
bool valid = true;
pid = (ParameterId_t)eprosima::fastdds::helpers::cdr_parse_u16(
(char*)&msg.buffer[msg.pos]);
msg.pos += 2;
uint16_t plength = eprosima::fastdds::helpers::cdr_parse_u16(
(char*)&msg.buffer[msg.pos]);
msg.pos += 2;

if (pid == PID_SENTINEL)
{
// PID_SENTINEL is always considered of length 0
plength = 0;
is_sentinel = true;
}

qos_size += (4 + plength);

// Align to 4 byte boundary and prepare for next iteration
qos_size = (qos_size + 3) & ~3;

if (!valid || ((msg.pos + plength) > msg.length))
{
return false;
}
else if (!is_sentinel)
{
if (PID_RPC_MORE_REPLIES == pid)
{
exists_pid_rpc_more_replies = true;
}
}
}
return true;
}

/**
* This test checks that PID_RPC_MORE_REPLIES is being sent as parameter when has_more_replies is
* specified in WriteParams, and that it is correctly passed into SampleInfo.
*/
TEST(DDSBasic, PidRpcMoreReplies)
{
PubSubWriter<HelloWorldPubSubType> reliable_writer(TEST_TOPIC_NAME);
PubSubReader<HelloWorldPubSubType> reliable_reader(TEST_TOPIC_NAME);

// Test transport will be used in order to filter inlineQoS
auto test_transport = std::make_shared<eprosima::fastdds::rtps::test_UDPv4TransportDescriptor>();
bool exists_pid_rpc_more_replies = false;

test_transport->drop_data_messages_filter_ =
[&exists_pid_rpc_more_replies]
(eprosima::fastdds::rtps::CDRMessage_t& msg)-> bool
{
bool ret = check_rpc_has_more_replies_field(msg, exists_pid_rpc_more_replies);
EXPECT_TRUE(ret);
return false;
};

reliable_writer.reliability(eprosima::fastdds::dds::RELIABLE_RELIABILITY_QOS)
.disable_builtin_transport()
.add_user_transport_to_pparams(test_transport)
.init();
ASSERT_TRUE(reliable_writer.isInitialized());

reliable_reader.reliability(eprosima::fastdds::dds::RELIABLE_RELIABILITY_QOS)
.disable_builtin_transport()
.add_user_transport_to_pparams(test_transport)
.init();
ASSERT_TRUE(reliable_reader.isInitialized());

reliable_writer.wait_discovery();
reliable_reader.wait_discovery();

DataWriter& native_writer = reliable_writer.get_native_writer();

HelloWorld data;
// Send reply associating it with the client request.
eprosima::fastdds::rtps::WriteParams write_params;
write_params.has_more_replies(true);

// Publish the value with has_more_replies set to true
ReturnCode_t write_ret = native_writer.write((void*)&data, write_params);
ASSERT_EQ(RETCODE_OK, write_ret);

DataReader& native_reader = reliable_reader.get_native_reader();

HelloWorld read_data;
eprosima::fastdds::dds::SampleInfo info;
eprosima::fastdds::dds::Duration_t timeout;
timeout.seconds = 2;
while (!native_reader.wait_for_unread_message(timeout))
{
}

ASSERT_EQ(RETCODE_OK,
native_reader.take_next_sample((void*)&read_data, &info));

ASSERT_TRUE(exists_pid_rpc_more_replies);

ASSERT_TRUE(info.has_more_replies);
}

/**
* This test checks that PID_RELATED_SAMPLE_IDENTITY and
* PID_CUSTOM_RELATED_SAMPLE_IDENTITY are being sent as parameter,
Expand Down
1 change: 1 addition & 0 deletions versions.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ Forthcoming
* Implemented Requester and Replier matching algorithm.
* Process key-only payloads:
* New `is_serialized_key` attribute in `SerializedPayload_t`.
* Add field `has_more_replies` to `WriteParams` and `SampleInfo`

Version v3.2.2
--------------
Expand Down