|
@@ -697,6 +697,30 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
result.webhook_response = node_result.output["_webhookResponse"];
|
|
result.webhook_response = node_result.output["_webhookResponse"];
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A node asking to stop ends the walk without failing the run. The
|
|
|
|
|
+ // distinction matters for anything that polls: a scheduled workflow
|
|
|
|
|
+ // that finds nothing to do has completed successfully, and marking
|
|
|
|
|
+ // it failed every five minutes buries real failures in noise.
|
|
|
|
|
+ if (node_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ node_result.output.contains("_stop")) {
|
|
|
|
|
+ const auto& stop = node_result.output["_stop"];
|
|
|
|
|
+
|
|
|
|
|
+ result.stop_requested = true;
|
|
|
|
|
+ result.stopped_node_id = node_id;
|
|
|
|
|
+ result.stop_reason = stop.value("reason", "");
|
|
|
|
|
+
|
|
|
|
|
+ // The marker is engine plumbing. What the node reports is that
|
|
|
|
|
+ // it stopped and why, not the mechanism that carried it.
|
|
|
|
|
+ auto& stored_node = result.node_results[node_id];
|
|
|
|
|
+ stored_node.output.erase("_stop");
|
|
|
|
|
+ stored_node.output["stopped"] = true;
|
|
|
|
|
+ stored_node.output["reason"] = result.stop_reason;
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Execution {} stopped at node {}: {}", result.execution_id, node_id,
|
|
|
|
|
+ result.stop_reason.empty() ? "no reason given" : result.stop_reason);
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// A node asking to pause ends this pass. The execution is stored as
|
|
// A node asking to pause ends this pass. The execution is stored as
|
|
|
// Waiting with everything computed so far, and a later resume picks
|
|
// Waiting with everything computed so far, and a later resume picks
|
|
|
// it up from here rather than starting again.
|
|
// it up from here rather than starting again.
|
|
@@ -773,6 +797,13 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
break;
|
|
break;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A body node asked to stop the run. The loop already unwound
|
|
|
|
|
+ // its own iterations; this ends the outer walk too, so nodes
|
|
|
|
|
+ // after the loop do not run.
|
|
|
|
|
+ if (result.stop_requested) {
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
continue; // Loop handles its own downstream execution
|
|
continue; // Loop handles its own downstream execution
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -791,8 +822,16 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
if (result.status == ExecutionStatus::Running) {
|
|
if (result.status == ExecutionStatus::Running) {
|
|
|
result.status = ExecutionStatus::Completed;
|
|
result.status = ExecutionStatus::Completed;
|
|
|
|
|
|
|
|
- // Get output from last executed node
|
|
|
|
|
- if (!execution_order.empty()) {
|
|
|
|
|
|
|
+ if (result.stop_requested) {
|
|
|
|
|
+ // The last node in the order never ran - the walk ended early on
|
|
|
|
|
+ // purpose - so the useful final output is the node that stopped
|
|
|
|
|
+ // it, not an entry that will not be found.
|
|
|
|
|
+ auto it = result.node_results.find(result.stopped_node_id);
|
|
|
|
|
+ if (it != result.node_results.end()) {
|
|
|
|
|
+ result.final_output = it->second.output;
|
|
|
|
|
+ }
|
|
|
|
|
+ } else if (!execution_order.empty()) {
|
|
|
|
|
+ // Get output from last executed node
|
|
|
auto it = result.node_results.find(execution_order.back());
|
|
auto it = result.node_results.find(execution_order.back());
|
|
|
if (it != result.node_results.end()) {
|
|
if (it != result.node_results.end()) {
|
|
|
result.final_output = it->second.output;
|
|
result.final_output = it->second.output;
|
|
@@ -1821,6 +1860,10 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
}
|
|
}
|
|
|
};
|
|
};
|
|
|
|
|
|
|
|
|
|
+ // Set when a body node asks to stop the whole run, so both the body-node
|
|
|
|
|
+ // loop and the per-item loop can unwind.
|
|
|
|
|
+ bool stopped_in_body = false;
|
|
|
|
|
+
|
|
|
// Execute body for each item
|
|
// Execute body for each item
|
|
|
for (size_t i = 0; i < ctx.items.size(); ++i) {
|
|
for (size_t i = 0; i < ctx.items.size(); ++i) {
|
|
|
ctx.current_index = i;
|
|
ctx.current_index = i;
|
|
@@ -2092,6 +2135,30 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
result.webhook_response = body_result.output["_webhookResponse"];
|
|
result.webhook_response = body_result.output["_webhookResponse"];
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A stop raised inside a loop body ends the whole run, not just the
|
|
|
|
|
+ // iteration - "stop the workflow" would be a strange thing to mean
|
|
|
|
|
+ // per-item. The flag travels out on the result because this walk
|
|
|
|
|
+ // cannot end the outer one itself.
|
|
|
|
|
+ if (body_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ body_result.output.contains("_stop")) {
|
|
|
|
|
+ const auto& stop = body_result.output["_stop"];
|
|
|
|
|
+
|
|
|
|
|
+ result.stop_requested = true;
|
|
|
|
|
+ result.stopped_node_id = body_node_id;
|
|
|
|
|
+ result.stop_reason = stop.value("reason", "");
|
|
|
|
|
+
|
|
|
|
|
+ auto& stored_body = result.node_results[result_key];
|
|
|
|
|
+ stored_body.output.erase("_stop");
|
|
|
|
|
+ stored_body.output["stopped"] = true;
|
|
|
|
|
+ stored_body.output["reason"] = result.stop_reason;
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Execution {} stopped at node {} during loop iteration {}: {}",
|
|
|
|
|
+ result.execution_id, body_node_id, i,
|
|
|
|
|
+ result.stop_reason.empty() ? "no reason given" : result.stop_reason);
|
|
|
|
|
+ stopped_in_body = true;
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
if (callback) {
|
|
if (callback) {
|
|
|
nlohmann::json event_data = {
|
|
nlohmann::json event_data = {
|
|
|
{"executionId", result.execution_id},
|
|
{"executionId", result.execution_id},
|
|
@@ -2139,6 +2206,12 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
if (iteration_failed && !ctx.continue_on_error) {
|
|
if (iteration_failed && !ctx.continue_on_error) {
|
|
|
break;
|
|
break;
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+ // A stop ends the remaining items too. Results collected so far are
|
|
|
|
|
+ // kept - they were produced before anyone asked to stop.
|
|
|
|
|
+ if (stopped_in_body) {
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// Copy results back (ctx was passed by const ref, but we used a mutable copy)
|
|
// Copy results back (ctx was passed by const ref, but we used a mutable copy)
|