feat: add engine API for fabric-bridge target_info exchange

- POST /api/fabric/create-target: add_writer to target node, poll for target_info
- POST /api/fabric/connect: get target_info from target, configure initiator node
This commit is contained in:
Johanness
2026-06-16 10:50:04 +03:00
parent 991620092e
commit ac3a506a5c
+147
View File
@@ -12,6 +12,7 @@
#include <cstring> #include <cstring>
#include <string> #include <string>
#include <thread>
namespace dmf_engine { namespace dmf_engine {
@@ -320,6 +321,152 @@ 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 == "/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");