|
@@ -828,6 +828,21 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
|
|
|
|
|
// Check for failure
|
|
// Check for failure
|
|
|
if (node_result.status == NodeStatus::Failed) {
|
|
if (node_result.status == NodeStatus::Failed) {
|
|
|
|
|
+ // A node interrupted because someone pressed Stop did not fail -
|
|
|
|
|
+ // it was stopped. Recording that as a failure would put a red
|
|
|
|
|
+ // mark on the list for something the user asked for, and would
|
|
|
|
|
+ // set off the error workflow for it.
|
|
|
|
|
+ bool was_cancelled = false;
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
|
+ was_cancelled = cancelled_executions_.contains(result.execution_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (was_cancelled) {
|
|
|
|
|
+ result.status = ExecutionStatus::Cancelled;
|
|
|
|
|
+ result.error = "Execution cancelled";
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
bool should_continue = workflow.settings.value("continueOnError", false);
|
|
bool should_continue = workflow.settings.value("continueOnError", false);
|
|
|
if (!should_continue) {
|
|
if (!should_continue) {
|
|
|
result.status = ExecutionStatus::Failed;
|
|
result.status = ExecutionStatus::Failed;
|
|
@@ -1074,8 +1089,23 @@ void WorkflowEngine::cancelExecution(const std::string& execution_id) {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- LOG_INFO("Cancellation requested for execution {} ({} execution(s) marked)",
|
|
|
|
|
- execution_id, cancelled);
|
|
|
|
|
|
|
+ // Interrupt whatever is running right now. Without this the flag is only
|
|
|
|
|
+ // noticed between nodes, and a node that takes minutes keeps the whole run
|
|
|
|
|
+ // going for minutes after Stop was pressed.
|
|
|
|
|
+ size_t interrupted = 0;
|
|
|
|
|
+ for (const auto& id : cancelled_executions_) {
|
|
|
|
|
+ auto engines = running_engines_.find(id);
|
|
|
|
|
+ if (engines == running_engines_.end()) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+ for (auto* running : engines->second) {
|
|
|
|
|
+ running->cancel();
|
|
|
|
|
+ ++interrupted;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Cancellation requested for execution {} ({} execution(s) marked, {} interrupted)",
|
|
|
|
|
+ execution_id, cancelled, interrupted);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
int WorkflowEngine::getActiveExecutionCount() const {
|
|
int WorkflowEngine::getActiveExecutionCount() const {
|
|
@@ -1534,7 +1564,42 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
result.retry_count = attempt;
|
|
result.retry_count = attempt;
|
|
|
|
|
|
|
|
auto* script_engine = script_pool_->acquire();
|
|
auto* script_engine = script_pool_->acquire();
|
|
|
|
|
+
|
|
|
|
|
+ // Registered before it runs, so a Stop arriving mid-node reaches the
|
|
|
|
|
+ // engine actually doing the work. A cancel that landed between the
|
|
|
|
|
+ // acquire and here would otherwise be missed entirely, so the flag is
|
|
|
|
|
+ // checked once more under the same lock.
|
|
|
|
|
+ bool already_cancelled = false;
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
|
+ if (cancelled_executions_.contains(execution_id)) {
|
|
|
|
|
+ already_cancelled = true;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ running_engines_[execution_id].insert(script_engine);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (already_cancelled) {
|
|
|
|
|
+ script_pool_->release(script_engine);
|
|
|
|
|
+ result.status = NodeStatus::Failed;
|
|
|
|
|
+ result.error = "Execution cancelled";
|
|
|
|
|
+ result.finished_at = TimeUtils::nowMs();
|
|
|
|
|
+ return result;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
auto script_result = script_engine->execute(node_def->code, ctx);
|
|
auto script_result = script_engine->execute(node_def->code, ctx);
|
|
|
|
|
+
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
|
+ auto it = running_engines_.find(execution_id);
|
|
|
|
|
+ if (it != running_engines_.end()) {
|
|
|
|
|
+ it->second.erase(script_engine);
|
|
|
|
|
+ if (it->second.empty()) {
|
|
|
|
|
+ running_engines_.erase(it);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
script_pool_->release(script_engine);
|
|
script_pool_->release(script_engine);
|
|
|
|
|
|
|
|
if (script_result.success) {
|
|
if (script_result.success) {
|
|
@@ -2656,7 +2721,17 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
if (body_result.status == NodeStatus::Failed) {
|
|
if (body_result.status == NodeStatus::Failed) {
|
|
|
iteration_failed = true;
|
|
iteration_failed = true;
|
|
|
all_succeeded = false;
|
|
all_succeeded = false;
|
|
|
- if (!ctx.continue_on_error) {
|
|
|
|
|
|
|
+
|
|
|
|
|
+ // A node interrupted by Stop is not a failure to carry on past.
|
|
|
|
|
+ // Continue-on-error would otherwise march through every
|
|
|
|
|
+ // remaining body node, failing each one the same way, before
|
|
|
|
|
+ // the between-iteration check finally noticed.
|
|
|
|
|
+ bool was_cancelled = false;
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
|
+ was_cancelled = cancelled_executions_.contains(result.execution_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (was_cancelled || !ctx.continue_on_error) {
|
|
|
break;
|
|
break;
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|