diff --git a/libs/dmf-engine/src/api_server.cpp b/libs/dmf-engine/src/api_server.cpp index e7601c2..fd5786f 100644 --- a/libs/dmf-engine/src/api_server.cpp +++ b/libs/dmf-engine/src/api_server.cpp @@ -12,6 +12,7 @@ #include #include +#include namespace dmf_engine { @@ -320,6 +321,152 @@ static void handle_request(const std::string& method, const std::string& path, } else { 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(); + 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().empty()) { + target_info = sr["data"]["target_info"].get(); + 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(); + auto initiator_node_id = req_body["initiator_node_id"].get(); + + 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().empty()) { + target_info = sr["data"]["target_info"].get(); + 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(); + } + } 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") { auto prefix = std::string("/api/graph/nodes/"); auto suffix_start = path.find("/command");