mctp: add rate-limiting cooldown for queries Add a rate-limiting cooldown of 5 seconds for back-to-back "Add" events on the same endpoint to prevent flooding the Object Mapper. An Add event triggered immediately after a Priming scan or a Remove event is not impacted by this cooldown. New unit tests are added to validate these rate-limiting behaviors. Misc: - Use TEST_SRC_DIR macro in tests to avoid hardcoded paths. - Increase Meson test timeout to 120 seconds for the test suite. Change-Id: I283acc82d6dc67fc0d09ba3187153845ed53801c Google-Bug-Id: 490106522 Signed-off-by: Hao Jiang <jianghao@google.com>
diff --git a/src/MctpUtil.cpp b/src/MctpUtil.cpp index 4b3ab4b..8d1e82c 100644 --- a/src/MctpUtil.cpp +++ b/src/MctpUtil.cpp
@@ -28,6 +28,8 @@ State state = State::None; bool querying = false; std::string emConfigPath; + std::chrono::steady_clock::time_point lastQueryTime = + std::chrono::steady_clock::time_point::min(); }; static std::map<std::string, EndpointState> endpointStates; @@ -38,14 +40,47 @@ static void performEmConfigQuery( const std::shared_ptr<sdbusplus::asio::connection>& conn, - const std::string& endpointPath, const std::string& emConfigPath) + const std::string& endpointPath, const std::string& emConfigPath, + bool bypassCooldown = false) { + auto& state = endpointStates[endpointPath]; + auto now = std::chrono::steady_clock::now(); + static constexpr auto cooldown = std::chrono::seconds(5); + + state.querying = true; // Protect the window! + + // Apply rate-limiting cooldown only if this endpoint has been queried + // before (to avoid subtracting time_point::min() which causes overflow), if + // we are not explicitly bypassing it (e.g. for post-priming queries), and + // if we are within the 5-second cooldown window. + if (state.lastQueryTime != std::chrono::steady_clock::time_point::min() && + !bypassCooldown && (now - state.lastQueryTime < cooldown)) + { + auto delay = cooldown - (now - state.lastQueryTime); + + auto timer = + std::make_shared<boost::asio::steady_timer>(conn->get_io_context()); + timer->expires_after(delay); + + timer->async_wait([conn, endpointPath, emConfigPath, + timer](const boost::system::error_code& ec) { + if (ec) + { + return; + } + performEmConfigQuery(conn, endpointPath, emConfigPath); + }); + return; + } + conn->async_method_call( [conn, endpointPath, emConfigPath](const boost::system::error_code& ec, const ManagedObjectType& objects) { auto& state = endpointStates[endpointPath]; state.querying = false; + state.lastQueryTime = std::chrono::steady_clock::now(); // Update time! + State lastState = state.state; state.state = State::None; // Reset for next query @@ -102,7 +137,8 @@ } static void triggerDeferredQueries( - const std::shared_ptr<sdbusplus::asio::connection>& conn) + const std::shared_ptr<sdbusplus::asio::connection>& conn, + bool isPriming = false) { for (auto& [path, state] : endpointStates) { @@ -113,9 +149,8 @@ if (state.state == State::Add) { - state.querying = true; state.state = State::None; - performEmConfigQuery(conn, path, state.emConfigPath); + performEmConfigQuery(conn, path, state.emConfigPath, isPriming); } else if (state.state == State::Remove) { @@ -283,7 +318,7 @@ if (*barrier == 0) { endpointStates["__global__"].querying = false; - triggerDeferredQueries(conn); + triggerDeferredQueries(conn, true); } };
diff --git a/tests/meson.build b/tests/meson.build index cc585af..455d538 100644 --- a/tests/meson.build +++ b/tests/meson.build
@@ -104,11 +104,12 @@ 'test_mctp_util', 'test_MctpUtil.cpp', '../src/MctpUtil.cpp', - cpp_args: ['-UBOOST_ASIO_NO_DEPRECATED', '-UBOOST_ASIO_DISABLE_THREADS', '-UBOOST_ASIO_HAS_IO_URING', '-DBUILDDIR='+ meson.current_build_dir()], + cpp_args: ['-UBOOST_ASIO_NO_DEPRECATED', '-UBOOST_ASIO_DISABLE_THREADS', '-UBOOST_ASIO_HAS_IO_URING', '-DBUILDDIR='+ meson.current_build_dir(), '-DTEST_SRC_DIR="' + meson.current_source_dir() + '"'], dependencies: [ut_deps_list, nlohmann_json], implicit_include_directories: false, include_directories: '../src', - ) + ), + timeout: 120 ) test( 'test_nvme_cache',
diff --git a/tests/test_MctpUtil.cpp b/tests/test_MctpUtil.cpp index 649d8ac..e939779 100644 --- a/tests/test_MctpUtil.cpp +++ b/tests/test_MctpUtil.cpp
@@ -75,21 +75,19 @@ // Clear global map before each test mctpEndpointConfigMap.clear(); - std::cout << "--- Starting Python Mock ---\n"; mockPid = fork(); if (mockPid == 0) { // Child process - std::string mockScript = std::string(xstr(BUILDDIR)) + - "/../../tests/mock_mctpd.py"; // clang-format off - execlp("python3", "python3", mockScript.c_str(), nullptr); // NOLINT(cppcoreguidelines-pro-type-vararg) + std::string scriptPath = TEST_SRC_DIR "/mock_mctpd.py"; + execlp("python3", "python3", scriptPath.c_str(), nullptr); // NOLINT(cppcoreguidelines-pro-type-vararg) // clang-format on exit(1); } // Give mock some time to start and claim names! - sleep(1); // NOLINT + sleep(3); // NOLINT } void TearDown() override @@ -310,7 +308,7 @@ runCmd(addCmd); // Wait for queries to complete with a timeout (up to 5 seconds) - for (int i = 0; i < 50; i++) + for (int i = 0; i < 70; i++) { io.poll(); usleep(100000); // NOLINT(cert-api50-c) @@ -631,6 +629,323 @@ EXPECT_FALSE(it->second.empty()); } +// 7. Add delay for a single ep Add to Add event. i.e. Add -> Remove -> Add +// (this add should be cooled down). +TEST_F(MctpUtilTest, Cooldown_AddRemoveAdd_IsCooledDown) +{ + std::string endpointPath = + "/au/com/codeconstruct/mctp1/networks/1/endpoints/test7"; + std::string emConfigPath = + "/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device"; + + setupMctpEndpointListener(conn); + + bool queryStarted = false; + int queryCompletedCount = 0; + + const std::string queryStartedSpec = + "type='signal',interface='com.example.Control',member='QueryStarted'"; + auto queryStartedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryStartedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryStarted received\n"; + queryStarted = true; + }); + + const std::string queryCompletedSpec = + "type='signal',interface='com.example.Control',member='QueryCompleted'"; + auto queryCompletedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryCompletedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryCompleted received\n"; + queryCompletedCount++; + }); + + // Trigger FIRST Add! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerAdd failed: " << ec.message() << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd", endpointPath, emConfigPath); + + // Wait for FIRST query to complete! + for (int i = 0; i < 50; i++) + { + while (io.poll() > 0) + {} // Drain all ready handlers! + usleep(100000); // 100ms + if (queryCompletedCount >= 1) + { + break; + } + } + ASSERT_EQ(queryCompletedCount, 1); + + std::cout << "[DEBUG] Triggering Remove and Add (racing!)\n"; + // Trigger Remove! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerRemove failed: " << ec.message() + << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerRemove", endpointPath); + + // Trigger SECOND Add! (Immediately after!). + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerAdd failed: " << ec.message() << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd", endpointPath, emConfigPath); + + // Wait for SECOND query to complete! + // We expect it to be DELAYED by 5 seconds! + auto startTime = std::chrono::steady_clock::now(); + for (int i = 0; i < 100; i++) + { + while (io.poll() > 0) + {} // Drain all ready handlers! + usleep(100000); // 100ms + if (queryCompletedCount >= 2) + { + break; + } + } + auto endTime = std::chrono::steady_clock::now(); + auto duration = + std::chrono::duration_cast<std::chrono::seconds>(endTime - startTime) + .count(); + + EXPECT_EQ(queryCompletedCount, 2); + EXPECT_GE(duration, 4); // Expect at least 4-5 seconds delay! + + // Verify map populated + auto it = mctpEndpointConfigMap.find(endpointPath); + ASSERT_NE(it, mctpEndpointConfigMap.end()); +} + +// 8. Priming scan does block its immediate following Add. i.e. Priming -> +// Remove -> Add (not cooldown) +TEST_F(MctpUtilTest, Cooldown_PrimingRemoveAdd_IsNormalized) +{ + std::string endpointPath = + "/au/com/codeconstruct/mctp1/networks/1/endpoints/test8"; + std::string emConfigPath = + "/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device"; + + // Pre-populate association by triggering it BEFORE listener starts! + auto m = conn->new_method_call("xyz.openbmc_project.Mctp", + "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd"); + m.append(endpointPath, emConfigPath); + conn->call(m); + + // Wait for mapperx to see it + bool found = false; + for (int i = 0; i < 300; i++) + { + if (checkPathInMapper(endpointPath + "/configured_by")) + { + found = true; + break; + } + usleep(100000); // 100ms + } + ASSERT_TRUE(found); + + // Now setup listener (this starts priming!) + setupMctpEndpointListener(conn); + + bool queryStarted = false; + int queryCompletedCount = 0; + + const std::string queryStartedSpec = + "type='signal',interface='com.example.Control',member='QueryStarted'"; + auto queryStartedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryStartedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryStarted received\n"; + queryStarted = true; + }); + + const std::string queryCompletedSpec = + "type='signal',interface='com.example.Control',member='QueryCompleted'"; + auto queryCompletedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryCompletedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryCompleted received\n"; + queryCompletedCount++; + }); + + // Wait for query to start (triggered by priming!) + while (!queryStarted) + { + io.poll(); + usleep(10000); // 10ms + } + + std::cout << "[DEBUG] Triggering Remove and Add (racing with priming!)\n"; + // Trigger Remove event! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerRemove failed: " << ec.message() + << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerRemove", endpointPath); + + // Trigger Add event! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerAdd failed: " << ec.message() << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd", endpointPath, emConfigPath); + + // Wait for queries to complete! + auto startTime = std::chrono::steady_clock::now(); + for (int i = 0; i < 100; i++) + { + while (io.poll() > 0) + {} // Drain all ready handlers! + usleep(100000); // 100ms + if (queryCompletedCount >= 2) + { + break; + } + } + auto endTime = std::chrono::steady_clock::now(); + auto duration = + std::chrono::duration_cast<std::chrono::seconds>(endTime - startTime) + .count(); + + EXPECT_EQ(queryCompletedCount, 2); + EXPECT_LT(duration, + 8); // Expect it to complete without additional 5s cooldown! + + // Verify map populated + auto it = mctpEndpointConfigMap.find(endpointPath); + ASSERT_NE(it, mctpEndpointConfigMap.end()); +} + +// 9. Remove has no cooldown (i.e. Remove -> Add -> Remove (no cooldown). +TEST_F(MctpUtilTest, Cooldown_RemoveAddRemove_NoCooldown) +{ + std::string endpointPath = + "/au/com/codeconstruct/mctp1/networks/1/endpoints/test9"; + std::string emConfigPath = + "/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device"; + + setupMctpEndpointListener(conn); + + bool queryStarted = false; + int queryCompletedCount = 0; + + const std::string queryStartedSpec = + "type='signal',interface='com.example.Control',member='QueryStarted'"; + auto queryStartedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryStartedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryStarted received\n"; + queryStarted = true; + }); + + const std::string queryCompletedSpec = + "type='signal',interface='com.example.Control',member='QueryCompleted'"; + auto queryCompletedMatch = std::make_unique<sdbusplus::bus::match_t>( + static_cast<sdbusplus::bus_t&>(*conn), queryCompletedSpec, + [&](sdbusplus::message_t&) { + std::cout << "[TEST] QueryCompleted received\n"; + queryCompletedCount++; + }); + + // Trigger FIRST Add! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerAdd failed: " << ec.message() << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd", endpointPath, emConfigPath); + + // Wait for FIRST query to complete! + for (int i = 0; i < 50; i++) + { + while (io.poll() > 0) + {} // Drain all ready handlers! + usleep(100000); // 100ms + if (queryCompletedCount >= 1) + { + break; + } + } + ASSERT_EQ(queryCompletedCount, 1); + + std::cout << "[DEBUG] Triggering Remove -> Add -> Remove\n"; + // Trigger Remove! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerRemove failed: " << ec.message() + << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerRemove", endpointPath); + + // Trigger Add! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerAdd failed: " << ec.message() << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerAdd", endpointPath, emConfigPath); + + // Trigger Remove! + conn->async_method_call( + [](const boost::system::error_code& ec, const std::string&) { + if (ec) + { + std::cout << "[TEST] TriggerRemove failed: " << ec.message() + << "\n"; + } + }, "xyz.openbmc_project.Mctp", "/au/com/codeconstruct/mctp1", + "com.example.Control", "TriggerRemove", endpointPath); + + // Wait for queries to complete! + for (int i = 0; i < 100; i++) + { + while (io.poll() > 0) + {} // Drain all ready handlers! + usleep(100000); // 100ms + if (queryCompletedCount >= 2) + { + break; + } + } + + EXPECT_EQ(queryCompletedCount, 2); + + // Verify map is EMPTY! + auto it = mctpEndpointConfigMap.find(endpointPath); + EXPECT_EQ(it, mctpEndpointConfigMap.end()); +} + int main(int argc, char** argv) { ::testing::InitGoogleTest(&argc, argv);