|
|
@@ -19,7 +19,8 @@ WorkflowController::WorkflowController(storage::StorageClient& storage,
|
|
|
DatabaseWatcher& db_watcher,
|
|
|
auth::AccessControl& access,
|
|
|
nodes::NodeStore& node_store,
|
|
|
- retention::RetentionService& retention)
|
|
|
+ retention::RetentionService& retention,
|
|
|
+ DispatchConfig dispatch_config)
|
|
|
: storage_(storage)
|
|
|
, middleware_(middleware)
|
|
|
, registry_(registry)
|
|
|
@@ -29,7 +30,12 @@ WorkflowController::WorkflowController(storage::StorageClient& storage,
|
|
|
, db_watcher_(db_watcher)
|
|
|
, access_(access)
|
|
|
, node_store_(node_store)
|
|
|
- , retention_(retention) {}
|
|
|
+ , retention_(retention)
|
|
|
+ , dispatch_pool_(std::move(dispatch_config)) {}
|
|
|
+
|
|
|
+void WorkflowController::stop() {
|
|
|
+ dispatch_pool_.stop();
|
|
|
+}
|
|
|
|
|
|
|
|
|
namespace {
|
|
|
@@ -730,9 +736,33 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
|
|
|
}
|
|
|
} catch (...) {}
|
|
|
|
|
|
- // Call runner to execute workflow
|
|
|
+ // The runner's ExecuteWorkflow RPC is a plain synchronous call: it runs
|
|
|
+ // the workflow to completion before returning anything at all -
|
|
|
+ // wait_for_completion only controls whether the response carries a
|
|
|
+ // result, it does not make the server return early (see the comment in
|
|
|
+ // webhook_controller.cpp's immediate-mode branch, which hit this same
|
|
|
+ // fact first). Holding this handler's thread on that call would hold the
|
|
|
+ // editor's "Execute" button for however long the run takes, and the
|
|
|
+ // editor's axios client times out at 30 seconds - a fast build for most
|
|
|
+ // workflows, but a lie against any run that legitimately takes longer,
|
|
|
+ // and a workflow taking an hour is legitimate. So the call runs on the
|
|
|
+ // controller's own bounded dispatch pool instead of this thread, exactly
|
|
|
+ // like the webhook path's immediate-mode form dispatch: nothing is left
|
|
|
+ // running against `this` after a restart (dispatch_pool_.stop() joins
|
|
|
+ // every worker before this controller is destroyed - see
|
|
|
+ // WorkflowController::stop()), and a burst of Execute presses is
|
|
|
+ // throttled rather than spawning threads without limit.
|
|
|
+ //
|
|
|
+ // The scheduler slot is claimed inside the queued task, not here, so a
|
|
|
+ // request that is refused for a full queue never claims one. Once
|
|
|
+ // claimed, it is released the same way any other dispatch's is: the
|
|
|
+ // runner reports execution.completed/failed/cancelled/waiting to
|
|
|
+ // /api/v1/internal/execution-event, and execution_controller.cpp
|
|
|
+ // releases it there.
|
|
|
+ const std::string execution_id = UUID::generatePrefixed("exec");
|
|
|
+ trigger_data["_assignedExecutionId"] = execution_id;
|
|
|
+
|
|
|
auto channel = load_balancer_.createChannel(runner->address);
|
|
|
- auto stub = proto::RunnerService::NewStub(channel);
|
|
|
|
|
|
proto::ExecuteWorkflowRequest grpc_req;
|
|
|
grpc_req.set_workflow_id(workflow_id);
|
|
|
@@ -743,40 +773,69 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
|
|
|
grpc_req.set_trigger_data(trigger_data.dump());
|
|
|
grpc_req.set_wait_for_completion(false);
|
|
|
|
|
|
- proto::ExecuteWorkflowResponse grpc_res;
|
|
|
- grpc::ClientContext grpc_ctx;
|
|
|
- grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
|
|
|
+ const std::string runner_id = runner->id;
|
|
|
+ bool accepted = dispatch_pool_.tryEnqueue(workflow_id,
|
|
|
+ [this, channel, workflow_id, execution_id, runner_id, grpc_req]() {
|
|
|
+ scheduler_.notifyExecutionStarted(workflow_id, execution_id);
|
|
|
+
|
|
|
+ auto bg_stub = proto::RunnerService::NewStub(channel);
|
|
|
+ proto::ExecuteWorkflowResponse bg_res;
|
|
|
+ grpc::ClientContext bg_ctx;
|
|
|
+ // Nobody is waiting on this response under normal operation, so the
|
|
|
+ // deadline only needs to bound a runner that never answers at all.
|
|
|
+ // At shutdown, stop() calls TryCancel() on this context directly
|
|
|
+ // instead of waiting for the deadline.
|
|
|
+ bg_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::hours(1));
|
|
|
+
|
|
|
+ dispatch_pool_.registerActiveTask(&bg_ctx, workflow_id);
|
|
|
+ auto bg_status = bg_stub->ExecuteWorkflow(&bg_ctx, grpc_req, &bg_res);
|
|
|
+ dispatch_pool_.unregisterActiveTask(&bg_ctx);
|
|
|
+
|
|
|
+ if (!bg_status.ok()) {
|
|
|
+ // Neither a lapsed deadline nor a shutdown-time TryCancel means
|
|
|
+ // the run stopped - both only end this process's own wait. The
|
|
|
+ // runner's ExecuteWorkflow handler never checks whether its
|
|
|
+ // caller cancelled, so the workflow keeps running server-side
|
|
|
+ // regardless and will still report its own
|
|
|
+ // execution.completed/failed/waiting event when it finishes -
|
|
|
+ // releasing the slot here would say idle while it is still
|
|
|
+ // working.
|
|
|
+ if (bg_status.error_code() != grpc::StatusCode::DEADLINE_EXCEEDED &&
|
|
|
+ bg_status.error_code() != grpc::StatusCode::CANCELLED) {
|
|
|
+ scheduler_.notifyExecutionFinished(workflow_id, execution_id);
|
|
|
+ }
|
|
|
+ LOG_ERROR("Editor execution failed for workflow {}: {}",
|
|
|
+ workflow_id, bg_status.error_message());
|
|
|
+ return;
|
|
|
+ }
|
|
|
|
|
|
- auto status = stub->ExecuteWorkflow(&grpc_ctx, grpc_req, &grpc_res);
|
|
|
+ if (bg_res.execution_id() != execution_id) {
|
|
|
+ LOG_WARN("Editor execution for workflow {} came back as {} but was dispatched as {} - "
|
|
|
+ "moving the schedule slot",
|
|
|
+ workflow_id, bg_res.execution_id(), execution_id);
|
|
|
+ scheduler_.notifyExecutionFinished(workflow_id, execution_id);
|
|
|
+ scheduler_.notifyExecutionStarted(workflow_id, bg_res.execution_id());
|
|
|
+ }
|
|
|
+ });
|
|
|
|
|
|
- if (!status.ok()) {
|
|
|
- LOG_ERROR("Failed to execute workflow: {}", status.error_message());
|
|
|
- sendError(res, "Failed to execute workflow: " + status.error_message(), 500);
|
|
|
+ if (!accepted) {
|
|
|
+ // The queue is full: refuse honestly rather than tell the editor
|
|
|
+ // "started" for a run that will never happen. No slot was claimed
|
|
|
+ // for this one.
|
|
|
+ LOG_WARN("Editor execute dispatch queue is full, refusing execution for workflow {}",
|
|
|
+ workflow_id);
|
|
|
+ sendError(res, "Too many executions in progress right now. Try again shortly.", 503);
|
|
|
return;
|
|
|
}
|
|
|
|
|
|
- // The scheduler counts what is running so an overlap policy of skip or
|
|
|
- // queue can hold a tick back. It only ever heard about runs it started
|
|
|
- // itself, so a run started from the editor was invisible: a workflow set to
|
|
|
- // never overlap would happily get a second scheduled run on top of the one
|
|
|
- // someone was watching.
|
|
|
- scheduler_.notifyExecutionStarted(workflow_id, grpc_res.execution_id());
|
|
|
-
|
|
|
nlohmann::json response;
|
|
|
- response["executionId"] = grpc_res.execution_id();
|
|
|
- response["status"] = grpc_res.status();
|
|
|
- response["runnerId"] = runner->id;
|
|
|
-
|
|
|
- // Broadcast execution started
|
|
|
- ws_server_.broadcast("executions." + grpc_res.execution_id() + ".started", {
|
|
|
- {"executionId", grpc_res.execution_id()},
|
|
|
- {"workflowId", workflow_id},
|
|
|
- {"runnerId", runner->id}
|
|
|
- });
|
|
|
+ response["executionId"] = execution_id;
|
|
|
+ response["status"] = "running";
|
|
|
+ response["runnerId"] = runner_id;
|
|
|
|
|
|
- LOG_INFO("Workflow {} execution started: {} on runner {}",
|
|
|
- workflow_id, grpc_res.execution_id(), runner->id);
|
|
|
- sendJson(res, response, 202);
|
|
|
+ LOG_INFO("Workflow {} execution dispatched: {} on runner {}",
|
|
|
+ workflow_id, execution_id, runner_id);
|
|
|
+ sendJson(res, response, 200);
|
|
|
}
|
|
|
|
|
|
|