Skip to content

simplifies EventsCBGExecutor::CBGScheduler by replacing its flags - #3290

Open
armaho wants to merge 1 commit into
ros2:rollingfrom
armaho:simplifying_events_cbg_executor
Open

armaho wants to merge 1 commit into
ros2:rollingfrom
armaho:simplifying_events_cbg_executor

Conversation

@armaho

@armaho armaho commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Description

Resolves #3289 by replacing the three flags inside CBGScheduler::CallbackGroupHandle (not_ready, idle and in_queue) with a simple in_scheduler_count field. This field keeps track of the number of entities that are currently in scheduler, either executing or just ready. The changes also make managing reentrant callback groups easier.

This PR also makes #3288 unnecessary since we don't handle reentrant callback groups in FIFOScheduler::get_next_ready_entity_intern methods anymore.

Is this user-facing behavior change?

No.

Did you use Generative AI?

No.

Additional Information

Aside from other test cases, I manually tested the changes to see its behavior regarding different types of CBGs. My test case is here (I'm using yaets for tracing):

#include "rclcpp/rclcpp.hpp"
#include <fstream>

#include "yaets/tracing.hpp"

#include "rclcpp/rclcpp.hpp"
#include "std_msgs/msg/int32.hpp"

using namespace std::chrono_literals;
using std::placeholders::_1;

yaets::TraceSession session("session1.log");

class ProducerNode : public rclcpp::Node {
public:
  ProducerNode() : Node("producer_node") {
    pub_1_ = create_publisher<std_msgs::msg::Int32>("topic_1", 100);
    pub_2_ = create_publisher<std_msgs::msg::Int32>("topic_2", 100);
    timer_ =
        create_wall_timer(1ms, std::bind(&ProducerNode::timer_callback, this));
  }

  void timer_callback() {
    RCLCPP_INFO(get_logger(), "Publishing");

    TRACE_EVENT(session);
    message_.data += 1;
    pub_1_->publish(message_);
    message_.data += 1;
    pub_2_->publish(message_);
  }

private:
  rclcpp::Publisher<std_msgs::msg::Int32>::SharedPtr pub_1_, pub_2_;
  rclcpp::TimerBase::SharedPtr timer_;
  std_msgs::msg::Int32 message_;
};

class ConsumerNode : public rclcpp::Node {
public:
  ConsumerNode() : Node("consumer_node") {
    auto cg = this->create_callback_group(
        rclcpp::CallbackGroupType::MutuallyExclusive);
    // auto cg = this->create_callback_group(
    //     rclcpp::CallbackGroupType::Reentrant);

    rclcpp::SubscriptionOptions opt;
    opt.callback_group = cg;

    sub_2_ = create_subscription<std_msgs::msg::Int32>(
        "topic_2", 100, std::bind(&ConsumerNode::cb_2, this, _1), opt);
    sub_1_ = create_subscription<std_msgs::msg::Int32>(
        "topic_1", 100, std::bind(&ConsumerNode::cb_1, this, _1), opt);

    timer_ =
        create_wall_timer(10ms, std::bind(&ConsumerNode::timer_callback, this));
  }

  void cb_1(const std_msgs::msg::Int32::SharedPtr msg) {
    TRACE_EVENT(session);

    RCLCPP_INFO(get_logger(), "Subscription executing");

    waste_time(500us);
  }

  void cb_2(const std_msgs::msg::Int32::SharedPtr msg) {
    TRACE_EVENT(session);

    RCLCPP_INFO(get_logger(), "Subscription executing");

    waste_time(500us);
  }

  void timer_callback() {
    TRACE_EVENT(session);

    waste_time(5ms);
  }

  void waste_time(const rclcpp::Duration &duration) {
    auto start = now();
    while (now() - start < duration)
      ;
  }

private:
  rclcpp::Subscription<std_msgs::msg::Int32>::SharedPtr sub_1_;
  rclcpp::Subscription<std_msgs::msg::Int32>::SharedPtr sub_2_;
  rclcpp::TimerBase::SharedPtr timer_;
};

int main(int argc, char *argv[]) {
  rclcpp::init(argc, argv);

  auto node_pub = std::make_shared<ProducerNode>();
  auto node_sub = std::make_shared<ConsumerNode>();

  // rclcpp::executors::SingleThreadedExecutor executor;
  // rclcpp::executors::MultiThreadedExecutor
  // executor(rclcpp::ExecutorOptions(), 8);
  rclcpp::executors::EventsCBGExecutor executor(rclcpp::ExecutorOptions(), 8);

  executor.add_node(node_pub);
  executor.add_node(node_sub);

  executor.spin();

  rclcpp::shutdown();
  return 0;
}

Here's the grant chart for the case of mutually exclusive subscribers:

gantt_chart_me

And for reentrant ones:

gantt_chart

@armaho

armaho commented Sep 29, 2026

Copy link
Copy Markdown
Contributor Author

@skyegalaxy @jmachowinski wdyt?

@github-actions

github-actions Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

ABI Compliance Check

❌ Verdict: incompatible

Library Verdict Summary
libcomponent_manager.so ✅ compatible No ABI changes detected.
librclcpp.so ❌ incompatible ABI-incompatible changes detected.
librclcpp_action.so ✅ compatible No ABI changes detected.
librclcpp_lifecycle.so ✅ compatible No ABI changes detected.
✅ libcomponent_manager.so — full abidiff report

Compared:

  • Base: lib-base/libcomponent_manager.so
  • Head: lib-pr/libcomponent_manager.so @ cd26180
(empty report — no differences printed by abidiff)
❌ librclcpp.so — full abidiff report

Compared:

  • Base: lib-base/librclcpp.so
  • Head: lib-pr/librclcpp.so @ cd26180
Functions changes summary: 1 Removed (4 filtered out), 2 Changed (133 filtered out), 1 Added functions
Variables changes summary: 0 Removed, 0 Changed, 0 Added variable

1 Removed function:

  [D] 'method void rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle::mark_as_skipped()'    {_ZN6rclcpp9executors12cbg_executor12CBGScheduler19CallbackGroupHandle15mark_as_skippedEv}

1 Added function:

  [A] 'method void rclcpp::executors::cbg_executor::CBGScheduler::mark_entity_as_executed(const rclcpp::executors::cbg_executor::CBGScheduler::ExecutableEntity&)'    {_ZN6rclcpp9executors12cbg_executor12CBGScheduler23mark_entity_as_executedERKNS2_16ExecutableEntityE}

2 functions with some indirect sub-type change:

  [C] 'method void rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle::mark_as_executed()' at scheduler.hpp:223:1 has some indirect sub-type changes:
    implicit parameter 0 of type 'rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle*' has sub-type changes:
      in pointed to type 'struct rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle' at scheduler.hpp:192:1:
        type size hasn't changed
        1 member function insertion:
          'method virtual rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle::~CallbackGroupHandle()' at scheduler.hpp:202:1
        no member function changes (10 filtered);
        2 data member deletions:
          'bool not_ready', at offset 512 (in bits) at scheduler.hpp:304:1
          'bool idle', at offset 520 (in bits) at scheduler.hpp:307:1
        2 data member changes (1 filtered):
          type of 'bool in_queue' changed:
            type name changed from 'bool' to 'int'
            type size changed from 8 to 32 (in bits)
          and name of 'rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle::in_queue' changed to 'rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle::in_scheduler_count' at scheduler.hpp:258:1
          'rclcpp::CallbackGroupType type' offset changed from 544 to 512 (in bits) (by -32 bits)

  [C] 'method rclcpp::executors::cbg_executor::FirstInFirstOutCallbackGroupHandle::FirstInFirstOutCallbackGroupHandle(rclcpp::executors::cbg_executor::CBGScheduler&, rclcpp::CallbackGroupType)' at first_in_first_out_scheduler.hpp:35:1 has some indirect sub-type changes:
    implicit parameter 0 of type 'rclcpp::executors::cbg_executor::FirstInFirstOutCallbackGroupHandle*' has sub-type changes:
      in pointed to type 'struct rclcpp::executors::cbg_executor::FirstInFirstOutCallbackGroupHandle' at first_in_first_out_scheduler.hpp:32:1:
        type size hasn't changed
        1 base class change:
          'struct rclcpp::executors::cbg_executor::CBGScheduler::CallbackGroupHandle' at scheduler.hpp:192:1 changed:
            details were reported earlier
        no member function changes (7 filtered);


✅ librclcpp_action.so — full abidiff report

Compared:

  • Base: lib-base/librclcpp_action.so
  • Head: lib-pr/librclcpp_action.so @ cd26180
(empty report — no differences printed by abidiff)
✅ librclcpp_lifecycle.so — full abidiff report

Compared:

  • Base: lib-base/librclcpp_lifecycle.so
  • Head: lib-pr/librclcpp_lifecycle.so @ cd26180
(empty report — no differences printed by abidiff)

Updated for commit cd26180 · suppressions: /home/runner/work/_temp/ros2-abi-suppressions.txt

@jmachowinski

Copy link
Copy Markdown
Collaborator

I got #3277 in flight and would like to merge it first. I highly expect merge conflicts from it.

@armaho

armaho commented Sep 29, 2026

Copy link
Copy Markdown
Contributor Author

@jmachowinski

Alright I would revisit this after that’s merged.

Signed-off-by: Arman Hosseini <armanhosseini878787@gmail.com>
@armaho
armaho force-pushed the simplifying_events_cbg_executor branch from 8d908eb to cd26180 Compare October 2, 2026 09:17
@armaho

armaho commented Oct 2, 2026

Copy link
Copy Markdown
Contributor Author

@jmachowinski

Now that #3277 is merged, we can say that there's no merge conflict. I also double checked the code and tests. Everything seems to hold up. This is ready to review.

@jmachowinski

Copy link
Copy Markdown
Collaborator

Even though this reduces the number of code lines, I think it makes the code harder to understand.
We swap explicitly named variables for a counter with an implicit meaning which I am not a big fan of.

I also fear that we are introducing lock order inversions. Did you check the code with something like helgrind (part of valgrind) or -fsanitize=thread ?

Note, that at some point there was a version that had a simpler logic, but that we needed shift code around to solve these locking issues.

@armaho

armaho commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor Author

@jmachowinski

I came up with this idea when I was trying to write a custom scheduler. It was hard for me to account for all 16 different combinations of not_ready, idle, in_queue and callback group's type, and I think it's easy to make mistakes while doing so. Providing a simple in_scheduler counter made that easier. It simply counts the number of entities that are already in the scheduler. Maybe a better name can help clear things up.

The only reason that it's a counter rather than a flag is that a reentrant callback group can have multiple entities in the scheduler at the same time and it's hard to keep a flag up to date when we remove a ready entity (we have to search the whole queue to find out whether we have another ready entity of the same group).

And about the locking order, by reading the code I think that we always take ready_callback_groups_mutex before CallbackGroupHandle::ready_mutex to avoid deadlock. I preserved that behavior here. I will also double check with the tools you mentioned and update the code if necessary.

@armaho

armaho commented Oct 2, 2026

Copy link
Copy Markdown
Contributor Author

@jmachowinski

Maybe we can replace the whole thing with a can_add_to_scheduler flag that's always true for reentrant callback groups and is false for mutually exclusive ones if they have an entity in the scheduler. Do you think that's a better approach?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Can we simplify the flags in CBGScheduler::CallbackGroupHandle?

2 participants