Compare commits
18 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 12896182c9 | |||
| d093d64ad1 | |||
| e6e439b16e | |||
| 0944f574ad | |||
| cca0cd810a | |||
| ac3a506a5c | |||
| 991620092e | |||
| da836785df | |||
| 0a49b31094 | |||
| 62e2dc49bf | |||
| 8fa058f37f | |||
| a6ac2bf1a1 | |||
| 17a8b23d03 | |||
| a9d4994504 | |||
| 919f61fc72 | |||
| 1ab41b6740 | |||
| 7568c8472d | |||
| e4f3498617 |
+13
-1
@@ -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)
|
||||||
|
|||||||
@@ -12,6 +12,7 @@
|
|||||||
|
|
||||||
#include <cstring>
|
#include <cstring>
|
||||||
#include <string>
|
#include <string>
|
||||||
|
#include <thread>
|
||||||
|
|
||||||
namespace dmf_engine {
|
namespace dmf_engine {
|
||||||
|
|
||||||
@@ -178,7 +179,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);
|
||||||
@@ -320,6 +321,172 @@ static void handle_request(const std::string& method, const std::string& path,
|
|||||||
} else {
|
} else {
|
||||||
error_resp(500, "Failed to send command to node");
|
error_resp(500, "Failed to send command to node");
|
||||||
}
|
}
|
||||||
|
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/register") != std::string::npos && method == "POST") {
|
||||||
|
auto prefix = std::string("/api/graph/nodes/");
|
||||||
|
auto suffix_start = path.find("/register");
|
||||||
|
auto node_id = path.substr(prefix.length(), suffix_start - prefix.length());
|
||||||
|
|
||||||
|
if (!req_body.contains("port")) {
|
||||||
|
error_resp(400, "Missing port");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
auto control_port = static_cast<uint16_t>(req_body["port"].get<int>());
|
||||||
|
cc.register_node(node_id, control_port);
|
||||||
|
|
||||||
|
auto* node = graph.get_node_mut(node_id);
|
||||||
|
if (node) {
|
||||||
|
node->control_port = control_port;
|
||||||
|
node->state = NodeState::Running;
|
||||||
|
}
|
||||||
|
|
||||||
|
ok({{"node_id", node_id}, {"control_port", control_port}});
|
||||||
|
} else if (path == "/api/fabric/create-target" && method == "POST") {
|
||||||
|
if (!req_body.contains("node_id")) {
|
||||||
|
error_resp(400, "Missing node_id");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
auto node_id = req_body["node_id"].get<std::string>();
|
||||||
|
auto port_id = req_body.value("port_id", "video_out");
|
||||||
|
auto flow_id = req_body.value("flow_id", fm.create_flow_id());
|
||||||
|
|
||||||
|
auto port_num = cc.get_port(node_id);
|
||||||
|
if (port_num == 0) {
|
||||||
|
error_resp(404, "Node not running: " + node_id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
int fps_num = 25, fps_den = 1;
|
||||||
|
int width = 1920, height = 1080;
|
||||||
|
|
||||||
|
nlohmann::json status_cmd;
|
||||||
|
status_cmd["cmd"] = "status";
|
||||||
|
auto status_resp = cc.send_command_with_response(port_num, status_cmd.dump());
|
||||||
|
if (!status_resp.empty()) {
|
||||||
|
try {
|
||||||
|
auto sr = nlohmann::json::parse(status_resp);
|
||||||
|
if (sr.contains("data")) {
|
||||||
|
auto& d = sr["data"];
|
||||||
|
if (d.contains("grain_rate")) {
|
||||||
|
fps_num = d["grain_rate"].value("numerator", fps_num);
|
||||||
|
fps_den = d["grain_rate"].value("denominator", fps_den);
|
||||||
|
}
|
||||||
|
width = d.value("width", width);
|
||||||
|
height = d.value("height", height);
|
||||||
|
}
|
||||||
|
} catch (...) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
auto flow_def = fm.create_v210_flow_def(flow_id, width, height, fps_num, fps_den);
|
||||||
|
|
||||||
|
nlohmann::json cmd;
|
||||||
|
cmd["cmd"] = "add_writer";
|
||||||
|
cmd["port_id"] = port_id;
|
||||||
|
cmd["flow_id"] = flow_id;
|
||||||
|
cmd["flow_def"] = flow_def;
|
||||||
|
if (!cc.send_command(port_num, cmd.dump())) {
|
||||||
|
error_resp(500, "Failed to add writer to target node");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::string target_info;
|
||||||
|
constexpr int max_poll_attempts = 30;
|
||||||
|
for (int i = 0; i < max_poll_attempts; ++i) {
|
||||||
|
auto resp = cc.send_command_with_response(port_num, status_cmd.dump());
|
||||||
|
if (!resp.empty()) {
|
||||||
|
try {
|
||||||
|
auto sr = nlohmann::json::parse(resp);
|
||||||
|
if (sr.contains("data") && sr["data"].contains("target_info") && !sr["data"]["target_info"].get<std::string>().empty()) {
|
||||||
|
target_info = sr["data"]["target_info"].get<std::string>();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
} catch (...) {}
|
||||||
|
}
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
||||||
|
}
|
||||||
|
|
||||||
|
if (target_info.empty()) {
|
||||||
|
error_resp(504, "Target node did not produce target_info in time");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
ok({
|
||||||
|
{"node_id", node_id},
|
||||||
|
{"flow_id", flow_id},
|
||||||
|
{"target_info", target_info}
|
||||||
|
});
|
||||||
|
} else if (path == "/api/fabric/connect" && method == "POST") {
|
||||||
|
if (!req_body.contains("target_node_id") || !req_body.contains("initiator_node_id")) {
|
||||||
|
error_resp(400, "Missing target_node_id/initiator_node_id");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
auto target_node_id = req_body["target_node_id"].get<std::string>();
|
||||||
|
auto initiator_node_id = req_body["initiator_node_id"].get<std::string>();
|
||||||
|
|
||||||
|
auto target_port = cc.get_port(target_node_id);
|
||||||
|
auto initiator_port = cc.get_port(initiator_node_id);
|
||||||
|
|
||||||
|
if (target_port == 0) {
|
||||||
|
error_resp(404, "Target node not running: " + target_node_id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (initiator_port == 0) {
|
||||||
|
error_resp(404, "Initiator node not running: " + initiator_node_id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
nlohmann::json status_cmd;
|
||||||
|
status_cmd["cmd"] = "status";
|
||||||
|
std::string target_info;
|
||||||
|
constexpr int max_poll_attempts = 30;
|
||||||
|
for (int i = 0; i < max_poll_attempts; ++i) {
|
||||||
|
auto resp = cc.send_command_with_response(target_port, status_cmd.dump());
|
||||||
|
if (!resp.empty()) {
|
||||||
|
try {
|
||||||
|
auto sr = nlohmann::json::parse(resp);
|
||||||
|
if (sr.contains("data") && sr["data"].contains("target_info") && !sr["data"]["target_info"].get<std::string>().empty()) {
|
||||||
|
target_info = sr["data"]["target_info"].get<std::string>();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
} catch (...) {}
|
||||||
|
}
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
||||||
|
}
|
||||||
|
|
||||||
|
if (target_info.empty()) {
|
||||||
|
error_resp(504, "Target node did not produce target_info in time");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
nlohmann::json configure_cmd;
|
||||||
|
configure_cmd["cmd"] = "configure";
|
||||||
|
configure_cmd["params"]["target_info"] = target_info;
|
||||||
|
if (!cc.send_command(initiator_port, configure_cmd.dump())) {
|
||||||
|
error_resp(500, "Failed to send target_info to initiator node");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(200));
|
||||||
|
|
||||||
|
bool initiator_running = false;
|
||||||
|
auto init_resp = cc.send_command_with_response(initiator_port, status_cmd.dump());
|
||||||
|
if (!init_resp.empty()) {
|
||||||
|
try {
|
||||||
|
auto ir = nlohmann::json::parse(init_resp);
|
||||||
|
if (ir.contains("data") && ir["data"].contains("running")) {
|
||||||
|
initiator_running = ir["data"]["running"].get<bool>();
|
||||||
|
}
|
||||||
|
} catch (...) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
ok({
|
||||||
|
{"target_node_id", target_node_id},
|
||||||
|
{"initiator_node_id", initiator_node_id},
|
||||||
|
{"target_info_length", target_info.size()},
|
||||||
|
{"initiator_running", initiator_running}
|
||||||
|
});
|
||||||
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/command") != std::string::npos && method == "POST") {
|
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/command") != std::string::npos && method == "POST") {
|
||||||
auto prefix = std::string("/api/graph/nodes/");
|
auto prefix = std::string("/api/graph/nodes/");
|
||||||
auto suffix_start = path.find("/command");
|
auto suffix_start = path.find("/command");
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
@@ -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,470 @@
|
|||||||
|
#include "fabric_bridge_node.hpp"
|
||||||
|
|
||||||
|
#include <mxl/mxl.h>
|
||||||
|
|
||||||
|
#include <spdlog/spdlog.h>
|
||||||
|
|
||||||
|
#include <sys/socket.h>
|
||||||
|
#include <netinet/in.h>
|
||||||
|
#include <arpa/inet.h>
|
||||||
|
#include <unistd.h>
|
||||||
|
#include <netdb.h>
|
||||||
|
#include <cerrno>
|
||||||
|
|
||||||
|
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");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (params.contains("target_host")) {
|
||||||
|
target_host_ = params["target_host"].get<std::string>();
|
||||||
|
target_port_ = params.value("target_port", static_cast<uint16_t>(9100));
|
||||||
|
spdlog::info("Fabric-bridge: will fetch target_info from {}:{}", target_host_, target_port_);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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) {
|
||||||
|
if (!target_info_str_.empty() || !target_host_.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::fetch_target_info() {
|
||||||
|
if (target_host_.empty()) return;
|
||||||
|
|
||||||
|
spdlog::info("Fabric-bridge: fetching target_info from {}:{}", target_host_, target_port_);
|
||||||
|
|
||||||
|
std::string body = R"({"cmd":"status"})";
|
||||||
|
std::string request = "POST /cmd HTTP/1.1\r\n"
|
||||||
|
"Host: " + target_host_ + "\r\n"
|
||||||
|
"Content-Type: application/json\r\n"
|
||||||
|
"Connection: close\r\n"
|
||||||
|
"Content-Length: " + std::to_string(body.size()) + "\r\n"
|
||||||
|
"\r\n" +
|
||||||
|
body;
|
||||||
|
|
||||||
|
for (int attempt = 0; attempt < 30 && target_info_str_.empty(); ++attempt) {
|
||||||
|
int sock = socket(AF_INET, SOCK_STREAM, 0);
|
||||||
|
if (sock < 0) {
|
||||||
|
spdlog::warn("Fabric-bridge: socket() failed, retry {}...", attempt + 1);
|
||||||
|
std::this_thread::sleep_for(std::chrono::seconds(1));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
struct timeval tv{};
|
||||||
|
tv.tv_sec = 2;
|
||||||
|
tv.tv_usec = 0;
|
||||||
|
setsockopt(sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
|
||||||
|
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
|
||||||
|
|
||||||
|
struct hostent* he = gethostbyname(target_host_.c_str());
|
||||||
|
if (!he) {
|
||||||
|
spdlog::warn("Fabric-bridge: cannot resolve {}, retry {}...", target_host_, attempt + 1);
|
||||||
|
close(sock);
|
||||||
|
std::this_thread::sleep_for(std::chrono::seconds(1));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
struct sockaddr_in addr{};
|
||||||
|
addr.sin_family = AF_INET;
|
||||||
|
addr.sin_port = htons(target_port_);
|
||||||
|
std::memcpy(&addr.sin_addr, he->h_addr_list[0], he->h_length);
|
||||||
|
|
||||||
|
if (connect(sock, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) < 0) {
|
||||||
|
spdlog::warn("Fabric-bridge: connect to {}:{} failed (errno={}), retry {}...", target_host_, target_port_, errno, attempt + 1);
|
||||||
|
close(sock);
|
||||||
|
std::this_thread::sleep_for(std::chrono::seconds(1));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (send(sock, request.c_str(), request.size(), 0) < 0) {
|
||||||
|
spdlog::warn("Fabric-bridge: send failed, retry {}...", attempt + 1);
|
||||||
|
close(sock);
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::string response;
|
||||||
|
char buf[4096];
|
||||||
|
ssize_t n;
|
||||||
|
while ((n = recv(sock, buf, sizeof(buf), 0)) > 0) {
|
||||||
|
response.append(buf, n);
|
||||||
|
}
|
||||||
|
close(sock);
|
||||||
|
|
||||||
|
spdlog::info("Fabric-bridge: received {} bytes from target", response.size());
|
||||||
|
|
||||||
|
auto body_pos = response.find("\r\n\r\n");
|
||||||
|
if (body_pos == std::string::npos) {
|
||||||
|
spdlog::warn("Fabric-bridge: no HTTP body in response, retry {}...", attempt + 1);
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::string response_body = response.substr(body_pos + 4);
|
||||||
|
|
||||||
|
try {
|
||||||
|
auto json = nlohmann::json::parse(response_body);
|
||||||
|
if (json.contains("data") && json["data"].contains("target_info")) {
|
||||||
|
auto ti = json["data"]["target_info"].get<std::string>();
|
||||||
|
if (!ti.empty()) {
|
||||||
|
target_info_str_ = ti;
|
||||||
|
spdlog::info("Fabric-bridge: fetched target_info ({} bytes)", target_info_str_.size());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
spdlog::debug("Fabric-bridge: response data: {}", json.contains("data") ? json["data"].dump() : "(no data)");
|
||||||
|
} catch (const nlohmann::json::exception& e) {
|
||||||
|
spdlog::warn("Fabric-bridge: JSON parse error: {}", e.what());
|
||||||
|
}
|
||||||
|
|
||||||
|
spdlog::info("Fabric-bridge: target_info not ready yet, retry {}...", attempt + 1);
|
||||||
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
||||||
|
}
|
||||||
|
|
||||||
|
if (target_info_str_.empty()) {
|
||||||
|
spdlog::error("Fabric-bridge: failed to fetch target_info after 30 attempts");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void FabricBridgeNode::process() {
|
||||||
|
if (!worker_thread_.has_value() && start_requested_.load()) {
|
||||||
|
if (mode_ == FabricMode::Initiator && reader_.has_value() && !target_host_.empty() && target_info_str_.empty()) {
|
||||||
|
fetch_target_info();
|
||||||
|
}
|
||||||
|
if (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_;
|
||||||
|
|
||||||
|
spdlog::info("Fabric-bridge: calling mxlFabricsInitiatorSetup on {}:{}", bind_node_, bind_service_);
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (target_info_str_.empty()) {
|
||||||
|
spdlog::error("Fabric-bridge: no target_info available, cannot add target");
|
||||||
|
mxlFabricsDestroyInitiator(fab_inst, initiator);
|
||||||
|
mxlFabricsDestroyInstance(fab_inst);
|
||||||
|
running_ = false;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
mxlFabricsTargetInfo target_info = nullptr;
|
||||||
|
spdlog::info("Fabric-bridge: parsing target_info ({} bytes)", target_info_str_.size());
|
||||||
|
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_);
|
||||||
|
|
||||||
|
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,68 @@
|
|||||||
|
#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_;
|
||||||
|
|
||||||
|
std::string target_host_;
|
||||||
|
uint16_t target_port_ = 9100;
|
||||||
|
|
||||||
|
void fetch_target_info();
|
||||||
|
|
||||||
|
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
|
||||||
@@ -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);
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user