|
|
@@ -15,6 +15,7 @@
|
|
|
#include "grpc/credential_service.hpp"
|
|
|
#include "credentials/credential_store.hpp"
|
|
|
#include "scheduler/workflow_scheduler.hpp"
|
|
|
+#include "common/time_utils.hpp"
|
|
|
#include "logging/logger.hpp"
|
|
|
#include <grpcpp/grpcpp.h>
|
|
|
#include "proto/runner.grpc.pb.h"
|
|
|
@@ -213,7 +214,12 @@ void WebServerService::setupRoutes() {
|
|
|
*storage_, *auth_middleware_, *ws_server_);
|
|
|
workflow_group_ctrl_->registerRoutes(server);
|
|
|
|
|
|
- execution_ctrl_ = std::make_unique<api::ExecutionController>(*storage_, *auth_middleware_, *ws_server_, *scheduler_);
|
|
|
+ execution_ctrl_ = std::make_unique<api::ExecutionController>(
|
|
|
+ *storage_, *auth_middleware_, *ws_server_, *scheduler_,
|
|
|
+ [this](const std::string& workflow_id, const std::string& execution_id,
|
|
|
+ const std::string& error) {
|
|
|
+ runErrorWorkflow(workflow_id, execution_id, error);
|
|
|
+ });
|
|
|
file_ctrl_ = std::make_unique<api::FileController>(*storage_, *auth_middleware_);
|
|
|
execution_ctrl_->registerRoutes(server);
|
|
|
file_ctrl_->registerRoutes(server);
|
|
|
@@ -331,6 +337,85 @@ void WebServerService::loadScheduledWorkflows() {
|
|
|
LOG_INFO("Loaded {} scheduled workflows", registered_count);
|
|
|
}
|
|
|
|
|
|
+void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
|
|
|
+ const std::string& failed_execution_id,
|
|
|
+ const std::string& error_message) {
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(handled_failures_mutex_);
|
|
|
+ if (!handled_failures_.insert(failed_execution_id).second) {
|
|
|
+ return; // already handled this failure
|
|
|
+ }
|
|
|
+ if (handled_failures_.size() > 512) {
|
|
|
+ handled_failures_.erase(handled_failures_.begin());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ auto failed = storage_->get("workflows", failed_workflow_id);
|
|
|
+ if (failed.failed()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ const auto settings = failed.value().value("settings", nlohmann::json::object());
|
|
|
+ const std::string handler_id = settings.value("errorWorkflowId", std::string());
|
|
|
+ if (handler_id.empty()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ // A handler that fails must not summon itself, which would run forever.
|
|
|
+ if (handler_id == failed_workflow_id) {
|
|
|
+ LOG_WARN("Workflow {} names itself as its error workflow; not running it",
|
|
|
+ failed_workflow_id);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto handler = storage_->get("workflows", handler_id);
|
|
|
+ if (handler.failed()) {
|
|
|
+ LOG_WARN("Workflow {} names error workflow {}, which no longer exists",
|
|
|
+ failed_workflow_id, handler_id);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto runner = load_balancer_->selectRunner();
|
|
|
+ if (!runner) {
|
|
|
+ LOG_ERROR("No runners available to run error workflow {}", handler_id);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ // The handler is told what failed rather than having to look it up, so it can
|
|
|
+ // notify or record without needing read access to the executions collection.
|
|
|
+ nlohmann::json trigger_data;
|
|
|
+ trigger_data["errorWorkflow"] = true;
|
|
|
+ trigger_data["failedWorkflowId"] = failed_workflow_id;
|
|
|
+ trigger_data["failedWorkflowName"] = failed.value().value("name", std::string());
|
|
|
+ trigger_data["failedExecutionId"] = failed_execution_id;
|
|
|
+ trigger_data["error"] = error_message;
|
|
|
+ trigger_data["failedAt"] = common::TimeUtils::nowMs();
|
|
|
+
|
|
|
+ auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
|
|
|
+ auto stub = proto::RunnerService::NewStub(channel);
|
|
|
+
|
|
|
+ proto::ExecuteWorkflowRequest request;
|
|
|
+ request.set_workflow_id(handler_id);
|
|
|
+ request.set_trigger_type("error-workflow");
|
|
|
+ request.set_trigger_data(trigger_data.dump());
|
|
|
+ request.set_wait_for_completion(false);
|
|
|
+
|
|
|
+ proto::ExecuteWorkflowResponse response;
|
|
|
+ ::grpc::ClientContext context;
|
|
|
+ context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
|
|
|
+
|
|
|
+ auto status = stub->ExecuteWorkflow(&context, request, &response);
|
|
|
+ if (!status.ok()) {
|
|
|
+ LOG_ERROR("Error workflow {} could not be started: {}", handler_id, status.error_message());
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ LOG_INFO("Error workflow {} started as {} after {} failed",
|
|
|
+ handler_id, response.execution_id(), failed_workflow_id);
|
|
|
+
|
|
|
+ scheduler_->notifyExecutionStarted(handler_id, response.execution_id());
|
|
|
+}
|
|
|
+
|
|
|
void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
|
|
|
const std::string& trigger_node_id,
|
|
|
const std::string& trigger_type) {
|