12 Commits

Author SHA1 Message Date
Johanness 991620092e Merge branch 'develop' 2026-06-16 10:42:28 +03:00
Johanness da836785df fix: set running_=true before starting fabric-bridge worker thread 2026-06-14 21:16:23 +03:00
Johanness 0a49b31094 fabric-bridge: flush logs on info, log worker thread start 2026-06-14 21:09:50 +03:00
Johanness 62e2dc49bf cmake: enable MXL tools build by default (testsrc, sink) 2026-06-14 17:03:50 +03:00
Johanness 8fa058f37f fix: add missing return values in add_writer/add_reader handlers (UB causing SIGILL) 2026-06-14 16:57:22 +03:00
Johanness a6ac2bf1a1 fabric-bridge: thread-safe startup, detailed Fabrics API logging 2026-06-14 16:54:37 +03:00
Johanness 17a8b23d03 fabric-bridge: add detailed logging before each Fabrics API call 2026-06-14 16:49:02 +03:00
Johanness a9d4994504 fix: base control port on engine port to avoid conflicts (port+100) 2026-06-14 16:41:52 +03:00
Johanness 919f61fc72 fix: fabric-bridge build - use NodeRunner::run, get writer config from node base class 2026-06-14 15:27:47 +03:00
Johanness 1ab41b6740 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
2026-06-14 11:26:00 +03:00
Johanness 7568c8472d Merge develop: Phase 3 DeckLink I/O + MXL update 2026-06-14 11:20:57 +03:00
Johanness e4f3498617 Merge develop: Phase 1 complete — Framework + Passthrough 2026-05-26 23:04:51 +03:00
9 changed files with 481 additions and 8 deletions
+13 -1
View File
@@ -12,8 +12,9 @@ set(CMAKE_EXPORT_COMPILE_COMMANDS ON)
list(APPEND CMAKE_MODULE_PATH "${CMAKE_CURRENT_SOURCE_DIR}/cmake") list(APPEND CMAKE_MODULE_PATH "${CMAKE_CURRENT_SOURCE_DIR}/cmake")
option(DMF_BUILD_TESTS "Build tests" ON) 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_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") 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) add_subdirectory(nodes/decklink-out)
endif() 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(fmt CONFIG REQUIRED)
find_package(spdlog CONFIG REQUIRED) find_package(spdlog CONFIG REQUIRED)
find_package(nlohmann_json CONFIG REQUIRED) find_package(nlohmann_json CONFIG REQUIRED)
@@ -29,6 +35,12 @@ find_package(Libwebsockets CONFIG REQUIRED)
find_package(Catch2 CONFIG QUIET) find_package(Catch2 CONFIG QUIET)
set(BUILD_TESTS OFF CACHE BOOL "" FORCE) 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(extern/mxl)
add_subdirectory(libs/dmf-node) add_subdirectory(libs/dmf-node)
add_subdirectory(libs/dmf-engine) add_subdirectory(libs/dmf-engine)
+1 -1
View File
@@ -178,7 +178,7 @@ static void handle_request(const std::string& method, const std::string& path,
} }
} else if (path == "/api/graph/start" && method == "POST") { } else if (path == "/api/graph/start" && method == "POST") {
auto nodes = graph.get_nodes(); auto nodes = graph.get_nodes();
uint16_t port = 9100; uint16_t port = g_impl->port + 100;
for (auto& node : nodes) { for (auto& node : nodes) {
pm.start_node(const_cast<GraphNode&>(node), g_impl->mxl_domain, port); pm.start_node(const_cast<GraphNode&>(node), g_impl->mxl_domain, port);
cc.register_node(node.id, port); cc.register_node(node.id, port);
+16 -2
View File
@@ -9,6 +9,7 @@
#include <nlohmann/json.hpp> #include <nlohmann/json.hpp>
#include <memory> #include <memory>
#include <unordered_map>
#include <vector> #include <vector>
namespace dmf_node { namespace dmf_node {
@@ -21,12 +22,25 @@ public:
virtual std::vector<PortDef> input_ports() const = 0; virtual std::vector<PortDef> input_ports() const = 0;
virtual std::vector<PortDef> output_ports() const = 0; virtual std::vector<PortDef> output_ports() const = 0;
virtual void configure(const nlohmann::json& params) = 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_writer(const std::string& port_id, mxlFlowWriter writer) {}
virtual void on_add_reader(const std::string& port_id, mxlFlowReader reader) = 0; 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_writer(const std::string& port_id) = 0;
virtual void on_remove_reader(const std::string& port_id) = 0; virtual void on_remove_reader(const std::string& port_id) = 0;
virtual void process() = 0; virtual void process() = 0;
virtual nlohmann::json status() const = 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<std::string, mxlFlowConfigInfo> writer_configs_;
std::unordered_map<std::string, mxlFlowConfigInfo> reader_configs_;
}; };
} // namespace dmf_node } // namespace dmf_node
+10 -3
View File
@@ -57,6 +57,7 @@ int NodeRunner::exec(std::unique_ptr<Node> node) {
std::signal(SIGINT, signal_handler); std::signal(SIGINT, signal_handler);
std::signal(SIGTERM, signal_handler); std::signal(SIGTERM, signal_handler);
spdlog::flush_on(spdlog::level::info);
spdlog::info("Starting node '{}' type='{}'", node_id_, node->type()); spdlog::info("Starting node '{}' type='{}'", node_id_, node->type());
if (!std::filesystem::exists(mxl_domain_)) { if (!std::filesystem::exists(mxl_domain_)) {
@@ -71,6 +72,8 @@ int NodeRunner::exec(std::unique_ptr<Node> node) {
} }
spdlog::info("MXL instance created on domain: {}", mxl_domain_); spdlog::info("MXL instance created on domain: {}", mxl_domain_);
node->set_mxl_instance(mxl_instance_);
mxlGarbageCollectFlows(mxl_instance_); mxlGarbageCollectFlows(mxl_instance_);
if (!config_str_.empty()) { if (!config_str_.empty()) {
@@ -114,10 +117,11 @@ int NodeRunner::exec(std::unique_ptr<Node> node) {
spdlog::error("Failed to create flow writer for flow {}: status={}", flow_id, static_cast<int>(status)); spdlog::error("Failed to create flow writer for flow {}: status={}", flow_id, static_cast<int>(status));
return {{"ok", false}, {"error", "failed to create flow writer"}}; 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}); flow_resources.push_back({port_id, writer, nullptr});
node->set_writer_config(port_id, config_info);
node->on_add_writer(port_id, writer); 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 { control_server->register_command("add_reader", [&](const nlohmann::json& msg) -> nlohmann::json {
@@ -130,8 +134,11 @@ int NodeRunner::exec(std::unique_ptr<Node> node) {
spdlog::error("Failed to create flow reader for flow {}: status={}", flow_id, static_cast<int>(status)); spdlog::error("Failed to create flow reader for flow {}: status={}", flow_id, static_cast<int>(status));
return {{"ok", false}, {"error", "failed to create flow reader"}}; 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}); 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); node->on_add_reader(port_id, reader);
return {{"ok", true}}; return {{"ok", true}};
}); });
+1 -1
Submodule mxl updated: 580abf71f4...ff0ece65e1
+28
View File
@@ -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,343 @@
#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>();
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<int>(status));
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;
spdlog::info("Fabric-bridge: creating initiator...");
mxlFabricsInitiator initiator = nullptr;
status = mxlFabricsCreateInitiator(fab_inst, &initiator);
spdlog::info("Fabric-bridge: mxlFabricsCreateInitiator returned {}", static_cast<int>(status));
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;
}
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<int>(status));
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;
spdlog::info("Fabric-bridge: creating target...");
mxlFabricsTarget tgt = nullptr;
status = mxlFabricsCreateTarget(fab_inst, &tgt);
spdlog::info("Fabric-bridge: mxlFabricsCreateTarget returned {}", static_cast<int>(status));
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;
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<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,63 @@
#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};
std::atomic<bool> start_requested_{false};
};
} // namespace dmf_node
+6
View File
@@ -0,0 +1,6 @@
#include "fabric_bridge_node.hpp"
#include <dmf-node/node_runner.hpp>
int main(int argc, char* argv[]) {
return dmf_node::NodeRunner::run<dmf_node::FabricBridgeNode>(argc, argv);
}