feat: add fabric-bridge node for RDMA/inter-host MXL flow sharing
- New fabric-bridge node: initiator mode reads local MXL flow and transfers via Fabrics (RDMA/TCP); target mode receives remote data and commits to local MXL flow writer - Add set_mxl_instance/get_mxl_instance to Node base class - CMake option DMF_BUILD_FABRICS=OFF by default (needs libfabric) - Target publishes target_info in status for out-of-band exchange - Initiator receives target_info via config parameter
This commit is contained in:
@@ -14,6 +14,7 @@ list(APPEND CMAKE_MODULE_PATH "${CMAKE_CURRENT_SOURCE_DIR}/cmake")
|
||||
option(DMF_BUILD_TESTS "Build tests" ON)
|
||||
option(DMF_BUILD_MXL_TOOLS "Build MXL tools (testsrc, sink)" OFF)
|
||||
option(DMF_BUILD_DECKLINK "Build DeckLink I/O nodes" ON)
|
||||
option(DMF_BUILD_FABRICS "Build Fabrics bridge nodes" OFF)
|
||||
|
||||
set(DECKLINK_SDK_DIR "" CACHE PATH "Path to Blackmagic DeckLink SDK root")
|
||||
|
||||
@@ -22,6 +23,11 @@ if(DMF_BUILD_DECKLINK AND DECKLINK_SDK_DIR)
|
||||
add_subdirectory(nodes/decklink-out)
|
||||
endif()
|
||||
|
||||
if(DMF_BUILD_FABRICS)
|
||||
set(MXL_ENABLE_FABRICS_OFI ON CACHE BOOL "" FORCE)
|
||||
add_subdirectory(nodes/fabric-bridge)
|
||||
endif()
|
||||
|
||||
find_package(fmt CONFIG REQUIRED)
|
||||
find_package(spdlog CONFIG REQUIRED)
|
||||
find_package(nlohmann_json CONFIG REQUIRED)
|
||||
|
||||
@@ -27,6 +27,12 @@ public:
|
||||
virtual void on_remove_reader(const std::string& port_id) = 0;
|
||||
virtual void process() = 0;
|
||||
virtual nlohmann::json status() const = 0;
|
||||
|
||||
void set_mxl_instance(mxlInstance instance) { mxl_instance_ = instance; }
|
||||
mxlInstance get_mxl_instance() const { return mxl_instance_; }
|
||||
|
||||
private:
|
||||
mxlInstance mxl_instance_ = nullptr;
|
||||
};
|
||||
|
||||
} // namespace dmf_node
|
||||
|
||||
@@ -71,6 +71,8 @@ int NodeRunner::exec(std::unique_ptr<Node> node) {
|
||||
}
|
||||
spdlog::info("MXL instance created on domain: {}", mxl_domain_);
|
||||
|
||||
node->set_mxl_instance(mxl_instance_);
|
||||
|
||||
mxlGarbageCollectFlows(mxl_instance_);
|
||||
|
||||
if (!config_str_.empty()) {
|
||||
|
||||
+1
-1
Submodule mxl updated: 580abf71f4...ff0ece65e1
@@ -0,0 +1,28 @@
|
||||
cmake_minimum_required(VERSION 3.24 FATAL_ERROR)
|
||||
|
||||
project(dmf-node-fabric-bridge VERSION 0.1.0 LANGUAGES CXX)
|
||||
|
||||
find_package(fmt CONFIG REQUIRED)
|
||||
find_package(spdlog CONFIG REQUIRED)
|
||||
find_package(nlohmann_json CONFIG REQUIRED)
|
||||
find_package(PkgConfig REQUIRED)
|
||||
pkg_check_modules(libfabric REQUIRED IMPORTED_TARGET libfabric)
|
||||
|
||||
add_executable(dmf-node-fabric-bridge
|
||||
src/fabric_bridge_node.cpp
|
||||
src/main.cpp
|
||||
)
|
||||
|
||||
target_link_libraries(dmf-node-fabric-bridge PRIVATE
|
||||
dmf-node
|
||||
mxl-fabrics-objects
|
||||
mxl-fabrics-headers
|
||||
PkgConfig::libfabric
|
||||
spdlog::spdlog
|
||||
fmt::fmt
|
||||
nlohmann_json::nlohmann_json
|
||||
)
|
||||
|
||||
target_include_directories(dmf-node-fabric-bridge PRIVATE
|
||||
${CMAKE_CURRENT_SOURCE_DIR}/src
|
||||
)
|
||||
@@ -0,0 +1,320 @@
|
||||
#include "fabric_bridge_node.hpp"
|
||||
|
||||
#include <mxl/mxl.h>
|
||||
|
||||
#include <spdlog/spdlog.h>
|
||||
|
||||
namespace dmf_node {
|
||||
|
||||
FabricBridgeNode::~FabricBridgeNode() {
|
||||
running_ = false;
|
||||
if (worker_thread_.has_value() && worker_thread_->joinable()) {
|
||||
worker_thread_->join();
|
||||
}
|
||||
}
|
||||
|
||||
std::vector<PortDef> FabricBridgeNode::input_ports() const {
|
||||
if (mode_ == FabricMode::Target) {
|
||||
return {};
|
||||
}
|
||||
return {{"video_in", PortDirection::Input, MediaType::VideoV210}};
|
||||
}
|
||||
|
||||
std::vector<PortDef> FabricBridgeNode::output_ports() const {
|
||||
if (mode_ == FabricMode::Initiator) {
|
||||
return {};
|
||||
}
|
||||
return {{"video_out", PortDirection::Output, MediaType::VideoV210}};
|
||||
}
|
||||
|
||||
void FabricBridgeNode::configure(const nlohmann::json& params) {
|
||||
if (params.contains("mode")) {
|
||||
auto mode_str = params["mode"].get<std::string>();
|
||||
if (mode_str == "initiator") mode_ = FabricMode::Initiator;
|
||||
else if (mode_str == "target") mode_ = FabricMode::Target;
|
||||
else spdlog::warn("Fabric-bridge: unknown mode '{}', defaulting to initiator", mode_str);
|
||||
}
|
||||
|
||||
if (params.contains("bind_node")) bind_node_ = params["bind_node"].get<std::string>();
|
||||
if (params.contains("bind_service")) bind_service_ = params["bind_service"].get<std::string>();
|
||||
if (params.contains("provider")) {
|
||||
auto prov_str = params["provider"].get<std::string>();
|
||||
mxlFabricsProviderFromString(prov_str.c_str(), &provider_);
|
||||
}
|
||||
|
||||
if (params.contains("target_info")) {
|
||||
target_info_str_ = params["target_info"].get<std::string>();
|
||||
}
|
||||
}
|
||||
|
||||
void FabricBridgeNode::on_add_writer(const std::string& port_id, mxlFlowWriter writer) {
|
||||
if (port_id == "video_out") {
|
||||
writer_ = writer;
|
||||
spdlog::info("Fabric-bridge: writer added (target mode)");
|
||||
}
|
||||
}
|
||||
|
||||
void FabricBridgeNode::on_add_reader(const std::string& port_id, mxlFlowReader reader) {
|
||||
if (port_id == "video_in") {
|
||||
reader_ = reader;
|
||||
spdlog::info("Fabric-bridge: reader added (initiator mode)");
|
||||
}
|
||||
}
|
||||
|
||||
void FabricBridgeNode::on_remove_writer(const std::string& port_id) {
|
||||
if (port_id == "video_out") {
|
||||
writer_.reset();
|
||||
spdlog::info("Fabric-bridge: writer removed");
|
||||
}
|
||||
}
|
||||
|
||||
void FabricBridgeNode::on_remove_reader(const std::string& port_id) {
|
||||
if (port_id == "video_in") {
|
||||
reader_.reset();
|
||||
spdlog::info("Fabric-bridge: reader removed");
|
||||
}
|
||||
}
|
||||
|
||||
void FabricBridgeNode::process() {
|
||||
if (!worker_thread_.has_value()) {
|
||||
if (mode_ == FabricMode::Initiator && reader_.has_value() && !target_info_str_.empty()) {
|
||||
running_ = true;
|
||||
worker_thread_ = std::thread(&FabricBridgeNode::run_initiator, this);
|
||||
} else if (mode_ == FabricMode::Target && writer_.has_value()) {
|
||||
running_ = true;
|
||||
worker_thread_ = std::thread(&FabricBridgeNode::run_target, this);
|
||||
}
|
||||
}
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
||||
}
|
||||
|
||||
void FabricBridgeNode::run_initiator() {
|
||||
spdlog::info("Fabric-bridge: starting initiator mode");
|
||||
|
||||
auto instance = get_mxl_instance();
|
||||
if (!instance) {
|
||||
spdlog::error("Fabric-bridge: no MXL instance");
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
mxlFabricsInstance fab_inst = nullptr;
|
||||
auto status = mxlFabricsCreateInstance(instance, nullptr, &fab_inst);
|
||||
if (status != MXL_STATUS_OK || !fab_inst) {
|
||||
spdlog::error("Fabric-bridge: failed to create fabrics instance: {}", static_cast<int>(status));
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
fabrics_instance_ = fab_inst;
|
||||
|
||||
mxlFabricsInitiator initiator = nullptr;
|
||||
status = mxlFabricsCreateInitiator(fab_inst, &initiator);
|
||||
if (status != MXL_STATUS_OK || !initiator) {
|
||||
spdlog::error("Fabric-bridge: failed to create initiator: {}", static_cast<int>(status));
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
initiator_ = initiator;
|
||||
|
||||
mxlFlowConfigInfo config{};
|
||||
mxlFlowReaderGetConfigInfo(*reader_, &config);
|
||||
grain_rate_ = config.common.grainRate;
|
||||
|
||||
mxlFabricsInitiatorConfig init_cfg{};
|
||||
init_cfg.version = MXL_FABRICS_API_VERSION;
|
||||
init_cfg.interface.version = MXL_FABRICS_API_VERSION;
|
||||
init_cfg.interface.provider = provider_;
|
||||
init_cfg.interface.address.node = bind_node_.c_str();
|
||||
init_cfg.interface.address.service = bind_service_.c_str();
|
||||
init_cfg.reader = *reader_;
|
||||
|
||||
status = mxlFabricsInitiatorSetup(initiator, &init_cfg, nullptr);
|
||||
if (status != MXL_STATUS_OK) {
|
||||
spdlog::error("Fabric-bridge: initiator setup failed: {}", static_cast<int>(status));
|
||||
mxlFabricsDestroyInitiator(fab_inst, initiator);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
mxlFabricsTargetInfo target_info = nullptr;
|
||||
status = mxlFabricsTargetInfoFromString(target_info_str_.c_str(), &target_info);
|
||||
if (status != MXL_STATUS_OK || !target_info) {
|
||||
spdlog::error("Fabric-bridge: failed to parse target_info: {}", static_cast<int>(status));
|
||||
mxlFabricsDestroyInitiator(fab_inst, initiator);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
status = mxlFabricsInitiatorAddTarget(initiator, target_info);
|
||||
if (status != MXL_STATUS_OK) {
|
||||
spdlog::error("Fabric-bridge: failed to add target: {}", static_cast<int>(status));
|
||||
mxlFabricsFreeTargetInfo(target_info);
|
||||
mxlFabricsDestroyInitiator(fab_inst, initiator);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
spdlog::info("Fabric-bridge: waiting for connection...");
|
||||
do {
|
||||
status = mxlFabricsInitiatorMakeProgressBlocking(initiator, 250);
|
||||
} while (status == MXL_ERR_NOT_READY && running_);
|
||||
spdlog::info("Fabric-bridge: connected");
|
||||
|
||||
auto now = mxlGetTime();
|
||||
read_index_ = mxlTimestampToIndex(&grain_rate_, now) - 2;
|
||||
|
||||
spdlog::info("Fabric-bridge: initiator started, grain_rate={}/{}", grain_rate_.numerator, grain_rate_.denominator);
|
||||
|
||||
while (running_) {
|
||||
auto deadline = mxlIndexToTimestamp(&grain_rate_, read_index_ + 1);
|
||||
mxlSleepUntil(deadline);
|
||||
|
||||
mxlGrainInfo grain_info{};
|
||||
uint8_t* payload = nullptr;
|
||||
auto grain_status = mxlFlowReaderGetGrain(*reader_, read_index_, 5000000ULL, &grain_info, &payload);
|
||||
if (grain_status != MXL_STATUS_OK) {
|
||||
if (grain_status == MXL_ERR_OUT_OF_RANGE_TOO_LATE || grain_status == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) {
|
||||
auto ts = mxlGetTime();
|
||||
read_index_ = mxlTimestampToIndex(&grain_rate_, ts) - 2;
|
||||
continue;
|
||||
}
|
||||
read_index_++;
|
||||
continue;
|
||||
}
|
||||
|
||||
auto transfer_status = mxlFabricsInitiatorTransferGrain(initiator, read_index_, 0, grain_info.validSlices);
|
||||
if (transfer_status != MXL_STATUS_OK) {
|
||||
spdlog::warn("Fabric-bridge: transfer grain {} failed: {}", read_index_, static_cast<int>(transfer_status));
|
||||
read_index_++;
|
||||
continue;
|
||||
}
|
||||
|
||||
do {
|
||||
status = mxlFabricsInitiatorMakeProgressBlocking(initiator, 10);
|
||||
} while (status == MXL_ERR_NOT_READY && running_);
|
||||
|
||||
read_index_++;
|
||||
grains_transferred_++;
|
||||
|
||||
if (grains_transferred_ == 1) {
|
||||
spdlog::info("Fabric-bridge: first grain transferred, index={}", read_index_ - 1);
|
||||
}
|
||||
}
|
||||
|
||||
mxlFabricsInitiatorRemoveTarget(initiator, target_info);
|
||||
mxlFabricsFreeTargetInfo(target_info);
|
||||
mxlFabricsDestroyInitiator(fab_inst, initiator);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
spdlog::info("Fabric-bridge: initiator stopped");
|
||||
}
|
||||
|
||||
void FabricBridgeNode::run_target() {
|
||||
spdlog::info("Fabric-bridge: starting target mode");
|
||||
|
||||
auto instance = get_mxl_instance();
|
||||
if (!instance) {
|
||||
spdlog::error("Fabric-bridge: no MXL instance");
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
mxlFabricsInstance fab_inst = nullptr;
|
||||
auto status = mxlFabricsCreateInstance(instance, nullptr, &fab_inst);
|
||||
if (status != MXL_STATUS_OK || !fab_inst) {
|
||||
spdlog::error("Fabric-bridge: failed to create fabrics instance: {}", static_cast<int>(status));
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
fabrics_instance_ = fab_inst;
|
||||
|
||||
mxlFabricsTarget tgt = nullptr;
|
||||
status = mxlFabricsCreateTarget(fab_inst, &tgt);
|
||||
if (status != MXL_STATUS_OK || !tgt) {
|
||||
spdlog::error("Fabric-bridge: failed to create target: {}", static_cast<int>(status));
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
target_ = tgt;
|
||||
|
||||
mxlFlowConfigInfo config{};
|
||||
mxlFlowWriterGetConfigInfo(*writer_, &config);
|
||||
grain_rate_ = config.common.grainRate;
|
||||
|
||||
mxlFabricsTargetConfig target_cfg{};
|
||||
target_cfg.version = MXL_FABRICS_API_VERSION;
|
||||
target_cfg.interface.version = MXL_FABRICS_API_VERSION;
|
||||
target_cfg.interface.provider = provider_;
|
||||
target_cfg.interface.address.node = bind_node_.c_str();
|
||||
target_cfg.interface.address.service = bind_service_.c_str();
|
||||
target_cfg.writer = *writer_;
|
||||
|
||||
mxlFabricsTargetInfo target_info = nullptr;
|
||||
status = mxlFabricsTargetSetup(tgt, &target_cfg, nullptr, &target_info);
|
||||
if (status != MXL_STATUS_OK || !target_info) {
|
||||
spdlog::error("Fabric-bridge: target setup failed: {}", static_cast<int>(status));
|
||||
mxlFabricsDestroyTarget(fab_inst, tgt);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
running_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
size_t info_size = 0;
|
||||
mxlFabricsTargetInfoToString(target_info, nullptr, &info_size);
|
||||
std::vector<char> info_buf(info_size);
|
||||
mxlFabricsTargetInfoToString(target_info, info_buf.data(), &info_size);
|
||||
target_info_str_ = std::string(info_buf.data(), info_size);
|
||||
|
||||
spdlog::info("Fabric-bridge: target ready, target_info length={}", target_info_str_.size());
|
||||
spdlog::info("Fabric-bridge: TARGET_INFO={}", target_info_str_);
|
||||
|
||||
while (running_) {
|
||||
uint64_t grain_index = 0;
|
||||
status = mxlFabricsTargetReadGrain(tgt, 200, &grain_index);
|
||||
if (status == MXL_ERR_TIMEOUT || status == MXL_ERR_NOT_READY) {
|
||||
continue;
|
||||
}
|
||||
if (status != MXL_STATUS_OK) {
|
||||
spdlog::warn("Fabric-bridge: target read grain failed: {}", static_cast<int>(status));
|
||||
continue;
|
||||
}
|
||||
|
||||
mxlGrainInfo grain_info{};
|
||||
uint8_t* dummy_payload = nullptr;
|
||||
auto open_status = mxlFlowWriterOpenGrain(*writer_, grain_index, &grain_info, &dummy_payload);
|
||||
if (open_status == MXL_STATUS_OK) {
|
||||
grain_info.validSlices = grain_info.totalSlices;
|
||||
mxlFlowWriterCommitGrain(*writer_, &grain_info);
|
||||
grains_transferred_++;
|
||||
if (grains_transferred_ == 1) {
|
||||
spdlog::info("Fabric-bridge: first grain committed, index={}", grain_index);
|
||||
}
|
||||
} else {
|
||||
spdlog::warn("Fabric-bridge: open grain {} failed: {}", grain_index, static_cast<int>(open_status));
|
||||
}
|
||||
}
|
||||
|
||||
mxlFabricsFreeTargetInfo(target_info);
|
||||
mxlFabricsDestroyTarget(fab_inst, tgt);
|
||||
mxlFabricsDestroyInstance(fab_inst);
|
||||
spdlog::info("Fabric-bridge: target stopped");
|
||||
}
|
||||
|
||||
nlohmann::json FabricBridgeNode::status() const {
|
||||
return {
|
||||
{"type", "fabric-bridge"},
|
||||
{"mode", mode_ == FabricMode::Initiator ? "initiator" : "target"},
|
||||
{"grains_transferred", grains_transferred_},
|
||||
{"running", running_.load()},
|
||||
{"has_reader", reader_.has_value()},
|
||||
{"has_writer", writer_.has_value()},
|
||||
{"target_info_length", target_info_str_.size()},
|
||||
{"target_info", target_info_str_},
|
||||
};
|
||||
}
|
||||
|
||||
} // namespace dmf_node
|
||||
@@ -0,0 +1,62 @@
|
||||
#pragma once
|
||||
|
||||
#include <dmf-node/node.hpp>
|
||||
#include <mxl/flow.h>
|
||||
#include <mxl/fabrics.h>
|
||||
#include <mxl/time.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <optional>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
|
||||
namespace dmf_node {
|
||||
|
||||
enum class FabricMode { Initiator, Target };
|
||||
|
||||
class FabricBridgeNode : public Node {
|
||||
public:
|
||||
FabricBridgeNode() = default;
|
||||
~FabricBridgeNode();
|
||||
|
||||
std::string type() const override { return "fabric-bridge"; }
|
||||
|
||||
std::vector<PortDef> input_ports() const override;
|
||||
std::vector<PortDef> output_ports() const override;
|
||||
|
||||
void configure(const nlohmann::json& params) override;
|
||||
void on_add_writer(const std::string& port_id, mxlFlowWriter writer) override;
|
||||
void on_add_reader(const std::string& port_id, mxlFlowReader reader) override;
|
||||
void on_remove_writer(const std::string& port_id) override;
|
||||
void on_remove_reader(const std::string& port_id) override;
|
||||
|
||||
void process() override;
|
||||
nlohmann::json status() const override;
|
||||
|
||||
private:
|
||||
void run_initiator();
|
||||
void run_target();
|
||||
|
||||
FabricMode mode_ = FabricMode::Initiator;
|
||||
|
||||
std::string bind_node_ = "0.0.0.0";
|
||||
std::string bind_service_ = "0";
|
||||
mxlFabricsProvider provider_ = MXL_FABRICS_PROVIDER_ANY;
|
||||
|
||||
std::optional<mxlFlowReader> reader_;
|
||||
std::optional<mxlFlowWriter> writer_;
|
||||
mxlRational grain_rate_{50, 1};
|
||||
uint64_t read_index_ = 0;
|
||||
uint64_t grains_transferred_ = 0;
|
||||
|
||||
std::string target_info_str_;
|
||||
|
||||
mxlFabricsInstance fabrics_instance_ = nullptr;
|
||||
mxlFabricsInitiator initiator_ = nullptr;
|
||||
mxlFabricsTarget target_ = nullptr;
|
||||
|
||||
std::optional<std::thread> worker_thread_;
|
||||
std::atomic<bool> running_{false};
|
||||
};
|
||||
|
||||
} // namespace dmf_node
|
||||
@@ -0,0 +1,11 @@
|
||||
#include "fabric_bridge_node.hpp"
|
||||
#include <dmf-node/node_runner.hpp>
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
dmf_node::NodeRunner runner;
|
||||
if (!runner.parse_args(argc, argv)) {
|
||||
return 1;
|
||||
}
|
||||
auto node = std::make_unique<dmf_node::FabricBridgeNode>();
|
||||
return runner.exec(std::move(node));
|
||||
}
|
||||
Reference in New Issue
Block a user