Selaa lähdekoodia

fix: raise the webserver's gRPC receive limit for runner responses

grpc::CreateChannel with no ChannelArguments left every webserver -> runner
channel at gRPC's 4 MB default for what it would accept back. The runner's
own server had already been raised to 64 MB (both send and receive) so a
large upload could get in, but nothing matched it on the client side that
reads the answer - so a wait-mode form whose respond-to-webhook body passed
4 MB failed with "Received message larger than max", even though the same
result would have fit through every other limit in the chain.

LoadBalancer::createChannel() is now the single place a webserver -> runner
channel gets built, sized from server.max_upload_mb (already the webserver's
one stated ceiling for request and response size) rather than a second
config knob that could drift out of step with it. Every existing call site -
node_controller, execution_controller (cancel and resume), workflow_controller,
webhook_controller, and the three in webserver_service.cpp - now goes through
it, so a limit raised once cannot be missed on a sibling the next time this
needs raising again.

Also completes the M5 cleanup noted in review: webhook_controller.cpp built a
stub before knowing whether the immediate-mode dispatch branch would use it;
it now only builds one where it is actually used.
fszontagh 1 kuukausi sitten
vanhempi
sitoutus
abda4ab403

+ 11 - 1
docs/nodes.md

@@ -358,7 +358,7 @@ failure rather than a silently dropped submission.
 
 ### Upload size
 
-Two limits apply, and a submission has to pass both:
+Two limits apply on the way in, and a submission has to pass both:
 
 - `server.max_upload_mb` in `config/webserver.json` bounds every request body
   the webserver accepts, uploads included - 32 MB by default. Anything larger
@@ -372,6 +372,16 @@ Two limits apply, and a submission has to pass both:
 A field's own `maxSizeMb` (see Fields, above) can tighten either of these
 further but never raise them.
 
+The way back is a separate, symmetric limit: every gRPC channel the webserver
+opens to a runner - including the one a `wait`-mode form's execution result
+travels back over - is sized off `server.max_upload_mb` as well (both send and
+receive), not off the runner's `max_message_size_mb`. Before this was wired
+up, that channel used gRPC's own 4 MB default for what it would accept back,
+so a `wait`-mode form whose `respond-to-webhook` body exceeded 4 MB failed
+with "Received message larger than max" even though the request that produced
+it was well inside every limit above. A large response is bounded by
+`server.max_upload_mb` the same as a large request is.
+
 ### Password
 
 `password` is optional; leave it empty and anyone with the link can submit the

+ 2 - 7
src/webserver/api/execution_controller.cpp

@@ -387,8 +387,7 @@ void ExecutionController::cancelExecution(const httplib::Request& req, httplib::
         return;
     }
 
-    // Qualified with the leading "::" for the same reason as the resume path.
-    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_.createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     proto::CancelExecutionRequest grpc_req;
@@ -517,11 +516,7 @@ void ExecutionController::resumeExecution(const httplib::Request& req, httplib::
         return;
     }
 
-    // Qualified with the leading "::" because smartbotic::webserver::grpc (a
-    // forward-declared namespace for the node sync / credential gRPC servers,
-    // pulled in via webserver_service.hpp) would otherwise shadow the real
-    // ::grpc namespace here.
-    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_.createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     proto::ResumeExecutionRequest grpc_req;

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

@@ -73,7 +73,7 @@ void NodeController::nodeOptions(const httplib::Request& req, httplib::Response&
     inline_workflow["connections"] = nlohmann::json::array();
     inline_workflow["settings"] = nlohmann::json::object();
 
-    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_.createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     proto::ExecuteWorkflowRequest grpc_req;

+ 7 - 3
src/webserver/api/webhook_controller.cpp

@@ -375,9 +375,12 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
     // Add client IP
     trigger_data["clientIp"] = req.remote_addr;
 
-    // Execute workflow via gRPC
-    auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
-    auto stub = proto::RunnerService::NewStub(channel);
+    // Execute workflow via gRPC. Only the channel is built here - the queued
+    // (!wait_for_run) branch below builds its own stub (bg_stub) off this
+    // same channel, and the synchronous wait branch further down builds
+    // its own right before using it, so a stub built unconditionally here
+    // would sit unused whenever wait_for_run is false.
+    auto channel = load_balancer_.createChannel(runner->address);
 
     // A schedule's overlap policy counts what is running, and a synchronous
     // dispatch blocks until the run has finished - so by the time it returns
@@ -503,6 +506,7 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
 
     scheduler_.notifyExecutionStarted(workflow_id, execution_id);
 
+    auto stub = proto::RunnerService::NewStub(channel);
     proto::ExecuteWorkflowResponse grpc_res;
     grpc::ClientContext grpc_ctx;
     grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(35));

+ 1 - 1
src/webserver/api/workflow_controller.cpp

@@ -725,7 +725,7 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
     } catch (...) {}
 
     // Call runner to execute workflow
-    auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_.createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     proto::ExecuteWorkflowRequest grpc_req;

+ 17 - 0
src/webserver/runners/load_balancer.cpp

@@ -1,6 +1,7 @@
 #include "load_balancer.hpp"
 #include <algorithm>
 #include <random>
