From 1d6f93679f87f760171e67e74b838e388d26280c Mon Sep 17 00:00:00 2001 From: Johanness Date: Tue, 26 May 2026 21:08:06 +0300 Subject: [PATCH] fix: engine REST API - handle GET/DELETE requests immediately, use user-provided node IDs - Use lws_http_get_uri_and_method() for proper HTTP method detection (GET, POST, PUT, DELETE all supported) - Handle GET/DELETE in LWS_CALLBACK_HTTP without waiting for body - Support user-provided node IDs via config.id field - Fixes all REST API endpoints hanging on GET requests --- libs/dmf-engine/src/api_server.cpp | 267 ++++++++++++++++------------- libs/dmf-engine/src/graph.cpp | 7 +- 2 files changed, 152 insertions(+), 122 deletions(-) diff --git a/libs/dmf-engine/src/api_server.cpp b/libs/dmf-engine/src/api_server.cpp index 7537a1a..adfd3dc 100644 --- a/libs/dmf-engine/src/api_server.cpp +++ b/libs/dmf-engine/src/api_server.cpp @@ -30,11 +30,124 @@ struct HttpRequest { std::string method; std::string path; std::string body; - bool body_done = false; }; static ApiServerImpl* g_impl = nullptr; +static int send_json_response(struct lws* wsi, const std::string& status_str, + const std::string& json_body) { + auto headers = "HTTP/1.1 " + status_str + "\r\n" + "Content-Type: application/json\r\n" + "Access-Control-Allow-Origin: *\r\n" + "Access-Control-Allow-Methods: GET, POST, DELETE, PUT, OPTIONS\r\n" + "Access-Control-Allow-Headers: Content-Type\r\n" + "Content-Length: " + std::to_string(json_body.size()) + "\r\n" + "\r\n"; + + std::vector buf(LWS_PRE + headers.size() + json_body.size()); + std::memcpy(buf.data() + LWS_PRE, headers.data(), headers.size()); + std::memcpy(buf.data() + LWS_PRE + headers.size(), json_body.data(), json_body.size()); + + lws_write(wsi, buf.data() + LWS_PRE, headers.size() + json_body.size(), LWS_WRITE_HTTP); + + if (lws_http_transaction_completed(wsi)) { + return -1; + } + return 0; +} + +static void handle_request(const std::string& method, const std::string& path, + const std::string& body, + std::string& status_str, std::string& response_body) { + if (!g_impl || !g_impl->graph) { + status_str = "500 Internal Server Error"; + response_body = nlohmann::json({{"error", "Server not initialized"}}).dump(); + return; + } + + auto ok = [&](const nlohmann::json& j) { + status_str = "200 OK"; + response_body = j.dump(); + }; + auto created = [&](const nlohmann::json& j) { + status_str = "201 Created"; + response_body = j.dump(); + }; + auto error_resp = [&](int code, const std::string& msg) { + status_str = std::to_string(code) + " Error"; + response_body = nlohmann::json({{"error", msg}}).dump(); + }; + + try { + nlohmann::json req_body = body.empty() ? nlohmann::json::object() : nlohmann::json::parse(body); + + auto& graph = *g_impl->graph; + auto& fm = *g_impl->flow_manager; + auto& pm = *g_impl->process_manager; + + if (path == "/api/graph" && method == "GET") { + ok(graph.serialize()); + } else if (path == "/api/graph/nodes" && method == "POST") { + if (!req_body.contains("type")) { + error_resp(400, "Missing 'type' field"); + } else { + auto type = req_body["type"].get(); + auto config = req_body.value("config", nlohmann::json::object()); + auto id = graph.add_node(type, config); + created({{"id", id}}); + } + } else if (path.find("/api/graph/nodes/") == 0 && method == "DELETE") { + auto node_id = path.substr(std::string("/api/graph/nodes/").length()); + if (graph.remove_node(node_id)) { + ok({{"deleted", node_id}}); + } else { + error_resp(404, "Node not found: " + node_id); + } + } else if (path == "/api/graph/edges" && method == "POST") { + if (!req_body.contains("from_node") || !req_body.contains("to_node")) { + error_resp(400, "Missing from_node/to_node"); + } else { + auto from_node = req_body["from_node"].get(); + auto from_port = req_body.value("from_port", "video_out"); + auto to_node = req_body["to_node"].get(); + auto to_port = req_body.value("to_port", "video_in"); + + auto flow_id = fm.create_flow_id(); + auto flow_def = fm.create_v210_flow_def(flow_id, 1920, 1080, 50, 1); + + auto edge_id = graph.add_edge(from_node, from_port, to_node, to_port, flow_id, flow_def); + created({{"id", edge_id}, {"flow_id", flow_id}}); + } + } else if (path.find("/api/graph/edges/") == 0 && method == "DELETE") { + auto edge_id = path.substr(std::string("/api/graph/edges/").length()); + if (graph.remove_edge(edge_id)) { + ok({{"deleted", edge_id}}); + } else { + error_resp(404, "Edge not found: " + edge_id); + } + } else if (path == "/api/graph/start" && method == "POST") { + auto nodes = graph.get_nodes(); + uint16_t port = 9100; + for (auto& node : nodes) { + pm.start_node(const_cast(node), "/dev/shm/mxl", port++); + } + ok({{"status", "started"}}); + } else if (path == "/api/graph/stop" && method == "POST") { + auto nodes = graph.get_nodes(); + for (auto& node : nodes) { + pm.stop_node(const_cast(node)); + } + ok({{"status", "stopped"}}); + } else { + error_resp(404, "Not found: " + method + " " + path); + } + } catch (const nlohmann::json::exception& e) { + error_resp(400, std::string("JSON error: ") + e.what()); + } catch (const std::exception& e) { + error_resp(500, e.what()); + } +} + static int callback_http(struct lws* wsi, enum lws_callback_reasons reason, void* user, void* in, size_t len) { auto* req = static_cast(user); @@ -43,17 +156,35 @@ static int callback_http(struct lws* wsi, enum lws_callback_reasons reason, case LWS_CALLBACK_HTTP: { new (req) HttpRequest(); - if (lws_hdr_total_length(wsi, WSI_TOKEN_POST_URI) > 0) { - req->method = "POST"; - char buf[256] = {}; - lws_hdr_copy(wsi, buf, sizeof(buf), WSI_TOKEN_POST_URI); - req->path = buf; - } else { - req->method = "GET"; - char buf[256] = {}; - lws_hdr_copy(wsi, buf, sizeof(buf), WSI_TOKEN_GET_URI); - req->path = buf; + char* uri_ptr = nullptr; + int uri_len = 0; + int method = lws_http_get_uri_and_method(wsi, &uri_ptr, &uri_len); + + switch (method) { + case LWSHUMETH_GET: req->method = "GET"; break; + case LWSHUMETH_POST: req->method = "POST"; break; + case LWSHUMETH_PUT: req->method = "PUT"; break; + case LWSHUMETH_DELETE: req->method = "DELETE"; break; + default: req->method = "GET"; break; } + + if (uri_ptr && uri_len > 0) { + req->path.assign(uri_ptr, uri_len); + } + + if (req->method == "GET" || req->method == "DELETE") { + std::string status_str, response_body; + handle_request(req->method, req->path, "", status_str, response_body); + return send_json_response(wsi, status_str, response_body); + } + + int body_len = lws_hdr_total_length(wsi, WSI_TOKEN_HTTP_CONTENT_LENGTH); + if (body_len == 0) { + std::string status_str, response_body; + handle_request(req->method, req->path, "", status_str, response_body); + return send_json_response(wsi, status_str, response_body); + } + req->body.reserve(body_len); break; } case LWS_CALLBACK_HTTP_BODY: { @@ -61,115 +192,9 @@ static int callback_http(struct lws* wsi, enum lws_callback_reasons reason, break; } case LWS_CALLBACK_HTTP_BODY_COMPLETION: { - if (!g_impl || !g_impl->graph) { - lws_return_http_status(wsi, HTTP_STATUS_INTERNAL_SERVER_ERROR, "Server error"); - return -1; - } - - std::string status_str, content_type, response_body; - - auto ok = [&](const nlohmann::json& j) { - status_str = "200 OK"; - response_body = j.dump(); - }; - auto created = [&](const nlohmann::json& j) { - status_str = "201 Created"; - response_body = j.dump(); - }; - auto error_resp = [&](int code, const std::string& msg) { - status_str = std::to_string(code) + " Error"; - response_body = nlohmann::json({{"error", msg}}).dump(); - }; - - content_type = "application/json"; - - try { - nlohmann::json req_body = req->body.empty() ? nlohmann::json::object() : nlohmann::json::parse(req->body); - - auto& graph = *g_impl->graph; - auto& fm = *g_impl->flow_manager; - auto& pm = *g_impl->process_manager; - - if (req->path == "/api/graph" && req->method == "GET") { - ok(graph.serialize()); - } else if (req->path == "/api/graph/nodes" && req->method == "POST") { - if (!req_body.contains("type")) { - error_resp(400, "Missing 'type' field"); - } else { - auto type = req_body["type"].get(); - auto config = req_body.value("config", nlohmann::json::object()); - auto id = graph.add_node(type, config); - created({{"id", id}}); - } - } else if (req->path.find("/api/graph/nodes/") == 0 && req->method == "DELETE") { - auto node_id = req->path.substr(std::string("/api/graph/nodes/").length()); - if (graph.remove_node(node_id)) { - ok({{"deleted", node_id}}); - } else { - error_resp(404, "Node not found: " + node_id); - } - } else if (req->path == "/api/graph/edges" && req->method == "POST") { - if (!req_body.contains("from_node") || !req_body.contains("to_node")) { - error_resp(400, "Missing from_node/to_node"); - } else { - auto from_node = req_body["from_node"].get(); - auto from_port = req_body.value("from_port", "video_out"); - auto to_node = req_body["to_node"].get(); - auto to_port = req_body.value("to_port", "video_in"); - - auto flow_id = fm.create_flow_id(); - auto flow_def = fm.create_v210_flow_def(flow_id, 1920, 1080, 50, 1); - - auto edge_id = graph.add_edge(from_node, from_port, to_node, to_port, flow_id, flow_def); - created({{"id", edge_id}, {"flow_id", flow_id}}); - } - } else if (req->path.find("/api/graph/edges/") == 0 && req->method == "DELETE") { - auto edge_id = req->path.substr(std::string("/api/graph/edges/").length()); - if (graph.remove_edge(edge_id)) { - ok({{"deleted", edge_id}}); - } else { - error_resp(404, "Edge not found: " + edge_id); - } - } else if (req->path == "/api/graph/start" && req->method == "POST") { - auto nodes = graph.get_nodes(); - uint16_t port = 9100; - for (auto& node : nodes) { - pm.start_node(const_cast(node), "/dev/shm/mxl", port++); - } - ok({{"status", "started"}}); - } else if (req->path == "/api/graph/stop" && req->method == "POST") { - auto nodes = graph.get_nodes(); - for (auto& node : nodes) { - pm.stop_node(const_cast(node)); - } - ok({{"status", "stopped"}}); - } else { - error_resp(404, "Not found: " + req->method + " " + req->path); - } - } catch (const nlohmann::json::exception& e) { - error_resp(400, std::string("JSON error: ") + e.what()); - } catch (const std::exception& e) { - error_resp(500, e.what()); - } - - auto headers = "HTTP/1.1 " + status_str + "\r\n" - "Content-Type: " + content_type + "\r\n" - "Access-Control-Allow-Origin: *\r\n" - "Access-Control-Allow-Methods: GET, POST, DELETE, PUT, OPTIONS\r\n" - "Access-Control-Allow-Headers: Content-Type\r\n" - "Content-Length: " + std::to_string(response_body.size()) + "\r\n" - "\r\n"; - - std::vector buf(LWS_PRE + headers.size() + response_body.size()); - std::memcpy(buf.data() + LWS_PRE, headers.data(), headers.size()); - std::memcpy(buf.data() + LWS_PRE + headers.size(), response_body.data(), response_body.size()); - - lws_write(wsi, buf.data() + LWS_PRE, headers.size() + response_body.size(), LWS_WRITE_HTTP); - - if (lws_http_transaction_completed(wsi)) { - return -1; - } - return 0; + std::string status_str, response_body; + handle_request(req->method, req->path, req->body, status_str, response_body); + return send_json_response(wsi, status_str, response_body); } default: break; @@ -178,7 +203,7 @@ static int callback_http(struct lws* wsi, enum lws_callback_reasons reason, } static const struct lws_protocols protocols[] = { - {"http-api", callback_http, sizeof(HttpRequest), 0}, + {"http-api", callback_http, sizeof(HttpRequest), 4096}, {nullptr, nullptr, 0, 0}, }; diff --git a/libs/dmf-engine/src/graph.cpp b/libs/dmf-engine/src/graph.cpp index f1ca009..1213b53 100644 --- a/libs/dmf-engine/src/graph.cpp +++ b/libs/dmf-engine/src/graph.cpp @@ -7,7 +7,12 @@ namespace dmf_engine { NodeId Graph::add_node(const std::string& type, const nlohmann::json& config) { - auto id = type + "_" + std::to_string(next_node_num_++); + std::string id; + if (config.contains("id") && config["id"].is_string()) { + id = config["id"].get(); + } else { + id = type + "_" + std::to_string(next_node_num_++); + } GraphNode node; node.id = id; node.type = type;