|
|
@@ -219,8 +219,8 @@ void WebServerService::setupRoutes() {
|
|
|
execution_ctrl_ = std::make_unique<api::ExecutionController>(
|
|
|
*storage_, *auth_middleware_, *ws_server_, *scheduler_, *load_balancer_,
|
|
|
[this](const std::string& workflow_id, const std::string& execution_id,
|
|
|
- const std::string& error) {
|
|
|
- runErrorWorkflow(workflow_id, execution_id, error);
|
|
|
+ bool failed, const std::string& error) {
|
|
|
+ noteExecutionOutcome(workflow_id, execution_id, failed, error);
|
|
|
});
|
|
|
file_ctrl_ = std::make_unique<api::FileController>(*storage_, *auth_middleware_);
|
|
|
execution_ctrl_->registerRoutes(server);
|
|
|
@@ -350,6 +350,68 @@ void WebServerService::loadScheduledWorkflows() {
|
|
|
LOG_INFO("Loaded {} scheduled workflows", registered_count);
|
|
|
}
|
|
|
|
|
|
+void WebServerService::noteExecutionOutcome(const std::string& workflow_id,
|
|
|
+ const std::string& execution_id,
|
|
|
+ bool failed,
|
|
|
+ const std::string& error) {
|
|
|
+ if (failed) {
|
|
|
+ runErrorWorkflow(workflow_id, execution_id, error);
|
|
|
+ }
|
|
|
+
|
|
|
+ auto stored = storage_->get("workflows", workflow_id);
|
|
|
+ if (stored.failed()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ const auto workflow = stored.value();
|
|
|
+ const auto settings = workflow.value("settings", nlohmann::json::object());
|
|
|
+ const int limit = settings.value("deactivateAfterFailures", 0);
|
|
|
+ const int streak = workflow.value("consecutiveFailures", 0);
|
|
|
+
|
|
|
+ if (!failed) {
|
|
|
+ // Only written when there is something to clear, so an ordinary run does
|
|
|
+ // not cost a write.
|
|
|
+ if (streak != 0) {
|
|
|
+ storage_->update("workflows", workflow_id,
|
|
|
+ {{"consecutiveFailures", 0}}, 0, true);
|
|
|
+ }
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ const int next = streak + 1;
|
|
|
+ nlohmann::json patch = {{"consecutiveFailures", next}};
|
|
|
+
|
|
|
+ // Counting is worth doing even when nothing is switched off - it is the
|
|
|
+ // number someone looks at when asking how long this has been going wrong.
|
|
|
+ if (limit <= 0 || next < limit || workflow.value("active", false) != true) {
|
|
|
+ storage_->update("workflows", workflow_id, patch, 0, true);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ // Switched off, and said out loud. A workflow that simply appeared
|
|
|
+ // "Inactive" one morning with no reason recorded is worse than one that
|
|
|
+ // kept failing, because nobody can tell which of the two happened.
|
|
|
+ patch["active"] = false;
|
|
|
+ patch["deactivatedReason"] = "Switched off after " + std::to_string(next) +
|
|
|
+ " failures in a row. The last one: " +
|
|
|
+ (error.empty() ? "no reason given" : error);
|
|
|
+ patch["deactivatedAt"] = common::TimeUtils::nowMs();
|
|
|
+ patch["updatedAt"] = common::TimeUtils::nowMs();
|
|
|
+ storage_->update("workflows", workflow_id, patch, 0, true);
|
|
|
+
|
|
|
+ // Off the schedule too, or it would keep firing while reading as inactive.
|
|
|
+ scheduler_->unregisterWorkflow(workflow_id);
|
|
|
+
|
|
|
+ LOG_WARN("Workflow {} deactivated after {} consecutive failures: {}",
|
|
|
+ workflow_id, next, error);
|
|
|
+
|
|
|
+ ws_server_->broadcast("workflows.deactivated", {
|
|
|
+ {"id", workflow_id},
|
|
|
+ {"reason", patch["deactivatedReason"]},
|
|
|
+ {"consecutiveFailures", next}
|
|
|
+ });
|
|
|
+}
|
|
|
+
|
|
|
void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
|
|
|
const std::string& failed_execution_id,
|
|
|
const std::string& error_message) {
|