فهرست منبع

Merge branch 'node-options-endpoint'

fszontagh 1 ماه پیش
والد
کامیت
b843803a31

+ 9 - 0
proto/workflow.proto

@@ -115,6 +115,15 @@ message ExecuteWorkflowRequest {
     string trigger_data = 3;  // JSON string
     bool wait_for_completion = 4;
     int32 timeout_ms = 5;
+    // A whole workflow, as JSON, to run instead of loading workflow_id from the
+    // database. Used to run one node on demand - the editor asking a node what
+    // options it can offer - where the node's config is being edited and has
+    // not been saved anywhere.
+    //
+    // workflow_id must still be set to the workflow this node belongs to.
+    // Credential access is checked against it, and an empty workflow id is
+    // treated as an administrative call with access to every credential.
+    string inline_workflow = 6;
 }
 
 // Execute workflow response

+ 40 - 0
src/runner/runner_service.cpp

@@ -47,6 +47,46 @@ grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
                                                 proto::ExecuteWorkflowResponse* response) {
     LOG_INFO("ExecuteWorkflow called for workflow: {}", request->workflow_id());
 
+    // An inline workflow runs as given, without being stored anywhere. This is
+    // how the editor asks a node what it can offer - which models a server has,
+    // say - while its config is still being edited.
+    if (!request->inline_workflow().empty()) {
+        nlohmann::json inline_doc;
+        try {
+            inline_doc = nlohmann::json::parse(request->inline_workflow());
+        } catch (const std::exception& e) {
+            response->set_status(proto::EXECUTION_STATUS_FAILED);
+            auto* error = response->mutable_error();
+            error->set_message(std::string("inline_workflow is not valid JSON: ") + e.what());
+            return grpc::Status::OK;
+        }
+
+        // The id decides which credentials this run may read, so it comes from
+        // the request rather than from the caller-supplied document.
+        inline_doc["_id"] = request->workflow_id();
+
+        auto inline_workflow = Workflow::fromJson(inline_doc);
+        nlohmann::json inline_trigger;
+        if (!request->trigger_data().empty()) {
+            try {
+                inline_trigger = nlohmann::json::parse(request->trigger_data());
+            } catch (...) {}
+        }
+
+        auto inline_outcome = engine_.execute(inline_workflow, "manual", inline_trigger, nullptr);
+        if (inline_outcome.failed()) {
+            response->set_status(proto::EXECUTION_STATUS_FAILED);
+            auto* error = response->mutable_error();
+            error->set_message(inline_outcome.error().message());
+            return grpc::Status::OK;
+        }
+        const auto& inline_result = inline_outcome.value();
+        response->set_execution_id(inline_result.execution_id);
+        response->set_status(toProtoStatus(inline_result.status));
+        response->set_result(inline_result.toJson().dump());
+        return grpc::Status::OK;
+    }
+
     // Get workflow from database
     auto workflow_result = storage_.get("workflows", request->workflow_id());
     if (workflow_result.failed()) {

+ 116 - 1
src/webserver/api/node_controller.cpp

@@ -1,4 +1,7 @@
 #include "node_controller.hpp"
+
+#include "proto/runner.grpc.pb.h"
+#include <grpcpp/grpcpp.h>
 #include "../grpc/node_sync_service.hpp"
 #include "logging/logger.hpp"
 
@@ -6,10 +9,122 @@ namespace smartbotic::webserver::api {
 
 NodeController::NodeController(nodes::NodeStore& node_store,
                                auth::AuthMiddleware& middleware,
+                               runners::LoadBalancer& load_balancer,
                                grpc::NodeSyncServiceImpl* sync_service)
-    : node_store_(node_store), middleware_(middleware), sync_service_(sync_service) {}
+    : node_store_(node_store), middleware_(middleware),
+      load_balancer_(load_balancer), sync_service_(sync_service) {}
+
+// Run a single node with the config the editor currently holds, and return what
+// it produced. A field whose choices only the far service knows - which models
+// a server has - is filled this way.
+//
+// The node runs on a runner, in its own code, with its own credential. That is
+// the only way to reach a service whose authentication is more than a fixed
+// header: sd.cpp exchanges a username and password for a token before it will
+// answer, and no generic proxy can do that. It also means the knowledge stays
+// in the node, where a new integration can add it without a rebuild.
+void NodeController::nodeOptions(const httplib::Request& req, httplib::Response& res,
+                                 const auth::AuthContext& ctx) {
+    (void)ctx;
+    const std::string node_type = req.matches[1];
+
+    nlohmann::json body;
+    try {
+        body = nlohmann::json::parse(req.body);
+    } catch (const std::exception& e) {
+        sendError(res, std::string("Invalid JSON: ") + e.what(), 400);
+        return;
+    }
+
+    // The workflow this node belongs to. Credential access is checked against
+    // it, and an empty id counts as an administrative call able to read every
+    // credential - so a request without one is refused rather than quietly
+    // granted more than the real workflow would have.
+    const std::string workflow_id = body.value("workflowId", "");
+    if (workflow_id.empty()) {
+        sendError(res, "workflowId is required - it decides which credentials this may use", 400);
+        return;
+    }
+
+    nlohmann::json config = body.value("config", nlohmann::json::object());
+    if (!config.is_object()) {
+        sendError(res, "config must be an object", 400);
+        return;
+    }
+
+    auto runner = load_balancer_.selectRunner();
+    if (!runner) {
+        sendError(res, "No runners available", 503);
+        return;
+    }
+
+    // One node, no connections. It runs with empty input, which is all a
+    // listing needs.
+    nlohmann::json node;
+    node["id"] = "options";
+    node["type"] = node_type;
+    node["name"] = "Options";
+    node["config"] = config;
+    node["position"] = {{"x", 0}, {"y", 0}};
+
+    nlohmann::json inline_workflow;
+    inline_workflow["name"] = "options:" + node_type;
+    inline_workflow["nodes"] = nlohmann::json::array({node});
+    inline_workflow["connections"] = nlohmann::json::array();
+    inline_workflow["settings"] = nlohmann::json::object();
+
+    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto stub = proto::RunnerService::NewStub(channel);
+
+    proto::ExecuteWorkflowRequest grpc_req;
+    grpc_req.set_workflow_id(workflow_id);
+    grpc_req.set_trigger_type("manual");
+    grpc_req.set_inline_workflow(inline_workflow.dump());
+    grpc_req.set_wait_for_completion(true);
+
+    proto::ExecuteWorkflowResponse grpc_res;
+    ::grpc::ClientContext grpc_ctx;
+    grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(60));
+
+    auto status = stub->ExecuteWorkflow(&grpc_ctx, grpc_req, &grpc_res);
+    if (!status.ok()) {
+        sendError(res, "Could not reach a runner: " + status.error_message(), 502);
+        return;
+    }
+
+    nlohmann::json execution;
+    try {
+        execution = nlohmann::json::parse(grpc_res.result());
+    } catch (...) {
+        sendError(res, "The runner returned something unreadable", 502);
+        return;
+    }
+
+    // Report the node's own failure as the failure, rather than a bare 500:
+    // "no credential" and "server refused" are things the person configuring
+    // this needs to read.
+    for (const auto& node_execution : execution.value("nodeExecutions", nlohmann::json::array())) {
+        if (node_execution.value("nodeId", "") != "options") {
+            continue;
+        }
+        if (node_execution.value("status", "") != "completed") {
+            sendError(res, node_execution.value("error", "The node did not complete"), 422);
+            return;
+        }
+        sendJson(res, {{"output", node_execution.value("output", nlohmann::json::object())}});
+        return;
+    }
+
+    sendError(res, execution.value("error", "The node did not run"), 422);
+}
 
 void NodeController::registerRoutes(httplib::Server& server) {
+    server.Post(R"(/api/v1/nodes/([^/]+)/options)", [this](const httplib::Request& req, httplib::Response& res) {
+        middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
+            nodeOptions(req, res, ctx);
+        });
+    });
+
     server.Get("/api/v1/nodes", [this](const httplib::Request& req, httplib::Response& res) {
         middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
             listNodes(req, res, ctx);

+ 7 - 0
src/webserver/api/node_controller.hpp

@@ -3,6 +3,7 @@
 #include <httplib.h>
 #include <nlohmann/json.hpp>
 #include "../auth/auth_middleware.hpp"
+#include "../runners/load_balancer.hpp"
 #include "../nodes/node_store.hpp"
 
 namespace smartbotic::webserver::grpc {
@@ -15,6 +16,7 @@ class NodeController {
 public:
     NodeController(nodes::NodeStore& node_store,
                    auth::AuthMiddleware& middleware,
+                   runners::LoadBalancer& load_balancer,
                    grpc::NodeSyncServiceImpl* sync_service = nullptr);
 
     void registerRoutes(httplib::Server& server);
@@ -34,12 +36,17 @@ private:
                     const auth::AuthContext& ctx);
     void migrateNodes(const httplib::Request& req, httplib::Response& res,
                       const auth::AuthContext& ctx);
+    // Run one node on demand and hand back its output, so a field can offer
+    // choices that only the far service knows - which models a server has.
+    void nodeOptions(const httplib::Request& req, httplib::Response& res,
+                     const auth::AuthContext& ctx);
 
     void sendJson(httplib::Response& res, const nlohmann::json& data, int status = 200);
     void sendError(httplib::Response& res, const std::string& message, int status);
 
     nodes::NodeStore& node_store_;
     auth::AuthMiddleware& middleware_;
+    runners::LoadBalancer& load_balancer_;
     grpc::NodeSyncServiceImpl* sync_service_;  // Optional: for notifying runners
 };
 

+ 1 - 1
src/webserver/webserver_service.cpp

@@ -230,7 +230,7 @@ void WebServerService::setupRoutes() {
     proxy_ctrl_->registerRoutes(server);
 
     node_ctrl_ = std::make_unique<api::NodeController>(
-        *node_store_, *auth_middleware_, &node_sync_server_->service());
+        *node_store_, *auth_middleware_, *load_balancer_, &node_sync_server_->service());
     node_ctrl_->registerRoutes(server);
 
     runner_ctrl_ = std::make_unique<api::RunnerController>(*runner_registry_, *auth_middleware_);