Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#ifndef RCLCPP__EXPERIMENTAL__BUFFERS__BUFFER_IMPLEMENTATION_BASE_HPP_
#define RCLCPP__EXPERIMENTAL__BUFFERS__BUFFER_IMPLEMENTATION_BASE_HPP_

#include <functional>
#include <vector>

namespace rclcpp
Expand All @@ -28,12 +29,14 @@ template<typename BufferT>
class BufferImplementationBase
{
public:
using ForEachFunc = std::function<void(const BufferT &)>;

virtual ~BufferImplementationBase() {}

virtual BufferT dequeue() = 0;
virtual void enqueue(BufferT request) = 0;

virtual std::vector<BufferT> get_all_data() = 0;
virtual void for_each(ForEachFunc && f) const = 0;

virtual void clear() = 0;
virtual bool has_data() const = 0;
Expand Down
104 changes: 53 additions & 51 deletions rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#ifndef RCLCPP__EXPERIMENTAL__BUFFERS__INTRA_PROCESS_BUFFER_HPP_
#define RCLCPP__EXPERIMENTAL__BUFFERS__INTRA_PROCESS_BUFFER_HPP_

#include <functional>
#include <memory>
#include <stdexcept>
#include <type_traits>
Expand Down Expand Up @@ -61,15 +62,20 @@ class IntraProcessBuffer : public IntraProcessBufferBase

using MessageUniquePtr = std::unique_ptr<MessageT, MessageDeleter>;
using MessageSharedPtr = std::shared_ptr<const MessageT>;
using ForEachSharedFunc = std::function<void(const MessageSharedPtr &)>;
// Takes ownership by value: the callee must be able to move from it, and a
// buffer storing unique_ptr can only ever hand out an independent copy, never
// a reference to the ring buffer's own stored element.
using ForEachUniqueFunc = std::function<void(MessageUniquePtr)>;

virtual void add_shared(MessageSharedPtr msg) = 0;
virtual void add_unique(MessageUniquePtr msg) = 0;

virtual MessageSharedPtr consume_shared() = 0;
virtual MessageUniquePtr consume_unique() = 0;

virtual std::vector<MessageSharedPtr> get_all_data_shared() = 0;
virtual std::vector<MessageUniquePtr> get_all_data_unique() = 0;
virtual void for_each_shared(ForEachSharedFunc && f) = 0;
virtual void for_each_unique(ForEachUniqueFunc && f) = 0;
};

template<
Expand All @@ -82,6 +88,7 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa
public:
RCLCPP_SMART_PTR_DEFINITIONS(TypedIntraProcessBuffer)

using Buffer = IntraProcessBuffer<MessageT, Alloc, MessageDeleter>;
using MessageAllocTraits = allocator::AllocRebind<MessageT, Alloc>;
using MessageAlloc = typename MessageAllocTraits::allocator_type;
using MessageUniquePtr = std::unique_ptr<MessageT, MessageDeleter>;
Expand Down Expand Up @@ -133,14 +140,14 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa
return consume_unique_impl<BufferT>();
}

std::vector<MessageSharedPtr> get_all_data_shared() override
void for_each_shared(typename Buffer::ForEachSharedFunc && f) override
{
return get_all_data_shared_impl();
for_each_shared_impl<BufferT>(std::move(f));
}

std::vector<MessageUniquePtr> get_all_data_unique() override
void for_each_unique(typename Buffer::ForEachUniqueFunc && f) override
{
return get_all_data_unique_impl();
for_each_unique_impl<BufferT>(std::move(f));
}

bool has_data() const override
Expand Down Expand Up @@ -260,67 +267,62 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa

// MessageSharedPtr to MessageSharedPtr
template<typename T = BufferT>
typename std::enable_if<
std::is_same<T, MessageSharedPtr>::value,
std::vector<MessageSharedPtr>
>::type
get_all_data_shared_impl()
typename std::enable_if<std::is_same<T, MessageSharedPtr>::value>::type
for_each_shared_impl(typename Buffer::ForEachSharedFunc && f)
{
return buffer_->get_all_data();
buffer_->for_each(std::move(f));
}

// MessageUniquePtr to MessageSharedPtr
template<typename T = BufferT>
typename std::enable_if<
std::is_same<T, MessageUniquePtr>::value,
std::vector<MessageSharedPtr>
>::type
get_all_data_shared_impl()
typename std::enable_if<std::is_same<T, MessageUniquePtr>::value>::type
for_each_shared_impl(typename Buffer::ForEachSharedFunc && f)
{
std::vector<MessageSharedPtr> 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;
// The buffer keeps ownership of each unique_ptr (needed for future transient-local
// replays), so an independent copy must be made for each shared_ptr view handed out.
buffer_->for_each(
[this, &f](const MessageUniquePtr & msg) {
auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1);
MessageAllocTraits::construct(*message_allocator_.get(), ptr, *msg);
MessageSharedPtr shared_msg(ptr, msg.get_deleter());
f(shared_msg);
});
}

