Преглед изворни кода

fix: reject a cached loop-body result that still carries a pause marker

The pause-rejection added for freshly executed loop body nodes only
covered that one code path. Single-node testing's cache branch built a
NodeExecutionResult straight from a pinned cache entry and stored it as
Completed without ever routing through the check, so a cache entry
recorded before this protection existed (any output cached before this
change shipped) would still replay with an orphaned _pause marker and
no error anywhere.

Factor the check into a small rejectPauseInLoopBody helper shared by
both the cache path and the fresh-execution path, so the failure and
its message cannot drift apart between them. Also thread the resulting
Failed status through the cache branch's continueOnError handling and
callback event, which previously assumed a cache hit could only ever
be Completed.
fszontagh пре 1 месец
родитељ
комит
2d00b322d3
1 измењених фајлова са 42 додато и 16 уклоњено
  1. 42 16
      src/runner/workflow_engine.cpp

+ 42 - 16
src/runner/workflow_engine.cpp

@@ -1715,6 +1715,23 @@ bool WorkflowEngine::executeLoopBody(
         LOG_INFO("  - Body node: {}", nid);
     }
 
+    // A completed body result - whether freshly executed or replayed from a
+    // pinned cache entry - that still carries a pause marker must not be
+    // replayed as Completed: the iteration state that would let it resume
+    // does not exist in this execution's stored data. Both the cache path and
+    // the fresh-execution path route through this single check so the
+    // failure and its message cannot drift apart between them.
+    auto rejectPauseInLoopBody = [](NodeExecutionResult& body_result) {
+        if (body_result.status == NodeStatus::Completed &&
+            body_result.output.contains("_pause")) {
+            body_result.status = NodeStatus::Failed;
+            body_result.error = "A node cannot pause inside a Loop body, because a "
+                                "paused loop iteration cannot be resumed. Collect the "
+                                "items first, approve once, then loop";
+            body_result.output.erase("_pause");
+        }
+    };
+
     // Execute body for each item
     for (size_t i = 0; i < ctx.items.size(); ++i) {
         ctx.current_index = i;
@@ -1873,10 +1890,13 @@ bool WorkflowEngine::executeLoopBody(
                     cached_result.finished_at = cached_result.started_at;
                     cached_result.from_cache = true;
 
+                    rejectPauseInLoopBody(cached_result);
+
                     iteration_results[body_node_id] = cached_result;
                     result.node_results[body_node_id + "_iter_" + std::to_string(i)] = cached_result;
 
-                    if (cached_result.output.contains("_webhookResponse")) {
+                    if (cached_result.status == NodeStatus::Completed &&
+                        cached_result.output.contains("_webhookResponse")) {
                         if (!result.webhook_response.is_null()) {
                             LOG_WARN("Node {} in a loop body overrides a webhook response already set",
                                      body_node_id);
@@ -1885,17 +1905,34 @@ bool WorkflowEngine::executeLoopBody(
                     }
 
                     if (callback) {
-                        callback("loop.node.completed", {
+                        nlohmann::json cache_event_data = {
                             {"executionId", result.execution_id},
                             {"nodeId", body_node_id},
                             {"iteration", i},
-                            {"status", "completed"},
+                            {"status", nodeStatusToString(cached_result.status)},
                             {"output", cached_result.output},
                             {"fromCache", true}
-                        });
+                        };
+                        if (!cached_result.error.empty()) {
+                            cache_event_data["error"] = cached_result.error;
+                        }
+                        callback("loop.node." + nodeStatusToString(cached_result.status), cache_event_data);
                     }
 
                     LOG_INFO("Using pinned output for body node {}", body_node_id);
+
+                    // A pinned output rejected above by rejectPauseInLoopBody must go
+                    // through the same continueOnError decision a freshly executed
+                    // failure does, rather than silently moving on to the next body
+                    // node as a plain cache hit would.
+                    if (cached_result.status == NodeStatus::Failed) {
+                        iteration_failed = true;
+                        all_succeeded = false;
+                        if (!ctx.continue_on_error) {
+                            break;
+                        }
+                    }
+
                     continue;
                 }
 
@@ -1939,18 +1976,7 @@ bool WorkflowEngine::executeLoopBody(
                 body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow);
             }
 
-            // A pause inside a loop body cannot be resumed: the iteration state that
-            // would have to be restored lives in locals this execution never stores.
-            // Failing the node is the only honest outcome - letting the marker through
-            // would strip it and continue, orphaning an approval nobody can answer.
-            if (body_result.status == NodeStatus::Completed &&
-                body_result.output.contains("_pause")) {
-                body_result.status = NodeStatus::Failed;
-                body_result.error = "A node cannot pause inside a Loop body, because a "
-                                    "paused loop iteration cannot be resumed. Collect the "
-                                    "items first, approve once, then loop";
-                body_result.output.erase("_pause");
-            }
+            rejectPauseInLoopBody(body_result);
 
             iteration_results[body_node_id] = body_result;