15 #include "rclcpp/experimental/intra_process_manager.hpp"
23 namespace experimental
26 static std::atomic<uint64_t> _next_unique_id {1};
28 IntraProcessManager::IntraProcessManager()
31 IntraProcessManager::~IntraProcessManager()
36 const rclcpp::PublisherBase::SharedPtr & publisher,
37 const rclcpp::experimental::buffers::IntraProcessBufferBase::SharedPtr & buffer)
39 std::unique_lock<std::shared_timed_mutex> lock(mutex_);
41 uint64_t pub_id = IntraProcessManager::get_next_unique_id();
43 publishers_[pub_id] = publisher;
44 if (publisher->is_durability_transient_local()) {
46 publisher_buffers_[pub_id] = buffer;
48 throw std::runtime_error(
49 "transient_local publisher needs to pass"
50 "a valid publisher buffer ptr when calling add_publisher()");
55 gid_to_publisher_info_[publisher->get_gid()] = {pub_id, publisher};
58 pub_to_subs_[pub_id] = SplittedSubscriptions();
61 for (
auto & pair : subscriptions_) {
62 auto subscription = pair.second.lock();
66 if (can_communicate(publisher, subscription)) {
67 uint64_t sub_id = pair.first;
68 insert_sub_id_for_pub(sub_id, pub_id, subscription->use_take_shared_method());
78 std::unique_lock<std::shared_timed_mutex> lock(mutex_);
80 subscriptions_.erase(intra_process_subscription_id);
82 for (
auto & pair : pub_to_subs_) {
83 pair.second.take_shared_subscriptions.erase(
85 pair.second.take_shared_subscriptions.begin(),
86 pair.second.take_shared_subscriptions.end(),
87 intra_process_subscription_id),
88 pair.second.take_shared_subscriptions.end());
90 pair.second.take_ownership_subscriptions.erase(
92 pair.second.take_ownership_subscriptions.begin(),
93 pair.second.take_ownership_subscriptions.end(),
94 intra_process_subscription_id),
95 pair.second.take_ownership_subscriptions.end());
102 std::unique_lock<std::shared_timed_mutex> lock(mutex_);
106 auto pub_it = publishers_.find(intra_process_publisher_id);
107 if (pub_it != publishers_.end()) {
108 auto publisher = pub_it->second.lock();
110 gid_to_publisher_info_.erase(publisher->get_gid());
113 for (
auto git = gid_to_publisher_info_.begin(); git != gid_to_publisher_info_.end(); ++git) {
114 if (git->second.pub_id == intra_process_publisher_id) {
115 gid_to_publisher_info_.erase(git);
122 publishers_.erase(intra_process_publisher_id);
123 publisher_buffers_.erase(intra_process_publisher_id);
124 pub_to_subs_.erase(intra_process_publisher_id);
130 std::shared_lock<std::shared_timed_mutex> lock(mutex_);
133 auto it = gid_to_publisher_info_.find(*
id);
134 if (it == gid_to_publisher_info_.end()) {
139 auto publisher = it->second.publisher.lock();
140 return publisher !=
nullptr;
146 std::shared_lock<std::shared_timed_mutex> lock(mutex_);
148 auto publisher_it = pub_to_subs_.find(intra_process_publisher_id);
149 if (publisher_it == pub_to_subs_.end()) {
153 "Calling get_subscription_count for invalid or no longer existing publisher id");
158 publisher_it->second.take_shared_subscriptions.size() +
159 publisher_it->second.take_ownership_subscriptions.size();
164 SubscriptionIntraProcessBase::SharedPtr
165 IntraProcessManager::get_subscription_intra_process(uint64_t intra_process_subscription_id)
167 std::shared_lock<std::shared_timed_mutex> lock(mutex_);
169 auto subscription_it = subscriptions_.find(intra_process_subscription_id);
170 if (subscription_it == subscriptions_.end()) {
173 auto subscription = subscription_it->second.lock();
177 subscriptions_.erase(subscription_it);
184 IntraProcessManager::get_next_unique_id()
186 auto next_id = _next_unique_id.fetch_add(1, std::memory_order_relaxed);
197 throw std::overflow_error(
198 "exhausted the unique id's for publishers and subscribers in this process "
199 "(congratulations your computer is either extremely fast or extremely old)");
206 IntraProcessManager::insert_sub_id_for_pub(
209 bool use_take_shared_method)
211 if (use_take_shared_method) {
212 pub_to_subs_[pub_id].take_shared_subscriptions.push_back(sub_id);
214 pub_to_subs_[pub_id].take_ownership_subscriptions.push_back(sub_id);
219 IntraProcessManager::can_communicate(
220 const rclcpp::PublisherBase::SharedPtr & pub,
221 const rclcpp::experimental::SubscriptionIntraProcessBase::SharedPtr & sub)
const
224 if (strcmp(pub->get_topic_name(), sub->get_topic_name()) != 0) {
229 if (check_result.compatibility == rclcpp::QoSCompatibility::Error) {
239 size_t capacity = std::numeric_limits<size_t>::max();
241 auto publisher_it = pub_to_subs_.find(intra_process_publisher_id);
242 if (publisher_it == pub_to_subs_.end()) {
246 "Calling lowest_available_capacity for invalid or no longer existing publisher id");
250 if (publisher_it->second.take_shared_subscriptions.empty() &&
251 publisher_it->second.take_ownership_subscriptions.empty())
257 auto available_capacity = [
this, &capacity](
const uint64_t intra_process_subscription_id)
259 auto subscription_it = subscriptions_.find(intra_process_subscription_id);
260 if (subscription_it != subscriptions_.end()) {
261 auto subscription = subscription_it->second.lock();
263 capacity = std::min(capacity, subscription->available_capacity());
269 "Calling available_capacity for invalid or no longer existing subscription id");
273 for (
const auto sub_id : publisher_it->second.take_shared_subscriptions) {
274 available_capacity(sub_id);
277 for (
const auto sub_id : publisher_it->second.take_ownership_subscriptions) {
278 available_capacity(sub_id);
RCLCPP_PUBLIC bool matches_any_publishers(const rmw_gid_t *id) const
Return true if the given rmw_gid_t matches any stored Publishers.
RCLCPP_PUBLIC void remove_subscription(uint64_t intra_process_subscription_id)
Unregister a subscription using the subscription's unique id.
RCLCPP_PUBLIC void remove_publisher(uint64_t intra_process_publisher_id)
Unregister a publisher using the publisher's unique id.
RCLCPP_PUBLIC size_t get_subscription_count(uint64_t intra_process_publisher_id) const
Return the number of intraprocess subscriptions that are matched with a given publisher id.
RCLCPP_PUBLIC size_t lowest_available_capacity(const uint64_t intra_process_publisher_id) const
Return the lowest available capacity for all subscription buffers for a publisher id.
RCLCPP_PUBLIC uint64_t add_publisher(const rclcpp::PublisherBase::SharedPtr &publisher, const rclcpp::experimental::buffers::IntraProcessBufferBase::SharedPtr &buffer=rclcpp::experimental::buffers::IntraProcessBufferBase::SharedPtr())
Register a publisher with the manager, returns the publisher unique id.
Versions of rosidl_typesupport_cpp::get_message_type_support_handle that handle adapted types.
RCLCPP_PUBLIC QoSCheckCompatibleResult qos_check_compatible(const QoS &publisher_qos, const QoS &subscription_qos)
Check if two QoS profiles are compatible.
RCLCPP_PUBLIC Logger get_logger(const std::string &name)
Return a named logger.