浏览代码

feat: let a node declare the HTTP response a webhook returns

fszontagh 1 月之前
父节点
当前提交
faae30c015

+ 7 - 1
src/runner/runner_service.cpp

@@ -122,7 +122,13 @@ grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
     response->set_status(static_cast<proto::ExecutionStatus>(result.value().status));
 
     if (request->wait_for_completion()) {
-        response->set_result(result.value().final_output.dump());
+        if (!result.value().webhook_response.is_null()) {
+            nlohmann::json envelope;
+            envelope["_webhookResponse"] = result.value().webhook_response;
+            response->set_result(envelope.dump());
+        } else {
+            response->set_result(result.value().final_output.dump());
+        }
     }
 
     return grpc::Status::OK;

+ 16 - 0
src/runner/workflow_engine.cpp

@@ -159,6 +159,10 @@ nlohmann::json ExecutionResult::toJson() const {
         j["workflowSnapshot"] = workflow_snapshot;
     }
 
+    if (!webhook_response.is_null()) {
+        j["webhookResponse"] = webhook_response;
+    }
+
     return j;
 }
 
@@ -600,6 +604,18 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                 }
             }
 
+            // Any node may declare the HTTP response, not just the last one to
+            // run, so appending a node to a workflow cannot silently change
+            // what its webhook returns.
+            if (node_result.status == NodeStatus::Completed &&
+                node_result.output.contains("_webhookResponse")) {
+                if (!result.webhook_response.is_null()) {
+                    LOG_WARN("Node {} overrides a webhook response already set by an earlier node",
+                             node_id);
+                }
+                result.webhook_response = node_result.output["_webhookResponse"];
+            }
+
             // Check for loop node
             if (node_result.status == NodeStatus::Completed &&
                 node_result.output.contains("_isLoop") &&

+ 1 - 0
src/runner/workflow_engine.hpp

@@ -101,6 +101,7 @@ struct ExecutionResult {
     std::string error;
     nlohmann::json final_output;
     nlohmann::json workflow_snapshot;  // Snapshot of workflow at execution time
+    nlohmann::json webhook_response;   // Set by a respond-to-webhook node, if any
 
     nlohmann::json toJson() const;
 };

+ 52 - 2
src/webserver/api/webhook_controller.cpp

@@ -3,6 +3,7 @@
 #include "common/time_utils.hpp"
 #include "proto/runner.grpc.pb.h"
 #include <grpcpp/grpcpp.h>
+#include <cctype>
 #include <optional>
 
 namespace smartbotic::webserver::api {
@@ -171,10 +172,59 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
 
     // Return result
     if (!grpc_res.result().empty()) {
+        nlohmann::json parsed;
+        bool parsed_ok = true;
         try {
-            auto result = nlohmann::json::parse(grpc_res.result());
-            sendJson(res, result);
+            parsed = nlohmann::json::parse(grpc_res.result());
         } catch (...) {
+            parsed_ok = false;
+        }
+
+        if (parsed_ok && parsed.is_object() && parsed.contains("_webhookResponse")) {
+            const auto& spec = parsed["_webhookResponse"];
+
+            int status_code = spec.value("status", 200);
+            if (status_code < 100 || status_code > 599) {
+                LOG_WARN("Webhook response asked for status {}, which is not a valid HTTP status; sending 500",
+                         status_code);
+                status_code = 500;
+            }
+
+            std::string content_type = "application/json";
+            if (spec.contains("headers") && spec["headers"].is_object()) {
+                for (auto it = spec["headers"].begin(); it != spec["headers"].end(); ++it) {
+                    if (!it.value().is_string()) {
+                        continue;
+                    }
+                    // Content-Type reaches httplib through set_content rather
+                    // than as a header, and setting both sends it twice.
+                    std::string name = it.key();
+                    std::string lowered;
+                    for (char c : name) {
+                        lowered += static_cast<char>(std::tolower(static_cast<unsigned char>(c)));
+                    }
+                    if (lowered == "content-type") {
+                        content_type = it.value().get<std::string>();
+                    } else {
+                        res.set_header(name.c_str(), it.value().get<std::string>().c_str());
+                    }
+                }
+            }
+
+            std::string body;
+            if (spec.contains("body")) {
+                const auto& value = spec["body"];
+                body = value.is_string() ? value.get<std::string>() : value.dump();
+            }
+
+            res.status = status_code;
+            res.set_content(body, content_type.c_str());
+            return;
+        }
+
+        if (parsed_ok) {
+            sendJson(res, parsed);
+        } else {
             res.set_content(grpc_res.result(), "text/plain");
         }
     } else {