From 98dc42124ea75dd7ac6bd9e0fbb34288caeb3589 Mon Sep 17 00:00:00 2001 From: Miguel Company Date: Wed, 8 Jul 2026 12:07:29 +0200 Subject: [PATCH 1/6] Add test for possible race in guard condition notification. Signed-off-by: Miguel Company --- .../test/test_wait_set.cpp | 90 +++++++++++++++++++ 1 file changed, 90 insertions(+) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index 6c4bef4c8..a956193d7 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -413,3 +413,93 @@ TEST_F(TestWaitSet, rmw_destroy_wait_set) ret = rmw_destroy_wait_set(wait_set); EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; } + +/** + * This test creates a wait set with multiple guard conditions and triggers them in a separate thread. + * Expects that all guard conditions are notified exactly once and that rmw_wait returns RMW_RET_TIMEOUT after all + * guard conditions have been triggered. + */ +TEST_F(TestWaitSet, rmw_wait_guard_conditions) +{ + constexpr size_t number_of_guard_conditions = 10u; + typedef rmw_guard_condition_t* guard_condition_ptr_t; + guard_condition_ptr_t guard_condition_ptrs[number_of_guard_conditions]; + + // Create the guard conditions + size_t ready_conditions_value = 0u; + for (size_t i = 0; i < number_of_guard_conditions; ++i) { + guard_condition_ptrs[i] = rmw_create_guard_condition(&context); + ready_conditions_value |= (1u << i); + ASSERT_NE(nullptr, guard_condition_ptrs[i]) << rcutils_get_error_string().str; + } + + // Create a wait set + rmw_wait_set_t* wait_set = rmw_create_wait_set(&context, number_of_guard_conditions); + ASSERT_NE(nullptr, wait_set) << rcutils_get_error_string().str; + + // Thread that triggers all guard conditions in reverse order after 100ms + std::thread trigger_thread([&]() { + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + for (size_t i = number_of_guard_conditions; i > 0; --i) { + rmw_ret_t ret = rmw_trigger_guard_condition(guard_condition_ptrs[i - 1]); + EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + } + }); + + // Prepare input arguments for rmw_wait + rmw_time_t timeout_argument = {0, 100000000}; // 100ms + rmw_guard_conditions_t guard_conditions; + INITIALIZE_ARRAY(guard_conditions, guard_condition, number_of_guard_conditions); + + // Wait up to number_of_guard_conditions times or until all guard conditions are triggered + for (size_t i = 0; i < number_of_guard_conditions; ++i) { + // Prepare the input array for rmw_wait + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; + } + // Wait for guard conditions to be triggered + rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, &timeout_argument); + if (RMW_RET_OK == ret) { + // Check which guard conditions are ready + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + if (guard_conditions.guard_conditions[j] != nullptr) { + // Check that condition was not already triggered + EXPECT_NE(ready_conditions_value & (1u << j), 0u); + // Mark this condition as triggered + ready_conditions_value &= ~(1u << j); + } + } + + if (ready_conditions_value == 0u) { + // All guard conditions have been triggered + break; + } + } + } + + // Join the trigger thread + trigger_thread.join(); + + // Check that all guard conditions were triggered + EXPECT_EQ(ready_conditions_value, 0u) << "Not all guard conditions were triggered"; + + // Calling rmw_wait shall now return RMW_RET_TIMEOUT since all guard conditions have been triggered + { + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; + } + rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, &timeout_argument); + EXPECT_EQ(ret, RMW_RET_TIMEOUT) << rcutils_get_error_string().str; + } + + // Clean up: destroy guard conditions and wait set + for (size_t i = 0; i < number_of_guard_conditions; ++i) { + rmw_ret_t ret = rmw_destroy_guard_condition(guard_condition_ptrs[i]); + EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + } + + { + rmw_ret_t ret = rmw_destroy_wait_set(wait_set); + EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + } +} From aed50c27ab0425360c0a5f41e175c7a2da66416c Mon Sep 17 00:00:00 2001 From: Janosch Machowinski Date: Wed, 8 Jul 2026 13:56:35 +0000 Subject: [PATCH 2/6] fix(rmw_destroy_wait_set): Wait for thread to finsh and run multiple times Signed-off-by: Janosch Machowinski --- .../test/test_wait_set.cpp | 118 +++++++++++------- 1 file changed, 73 insertions(+), 45 deletions(-) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index a956193d7..7b4db5479 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -13,6 +13,7 @@ // limitations under the License. #include +#include #include "osrf_testing_tools_cpp/scope_exit.hpp" @@ -421,85 +422,112 @@ TEST_F(TestWaitSet, rmw_destroy_wait_set) */ TEST_F(TestWaitSet, rmw_wait_guard_conditions) { - constexpr size_t number_of_guard_conditions = 10u; - typedef rmw_guard_condition_t* guard_condition_ptr_t; - guard_condition_ptr_t guard_condition_ptrs[number_of_guard_conditions]; + constexpr size_t number_of_guard_conditions = 100u; + typedef rmw_guard_condition_t * guard_condition_ptr_t; + std::array guard_condition_ptrs; // Create the guard conditions - size_t ready_conditions_value = 0u; + std::vector guard_conditions_triggered(number_of_guard_conditions, false); + for (size_t i = 0; i < number_of_guard_conditions; ++i) { guard_condition_ptrs[i] = rmw_create_guard_condition(&context); - ready_conditions_value |= (1u << i); ASSERT_NE(nullptr, guard_condition_ptrs[i]) << rcutils_get_error_string().str; } // Create a wait set - rmw_wait_set_t* wait_set = rmw_create_wait_set(&context, number_of_guard_conditions); + rmw_wait_set_t * wait_set = rmw_create_wait_set(&context, number_of_guard_conditions); ASSERT_NE(nullptr, wait_set) << rcutils_get_error_string().str; - // Thread that triggers all guard conditions in reverse order after 100ms - std::thread trigger_thread([&]() { - std::this_thread::sleep_for(std::chrono::milliseconds(100)); - for (size_t i = number_of_guard_conditions; i > 0; --i) { - rmw_ret_t ret = rmw_trigger_guard_condition(guard_condition_ptrs[i - 1]); - EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; - } - }); + for(size_t runs = 0; runs < 100; runs++) { + guard_conditions_triggered = std::vector(number_of_guard_conditions, false); - // Prepare input arguments for rmw_wait - rmw_time_t timeout_argument = {0, 100000000}; // 100ms - rmw_guard_conditions_t guard_conditions; - INITIALIZE_ARRAY(guard_conditions, guard_condition, number_of_guard_conditions); + // Prepare input arguments for rmw_wait + rmw_time_t timeout_argument = {0, 100000000}; // 100ms + rmw_guard_conditions_t guard_conditions; + INITIALIZE_ARRAY(guard_conditions, guard_condition, number_of_guard_conditions); - // Wait up to number_of_guard_conditions times or until all guard conditions are triggered - for (size_t i = 0; i < number_of_guard_conditions; ++i) { // Prepare the input array for rmw_wait for (size_t j = 0; j < number_of_guard_conditions; ++j) { guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; } - // Wait for guard conditions to be triggered - rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, &timeout_argument); - if (RMW_RET_OK == ret) { - // Check which guard conditions are ready + + // Thread that triggers all guard conditions in reverse order after 100ms + std::thread trigger_thread([&]() { + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + for (size_t i = number_of_guard_conditions; i > 0; --i) { + rmw_ret_t ret = rmw_trigger_guard_condition(guard_condition_ptrs[i - 1]); + ASSERT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + } + }); + + // wait until we receive the first guard condition trigger + rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, + &timeout_argument); + ASSERT_EQ(ret, RMW_RET_OK) << "Failed to receive any guard condition trigger"; + + // Join the trigger thread + trigger_thread.join(); + + // mark triggered guard conditions as ready + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + if (guard_conditions.guard_conditions[j] != nullptr) { + // Check that condition was not already triggered + ASSERT_NE(guard_conditions_triggered[j], true); + // Mark this condition as triggered + guard_conditions_triggered[j] = true; + } + } + + // set timeout to almost zero there is not point in waiting longer from here on + timeout_argument.nsec = 1; + + if(std::ranges::any_of(guard_conditions_triggered, [](bool triggered) {return !triggered;})) { + // Prepare the input array for rmw_wait + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; + } + + // Wait for guard conditions to be triggered + rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, + &timeout_argument); + ASSERT_EQ(ret, RMW_RET_OK) << "Expected guard condition trigger but got timeout"; + + // This should collect the outstanding triggers now for (size_t j = 0; j < number_of_guard_conditions; ++j) { if (guard_conditions.guard_conditions[j] != nullptr) { // Check that condition was not already triggered - EXPECT_NE(ready_conditions_value & (1u << j), 0u); + ASSERT_NE(guard_conditions_triggered[j], true); // Mark this condition as triggered - ready_conditions_value &= ~(1u << j); + guard_conditions_triggered[j] = true; } } - - if (ready_conditions_value == 0u) { - // All guard conditions have been triggered - break; - } } - } - - // Join the trigger thread - trigger_thread.join(); - // Check that all guard conditions were triggered - EXPECT_EQ(ready_conditions_value, 0u) << "Not all guard conditions were triggered"; + // Check that all guard conditions were triggered + ASSERT_FALSE(std::ranges::any_of(guard_conditions_triggered, [](bool triggered) { + return !triggered; + })) << "Not all guard conditions were triggered"; - // Calling rmw_wait shall now return RMW_RET_TIMEOUT since all guard conditions have been triggered - { - for (size_t j = 0; j < number_of_guard_conditions; ++j) { - guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; + // Calling rmw_wait shall now return RMW_RET_TIMEOUT since all guard conditions + // have been triggered + { + for (size_t j = 0; j < number_of_guard_conditions; ++j) { + guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; + } + rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, + &timeout_argument); + ASSERT_EQ(ret, RMW_RET_TIMEOUT) << rcutils_get_error_string().str; } - rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, &timeout_argument); - EXPECT_EQ(ret, RMW_RET_TIMEOUT) << rcutils_get_error_string().str; } // Clean up: destroy guard conditions and wait set for (size_t i = 0; i < number_of_guard_conditions; ++i) { rmw_ret_t ret = rmw_destroy_guard_condition(guard_condition_ptrs[i]); - EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + ASSERT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; } { rmw_ret_t ret = rmw_destroy_wait_set(wait_set); - EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + ASSERT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; } } From 8a7a7771afe9277dae584dadbb891773a586b27f Mon Sep 17 00:00:00 2001 From: Miguel Company Date: Mon, 13 Jul 2026 08:00:05 +0200 Subject: [PATCH 3/6] Apply changes from review. Signed-off-by: Miguel Company --- test_rmw_implementation/test/test_wait_set.cpp | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index 7b4db5479..b571fc290 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -451,19 +451,19 @@ TEST_F(TestWaitSet, rmw_wait_guard_conditions) guard_conditions.guard_conditions[j] = guard_condition_ptrs[j]->data; } - // Thread that triggers all guard conditions in reverse order after 100ms + // Thread that triggers all guard conditions in reverse order after 10ms std::thread trigger_thread([&]() { std::this_thread::sleep_for(std::chrono::milliseconds(10)); for (size_t i = number_of_guard_conditions; i > 0; --i) { rmw_ret_t ret = rmw_trigger_guard_condition(guard_condition_ptrs[i - 1]); - ASSERT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; + EXPECT_EQ(ret, RMW_RET_OK) << rcutils_get_error_string().str; } }); // wait until we receive the first guard condition trigger rmw_ret_t ret = rmw_wait(nullptr, &guard_conditions, nullptr, nullptr, nullptr, wait_set, &timeout_argument); - ASSERT_EQ(ret, RMW_RET_OK) << "Failed to receive any guard condition trigger"; + EXPECT_EQ(ret, RMW_RET_OK) << "Failed to receive any guard condition trigger"; // Join the trigger thread trigger_thread.join(); From 8b35d753ab61e1d2cfd15de658ca6f7f852572e8 Mon Sep 17 00:00:00 2001 From: Janosch Machowinski Date: Fri, 10 Jul 2026 09:53:57 +0000 Subject: [PATCH 4/6] fix: added missing header Signed-off-by: Janosch Machowinski --- test_rmw_implementation/test/test_wait_set.cpp | 1 + 1 file changed, 1 insertion(+) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index b571fc290..cc3e416a1 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -14,6 +14,7 @@ #include #include +#include #include "osrf_testing_tools_cpp/scope_exit.hpp" From 093973136c72f48953c1782ef0bb1b853286827e Mon Sep 17 00:00:00 2001 From: Janosch Machowinski Date: Fri, 17 Jul 2026 14:48:18 +0200 Subject: [PATCH 5/6] chore: Comments from review Signed-off-by: Janosch Machowinski --- test_rmw_implementation/test/test_wait_set.cpp | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index cc3e416a1..18c348735 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -428,7 +428,8 @@ TEST_F(TestWaitSet, rmw_wait_guard_conditions) std::array guard_condition_ptrs; // Create the guard conditions - std::vector guard_conditions_triggered(number_of_guard_conditions, false); + std::array guard_conditions_triggered; + guard_conditions_triggered.fill(false); for (size_t i = 0; i < number_of_guard_conditions; ++i) { guard_condition_ptrs[i] = rmw_create_guard_condition(&context); @@ -440,10 +441,10 @@ TEST_F(TestWaitSet, rmw_wait_guard_conditions) ASSERT_NE(nullptr, wait_set) << rcutils_get_error_string().str; for(size_t runs = 0; runs < 100; runs++) { - guard_conditions_triggered = std::vector(number_of_guard_conditions, false); + guard_conditions_triggered.fill(false); // Prepare input arguments for rmw_wait - rmw_time_t timeout_argument = {0, 100000000}; // 100ms + rmw_time_t timeout_argument = {2, 0}; // 2 seconds rmw_guard_conditions_t guard_conditions; INITIALIZE_ARRAY(guard_conditions, guard_condition, number_of_guard_conditions); From d85401d8def88b78dfff405807d701b9ba8084d5 Mon Sep 17 00:00:00 2001 From: Janosch Machowinski Date: Fri, 17 Jul 2026 19:04:49 +0200 Subject: [PATCH 6/6] fix: Fix timeout in collect phase Signed-off-by: Janosch Machowinski --- test_rmw_implementation/test/test_wait_set.cpp | 1 + 1 file changed, 1 insertion(+) diff --git a/test_rmw_implementation/test/test_wait_set.cpp b/test_rmw_implementation/test/test_wait_set.cpp index 18c348735..fd962c90a 100644 --- a/test_rmw_implementation/test/test_wait_set.cpp +++ b/test_rmw_implementation/test/test_wait_set.cpp @@ -481,6 +481,7 @@ TEST_F(TestWaitSet, rmw_wait_guard_conditions) } // set timeout to almost zero there is not point in waiting longer from here on + timeout_argument.sec = 0; timeout_argument.nsec = 1; if(std::ranges::any_of(guard_conditions_triggered, [](bool triggered) {return !triggered;})) {