Skip to content
Open
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
233 changes: 50 additions & 183 deletions rclcpp/include/rclcpp/experimental/buffers/intra_process_buffer.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,14 @@
#include <stdexcept>
#include <type_traits>
#include <utility>
#include <variant>
#include <vector>

#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"

Expand All @@ -44,13 +47,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<void>,
typename MessageDeleter = std::default_delete<MessageT>>
class IntraProcessBuffer : public IntraProcessBufferBase
{
Expand All @@ -59,47 +61,38 @@ class IntraProcessBuffer : public IntraProcessBufferBase

virtual ~IntraProcessBuffer() {}

using MessageUniquePtr = std::unique_ptr<MessageT, MessageDeleter>;
using MessageSharedPtr = std::shared_ptr<const MessageT>;
using Data = IntraProcessBufferData<MessageT, MessageDeleter>;

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

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

virtual std::vector<MessageSharedPtr> get_all_data_shared() = 0;
virtual std::vector<MessageUniquePtr> get_all_data_unique() = 0;
virtual std::vector<Data> get_all_data() = 0;
};

template<
typename MessageT,
typename Alloc = std::allocator<void>,
typename MessageDeleter = std::default_delete<MessageT>,
typename BufferT = std::unique_ptr<MessageT>>
class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, MessageDeleter>
IntraProcessBufferType BufferType = IntraProcessBufferType::CallbackDefault>
class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, MessageDeleter>
{
public:
RCLCPP_SMART_PTR_DEFINITIONS(TypedIntraProcessBuffer)

using Buffer = IntraProcessBuffer<MessageT, MessageDeleter>;
using MessageAllocTraits = allocator::AllocRebind<MessageT, Alloc>;
using MessageAlloc = typename MessageAllocTraits::allocator_type;
using MessageUniquePtr = std::unique_ptr<MessageT, MessageDeleter>;
using MessageSharedPtr = std::shared_ptr<const MessageT>;
using Data = typename Buffer::Data;
using MessageSharedPtr = typename Data::MessageSharedPtr;
using MessageUniquePtr = typename Data::MessageUniquePtr;

explicit
TypedIntraProcessBuffer(
std::unique_ptr<BufferImplementationBase<BufferT>> buffer_impl,
std::unique_ptr<BufferImplementationBase<Data>> buffer_impl,
std::shared_ptr<Alloc> allocator = nullptr)
: buffer_(std::move(buffer_impl))
{
bool valid_type = (std::is_same<BufferT, MessageSharedPtr>::value ||
std::is_same<BufferT, MessageUniquePtr>::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<const void *>(buffer_.get()),
Expand All @@ -113,34 +106,19 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa

virtual ~TypedIntraProcessBuffer() {}

void add_shared(MessageSharedPtr msg) override
{
add_shared_impl<BufferT>(std::move(msg));
}

void add_unique(MessageUniquePtr msg) override
{
buffer_->enqueue(std::move(msg));
}

MessageSharedPtr consume_shared() override
void add(Data data) override
{
return consume_shared_impl<BufferT>();
add_impl<BufferType>(std::move(data));
}

MessageUniquePtr consume_unique() override
Data consume() override
{
return consume_unique_impl<BufferT>();
}

std::vector<MessageSharedPtr> get_all_data_shared() override
{
return get_all_data_shared_impl();
return buffer_->dequeue();
}

std::vector<MessageUniquePtr> get_all_data_unique() override
std::vector<Data> get_all_data() override
{
return get_all_data_unique_impl();
return buffer_->get_all_data();
}

bool has_data() const override
Expand All @@ -153,9 +131,9 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa
buffer_->clear();
}

bool use_take_shared_method() const override
IntraProcessBufferType buffer_type() const override
{
return std::is_same<BufferT, MessageSharedPtr>::value;
return BufferType;
}

size_t available_capacity() const override
Expand All @@ -164,163 +142,52 @@ class TypedIntraProcessBuffer : public IntraProcessBuffer<MessageT, Alloc, Messa
}

private:
std::unique_ptr<BufferImplementationBase<BufferT>> buffer_;
std::unique_ptr<BufferImplementationBase<Data>> buffer_;

std::shared_ptr<MessageAlloc> message_allocator_;

// MessageSharedPtr to MessageSharedPtr
template<typename DestinationT>
typename std::enable_if<
std::is_same<DestinationT, MessageSharedPtr>::value
>::type
add_shared_impl(MessageSharedPtr shared_msg)
template<IntraProcessBufferType BufferT>
typename std::enable_if_t<
BufferT == IntraProcessBufferType::CallbackDefault>
add_impl(Data data)
{
buffer_->enqueue(std::move(shared_msg));
buffer_->enqueue(std::move(data));
}

// MessageSharedPtr to MessageUniquePtr
template<typename DestinationT>
typename std::enable_if<
std::is_same<DestinationT, MessageUniquePtr>::value
>::type
add_shared_impl(MessageSharedPtr shared_msg)
template<IntraProcessBufferType BufferT>
typename std::enable_if_t<
BufferT == IntraProcessBufferType::SharedPtr>
add_impl(Data data)
{
// 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<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);
if (std::holds_alternative<MessageSharedPtr>(data.message)) {
buffer_->enqueue(std::move(data));
} else {
unique_msg = MessageUniquePtr(ptr);
// Promote to a shared pointer
auto unique_msg = std::move(std::get<MessageUniquePtr>(data.message));
data.message = MessageSharedPtr(unique_msg.release());
buffer_->enqueue(std::move(data));
}

buffer_->enqueue(std::move(unique_msg));
}

// MessageSharedPtr to MessageSharedPtr
template<typename OriginT>
typename std::enable_if<
std::is_same<OriginT, MessageSharedPtr>::value,
MessageSharedPtr
>::type
consume_shared_impl()
template<IntraProcessBufferType BufferT>
typename std::enable_if_t<
BufferT == IntraProcessBufferType::UniquePtr>
add_impl(Data data)
{
return buffer_->dequeue();
}

// MessageUniquePtr to MessageSharedPtr
template<typename OriginT>
typename std::enable_if<
(std::is_same<OriginT, MessageUniquePtr>::value),
MessageSharedPtr
>::type
consume_shared_impl()
{
// automatic cast from unique ptr to shared ptr
return buffer_->dequeue();
}

// MessageSharedPtr to MessageUniquePtr
template<typename OriginT>
typename std::enable_if<
(std::is_same<OriginT, MessageSharedPtr>::value),
MessageUniquePtr
>::type
consume_unique_impl()
{
MessageSharedPtr buffer_msg = buffer_->dequeue();

MessageUniquePtr unique_msg;
MessageDeleter * deleter = std::get_deleter<MessageDeleter, const MessageT>(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<MessageUniquePtr>(data.message)) {
buffer_->enqueue(std::move(data));
} else {
unique_msg = MessageUniquePtr(ptr);
}

return unique_msg;
}

// MessageUniquePtr to MessageUniquePtr
template<typename OriginT>
typename std::enable_if<
(std::is_same<OriginT, MessageUniquePtr>::value),
MessageUniquePtr
>::type
consume_unique_impl()
{
return buffer_->dequeue();
}

// 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()
{
return buffer_->get_all_data();
}

// 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()
{
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;
}

// 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()
{
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;
auto shared_msg = std::move(std::get<MessageSharedPtr>(data.message));
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);
data.message = MessageUniquePtr(ptr, *deleter);
} else {
unique_msg = MessageUniquePtr(ptr);
data.message = MessageUniquePtr(ptr);
}
result.push_back(std::move(unique_msg));
buffer_->enqueue(std::move(data));
}
return result;
}

// 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()
{
return buffer_->get_all_data();
}
};

Expand Down
Loading