|
|
@@ -240,6 +240,10 @@ WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path&
|
|
|
config.form_dispatch_threads = cfg.getOr<int>("server.form_dispatch_threads", 4);
|
|
|
config.form_dispatch_queue_capacity =
|
|
|
cfg.getOr<int>("server.form_dispatch_queue_capacity", 32);
|
|
|
+ config.workflow_dispatch_threads =
|
|
|
+ cfg.getOr<int>("server.workflow_dispatch_threads", 4);
|
|
|
+ config.workflow_dispatch_queue_capacity =
|
|
|
+ cfg.getOr<int>("server.workflow_dispatch_queue_capacity", 32);
|
|
|
|
|
|
// Runner config
|
|
|
config.runner_config.heartbeat_timeout_sec =
|
|
|
@@ -371,12 +375,16 @@ void WebServerService::stop() {
|
|
|
LOG_INFO("Stopping WebServer service...");
|
|
|
|
|
|
http_server_->stop();
|
|
|
- // No new requests can arrive once the HTTP server is stopped, but the
|
|
|
- // dispatch pool's own workers may still be mid-run - stop() joins them,
|
|
|
- // so nothing is left running against a controller about to be destroyed.
|
|
|
+ // No new requests can arrive once the HTTP server is stopped, but each
|
|
|
+ // controller's own dispatch pool workers may still be mid-run - stop()
|
|
|
+ // joins them, so nothing is left running against a controller about to
|
|
|
+ // be destroyed.
|
|
|
if (webhook_ctrl_) {
|
|
|
webhook_ctrl_->stop();
|
|
|
}
|
|
|
+ if (workflow_ctrl_) {
|
|
|
+ workflow_ctrl_->stop();
|
|
|
+ }
|
|
|
ws_server_->stop();
|
|
|
credential_server_->stop();
|
|
|
node_sync_server_->stop();
|
|
|
@@ -423,9 +431,15 @@ void WebServerService::setupRoutes() {
|
|
|
*storage_, *auth_middleware_, *access_, *auth_store_, *retention_);
|
|
|
project_ctrl_->registerRoutes(server);
|
|
|
|
|
|
+ api::WorkflowController::DispatchConfig workflow_dispatch_config;
|
|
|
+ workflow_dispatch_config.threads = clampDispatchSetting(
|
|
|
+ "server.workflow_dispatch_threads", config_.workflow_dispatch_threads, 1, 64);
|
|
|
+ workflow_dispatch_config.queue_capacity = clampDispatchSetting(
|
|
|
+ "server.workflow_dispatch_queue_capacity", config_.workflow_dispatch_queue_capacity, 1, 10000);
|
|
|
+ workflow_dispatch_config.label_kind = "editor execute dispatch";
|
|
|
workflow_ctrl_ = std::make_unique<api::WorkflowController>(
|
|
|
*storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
|
|
|
- *scheduler_, *db_watcher_, *access_, *node_store_, *retention_);
|
|
|
+ *scheduler_, *db_watcher_, *access_, *node_store_, *retention_, workflow_dispatch_config);
|
|
|
workflow_ctrl_->registerRoutes(server);
|
|
|
|
|
|
workflow_group_ctrl_ = std::make_unique<api::WorkflowGroupController>(
|
|
|
@@ -502,10 +516,18 @@ void WebServerService::ensureExecutionsSummaryView() {
|
|
|
// A stopped run has status "completed" and an empty error, so without this
|
|
|
// field the list cannot tell a run that finished its work from one that
|
|
|
// deliberately ended early.
|
|
|
+ //
|
|
|
+ // toleratedErrorCount is the same problem one level down: a node can
|
|
|
+ // swallow its own failure (e.g. skipOnError) so the run keeps going, and
|
|
|
+ // the whole execution still ends up status "completed" with an empty
|
|
|
+ // top-level error. Without this field a run that quietly tolerated a
|
|
|
+ // failed feed is indistinguishable, in the list, from a run where nothing
|
|
|
+ // went wrong.
|
|
|
auto result = storage_->createView(kExecutionsSummaryView, "executions",
|
|
|
{"workflowId", "workflowName", "status", "triggerType",
|
|
|
"startedAt", "finishedAt", "error", "runnerId",
|
|
|
- "stopped", "stopReason", "stoppedNodeId"});
|
|
|
+ "stopped", "stopReason", "stoppedNodeId",
|
|
|
+ "toleratedErrorCount"});
|
|
|
if (result.failed()) {
|
|
|
LOG_ERROR("Could not create the executions summary view: {}. The executions "
|
|
|
"listing will fail until this is resolved.", result.error().message());
|