blob: c1a3f4885584deb5bc33676136602f4524b328fd [file]
#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();
}