Skip to content
Merged
Show file tree
Hide file tree
Changes from 18 commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
95811f6
Refs #22901. Add `is_serialized_key` to `SerializedPayload_t`.
MiguelCompany Mar 4, 2025
19ad49d
Refs #22901. Fill `is_serialized_key` in `MessageReceiver`.
MiguelCompany Mar 4, 2025
22f0acd
Refs #22901. Handle `is_serialized_key` in `SerializedPayload` methods.
MiguelCompany Mar 4, 2025
32f0114
Refs #22901. Never filter out key-only payloads
MiguelCompany Mar 4, 2025
c3b4004
Refs #22901. Handle `is_serialized_key` on CacheChange
MiguelCompany Mar 4, 2025
251bf9f
Refs #22901. Handle `is_serialized_key` on payload pools.
MiguelCompany Mar 4, 2025
4b8edd3
Refs #22901. Handle `is_serialized_key` in DDSFilterExpression.
MiguelCompany Jun 19, 2025
c29c6c4
Refs #22901. Remove TODO.
MiguelCompany Jun 19, 2025
0de8af2
Refs #22901. Add BB test.
MiguelCompany Jun 24, 2025
1f41964
Refs #22901. Handle key computation failure in DataReaderHistory.
MiguelCompany Jun 25, 2025
d5dec35
Refs #22901. Fix encapsulation header in test.
MiguelCompany Jun 30, 2025
bc7ac3b
Refs #22901. Fix comment in test.
MiguelCompany Jun 30, 2025
6b9828e
Refs #22901. Use CV in test.
MiguelCompany Jun 30, 2025
0417cf4
Refs #22901. Extend test with fragmented data.
MiguelCompany Jun 30, 2025
a0a2c32
Refs #22901. Extend fragment management unit test.
MiguelCompany Jun 30, 2025
e8c4da5
Refs #22901. Add unit test in DDSSQLFilterValueTests.
MiguelCompany Jun 30, 2025
c5a5016
Refs #22901. Bonus point: fix change kind in data frag processing.
MiguelCompany Jun 30, 2025
7418f24
Refs #22901. Add feature to versions.md.
MiguelCompany Jul 1, 2025
753676f
Apply suggestion
MiguelCompany Jul 1, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions include/fastdds/rtps/common/CacheChange.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ struct FASTDDS_EXPORTED_API CacheChange_t

// Copy certain values from serializedPayload
serializedPayload.encapsulation = ch_ptr->serializedPayload.encapsulation;
serializedPayload.is_serialized_key = ch_ptr->serializedPayload.is_serialized_key;

// Copy fragment size and calculate fragment count
setFragmentSize(ch_ptr->fragment_size_, false);
Expand Down Expand Up @@ -287,6 +288,12 @@ struct FASTDDS_EXPORTED_API CacheChange_t
uint32_t incoming_length = fragment_size_ * fragments_in_submessage;
uint32_t last_fragment_index = fragment_starting_num + fragments_in_submessage - 1;

// Validate payload types
if (serializedPayload.is_serialized_key != incoming_data.is_serialized_key)
{
return false;
}

