15 #include "rclcpp/publisher_base.hpp"
17 #include <rmw/error_handling.h>
26 #include <unordered_map>
30 #include "rcutils/logging_macros.h"
31 #include "rmw/impl/cpp/demangle.hpp"
33 #include "rclcpp/allocator/allocator_common.hpp"
34 #include "rclcpp/allocator/allocator_deleter.hpp"
35 #include "rclcpp/exceptions.hpp"
36 #include "rclcpp/expand_topic_or_service_name.hpp"
37 #include "rclcpp/experimental/intra_process_manager.hpp"
38 #include "rclcpp/logging.hpp"
39 #include "rclcpp/macros.hpp"
40 #include "rclcpp/network_flow_endpoint.hpp"
41 #include "rclcpp/node.hpp"
42 #include "rclcpp/event_handler.hpp"
46 PublisherBase::PublisherBase(
48 const std::string & topic,
49 const rosidl_message_type_support_t & type_support,
52 bool use_default_callbacks)
53 : rcl_node_handle_(node_base->get_shared_rcl_node_handle()),
54 intra_process_is_enabled_(false),
55 intra_process_publisher_id_(0),
56 type_support_(type_support),
57 event_callbacks_(event_callbacks)
59 auto custom_deleter = [node_handle = this->rcl_node_handle_](
rcl_publisher_t * rcl_pub)
64 "Error in destruction of rcl publisher handle: %s",
65 rcl_get_error_string().str);
71 publisher_handle_ = std::shared_ptr<rcl_publisher_t>(
76 publisher_handle_.get(),
77 rcl_node_handle_.get(),
83 auto rcl_node_handle = rcl_node_handle_.get();
92 rclcpp::exceptions::throw_from_rcl_error(ret,
"could not create publisher");
96 if (!publisher_rmw_handle) {
97 auto msg = std::string(
"failed to get rmw handle: ") + rcl_get_error_string().str;
99 throw std::runtime_error(msg);
101 if (rmw_get_gid_for_publisher(publisher_rmw_handle, &rmw_gid_) != RMW_RET_OK) {
102 auto msg = std::string(
"failed to get publisher gid: ") + rmw_get_error_string().str;
104 throw std::runtime_error(msg);
110 PublisherBase::~PublisherBase()
113 event_handlers_.clear();
115 auto ipm = weak_ipm_.lock();
117 if (!intra_process_is_enabled_) {
124 "Intra process manager died before a publisher.");
127 ipm->remove_publisher(intra_process_publisher_id_);
147 if (event_callbacks.deadline_callback) {
148 this->add_event_handler(
149 event_callbacks.deadline_callback,
150 RCL_PUBLISHER_OFFERED_DEADLINE_MISSED);
155 "Failed to add event handler for deadline; not supported");
159 if (event_callbacks.liveliness_callback) {
160 this->add_event_handler(
161 event_callbacks.liveliness_callback,
162 RCL_PUBLISHER_LIVELINESS_LOST);
167 "Failed to add event handler for liveliness; not supported");
170 QOSOfferedIncompatibleQoSCallbackType incompatible_qos_cb;
171 if (event_callbacks.incompatible_qos_callback) {
172 incompatible_qos_cb = event_callbacks.incompatible_qos_callback;
173 }
else if (use_default_callbacks) {
175 incompatible_qos_cb = [
this](QOSOfferedIncompatibleQoSInfo & info) {
176 this->default_incompatible_qos_callback(info);
180 if (incompatible_qos_cb) {
181 this->add_event_handler(incompatible_qos_cb, RCL_PUBLISHER_OFFERED_INCOMPATIBLE_QOS);
186 "Failed to add event handler for incompatible qos; not supported");
189 IncompatibleTypeCallbackType incompatible_type_cb;
190 if (event_callbacks.incompatible_type_callback) {
191 incompatible_type_cb = event_callbacks.incompatible_type_callback;
192 }
else if (use_default_callbacks) {
194 incompatible_type_cb = [
this](IncompatibleTypeInfo & info) {
195 this->default_incompatible_type_callback(info);
199 if (incompatible_type_cb) {
200 this->add_event_handler(incompatible_type_cb, RCL_PUBLISHER_INCOMPATIBLE_TYPE);
205 "Failed to add event handler for incompatible type; not supported");
209 if (event_callbacks.matched_callback) {
210 this->add_event_handler(
211 event_callbacks.matched_callback,
212 RCL_PUBLISHER_MATCHED);
217 "Failed to add event handler for matched; not supported");
225 publisher_handle_.get());
226 if (!publisher_options) {
227 auto msg = std::string(
"failed to get publisher options: ") + rcl_get_error_string().str;
229 throw std::runtime_error(msg);
231 return publisher_options->
qos.depth;
240 std::shared_ptr<rcl_publisher_t>
243 return publisher_handle_;
246 std::shared_ptr<const rcl_publisher_t>
249 return publisher_handle_;
253 std::unordered_map<rcl_publisher_event_type_t, std::shared_ptr<rclcpp::EventHandlerBase>> &
256 return event_handlers_;
262 size_t inter_process_subscription_count = 0;
265 publisher_handle_.get(),
266 &inter_process_subscription_count);
279 rclcpp::exceptions::throw_from_rcl_error(status,
"failed to get get subscription count");
281 return inter_process_subscription_count;
287 auto ipm = weak_ipm_.lock();
288 if (!intra_process_is_enabled_) {
294 throw std::runtime_error(
295 "intra process subscriber count called after "
296 "destruction of intra process manager");
298 return ipm->get_subscription_count(intra_process_publisher_id_);
305 RMW_QOS_POLICY_DURABILITY_TRANSIENT_LOCAL;
313 auto msg = std::string(
"failed to get qos settings: ") + rcl_get_error_string().str;
315 throw std::runtime_error(msg);
336 return *
this == &gid;
343 auto ret = rmw_compare_gids_equal(gid, &this->
get_gid(), &result);
344 if (ret != RMW_RET_OK) {
345 auto msg = std::string(
"failed to compare gids: ") + rmw_get_error_string().str;
347 throw std::runtime_error(msg);
354 uint64_t intra_process_publisher_id,
355 const IntraProcessManagerSharedPtr & ipm)
357 intra_process_publisher_id_ = intra_process_publisher_id;
359 intra_process_is_enabled_ =
true;
363 PublisherBase::default_incompatible_qos_callback(
364 rclcpp::QOSOfferedIncompatibleQoSInfo & event)
const
366 std::string policy_name = qos_policy_name_from_kind(event.last_policy_kind);
369 "New subscription discovered on topic '%s', requesting incompatible QoS. "
370 "No messages will be sent to it. "
371 "Last incompatible policy: %s",
373 policy_name.c_str());
377 PublisherBase::default_incompatible_type_callback(
378 [[maybe_unused]] rclcpp::IncompatibleTypeInfo & event)
const
382 "Incompatible type on topic '%s', no messages will be sent to it.",
get_topic_name());
387 rcutils_allocator_t allocator = rcutils_get_default_allocator();
388 rcl_network_flow_endpoint_array_t network_flow_endpoint_array =
389 rcl_get_zero_initialized_network_flow_endpoint_array();
390 rcl_ret_t ret = rcl_publisher_get_network_flow_endpoints(
391 publisher_handle_.get(), &allocator, &network_flow_endpoint_array);
393 auto error_msg = std::string(
"error obtaining network flows of publisher: ") +
394 rcl_get_error_string().str;
397 rcl_network_flow_endpoint_array_fini(&network_flow_endpoint_array))
399 error_msg += std::string(
", also error cleaning up network flow array: ") +
400 rcl_get_error_string().str;
403 rclcpp::exceptions::throw_from_rcl_error(ret, error_msg);
406 std::vector<rclcpp::NetworkFlowEndpoint> network_flow_endpoint_vector;
407 network_flow_endpoint_vector.reserve(network_flow_endpoint_array.size);
408 for (
size_t i = 0; i < network_flow_endpoint_array.size; ++i) {
409 network_flow_endpoint_vector.emplace_back(
410 network_flow_endpoint_array.network_flow_endpoint[i]);
413 ret = rcl_network_flow_endpoint_array_fini(&network_flow_endpoint_array);
415 rclcpp::exceptions::throw_from_rcl_error(ret,
"error cleaning up network flow array");
418 return network_flow_endpoint_vector;
423 if (!intra_process_is_enabled_) {
427 auto ipm = weak_ipm_.lock();
433 "Intra process manager died for a publisher.");
437 return ipm->lowest_available_capacity(intra_process_publisher_id_);
442 const std::function<
void(
size_t)> & callback,
445 if (event_handlers_.count(event_type) == 0) {
448 "Calling set_on_new_qos_event_callback for non registered publisher event_type");
453 throw std::invalid_argument(
454 "The callback passed to set_on_new_qos_event_callback "
461 std::function<void(
size_t,
int)> new_callback = [callback] (
size_t nr, int) {callback(nr);};
462 event_handlers_[event_type]->set_on_ready_callback(new_callback);
468 if (event_handlers_.count(event_type) == 0) {
471 "Calling clear_on_new_qos_event_callback for non registered event_type");
475 event_handlers_[event_type]->clear_on_ready_callback();
RCLCPP_PUBLIC Logger get_child(const std::string &suffix)
Return a logger that is a descendant of this logger.
RCLCPP_PUBLIC const rmw_gid_t & get_gid() const
Get the global identifier for this publisher (used in rmw and by DDS).
RCLCPP_PUBLIC void set_on_new_qos_event_callback(const std::function< void(size_t)> &callback, rcl_publisher_event_type_t event_type)
Set a callback to be called when each new qos event instance occurs.
RCLCPP_PUBLIC std::shared_ptr< rcl_publisher_t > get_publisher_handle()
Get the rcl publisher handle.
RCLCPP_PUBLIC size_t get_intra_process_subscription_count() const
Get intraprocess subscription count.
RCLCPP_PUBLIC void clear_on_new_qos_event_callback(rcl_publisher_event_type_t event_type)
Unset the callback registered for new qos events, if any.
RCLCPP_PUBLIC void bind_event_callbacks(const PublisherEventCallbacks &event_callbacks, bool use_default_callbacks)
Add event handlers for passed in event_callbacks.
RCLCPP_PUBLIC rclcpp::QoS get_actual_qos() const
Get the actual QoS settings, after the defaults have been determined.
RCLCPP_PUBLIC const char * get_topic_name() const
Get the topic that this publisher publishes on.
RCLCPP_PUBLIC void setup_intra_process(uint64_t intra_process_publisher_id, const IntraProcessManagerSharedPtr &ipm)
Implementation utility function used to setup intra process publishing after creation.
RCLCPP_PUBLIC bool operator==(const rmw_gid_t &gid) const
Compare this publisher to a gid.
RCLCPP_PUBLIC size_t get_queue_size() const
Get the queue size for this publisher.
RCLCPP_PUBLIC bool can_loan_messages() const
Check if publisher instance can loan messages.
RCLCPP_PUBLIC bool is_durability_transient_local() const
Get if durability is transient local.
RCLCPP_PUBLIC size_t get_subscription_count() const
Get subscription count.
RCLCPP_PUBLIC std::vector< rclcpp::NetworkFlowEndpoint > get_network_flow_endpoints() const
Get network flow endpoints.
RCLCPP_PUBLIC RCUTILS_WARN_UNUSED bool assert_liveliness() const
Manually assert that this Publisher is alive (for RMW_QOS_POLICY_LIVELINESS_MANUAL_BY_TOPIC).
RCLCPP_PUBLIC const std::unordered_map< rcl_publisher_event_type_t, std::shared_ptr< rclcpp::EventHandlerBase > > & get_event_handlers() const
Get all the QoS event handlers associated with this publisher.
RCLCPP_PUBLIC size_t lowest_available_ipm_capacity() const
Return the lowest available capacity for all subscription buffers.
static RCLCPP_PUBLIC bool event_type_is_supported(const rcl_publisher_event_type_t event_type)
Check if a publisher event type is supported by the active RMW implementation.
Encapsulation of Quality of Service settings.
Pure virtual interface class for the NodeBase part of the Node API.
RCL_PUBLIC RCL_WARN_UNUSED bool rcl_context_is_valid(const rcl_context_t *context)
Return true if the given context is currently valid, otherwise false.
enum rcl_publisher_event_type_e rcl_publisher_event_type_t
Enumeration of all of the publisher events that may fire.
RCL_PUBLIC RCL_WARN_UNUSED bool rcl_publisher_event_type_is_supported(const rcl_publisher_event_type_t event_type)
Check if a publisher event type is supported by the active RMW implementation.
RCLCPP_PUBLIC std::string expand_topic_or_service_name(const std::string &name, const std::string &node_name, const std::string &namespace_, bool is_service=false)
Expand a topic or service name and throw if it is not valid.
RCLCPP_PUBLIC Logger get_node_logger(const rcl_node_t *node)
Return a named logger using an rcl_node_t.
RCLCPP_PUBLIC Logger get_logger(const std::string &name)
Return a named logger.
RCL_PUBLIC RCL_WARN_UNUSED const char * rcl_node_get_name(const rcl_node_t *node)
Return the name of the node.
RCL_PUBLIC RCL_WARN_UNUSED const char * rcl_node_get_namespace(const rcl_node_t *node)
Return the namespace of the node.
RCL_PUBLIC RCL_WARN_UNUSED const char * rcl_node_get_logger_name(const rcl_node_t *node)
Return the logger name of the node.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_publisher_init(rcl_publisher_t *publisher, const rcl_node_t *node, const rosidl_message_type_support_t *type_support, const char *topic_name, const rcl_publisher_options_t *options)
Initialize a rcl publisher.
RCL_PUBLIC RCL_WARN_UNUSED rcl_context_t * rcl_publisher_get_context(const rcl_publisher_t *publisher)
Return the context associated with this publisher.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_publisher_get_subscription_count(const rcl_publisher_t *publisher, size_t *subscription_count)
Get the number of subscriptions matched to a publisher.
RCL_PUBLIC RCL_WARN_UNUSED rmw_publisher_t * rcl_publisher_get_rmw_handle(const rcl_publisher_t *publisher)
Return the rmw publisher handle.
RCL_PUBLIC RCL_WARN_UNUSED const char * rcl_publisher_get_topic_name(const rcl_publisher_t *publisher)
Get the topic name for the publisher.
RCL_PUBLIC RCL_WARN_UNUSED const rmw_qos_profile_t * rcl_publisher_get_actual_qos(const rcl_publisher_t *publisher)
Get the actual qos settings of the publisher.
RCL_PUBLIC bool rcl_publisher_is_valid_except_context(const rcl_publisher_t *publisher)
Return true if the publisher is valid except the context, otherwise false.
RCL_PUBLIC RCL_WARN_UNUSED const rcl_publisher_options_t * rcl_publisher_get_options(const rcl_publisher_t *publisher)
Return the rcl publisher options.
RCL_PUBLIC RCL_WARN_UNUSED rcl_publisher_t rcl_get_zero_initialized_publisher(void)
Return a rcl_publisher_t struct with members set to NULL.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_publisher_fini(rcl_publisher_t *publisher, rcl_node_t *node)
Finalize a rcl_publisher_t.
RCL_PUBLIC bool rcl_publisher_can_loan_messages(const rcl_publisher_t *publisher)
Check if publisher instance can loan messages.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_publisher_assert_liveliness(const rcl_publisher_t *publisher)
Manually assert that this Publisher is alive (for RMW_QOS_POLICY_LIVELINESS_MANUAL_BY_TOPIC)
Encapsulates the non-global state of an init/shutdown cycle.
Options available for a rcl publisher.
rmw_qos_profile_t qos
Middleware quality of service settings for the publisher.
Structure which encapsulates a ROS Publisher.
Contains callbacks for various types of events a Publisher can receive from the middleware.
static QoSInitialization from_rmw(const rmw_qos_profile_t &rmw_qos)
Create a QoSInitialization from an existing rmw_qos_profile_t, using its history and depth.
#define RCL_RET_OK
Success return code.
#define RCL_RET_TOPIC_NAME_INVALID
Topic name does not pass validation.
rmw_ret_t rcl_ret_t
The type that holds an rcl return code.
#define RCL_RET_PUBLISHER_INVALID
Invalid rcl_publisher_t given return code.