tests: Add unit tests for MctpUtil to handle concurrent D-Bus events

Add unit tests for MctpUtil to verify that concurrent Add and Remove
events arriving while a D-Bus query is in-flight are handled correctly.
The tests catch race conditions where state changes occur while the
application is asynchronously waiting for Entity Manager configurations:

* Remove event during an active query cancels the operation.
* Rapid succession of Add/Remove/Add events
  correctly schedules a new query.

Change-Id: I553bdf8cc656f8cc5af1898370536756f55253d0
Google-Bug-Id: 490106522
Signed-off-by: Hao Jiang <jianghao@google.com>
diff --git a/meson.build b/meson.build
index b31ce7c..185cb0c 100644
--- a/meson.build
+++ b/meson.build
@@ -128,7 +128,6 @@
 if get_option('nvme_shmem').enabled()
     add_project_arguments('-DNVMED_ENABLE_SHMEM', language: 'cpp')
     add_project_arguments('-fconstexpr-ops-limit=134217728', language: 'cpp')
-
     shared_mem_dep = dependency('shared_memory', required: true)
     absl_status_dep = dependency('absl_status', required: true)
 endif
diff --git a/src/MctpUtil.cpp b/src/MctpUtil.cpp
index add107a..299f0fe 100644
--- a/src/MctpUtil.cpp
+++ b/src/MctpUtil.cpp
@@ -12,7 +12,7 @@
 
 #include <chrono>
 #include <filesystem>
-#include <iostream>
+#include <string>
 
 std::map<std::string, SensorData> mctpEndpointConfigMap;
 
@@ -39,9 +39,6 @@
     const std::shared_ptr<sdbusplus::asio::connection>& conn,
     const std::string& endpointPath, const std::string& emConfigPath)
 {
-    std::filesystem::path p(emConfigPath);
-    std::string parentPath = p.parent_path().string();
-
     conn->async_method_call(
         [conn, endpointPath, emConfigPath](const boost::system::error_code& ec,
                                            const ManagedObjectType& objects) {
@@ -72,6 +69,7 @@
 
         auto objIt =
             objects.find(sdbusplus::message::object_path(emConfigPath));
+
         if (objIt == objects.end())
         {
             lg2::warning("EM config path {PATH} not found in managed objects",
@@ -85,15 +83,20 @@
                 "Recorded MCTP endpoint {ENDPOINT} with config from {PATH}",
                 "ENDPOINT", endpointPath, "PATH", emConfigPath);
             mctpEndpointConfigMap[endpointPath] = objIt->second;
+            lg2::info("DEBUG: Map populated, new size: {SIZE}", "SIZE",
+                      mctpEndpointConfigMap.size());
         }
 
         if (lastState == State::Add)
         {
+            lg2::info(
+                "DEBUG: Triggering deferred query for endpoint {ENDPOINT}",
+                "ENDPOINT", endpointPath);
             state.querying = true;
             performEmConfigQuery(conn, endpointPath, emConfigPath);
         }
     },
-        "xyz.openbmc_project.EntityManager", parentPath,
+        "xyz.openbmc_project.EntityManager", "/xyz/openbmc_project/inventory",
         "org.freedesktop.DBus.ObjectManager", "GetManagedObjects");
 }
 
@@ -102,7 +105,7 @@
 {
     // Match 1: Listen for Associations from mctp-reactor
     const std::string associationMatchSpec =
-        "type='signal',interface='org.freedesktop.DBus.ObjectManager',member='InterfacesAdded',path_namespace='/au/com/codeconstruct/mctp1'";
+        "type='signal',interface='org.freedesktop.DBus.ObjectManager',member='InterfacesAdded',arg0path='/au/com/codeconstruct/mctp1/'";
 
     associationMatch = std::make_unique<sdbusplus::bus::match_t>(
         static_cast<sdbusplus::bus_t&>(*conn), associationMatchSpec,
@@ -115,6 +118,8 @@
                 std::variant<std::string, std::vector<std::string>>>>
             interfaces;
         msg.read(path, interfaces);
+        lg2::info("DEBUG: InterfacesAdded callback for {PATH}", "PATH",
+                  path.str);
 
         auto it = interfaces.find("xyz.openbmc_project.Association");
         if (it == interfaces.end())
@@ -136,6 +141,7 @@
         }
 
         std::string emConfigPath = endpoints->front();
+
         std::string endpointPath = path.str;
         constexpr std::string_view suffix = "/configured_by";
         if (endpointPath.ends_with("/configured_by"))
@@ -156,9 +162,8 @@
         performEmConfigQuery(conn, endpointPath, emConfigPath);
     });
 
