diff --git a/rclcpp/src/rclcpp/executors/events_cbg_executor/first_in_first_out_scheduler.cpp b/rclcpp/src/rclcpp/executors/events_cbg_executor/first_in_first_out_scheduler.cpp index c758982c4c..112625cd75 100644 --- a/rclcpp/src/rclcpp/executors/events_cbg_executor/first_in_first_out_scheduler.cpp +++ b/rclcpp/src/rclcpp/executors/events_cbg_executor/first_in_first_out_scheduler.cpp @@ -107,13 +107,9 @@ get_next_ready_entity() continue; } - mark_as_executing(); - return CBGScheduler::ExecutableEntity{exec_fun, this}; } - mark_as_skipped(); - return std::nullopt; } @@ -135,13 +131,9 @@ get_next_ready_entity(GlobalEventIdProvider::MonotonicId max_id) continue; } - mark_as_executing(); - return CBGScheduler::ExecutableEntity{exec_fun, this}; } - mark_as_skipped(); - return std::nullopt; } @@ -158,19 +150,16 @@ CBGScheduler::ExecutableEntityWithInfo FirstInFirstOutScheduler::get_next_ready_ FirstInFirstOutCallbackGroupHandle *ready_cbg = static_cast(ready_callback_groups.front()); ready_callback_groups.pop_front(); - ready_cbg->in_queue = false; std::optional ret = ready_cbg->get_next_ready_entity(); - if (ready_cbg->get_type() == CallbackGroupType::Reentrant && ready_cbg->has_ready_entities()) { - ready_callback_groups.push_back(ready_cbg); - ready_cbg->in_queue = true; - } - if(ret) { return CBGScheduler::ExecutableEntityWithInfo{.entity = std::move(ret), .moreEntitiesReady = !ready_callback_groups.empty()}; + } else { + std::lock_guard lock(ready_cbg->get_ready_mutex()); + ready_cbg->in_scheduler_count--; } } @@ -191,14 +180,7 @@ CBGScheduler::ExecutableEntityWithInfo FirstInFirstOutScheduler::get_next_ready_ ready_cbg->get_next_ready_entity(max_id); if(ret) { ready_callback_groups.erase(it); - ready_cbg->in_queue = false; - - if ( - ready_cbg->get_type() == CallbackGroupType::Reentrant && ready_cbg->has_ready_entities()) - { - ready_callback_groups.push_back(ready_cbg); - ready_cbg->in_queue = true; - } + return CBGScheduler::ExecutableEntityWithInfo{ .entity = std::move(ret), .moreEntitiesReady = !ready_callback_groups.empty()}; } diff --git a/rclcpp/src/rclcpp/executors/events_cbg_executor/scheduler.hpp b/rclcpp/src/rclcpp/executors/events_cbg_executor/scheduler.hpp index 48a4cec116..deb7f9ed96 100644 --- a/rclcpp/src/rclcpp/executors/events_cbg_executor/scheduler.hpp +++ b/rclcpp/src/rclcpp/executors/events_cbg_executor/scheduler.hpp @@ -222,15 +222,20 @@ class CBGScheduler */ void mark_as_executed() { + // for reentrant callback groups, we always inform the scheduler + // about a ready entity in add_ready_entity. + if (type == CallbackGroupType::Reentrant) { + return; + } + { std::lock_guard l(ready_mutex); - not_ready = false; - if(!has_ready_entities()) { - idle = true; + if (!has_ready_entities()) { return; } } + // inform scheduler that we have more work scheduler.callback_group_ready(this, false); } @@ -239,8 +244,18 @@ class CBGScheduler bool is_ready(); - // true if this cbg is inside the scheduler's queue - bool in_queue = false; + std::mutex & get_ready_mutex() + { + return ready_mutex; + } + + // number of entities in scheduler, either executing or in queue. + // + // modify while holding both CBGScheduler::ready_callback_groups_mutex and + // CallbackGroupHandle::ready_mutex, since you probably need to also change + // the scheduler's ready queue. reads can happen while holding only + // CallbackGroupHandle::ready_mutex. + int in_scheduler_count = 0; protected: CBGScheduler & scheduler; @@ -265,48 +280,22 @@ class CBGScheduler fun(); - if(not_ready || !idle) { + if (type == CallbackGroupType::MutuallyExclusive && + in_scheduler_count != 0) + { + // There's another ready entity of this group in the scheduler. + // We will add this to the scheduler's queue when it calls + // mark_as_executed for that entity. return; } - - idle = false; } - // If we reached this point, we were idle and now have work, - // therefore we need to move this callback group into the list - // of ready callback groups. scheduler.callback_group_ready(this, true); } - void mark_as_skipped() - { - if(!has_ready_entities()) { - idle = true; - } - } - - /** - * Must be called by derived classes if a ready entity is - * returned. This call must happen under a lock holding the - * ready_mutex. - */ - void mark_as_executing() - { - if (type != CallbackGroupType::Reentrant) { - not_ready = true; - } - } - std::mutex ready_mutex; private: - // will be set if cbg is mutual exclusive and something is executing - bool not_ready = false; - - // true, if nothing is beeing executed, and there are no pending events - bool idle = true; - - // type of the underlying callback group CallbackGroupType type; }; @@ -358,24 +347,23 @@ class CBGScheduler }); } - /** Will be called, by CallbackGroupHandle if any entity in the cb group is ready for execution - * and the cb group was idle before - * @param callback_group_was_idle Is false, if no entity of the callback group was executed, - * before this call was made. This means we need to wakeup a - * a new thread. + /** Will be called, by CallbackGroupHandle if any entity in the cb group is ready for execution. + * @param unblock_thread If true, that means that an idle thread can + * start executing the new entity. This can happen if + * any entity in a reentrant cbg becomes ready, or if + * a previously idle mutually exclusive cbg becomes ready. */ - void callback_group_ready(CallbackGroupHandle *handle, bool callback_group_was_idle) + void callback_group_ready(CallbackGroupHandle *handle, bool unblock_thread) { { - std::lock_guard l(ready_callback_groups_mutex); + std::lock_guard lock1(ready_callback_groups_mutex); + std::lock_guard lock2(handle->get_ready_mutex()); - if (!handle->in_queue) { - ready_callback_groups.push_back(handle); - handle->in_queue = true; - } + ready_callback_groups.push_back(handle); + handle->in_scheduler_count++; } - if(callback_group_was_idle) { + if(unblock_thread) { unblock_one_worker_thread(); } } @@ -461,6 +449,13 @@ class CBGScheduler void mark_entity_as_executed(const ExecutableEntity & e) { if(e.callback_handle != nullptr) { + { + // We do not hold CBGScheduler::ready_callback_groups_mutex here + // since we're not modifying the queue. This entity is already removed. + std::lock_guard lock(e.callback_handle->get_ready_mutex()); + e.callback_handle->in_scheduler_count--; + } + e.callback_handle->mark_as_executed(); } }