// MessageSharedPtr to MessageUniquePtr
template<typename T = BufferT>
typename std::enable_if<
std::is_same<T, MessageSharedPtr>::value,
std::vector<MessageUniquePtr>
>::type
get_all_data_unique_impl()
typename std::enable_if<std::is_same<T, MessageSharedPtr>::value>::type
for_each_unique_impl(typename Buffer::ForEachUniqueFunc && f)
{
std::vector<MessageUniquePtr> 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;
MessageDeleter * deleter = std::get_deleter<MessageDeleter, const MessageT>(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);
} else {
unique_msg = MessageUniquePtr(ptr);
}
result.push_back(std::move(unique_msg));
}
return result;
buffer_->for_each(
[this, &f](const MessageSharedPtr & shared_msg) {
MessageUniquePtr unique_msg;
MessageDeleter * deleter = std::get_deleter<MessageDeleter, const MessageT>(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);
} else {
unique_msg = MessageUniquePtr(ptr);
}
f(std::move(unique_msg));
});
}

// MessageUniquePtr to MessageUniquePtr
template<typename T = BufferT>
typename std::enable_if<
std::is_same<T, MessageUniquePtr>::value,
std::vector<MessageUniquePtr>
>::type
get_all_data_unique_impl()
typename std::enable_if<std::is_same<T, MessageUniquePtr>::value>::type
for_each_unique_impl(typename Buffer::ForEachUniqueFunc && f)
{
return buffer_->get_all_data();
// The buffer keeps ownership of each unique_ptr (needed for future transient-local
// replays), so an independent copy must be made for each unique_ptr handed out.
buffer_->for_each(
[this, &f](const MessageUniquePtr & msg) {
auto ptr = MessageAllocTraits::allocate(*message_allocator_.get(), 1);
MessageAllocTraits::construct(*message_allocator_.get(), ptr, *msg);
MessageUniquePtr unique_msg(ptr, msg.get_deleter());
f(std::move(unique_msg));
});
}
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ template<typename BufferT>
class RingBufferImplementation : public BufferImplementationBase<BufferT>
{
public:
using ImplementationBase = BufferImplementationBase<BufferT>;

explicit RingBufferImplementation(size_t capacity)
: capacity_(capacity),
ring_buffer_(capacity),
Expand Down Expand Up @@ -114,15 +116,16 @@ class RingBufferImplementation : public BufferImplementationBase<BufferT>
return request;
}

/// Get all the elements from the ring buffer
/// Iterates over all the elements in the ring buffer
/**
* This member function is thread-safe.
*
* \return a vector containing all the elements from the ring buffer
* \param f a function that is called for every element in the ring buffer
*/
std::vector<BufferT> get_all_data() override
void for_each(typename ImplementationBase::ForEachFunc && f) const override
{
return get_all_data_impl();
std::lock_guard<std::mutex> lock(mutex_);
return for_each_(std::move(f));
}

/// Get the next index value for the ring buffer
Expand Down Expand Up @@ -236,73 +239,17 @@ class RingBufferImplementation : public BufferImplementationBase<BufferT>
write_index_ = capacity_ - 1;
}

/// Traits for checking if a type is std::unique_ptr
template<typename ...>
struct is_std_unique_ptr final : std::false_type {};
template<class T, typename ... Args>
struct is_std_unique_ptr<std::unique_ptr<T, Args...>> final : std::true_type
{
typedef T Ptr_type;
};

/// Get all the elements from the ring buffer
/// Iterates over all the elements in the ring buffer
/**
* This member function is thread-safe.
* Two versions for the implementation of the function.
* One for buffer containing unique_ptr and the other for other types
* This member function is not thread-safe.
*
* \return a vector containing all the elements from the ring buffer
* \param f a function that is called for every element in the ring buffer
*/
template<typename T = BufferT, std::enable_if_t<is_std_unique_ptr<T>::value &&
std::is_copy_constructible<
typename is_std_unique_ptr<T>::Ptr_type
>::value,
void> * = nullptr>
std::vector<BufferT> get_all_data_impl()
void for_each_(typename ImplementationBase::ForEachFunc && f) const
{
std::lock_guard<std::mutex> lock(mutex_);
std::vector<BufferT> result_vtr;
result_vtr.reserve(size_);
for (size_t id = 0; id < size_; ++id) {
const auto & elem(ring_buffer_[(read_index_ + id) % capacity_]);
if (elem != nullptr) {
result_vtr.emplace_back(new typename is_std_unique_ptr<T>::Ptr_type(
*elem));
} else {
result_vtr.emplace_back(nullptr);
}
f(ring_buffer_[(read_index_ + id) % capacity_]);
}
return result_vtr;
}