-    // Match 2: Listen for Associations (Removed)
     const std::string associationRemoveMatchSpec =
-        "type='signal',interface='org.freedesktop.DBus.ObjectManager',member='InterfacesRemoved',path_namespace='/au/com/codeconstruct/mctp1'";
+        "type='signal',interface='org.freedesktop.DBus.ObjectManager',member='InterfacesRemoved',arg0path='/au/com/codeconstruct/mctp1/'";
 
     associationRemoveMatch = std::make_unique<sdbusplus::bus::match_t>(
         static_cast<sdbusplus::bus_t&>(*conn), associationRemoveMatchSpec,
@@ -303,3 +308,9 @@
     }
     return busInfo;
 }
+
+void cleanupMctpEndpointListener()
+{
+    associationMatch.reset();
+    associationRemoveMatch.reset();
+}
diff --git a/src/MctpUtil.hpp b/src/MctpUtil.hpp
index b34fd86..282e7e2 100644
--- a/src/MctpUtil.hpp
+++ b/src/MctpUtil.hpp
@@ -18,5 +18,6 @@
 
 void setupMctpEndpointListener(
     const std::shared_ptr<sdbusplus::asio::connection>& conn);
+void cleanupMctpEndpointListener();
 
 BusInfo extractBusInfo(const SensorData& configData);
