diff --git a/CMakeLists.txt b/CMakeLists.txt index 12da608..d074197 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -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) diff --git a/libs/dmf-node/include/dmf-node/node.hpp b/libs/dmf-node/include/dmf-node/node.hpp index 52a9871..b8a929f 100644 --- a/libs/dmf-node/include/dmf-node/node.hpp +++ b/libs/dmf-node/include/dmf-node/node.hpp @@ -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 diff --git a/libs/dmf-node/src/node_runner.cpp b/libs/dmf-node/src/node_runner.cpp index 51629f8..fd93b62 100644 --- a/libs/dmf-node/src/node_runner.cpp +++ b/libs/dmf-node/src/node_runner.cpp @@ -71,6 +71,8 @@ int NodeRunner::exec(std::unique_ptr node) { } spdlog::info("MXL instance created on domain: {}", mxl_domain_); + node->set_mxl_instance(mxl_instance_); + mxlGarbageCollectFlows(mxl_instance_); if (!config_str_.empty()) { diff --git a/mxl b/mxl index 580abf7..ff0ece6 160000 --- a/mxl +++ b/mxl @@ -1 +1 @@ -Subproject commit 580abf71f4a35c2cbbcfe532e4f8e3f0803af45a +Subproject commit ff0ece65e1f5f106122ecf315156cc84b5467b99 diff --git a/nodes/fabric-bridge/CMakeLists.txt b/nodes/fabric-bridge/CMakeLists.txt new file mode 100644 index 0000000..a06f0c0 --- /dev/null +++ b/nodes/fabric-bridge/CMakeLists.txt @@ -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 +) \ No newline at end of file diff --git a/nodes/fabric-bridge/src/fabric_bridge_node.cpp b/nodes/fabric-bridge/src/fabric_bridge_node.cpp new file mode 100644 index 0000000..9ad6967 --- /dev/null +++ b/nodes/fabric-bridge/src/fabric_bridge_node.cpp @@ -0,0 +1,320 @@ +#include "fabric_bridge_node.hpp" + +#include + +#include + +namespace dmf_node { + +FabricBridgeNode::~FabricBridgeNode() { + running_ = false; + if (worker_thread_.has_value() && worker_thread_->joinable()) { + worker_thread_->join(); + } +} + +std::vector FabricBridgeNode::input_ports() const { + if (mode_ == FabricMode::Target) { + return {}; + } + return {{"video_in", PortDirection::Input, MediaType::VideoV210}}; +} + +std::vector 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(); + 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(); + if (params.contains("bind_service")) bind_service_ = params["bind_service"].get(); + if (params.contains("provider")) { + auto prov_str = params["provider"].get(); + mxlFabricsProviderFromString(prov_str.c_str(), &provider_); + } + + if (params.contains("target_info")) { + target_info_str_ = params["target_info"].get(); + } +} + +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(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(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(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(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(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(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(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(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(status)); + mxlFabricsDestroyTarget(fab_inst, tgt); + mxlFabricsDestroyInstance(fab_inst); + running_ = false; + return; + } + + size_t info_size = 0; + mxlFabricsTargetInfoToString(target_info, nullptr, &info_size); + std::vector 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(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(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 \ No newline at end of file diff --git a/nodes/fabric-bridge/src/fabric_bridge_node.hpp b/nodes/fabric-bridge/src/fabric_bridge_node.hpp new file mode 100644 index 0000000..3a39a06 --- /dev/null +++ b/nodes/fabric-bridge/src/fabric_bridge_node.hpp @@ -0,0 +1,62 @@ +#pragma once + +#include +#include +#include +#include + +#include +#include +#include +#include + +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 input_ports() const override; + std::vector 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 reader_; + std::optional 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 worker_thread_; + std::atomic running_{false}; +}; + +} // namespace dmf_node \ No newline at end of file diff --git a/nodes/fabric-bridge/src/main.cpp b/nodes/fabric-bridge/src/main.cpp new file mode 100644 index 0000000..250abb2 --- /dev/null +++ b/nodes/fabric-bridge/src/main.cpp @@ -0,0 +1,11 @@ +#include "fabric_bridge_node.hpp" +#include + +int main(int argc, char* argv[]) { + dmf_node::NodeRunner runner; + if (!runner.parse_args(argc, argv)) { + return 1; + } + auto node = std::make_unique(); + return runner.exec(std::move(node)); +} \ No newline at end of file