// Validate fragment indexes
if (last_fragment_index > fragment_count_)
{
Expand Down
2 changes: 2 additions & 0 deletions include/fastdds/rtps/common/SerializedPayload.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ struct FASTDDS_EXPORTED_API SerializedPayload_t
uint32_t pos;
//!Pool that created the payload
IPayloadPool* payload_owner = nullptr;
//!Whether the payload contains a serialized key, or the whole data
bool is_serialized_key = false;

//!Default constructor
SerializedPayload_t()
Expand Down
20 changes: 16 additions & 4 deletions src/cpp/fastdds/subscriber/history/DataReaderHistory.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -215,8 +215,14 @@ bool DataReaderHistory::received_change_keep_all(
{
if (!compute_key_for_change_fn_(a_change))
{
// Store the sample temporally only in ReaderHistory. When completed it will be stored in DataReaderHistory too.
return add_to_reader_history_if_not_full(a_change, rejection_reason);
if (!a_change->is_fully_assembled())
{
// Store the sample temporally only in ReaderHistory. When completed it will be stored in SubscriberHistory too.
return add_to_reader_history_if_not_full(a_change, rejection_reason);
}

rejection_reason = REJECTED_BY_INSTANCES_LIMIT;
return false;
}

bool ret_value = false;
Expand Down Expand Up @@ -251,8 +257,14 @@ bool DataReaderHistory::received_change_keep_last(
{
if (!compute_key_for_change_fn_(a_change))
{
// Store the sample temporally only in ReaderHistory. When completed it will be stored in SubscriberHistory too.
return add_to_reader_history_if_not_full(a_change, rejection_reason);
if (!a_change->is_fully_assembled())
{
// Store the sample temporally only in ReaderHistory. When completed it will be stored in SubscriberHistory too.
return add_to_reader_history_if_not_full(a_change, rejection_reason);
}

rejection_reason = REJECTED_BY_INSTANCES_LIMIT;
return false;
}

bool ret_value = false;
Expand Down
6 changes: 6 additions & 0 deletions src/cpp/fastdds/topic/DDSSQLFilter/DDSFilterExpression.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,12 @@ bool DDSFilterExpression::evaluate(
using namespace eprosima::fastdds::dds::xtypes;
using namespace eprosima::fastcdr;

// Always pass filter for key-only payloads
if (payload.is_serialized_key)
{
return true;
}

dyn_data_->clear_all_values();
try
{
Expand Down
1 change: 1 addition & 0 deletions src/cpp/rtps/DataSharing/DataSharingPayloadPool.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ bool DataSharingPayloadPool::release_payload(
payload.pos = 0;
payload.max_size = 0;
payload.data = nullptr;
payload.is_serialized_key = false;
payload.payload_owner = nullptr;
return true;
}
Expand Down
1 change: 1 addition & 0 deletions src/cpp/rtps/DataSharing/ReaderPool.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ class ReaderPool : public DataSharingPayloadPool
payload.data = data.data;
payload.length = data.length;
payload.max_size = data.length;
payload.is_serialized_key = data.is_serialized_key;
payload.payload_owner = this;
return true;
}
Expand Down
1 change: 1 addition & 0 deletions src/cpp/rtps/DataSharing/WriterPool.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@ class WriterPool : public DataSharingPayloadPool
payload.data = data.data;
payload.length = data.length;
payload.max_size = data.length;
payload.is_serialized_key = data.is_serialized_key;
payload.payload_owner = this;
return true;
}
Expand Down
6 changes: 6 additions & 0 deletions src/cpp/rtps/common/SerializedPayload.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -36,13 +36,15 @@ SerializedPayload_t& SerializedPayload_t::operator = (
max_size = other.max_size;
pos = other.pos;
payload_owner = other.payload_owner;
is_serialized_key = other.is_serialized_key;

other.encapsulation = CDR_BE;
other.length = 0;
other.data = nullptr;
other.max_size = 0;
other.pos = 0;
other.payload_owner = nullptr;
other.is_serialized_key = false;

return *this;
}
Expand All @@ -60,6 +62,7 @@ bool SerializedPayload_t::operator == (
const SerializedPayload_t& other) const
{
return ((encapsulation == other.encapsulation) &&
(is_serialized_key == other.is_serialized_key) &&
(length == other.length) &&
(0 == memcmp(data, other.data, length)));
}
Expand All @@ -82,6 +85,7 @@ bool SerializedPayload_t::copy(
}
}
encapsulation = serData->encapsulation;
is_serialized_key = serData->is_serialized_key;
if (length == 0)
{
return true;
Expand All @@ -96,6 +100,7 @@ bool SerializedPayload_t::reserve_fragmented(
length = serData->length;
max_size = serData->length;
encapsulation = serData->encapsulation;
is_serialized_key = serData->is_serialized_key;
data = (octet*)calloc(length, sizeof(octet));
return true;
}
Expand All @@ -112,6 +117,7 @@ void SerializedPayload_t::empty()
free(data);
}
data = nullptr;
is_serialized_key = false;
}

void SerializedPayload_t::reserve(
Expand Down
2 changes: 2 additions & 0 deletions src/cpp/rtps/history/TopicPayloadPool.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ bool TopicPayloadPool::get_payload(
payload.data = data.data;
payload.length = data.length;
payload.max_size = PayloadNode::data_size(data.data);
payload.is_serialized_key = data.is_serialized_key;
payload.payload_owner = this;
return true;
}
Expand Down Expand Up @@ -136,6 +137,7 @@ bool TopicPayloadPool::release_payload(
payload.pos = 0;
payload.max_size = 0;
payload.data = nullptr;
payload.is_serialized_key = false;
payload.payload_owner = nullptr;
return true;
}
Expand Down
82 changes: 20 additions & 62 deletions src/cpp/rtps/messages/MessageReceiver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -867,28 +867,10 @@ bool MessageReceiver::proc_Submsg_Data(
uint32_t next_pos = msg->pos + payload_size;
if (msg->length >= next_pos && payload_size > 0)
{
FASTDDS_TODO_BEFORE(3, 3, "Pass keyFlag in serializedPayload, and always pass input data upwards");
if (dataFlag)
{
ch.serializedPayload.data = &msg->buffer[msg->pos];
ch.serializedPayload.length = payload_size;
ch.serializedPayload.max_size = payload_size;
}
else // keyFlag would be true since we are inside an if (dataFlag || keyFlag)
{
if (payload_size <= PARAMETER_KEY_HASH_LENGTH)
{
if (!ch.instanceHandle.isDefined())
{
memcpy(ch.instanceHandle.value, &msg->buffer[msg->pos], payload_size);
}
}
else
{
EPROSIMA_LOG_WARNING(RTPS_MSG_IN, IDSTRING "Ignoring Serialized Payload for too large key-only data (" <<
payload_size << ")");
}
}
ch.serializedPayload.data = &msg->buffer[msg->pos];
ch.serializedPayload.length = payload_size;
ch.serializedPayload.max_size = payload_size;
ch.serializedPayload.is_serialized_key = keyFlag;
msg->pos = next_pos;
}
else
Expand Down Expand Up @@ -977,6 +959,7 @@ bool MessageReceiver::proc_Submsg_DataFrag(
//FOUND THE READER.
//We ask the reader for a cachechange to store the information.
CacheChange_t ch;
ch.kind = ALIVE;
ch.writerGUID.guidPrefix = source_guid_prefix_;
valid &= CDRMessage::readEntityId(msg, &ch.writerGUID.entityId);

Expand Down Expand Up @@ -1046,49 +1029,24 @@ bool MessageReceiver::proc_Submsg_DataFrag(

// Validations??? XXX TODO

if (!keyFlag)
uint32_t next_pos = msg->pos + payload_size;
if (msg->length >= next_pos && payload_size > 0)
{
uint32_t next_pos = msg->pos + payload_size;
if (msg->length >= next_pos && payload_size > 0)
{
ch.kind = ALIVE;
ch.serializedPayload.data = &msg->buffer[msg->pos];
ch.serializedPayload.length = payload_size;
ch.serializedPayload.max_size = payload_size;
ch.setFragmentSize(fragmentSize);
ch.serializedPayload.data = &msg->buffer[msg->pos];
ch.serializedPayload.length = payload_size;
ch.serializedPayload.max_size = payload_size;
ch.serializedPayload.is_serialized_key = keyFlag;
ch.setFragmentSize(fragmentSize);

msg->pos = next_pos;
}
else
{
EPROSIMA_LOG_WARNING(RTPS_MSG_IN, IDSTRING "Serialized Payload value invalid or larger than maximum allowed size "
"(" << payload_size << "/" << (msg->length - msg->pos) << ")");
ch.serializedPayload.data = nullptr;
ch.inline_qos.data = nullptr;
return false;
}
msg->pos = next_pos;
}
else if (keyFlag)
{
/* XXX TODO
Endianness_t previous_endian = msg->msg_endian;
if (ch->serializedPayload.encapsulation == PL_CDR_BE)
msg->msg_endian = BIGEND;
else if (ch->serializedPayload.encapsulation == PL_CDR_LE)
msg->msg_endian = LITTLEEND;
else
{
EPROSIMA_LOG_ERROR(RTPS_MSG_IN, IDSTRING"Bad encapsulation for KeyHash and status parameter list");
return false;
}
//uint32_t param_size;
if (ParameterList::readParameterListfromCDRMsg(msg, &m_ParamList, ch, false) <= 0)
{
EPROSIMA_LOG_INFO(RTPS_MSG_IN, IDSTRING"SubMessage Data ERROR, keyFlag ParameterList");
return false;
}
msg->msg_endian = previous_endian;
*/
else
{
EPROSIMA_LOG_WARNING(RTPS_MSG_IN, IDSTRING "Serialized Payload value invalid or larger than maximum allowed size "
"(" << payload_size << "/" << (msg->length - msg->pos) << ")");
ch.serializedPayload.data = nullptr;
ch.inline_qos.data = nullptr;
return false;
}

// Set sourcetimestamp
Expand Down
6 changes: 6 additions & 0 deletions src/cpp/rtps/reader/reader_utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,12 @@ bool change_is_relevant_for_filter(
{
bool ret = true;

// If the change contains only a serialized key, it is always relevant
if (change.serializedPayload.is_serialized_key)
{
return true;
}

// If the change has no payload, it should have an instanceHandle.
// This is only allowed for UNREGISTERED and DISPOSED changes, where the instanceHandle is used to identify the
// instance to unregister or dispose.
Expand Down
Loading