| #include "MctpUtil.hpp" |
| |
| #include "absl/strings/match.h" |
| |
| #include <boost/asio/spawn.hpp> |
| #include <boost/asio/steady_timer.hpp> |
| #include <boost/container/flat_map.hpp> |
| #include <phosphor-logging/lg2.hpp> |
| #include <sdbusplus/bus/match.hpp> |
| |
| #include <chrono> |
| #include <cstdint> |
| #include <filesystem> |
| #include <ranges> |
| |
| std::map<std::string, std::vector<SensorData>> mctpEndpointConfigMap; |
| |
| enum class State : std::uint8_t |
| { |
| None, |
| Add, |
| Remove |
| }; |
| |
| struct EndpointState |
| { |
| 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; |
| static uint8_t filterMsgType = 0; |
| |
| static std::map<MctpCallbackToken, MctpEndpointCallback> mctpEndpointCallbacks; |
| static MctpCallbackToken nextCallbackToken = 1; |
| |
| static std::unique_ptr<sdbusplus::bus::match_t> associationMatch = nullptr; |
| static std::unique_ptr<sdbusplus::bus::match_t> associationRemoveMatch = |
| nullptr; |
| |
| static void performEmConfigQuery( |
| const std::shared_ptr<sdbusplus::asio::connection>& conn, |
| const std::string& endpointPath, const std::string& emConfigPath, |
| bool bypassCooldown, bool bypassRetry, |
| const boost::asio::yield_context& yield) |
| { |
| auto& state = endpointStates[endpointPath]; |
| static constexpr auto cooldown = std::chrono::seconds(5); |
| |
| state.querying = true; // Protect the window! |
| |
| // Use shared_ptr constructor directly with a null pointer to act as a scope |
| // guard |
| std::shared_ptr<void> cleanupTrap(nullptr, [endpointPath](void*) { |
| auto it = endpointStates.find(endpointPath); |
| if (it != endpointStates.end()) |
| { |
| it->second.querying = false; |
| } |
| }); |
| |
| while (true) |
| { |
| auto now = std::chrono::steady_clock::now(); |
| |
| // 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); |
| |
| boost::asio::steady_timer timer(conn->get_io_context()); |
| timer.expires_after(delay); |
| |
| boost::system::error_code ec; |
| timer.async_wait(yield[ec]); |
| if (ec) |
| { |
| return; |
| } |
| if (state.state == State::Remove) |
| { |
| break; |
| } |
| } |
| bypassCooldown = false; |
| |
| state.lastQueryTime = |
| std::chrono::steady_clock::now(); // Update time for EVERY attempt! |
| bool success = false; |
| // Use a single-trip do-while loop to flatten error handling |
| // and allow early exits straight to the deferred state evaluation |
| do |
| { |
| // Step 1: Query SupportedMessageTypes directly from mctpd |
| boost::system::error_code ec; |
| std::variant<std::vector<uint8_t>> value; |
| |
| try |
| { |
| value = |
| conn->yield_method_call<std::variant<std::vector<uint8_t>>>( |
| yield, ec, "au.com.codeconstruct.MCTP1", endpointPath, |
| "org.freedesktop.DBus.Properties", "Get", |
| "xyz.openbmc_project.MCTP.Endpoint", |
| "SupportedMessageTypes"); |
| if (state.state == State::Remove) |
| { |
| break; |
| } |
| } |
| catch (const std::exception& e) |
| { |
| lg2::error( |
| "Exception getting SupportedMessageTypes for {PATH}: {ERROR}", |
| "PATH", endpointPath, "ERROR", e.what()); |
| break; |
| } |
| |
| if (ec) |
| { |
| lg2::error( |
| "Failed to get SupportedMessageTypes for {PATH}: {ERROR}", |
| "PATH", endpointPath, "ERROR", ec.message()); |
| break; |
| } |
| |
| const auto* typesPtr = std::get_if<std::vector<uint8_t>>(&value); |
| if (typesPtr == nullptr) |
| { |
| lg2::error("Invalid SupportedMessageTypes type for {PATH}", |
| "PATH", endpointPath); |
| break; |
| } |
| if (std::find(typesPtr->begin(), typesPtr->end(), filterMsgType) == |
| typesPtr->end()) |
| { |
| lg2::info( |
| "Endpoint {PATH} does not support message type {TYPE}, ignoring", |
| "PATH", endpointPath, "TYPE", filterMsgType); |
| return; |
| } |
| |
| // Step 2: Proceed to EM query |
| ManagedObjectType objects; |
| try |
| { |
| objects = conn->yield_method_call<ManagedObjectType>( |
| yield, ec, "xyz.openbmc_project.EntityManager", |
| "/xyz/openbmc_project/inventory", |
| "org.freedesktop.DBus.ObjectManager", "GetManagedObjects"); |
| if (state.state == State::Remove) |
| { |
| break; |
| } |
| } |
| catch (const std::exception& e) |
| { |
| lg2::error("Exception getting managed objects: {ERROR}", |
| "ERROR", e.what()); |
| break; |
| } |
| |
| if (state.state != State::Remove) |
| { |
| if (ec) |
| { |
| lg2::error( |
| "Failed to get managed objects for {PATH}: {ERROR}", |
| "PATH", emConfigPath, "ERROR", ec.message()); |
| break; |
| } |
| |
| auto objIt = |
| objects.find(sdbusplus::message::object_path(emConfigPath)); |
| |
| if (objIt == objects.end()) |
| { |
| lg2::warning( |
| "EM config path {PATH} not found in managed objects", |
| "PATH", emConfigPath); |
| } |
| else if (!objIt->second.empty()) |
| { |
| lg2::info( |
| "Recorded MCTP endpoint {ENDPOINT} with config from {PATH}", |
| "ENDPOINT", endpointPath, "PATH", emConfigPath); |
| |
| std::vector<SensorData> configChain; |
| configChain.push_back(objIt->second); |
| |
| // Resolve bridge chain if needed |
| std::string currentPath = emConfigPath; |
| bool resolutionFailed = false; |
| std::vector<std::string> visitedPaths = {currentPath}; |
| |
| while (true) |
| { |
| const auto& currentConfig = configChain.back(); |
| auto bridgeIt = currentConfig.find( |
| "xyz.openbmc_project.Configuration.MCTPBridgeDownstreamDevice"); |
| if (bridgeIt == currentConfig.end()) |
| { |
| break; // Not a bridged endpoint, or reached the end |
| // of chain |
| } |
| |
| auto propIt = bridgeIt->second.find("BridgeInventory"); |
| if (propIt == bridgeIt->second.end()) |
| { |
| lg2::error( |
| "BridgeInventory property missing for {PATH}", |
| "PATH", currentPath); |
| resolutionFailed = true; |
| break; |
| } |
| |
| const auto* bridgePathPtr = |
| std::get_if<std::string>(&propIt->second); |
| if (bridgePathPtr == nullptr) |
| { |
| lg2::error( |
| "Invalid BridgeInventory property type for {PATH}", |
| "PATH", currentPath); |
| resolutionFailed = true; |
| break; |
| } |
| |
| std::string nextPath = *bridgePathPtr; |
| |
| if (std::find(visitedPaths.begin(), visitedPaths.end(), |
| nextPath) != visitedPaths.end()) |
| { |
| lg2::error( |
| "Circular dependency detected in bridge chain at {PATH}", |
| "PATH", nextPath); |
| resolutionFailed = true; |
| break; |
| } |
| visitedPaths.push_back(nextPath); |
| |
| auto nextObjIt = objects.find( |
| sdbusplus::message::object_path(nextPath)); |
| |
| if (nextObjIt == objects.end()) |
| { |
| lg2::error("Bridge config {PATH} not found", "PATH", |
| nextPath); |
| resolutionFailed = true; |
| break; |
| } |
| |
| // Verify it contains a supported bridge interface |
| if (nextObjIt->second.find( |
| "xyz.openbmc_project.Configuration.MCTPBridgeDownstreamDevice") == |
| nextObjIt->second.end() && |
| nextObjIt->second.find( |
| "xyz.openbmc_project.Configuration.MCTPUSBDevice") == |
| nextObjIt->second.end()) |
| { |
| lg2::error("Object {PATH} is not a valid bridge", |
| "PATH", nextPath); |
| resolutionFailed = true; |
| break; |
| } |
| |
| configChain.push_back(nextObjIt->second); |
| currentPath = nextPath; |
| } |
| |
| if (resolutionFailed) |
| { |
| break; // Abort resolution for this endpoint (exits the |
| // do-while(false) loop) |
| } |
| |
| mctpEndpointConfigMap[endpointPath] = configChain; |
| |
| for (const auto& [token, cb] : mctpEndpointCallbacks) |
| { |
| cb(endpointPath, configChain, false); |
| } |
| success = true; |
| } |
| } |
| } while (false); |
| |
| if (!success) |
| { |
| if (bypassRetry) |
| { |
| lg2::info( |
| "Query failed for endpoint {ENDPOINT} during priming, deferring.", |
| "ENDPOINT", endpointPath); |
| state.state = State::Add; // Mark for retry later |
| break; // Exit loop |
| } |
| lg2::info("Query failed for endpoint {ENDPOINT}, retrying...", |
| "ENDPOINT", endpointPath); |
| continue; |
| } |
| |
| State lastState = state.state; |
| state.state = State::None; // Reset for next query |
| |
| if (lastState == State::Remove) |
| { |
| lg2::info("Ignoring query reply for removed endpoint {ENDPOINT}", |
| "ENDPOINT", endpointPath); |
| auto it = mctpEndpointConfigMap.find(endpointPath); |
| if (it != mctpEndpointConfigMap.end()) |
| { |
| mctpEndpointConfigMap.erase(it); |
| for (const auto& [token, cb] : mctpEndpointCallbacks) |
| { |
| cb(endpointPath, {}, true); |
| } |
| } |
| break; |
| } |
| if (lastState == State::Add) |
| { |
| lg2::info( |
| "DEBUG: Triggering deferred query for endpoint {ENDPOINT}", |
| "ENDPOINT", endpointPath); |
| continue; |
| } |
| |
| break; |
| } |
| } |
| |
| static void triggerDeferredQueries( |
| const std::shared_ptr<sdbusplus::asio::connection>& conn, bool isPriming) |
| { |
| for (auto& [path, state] : endpointStates) |
| { |
| if (path == "__global__") |
| { |
| continue; |
| } |
| |
| if (state.state == State::Add) |
| { |
| std::string emConfigPath = state.emConfigPath; |
| std::string endpointPath = path; |
| |
| boost::asio::spawn(conn->get_io_context(), |
| [conn, endpointPath, emConfigPath, isPriming]( |
| const boost::asio::yield_context& yield) { |
| auto& state = endpointStates[endpointPath]; |
| if (state.querying) |
| { |
| return; // Drop if already querying |
| } |
| if (state.state != State::Add) |
| { |
| return; // Drop if state changed across boundary |
| } |
| state.state = State::None; |
| performEmConfigQuery(conn, endpointPath, emConfigPath, |
| isPriming, false, yield); |
| }); |
| } |
| else if (state.state == State::Remove) |
| { |
| state.state = State::None; |
| auto it = mctpEndpointConfigMap.find(path); |
| if (it != mctpEndpointConfigMap.end()) |
| { |
| mctpEndpointConfigMap.erase(it); |
| for (const auto& [token, cb] : mctpEndpointCallbacks) |
| { |
| cb(path, {}, true); |
| } |
| } |
| } |
| } |
| } |
| |
| void setupMctpEndpointListener( |
| const std::shared_ptr<sdbusplus::asio::connection>& conn, |
| MctpMessageType msgType) |
| { |
| filterMsgType = static_cast<uint8_t>(msgType); |
| // Match 1: Listen for Associations from mctp-reactor |
| const std::string associationMatchSpec = |
| "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, |
| [conn](sdbusplus::message_t& msg) { |
| sdbusplus::message::object_path path; |
| boost::container::flat_map< |
| std::string, |
| boost::container::flat_map< |
| std::string, |
| std::variant<std::string, std::vector<std::string>>>> |
| interfaces; |
| msg.read(path, interfaces); |
| |
| auto it = interfaces.find("xyz.openbmc_project.Association"); |
| if (it == interfaces.end()) |
| { |
| return; |
| } |
| |
| auto propIt = it->second.find("endpoints"); |
| if (propIt == it->second.end()) |
| { |
| return; |
| } |
| |
| const auto* endpointsPtr = |
| std::get_if<std::vector<std::string>>(&propIt->second); |
| if (!endpointsPtr) |
| { |
| return; |
| } |
| |
| std::string emConfigPath = endpointsPtr->front(); |
| |
| std::string endpointPath = path.str; |
| constexpr std::string_view suffix = "/configured_by"; |
| if (endpointPath.ends_with("/configured_by")) |
| { |
| endpointPath = |
| endpointPath.substr(0, endpointPath.size() - suffix.length()); |
| } |
| |
| boost::asio::spawn(conn->get_io_context(), |
| [conn, endpointPath, emConfigPath]( |
| const boost::asio::yield_context& yield) { |
| auto& state = endpointStates[endpointPath]; |
| state.emConfigPath = emConfigPath; |
| bool wasQuerying = endpointStates["__global__"].querying || |
| state.querying; |
| |
| if (wasQuerying) |
| { |
| lg2::info("Deferring Add for endpoint {PATH}", "PATH", |
| endpointPath); |
| state.state = State::Add; |
| } |
| else |
| { |
| performEmConfigQuery(conn, endpointPath, emConfigPath, false, |
| false, yield); |
| } |
| }); |
| }); |
| |
| const std::string associationRemoveMatchSpec = |
| "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, |
| [](sdbusplus::message_t& msg) { |
| sdbusplus::message::object_path path; |
| std::vector<std::string> interfaces; |
| msg.read(path, interfaces); |
| |
| std::string endpointPath = path.str; |
| constexpr std::string_view suffix = "/configured_by"; |
| if (endpointPath.ends_with("/configured_by")) |
| { |
| endpointPath = |
| endpointPath.substr(0, endpointPath.size() - suffix.length()); |
| } |
| |
| auto& state = endpointStates[endpointPath]; |
| |
| if (endpointStates["__global__"].querying || state.querying) |
| { |
| state.state = State::Remove; |
| } |
| else |
| { |
| auto it = mctpEndpointConfigMap.find(endpointPath); |
| if (it != mctpEndpointConfigMap.end()) |
| { |
| lg2::info("Removing MCTP endpoint {ENDPOINT} from map", |
| "ENDPOINT", endpointPath); |
| mctpEndpointConfigMap.erase(it); |
| |
| for (const auto& [token, cb] : mctpEndpointCallbacks) |
| { |
| cb(endpointPath, {}, true); |
| } |
| } |
| } |
| }); |
| |
| // Prime the cache: Query existing Associations |
| endpointStates["__global__"].querying = true; |
| |
| boost::asio::spawn(conn->get_io_context(), |
| [conn](const boost::asio::yield_context& yield) { |
| #ifdef UNIT_TEST |
| // Inject 5 second delay ONLY for tests to hit the race condition |
| auto timer = |
| std::make_shared<boost::asio::steady_timer>(conn->get_io_context()); |
| timer->expires_after(std::chrono::seconds(5)); |
| boost::system::error_code ec_timer; |
| timer->async_wait(yield[ec_timer]); |
| if (ec_timer) |
| return; |
| #endif |
| constexpr auto associationInterfaces = |
| std::to_array({"xyz.openbmc_project.Association"}); |
| boost::system::error_code ec; |
| GetSubTreeType subtree; |
| |
| try |
| { |
| subtree = conn->yield_method_call<GetSubTreeType>( |
| yield, ec, "xyz.openbmc_project.ObjectMapper", |
| "/xyz/openbmc_project/object_mapper", |
| "xyz.openbmc_project.ObjectMapper", "GetSubTree", |
| "/au/com/codeconstruct/mctp1", 0, associationInterfaces); |
| } |
| catch (const std::exception& e) |
| { |
| lg2::error("Exception getting managed objects: {ERROR}", "ERROR", |
| e.what()); |
| endpointStates["__global__"].querying = false; |
| return; |
| } |
| |
| if (ec) |
| { |
| lg2::error("Failed to get associations from Object Mapper: {ERROR}", |
| "ERROR", ec.message()); |
| endpointStates["__global__"].querying = false; |
| return; |
| } |
| |
| if (subtree.empty()) |
| { |
| endpointStates["__global__"].querying = false; |
| return; |
| } |
| |
| for (const auto& [path, services] : subtree) |
| { |
| if (!path.ends_with("/configured_by")) |
| { |
| continue; |
| } |
| |
| std::string endpointPath = path; |
| constexpr std::string_view suffix = "/configured_by"; |
| endpointPath = |
| endpointPath.substr(0, endpointPath.size() - suffix.length()); |
| |
| std::variant<std::vector<std::string>> value; |
| try |
| { |
| value = conn->yield_method_call< |
| std::variant<std::vector<std::string>>>( |
| yield, ec, "xyz.openbmc_project.ObjectMapper", path, |
| "org.freedesktop.DBus.Properties", "Get", |
| "xyz.openbmc_project.Association", "endpoints"); |
| } |
| catch (const std::exception& e) |
| { |
| lg2::error("Failed to get endpoints for {PATH}: {ERROR}", |
| "PATH", endpointPath, "ERROR", e.what()); |
| continue; |
| } |
| |
| if (ec) |
| { |
| lg2::error("Failed to get endpoints for {PATH}: {ERROR}", |
| "PATH", endpointPath, "ERROR", ec.message()); |
| continue; |
| } |
| |
| const auto* endpoints = |
| std::get_if<std::vector<std::string>>(&value); |
| if (!endpoints || endpoints->empty()) |
| { |
| continue; |
| } |
| |
| std::string emConfigPath = endpoints->front(); |
| |
| auto& state = endpointStates[endpointPath]; |
| state.emConfigPath = emConfigPath; |
| performEmConfigQuery(conn, endpointPath, emConfigPath, true, true, |
| yield); |
| } |
| |
| endpointStates["__global__"].querying = false; |
| triggerDeferredQueries(conn, true); |
| }); |
| } |
| |
| MctpReactorBusInfo::MctpReactorBusInfo( |
| const std::vector<SensorData>& configDataChain) |
| { |
| auto extractProp = [](const auto& dbusProps, |
| std::map<std::string, std::string>& info, |
| const std::string& key) { |
| auto it = dbusProps.find(key); |
| if (it != dbusProps.end()) |
| { |
| info[key] = std::visit(VariantToStringVisitor(), it->second); |
| } |
| }; |
| |
| // For a MCTP device behind (multiple layers of) MCTP bridge, |
| // configDataChain include all config/busInfo from its bottom-most bridge |
| // routing info (a.k.a BridgeDownsteamMCTPDDevice) to the top-most port info |
| // (a.k.a MCTPUSBDevice/etc). |
| |
| // And MctpReactorBusInfo extracts the config for |
| // both routing and the top most port info and put them into a reverse order |
| // (from top most port to the bottom most routing) to accomodate the |
| // NVMe1000 bus config |
| for (const auto& configData : std::views::reverse(configDataChain)) |
| { |
| std::string type; |
| std::map<std::string, std::string> properties; |
| bool levelValid = false; |
| |
| for (const auto& [intf, props] : configData) |
| { |
| if (absl::StrContains( |
| intf, "xyz.openbmc_project.Configuration.MCTPUSBDevice")) |
| { |
| type = "USB"; |
| levelValid = true; |
| properties["BusType"] = "USB"; |
| |
| extractProp(props, properties, "RootHubPath"); |
| extractProp(props, properties, "Port"); |
| extractProp(props, properties, "InterfaceNum"); |
| extractProp(props, properties, "Configuration"); |
| break; |
| } |
| if (absl::StrContains( |
| intf, "xyz.openbmc_project.Configuration.MCTPI2CTarget")) |
| { |
| type = "I2C"; |
| levelValid = true; |
| properties["BusType"] = "I2C"; |
| |
| extractProp(props, properties, "Bus"); |
| extractProp(props, properties, "Address"); |
| break; |
| } |
| if (absl::StrContains( |
| intf, "xyz.openbmc_project.Configuration.MCTPI3CTarget")) |
| { |
| type = "I3C"; |
| levelValid = true; |
| properties["BusType"] = "I3C"; |
| |
| extractProp(props, properties, "Bus"); |
| extractProp(props, properties, "Address"); |
| break; |
| } |
| if (absl::StrContains( |
| intf, |
| "xyz.openbmc_project.Configuration.MCTPBridgeDownstreamDevice")) |
| { |
| type = "MCTP Bridge"; |
| levelValid = true; |
| properties["BusType"] = "MCTP Bridge"; |
| |
| extractProp(props, properties, "EIDOffset"); |
| break; |
| } |
| } |
| |
| if (levelValid) |
| { |
| types.push_back(type); |
| propertiesList.push_back(properties); |
| valid = true; |
| } |
| } |
| } |
| |
| bool MctpReactorBusInfo::operator==( |
| const std::vector<BusInfo>& targetBusInfo) const |
| { |
| if (propertiesList.size() != targetBusInfo.size()) |
| { |
| return false; |
| } |
| |
| for (size_t i = 0; i < targetBusInfo.size(); ++i) |
| { |
| const auto& configBusInfo = targetBusInfo[i]; |
| const auto& epProps = propertiesList[i]; |
| |
| for (const auto& [key, val] : configBusInfo) |
| { |
| auto it = epProps.find(key); |
| if (it == epProps.end() || it->second != val) |
| { |
| return false; |
| } |
| } |
| } |
| return true; |
| } |
| |
| void cleanupMctpEndpointListener() |
| { |
| associationMatch.reset(); |
| associationRemoveMatch.reset(); |
| endpointStates.clear(); // Prevent state leaking between unit tests |
| } |
| |
| MctpCallbackToken registerMctpEndpointCallback(MctpEndpointCallback&& cb) |
| { |
| MctpCallbackToken token = nextCallbackToken++; |
| mctpEndpointCallbacks[token] = std::move(cb); |
| return token; |
| } |
| |
| void unregisterMctpEndpointCallback(MctpCallbackToken token) |
| { |
| mctpEndpointCallbacks.erase(token); |
| } |
| |
| void unregisterAllMctpEndpointCallbacks() |
| { |
| mctpEndpointCallbacks.clear(); |
| } |