From c3b20c014e0ca0741c733abcf5d461b0132229df Mon Sep 17 00:00:00 2001 From: Aaravanand Date: Tue, 29 Sep 2026 04:51:04 +0000 Subject: [PATCH 1/2] Fix loaned message ownership transfer Signed-off-by: Aaravanand --- rclcpp/include/rclcpp/subscription.hpp | 25 ++++++++++++++++----- rclcpp/include/rclcpp/subscription_base.hpp | 11 +++++++++ rclcpp/src/rclcpp/executor.cpp | 15 ++----------- rclcpp/src/rclcpp/subscription_base.cpp | 18 +++------------ rclcpp/test/rclcpp/test_subscription.cpp | 8 ------- 5 files changed, 36 insertions(+), 41 deletions(-) diff --git a/rclcpp/include/rclcpp/subscription.hpp b/rclcpp/include/rclcpp/subscription.hpp index 1f0a1c51a7..101d2e5abe 100644 --- a/rclcpp/include/rclcpp/subscription.hpp +++ b/rclcpp/include/rclcpp/subscription.hpp @@ -390,17 +390,32 @@ class Subscription : public SubscriptionBase void * loaned_message, const rclcpp::MessageInfo & message_info) override { + auto typed_message = static_cast(loaned_message); + + auto sub_handle = this->get_subscription_handle(); + auto loan_mutex = this->loaned_message_mutex_; + auto sptr = std::shared_ptr( + typed_message, [sub_handle, loan_mutex](ROSMessageType * msg) { + if (msg) { + std::lock_guard lock(*loan_mutex); + rcl_ret_t ret = rcl_return_loaned_message_from_subscription( + sub_handle.get(), static_cast(msg)); + if (RCL_RET_OK != ret) { + RCLCPP_ERROR( + rclcpp::get_logger("rclcpp"), + "rcl_return_loaned_message_from_subscription() failed: %s", + rcl_get_error_string().str); + rcl_reset_error(); + } + } + }); + if (matches_any_intra_process_publishers(&message_info.get_rmw_message_info().publisher_gid)) { // In this case, the message will be delivered via intra process and // we should ignore this copy of the message. return; } - auto typed_message = static_cast(loaned_message); - // message is loaned, so we have to make sure that the deleter does not deallocate the message - auto sptr = std::shared_ptr( - typed_message, [](ROSMessageType * msg) {(void) msg;}); - std::chrono::time_point now; if (subscription_topic_statistics_) { // get current time before executing callback to diff --git a/rclcpp/include/rclcpp/subscription_base.hpp b/rclcpp/include/rclcpp/subscription_base.hpp index ccf5cca1d5..c2f8395cc5 100644 --- a/rclcpp/include/rclcpp/subscription_base.hpp +++ b/rclcpp/include/rclcpp/subscription_base.hpp @@ -650,8 +650,16 @@ class SubscriptionBase : public std::enable_shared_from_this take_dynamic_message( rclcpp::dynamic_typesupport::DynamicMessage & message_out, rclcpp::MessageInfo & message_info_out); + + RCLCPP_PUBLIC + std::shared_ptr + get_loaned_message_mutex() const + { + return loaned_message_mutex_; + } // =============================================================================================== + protected: template void @@ -688,6 +696,9 @@ class SubscriptionBase : public std::enable_shared_from_this std::shared_ptr node_handle_; std::recursive_mutex on_new_message_callback_mutex_; + + /// Mutex to protect rcl_take_loaned_message and rcl_return_loaned_message_from_subscription + std::shared_ptr loaned_message_mutex_; // It is important to declare on_new_message_callback_ before // subscription_handle_, so on destruction the subscription is // destroyed first. Otherwise, the rmw subscription callback diff --git a/rclcpp/src/rclcpp/executor.cpp b/rclcpp/src/rclcpp/executor.cpp index 90f3277bc7..6e19a03c6d 100644 --- a/rclcpp/src/rclcpp/executor.cpp +++ b/rclcpp/src/rclcpp/executor.cpp @@ -575,6 +575,8 @@ Executor::execute_subscription(const rclcpp::SubscriptionBase::SharedPtr & subsc subscription->get_topic_name(), [&]() { + std::shared_ptr loan_mutex = subscription->get_loaned_message_mutex(); + std::lock_guard lock(*loan_mutex); rcl_ret_t ret = rcl_take_loaned_message( subscription->get_subscription_handle().get(), &loaned_msg, @@ -589,19 +591,6 @@ Executor::execute_subscription(const rclcpp::SubscriptionBase::SharedPtr & subsc return true; }, [&]() {subscription->handle_loaned_message(loaned_msg, message_info);}); - if (nullptr != loaned_msg) { - rcl_ret_t ret = rcl_return_loaned_message_from_subscription( - subscription->get_subscription_handle().get(), loaned_msg); - if (RCL_RET_OK != ret) { - RCLCPP_ERROR( - rclcpp::get_logger("rclcpp"), - "rcl_return_loaned_message_from_subscription() failed for subscription on topic " - "'%s': %s", - subscription->get_topic_name(), rcl_get_error_string().str); - rcl_reset_error(); - } - loaned_msg = nullptr; - } } else { // This case is taking a copy of the message data from the middleware via // inter-process communication. diff --git a/rclcpp/src/rclcpp/subscription_base.cpp b/rclcpp/src/rclcpp/subscription_base.cpp index d286358cfb..547d4b5ce4 100644 --- a/rclcpp/src/rclcpp/subscription_base.cpp +++ b/rclcpp/src/rclcpp/subscription_base.cpp @@ -53,7 +53,8 @@ SubscriptionBase::SubscriptionBase( intra_process_subscription_id_(0), event_callbacks_(event_callbacks), type_support_(type_support_handle), - delivered_message_kind_(delivered_message_kind) + delivered_message_kind_(delivered_message_kind), + loaned_message_mutex_(std::make_shared()) { auto custom_deletor = [node_handle = this->node_handle_](rcl_subscription_t * rcl_subs) { @@ -337,20 +338,7 @@ SubscriptionBase::setup_intra_process( bool SubscriptionBase::can_loan_messages() const { - bool retval = rcl_subscription_can_loan_messages(subscription_handle_.get()); - if (retval) { - // TODO(clalancette): The loaned message interface is currently not safe to use with - // shared_ptr callbacks. If a user takes a copy of the shared_ptr, it can get freed from - // underneath them via rcl_return_loaned_message_from_subscription(). The correct solution is - // to return the loaned message in a custom deleter, but that needs to be carefully handled - // with locking. Warn the user about this until we fix it. - RCLCPP_WARN_ONCE( - this->node_logger_, - "Loaned messages are only safe with const ref subscription callbacks. " - "If you are using any other kind of subscriptions, " - "set the ROS_DISABLE_LOANED_MESSAGES environment variable to 1 (the default)."); - } - return retval; + return rcl_subscription_can_loan_messages(subscription_handle_.get()); } rclcpp::Waitable::SharedPtr diff --git a/rclcpp/test/rclcpp/test_subscription.cpp b/rclcpp/test/rclcpp/test_subscription.cpp index c7dc473061..5648af3750 100644 --- a/rclcpp/test/rclcpp/test_subscription.cpp +++ b/rclcpp/test/rclcpp/test_subscription.cpp @@ -351,15 +351,7 @@ TEST_F(TestSubscription, rcl_subscription_get_publisher_count_error) { EXPECT_THROW(sub->get_publisher_count(), rclcpp::exceptions::RCLError); } -TEST_F(TestSubscription, handle_loaned_message) { - initialize(); - auto callback = [](std::shared_ptr) {}; - auto sub = node_->create_subscription("topic", 10, callback); - test_msgs::msg::Empty msg; - rclcpp::MessageInfo message_info; - EXPECT_NO_THROW(sub->handle_loaned_message(&msg, message_info)); -} /* Testing on_new_message callbacks. From f577e7fe4ed560b7887d6fdf11959b973a220bb6 Mon Sep 17 00:00:00 2001 From: Aaravanand Date: Mon, 5 Oct 2026 18:32:39 +0000 Subject: [PATCH 2/2] Address review feedback: hide mutex and fix callback lifetime Signed-off-by: Aaravanand --- rclcpp/include/rclcpp/subscription_base.hpp | 9 ++++----- rclcpp/src/rclcpp/executor.cpp | 8 ++------ rclcpp/src/rclcpp/subscription_base.cpp | 14 ++++++++++++++ 3 files changed, 20 insertions(+), 11 deletions(-) diff --git a/rclcpp/include/rclcpp/subscription_base.hpp b/rclcpp/include/rclcpp/subscription_base.hpp index c2f8395cc5..24254fd0f1 100644 --- a/rclcpp/include/rclcpp/subscription_base.hpp +++ b/rclcpp/include/rclcpp/subscription_base.hpp @@ -652,11 +652,10 @@ class SubscriptionBase : public std::enable_shared_from_this rclcpp::MessageInfo & message_info_out); RCLCPP_PUBLIC - std::shared_ptr - get_loaned_message_mutex() const - { - return loaned_message_mutex_; - } + rcl_ret_t + take_loaned_message( + void ** loaned_message, + rmw_message_info_t * message_info_out); // =============================================================================================== diff --git a/rclcpp/src/rclcpp/executor.cpp b/rclcpp/src/rclcpp/executor.cpp index 6e19a03c6d..e9e4320ac3 100644 --- a/rclcpp/src/rclcpp/executor.cpp +++ b/rclcpp/src/rclcpp/executor.cpp @@ -575,13 +575,9 @@ Executor::execute_subscription(const rclcpp::SubscriptionBase::SharedPtr & subsc subscription->get_topic_name(), [&]() { - std::shared_ptr loan_mutex = subscription->get_loaned_message_mutex(); - std::lock_guard lock(*loan_mutex); - rcl_ret_t ret = rcl_take_loaned_message( - subscription->get_subscription_handle().get(), + rcl_ret_t ret = subscription->take_loaned_message( &loaned_msg, - &message_info.get_rmw_message_info(), - nullptr); + &message_info.get_rmw_message_info()); TRACETOOLS_TRACEPOINT(rclcpp_take, static_cast(loaned_msg)); if (RCL_RET_SUBSCRIPTION_TAKE_FAILED == ret) { return false; diff --git a/rclcpp/src/rclcpp/subscription_base.cpp b/rclcpp/src/rclcpp/subscription_base.cpp index 547d4b5ce4..d09a26b8d8 100644 --- a/rclcpp/src/rclcpp/subscription_base.cpp +++ b/rclcpp/src/rclcpp/subscription_base.cpp @@ -96,6 +96,7 @@ SubscriptionBase::SubscriptionBase( SubscriptionBase::~SubscriptionBase() { + this->clear_on_new_message_callback(); if (!use_intra_process_) { return; } @@ -579,6 +580,19 @@ SubscriptionBase::take_dynamic_message( return false; } +rcl_ret_t +SubscriptionBase::take_loaned_message( + void ** loaned_message, + rmw_message_info_t * message_info_out) +{ + std::lock_guard lock(*loaned_message_mutex_); + return rcl_take_loaned_message( + subscription_handle_.get(), + loaned_message, + message_info_out, + nullptr); +} + void SubscriptionBase::disable_callbacks() {