diff --git a/tests/meson.build b/tests/meson.build
index 69e22b7..cc585af 100644
--- a/tests/meson.build
+++ b/tests/meson.build
@@ -99,6 +99,18 @@
         env:'LD_LIBRARY_PATH=/usr/local/lib'
     )
     test(
+        'test_mctp_util',
+        executable(
+            '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()],
+            dependencies: [ut_deps_list, nlohmann_json],
+            implicit_include_directories: false,
+            include_directories: '../src',
+        )
+    )
+    test(
         'test_nvme_cache',
         executable(
             'test_nvme_cache',
diff --git a/tests/mock_mctpd.py b/tests/mock_mctpd.py
new file mode 100644
index 0000000..1d99a03
--- /dev/null
+++ b/tests/mock_mctpd.py
@@ -0,0 +1,164 @@
+"""
+mock_mctpd.py - Mock server for mctpd and Entity Manager in MctpUtil tests.
+
+This script simulates:
+1. mctpd emitting InterfacesAdded signals for MCTP endpoints.
+2. Entity Manager responding to GetManagedObjects queries.
+
+It also provides a control interface to trigger events from the test:
+- TriggerAdd(endpoint_path, em_config_path)
+- TriggerRemove(endpoint_path)
+
+To support precise synchronization without modifying product code, this mock
+server emits custom D-Bus signals on the 'com.example.Control' interface:
+- QueryStarted: Emitted when GetManagedObjects is called.
+- QueryCompleted: Emitted when GetManagedObjects returns.
+
+Usage in tests:
+1. Start this script as a subprocess in SetUp().
+2. Use `busctl call` to trigger events via TriggerAdd/TriggerRemove.
+3. Listen for QueryStarted/QueryCompleted to synchronize test steps.
+"""
+import dbus
+import dbus.service
+import dbus.mainloop.glib
+from gi.repository import GLib
+import sys
+import time
+import signal
+
+
+class MockEndpoint(dbus.service.Object):
+    def __init__(self, bus, path, em_config_path):
+        dbus.service.Object.__init__(self, bus, path)
+        self.associations = dbus.Array([
+            dbus.Struct(('configured_by', 'endpoint', em_config_path), signature='(sss)')
+        ], signature='(sss)')
+
+    @dbus.service.method('org.freedesktop.DBus.Properties', in_signature='ss', out_signature='v')
+    def Get(self, interface_name, property_name):
+        if interface_name == 'xyz.openbmc_project.Association.Definitions' and property_name == 'Associations':
+            return dbus.LowLevel.Variant(self.associations, 'a(sss)')
+        raise dbus.exceptions.DBusException('org.freedesktop.DBus.Error.InvalidArgs')
+
+    @dbus.service.method('org.freedesktop.DBus.Properties', in_signature='s', out_signature='a{sv}')
+    def GetAll(self, interface_name):
+        if interface_name == 'xyz.openbmc_project.Association.Definitions':
+            return {'Associations': dbus.LowLevel.Variant(self.associations, 'a(sss)')}
+        return {}
+
+
+class MockTarget(dbus.service.Object):
+    def __init__(self, bus):
+        dbus.service.Object.__init__(self, bus, '/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device')
+
+    @dbus.service.method('org.freedesktop.DBus.Properties', in_signature='s', out_signature='a{sv}')
+    def GetAll(self, interface_name):
+        if interface_name == 'xyz.openbmc_project.Configuration.BusInfo':
+            return {
+                'BusType': 'USB',
+                'Port': '1.2.5',
+                'Configuration': dbus.UInt32(1)
+            }
+        return {}
+
+
+class MockMctpd(dbus.service.Object):
+    def __init__(self, bus):
+        dbus.service.Object.__init__(self, bus, '/au/com/codeconstruct/mctp1')
+        self.bus = bus
+        self.endpoints = {}
+
+    @dbus.service.signal('org.freedesktop.DBus.ObjectManager', signature='oa{sa{sv}}')
+    def InterfacesAdded(self, object_path, interfaces_and_properties):
+        print(f"Emitting InterfacesAdded for {object_path}")
+        sys.stdout.flush()
+
+    @dbus.service.signal('org.freedesktop.DBus.ObjectManager', signature='oas')
+    def InterfacesRemoved(self, object_path, interfaces):
+        pass
+
+    @dbus.service.method('com.example.Control', in_signature='ss')
+    def TriggerAdd(self, endpoint_path, em_config_path):
+        if endpoint_path not in self.endpoints:
+            self.endpoints[endpoint_path] = MockEndpoint(self.bus, endpoint_path, em_config_path)
+
+        self.InterfacesAdded(dbus.ObjectPath(endpoint_path), {
+            'xyz.openbmc_project.Association.Definitions': {
+                'Associations': dbus.Array([
+                    dbus.Struct(('configured_by', 'endpoint', em_config_path), signature='(sss)')
+                ], signature='(sss)')
+            }
+        })
+        return "OK"
+
+    @dbus.service.method('com.example.Control', in_signature='s')
+    def TriggerRemove(self, endpoint_path):
+        print("TriggerRemove called for " + endpoint_path)
+        sys.stdout.flush()
+        if endpoint_path in self.endpoints:
+            self.InterfacesRemoved(dbus.ObjectPath(endpoint_path), ['xyz.openbmc_project.Association.Definitions'])
+            self.endpoints[endpoint_path].remove_from_connection()
+            del self.endpoints[endpoint_path]
+            return "OK"
+        return "Not Found"
+
+
+class MockEM(dbus.service.Object):
+    def __init__(self, bus):
+        dbus.service.Object.__init__(self, bus, '/xyz/openbmc_project/inventory')
+
+    @dbus.service.signal('com.example.Control', signature='')
+    def QueryStarted(self):
+        pass
+
+    @dbus.service.signal('com.example.Control', signature='')
+    def QueryCompleted(self):
+        pass
+
+    @dbus.service.method('org.freedesktop.DBus.ObjectManager', out_signature='a{oa{sa{sv}}}')
+    def GetManagedObjects(self):
+        self.QueryStarted()
+        print("GetManagedObjects called, delaying reply...")
+        sys.stdout.flush()
+        time.sleep(1)
+        self.QueryCompleted()
+        return {
+            dbus.ObjectPath('/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device'): {
+                'xyz.openbmc_project.Configuration.BusInfo': {
+                    'BusType': 'USB',
+                    'Port': '1.2.5',
+                    'Configuration': dbus.UInt32(1)
+                }
+            }
+        }
+
+
+def main():
+    dbus.mainloop.glib.DBusGMainLoop(set_as_default=True)
+    import os
+    if 'DBUS_SESSION_BUS_ADDRESS' in os.environ:
+        bus = dbus.SessionBus()
+    else:
+        bus = dbus.SystemBus()
+    loop = GLib.MainLoop()
+
+    def signal_handler(sig, frame):
+        print("Received SIGTERM, shutting down gracefully...")
+        loop.quit()
+
+    signal.signal(signal.SIGTERM, signal_handler)
+
+    _target = MockTarget(bus)  # noqa: F841
+    _mctpd = MockMctpd(bus)  # noqa: F841
+    _em = MockEM(bus)  # noqa: F841
+
+    _name_em = dbus.service.BusName("xyz.openbmc_project.EntityManager", bus)  # noqa: F841
+    _name_mctp = dbus.service.BusName("xyz.openbmc_project.Mctp", bus)  # noqa: F841
+    print("Mock mctpd and EM running...")
+    sys.stdout.flush()
+    loop.run()
+
+
+if __name__ == '__main__':
+    main()
diff --git a/tests/setup.sh b/tests/setup.sh
index b394faa..a593378 100644
--- a/tests/setup.sh
+++ b/tests/setup.sh
@@ -33,3 +33,4 @@
 
 sleep 1
 busctl $USER_FLAG tree xyz.openbmc_project.ObjectMapper
+busctl $USER_FLAG monitor xyz.openbmc_project.ObjectMapper > /tmp/mapper_monitor.log 2>&1 < /dev/null &
diff --git a/tests/test_MctpUtil.cpp b/tests/test_MctpUtil.cpp
new file mode 100644
index 0000000..f643a1a
--- /dev/null
+++ b/tests/test_MctpUtil.cpp
@@ -0,0 +1,339 @@
+/**
+ * @file test_MctpUtil.cpp
+ * @brief Unit tests for MctpUtil.cpp
+ *
+ * These tests verify the deferred state handling in MctpUtil when processing
+ * MCTP endpoint discovery events. Specifically, it tests that:
+ * 1. A single Add event populates the map.
+ * 2. A Remove event arriving while a query is in progress correctly cancels
+ *    or rolls back the operation, leaving the map empty.
+ * 3. An Add -> Remove -> Add sequence in rapid succession correctly schedules
+ *    a new query and leaves the map populated at the end.
+ *
+ * Synchronization Design:
+ * To avoid impacting production code, the tests use custom D-Bus signals
+ * emitted by the mock server (mock_mctpd.py) to know exactly when a query
+ * starts and completes. This allows precise orchestration of events without
+ * sleeps or production code hooks.
+ *
+ * The tests run in a single thread using `io.poll()` to match the production
+ * execution model.
+ */
+#include "MctpUtil.hpp"
+#include "Utils.hpp"
+
+#include <sys/types.h>
+#include <sys/wait.h>
+#include <unistd.h>
+
+#include <boost/asio/io_context.hpp>
+#include <sdbusplus/asio/connection.hpp>
+#include <sdbusplus/asio/object_server.hpp>
+#include <sdbusplus/bus/match.hpp>
+
+#include <array>
+#include <cstdlib>
+#include <iostream>
+#include <memory>
+#include <string>
+#include <thread>
+#include <vector>
+
+#include <gtest/gtest.h>
+
+#define xstr(s) str(s) // NOLINT
+#define str(s) #s      // NOLINT
+
+static void runCmd(const std::string& cmd)
+{
+    std::string finalCmd = cmd;
+    const char* dbusAddr = std::getenv("DBUS_SESSION_BUS_ADDRESS");
+    if (dbusAddr != nullptr)
+    {
+        size_t pos = finalCmd.find("busctl ");
+        if (pos != std::string::npos)
+        {
+            finalCmd.replace(pos, 7,
+                             "busctl --address=" + std::string(dbusAddr) + " ");
+        }
+    }
+    int rc = system(finalCmd.c_str()); // NOLINT(cert-env33-c)
+    ASSERT_EQ(rc, 0);
+}
+
+class MctpUtilTest : public ::testing::Test
+{
+  protected:
+    boost::asio::io_context io;
+    std::shared_ptr<sdbusplus::asio::connection> conn;
+    pid_t mockPid = -1;
+
+    void SetUp() override
+    {
+        io.restart();
+        conn = std::make_shared<sdbusplus::asio::connection>(io);
+
+        // 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)
+            // clang-format on
+            exit(1);
+        }
+        // Give mock some time to start and claim names!
+        sleep(1); // NOLINT
+    }
+
+    void TearDown() override
+    {
+        cleanupMctpEndpointListener();
+        if (mockPid > 0)
+        {
+            kill(mockPid, SIGTERM);
+            waitpid(mockPid, nullptr, 0);
+        }
+    }
+
+    static bool checkPathInMapper(const std::string& path)
+    {
+        std::string cmd = "busctl tree xyz.openbmc_project.ObjectMapper";
+        const char* dbusAddr = std::getenv("DBUS_SESSION_BUS_ADDRESS");
+        if (dbusAddr != nullptr)
+        {
+            cmd = "busctl --address=" + std::string(dbusAddr) +
+                  " tree xyz.openbmc_project.ObjectMapper";
+        }
+        FILE* pipe = popen(cmd.c_str(), "r"); // NOLINT(cert-env33-c)
+        if (pipe == nullptr)
+        {
+            return false;
+        }
+        std::array<char, 128> buffer{};
+        std::string result;
+        while (fgets(buffer.data(), buffer.size(), pipe) != nullptr)
+        {
+            result += buffer.data();
+        }
+        pclose(pipe);
+        return result.find(path) != std::string::npos;
+    }
+};
+
+// 1. A single Add event will trigger only one O-M query. And add the correct
+// result in mctpEndpointConfigMap
+TEST_F(MctpUtilTest, SingleAdd_PopulatesMap)
+{
+    std::string endpointPath =
+        "/au/com/codeconstruct/mctp1/networks/1/endpoints/test1";
+    std::string emConfigPath =
+        "/xyz/openbmc_project/inventory/system/board/MockBoard/mctp_device";
+
+    // 2.a setupMctpEndpointListener() first in the test.
+    setupMctpEndpointListener(conn);
+
+    // 2.b setup the busctl event listener to all O-M event.
+    // We use a C++ match listener for all InterfacesAdded signals!
+    int queryCompletedCount = 0;
+
+    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&) { queryCompletedCount++; });
+
+    // Wait for priming loop to complete (finding nothing)
+    io.poll();
+    usleep(100000); // 100ms
+    io.poll();
+
+    // Trigger Add event via D-Bus!
+    std::string addCmd =
+        "busctl call xyz.openbmc_project.Mctp /au/com/codeconstruct/mctp1 com.example.Control TriggerAdd ss \"" +
+        endpointPath + "\" \"" + emConfigPath + "\"";
+    runCmd(addCmd);
+
+    // Wait for query to complete with a timeout (up to 5 seconds)
+    for (int i = 0; i < 50; i++)
+    {
+        io.poll();
+        usleep(100000); // NOLINT(cert-api50-c)
+        if (queryCompletedCount >= 1)
+        {
+            break;
+        }
+    }
+
+    // Wait a bit more to see if any extra signals arrive
+    usleep(500000); // 500ms
+    io.poll();
+
+    EXPECT_EQ(queryCompletedCount, 1);
+
+    // Wait a bit more for MctpUtil to process the reply!
+    usleep(100000); // 100ms
+
+    // Verify map populated
+    auto it = mctpEndpointConfigMap.find(endpointPath);
+    ASSERT_NE(it, mctpEndpointConfigMap.end());
+    EXPECT_FALSE(it->second.empty());
+}
+
+// 2. Add -> Remove during query
+TEST_F(MctpUtilTest, AddRemove_ClearsMap)
+{
+    std::string endpointPath =
+        "/au/com/codeconstruct/mctp1/networks/1/endpoints/test2";
+    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&) { 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&) { queryCompletedCount++; });
+
+    // Wait for priming loop to complete (finding nothing)
+    io.poll();
+    usleep(100000); // 100ms
+    io.poll();
+
+    // Trigger Add event
+    std::string addCmd =
+        "busctl call xyz.openbmc_project.Mctp /au/com/codeconstruct/mctp1 com.example.Control TriggerAdd ss \"" +
+        endpointPath + "\" \"" + emConfigPath + "\"";
+    runCmd(addCmd);
+
+    // Wait for query to start
+    while (!queryStarted)
+    {
+        io.poll();
+        usleep(10000); // 10ms
+    }
+
+    // IMMEDIATELY trigger Remove event (while query is in progress!)
+    std::string removeCmd =
+        "busctl call xyz.openbmc_project.Mctp /au/com/codeconstruct/mctp1 com.example.Control TriggerRemove s \"" +
+        endpointPath + "\"";
+    runCmd(removeCmd);
+
+    // Wait for query to complete with a timeout (up to 5 seconds)
+    for (int i = 0; i < 50; i++)
+    {
+        io.poll();
+        usleep(100000); // NOLINT(cert-api50-c)
+        if (queryCompletedCount >= 1)
+        {
+            break;
+        }
+    }
+
+    // Wait a bit more to see if any extra signals arrive
+    usleep(500000); // 500ms
+    io.poll();
+
+    EXPECT_EQ(queryCompletedCount, 1);
+
+    // Verify map is empty
+    auto it = mctpEndpointConfigMap.find(endpointPath);
+    EXPECT_TRUE(it == mctpEndpointConfigMap.end());
+}
+
+// 3. Add -> Remove -> Add during query
+TEST_F(MctpUtilTest, AddRemoveAdd_PopulatesMap)
+{
+    std::string endpointPath =
+        "/au/com/codeconstruct/mctp1/networks/1/endpoints/test3";
+    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&) { 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&) { queryCompletedCount++; });
+
+    // Wait for priming loop to complete (finding nothing)
+    io.poll();
+    usleep(100000); // 100ms
+    io.poll();
+
+    // Trigger Add event
+    std::string addCmd =
+        "busctl call xyz.openbmc_project.Mctp /au/com/codeconstruct/mctp1 com.example.Control TriggerAdd ss \"" +
+        endpointPath + "\" \"" + emConfigPath + "\"";
+    runCmd(addCmd);
+
+    // Wait for query to start
+    while (!queryStarted)
+    {
+        io.poll();
+        usleep(10000); // 10ms
+    }
+
+    // IMMEDIATELY trigger Remove event
+    std::string removeCmd =
+        "busctl call xyz.openbmc_project.Mctp /au/com/codeconstruct/mctp1 com.example.Control TriggerRemove s \"" +
+        endpointPath + "\"";
+    runCmd(removeCmd);
+
+    // IMMEDIATELY trigger Add event again!
+    runCmd(addCmd);
+
+    // Wait for queries to complete with a timeout (up to 5 seconds)
+    for (int i = 0; i < 50; i++)
+    {
+        io.poll();
+        usleep(100000); // NOLINT(cert-api50-c)
+        if (queryCompletedCount >= 2)
+        {
+            break;
+        }
+    }
+
+    // Wait a bit more to see if any extra signals arrive
+    usleep(500000); // 500ms
+    io.poll();
+
+    EXPECT_EQ(queryCompletedCount, 2);
+
+    // Verify map is populated at the end
+    auto it = mctpEndpointConfigMap.find(endpointPath);
+    ASSERT_NE(it, mctpEndpointConfigMap.end());
+    EXPECT_FALSE(it->second.empty());
+}
+
+int main(int argc, char** argv)
+{
+    ::testing::InitGoogleTest(&argc, argv);
+    return RUN_ALL_TESTS();
+}