15 #include "rclcpp/experimental/timers_manager.hpp"
24 #include "rclcpp/utilities.hpp"
25 #include "rcpputils/scope_exit.hpp"
29 TimersManager::TimersManager(
30 std::shared_ptr<rclcpp::Context> context,
31 std::function<
void(
const rclcpp::TimerBase *,
const std::shared_ptr<void> &)> on_ready_callback)
32 : on_ready_callback_(on_ready_callback),
33 context_(std::move(context))
49 throw std::invalid_argument(
"TimersManager::add_timer() trying to add nullptr timer");
54 std::unique_lock<std::mutex> lock(timers_mutex_);
55 added = weak_timers_heap_.add_timer(timer);
56 timers_updated_ = timers_updated_ || added;
59 timer->set_on_reset_callback(
63 std::unique_lock<std::mutex> lock(timers_mutex_);
64 timers_updated_ =
true;
66 timers_cv_.notify_one();
71 timers_cv_.notify_one();
78 if (running_.exchange(
true)) {
79 throw std::runtime_error(
"TimersManager::start() can't start timers thread as already running");
82 timers_thread_ = std::thread(&TimersManager::run_timers,
this);
88 std::unique_lock<std::mutex> lock(stop_mutex_);
93 std::unique_lock<std::mutex> lock(timers_mutex_);
94 timers_updated_ =
true;
96 timers_cv_.notify_one();
99 if (timers_thread_.joinable()) {
100 timers_thread_.join();
108 throw std::runtime_error(
109 "get_head_timeout() can't be used while timers thread is running");
112 std::unique_lock<std::mutex> lock(timers_mutex_);
113 return this->get_head_timeout_unsafe();
120 throw std::runtime_error(
121 "get_number_ready_timers() can't be used while timers thread is running");
124 std::unique_lock<std::mutex> lock(timers_mutex_);
125 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
126 return locked_heap.get_number_ready_timers();
133 throw std::runtime_error(
134 "execute_head_timer() can't be used while timers thread is running");
137 std::unique_lock<std::mutex> lock(timers_mutex_);
139 TimersHeap timers_heap = weak_timers_heap_.validate_and_lock();
142 if (timers_heap.empty()) {
146 TimerPtr head_timer = timers_heap.front();
148 const bool timer_ready = head_timer->is_ready();
152 auto data = head_timer->call();
157 head_timer->execute_callback(data);
158 timers_heap.heapify_root();
159 weak_timers_heap_.store(timers_heap);
167 const std::shared_ptr<void> & data)
169 TimerPtr ready_timer;
171 std::unique_lock<std::mutex> lock(timers_mutex_);
172 ready_timer = weak_timers_heap_.get_timer(timer_id);
175 ready_timer->execute_callback(data);
179 std::optional<std::chrono::nanoseconds> TimersManager::get_head_timeout_unsafe()
182 if (weak_timers_heap_.empty()) {
183 return std::chrono::nanoseconds::max();
187 TimerPtr head_timer = weak_timers_heap_.front().lock();
192 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
196 if (locked_heap.empty()) {
197 return std::chrono::nanoseconds::max();
199 head_timer = locked_heap.front();
201 if (head_timer->is_canceled()) {
204 return head_timer->time_until_trigger();
207 void TimersManager::execute_ready_timers_unsafe()
210 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
213 if (locked_heap.empty()) {
221 TimerPtr head_timer = locked_heap.front();
222 const size_t number_ready_timers = locked_heap.get_number_ready_timers();
223 size_t executed_timers = 0;
224 while (executed_timers < number_ready_timers && head_timer->is_ready()) {
225 auto data = head_timer->call();
227 if (on_ready_callback_) {
228 on_ready_callback_(head_timer.get(), data);
230 head_timer->execute_callback(data);
239 locked_heap.heapify_root();
241 head_timer = locked_heap.front();
246 weak_timers_heap_.store(locked_heap);
249 void TimersManager::run_timers()
253 RCPPUTILS_SCOPE_EXIT(this->running_.store(
false); );
257 std::unique_lock<std::mutex> lock(timers_mutex_);
259 std::optional<std::chrono::nanoseconds> time_to_sleep = get_head_timeout_unsafe();
264 if (!time_to_sleep.has_value()) {
267 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
268 locked_heap.heapify();
269 weak_timers_heap_.store(locked_heap);
270 time_to_sleep = get_head_timeout_unsafe();
274 if (!time_to_sleep.has_value() || (time_to_sleep.value() == std::chrono::nanoseconds::max()) ) {
276 timers_cv_.wait(lock, [
this]() {
return timers_updated_;});
280 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
281 locked_heap.heapify();
282 weak_timers_heap_.store(locked_heap);
283 }
else if (time_to_sleep.value() != std::chrono::nanoseconds::zero()) {
286 timers_cv_.wait_for(lock, time_to_sleep.value(), [
this]() {return timers_updated_;});
290 timers_updated_ =
false;
293 this->execute_ready_timers_unsafe();
301 std::unique_lock<std::mutex> lock(timers_mutex_);
303 TimersHeap locked_heap = weak_timers_heap_.validate_and_lock();
304 locked_heap.clear_timers_on_reset_callbacks();
306 weak_timers_heap_.clear();
308 timers_updated_ =
true;
312 timers_cv_.notify_one();
317 bool removed =
false;
319 std::unique_lock<std::mutex> lock(timers_mutex_);
320 removed = weak_timers_heap_.remove_timer(timer);
322 timers_updated_ = timers_updated_ || removed;
327 timers_cv_.notify_one();
328 timer->clear_on_reset_callback();
This class provides a way for storing and executing timer objects. It provides APIs to suit the needs...
RCLCPP_PUBLIC void start()
Starts a thread that takes care of executing the timers stored in this object. Function will throw an...
RCLCPP_PUBLIC ~TimersManager()
Destruct the TimersManager object making sure to stop thread and release memory.
RCLCPP_PUBLIC std::optional< std::chrono::nanoseconds > get_head_timeout()
Get the amount of time before the next timer triggers. This function is thread safe.
RCLCPP_PUBLIC void add_timer(const rclcpp::TimerBase::SharedPtr &timer)
Adds a new timer to the storage, maintaining weak ownership of it. Function is thread safe and it can...
RCLCPP_PUBLIC bool execute_head_timer()
Executes head timer if ready. This function is thread safe. This function will try to execute the tim...
RCLCPP_PUBLIC void clear()
Remove all the timers stored in the object. Function is thread safe and it can be called regardless o...
RCLCPP_PUBLIC size_t get_number_ready_timers()
Get the number of timers that are currently ready. This function is thread safe.
RCLCPP_PUBLIC void remove_timer(const rclcpp::TimerBase::SharedPtr &timer)
Remove a single timer from the object storage. Will do nothing if the timer was not being stored here...
RCLCPP_PUBLIC void stop()
Stops the timers thread. Will do nothing if the timer thread was not running.
RCLCPP_PUBLIC void execute_ready_timer(const rclcpp::TimerBase *timer_id, const std::shared_ptr< void > &data)
Executes timer identified by its ID. This function is thread safe. This function will try to execute ...
RCLCPP_PUBLIC bool ok(const rclcpp::Context::SharedPtr &context=rclcpp::contexts::get_global_default_context())
Check rclcpp's status.