15 #ifndef RCLCPP__CLIENT_HPP_
16 #define RCLCPP__CLIENT_HPP_
27 #include <unordered_map>
33 #include "rcl/error_handling.h"
34 #include "rcl/event_callback.h"
35 #include "rcl/service_introspection.h"
38 #include "rclcpp/clock.hpp"
39 #include "rclcpp/detail/cpp_callback_trampoline.hpp"
40 #include "rclcpp/exceptions.hpp"
41 #include "rclcpp/expand_topic_or_service_name.hpp"
42 #include "rclcpp/function_traits.hpp"
43 #include "rclcpp/logging.hpp"
44 #include "rclcpp/macros.hpp"
45 #include "rclcpp/node_interfaces/node_graph_interface.hpp"
46 #include "rclcpp/qos.hpp"
47 #include "rclcpp/type_support_decl.hpp"
48 #include "rclcpp/utilities.hpp"
49 #include "rclcpp/visibility_control.hpp"
51 #include "rmw/error_handling.h"
52 #include "rmw/impl/cpp/demangle.hpp"
60 template<
typename FutureT>
67 : future(std::move(impl)), request_id(req_id)
71 operator FutureT &() {
return this->future;}
76 auto get() {
return this->future.get();}
78 bool valid() const noexcept {
return this->future.valid();}
80 void wait()
const {
return this->future.wait();}
82 template<
class Rep,
class Period>
84 const std::chrono::duration<Rep, Period> & timeout_duration)
const
86 return this->future.wait_for(timeout_duration);
89 template<
class Clock,
class Duration>
91 const std::chrono::time_point<Clock, Duration> & timeout_time)
const
93 return this->future.wait_until(timeout_time);
111 template<
typename PendingRequestsT,
typename AllocatorT = std::allocator<
int64_t>>
113 prune_requests_older_than_impl(
114 PendingRequestsT & pending_requests,
115 std::mutex & pending_requests_mutex,
116 std::chrono::time_point<std::chrono::system_clock> time_point,
117 std::vector<int64_t, AllocatorT> * pruned_requests =
nullptr)
119 std::lock_guard guard(pending_requests_mutex);
120 auto old_size = pending_requests.size();
121 for (
auto it = pending_requests.begin(), last = pending_requests.end(); it != last; ) {
122 if (it->second.first < time_point) {
123 if (pruned_requests) {
124 pruned_requests->push_back(it->first);
126 it = pending_requests.erase(it);
131 return old_size - pending_requests.size();
135 namespace node_interfaces
137 class NodeBaseInterface;
143 RCLCPP_SMART_PTR_DEFINITIONS_NOT_COPYABLE(
ClientBase)
148 const rclcpp::node_interfaces::NodeGraphInterface::SharedPtr & node_graph);
186 std::shared_ptr<rcl_client_t>
195 std::shared_ptr<const rcl_client_t>
211 template<
typename RepT =
int64_t,
typename RatioT = std::milli>
214 std::chrono::duration<RepT, RatioT> timeout = std::chrono::duration<RepT, RatioT>(-1))
216 return wait_for_service_nanoseconds(
217 std::chrono::duration_cast<std::chrono::nanoseconds>(timeout)
221 virtual std::shared_ptr<void> create_response() = 0;
222 virtual std::shared_ptr<rmw_request_id_t> create_request_header() = 0;
223 virtual void handle_response(
224 const std::shared_ptr<rmw_request_id_t> & request_header,
225 const std::shared_ptr<void> & response) = 0;
303 throw std::invalid_argument(
304 "The callback passed to set_on_new_response_callback "
309 [callback,
this](
size_t number_of_responses) {
311 callback(number_of_responses);
312 }
catch (
const std::exception & exception) {
315 "rclcpp::ClientBase@" <<
this <<
316 " caught " << rmw::impl::cpp::demangle(exception) <<
317 " exception in user-provided callback for the 'on new response' callback: " <<
322 "rclcpp::ClientBase@" <<
this <<
323 " caught unhandled exception in user-provided callback " <<
324 "for the 'on new response' callback");
328 std::lock_guard<std::recursive_mutex> lock(callback_mutex_);
334 rclcpp::detail::cpp_callback_trampoline<decltype(new_callback),
const void *,
size_t>,
335 static_cast<const void *
>(&new_callback));
338 on_new_response_callback_ = new_callback;
342 rclcpp::detail::cpp_callback_trampoline<
343 decltype(on_new_response_callback_),
const void *,
size_t>,
344 static_cast<const void *
>(&on_new_response_callback_));
351 std::lock_guard<std::recursive_mutex> lock(callback_mutex_);
352 if (on_new_response_callback_) {
354 on_new_response_callback_ =
nullptr;
363 wait_for_service_nanoseconds(std::chrono::nanoseconds timeout);
367 get_rcl_node_handle();
371 get_rcl_node_handle()
const;
377 rclcpp::node_interfaces::NodeGraphInterface::WeakPtr node_graph_;
378 std::shared_ptr<rcl_node_t> node_handle_;
379 std::shared_ptr<rclcpp::Context> context_;
382 std::recursive_mutex callback_mutex_;
387 std::function<void(
size_t)> on_new_response_callback_{
nullptr};
389 std::shared_ptr<rcl_client_t> client_handle_;
391 std::atomic<bool> in_use_by_wait_set_{
false};
394 template<
typename ServiceT>
398 using Request =
typename ServiceT::Request;
399 using Response =
typename ServiceT::Response;
401 using SharedRequest =
typename ServiceT::Request::SharedPtr;
402 using SharedResponse =
typename ServiceT::Response::SharedPtr;
404 using Promise = std::promise<SharedResponse>;
405 using PromiseWithRequest = std::promise<std::pair<SharedRequest, SharedResponse>>;
407 using SharedPromise = std::shared_ptr<Promise>;
408 using SharedPromiseWithRequest = std::shared_ptr<PromiseWithRequest>;
410 using Future = std::future<SharedResponse>;
411 using SharedFuture = std::shared_future<SharedResponse>;
412 using SharedFutureWithRequest = std::shared_future<std::pair<SharedRequest, SharedResponse>>;
414 using CallbackType = std::function<void (SharedFuture)>;
415 using CallbackWithRequestType = std::function<void (SharedFutureWithRequest)>;
417 RCLCPP_SMART_PTR_DEFINITIONS(
Client)
435 SharedFuture
share() noexcept {
return this->future.share();}
464 std::shared_future<std::pair<SharedRequest, SharedResponse>>
481 const rclcpp::node_interfaces::NodeGraphInterface::SharedPtr & node_graph,
482 const std::string & service_name,
485 srv_type_support_handle_(rosidl_typesupport_cpp::get_service_type_support_handle<ServiceT>())
489 this->get_rcl_node_handle(),
490 srv_type_support_handle_,
491 service_name.c_str(),
495 auto rcl_node_handle = this->get_rcl_node_handle();
504 rclcpp::exceptions::throw_from_rcl_error(ret,
"could not create client");
526 take_response(
typename ServiceT::Response & response_out, rmw_request_id_t & request_header_out)
535 std::shared_ptr<void>
538 return std::shared_ptr<void>(
new typename ServiceT::Response());
545 std::shared_ptr<rmw_request_id_t>
550 return std::shared_ptr<rmw_request_id_t>(
new rmw_request_id_t);
560 const std::shared_ptr<rmw_request_id_t> & request_header,
561 const std::shared_ptr<void> & response)
override
563 std::optional<CallbackInfoVariant>
564 optional_pending_request = this->get_and_erase_pending_request(request_header->sequence_number);
565 if (!optional_pending_request) {
568 auto & value = *optional_pending_request;
569 auto typed_response = std::static_pointer_cast<typename ServiceT::Response>(
571 if (std::holds_alternative<Promise>(value)) {
572 auto & promise = std::get<Promise>(value);
573 promise.set_value(std::move(typed_response));
574 }
else if (std::holds_alternative<CallbackTypeValueVariant>(value)) {
575 auto & inner = std::get<CallbackTypeValueVariant>(value);
576 const auto & callback = std::get<CallbackType>(inner);
577 auto & promise = std::get<Promise>(inner);
578 auto & future = std::get<SharedFuture>(inner);
579 promise.set_value(std::move(typed_response));
580 callback(std::move(future));
581 }
else if (std::holds_alternative<CallbackWithRequestTypeValueVariant>(value)) {
582 auto & inner = std::get<CallbackWithRequestTypeValueVariant>(value);
583 const auto & callback = std::get<CallbackWithRequestType>(inner);
584 auto & promise = std::get<PromiseWithRequest>(inner);
585 auto & future = std::get<SharedFutureWithRequest>(inner);
586 auto & request = std::get<SharedRequest>(inner);
587 promise.set_value(std::make_pair(std::move(request), std::move(typed_response)));
588 callback(std::move(future));
624 auto future = promise.get_future();
625 auto req_id = async_send_request_impl(
648 typename std::enable_if<
655 SharedFutureAndRequestId
659 auto shared_future = promise.get_future().share();
660 auto req_id = async_send_request_impl(
663 CallbackType{std::forward<CallbackT>(cb)},
665 std::move(promise)));
679 typename std::enable_if<
682 CallbackWithRequestType
686 SharedFutureWithRequestAndRequestId
689 PromiseWithRequest promise;
690 auto shared_future = promise.get_future().share();
691 auto req_id = async_send_request_impl(
694 CallbackWithRequestType{std::forward<CallbackT>(cb)},
697 std::move(promise)));
715 std::lock_guard guard(pending_requests_mutex_);
716 return pending_requests_.erase(request_id) != 0u;
762 std::lock_guard guard(pending_requests_mutex_);
763 auto ret = pending_requests_.size();
764 pending_requests_.clear();
775 template<
typename AllocatorT = std::allocator<
int64_t>>
778 std::chrono::time_point<std::chrono::system_clock> time_point,
779 std::vector<int64_t, AllocatorT> * pruned_requests =
nullptr)
781 return detail::prune_requests_older_than_impl(
783 pending_requests_mutex_,
799 const Clock::SharedPtr & clock,
const QoS & qos_service_event_pub,
800 rcl_service_introspection_state_t introspection_state)
806 client_handle_.get(),
808 clock->get_clock_handle(),
809 srv_type_support_handle_,
811 introspection_state);
814 rclcpp::exceptions::throw_from_rcl_error(ret,
"failed to configure client introspection");
819 using CallbackTypeValueVariant = std::tuple<CallbackType, SharedFuture, Promise>;
820 using CallbackWithRequestTypeValueVariant = std::tuple<
821 CallbackWithRequestType, SharedRequest, SharedFutureWithRequest, PromiseWithRequest>;
823 using CallbackInfoVariant = std::variant<
824 std::promise<SharedResponse>,
825 CallbackTypeValueVariant,
826 CallbackWithRequestTypeValueVariant>;
829 async_send_request_impl(
const Request & request, CallbackInfoVariant value)
831 int64_t sequence_number;
832 std::lock_guard<std::mutex> lock(pending_requests_mutex_);
835 rclcpp::exceptions::throw_from_rcl_error(ret,
"failed to send request");
837 pending_requests_.try_emplace(
839 std::make_pair(std::chrono::system_clock::now(), std::move(value)));
840 return sequence_number;
843 std::optional<CallbackInfoVariant>
844 get_and_erase_pending_request(int64_t request_number)
846 std::unique_lock<std::mutex> lock(pending_requests_mutex_);
847 auto it = this->pending_requests_.find(request_number);
848 if (it == this->pending_requests_.end()) {
849 RCUTILS_LOG_DEBUG_NAMED(
851 "Received invalid sequence number. Ignoring...");
854 std::optional<CallbackInfoVariant> value = std::move(it->second.second);
855 this->pending_requests_.erase(request_number);
859 RCLCPP_DISABLE_COPY(
Client)
864 std::chrono::time_point<std::chrono::system_clock>,
865 CallbackInfoVariant>>
867 std::mutex pending_requests_mutex_;
870 const rosidl_service_type_support_t * srv_type_support_handle_;
RCLCPP_PUBLIC bool exchange_in_use_by_wait_set_state(bool in_use_state)
Exchange the "in use by wait set" state for this client.
RCLCPP_PUBLIC rclcpp::QoS get_request_publisher_actual_qos() const
Get the actual request publsher QoS settings, after the defaults have been determined.
RCLCPP_PUBLIC std::shared_ptr< rcl_client_t > get_client_handle()
Return the rcl_client_t client handle in a std::shared_ptr.
RCLCPP_PUBLIC bool take_type_erased_response(void *response_out, rmw_request_id_t &request_header_out)
Take the next response for this client as a type erased pointer.
void set_on_new_response_callback(const std::function< void(size_t)> &callback)
Set a callback to be called when each new response is received.
bool wait_for_service(std::chrono::duration< RepT, RatioT > timeout=std::chrono::duration< RepT, RatioT >(-1))
Wait for a service to be ready.
RCLCPP_PUBLIC const char * get_service_name() const
Return the name of the service.
void clear_on_new_response_callback()
Unset the callback registered for new responses, if any.
RCLCPP_PUBLIC rclcpp::QoS get_response_subscription_actual_qos() const
Get the actual response subscription QoS settings, after the defaults have been determined.
RCLCPP_PUBLIC bool service_is_ready() const
Return if the service is ready.
bool take_response(typename ServiceT::Response &response_out, rmw_request_id_t &request_header_out)
Take the next response for this client.
std::shared_ptr< rmw_request_id_t > create_request_header() override
Create a shared pointer with a rmw_request_id_t.
bool remove_pending_request(int64_t request_id)
Cleanup a pending request.
std::shared_ptr< void > create_response() override
Create a shared pointer with the response type.
size_t prune_requests_older_than(std::chrono::time_point< std::chrono::system_clock > time_point, std::vector< int64_t, AllocatorT > *pruned_requests=nullptr)
Clean all pending requests older than a time_point.
void configure_introspection(const Clock::SharedPtr &clock, const QoS &qos_service_event_pub, rcl_service_introspection_state_t introspection_state)
Configure client introspection.
FutureAndRequestId async_send_request(const SharedRequest &request)
Send a request to the service server.
Client(rclcpp::node_interfaces::NodeBaseInterface *node_base, const rclcpp::node_interfaces::NodeGraphInterface::SharedPtr &node_graph, const std::string &service_name, rcl_client_options_t &client_options)
Default constructor.
size_t prune_pending_requests()
Clean all pending requests.
bool remove_pending_request(const SharedFutureAndRequestId &future)
Cleanup a pending request.
SharedFutureWithRequestAndRequestId async_send_request(const SharedRequest &request, CallbackT &&cb)
Send a request to the service server and schedule a callback in the executor.
bool remove_pending_request(const SharedFutureWithRequestAndRequestId &future)
Cleanup a pending request.
void handle_response(const std::shared_ptr< rmw_request_id_t > &request_header, const std::shared_ptr< void > &response) override
Handle a server response.
SharedFutureAndRequestId async_send_request(const SharedRequest &request, CallbackT &&cb)
Send a request to the service server and schedule a callback in the executor.
bool remove_pending_request(const FutureAndRequestId &future)
Cleanup a pending request.
Encapsulation of Quality of Service settings.
rmw_qos_profile_t & get_rmw_qos_profile()
Return the rmw qos profile.
Pure virtual interface class for the NodeBase part of the Node API.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_client_configure_service_introspection(rcl_client_t *client, rcl_node_t *node, rcl_clock_t *clock, const rosidl_service_type_support_t *type_support, const rcl_publisher_options_t publisher_options, rcl_service_introspection_state_t introspection_state)
Configures service introspection features for the client.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_client_init(rcl_client_t *client, const rcl_node_t *node, const rosidl_service_type_support_t *type_support, const char *service_name, const rcl_client_options_t *options)
Initialize a rcl client.
RCL_PUBLIC RCL_WARN_UNUSED rcl_ret_t rcl_send_request(const rcl_client_t *client, const void *ros_request, int64_t *sequence_number)
Send a ROS request using a client.
Versions of rosidl_typesupport_cpp::get_message_type_support_handle that handle adapted types.
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.
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 rcl_publisher_options_t rcl_publisher_get_default_options(void)
Return the default publisher options in a rcl_publisher_options_t.
Options available for a rcl_client_t.
Structure which encapsulates a ROS Node.
Options available for a rcl publisher.
rmw_qos_profile_t qos
Middleware quality of service settings for the publisher.
A convenient Client::Future and request id pair.
SharedFuture share() noexcept
See std::future::share().
A convenient Client::SharedFuture and request id pair.
A convenient Client::SharedFutureWithRequest and request id pair.
std::future_status wait_for(const std::chrono::duration< Rep, Period > &timeout_duration) const
See std::future::wait_for().
std::future_status wait_until(const std::chrono::time_point< Clock, Duration > &timeout_time) const
See std::future::wait_until().
FutureAndRequestId & operator=(FutureAndRequestId &&other) noexcept=default
Move assignment.
FutureAndRequestId & operator=(const FutureAndRequestId &other)=delete
Deleted copy assignment, each instance is a unique owner of the future.
FutureAndRequestId(const FutureAndRequestId &other)=delete
Deleted copy constructor, each instance is a unique owner of the future.
void wait() const
See std::future::wait().
FutureAndRequestId(FutureAndRequestId &&other) noexcept=default
Move constructor.
~FutureAndRequestId()=default
Destructor.
auto get()
See std::future::get().
bool valid() const noexcept
See std::future::valid().
#define RCL_RET_SERVICE_NAME_INVALID
Service name (same as topic name) does not pass validation.
#define RCL_RET_OK
Success return code.
rmw_ret_t rcl_ret_t
The type that holds an rcl return code.