|
|
@@ -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);
|