+#include <grpcpp/grpcpp.h>
 
 namespace smartbotic::webserver::runners {
 
@@ -25,6 +26,22 @@ std::string loadBalancingStrategyToString(LoadBalancingStrategy strategy) {
 LoadBalancer::LoadBalancer(RunnerRegistry& registry, const LoadBalancerConfig& config)
     : registry_(registry), config_(config) {}
 
+std::shared_ptr<::grpc::Channel> LoadBalancer::createChannel(const std::string& address) const {
+    ::grpc::ChannelArguments args;
+    const int max_bytes = config_.max_message_size_mb * 1024 * 1024;
+    // Both directions: the request can carry a large upload (form file
+    // fields, arbitrary POST bodies) and the response can carry a large
+    // result (a wait-mode form's respond-to-webhook body, a big node
+    // output). gRPC's default is 4 MB for the receive side and unlimited for
+    // send, but leaving send unset here would mean this webserver could ask
+    // a runner for something the runner's own receive limit then rejects -
+    // pinning both to the same number keeps this channel's promise
+    // symmetric.
+    args.SetMaxReceiveMessageSize(max_bytes);
+    args.SetMaxSendMessageSize(max_bytes);
+    return ::grpc::CreateCustomChannel(address, ::grpc::InsecureChannelCredentials(), args);
+}
+
 std::optional<Runner> LoadBalancer::selectRunner(const std::string& required_node_type) {
     std::vector<Runner> candidates;
 

+ 25 - 0
src/webserver/runners/load_balancer.hpp

@@ -3,8 +3,13 @@
 #include <string>
 #include <optional>
 #include <atomic>
+#include <memory>
 #include "runner_registry.hpp"
 
+namespace grpc {
+class Channel;
+}
+
 namespace smartbotic::webserver::runners {
 
 // Load balancing strategy
@@ -22,6 +27,19 @@ std::string loadBalancingStrategyToString(LoadBalancingStrategy strategy);
 struct LoadBalancerConfig {
     LoadBalancingStrategy strategy = LoadBalancingStrategy::LeastConnections;
     double busy_threshold = 0.8;  // Runner considered busy above this load
+
+    // Ceiling for gRPC messages exchanged with a runner over a channel this
+    // class creates - both directions, request and response. Sourced from
+    // server.max_upload_mb (see WebServerService::loadConfig): that is
+    // already the webserver's single stated ceiling for how big a request or
+    // its answer is allowed to be, so reusing it here keeps one size promise
+    // across the HTTP layer and the gRPC hop behind it, instead of a second
+    // knob that could drift out of step with the first. Every webserver ->
+    // runner channel goes through createChannel() precisely so this one
+    // number is what all of them enforce - see the form-response bug this
+    // was written to fix, where only one caller had been widened and its
+    // neighbours quietly kept gRPC's 4 MB default.
+    int max_message_size_mb = 32;
 };
 
 // Load balancer - selects runner for workflow execution
@@ -43,6 +61,13 @@ public:
         return registry_.getRunner(runner_id);
     }
 
+    // Creates a gRPC channel to a runner with both message-size limits set
+    // from config_.max_message_size_mb, instead of gRPC's 4 MB default. Every
+    // webserver -> runner channel must be created through this, not
+    // grpc::CreateChannel directly, so a limit raised once cannot be missed
+    // on a sibling call site the next time this needs raising again.
+    std::shared_ptr<::grpc::Channel> createChannel(const std::string& address) const;
+
     // Get current strategy
     LoadBalancingStrategy getStrategy() const { return config_.strategy; }
 

+ 8 - 3
src/webserver/webserver_service.cpp

@@ -217,6 +217,11 @@ WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path&
         auto strategy_str = cfg.getOr<std::string>("runners.load_balancing", "least-connections");
         config.load_balancer_config.strategy =
             runners::loadBalancingStrategyFromString(strategy_str);
+        // Every webserver -> runner gRPC channel is sized off this, not a
+        // separate knob - see the comment on LoadBalancerConfig for why
+        // server.max_upload_mb (already loaded above) is the right source of
+        // truth rather than a second config key that could drift from it.
+        config.load_balancer_config.max_message_size_mb = config.max_upload_mb;
 
         // JWT config
         config.jwt_config.secret = cfg.getOr<std::string>("auth.jwt_secret",
@@ -625,7 +630,7 @@ void WebServerService::reconcileOrphanedExecutions(const std::string& runner_id,
     // runner knows which those are.
     std::unordered_set<std::string> still_running;
     {
-        auto channel = ::grpc::CreateChannel(address, ::grpc::InsecureChannelCredentials());
+        auto channel = load_balancer_->createChannel(address);
         auto stub = proto::RunnerService::NewStub(channel);
 
         proto::ListActiveExecutionsRequest request;
@@ -740,7 +745,7 @@ void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
     trigger_data["error"] = error_message;
     trigger_data["failedAt"] = common::TimeUtils::nowMs();
 
-    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_->createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     proto::ExecuteWorkflowRequest request;
@@ -779,7 +784,7 @@ void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
     }
 
     // Create gRPC channel and stub
-    auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
+    auto channel = load_balancer_->createChannel(runner->address);
     auto stub = proto::RunnerService::NewStub(channel);
 
     // Prepare request