From 6f44d2cc9ab65879dc4ad94a109e295067964d8e Mon Sep 17 00:00:00 2001 From: Thomas Moore Date: Fri, 28 Jun 2024 17:04:23 -0400 Subject: [PATCH 1/3] Populate message info for intra-process messages (cherry picked from commit 8f255f3da12d5d0cc8aa5a49f69bb74339226820) (cherry picked from commit 94a3d5856c046e3a7b6bda942394cd8662926c28) Signed-off-by: Thomas Moore --- .../buffers/intra_process_buffer.hpp | 243 +++++------------ .../create_intra_process_buffer.hpp | 54 ++-- .../experimental/intra_process_manager.hpp | 248 +++++++++++++----- .../ros_message_intra_process_buffer.hpp | 7 +- .../subscription_intra_process.hpp | 69 ++--- .../subscription_intra_process_buffer.hpp | 60 ++--- rclcpp/include/rclcpp/publisher.hpp | 56 ++-- rclcpp/src/rclcpp/intra_process_manager.cpp | 13 +- 8 files changed, 370 insertions(+), 380 deletions(-) diff --git a/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp index 268c3f6649..9e27c1481d 100644 --- a/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp @@ -19,11 +19,13 @@ #include #include #include +#include #include #include "rclcpp/allocator/allocator_common.hpp" #include "rclcpp/allocator/allocator_deleter.hpp" #include "rclcpp/experimental/buffers/buffer_implementation_base.hpp" +#include "rclcpp/intra_process_buffer_type.hpp" #include "rclcpp/macros.hpp" #include "tracetools/tracetools.h" @@ -34,6 +36,18 @@ namespace experimental namespace buffers { +template< + typename MessageT, + typename MessageDeleter = std::default_delete> +struct IntraProcessBufferNode +{ + using MessageUniquePtr = std::unique_ptr; + using MessageSharedPtr = std::shared_ptr; + + std::variant message; + rmw_message_info_t message_info; +}; + class IntraProcessBufferBase { public: @@ -44,13 +58,12 @@ class IntraProcessBufferBase virtual void clear() = 0; virtual bool has_data() const = 0; - virtual bool use_take_shared_method() const = 0; + virtual IntraProcessBufferType buffer_type() const = 0; virtual size_t available_capacity() const = 0; }; template< typename MessageT, - typename Alloc = std::allocator, typename MessageDeleter = std::default_delete> class IntraProcessBuffer : public IntraProcessBufferBase { @@ -59,47 +72,37 @@ class IntraProcessBuffer : public IntraProcessBufferBase virtual ~IntraProcessBuffer() {} - using MessageUniquePtr = std::unique_ptr; - using MessageSharedPtr = std::shared_ptr; + using Node = IntraProcessBufferNode; - virtual void add_shared(MessageSharedPtr msg) = 0; - virtual void add_unique(MessageUniquePtr msg) = 0; + virtual void add(Node node) = 0; - virtual MessageSharedPtr consume_shared() = 0; - virtual MessageUniquePtr consume_unique() = 0; + virtual Node consume() = 0; - virtual std::vector get_all_data_shared() = 0; - virtual std::vector get_all_data_unique() = 0; + virtual std::vector get_all_data() = 0; }; template< typename MessageT, typename Alloc = std::allocator, typename MessageDeleter = std::default_delete, - typename BufferT = std::unique_ptr> -class TypedIntraProcessBuffer : public IntraProcessBuffer + IntraProcessBufferType BufferType = IntraProcessBufferType::CallbackDefault> +class TypedIntraProcessBuffer : public IntraProcessBuffer { public: RCLCPP_SMART_PTR_DEFINITIONS(TypedIntraProcessBuffer) using MessageAllocTraits = allocator::AllocRebind; using MessageAlloc = typename MessageAllocTraits::allocator_type; - using MessageUniquePtr = std::unique_ptr; - using MessageSharedPtr = std::shared_ptr; + using Node = IntraProcessBufferNode; + using MessageSharedPtr = typename Node::MessageSharedPtr; + using MessageUniquePtr = typename Node::MessageUniquePtr; explicit TypedIntraProcessBuffer( - std::unique_ptr> buffer_impl, + std::unique_ptr> buffer_impl, std::shared_ptr allocator = nullptr) + : buffer_(std::move(buffer_impl)) { - bool valid_type = (std::is_same::value || - std::is_same::value); - if (!valid_type) { - throw std::runtime_error("Creating TypedIntraProcessBuffer with not valid BufferT"); - } - - buffer_ = std::move(buffer_impl); - TRACETOOLS_TRACEPOINT( rclcpp_buffer_to_ipb, static_cast(buffer_.get()), @@ -113,34 +116,19 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer(std::move(msg)); - } - - void add_unique(MessageUniquePtr msg) override - { - buffer_->enqueue(std::move(msg)); - } - - MessageSharedPtr consume_shared() override - { - return consume_shared_impl(); - } - - MessageUniquePtr consume_unique() override + void add(Node node) override { - return consume_unique_impl(); + add_impl(std::move(node)); } - std::vector get_all_data_shared() override + Node consume() override { - return get_all_data_shared_impl(); + return buffer_->dequeue(); } - std::vector get_all_data_unique() override + std::vector get_all_data() override { - return get_all_data_unique_impl(); + return buffer_->get_all_data(); } bool has_data() const override @@ -153,9 +141,9 @@ class TypedIntraProcessBuffer : public IntraProcessBufferclear(); } - bool use_take_shared_method() const override + IntraProcessBufferType buffer_type() const override { - return std::is_same::value; + return BufferType; } size_t available_capacity() const override @@ -164,163 +152,52 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer> buffer_; + std::unique_ptr> buffer_; std::shared_ptr message_allocator_; - // MessageSharedPtr to MessageSharedPtr - template - typename std::enable_if< - std::is_same::value - >::type - add_shared_impl(MessageSharedPtr shared_msg) + template + typename std::enable_if_t< + BufferT == IntraProcessBufferType::CallbackDefault> + add_impl(Node node) { - buffer_->enqueue(std::move(shared_msg)); + buffer_->enqueue(std::move(node)); } - // MessageSharedPtr to MessageUniquePtr - template - typename std::enable_if< - std::is_same::value - >::type - add_shared_impl(MessageSharedPtr shared_msg) + template + typename std::enable_if_t< + BufferT == IntraProcessBufferType::SharedPtr> + add_impl(Node node) { - // This should not happen: here a copy is unconditionally made, while the intra-process manager - // can decide whether a copy is needed depending on the number and the type of buffers - - MessageUniquePtr unique_msg; - MessageDeleter * deleter = std::get_deleter(shared_msg); - auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1); - MessageAllocTraits::construct(*message_allocator_.get(), ptr, *shared_msg); - if (deleter) { - unique_msg = MessageUniquePtr(ptr, *deleter); + if (std::holds_alternative(node.message)) { + buffer_->enqueue(std::move(node)); } else { - unique_msg = MessageUniquePtr(ptr); + // Promote to a shared pointer + auto unique_msg = std::move(std::get(node.message)); + node.message = MessageSharedPtr(unique_msg.release()); + buffer_->enqueue(std::move(node)); } - - buffer_->enqueue(std::move(unique_msg)); } - // MessageSharedPtr to MessageSharedPtr - template - typename std::enable_if< - std::is_same::value, - MessageSharedPtr - >::type - consume_shared_impl() + template + typename std::enable_if_t< + BufferT == IntraProcessBufferType::UniquePtr> + add_impl(Node node) { - return buffer_->dequeue(); - } - - // MessageUniquePtr to MessageSharedPtr - template - typename std::enable_if< - (std::is_same::value), - MessageSharedPtr - >::type - consume_shared_impl() - { - // automatic cast from unique ptr to shared ptr - return buffer_->dequeue(); - } - - // MessageSharedPtr to MessageUniquePtr - template - typename std::enable_if< - (std::is_same::value), - MessageUniquePtr - >::type - consume_unique_impl() - { - MessageSharedPtr buffer_msg = buffer_->dequeue(); - - MessageUniquePtr unique_msg; - MessageDeleter * deleter = std::get_deleter(buffer_msg); - auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1); - MessageAllocTraits::construct(*message_allocator_.get(), ptr, *buffer_msg); - if (deleter) { - unique_msg = MessageUniquePtr(ptr, *deleter); + if (std::holds_alternative(node.message)) { + buffer_->enqueue(std::move(node)); } else { - unique_msg = MessageUniquePtr(ptr); - } - - return unique_msg; - } - - // MessageUniquePtr to MessageUniquePtr - template - typename std::enable_if< - (std::is_same::value), - MessageUniquePtr - >::type - consume_unique_impl() - { - return buffer_->dequeue(); - } - - // MessageSharedPtr to MessageSharedPtr - template - typename std::enable_if< - std::is_same::value, - std::vector - >::type - get_all_data_shared_impl() - { - return buffer_->get_all_data(); - } - - // MessageUniquePtr to MessageSharedPtr - template - typename std::enable_if< - std::is_same::value, - std::vector - >::type - get_all_data_shared_impl() - { - std::vector result; - auto uni_ptr_vec = buffer_->get_all_data(); - result.reserve(uni_ptr_vec.size()); - for (MessageUniquePtr & uni_ptr : uni_ptr_vec) { - result.emplace_back(std::move(uni_ptr)); - } - return result; - } - - // MessageSharedPtr to MessageUniquePtr - template - typename std::enable_if< - std::is_same::value, - std::vector - >::type - get_all_data_unique_impl() - { - std::vector result; - auto shared_ptr_vec = buffer_->get_all_data(); - result.reserve(shared_ptr_vec.size()); - for (MessageSharedPtr shared_msg : shared_ptr_vec) { - MessageUniquePtr unique_msg; + auto shared_msg = std::move(std::get(node.message)); MessageDeleter * deleter = std::get_deleter(shared_msg); auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1); MessageAllocTraits::construct(*message_allocator_.get(), ptr, *shared_msg); if (deleter) { - unique_msg = MessageUniquePtr(ptr, *deleter); + node.message = MessageUniquePtr(ptr, *deleter); } else { - unique_msg = MessageUniquePtr(ptr); + node.message = MessageUniquePtr(ptr); } - result.push_back(std::move(unique_msg)); + buffer_->enqueue(std::move(node)); } - return result; - } - - // MessageUniquePtr to MessageUniquePtr - template - typename std::enable_if< - std::is_same::value, - std::vector - >::type - get_all_data_unique_impl() - { - return buffer_->get_all_data(); } }; diff --git a/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp index f292c22448..40712a0b7e 100644 --- a/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp @@ -16,7 +16,6 @@ #define RCLCPP__EXPERIMENTAL__CREATE_INTRA_PROCESS_BUFFER_HPP_ #include -#include #include #include "rclcpp/experimental/buffers/intra_process_buffer.hpp" @@ -33,58 +32,59 @@ template< typename MessageT, typename Alloc = std::allocator, typename Deleter = std::default_delete> -typename rclcpp::experimental::buffers::IntraProcessBuffer::UniquePtr +typename rclcpp::experimental::buffers::IntraProcessBuffer::UniquePtr create_intra_process_buffer( IntraProcessBufferType buffer_type, const rclcpp::QoS & qos, std::shared_ptr allocator) { - using MessageSharedPtr = std::shared_ptr; - using MessageUniquePtr = std::unique_ptr; - size_t buffer_size = qos.depth(); using rclcpp::experimental::buffers::IntraProcessBuffer; - typename IntraProcessBuffer::UniquePtr buffer; + using rclcpp::experimental::buffers::IntraProcessBufferNode; + using rclcpp::experimental::buffers::TypedIntraProcessBuffer; + using rclcpp::experimental::buffers::RingBufferImplementation; + + using BufferT = IntraProcessBufferNode; + using BufferImplT = RingBufferImplementation; + + auto buffer_implementation = std::make_unique(buffer_size); + + typename IntraProcessBuffer::UniquePtr buffer; switch (buffer_type) { case IntraProcessBufferType::SharedPtr: { - using BufferT = MessageSharedPtr; - - auto buffer_implementation = - std::make_unique>( - buffer_size); + using IntraProcessBufferT = TypedIntraProcessBuffer< + MessageT, Alloc, Deleter, IntraProcessBufferType::SharedPtr>; // Construct the intra_process_buffer - buffer = - std::make_unique>( - std::move(buffer_implementation), - allocator); + buffer = std::make_unique( + std::move(buffer_implementation), allocator); break; } case IntraProcessBufferType::UniquePtr: { - using BufferT = MessageUniquePtr; - - auto buffer_implementation = - std::make_unique>( - buffer_size); + using IntraProcessBufferT = TypedIntraProcessBuffer< + MessageT, Alloc, Deleter, IntraProcessBufferType::UniquePtr>; // Construct the intra_process_buffer - buffer = - std::make_unique>( - std::move(buffer_implementation), - allocator); + buffer = std::make_unique( + std::move(buffer_implementation), allocator); break; } case IntraProcessBufferType::CallbackDefault: { - throw std::runtime_error("IntraProcessBufferType::CallbackDefault is not allowed"); + using IntraProcessBufferT = TypedIntraProcessBuffer< + MessageT, Alloc, Deleter, IntraProcessBufferType::CallbackDefault>; + + // Construct the intra_process_buffer + buffer = std::make_unique( + std::move(buffer_implementation), allocator); + + break; } } diff --git a/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp b/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp index a8eba4baf7..e0f7dd3cbc 100644 --- a/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp +++ b/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp @@ -125,11 +125,11 @@ class IntraProcessManager uint64_t sub_id = IntraProcessManager::get_next_unique_id(); - subscriptions_[sub_id] = subscription; + subscriptions_[sub_id] = {subscription}; // adds the subscription id to all the matchable publishers for (auto & pair : publishers_) { - auto publisher = pair.second.lock(); + auto publisher = pair.second.weak_publisher.lock(); if (!publisher) { continue; } @@ -170,7 +170,7 @@ class IntraProcessManager * In addition this generates a unique intra process id for the publisher. * * \param publisher publisher to be registered with the manager. - * \param buffer publisher's buffer to be stored if its duability is transient local. + * \param buffer publisher's buffer to be stored if its durability is transient local. * \return an unsigned 64-bit integer which is the publisher's unique id. */ RCLCPP_PUBLIC @@ -224,57 +224,92 @@ class IntraProcessManager do_intra_process_publish( uint64_t intra_process_publisher_id, std::unique_ptr message, - typename allocator::AllocRebind::allocator_type & allocator) + typename allocator::AllocRebind::allocator_type & allocator, + rmw_message_info_t * message_info_out = nullptr) { using MessageAllocTraits = allocator::AllocRebind; using MessageAllocatorT = typename MessageAllocTraits::allocator_type; std::shared_lock lock(mutex_); - auto publisher_it = pub_to_subs_.find(intra_process_publisher_id); - if (publisher_it == pub_to_subs_.end()) { + auto pub_to_subs_it = pub_to_subs_.find(intra_process_publisher_id); + if (pub_to_subs_it == pub_to_subs_.end()) { // Publisher is either invalid or no longer exists. RCLCPP_WARN( rclcpp::get_logger("rclcpp"), "Calling do_intra_process_publish for invalid or no longer existing publisher id"); return; } - const auto & sub_ids = publisher_it->second; - if (sub_ids.take_ownership_subscriptions.empty()) { + auto publisher_it = publishers_.find(intra_process_publisher_id); + if (publisher_it == publishers_.end()) { + throw std::runtime_error("publisher has unexpectedly gone out of scope"); + } + auto publisher = publisher_it->second.weak_publisher.lock(); + if (!publisher) { + throw std::runtime_error("publisher has unexpectedly gone out of scope"); + } + + rmw_message_info_t message_info{}; + message_info.from_intra_process = true; + message_info.publication_sequence_number = publisher_it->second.publication_sequence_number++; + message_info.publisher_gid = publisher->get_gid(); + + rcutils_time_point_value_t now; + if (rcutils_system_time_now(&now) == RCUTILS_RET_OK) { + message_info.source_timestamp = now; + message_info.received_timestamp = now; + } + + const auto & take_ownership_subscriptions = pub_to_subs_it->second.take_ownership_subscriptions; + const auto & take_shared_subscriptions = pub_to_subs_it->second.take_shared_subscriptions; + + if (take_ownership_subscriptions.empty()) { // None of the buffers require ownership, so we promote the pointer std::shared_ptr msg = std::move(message); this->template add_shared_msg_to_buffers( - msg, sub_ids.take_shared_subscriptions); - } else if (!sub_ids.take_ownership_subscriptions.empty() && // NOLINT - sub_ids.take_shared_subscriptions.size() <= 1) + msg, + message_info, + take_shared_subscriptions); + } else if (!take_ownership_subscriptions.empty() && // NOLINT + take_shared_subscriptions.size() <= 1) { // There is at maximum 1 buffer that does not require ownership. // So this case is equivalent to all the buffers requiring ownership // Merge the two vector of ids into a unique one std::vector concatenated_vector( - sub_ids.take_shared_subscriptions.begin(), sub_ids.take_shared_subscriptions.end()); + take_shared_subscriptions.begin(), take_shared_subscriptions.end()); concatenated_vector.insert( concatenated_vector.end(), - sub_ids.take_ownership_subscriptions.begin(), - sub_ids.take_ownership_subscriptions.end()); + take_ownership_subscriptions.begin(), + take_ownership_subscriptions.end()); this->template add_owned_msg_to_buffers( std::move(message), + message_info, concatenated_vector, allocator); - } else if (!sub_ids.take_ownership_subscriptions.empty() && // NOLINT - sub_ids.take_shared_subscriptions.size() > 1) + } else if (!take_ownership_subscriptions.empty() && // NOLINT + take_shared_subscriptions.size() > 1) { // Construct a new shared pointer from the message // for the buffers that do not require ownership auto shared_msg = std::allocate_shared(allocator, *message); this->template add_shared_msg_to_buffers( - shared_msg, sub_ids.take_shared_subscriptions); + shared_msg, + message_info, + take_shared_subscriptions); this->template add_owned_msg_to_buffers( - std::move(message), sub_ids.take_ownership_subscriptions, allocator); + std::move(message), + message_info, + take_ownership_subscriptions, + allocator); + } + + if (message_info_out) { + *message_info_out = message_info; } } @@ -284,7 +319,7 @@ class IntraProcessManager typename Alloc, typename Deleter = std::default_delete > - std::shared_ptr + std::pair, rmw_message_info_t> do_intra_process_publish_and_return_shared( uint64_t intra_process_publisher_id, std::unique_ptr message, @@ -295,41 +330,67 @@ class IntraProcessManager std::shared_lock lock(mutex_); - auto publisher_it = pub_to_subs_.find(intra_process_publisher_id); - if (publisher_it == pub_to_subs_.end()) { + auto pub_to_subs_it = pub_to_subs_.find(intra_process_publisher_id); + if (pub_to_subs_it == pub_to_subs_.end()) { // Publisher is either invalid or no longer exists. RCLCPP_WARN( rclcpp::get_logger("rclcpp"), "Calling do_intra_process_publish for invalid or no longer existing publisher id"); - return nullptr; + return {}; + } + + auto publisher_it = publishers_.find(intra_process_publisher_id); + if (publisher_it == publishers_.end()) { + throw std::runtime_error("publisher has unexpectedly gone out of scope"); + } + auto publisher = publisher_it->second.weak_publisher.lock(); + if (!publisher) { + throw std::runtime_error("publisher has unexpectedly gone out of scope"); } - const auto & sub_ids = publisher_it->second; - if (sub_ids.take_ownership_subscriptions.empty()) { + rmw_message_info_t message_info{}; + message_info.from_intra_process = true; + message_info.publication_sequence_number = publisher_it->second.publication_sequence_number++; + message_info.publisher_gid = publisher->get_gid(); + + rcutils_time_point_value_t now; + if (rcutils_system_time_now(&now) == RCUTILS_RET_OK) { + message_info.source_timestamp = now; + message_info.received_timestamp = now; + } + + const auto & take_ownership_subscriptions = pub_to_subs_it->second.take_ownership_subscriptions; + const auto & take_shared_subscriptions = pub_to_subs_it->second.take_shared_subscriptions; + + if (take_ownership_subscriptions.empty()) { // If there are no owning, just convert to shared. std::shared_ptr shared_msg = std::move(message); - if (!sub_ids.take_shared_subscriptions.empty()) { + if (!take_shared_subscriptions.empty()) { this->template add_shared_msg_to_buffers( - shared_msg, sub_ids.take_shared_subscriptions); + shared_msg, + message_info, + take_shared_subscriptions); } - return shared_msg; + return {shared_msg, message_info}; } else { // Construct a new shared pointer from the message for the buffers that // do not require ownership and to return. auto shared_msg = std::allocate_shared(allocator, *message); - if (!sub_ids.take_shared_subscriptions.empty()) { + if (!take_shared_subscriptions.empty()) { this->template add_shared_msg_to_buffers( shared_msg, - sub_ids.take_shared_subscriptions); + message_info, + take_shared_subscriptions); } - if (!sub_ids.take_ownership_subscriptions.empty()) { + if (!take_ownership_subscriptions.empty()) { this->template add_owned_msg_to_buffers( std::move(message), - sub_ids.take_ownership_subscriptions, + message_info, + take_ownership_subscriptions, allocator); } - return shared_msg; + return {shared_msg, message_info}; } } @@ -341,9 +402,11 @@ class IntraProcessManager void add_shared_msg_to_buffer( std::shared_ptr message, + const rmw_message_info_t & message_info, uint64_t subscription_id) { - add_shared_msg_to_buffers(message, {subscription_id}); + add_shared_msg_to_buffers(message, message_info, + {subscription_id}); } template< @@ -354,11 +417,12 @@ class IntraProcessManager void add_owned_msg_to_buffer( std::unique_ptr message, + const rmw_message_info_t & message_info, uint64_t subscription_id, typename allocator::AllocRebind::allocator_type & allocator) { add_owned_msg_to_buffers( - std::move(message), {subscription_id}, allocator); + std::move(message), message_info, {subscription_id}, allocator); } /// Return true if the given rmw_gid_t matches any stored Publishers. @@ -381,6 +445,18 @@ class IntraProcessManager lowest_available_capacity(const uint64_t intra_process_publisher_id) const; private: + struct SubscriptionData + { + rclcpp::experimental::SubscriptionIntraProcessBase::WeakPtr weak_subscription; + uint64_t reception_sequence_number{0}; + }; + + struct PublisherData + { + rclcpp::PublisherBase::WeakPtr weak_publisher; + uint64_t publication_sequence_number{0}; + }; + struct SplittedSubscriptions { std::vector take_shared_subscriptions; @@ -421,10 +497,10 @@ class IntraProcessManager }; using SubscriptionMap = - std::unordered_map; + std::unordered_map; using PublisherMap = - std::unordered_map; + std::unordered_map; using PublisherBufferMap = std::unordered_map; @@ -468,18 +544,18 @@ class IntraProcessManager using ROSMessageTypeAllocatorTraits = allocator::AllocRebind; using ROSMessageTypeAllocator = typename ROSMessageTypeAllocatorTraits::allocator_type; using ROSMessageTypeDeleter = allocator::Deleter; + using IntraProcessBuffer = rclcpp::experimental::buffers::IntraProcessBuffer< + ROSMessageType, + ROSMessageTypeDeleter + >; + using ROSMessageSharedPtr = typename IntraProcessBuffer::Node::MessageSharedPtr; + using ROSMessageUniquePtr = typename IntraProcessBuffer::Node::MessageUniquePtr; auto publisher_buffer = publisher_buffers_[pub_id].lock(); if (!publisher_buffer) { throw std::runtime_error("publisher buffer has unexpectedly gone out of scope"); } - auto buffer = std::dynamic_pointer_cast< - rclcpp::experimental::buffers::IntraProcessBuffer< - ROSMessageType, - ROSMessageTypeAllocator, - ROSMessageTypeDeleter - > - >(publisher_buffer); + auto buffer = std::dynamic_pointer_cast(publisher_buffer); if (!buffer) { throw std::runtime_error( "failed to dynamic cast publisher's IntraProcessBufferBase to " @@ -487,20 +563,52 @@ class IntraProcessManager "ROSMessageTypeDeleter> which can happen when the publisher and " "subscription use different allocator types, which is not supported"); } - if (use_take_shared_method) { - auto data_vec = buffer->get_all_data_shared(); - for (auto shared_data : data_vec) { + auto data_vec = buffer->get_all_data(); + for (auto & data : data_vec) { + // The buffer's own storage (shared vs unique) is independent of what this + // particular new subscription requests, so the variant's actual alternative may + // not match use_take_shared_method: convert as needed. + if (use_take_shared_method) { + ROSMessageSharedPtr shared_ptr; + + std::visit( + [&shared_ptr](auto && message) { + using T = std::decay_t; + if constexpr (std::is_same_v) { + shared_ptr = message; + } else if constexpr (std::is_same_v) { + auto allocator = ROSMessageTypeAllocator(); + shared_ptr = std::allocate_shared( + allocator, *message); + } + }, data.message); + this->template add_shared_msg_to_buffer< ROSMessageType, ROSMessageTypeAllocator, ROSMessageTypeDeleter, ROSMessageType>( - shared_data, sub_id); - } - } else { - auto data_vec = buffer->get_all_data_unique(); - for (auto & owned_data : data_vec) { + shared_ptr, data.message_info, sub_id); + } else { + ROSMessageUniquePtr unique_ptr; auto allocator = ROSMessageTypeAllocator(); + + std::visit( + [&unique_ptr, &allocator](auto && message) { + ROSMessageTypeDeleter deleter; + auto ptr = ROSMessageTypeAllocatorTraits::allocate(allocator, 1); + ROSMessageTypeAllocatorTraits::construct(allocator, ptr, *message); + + using T = std::decay_t; + if constexpr (std::is_same_v) { + allocator::set_allocator_for_deleter(&deleter, &allocator); + } else if constexpr (std::is_same_v) { + deleter = message.get_deleter(); + } + + unique_ptr = ROSMessageUniquePtr(ptr, deleter); + }, data.message); + this->template add_owned_msg_to_buffer< ROSMessageType, ROSMessageTypeAllocator, ROSMessageTypeDeleter, ROSMessageType>( - std::move(owned_data), sub_id, allocator); + std::move(unique_ptr), data.message_info, sub_id, allocator); } } } @@ -513,7 +621,8 @@ class IntraProcessManager void add_shared_msg_to_buffers( std::shared_ptr message, - std::vector subscription_ids) + rmw_message_info_t message_info, + const std::vector & subscription_ids) { using ROSMessageTypeAllocatorTraits = allocator::AllocRebind; using ROSMessageTypeAllocator = typename ROSMessageTypeAllocatorTraits::allocator_type; @@ -529,18 +638,20 @@ class IntraProcessManager if (subscription_it == subscriptions_.end()) { throw std::runtime_error("subscription has unexpectedly gone out of scope"); } - auto subscription_base = subscription_it->second.lock(); + auto subscription_base = subscription_it->second.weak_subscription.lock(); if (subscription_base == nullptr) { subscriptions_.erase(id); continue; } + message_info.reception_sequence_number = subscription_it->second.reception_sequence_number++; + auto subscription = std::dynamic_pointer_cast< rclcpp::experimental::SubscriptionIntraProcessBuffer >(subscription_base); if (subscription != nullptr) { - subscription->provide_intra_process_data(message); + subscription->provide_intra_process_data(message, message_info); continue; } @@ -561,10 +672,11 @@ class IntraProcessManager ROSMessageType ros_msg; rclcpp::TypeAdapter::convert_to_ros_message(*message, ros_msg); ros_message_subscription->provide_intra_process_message( - std::make_shared(ros_msg)); + std::make_shared(ros_msg), message_info); } else { if constexpr (std::is_same::value) { - ros_message_subscription->provide_intra_process_message(message); + ros_message_subscription->provide_intra_process_message(message, + message_info); } else { if constexpr (std::is_same::ros_message_type, ROSMessageType>::value) @@ -573,7 +685,7 @@ class IntraProcessManager rclcpp::TypeAdapter::convert_to_ros_message( *message, ros_msg); ros_message_subscription->provide_intra_process_message( - std::make_shared(ros_msg)); + std::make_shared(ros_msg), message_info); } } } @@ -588,7 +700,8 @@ class IntraProcessManager void add_owned_msg_to_buffers( std::unique_ptr message, - std::vector subscription_ids, + rmw_message_info_t message_info, + const std::vector & subscription_ids, typename allocator::AllocRebind::allocator_type & allocator) { using MessageAllocTraits = allocator::AllocRebind; @@ -608,12 +721,14 @@ class IntraProcessManager if (subscription_it == subscriptions_.end()) { throw std::runtime_error("subscription has unexpectedly gone out of scope"); } - auto subscription_base = subscription_it->second.lock(); + auto subscription_base = subscription_it->second.weak_subscription.lock(); if (subscription_base == nullptr) { subscriptions_.erase(subscription_it); continue; } + message_info.reception_sequence_number = subscription_it->second.reception_sequence_number++; + auto subscription = std::dynamic_pointer_cast< rclcpp::experimental::SubscriptionIntraProcessBuffer @@ -621,7 +736,7 @@ class IntraProcessManager if (subscription != nullptr) { if (std::next(it) == subscription_ids.end()) { // If this is the last subscription, give up ownership - subscription->provide_intra_process_data(std::move(message)); + subscription->provide_intra_process_data(std::move(message), message_info); // Last message delivered, break from for loop break; } else { @@ -630,7 +745,8 @@ class IntraProcessManager auto ptr = MessageAllocTraits::allocate(allocator, 1); MessageAllocTraits::construct(allocator, ptr, *message); - subscription->provide_intra_process_data(MessageUniquePtr(ptr, deleter)); + subscription->provide_intra_process_data(MessageUniquePtr(ptr, deleter), + message_info); } continue; @@ -657,12 +773,14 @@ class IntraProcessManager allocator::set_allocator_for_deleter(&deleter, &allocator); rclcpp::TypeAdapter::convert_to_ros_message(*message, *ptr); auto ros_msg = std::unique_ptr(ptr, deleter); - ros_message_subscription->provide_intra_process_message(std::move(ros_msg)); + ros_message_subscription->provide_intra_process_message(std::move(ros_msg), + message_info); } else { if constexpr (std::is_same::value) { if (std::next(it) == subscription_ids.end()) { // If this is the last subscription, give up ownership - ros_message_subscription->provide_intra_process_message(std::move(message)); + ros_message_subscription->provide_intra_process_message(std::move(message), + message_info); // Last message delivered, break from for loop break; } else { @@ -673,7 +791,7 @@ class IntraProcessManager MessageAllocTraits::construct(allocator, ptr, *message); ros_message_subscription->provide_intra_process_message( - MessageUniquePtr(ptr, deleter)); + MessageUniquePtr(ptr, deleter), message_info); } } } diff --git a/rclcpp/include/rclcpp/experimental/ros_message_intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/ros_message_intra_process_buffer.hpp index 3200d34495..314efca0ea 100644 --- a/rclcpp/include/rclcpp/experimental/ros_message_intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/ros_message_intra_process_buffer.hpp @@ -56,10 +56,9 @@ class SubscriptionROSMsgIntraProcessBuffer : public SubscriptionIntraProcessBase {} virtual void - provide_intra_process_message(ConstMessageSharedPtr message) = 0; - - virtual void - provide_intra_process_message(MessageUniquePtr message) = 0; + provide_intra_process_message( + std::variant message, + const rmw_message_info_t & message_info) = 0; }; } // namespace experimental diff --git a/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp b/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp index 6fbbf5503f..0cd221a842 100644 --- a/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp +++ b/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp @@ -17,7 +17,6 @@ #include -#include #include #include #include @@ -31,7 +30,6 @@ #include "rclcpp/context.hpp" #include "rclcpp/experimental/buffers/intra_process_buffer.hpp" #include "rclcpp/experimental/subscription_intra_process_buffer.hpp" -#include "rclcpp/logging.hpp" #include "rclcpp/qos.hpp" #include "rclcpp/time.hpp" #include "rclcpp/type_support_decl.hpp" @@ -128,19 +126,9 @@ class SubscriptionIntraProcess std::shared_ptr take_data() override { - ConstMessageSharedPtr shared_msg; - MessageUniquePtr unique_msg; - - if (any_callback_.use_take_shared_method()) { - shared_msg = this->buffer_->consume_shared(); - if (!shared_msg) { - return nullptr; - } - } else { - unique_msg = this->buffer_->consume_unique(); - if (!unique_msg) { - return nullptr; - } + auto node = this->buffer_->consume(); + if (node.message.index() == std::variant_npos) { + return nullptr; } if (this->buffer_->has_data()) { @@ -150,9 +138,9 @@ class SubscriptionIntraProcess } return std::static_pointer_cast( - std::make_shared>( - std::pair( - shared_msg, std::move(unique_msg))) + std::make_shared< + typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Node>( + std::move(node)) ); } @@ -184,6 +172,16 @@ class SubscriptionIntraProcess any_callback_.enable(); } + bool + use_take_shared_method() const override + { + if (this->buffer_->buffer_type() == IntraProcessBufferType::CallbackDefault) { + return any_callback_.use_take_shared_method(); + } else { + return this->buffer_->buffer_type() == IntraProcessBufferType::SharedPtr; + } + } + protected: template typename std::enable_if::value, void>::type @@ -200,36 +198,23 @@ class SubscriptionIntraProcess return; } - rmw_message_info_t msg_info; - msg_info.publisher_gid = {0, {0}}; - msg_info.from_intra_process = true; + auto shared_ptr = std::static_pointer_cast< + typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Node>( + data); - const auto nanos = std::chrono::time_point_cast( - std::chrono::system_clock::now()); - if (stats_handler_) { - RCLCPP_WARN_ONCE( - rclcpp::get_logger("rclcpp"), - "Intra-process communication does not support accurate message age statistics"); - // Set source_timestamp to "now" so that message_age reports 0ms rather than - // an invalid value taken from an un-initialised timestamp. IPC delivery - // has little/no transport latency by definition, so near-zero age is expected. - msg_info.source_timestamp = nanos.time_since_epoch().count(); - } + // Copy the message info out before the callback (potentially) moves the message, since + // the stats handler below is invoked after the callback has run. + const rmw_message_info_t message_info = shared_ptr->message_info; - auto shared_ptr = std::static_pointer_cast>( - data); + std::visit( + [&shared_ptr, this](auto && msg) { + any_callback_.dispatch_intra_process(std::move(msg), shared_ptr->message_info); + }, shared_ptr->message); - if (any_callback_.use_take_shared_method()) { - ConstMessageSharedPtr shared_msg = shared_ptr->first; - any_callback_.dispatch_intra_process(shared_msg, msg_info); - } else { - MessageUniquePtr unique_msg = std::move(shared_ptr->second); - any_callback_.dispatch_intra_process(std::move(unique_msg), msg_info); - } shared_ptr.reset(); if (stats_handler_) { - stats_handler_(msg_info, rclcpp::Time(nanos.time_since_epoch().count())); + stats_handler_(message_info, rclcpp::Time(message_info.source_timestamp)); } } diff --git a/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp index fce7bd452b..a8a78bd06b 100644 --- a/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp @@ -69,11 +69,11 @@ class SubscriptionIntraProcessBuffer : public SubscriptionROSMsgIntraProcessBuff using ConstDataSharedPtr = std::shared_ptr; using SubscribedTypeUniquePtr = std::unique_ptr; - using BufferUniquePtr = typename rclcpp::experimental::buffers::IntraProcessBuffer< + using IntraProcessBuffer = typename rclcpp::experimental::buffers::IntraProcessBuffer< SubscribedType, - Alloc, SubscribedTypeDeleter - >::UniquePtr; + >; + using BufferUniquePtr = typename IntraProcessBuffer::UniquePtr; SubscriptionIntraProcessBuffer( std::shared_ptr allocator, @@ -89,8 +89,8 @@ class SubscriptionIntraProcessBuffer : public SubscriptionROSMsgIntraProcessBuff allocator::set_allocator_for_deleter(&subscribed_type_deleter_, &subscribed_type_allocator_); // Create the intra-process buffer. - buffer_ = rclcpp::experimental::create_intra_process_buffer( + buffer_ = rclcpp::experimental::create_intra_process_buffer( buffer_type, qos_profile, std::make_shared(subscribed_type_allocator_)); @@ -131,53 +131,39 @@ class SubscriptionIntraProcessBuffer : public SubscriptionROSMsgIntraProcessBuff } void - provide_intra_process_message(ConstMessageSharedPtr message) override + provide_intra_process_message( + std::variant message, + const rmw_message_info_t & message_info) override { + typename IntraProcessBuffer::Node node; + node.message_info = message_info; if constexpr (std::is_same::value) { - buffer_->add_shared(std::move(message)); + node.message = std::move(message); trigger_guard_condition(); } else { - buffer_->add_shared(convert_ros_message_to_subscribed_type_unique_ptr(*message)); + std::visit( + [this, &node](auto && msg) { + node.message = convert_ros_message_to_subscribed_type_unique_ptr(*msg); + }, message); trigger_guard_condition(); } + buffer_->add(std::move(node)); this->invoke_on_new_message(); } void - provide_intra_process_message(MessageUniquePtr message) override + provide_intra_process_data( + std::variant message, + const rmw_message_info_t & message_info) { - if constexpr (std::is_same::value) { - buffer_->add_unique(std::move(message)); - trigger_guard_condition(); - } else { - buffer_->add_unique(convert_ros_message_to_subscribed_type_unique_ptr(*message)); - trigger_guard_condition(); - } - this->invoke_on_new_message(); - } - - void - provide_intra_process_data(ConstDataSharedPtr message) - { - buffer_->add_shared(std::move(message)); + typename IntraProcessBuffer::Node node; + node.message_info = message_info; + node.message = std::move(message); + buffer_->add(std::move(node)); trigger_guard_condition(); this->invoke_on_new_message(); } - void - provide_intra_process_data(SubscribedTypeUniquePtr message) - { - buffer_->add_unique(std::move(message)); - trigger_guard_condition(); - this->invoke_on_new_message(); - } - - bool - use_take_shared_method() const override - { - return buffer_->use_take_shared_method(); - } - size_t available_capacity() const override { return buffer_->available_capacity(); diff --git a/rclcpp/include/rclcpp/publisher.hpp b/rclcpp/include/rclcpp/publisher.hpp index 6ea50b67a2..9afdd965e7 100644 --- a/rclcpp/include/rclcpp/publisher.hpp +++ b/rclcpp/include/rclcpp/publisher.hpp @@ -94,11 +94,11 @@ class Publisher : public PublisherBase using ROSMessageTypeAllocator = typename ROSMessageTypeAllocatorTraits::allocator_type; using ROSMessageTypeDeleter = allocator::Deleter; - using BufferSharedPtr = typename rclcpp::experimental::buffers::IntraProcessBuffer< + using IntraProcessBuffer = typename rclcpp::experimental::buffers::IntraProcessBuffer< ROSMessageType, - ROSMessageTypeAllocator, ROSMessageTypeDeleter - >::SharedPtr; + >; + using BufferSharedPtr = typename IntraProcessBuffer::SharedPtr; RCLCPP_SMART_PTR_DEFINITIONS(Publisher) @@ -238,17 +238,23 @@ class Publisher : public PublisherBase get_subscription_count() > get_intra_process_subscription_count() || buffer_; if (inter_process_publish_needed) { - auto shared_msg = + auto msg_info_pair = this->do_intra_process_ros_message_publish_and_return_shared(std::move(msg)); if (buffer_) { - buffer_->add_shared(shared_msg); + typename IntraProcessBuffer::Node node; + node.message = msg_info_pair.first; + node.message_info = msg_info_pair.second; + buffer_->add(std::move(node)); } - this->do_inter_process_publish(*shared_msg); + this->do_inter_process_publish(*msg_info_pair.first); } else { if (buffer_) { - auto shared_msg = + auto msg_info_pair = this->do_intra_process_ros_message_publish_and_return_shared(std::move(msg)); - buffer_->add_shared(shared_msg); + typename IntraProcessBuffer::Node node; + node.message = msg_info_pair.first; + node.message_info = msg_info_pair.second; + buffer_->add(std::move(node)); } else { this->do_intra_process_ros_message_publish(std::move(msg)); } @@ -321,18 +327,28 @@ class Publisher : public PublisherBase if (inter_process_publish_needed) { auto ros_msg_ptr = std::make_shared(); rclcpp::TypeAdapter::convert_to_ros_message(*msg, *ros_msg_ptr); - this->do_intra_process_publish(std::move(msg)); + rmw_message_info_t message_info; + this->do_intra_process_publish(std::move(msg), &message_info); this->do_inter_process_publish(*ros_msg_ptr); if (buffer_) { - buffer_->add_shared(ros_msg_ptr); + typename IntraProcessBuffer::Node node; + node.message = std::move(ros_msg_ptr); + node.message_info = std::move(message_info); + buffer_->add(std::move(node)); } } else { if (buffer_) { auto ros_msg_ptr = std::make_shared(); rclcpp::TypeAdapter::convert_to_ros_message(*msg, *ros_msg_ptr); - buffer_->add_shared(ros_msg_ptr); + rmw_message_info_t message_info; + this->do_intra_process_publish(std::move(msg), &message_info); + typename IntraProcessBuffer::Node node; + node.message = std::move(ros_msg_ptr); + node.message_info = std::move(message_info); + buffer_->add(std::move(node)); + } else { + this->do_intra_process_publish(std::move(msg)); } - this->do_intra_process_publish(std::move(msg)); } } @@ -482,7 +498,9 @@ class Publisher : public PublisherBase } void - do_intra_process_publish(std::unique_ptr msg) + do_intra_process_publish( + std::unique_ptr msg, + rmw_message_info_t * message_info_out = nullptr) { auto ipm = weak_ipm_.lock(); if (!ipm) { @@ -500,11 +518,14 @@ class Publisher : public PublisherBase ipm->template do_intra_process_publish( intra_process_publisher_id_, std::move(msg), - published_type_allocator_); + published_type_allocator_, + message_info_out); } void - do_intra_process_ros_message_publish(std::unique_ptr msg) + do_intra_process_ros_message_publish( + std::unique_ptr msg, + rmw_message_info_t * message_info_out = nullptr) { auto ipm = weak_ipm_.lock(); if (!ipm) { @@ -522,10 +543,11 @@ class Publisher : public PublisherBase ipm->template do_intra_process_publish( intra_process_publisher_id_, std::move(msg), - ros_message_type_allocator_); + ros_message_type_allocator_, + message_info_out); } - std::shared_ptr + std::pair, rmw_message_info_t> do_intra_process_ros_message_publish_and_return_shared( std::unique_ptr msg) { diff --git a/rclcpp/src/rclcpp/intra_process_manager.cpp b/rclcpp/src/rclcpp/intra_process_manager.cpp index 2cfc5386d8..76044f0622 100644 --- a/rclcpp/src/rclcpp/intra_process_manager.cpp +++ b/rclcpp/src/rclcpp/intra_process_manager.cpp @@ -40,7 +40,8 @@ IntraProcessManager::add_publisher( uint64_t pub_id = IntraProcessManager::get_next_unique_id(); - publishers_[pub_id] = publisher; + PublisherData data; + data.weak_publisher = publisher; if (publisher->is_durability_transient_local()) { if (buffer) { publisher_buffers_[pub_id] = buffer; @@ -51,6 +52,8 @@ IntraProcessManager::add_publisher( } } + publishers_[pub_id] = data; + // Add GID to publisher info mapping for fast lookups (stores both ID and weak_ptr) gid_to_publisher_info_[publisher->get_gid()] = {pub_id, publisher}; @@ -59,7 +62,7 @@ IntraProcessManager::add_publisher( // create an entry for the publisher id and populate with already existing subscriptions for (auto & pair : subscriptions_) { - auto subscription = pair.second.lock(); + auto subscription = pair.second.weak_subscription.lock(); if (!subscription) { continue; } @@ -105,7 +108,7 @@ IntraProcessManager::remove_publisher(uint64_t intra_process_publisher_id) // First try via the publisher's own GID (fast path). auto pub_it = publishers_.find(intra_process_publisher_id); if (pub_it != publishers_.end()) { - auto publisher = pub_it->second.lock(); + auto publisher = pub_it->second.weak_publisher.lock(); if (publisher) { gid_to_publisher_info_.erase(publisher->get_gid()); } else { @@ -170,7 +173,7 @@ IntraProcessManager::get_subscription_intra_process(uint64_t intra_process_subsc if (subscription_it == subscriptions_.end()) { return nullptr; } else { - auto subscription = subscription_it->second.lock(); + auto subscription = subscription_it->second.weak_subscription.lock(); if (subscription) { return subscription; } else { @@ -258,7 +261,7 @@ IntraProcessManager::lowest_available_capacity(const uint64_t intra_process_publ { auto subscription_it = subscriptions_.find(intra_process_subscription_id); if (subscription_it != subscriptions_.end()) { - auto subscription = subscription_it->second.lock(); + auto subscription = subscription_it->second.weak_subscription.lock(); if (subscription) { capacity = std::min(capacity, subscription->available_capacity()); } From 05f5e585c307573ecf7d2abe5a9dd7797d3fe902 Mon Sep 17 00:00:00 2001 From: Thomas Moore Date: Tue, 2 Jul 2024 11:10:41 -0400 Subject: [PATCH 2/3] Move IntraProcessBufferData to its own header (cherry picked from commit f2302935bd1b41a75f56c6a9f4ef5b7fbe4d7e9b) (cherry picked from commit 8a875ee1df60afcdb7ed68ceb27f097acdcdacf0) Signed-off-by: Thomas Moore --- .../buffers/intra_process_buffer.hpp | 70 +++++++-------- .../buffers/intra_process_buffer_data.hpp | 90 +++++++++++++++++++ .../create_intra_process_buffer.hpp | 4 +- .../experimental/intra_process_manager.hpp | 14 +-- .../subscription_intra_process.hpp | 10 +-- .../subscription_intra_process_buffer.hpp | 20 ++--- rclcpp/include/rclcpp/publisher.hpp | 32 +++---- 7 files changed, 160 insertions(+), 80 deletions(-) create mode 100644 rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer_data.hpp diff --git a/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp index 9e27c1481d..2093f6caf6 100644 --- a/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp @@ -25,6 +25,7 @@ #include "rclcpp/allocator/allocator_common.hpp" #include "rclcpp/allocator/allocator_deleter.hpp" #include "rclcpp/experimental/buffers/buffer_implementation_base.hpp" +#include "rclcpp/experimental/buffers/intra_process_buffer_data.hpp" #include "rclcpp/intra_process_buffer_type.hpp" #include "rclcpp/macros.hpp" #include "tracetools/tracetools.h" @@ -36,18 +37,6 @@ namespace experimental namespace buffers { -template< - typename MessageT, - typename MessageDeleter = std::default_delete> -struct IntraProcessBufferNode -{ - using MessageUniquePtr = std::unique_ptr; - using MessageSharedPtr = std::shared_ptr; - - std::variant message; - rmw_message_info_t message_info; -}; - class IntraProcessBufferBase { public: @@ -72,13 +61,13 @@ class IntraProcessBuffer : public IntraProcessBufferBase virtual ~IntraProcessBuffer() {} - using Node = IntraProcessBufferNode; + using Data = IntraProcessBufferData; - virtual void add(Node node) = 0; + virtual void add(Data data) = 0; - virtual Node consume() = 0; + virtual Data consume() = 0; - virtual std::vector get_all_data() = 0; + virtual std::vector get_all_data() = 0; }; template< @@ -91,15 +80,16 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer; using MessageAllocTraits = allocator::AllocRebind; using MessageAlloc = typename MessageAllocTraits::allocator_type; - using Node = IntraProcessBufferNode; - using MessageSharedPtr = typename Node::MessageSharedPtr; - using MessageUniquePtr = typename Node::MessageUniquePtr; + using Data = typename Buffer::Data; + using MessageSharedPtr = typename Data::MessageSharedPtr; + using MessageUniquePtr = typename Data::MessageUniquePtr; explicit TypedIntraProcessBuffer( - std::unique_ptr> buffer_impl, + std::unique_ptr> buffer_impl, std::shared_ptr allocator = nullptr) : buffer_(std::move(buffer_impl)) { @@ -116,17 +106,17 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer(std::move(node)); + add_impl(std::move(data)); } - Node consume() override + Data consume() override { return buffer_->dequeue(); } - std::vector get_all_data() override + std::vector get_all_data() override { return buffer_->get_all_data(); } @@ -152,51 +142,51 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer> buffer_; + std::unique_ptr> buffer_; std::shared_ptr message_allocator_; template typename std::enable_if_t< BufferT == IntraProcessBufferType::CallbackDefault> - add_impl(Node node) + add_impl(Data data) { - buffer_->enqueue(std::move(node)); + buffer_->enqueue(std::move(data)); } template typename std::enable_if_t< BufferT == IntraProcessBufferType::SharedPtr> - add_impl(Node node) + add_impl(Data data) { - if (std::holds_alternative(node.message)) { - buffer_->enqueue(std::move(node)); + if (std::holds_alternative(data.message)) { + buffer_->enqueue(std::move(data)); } else { // Promote to a shared pointer - auto unique_msg = std::move(std::get(node.message)); - node.message = MessageSharedPtr(unique_msg.release()); - buffer_->enqueue(std::move(node)); + auto unique_msg = std::move(std::get(data.message)); + data.message = MessageSharedPtr(unique_msg.release()); + buffer_->enqueue(std::move(data)); } } template typename std::enable_if_t< BufferT == IntraProcessBufferType::UniquePtr> - add_impl(Node node) + add_impl(Data data) { - if (std::holds_alternative(node.message)) { - buffer_->enqueue(std::move(node)); + if (std::holds_alternative(data.message)) { + buffer_->enqueue(std::move(data)); } else { - auto shared_msg = std::move(std::get(node.message)); + auto shared_msg = std::move(std::get(data.message)); MessageDeleter * deleter = std::get_deleter(shared_msg); auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1); MessageAllocTraits::construct(*message_allocator_.get(), ptr, *shared_msg); if (deleter) { - node.message = MessageUniquePtr(ptr, *deleter); + data.message = MessageUniquePtr(ptr, *deleter); } else { - node.message = MessageUniquePtr(ptr); + data.message = MessageUniquePtr(ptr); } - buffer_->enqueue(std::move(node)); + buffer_->enqueue(std::move(data)); } } }; diff --git a/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer_data.hpp b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer_data.hpp new file mode 100644 index 0000000000..00068fb2d6 --- /dev/null +++ b/rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer_data.hpp @@ -0,0 +1,90 @@ +// Copyright 2024 Open Source Robotics Foundation, Inc. +// +// 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. + +#ifndef RCLCPP__EXPERIMENTAL__BUFFERS__INTRA_PROCESS_BUFFER_DATA_HPP_ +#define RCLCPP__EXPERIMENTAL__BUFFERS__INTRA_PROCESS_BUFFER_DATA_HPP_ + +#include +#include +#include + +#include "rmw/types.h" + +namespace rclcpp +{ +namespace experimental +{ +namespace buffers +{ + +template< + typename MessageT, + typename MessageDeleter = std::default_delete> +struct IntraProcessBufferData +{ + using MessageUniquePtr = std::unique_ptr; + using MessageSharedPtr = std::shared_ptr; + using MessageVariant = std::variant; + + MessageVariant message; + rmw_message_info_t message_info; + + IntraProcessBufferData() = default; + IntraProcessBufferData(IntraProcessBufferData &&) = default; + IntraProcessBufferData & operator=(IntraProcessBufferData &&) = default; + + IntraProcessBufferData(MessageVariant message_in, rmw_message_info_t message_info_in) + : message(std::move(message_in)), message_info(message_info_in) + {} + + /// Deep-copies the unique_ptr alternative. + /** + * std::variant's implicit copy constructor is deleted whenever one of its + * alternatives (here, MessageUniquePtr) is not copy-constructible, so this must be + * done explicitly to keep IntraProcessBufferData itself copyable. A null + * unique_ptr copies to another null unique_ptr; otherwise a fresh copy of the + * pointee is allocated, since the original must remain owned by the buffer (e.g. + * for future transient-local replays to other subscriptions). + */ + IntraProcessBufferData(const IntraProcessBufferData & other) + : message_info(other.message_info) + { + if (std::holds_alternative(other.message)) { + message = std::get(other.message); + } else { + const auto & other_unique_msg = std::get(other.message); + if (other_unique_msg) { + message = MessageUniquePtr( + new MessageT(*other_unique_msg), other_unique_msg.get_deleter()); + } else { + message = MessageUniquePtr(nullptr, other_unique_msg.get_deleter()); + } + } + } + + IntraProcessBufferData & operator=(const IntraProcessBufferData & other) + { + if (this != &other) { + *this = IntraProcessBufferData(other); + } + return *this; + } +}; + +} // namespace buffers +} // namespace experimental +} // namespace rclcpp + + +#endif // RCLCPP__EXPERIMENTAL__BUFFERS__INTRA_PROCESS_BUFFER_DATA_HPP_ diff --git a/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp index 40712a0b7e..9f033299ef 100644 --- a/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/create_intra_process_buffer.hpp @@ -41,11 +41,11 @@ create_intra_process_buffer( size_t buffer_size = qos.depth(); using rclcpp::experimental::buffers::IntraProcessBuffer; - using rclcpp::experimental::buffers::IntraProcessBufferNode; + using rclcpp::experimental::buffers::IntraProcessBufferData; using rclcpp::experimental::buffers::TypedIntraProcessBuffer; using rclcpp::experimental::buffers::RingBufferImplementation; - using BufferT = IntraProcessBufferNode; + using BufferT = IntraProcessBufferData; using BufferImplT = RingBufferImplementation; auto buffer_implementation = std::make_unique(buffer_size); diff --git a/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp b/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp index e0f7dd3cbc..25c9ba325c 100644 --- a/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp +++ b/rclcpp/include/rclcpp/experimental/intra_process_manager.hpp @@ -544,18 +544,18 @@ class IntraProcessManager using ROSMessageTypeAllocatorTraits = allocator::AllocRebind; using ROSMessageTypeAllocator = typename ROSMessageTypeAllocatorTraits::allocator_type; using ROSMessageTypeDeleter = allocator::Deleter; - using IntraProcessBuffer = rclcpp::experimental::buffers::IntraProcessBuffer< - ROSMessageType, - ROSMessageTypeDeleter - >; - using ROSMessageSharedPtr = typename IntraProcessBuffer::Node::MessageSharedPtr; - using ROSMessageUniquePtr = typename IntraProcessBuffer::Node::MessageUniquePtr; + using BufferType = rclcpp::experimental::buffers::IntraProcessBuffer< + ROSMessageType, + ROSMessageTypeDeleter + >; + using ROSMessageSharedPtr = typename BufferType::Data::MessageSharedPtr; + using ROSMessageUniquePtr = typename BufferType::Data::MessageUniquePtr; auto publisher_buffer = publisher_buffers_[pub_id].lock(); if (!publisher_buffer) { throw std::runtime_error("publisher buffer has unexpectedly gone out of scope"); } - auto buffer = std::dynamic_pointer_cast(publisher_buffer); + auto buffer = std::dynamic_pointer_cast(publisher_buffer); if (!buffer) { throw std::runtime_error( "failed to dynamic cast publisher's IntraProcessBufferBase to " diff --git a/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp b/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp index 0cd221a842..6eb2a291de 100644 --- a/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp +++ b/rclcpp/include/rclcpp/experimental/subscription_intra_process.hpp @@ -126,8 +126,8 @@ class SubscriptionIntraProcess std::shared_ptr take_data() override { - auto node = this->buffer_->consume(); - if (node.message.index() == std::variant_npos) { + auto data = this->buffer_->consume(); + if (data.message.index() == std::variant_npos) { return nullptr; } @@ -139,8 +139,8 @@ class SubscriptionIntraProcess return std::static_pointer_cast( std::make_shared< - typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Node>( - std::move(node)) + typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Data>( + std::move(data)) ); } @@ -199,7 +199,7 @@ class SubscriptionIntraProcess } auto shared_ptr = std::static_pointer_cast< - typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Node>( + typename SubscriptionIntraProcessBufferT::IntraProcessBuffer::Data>( data); // Copy the message info out before the callback (potentially) moves the message, since diff --git a/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp b/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp index a8a78bd06b..a5fb27ab33 100644 --- a/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp +++ b/rclcpp/include/rclcpp/experimental/subscription_intra_process_buffer.hpp @@ -135,19 +135,19 @@ class SubscriptionIntraProcessBuffer : public SubscriptionROSMsgIntraProcessBuff std::variant message, const rmw_message_info_t & message_info) override { - typename IntraProcessBuffer::Node node; - node.message_info = message_info; + typename IntraProcessBuffer::Data data; + data.message_info = message_info; if constexpr (std::is_same::value) { - node.message = std::move(message); + data.message = std::move(message); trigger_guard_condition(); } else { std::visit( - [this, &node](auto && msg) { - node.message = convert_ros_message_to_subscribed_type_unique_ptr(*msg); + [this, &data](auto && msg) { + data.message = convert_ros_message_to_subscribed_type_unique_ptr(*msg); }, message); trigger_guard_condition(); } - buffer_->add(std::move(node)); + buffer_->add(std::move(data)); this->invoke_on_new_message(); } @@ -156,10 +156,10 @@ class SubscriptionIntraProcessBuffer : public SubscriptionROSMsgIntraProcessBuff std::variant message, const rmw_message_info_t & message_info) { - typename IntraProcessBuffer::Node node; - node.message_info = message_info; - node.message = std::move(message); - buffer_->add(std::move(node)); + typename IntraProcessBuffer::Data data; + data.message_info = message_info; + data.message = std::move(message); + buffer_->add(std::move(data)); trigger_guard_condition(); this->invoke_on_new_message(); } diff --git a/rclcpp/include/rclcpp/publisher.hpp b/rclcpp/include/rclcpp/publisher.hpp index 9afdd965e7..afad030570 100644 --- a/rclcpp/include/rclcpp/publisher.hpp +++ b/rclcpp/include/rclcpp/publisher.hpp @@ -241,20 +241,20 @@ class Publisher : public PublisherBase auto msg_info_pair = this->do_intra_process_ros_message_publish_and_return_shared(std::move(msg)); if (buffer_) { - typename IntraProcessBuffer::Node node; - node.message = msg_info_pair.first; - node.message_info = msg_info_pair.second; - buffer_->add(std::move(node)); + typename IntraProcessBuffer::Data data; + data.message = msg_info_pair.first; + data.message_info = msg_info_pair.second; + buffer_->add(std::move(data)); } this->do_inter_process_publish(*msg_info_pair.first); } else { if (buffer_) { auto msg_info_pair = this->do_intra_process_ros_message_publish_and_return_shared(std::move(msg)); - typename IntraProcessBuffer::Node node; - node.message = msg_info_pair.first; - node.message_info = msg_info_pair.second; - buffer_->add(std::move(node)); + typename IntraProcessBuffer::Data data; + data.message = msg_info_pair.first; + data.message_info = msg_info_pair.second; + buffer_->add(std::move(data)); } else { this->do_intra_process_ros_message_publish(std::move(msg)); } @@ -331,10 +331,10 @@ class Publisher : public PublisherBase this->do_intra_process_publish(std::move(msg), &message_info); this->do_inter_process_publish(*ros_msg_ptr); if (buffer_) { - typename IntraProcessBuffer::Node node; - node.message = std::move(ros_msg_ptr); - node.message_info = std::move(message_info); - buffer_->add(std::move(node)); + typename IntraProcessBuffer::Data data; + data.message = std::move(ros_msg_ptr); + data.message_info = std::move(message_info); + buffer_->add(std::move(data)); } } else { if (buffer_) { @@ -342,10 +342,10 @@ class Publisher : public PublisherBase rclcpp::TypeAdapter::convert_to_ros_message(*msg, *ros_msg_ptr); rmw_message_info_t message_info; this->do_intra_process_publish(std::move(msg), &message_info); - typename IntraProcessBuffer::Node node; - node.message = std::move(ros_msg_ptr); - node.message_info = std::move(message_info); - buffer_->add(std::move(node)); + typename IntraProcessBuffer::Data data; + data.message = std::move(ros_msg_ptr); + data.message_info = std::move(message_info); + buffer_->add(std::move(data)); } else { this->do_intra_process_publish(std::move(msg)); } From bef63b7becf4c05ed8ca826476282982270d1eed Mon Sep 17 00:00:00 2001 From: Thomas Moore Date: Thu, 3 Sep 2026 19:24:19 +0000 Subject: [PATCH 3/3] Update intra-process unit tests Signed-off-by: Thomas Moore --- .../test/rclcpp/test_intra_process_buffer.cpp | 364 ++++++++++-------- .../rclcpp/test_intra_process_manager.cpp | 221 ++++++----- 2 files changed, 335 insertions(+), 250 deletions(-) diff --git a/rclcpp/test/rclcpp/test_intra_process_buffer.cpp b/rclcpp/test/rclcpp/test_intra_process_buffer.cpp index 1cd4d37bb6..d76ffe923b 100644 --- a/rclcpp/test/rclcpp/test_intra_process_buffer.cpp +++ b/rclcpp/test/rclcpp/test_intra_process_buffer.cpp @@ -18,8 +18,23 @@ #include "gtest/gtest.h" -#include "rclcpp/experimental/buffers/intra_process_buffer.hpp" -#include "rclcpp/experimental/buffers/ring_buffer_implementation.hpp" +#include "rclcpp/experimental/buffers/intra_process_buffer_data.hpp" +#include "rclcpp/intra_process_buffer_type.hpp" +#include "rclcpp/rclcpp.hpp" + +template +static constexpr rmw_message_info_t generate_message_info() +{ + rmw_message_info_t message_info{}; + message_info.source_timestamp = static_cast(10 * idx + 1); + message_info.received_timestamp = static_cast(10 * idx + 2); + message_info.publication_sequence_number = static_cast(10 * idx + 3); + message_info.reception_sequence_number = static_cast(10 * idx + 4); + message_info.publisher_gid = {"test", {idx, idx, idx, idx, idx, idx, idx, + idx, idx, idx, idx, idx, idx, idx, idx, idx}}; + message_info.from_intra_process = true; + return message_info; +} /* Constructor @@ -28,26 +43,24 @@ TEST(TestIntraProcessBuffer, constructor) { using MessageT = char; using Alloc = std::allocator; using Deleter = std::default_delete; - using SharedMessageT = std::shared_ptr; - using UniqueMessageT = std::unique_ptr; using SharedIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, SharedMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::SharedPtr>; using UniqueIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, UniqueMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::UniquePtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; - auto shared_buffer_impl = - std::make_unique>(2); + auto shared_buffer_impl = std::make_unique(2); SharedIntraProcessBufferT shared_intra_process_buffer(std::move(shared_buffer_impl)); - EXPECT_EQ(true, shared_intra_process_buffer.use_take_shared_method()); + EXPECT_EQ(rclcpp::IntraProcessBufferType::SharedPtr, shared_intra_process_buffer.buffer_type()); - auto unique_buffer_impl = - std::make_unique>(2); + auto unique_buffer_impl = std::make_unique(2); UniqueIntraProcessBufferT unique_intra_process_buffer(std::move(unique_buffer_impl)); - EXPECT_EQ(false, unique_intra_process_buffer.use_take_shared_method()); + EXPECT_EQ(rclcpp::IntraProcessBufferType::UniquePtr, unique_intra_process_buffer.buffer_type()); } /* @@ -62,40 +75,58 @@ TEST(TestIntraProcessBuffer, shared_buffer_add) { using Deleter = std::default_delete; using SharedMessageT = std::shared_ptr; using SharedIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, SharedMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::SharedPtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; - auto buffer_impl = - std::make_unique>(2); + auto buffer_impl = std::make_unique(2); SharedIntraProcessBufferT intra_process_buffer(std::move(buffer_impl)); auto original_shared_msg = std::make_shared('a'); auto original_message_pointer = reinterpret_cast(original_shared_msg.get()); + rmw_message_info_t original_message_info = generate_message_info<0>(); - intra_process_buffer.add_shared(original_shared_msg); + intra_process_buffer.add({original_shared_msg, original_message_info}); EXPECT_EQ(2L, original_shared_msg.use_count()); SharedMessageT popped_shared_msg; - popped_shared_msg = intra_process_buffer.consume_shared(); + rmw_message_info_t popped_message_info; + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_shared_msg = std::get(popped_data.message); + popped_message_info = popped_data.message_info; + } auto popped_message_pointer = reinterpret_cast(popped_shared_msg.get()); EXPECT_EQ(original_shared_msg.use_count(), popped_shared_msg.use_count()); EXPECT_EQ(*original_shared_msg, *popped_shared_msg); EXPECT_EQ(original_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); auto original_unique_msg = std::make_unique('b'); original_message_pointer = reinterpret_cast(original_unique_msg.get()); auto original_value = *original_unique_msg; + original_message_info = generate_message_info<1>(); - intra_process_buffer.add_unique(std::move(original_unique_msg)); + intra_process_buffer.add({std::move(original_unique_msg), original_message_info}); - popped_shared_msg = intra_process_buffer.consume_shared(); + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_shared_msg = std::get(popped_data.message); + popped_message_info = popped_data.message_info; + } popped_message_pointer = reinterpret_cast(popped_shared_msg.get()); EXPECT_EQ(1L, popped_shared_msg.use_count()); EXPECT_EQ(original_value, *popped_shared_msg); EXPECT_EQ(original_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); } /* @@ -110,188 +141,176 @@ TEST(TestIntraProcessBuffer, unique_buffer_add) { using Deleter = std::default_delete; using UniqueMessageT = std::unique_ptr; using UniqueIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, UniqueMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::UniquePtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; - auto buffer_impl = - std::make_unique>(2); + auto buffer_impl = std::make_unique(2); UniqueIntraProcessBufferT intra_process_buffer(std::move(buffer_impl)); auto original_shared_msg = std::make_shared('a'); auto original_message_pointer = reinterpret_cast(original_shared_msg.get()); + rmw_message_info_t original_message_info = generate_message_info<0>(); - intra_process_buffer.add_shared(original_shared_msg); + intra_process_buffer.add({original_shared_msg, original_message_info}); EXPECT_EQ(1L, original_shared_msg.use_count()); UniqueMessageT popped_unique_msg; - popped_unique_msg = intra_process_buffer.consume_unique(); + rmw_message_info_t popped_message_info; + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_unique_msg = std::move(std::get(popped_data.message)); + popped_message_info = popped_data.message_info; + } auto popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); EXPECT_EQ(*original_shared_msg, *popped_unique_msg); EXPECT_NE(original_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); auto original_unique_msg = std::make_unique('b'); original_message_pointer = reinterpret_cast(original_unique_msg.get()); auto original_value = *original_unique_msg; + original_message_info = generate_message_info<1>(); - intra_process_buffer.add_unique(std::move(original_unique_msg)); + intra_process_buffer.add({std::move(original_unique_msg), original_message_info}); - popped_unique_msg = intra_process_buffer.consume_unique(); + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_unique_msg = std::move(std::get(popped_data.message)); + popped_message_info = popped_data.message_info; + } popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); EXPECT_EQ(original_value, *popped_unique_msg); EXPECT_EQ(original_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); } /* - Consume data from an intra-process buffer with an implementations that stores shared_ptr - Messages are inserted using the same data as the implementation, i.e. shared_ptr - - Request shared_ptr no copies are expected - - Request unique_ptr a copy is expected + Get all data from an intra-process buffer with an implementations that stores shared_ptr + Messages are inserted as shared_ptr and then as unique_ptr + - Add shared_ptr no copies are expected + - Add unique_ptr no copies are expected */ -TEST(TestIntraProcessBuffer, shared_buffer_consume) { +TEST(TestIntraProcessBuffer, shared_buffer_get_all_data) { using MessageT = char; using Alloc = std::allocator; using Deleter = std::default_delete; using SharedMessageT = std::shared_ptr; - using UniqueMessageT = std::unique_ptr; using SharedIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, SharedMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::SharedPtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; - auto buffer_impl = - std::make_unique>(2); + auto buffer_impl = std::make_unique(2); SharedIntraProcessBufferT intra_process_buffer(std::move(buffer_impl)); auto original_shared_msg = std::make_shared('a'); auto original_message_pointer = reinterpret_cast(original_shared_msg.get()); + auto original_value = *original_shared_msg; + rmw_message_info_t original_message_info = generate_message_info<0>(); + intra_process_buffer.add({original_shared_msg, original_message_info}); - intra_process_buffer.add_shared(original_shared_msg); - - EXPECT_EQ(2L, original_shared_msg.use_count()); - - SharedMessageT popped_shared_msg; - popped_shared_msg = intra_process_buffer.consume_shared(); - auto popped_message_pointer = reinterpret_cast(popped_shared_msg.get()); - - EXPECT_EQ(original_shared_msg.use_count(), popped_shared_msg.use_count()); - EXPECT_EQ(*original_shared_msg, *popped_shared_msg); - EXPECT_EQ(original_message_pointer, popped_message_pointer); - - original_shared_msg = std::make_shared('b'); - original_message_pointer = reinterpret_cast(original_shared_msg.get()); - - intra_process_buffer.add_shared(original_shared_msg); - - UniqueMessageT popped_unique_msg; - popped_unique_msg = intra_process_buffer.consume_unique(); - popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); - - EXPECT_EQ(1L, original_shared_msg.use_count()); - EXPECT_EQ(*original_shared_msg, *popped_unique_msg); - EXPECT_NE(original_message_pointer, popped_message_pointer); - - original_shared_msg = std::make_shared('c'); - original_message_pointer = reinterpret_cast(original_shared_msg.get()); - auto original_shared_msg_2 = std::make_shared('d'); - auto original_message_pointer_2 = reinterpret_cast(original_shared_msg_2.get()); - intra_process_buffer.add_shared(original_shared_msg); - intra_process_buffer.add_shared(original_shared_msg_2); - - auto shared_data_vec = intra_process_buffer.get_all_data_shared(); - EXPECT_EQ(2L, shared_data_vec.size()); - EXPECT_EQ(3L, original_shared_msg.use_count()); - EXPECT_EQ(original_shared_msg.use_count(), shared_data_vec[0].use_count()); - EXPECT_EQ(*original_shared_msg, *shared_data_vec[0]); - EXPECT_EQ(original_message_pointer, reinterpret_cast(shared_data_vec[0].get())); - EXPECT_EQ(3L, original_shared_msg_2.use_count()); - EXPECT_EQ(original_shared_msg_2.use_count(), shared_data_vec[1].use_count()); - EXPECT_EQ(*original_shared_msg_2, *shared_data_vec[1]); - EXPECT_EQ(original_message_pointer_2, reinterpret_cast(shared_data_vec[1].get())); - - auto unique_data_vec = intra_process_buffer.get_all_data_unique(); - EXPECT_EQ(2L, unique_data_vec.size()); - EXPECT_EQ(3L, original_shared_msg.use_count()); - EXPECT_EQ(*original_shared_msg, *unique_data_vec[0]); - EXPECT_NE(original_message_pointer, reinterpret_cast(unique_data_vec[0].get())); - EXPECT_EQ(3L, original_shared_msg_2.use_count()); - EXPECT_EQ(*original_shared_msg_2, *unique_data_vec[1]); - EXPECT_NE(original_message_pointer_2, reinterpret_cast(unique_data_vec[1].get())); + auto original_unique_msg_2 = std::make_unique('b'); + auto original_message_pointer_2 = reinterpret_cast(original_unique_msg_2.get()); + auto original_value_2 = *original_unique_msg_2; + rmw_message_info_t original_message_info_2 = generate_message_info<1>(); + intra_process_buffer.add({std::move(original_unique_msg_2), original_message_info_2}); + + auto all_data = intra_process_buffer.get_all_data(); + ASSERT_EQ(2u, all_data.size()); + + ASSERT_TRUE(std::holds_alternative(all_data[0].message)); + { + const auto & shared_msg = std::get(all_data[0].message); + // 3 owners: original_shared_msg, the buffer's own stored copy, and this + // get_all_data() snapshot's copy. + EXPECT_EQ(3L, original_shared_msg.use_count()); + EXPECT_EQ(original_shared_msg.use_count(), shared_msg.use_count()); + EXPECT_EQ(original_value, *shared_msg); + EXPECT_EQ(original_message_pointer, reinterpret_cast(shared_msg.get())); + EXPECT_EQ( + std::memcmp(&original_message_info, &all_data[0].message_info, sizeof(rmw_message_info_t)), + 0); + } + + ASSERT_TRUE(std::holds_alternative(all_data[1].message)); + { + const auto & shared_msg = std::get(all_data[1].message); + EXPECT_EQ(original_value_2, *shared_msg); + EXPECT_EQ(original_message_pointer_2, reinterpret_cast(shared_msg.get())); + EXPECT_EQ( + std::memcmp( + &original_message_info_2, &all_data[1].message_info, sizeof(rmw_message_info_t)), + 0); + } } /* - Consume data from an intra-process buffer with an implementations that stores unique_ptr - Messages are inserted using the same data as the implementation, i.e. unique_ptr - - Request shared_ptr no copies are expected - - Request unique_ptr no copies are expected + Get all data from an intra-process buffer with an implementations that stores unique_ptr + Messages are inserted as shared_ptr and then as unique_ptr + - Add shared_ptr a copy is expected + - Add unique_ptr no copies are expected */ -TEST(TestIntraProcessBuffer, unique_buffer_consume) { +TEST(TestIntraProcessBuffer, unique_buffer_get_all_data) { using MessageT = char; using Alloc = std::allocator; using Deleter = std::default_delete; - using SharedMessageT = std::shared_ptr; using UniqueMessageT = std::unique_ptr; using UniqueIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, UniqueMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::UniquePtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; - auto buffer_impl = - std::make_unique>(2); + auto buffer_impl = std::make_unique(2); UniqueIntraProcessBufferT intra_process_buffer(std::move(buffer_impl)); - auto original_unique_msg = std::make_unique('a'); - auto original_message_pointer = reinterpret_cast(original_unique_msg.get()); - auto original_value = *original_unique_msg; - - intra_process_buffer.add_unique(std::move(original_unique_msg)); - - SharedMessageT popped_shared_msg; - popped_shared_msg = intra_process_buffer.consume_shared(); - auto popped_message_pointer = reinterpret_cast(popped_shared_msg.get()); - - EXPECT_EQ(original_value, *popped_shared_msg); - EXPECT_EQ(original_message_pointer, popped_message_pointer); - - original_unique_msg = std::make_unique('b'); - original_message_pointer = reinterpret_cast(original_unique_msg.get()); - original_value = *original_unique_msg; - - intra_process_buffer.add_unique(std::move(original_unique_msg)); - - UniqueMessageT popped_unique_msg; - popped_unique_msg = intra_process_buffer.consume_unique(); - popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); - - EXPECT_EQ(original_value, *popped_unique_msg); - EXPECT_EQ(original_message_pointer, popped_message_pointer); + auto original_shared_msg = std::make_shared('a'); + auto original_message_pointer = reinterpret_cast(original_shared_msg.get()); + auto original_value = *original_shared_msg; + rmw_message_info_t original_message_info = generate_message_info<0>(); + intra_process_buffer.add({original_shared_msg, original_message_info}); - original_unique_msg = std::make_unique('c'); - original_message_pointer = reinterpret_cast(original_unique_msg.get()); - original_value = *original_unique_msg; - auto original_unique_msg_2 = std::make_unique('d'); - auto original_message_pointer_2 = reinterpret_cast(original_unique_msg.get()); + auto original_unique_msg_2 = std::make_unique('b'); + auto original_message_pointer_2 = reinterpret_cast(original_unique_msg_2.get()); auto original_value_2 = *original_unique_msg_2; - intra_process_buffer.add_unique(std::move(original_unique_msg)); - intra_process_buffer.add_unique(std::move(original_unique_msg_2)); - - auto shared_data_vec = intra_process_buffer.get_all_data_shared(); - EXPECT_EQ(2L, shared_data_vec.size()); - EXPECT_EQ(1L, shared_data_vec[0].use_count()); - EXPECT_EQ(original_value, *shared_data_vec[0]); - EXPECT_NE(original_message_pointer, reinterpret_cast(shared_data_vec[0].get())); - EXPECT_EQ(1L, shared_data_vec[1].use_count()); - EXPECT_EQ(original_value_2, *shared_data_vec[1]); - EXPECT_NE(original_message_pointer_2, reinterpret_cast(shared_data_vec[1].get())); - - auto unique_data_vec = intra_process_buffer.get_all_data_unique(); - EXPECT_EQ(2L, unique_data_vec.size()); - EXPECT_EQ(1L, shared_data_vec[0].use_count()); - EXPECT_EQ(original_value, *unique_data_vec[0]); - EXPECT_NE(original_message_pointer, reinterpret_cast(unique_data_vec[0].get())); - EXPECT_EQ(1L, shared_data_vec[1].use_count()); - EXPECT_EQ(original_value_2, *unique_data_vec[1]); - EXPECT_NE(original_message_pointer_2, reinterpret_cast(unique_data_vec[1].get())); + rmw_message_info_t original_message_info_2 = generate_message_info<1>(); + intra_process_buffer.add({std::move(original_unique_msg_2), original_message_info_2}); + + auto all_data = intra_process_buffer.get_all_data(); + ASSERT_EQ(2u, all_data.size()); + + ASSERT_TRUE(std::holds_alternative(all_data[0].message)); + { + const auto & unique_msg = std::get(all_data[0].message); + EXPECT_EQ(1L, original_shared_msg.use_count()); + EXPECT_EQ(original_value, *unique_msg); + EXPECT_NE(original_message_pointer, reinterpret_cast(unique_msg.get())); + EXPECT_EQ( + std::memcmp(&original_message_info, &all_data[0].message_info, sizeof(rmw_message_info_t)), + 0); + } + + ASSERT_TRUE(std::holds_alternative(all_data[1].message)); + { + const auto & unique_msg = std::get(all_data[1].message); + EXPECT_EQ(original_value_2, *unique_msg); + // get_all_data() always returns an independent copy of the buffer's own storage. + EXPECT_NE(original_message_pointer_2, reinterpret_cast(unique_msg.get())); + EXPECT_EQ( + std::memcmp( + &original_message_info_2, &all_data[1].message_info, sizeof(rmw_message_info_t)), + 0); + } } /* @@ -305,16 +324,15 @@ TEST(TestIntraProcessBuffer, available_capacity) { using MessageT = char; using Alloc = std::allocator; using Deleter = std::default_delete; - using SharedMessageT = std::shared_ptr; using UniqueMessageT = std::unique_ptr; using UniqueIntraProcessBufferT = rclcpp::experimental::buffers::TypedIntraProcessBuffer< - MessageT, Alloc, Deleter, UniqueMessageT>; + MessageT, Alloc, Deleter, rclcpp::IntraProcessBufferType::UniquePtr>; + using BufferT = rclcpp::experimental::buffers::IntraProcessBufferData; + using BufferImplT = rclcpp::experimental::buffers::RingBufferImplementation; constexpr auto history_depth = 5u; - auto buffer_impl = - std::make_unique>( - history_depth); + auto buffer_impl = std::make_unique(history_depth); UniqueIntraProcessBufferT intra_process_buffer(std::move(buffer_impl)); @@ -323,45 +341,69 @@ TEST(TestIntraProcessBuffer, available_capacity) { auto original_unique_msg = std::make_unique('a'); auto original_message_pointer = reinterpret_cast(original_unique_msg.get()); auto original_value = *original_unique_msg; + rmw_message_info_t original_message_info = generate_message_info<0>(); - intra_process_buffer.add_unique(std::move(original_unique_msg)); + intra_process_buffer.add({std::move(original_unique_msg), original_message_info}); EXPECT_EQ(history_depth - 1u, intra_process_buffer.available_capacity()); - SharedMessageT popped_shared_msg; - popped_shared_msg = intra_process_buffer.consume_shared(); - auto popped_message_pointer = reinterpret_cast(popped_shared_msg.get()); + UniqueMessageT popped_unique_msg; + rmw_message_info_t popped_message_info; + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_unique_msg = std::move(std::get(popped_data.message)); + popped_message_info = popped_data.message_info; + } + auto popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); EXPECT_EQ(history_depth, intra_process_buffer.available_capacity()); - EXPECT_EQ(original_value, *popped_shared_msg); + EXPECT_EQ(original_value, *popped_unique_msg); EXPECT_EQ(original_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); original_unique_msg = std::make_unique('b'); original_message_pointer = reinterpret_cast(original_unique_msg.get()); original_value = *original_unique_msg; + original_message_info = generate_message_info<1>(); - intra_process_buffer.add_unique(std::move(original_unique_msg)); + intra_process_buffer.add({std::move(original_unique_msg), original_message_info}); auto second_unique_msg = std::make_unique('c'); auto second_message_pointer = reinterpret_cast(second_unique_msg.get()); auto second_value = *second_unique_msg; + rmw_message_info_t second_message_info = generate_message_info<2>(); - intra_process_buffer.add_unique(std::move(second_unique_msg)); + intra_process_buffer.add({std::move(second_unique_msg), second_message_info}); EXPECT_EQ(history_depth - 2u, intra_process_buffer.available_capacity()); - UniqueMessageT popped_unique_msg; - popped_unique_msg = intra_process_buffer.consume_unique(); + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_unique_msg = std::move(std::get(popped_data.message)); + popped_message_info = popped_data.message_info; + } popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); EXPECT_EQ(history_depth - 1u, intra_process_buffer.available_capacity()); EXPECT_EQ(original_value, *popped_unique_msg); EXPECT_EQ(original_message_pointer, popped_message_pointer); - - popped_unique_msg = intra_process_buffer.consume_unique(); + EXPECT_EQ(std::memcmp(&original_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); + + { + auto popped_data = intra_process_buffer.consume(); + ASSERT_TRUE(std::holds_alternative(popped_data.message)); + popped_unique_msg = std::move(std::get(popped_data.message)); + popped_message_info = popped_data.message_info; + } popped_message_pointer = reinterpret_cast(popped_unique_msg.get()); EXPECT_EQ(history_depth, intra_process_buffer.available_capacity()); EXPECT_EQ(second_value, *popped_unique_msg); EXPECT_EQ(second_message_pointer, popped_message_pointer); + EXPECT_EQ(std::memcmp(&second_message_info, &popped_message_info, sizeof(rmw_message_info_t)), + 0); } diff --git a/rclcpp/test/rclcpp/test_intra_process_manager.cpp b/rclcpp/test/rclcpp/test_intra_process_manager.cpp index 838b8298b8..844af0d398 100644 --- a/rclcpp/test/rclcpp/test_intra_process_manager.cpp +++ b/rclcpp/test/rclcpp/test_intra_process_manager.cpp @@ -23,6 +23,7 @@ #define RCLCPP_BUILDING_LIBRARY 1 #include "rclcpp/allocator/allocator_common.hpp" #include "rclcpp/context.hpp" +#include "rclcpp/experimental/buffers/intra_process_buffer_data.hpp" #include "rclcpp/macros.hpp" #include "rclcpp/qos.hpp" #include "rmw/types.h" @@ -61,40 +62,38 @@ namespace buffers { namespace mock { + template< typename MessageT, - typename Alloc = std::allocator, typename MessageDeleter = std::default_delete> class IntraProcessBuffer : public IntraProcessBufferBase { public: - using ConstMessageSharedPtr = std::shared_ptr; - using MessageUniquePtr = std::unique_ptr; + using Data = IntraProcessBufferData; RCLCPP_SMART_PTR_DEFINITIONS(IntraProcessBuffer) IntraProcessBuffer() {} - void add(ConstMessageSharedPtr msg) - { - message_ptr = reinterpret_cast(msg.get()); - shared_msg = msg; - ++num_msgs; - } - - void add(MessageUniquePtr msg) + void add(Data d) { - message_ptr = reinterpret_cast(msg.get()); - unique_msg = std::move(msg); + data = std::move(d); + std::visit( + [this](auto && msg) { + message_ptr = reinterpret_cast(msg.get()); + }, data->message); + message_info = data->message_info; ++num_msgs; } - void pop(std::uintptr_t & msg_ptr) + std::pair pop() { - msg_ptr = message_ptr; + std::pair ret = {message_ptr, message_info}; message_ptr = 0; + message_info = rmw_message_info_t{}; --num_msgs; + return ret; } size_t size() const @@ -102,33 +101,20 @@ class IntraProcessBuffer : public IntraProcessBufferBase return num_msgs; } - std::vector get_all_data_shared() + std::vector get_all_data() const { - if (shared_msg) { - return {shared_msg}; - } else if (unique_msg) { - return {std::make_shared(*unique_msg)}; + if (data) { + return {*data}; } return {}; } - std::vector get_all_data_unique() - { - std::vector result; - if (shared_msg) { - result.push_back(std::make_unique(*shared_msg)); - } else if (unique_msg) { - result.push_back(std::make_unique(*unique_msg)); - } - return result; - } - private: // need to store the messages somewhere otherwise the memory address will be reused - ConstMessageSharedPtr shared_msg; - MessageUniquePtr unique_msg; + std::optional data; - std::uintptr_t message_ptr; + std::uintptr_t message_ptr{}; + rmw_message_info_t message_info{}; // count add and pop size_t num_msgs = 0u; }; @@ -199,6 +185,12 @@ class PublisherBase return qos_profile.durability() == rclcpp::DurabilityPolicy::TransientLocal; } + void + set_gid(const rmw_gid_t & gid) + { + gid_ = gid; + } + const rmw_gid_t & get_gid() const { @@ -206,15 +198,15 @@ class PublisherBase } bool - operator==([[maybe_unused]] const rmw_gid_t & gid) const + operator==(const rmw_gid_t & gid) const { - return false; + return std::memcmp(&gid_, &gid, sizeof(rmw_gid_t)) == 0; } bool - operator==([[maybe_unused]] const rmw_gid_t * gid) const + operator==(const rmw_gid_t * gid) const { - return false; + return std::memcmp(&gid_, gid, sizeof(rmw_gid_t)) == 0; } uint64_t intra_process_publisher_id_; @@ -322,42 +314,52 @@ class SubscriptionIntraProcessBuffer : public SubscriptionIntraProcessBase public: RCLCPP_SMART_PTR_DEFINITIONS(SubscriptionIntraProcessBuffer) + using ConstMessageSharedPtr = std::shared_ptr; + using MessageUniquePtr = std::unique_ptr; + + using IntraProcessBuffer = typename rclcpp::experimental::buffers::mock::IntraProcessBuffer< + MessageT, + Deleter + >; + explicit SubscriptionIntraProcessBuffer(const std::string & topic, const rclcpp::QoS & qos) : SubscriptionIntraProcessBase(nullptr, topic, qos), take_shared_method(false) { - buffer = std::make_unique>(); - } - - void - provide_intra_process_message(std::shared_ptr msg) - { - buffer->add(msg); + buffer = std::make_unique(); } - void - provide_intra_process_message(std::unique_ptr msg) + explicit SubscriptionIntraProcessBuffer(const rclcpp::QoS & qos = rclcpp::QoS(10)) + : SubscriptionIntraProcessBase(nullptr, "topic", qos), take_shared_method(false) { - buffer->add(std::move(msg)); + buffer = std::make_unique(); } void - provide_intra_process_data(std::shared_ptr msg) + provide_intra_process_message( + std::variant message, + const rmw_message_info_t & message_info) { - buffer->add(msg); + typename IntraProcessBuffer::Data data; + data.message_info = message_info; + data.message = std::move(message); + buffer->add(std::move(data)); } void - provide_intra_process_data(std::unique_ptr msg) + provide_intra_process_data( + std::variant message, + const rmw_message_info_t & message_info) { - buffer->add(std::move(msg)); + typename IntraProcessBuffer::Data data; + data.message_info = message_info; + data.message = std::move(message); + buffer->add(std::move(data)); } - std::uintptr_t + std::pair pop() { - std::uintptr_t ptr; - buffer->pop(ptr); - return ptr; + return buffer->pop(); } bool @@ -394,6 +396,11 @@ class SubscriptionIntraProcess : public SubscriptionIntraProcessBuffer< : SubscriptionIntraProcessBuffer(topic, qos) { } + + explicit SubscriptionIntraProcess(const rclcpp::QoS & qos = rclcpp::QoS(10)) + : SubscriptionIntraProcessBuffer(qos) + { + } }; } // namespace mock @@ -443,11 +450,14 @@ void Publisher::publish(MessageUniquePtr msg) } if (buffer) { - auto shared_msg = ipm->template do_intra_process_publish_and_return_shared( + auto pair = ipm->template do_intra_process_publish_and_return_shared( intra_process_publisher_id_, std::move(msg), *message_allocator_); - buffer->add(shared_msg); + experimental::buffers::IntraProcessBufferData data; + data.message = pair.first; + data.message_info = pair.second; + buffer->add(std::move(data)); } else { ipm->template do_intra_process_publish( intra_process_publisher_id_, @@ -564,15 +574,29 @@ TEST(TestIntraProcessManager, single_subscription) { auto p1_id = ipm->add_publisher(p1); p1->set_intra_process_manager(p1_id, ipm); - auto s1 = std::make_shared("topic", rclcpp::QoS(10)); + rmw_gid_t p1_gid = {"test", {1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16}}; + p1->set_gid(p1_gid); + + auto s1 = std::make_shared(); s1->take_shared_method = false; auto s1_id = ipm->template add_subscription(s1); auto unique_msg = std::make_unique(); auto original_message_pointer = reinterpret_cast(unique_msg.get()); + rmw_time_point_value_t last_timestamp; p1->publish(std::move(unique_msg)); - auto received_message_pointer_1 = s1->pop(); - ASSERT_EQ(original_message_pointer, received_message_pointer_1); + { + auto [received_message_pointer_1, received_message_info_1] = s1->pop(); + ASSERT_EQ(original_message_pointer, received_message_pointer_1); + ASSERT_TRUE(received_message_info_1.from_intra_process); + ASSERT_EQ(received_message_info_1.publication_sequence_number, 0L); + ASSERT_EQ(received_message_info_1.reception_sequence_number, 0L); + ASSERT_NE(received_message_info_1.source_timestamp, 0L); + ASSERT_NE(received_message_info_1.received_timestamp, 0L); + ASSERT_EQ(received_message_info_1.source_timestamp, received_message_info_1.received_timestamp); + ASSERT_EQ(memcmp(&p1_gid, &received_message_info_1.publisher_gid, sizeof(rmw_gid_t)), 0); + last_timestamp = received_message_info_1.source_timestamp; + } ipm->remove_subscription(s1_id); auto s2 = std::make_shared("topic", rclcpp::QoS(10)); @@ -583,16 +607,35 @@ TEST(TestIntraProcessManager, single_subscription) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - received_message_pointer_1 = s1->pop(); - auto received_message_pointer_2 = s2->pop(); - ASSERT_EQ(original_message_pointer, received_message_pointer_2); - ASSERT_EQ(0u, received_message_pointer_1); + { + auto [received_message_pointer_1, received_message_info_1] = s1->pop(); + auto [received_message_pointer_2, received_message_info_2] = s2->pop(); + ASSERT_EQ(original_message_pointer, received_message_pointer_2); + ASSERT_TRUE(received_message_info_2.from_intra_process); + ASSERT_EQ(received_message_info_2.publication_sequence_number, 1L); + ASSERT_EQ(received_message_info_2.reception_sequence_number, 0L); + ASSERT_GT(received_message_info_2.source_timestamp, last_timestamp); + ASSERT_GT(received_message_info_2.received_timestamp, last_timestamp); + ASSERT_EQ(received_message_info_2.source_timestamp, received_message_info_2.received_timestamp); + ASSERT_EQ(memcmp(&p1_gid, &received_message_info_2.publisher_gid, sizeof(rmw_gid_t)), 0); + ASSERT_EQ(0u, received_message_pointer_1); + last_timestamp = received_message_info_2.source_timestamp; + } unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - received_message_pointer_2 = s2->pop(); - ASSERT_EQ(original_message_pointer, received_message_pointer_2); + { + auto [received_message_pointer_2, received_message_info_2] = s2->pop(); + ASSERT_EQ(original_message_pointer, received_message_pointer_2); + ASSERT_TRUE(received_message_info_2.from_intra_process); + ASSERT_EQ(received_message_info_2.publication_sequence_number, 2L); + ASSERT_EQ(received_message_info_2.reception_sequence_number, 1L); + ASSERT_GT(received_message_info_2.source_timestamp, last_timestamp); + ASSERT_GT(received_message_info_2.received_timestamp, last_timestamp); + ASSERT_EQ(received_message_info_2.source_timestamp, received_message_info_2.received_timestamp); + ASSERT_EQ(memcmp(&p1_gid, &received_message_info_2.publisher_gid, sizeof(rmw_gid_t)), 0); + } } /* @@ -629,8 +672,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_same_type) { auto unique_msg = std::make_unique(); auto original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - bool received_original_1 = s1->pop() == original_message_pointer; - bool received_original_2 = s2->pop() == original_message_pointer; + bool received_original_1 = s1->pop().first == original_message_pointer; + bool received_original_2 = s2->pop().first == original_message_pointer; std::vector received_original_vec = {received_original_1, received_original_2}; ASSERT_THAT(received_original_vec, UnorderedElementsAre(true, false)); @@ -649,8 +692,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_same_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_3 = s3->pop(); - auto received_message_pointer_4 = s4->pop(); + auto received_message_pointer_3 = s3->pop().first; + auto received_message_pointer_4 = s4->pop().first; ASSERT_EQ(original_message_pointer, received_message_pointer_3); ASSERT_EQ(original_message_pointer, received_message_pointer_4); @@ -668,8 +711,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_same_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_5 = s5->pop(); - auto received_message_pointer_6 = s6->pop(); + auto received_message_pointer_5 = s5->pop().first; + auto received_message_pointer_6 = s6->pop().first; ASSERT_NE(original_message_pointer, received_message_pointer_5); // Someone gets the original unique_ptr, the last one to take. ASSERT_EQ(original_message_pointer, received_message_pointer_6); @@ -690,8 +733,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_same_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_7 = s7->pop(); - auto received_message_pointer_8 = s8->pop(); + auto received_message_pointer_7 = s7->pop().first; + auto received_message_pointer_8 = s8->pop().first; ASSERT_EQ(original_message_pointer, received_message_pointer_7); ASSERT_EQ(original_message_pointer, received_message_pointer_8); } @@ -735,8 +778,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_different_type) { auto unique_msg = std::make_unique(); auto original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_1 = s1->pop(); - auto received_message_pointer_2 = s2->pop(); + auto received_message_pointer_1 = s1->pop().first; + auto received_message_pointer_2 = s2->pop().first; ASSERT_NE(original_message_pointer, received_message_pointer_1); ASSERT_EQ(original_message_pointer, received_message_pointer_2); @@ -758,9 +801,9 @@ TEST(TestIntraProcessManager, multiple_subscriptions_different_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_3 = s3->pop(); - auto received_message_pointer_4 = s4->pop(); - auto received_message_pointer_5 = s5->pop(); + auto received_message_pointer_3 = s3->pop().first; + auto received_message_pointer_4 = s4->pop().first; + auto received_message_pointer_5 = s5->pop().first; bool received_original_3 = received_message_pointer_3 == original_message_pointer; bool received_original_4 = received_message_pointer_4 == original_message_pointer; bool received_original_5 = received_message_pointer_5 == original_message_pointer; @@ -794,10 +837,10 @@ TEST(TestIntraProcessManager, multiple_subscriptions_different_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_6 = s6->pop(); - auto received_message_pointer_7 = s7->pop(); - auto received_message_pointer_8 = s8->pop(); - auto received_message_pointer_9 = s9->pop(); + auto received_message_pointer_6 = s6->pop().first; + auto received_message_pointer_7 = s7->pop().first; + auto received_message_pointer_8 = s8->pop().first; + auto received_message_pointer_9 = s9->pop().first; bool received_original_8 = received_message_pointer_8 == original_message_pointer; bool received_original_9 = received_message_pointer_9 == original_message_pointer; received_original_vec = {received_original_8, received_original_9}; @@ -826,8 +869,8 @@ TEST(TestIntraProcessManager, multiple_subscriptions_different_type) { unique_msg = std::make_unique(); original_message_pointer = reinterpret_cast(unique_msg.get()); p1->publish(std::move(unique_msg)); - auto received_message_pointer_10 = s10->pop(); - auto received_message_pointer_11 = s11->pop(); + auto received_message_pointer_10 = s10->pop().first; + auto received_message_pointer_11 = s11->pop().first; EXPECT_EQ(original_message_pointer, received_message_pointer_10); EXPECT_NE(original_message_pointer, received_message_pointer_11); } @@ -985,9 +1028,9 @@ TEST(TestIntraProcessManager, transient_local) { ipm->template add_subscription(s2); ipm->template add_subscription(s3); - auto received_message_pointer_1 = s1->pop(); - auto received_message_pointer_2 = s2->pop(); - auto received_message_pointer_3 = s3->pop(); + auto received_message_pointer_1 = s1->pop().first; + auto received_message_pointer_2 = s2->pop().first; + auto received_message_pointer_3 = s3->pop().first; ASSERT_NE(0u, received_message_pointer_1); ASSERT_NE(0u, received_message_pointer_2); ASSERT_NE(0u, received_message_pointer_3);