diff --git a/CMakeLists.txt b/CMakeLists.txt index 12da608..9d8a9d3 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -12,8 +12,9 @@ set(CMAKE_EXPORT_COMPILE_COMMANDS ON) 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_MXL_TOOLS "Build MXL tools (testsrc, sink)" ON) 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) @@ -29,6 +35,12 @@ find_package(Libwebsockets CONFIG REQUIRED) find_package(Catch2 CONFIG QUIET) set(BUILD_TESTS OFF CACHE BOOL "" FORCE) + +if(DMF_BUILD_MXL_TOOLS) + set(BUILD_TOOLS ON CACHE BOOL "" FORCE) + set(BUILD_UTILS ON CACHE BOOL "" FORCE) +endif() + add_subdirectory(extern/mxl) add_subdirectory(libs/dmf-node) add_subdirectory(libs/dmf-engine) diff --git a/libs/dmf-engine/src/api_server.cpp b/libs/dmf-engine/src/api_server.cpp index 7d44fcb..e7601c2 100644 --- a/libs/dmf-engine/src/api_server.cpp +++ b/libs/dmf-engine/src/api_server.cpp @@ -178,7 +178,7 @@ static void handle_request(const std::string& method, const std::string& path, } } else if (path == "/api/graph/start" && method == "POST") { auto nodes = graph.get_nodes(); - uint16_t port = 9100; + uint16_t port = g_impl->port + 100; for (auto& node : nodes) { pm.start_node(const_cast(node), g_impl->mxl_domain, port); cc.register_node(node.id, port); diff --git a/libs/dmf-node/include/dmf-node/node.hpp b/libs/dmf-node/include/dmf-node/node.hpp index 52a9871..e028efa 100644 --- a/libs/dmf-node/include/dmf-node/node.hpp +++ b/libs/dmf-node/include/dmf-node/node.hpp @@ -9,6 +9,7 @@ #include #include +#include #include namespace dmf_node { @@ -21,12 +22,25 @@ public: virtual std::vector input_ports() const = 0; virtual std::vector output_ports() const = 0; virtual void configure(const nlohmann::json& params) = 0; - virtual void on_add_writer(const std::string& port_id, mxlFlowWriter writer) = 0; - virtual void on_add_reader(const std::string& port_id, mxlFlowReader reader) = 0; + virtual void on_add_writer(const std::string& port_id, mxlFlowWriter writer) {} + virtual void on_add_reader(const std::string& port_id, mxlFlowReader reader) {} virtual void on_remove_writer(const std::string& port_id) = 0; 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_; } + + void set_writer_config(const std::string& port_id, const mxlFlowConfigInfo& config) { writer_configs_[port_id] = config; } + void set_reader_config(const std::string& port_id, const mxlFlowConfigInfo& config) { reader_configs_[port_id] = config; } + const mxlFlowConfigInfo* get_writer_config(const std::string& port_id) const { auto it = writer_configs_.find(port_id); return it != writer_configs_.end() ? &it->second : nullptr; } + const mxlFlowConfigInfo* get_reader_config(const std::string& port_id) const { auto it = reader_configs_.find(port_id); return it != reader_configs_.end() ? &it->second : nullptr; } + +private: + mxlInstance mxl_instance_ = nullptr; + std::unordered_map writer_configs_; + std::unordered_map reader_configs_; }; } // namespace dmf_node diff --git a/libs/dmf-node/src/node_runner.cpp b/libs/dmf-node/src/node_runner.cpp index 51629f8..603ba30 100644 --- a/libs/dmf-node/src/node_runner.cpp +++ b/libs/dmf-node/src/node_runner.cpp @@ -57,6 +57,7 @@ int NodeRunner::exec(std::unique_ptr node) { std::signal(SIGINT, signal_handler); std::signal(SIGTERM, signal_handler); + spdlog::flush_on(spdlog::level::info); spdlog::info("Starting node '{}' type='{}'", node_id_, node->type()); if (!std::filesystem::exists(mxl_domain_)) { @@ -71,6 +72,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()) { @@ -114,10 +117,11 @@ int NodeRunner::exec(std::unique_ptr node) { spdlog::error("Failed to create flow writer for flow {}: status={}", flow_id, static_cast(status)); return {{"ok", false}, {"error", "failed to create flow writer"}}; } - spdlog::info("Created flow writer on port '{}' flow {} (created={})", port_id, flow_id, created); +spdlog::info("Created flow writer on port '{}' flow {} (created={})", port_id, flow_id, created); flow_resources.push_back({port_id, writer, nullptr}); + node->set_writer_config(port_id, config_info); node->on_add_writer(port_id, writer); - return {{"ok", true}}; + return {{"ok", true}, {"created", created}}; }); control_server->register_command("add_reader", [&](const nlohmann::json& msg) -> nlohmann::json { @@ -130,8 +134,11 @@ int NodeRunner::exec(std::unique_ptr node) { spdlog::error("Failed to create flow reader for flow {}: status={}", flow_id, static_cast(status)); return {{"ok", false}, {"error", "failed to create flow reader"}}; } - spdlog::info("Created flow reader on port '{}' flow {}", port_id, flow_id); +spdlog::info("Created flow reader on port '{}' flow {}", port_id, flow_id); flow_resources.push_back({port_id, nullptr, reader}); + mxlFlowConfigInfo config{}; + mxlFlowReaderGetConfigInfo(reader, &config); + node->set_reader_config(port_id, config); node->on_add_reader(port_id, reader); return {{"ok", true}}; }); 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..127d446 --- /dev/null +++ b/nodes/fabric-bridge/src/fabric_bridge_node.cpp @@ -0,0 +1,343 @@ +#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(); + if (mode_ == FabricMode::Initiator && reader_.has_value() && !target_info_str_.empty()) { + start_requested_.store(true); + spdlog::info("Fabric-bridge: target_info set, signaling startup"); + } + } +} + +void FabricBridgeNode::on_add_writer(const std::string& port_id, mxlFlowWriter writer) { + if (port_id == "video_out") { + writer_ = writer; + start_requested_.store(true); + spdlog::info("Fabric-bridge: writer added (target mode), signaling startup"); + } +} + +void FabricBridgeNode::on_add_reader(const std::string& port_id, mxlFlowReader reader) { + if (port_id == "video_in") { + reader_ = reader; + if (mode_ == FabricMode::Initiator && !target_info_str_.empty()) { + start_requested_.store(true); + spdlog::info("Fabric-bridge: reader added (initiator mode), signaling startup"); + } else { + spdlog::info("Fabric-bridge: reader added (initiator mode), waiting for target_info"); + } + } +} + +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() && start_requested_.load()) { + running_ = true; + spdlog::info("Fabric-bridge: starting worker thread (mode={})", mode_ == FabricMode::Initiator ? "initiator" : "target"); + if (mode_ == FabricMode::Initiator) { + worker_thread_ = std::thread(&FabricBridgeNode::run_initiator, this); + } else { + 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; + } + + spdlog::info("Fabric-bridge: creating fabrics instance..."); + mxlFabricsInstance fab_inst = nullptr; + auto status = mxlFabricsCreateInstance(instance, nullptr, &fab_inst); + spdlog::info("Fabric-bridge: mxlFabricsCreateInstance returned {}", static_cast(status)); + 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; + + spdlog::info("Fabric-bridge: creating initiator..."); + mxlFabricsInitiator initiator = nullptr; + status = mxlFabricsCreateInitiator(fab_inst, &initiator); + spdlog::info("Fabric-bridge: mxlFabricsCreateInitiator returned {}", static_cast(status)); + 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; + } + + spdlog::info("Fabric-bridge: creating fabrics instance..."); + mxlFabricsInstance fab_inst = nullptr; + auto status = mxlFabricsCreateInstance(instance, nullptr, &fab_inst); + spdlog::info("Fabric-bridge: mxlFabricsCreateInstance returned {}", static_cast(status)); + 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; + + spdlog::info("Fabric-bridge: creating target..."); + mxlFabricsTarget tgt = nullptr; + status = mxlFabricsCreateTarget(fab_inst, &tgt); + spdlog::info("Fabric-bridge: mxlFabricsCreateTarget returned {}", static_cast(status)); + 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; + + auto* config = get_writer_config("video_out"); + if (config) { + grain_rate_ = config->common.grainRate; + spdlog::info("Fabric-bridge: target grain_rate={}/{}", grain_rate_.numerator, grain_rate_.denominator); + } else { + spdlog::warn("Fabric-bridge: no writer config, using default grain_rate"); + } + + 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_; + + spdlog::info("Fabric-bridge: calling mxlFabricsTargetSetup on {}:{}", bind_node_, bind_service_); + 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..ed52c58 --- /dev/null +++ b/nodes/fabric-bridge/src/fabric_bridge_node.hpp @@ -0,0 +1,63 @@ +#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}; + std::atomic start_requested_{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..5e12f59 --- /dev/null +++ b/nodes/fabric-bridge/src/main.cpp @@ -0,0 +1,6 @@ +#include "fabric_bridge_node.hpp" +#include + +int main(int argc, char* argv[]) { + return dmf_node::NodeRunner::run(argc, argv); +} \ No newline at end of file