template<typename T = BufferT, std::enable_if_t<
std::is_copy_constructible<T>::value, void> * = nullptr>
std::vector<BufferT> get_all_data_impl()
{
std::lock_guard<std::mutex> lock(mutex_);
std::vector<BufferT> result_vtr;
result_vtr.reserve(size_);
for (size_t id = 0; id < size_; ++id) {
result_vtr.emplace_back(ring_buffer_[(read_index_ + id) % capacity_]);
}
return result_vtr;
}

template<typename T = BufferT, std::enable_if_t<!is_std_unique_ptr<T>::value &&
!std::is_copy_constructible<T>::value, void> * = nullptr>
std::vector<BufferT> get_all_data_impl()
{
throw std::logic_error("Underlined type results in invalid get_all_data_impl()");
return {};
}

template<typename T = BufferT, std::enable_if_t<is_std_unique_ptr<T>::value &&
!std::is_copy_constructible<typename is_std_unique_ptr<T>::Ptr_type>::value,
void> * = nullptr>
std::vector<BufferT> get_all_data_impl()
{
throw std::logic_error("Underlined type in unique_ptr results in invalid get_all_data_impl()");
return {};
}

const size_t capacity_;
Expand Down
48 changes: 20 additions & 28 deletions rclcpp/include/rclcpp/experimental/intra_process_manager.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -247,26 +247,14 @@ class IntraProcessManager

this->template add_shared_msg_to_buffers<MessageT, Alloc, Deleter, ROSMessageType>(
msg, sub_ids.take_shared_subscriptions);
} else if (!sub_ids.take_ownership_subscriptions.empty() && // NOLINT
sub_ids.take_shared_subscriptions.size() <= 1)
{
} else if (!sub_ids.concatenated_take_ownership_subscriptions.empty()) {
// 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<uint64_t> concatenated_vector(
sub_ids.take_shared_subscriptions.begin(), sub_ids.take_shared_subscriptions.end());
concatenated_vector.insert(
concatenated_vector.end(),
sub_ids.take_ownership_subscriptions.begin(),
sub_ids.take_ownership_subscriptions.end());
this->template add_owned_msg_to_buffers<MessageT, Alloc, Deleter, ROSMessageType>(
std::move(message),
concatenated_vector,
sub_ids.concatenated_take_ownership_subscriptions,
allocator);
} else if (!sub_ids.take_ownership_subscriptions.empty() && // NOLINT
sub_ids.take_shared_subscriptions.size() > 1)
{
} else {
// Construct a new shared pointer from the message
// for the buffers that do not require ownership
auto shared_msg = std::allocate_shared<MessageT, MessageAllocatorT>(allocator, *message);
Expand Down Expand Up @@ -385,6 +373,10 @@ class IntraProcessManager
{
std::vector<uint64_t> take_shared_subscriptions;
std::vector<uint64_t> take_ownership_subscriptions;
// If there is at maximum 1 buffer that does not require ownership,
// this case is equivalent to all the buffers requiring ownership
// and this vector will contain all of the subscriptions.
std::vector<uint64_t> concatenated_take_ownership_subscriptions;
};

/// Hash function for rmw_gid_t to enable use in unordered_map
Expand Down Expand Up @@ -488,20 +480,20 @@ class IntraProcessManager
"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) {
this->template add_shared_msg_to_buffer<
ROSMessageType, ROSMessageTypeAllocator, ROSMessageTypeDeleter, ROSMessageType>(
shared_data, sub_id);
}
buffer->for_each_shared(
[this, sub_id](const auto & shared_data) {
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) {
auto allocator = ROSMessageTypeAllocator();
this->template add_owned_msg_to_buffer<
ROSMessageType, ROSMessageTypeAllocator, ROSMessageTypeDeleter, ROSMessageType>(
std::move(owned_data), sub_id, allocator);
}
buffer->for_each_unique(
[this, sub_id](auto owned_data) {
auto allocator = ROSMessageTypeAllocator();
this->template add_owned_msg_to_buffer<
ROSMessageType, ROSMessageTypeAllocator, ROSMessageTypeDeleter, ROSMessageType>(
std::move(owned_data), sub_id, allocator);
});
}
}

Expand Down
Loading