فهرست منبع

refactor: one place decides what a node's markers mean

A node asks the engine for things through its output - run another workflow,
name the workflow's result, set a webhook response, stop, pause. There are
two walks over nodes, the main one and the loop body, and each carried its
own copy of that handling. Seven-plus bugs on this branch have been one walk
gaining something the other did not, every one caught by review rather than
by a test, because a fixture that exercises a marker outside a loop says
nothing about the same marker inside one.

applyNodeMarkers is now the only place that reads them, and both walks call
it - along with the two cached-result paths, which had partial copies of
their own. A marker added there reaches both walks or neither.

It exposed a live bug rather than only removing a risk: a Workflow Output
node inside a loop body was ignored entirely. The loop had no handling for
_workflowOutput, so the marker stayed in the output and the caller received
the raw {"_workflowOutput": ...} instead of the named result. Measured
against the previous binary: the new fixture fails there with the marker
visible in the caller's output, and passes here. Two smaller ones went the
same way - a cached result no longer skips a stop the original run made, and
a stop reached through the loop's pinned-output path now ends the run.

The named output moved onto ExecutionResult, because it was a local of the
main walk and a loop body had nowhere to put one. That was the mechanical
reason the bug existed.

62/62, one new fixture: a sub-workflow that names its result from inside a
loop, called by another workflow, which is where "the last node" is ambiguous
and a named result matters most.
fszontagh 1 ماه پیش
والد
کامیت
08a03db482
3فایلهای تغییر یافته به همراه238 افزوده شده و 167 حذف شده
  1. 168 167
      src/runner/workflow_engine.cpp
  2. 30 0
      src/runner/workflow_engine.hpp
  3. 40 0
      tests/nodes/workflow-output-in-loop.json

+ 168 - 167
src/runner/workflow_engine.cpp

@@ -396,10 +396,6 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
              result.execution_id, workflow.id);
 
     try {
-        // A Workflow Output node's value, when the workflow named one.
-        nlohmann::json explicit_output;
-        bool has_explicit_output = false;
-
         // Build execution graph
         std::unordered_map<std::string, std::vector<std::string>> dependencies;
         std::unordered_map<std::string, std::vector<std::string>> dependents;
@@ -679,6 +675,11 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
             }
 
             NodeExecutionResult node_result;
+            // What this node's markers asked for. Set in both branches below,
+            // because a cached output carries the same markers a fresh one does
+            // - a run replayed from cache has to stop where the original did.
+            MarkerOutcome marker_outcome = MarkerOutcome::Continue;
+
             if (use_cached) {
                 // Use cached output instead of executing
                 node_result.node_id = node_id;
@@ -689,7 +690,8 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                 node_result.finished_at = node_result.started_at;
                 node_result.from_cache = true;
 
-                result.node_results[node_id] = node_result;
+                marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
+                                                  nlohmann::json::object(), {}, -1);
 
                 if (callback) {
                     callback("node.completed", {
@@ -744,26 +746,8 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                                                node_def_for_defaults);
                 }
 
-                // A node asking to run another workflow has it run here: the
-                // node itself has no way to reach the engine, and this is the
-                // point where the result can still become its output.
-                if (node_result.status == NodeStatus::Completed &&
-                    node_result.output.contains("_callWorkflow")) {
-                    runSubWorkflow(node_result, result.call_depth, result.execution_id);
-                }
-
-                // Record what was applied. Once a node's settings can come from
-                // elsewhere, its stored config no longer says what it ran with,
-                // and this is the only place that closes that gap.
-                if (!overlay.empty() && node_result.status == NodeStatus::Completed) {
-                    node_result.output["_appliedConfig"] = overlay;
-                    if (!overlay_ignored.empty()) {
-                        node_result.output["_ignoredConfigKeys"] = overlay_ignored;
-                        LOG_WARN("Node {} was given {} setting(s) it does not have",
-                                 node_id, overlay_ignored.size());
-                    }
-                }
-                result.node_results[node_id] = node_result;
+                marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
+                                                  overlay, overlay_ignored, -1);
 
                 if (callback) {
                     nlohmann::json event_data = {
@@ -779,79 +763,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                 }
             }
 
-            // Any node may declare the HTTP response, not just the last one to
-            // run, so appending a node to a workflow cannot silently change
-            // what its webhook returns.
-            if (node_result.status == NodeStatus::Completed &&
-                node_result.output.contains("_webhookResponse")) {
-                if (!result.webhook_response.is_null()) {
-                    LOG_WARN("Node {} overrides a webhook response already set by an earlier node",
-                             node_id);
-                }
-                result.webhook_response = node_result.output["_webhookResponse"];
-            }
-
-            // A workflow can name its own result instead of leaving it to be
-            // whatever the last node happened to produce. It matters most when
-            // the workflow is called by another one, where "the last node" is
-            // ambiguous in anything that branches.
-            if (node_result.status == NodeStatus::Completed &&
-                node_result.output.contains("_workflowOutput")) {
-                explicit_output = node_result.output["_workflowOutput"];
-                has_explicit_output = true;
-                auto& stored_node = result.node_results[node_id];
-                stored_node.output.erase("_workflowOutput");
-                stored_node.output["output"] = explicit_output;
-            }
-
-            // 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
-            // Waiting with everything computed so far, and a later resume picks
-            // it up from here rather than starting again.
-            if (node_result.status == NodeStatus::Completed &&
-                node_result.output.contains("_pause")) {
-                const auto& pause = node_result.output["_pause"];
-
-                result.status = ExecutionStatus::Waiting;
-                result.paused_node_id = node_id;
-                result.pause_token = pause.value("token", "");
-                result.pause_expires_at = pause.value("expiresAt", static_cast<int64_t>(0));
-
-                // The marker is engine plumbing. What is stored is the request a
-                // person has to answer, not the mechanism that carried it.
-                auto& stored_node = result.node_results[node_id];
-                stored_node.output.erase("_pause");
-
-                // No callback here - the generic terminal callback below (after
-                // storeExecution) already emits "execution.waiting" once this
-                // pass ends, and it carries nodeId/expiresAt too. Emitting here
-                // as well produced the same event twice with two different
-                // payload shapes.
-
-                LOG_INFO("Execution {} paused at node {}", result.execution_id, node_id);
+            if (marker_outcome == MarkerOutcome::Stop || marker_outcome == MarkerOutcome::Pause) {
                 break;
             }
 
@@ -929,10 +841,10 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
         if (result.status == ExecutionStatus::Running) {
             result.status = ExecutionStatus::Completed;
 
-            if (has_explicit_output) {
+            if (result.has_explicit_output) {
                 // A Workflow Output node said what this workflow returns, which
                 // beats guessing from whichever node ran last.
-                result.final_output = explicit_output;
+                result.final_output = result.explicit_output;
             } else 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
@@ -1697,6 +1609,142 @@ static std::vector<std::string> applyConfigOverlay(nlohmann::json& config,
 // the failure would look like a hang rather than a mistake in the workflow.
 static constexpr int kMaxCallDepth = 5;
 
+// See the header for why this is one function rather than a block in each walk.
+// iteration is -1 outside a loop body, and the item index inside one.
+WorkflowEngine::MarkerOutcome WorkflowEngine::applyNodeMarkers(
+    NodeExecutionResult& node_result,
+    const std::string& node_id,
+    const std::string& result_key,
+    ExecutionResult& result,
+    const nlohmann::json& overlay,
+    const std::vector<std::string>& overlay_ignored,
+    int iteration) {
+
+    const bool in_loop = iteration >= 0;
+    const std::string where = in_loop
+        ? (" during loop iteration " + std::to_string(iteration))
+        : std::string();
+
+    if (node_result.status != NodeStatus::Completed) {
+        // A node that failed or was skipped asked for nothing. Its result is
+        // still stored, so the caller can see what happened.
+        result.node_results[result_key] = node_result;
+        return MarkerOutcome::Continue;
+    }
+
+    // A node asking to run another workflow has it run here: the node itself
+    // has no way to reach the engine, and this is the point where the result
+    // can still become its output.
+    if (node_result.output.contains("_callWorkflow")) {
+        runSubWorkflow(node_result, result.call_depth, result.execution_id);
+    }
+
+    // A pause cannot work inside a loop: there is no way to resume one item of
+    // an iteration. Better to say so than to park a run nobody can restart.
+    if (in_loop && node_result.status == NodeStatus::Completed &&
+        node_result.output.contains("_pause")) {
+        node_result.status = NodeStatus::Failed;
+        node_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";
+        node_result.output.erase("_pause");
+    }
+
+    // Record what was applied. Once a node's settings can come from elsewhere,
+    // its stored config no longer says what it ran with, and this is the only
+    // place that closes that gap.
+    if (!overlay.empty() && node_result.status == NodeStatus::Completed) {
+        node_result.output["_appliedConfig"] = overlay;
+        if (!overlay_ignored.empty()) {
+            node_result.output["_ignoredConfigKeys"] = overlay_ignored;
+            LOG_WARN("Node {} was given {} setting(s) it does not have",
+                     node_id, overlay_ignored.size());
+        }
+    }
+
+    result.node_results[result_key] = node_result;
+
+    if (node_result.status != NodeStatus::Completed) {
+        // Only reachable when the pause rejection above failed it.
+        return MarkerOutcome::Continue;
+    }
+
+    // Any node may declare the HTTP response, not just the last one to run, so
+    // appending a node to a workflow cannot silently change what its webhook
+    // returns.
+    if (node_result.output.contains("_webhookResponse")) {
+        if (!result.webhook_response.is_null()) {
+            LOG_WARN("Node {} overrides a webhook response already set by an earlier node",
+                     node_id);
+        }
+        result.webhook_response = node_result.output["_webhookResponse"];
+    }
+
+    // A workflow can name its own result instead of leaving it to be whatever
+    // the last node happened to produce. It matters most when the workflow is
+    // called by another one, where "the last node" is ambiguous in anything
+    // that branches - and a loop body is exactly such a case, which is why this
+    // has to work there too.
+    if (node_result.output.contains("_workflowOutput")) {
+        result.explicit_output = node_result.output["_workflowOutput"];
+        result.has_explicit_output = true;
+        auto& stored = result.node_results[result_key];
+        stored.output.erase("_workflowOutput");
+        stored.output["output"] = result.explicit_output;
+    }
+
+    // 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. Raised inside a loop it
+    // ends the whole run, not just the iteration - "stop the workflow" would be
+    // a strange thing to mean per-item.
+    if (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 = result.node_results[result_key];
+        stored.output.erase("_stop");
+        stored.output["stopped"] = true;
+        stored.output["reason"] = result.stop_reason;
+
+        LOG_INFO("Execution {} stopped at node {}{}: {}", result.execution_id, node_id, where,
+                 result.stop_reason.empty() ? "no reason given" : result.stop_reason);
+        return MarkerOutcome::Stop;
+    }
+
+    // A node asking to pause ends this pass. The execution is stored as Waiting
+    // with everything computed so far, and a later resume picks it up from here
+    // rather than starting again. Unreachable inside a loop, where the pause was
+    // turned into a failure above.
+    if (node_result.output.contains("_pause")) {
+        const auto& pause = node_result.output["_pause"];
+
+        result.status = ExecutionStatus::Waiting;
+        result.paused_node_id = node_id;
+        result.pause_token = pause.value("token", "");
+        result.pause_expires_at = pause.value("expiresAt", static_cast<int64_t>(0));
+
+        auto& stored = result.node_results[result_key];
+        stored.output.erase("_pause");
+
+        // No callback here - the generic terminal callback (after
+        // storeExecution) already emits "execution.waiting" once this pass
+        // ends, and it carries nodeId/expiresAt too. Emitting here as well
+        // produced the same event twice with two different payload shapes.
+
+        LOG_INFO("Execution {} paused at node {}", result.execution_id, node_id);
+        return MarkerOutcome::Pause;
+    }
+
+    return MarkerOutcome::Continue;
+}
+
 bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_depth,
                                     const std::string& parent_execution_id) {
     const auto call = node_result.output["_callWorkflow"];
@@ -2237,17 +2285,6 @@ bool WorkflowEngine::executeLoopBody(
     // 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");
-        }
-    };
-
     // 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;
@@ -2460,19 +2497,15 @@ bool WorkflowEngine::executeLoopBody(
                     cached_result.finished_at = cached_result.started_at;
                     cached_result.from_cache = true;
 
-                    rejectPauseInLoopBody(cached_result);
+                    // A cached output carries the same markers a fresh one
+                    // does, and a replayed run has to stop where the original
+                    // did - so it goes through the same handling.
+                    const MarkerOutcome cached_outcome = applyNodeMarkers(
+                        cached_result, body_node_id,
+                        body_node_id + "_iter_" + std::to_string(i), result,
+                        nlohmann::json::object(), {}, static_cast<int>(i));
 
                     iteration_results[body_node_id] = cached_result;
-                    result.node_results[body_node_id + "_iter_" + std::to_string(i)] = cached_result;
-
-                    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);
-                        }
-                        result.webhook_response = cached_result.output["_webhookResponse"];
-                    }
 
                     if (callback) {
                         nlohmann::json cache_event_data = {
@@ -2503,6 +2536,11 @@ bool WorkflowEngine::executeLoopBody(
                         }
                     }
 
+                    if (cached_outcome == MarkerOutcome::Stop) {
+                        stopped_in_body = true;
+                        break;
+                    }
+
                     continue;
                 }
 
@@ -2577,19 +2615,12 @@ bool WorkflowEngine::executeLoopBody(
                                            body_node_def_for_defaults);
             }
 
-            if (!body_overlay.empty() && body_result.status == NodeStatus::Completed) {
-                body_result.output["_appliedConfig"] = body_overlay;
-                if (!body_overlay_ignored.empty()) {
-                    body_result.output["_ignoredConfigKeys"] = body_overlay_ignored;
-                }
-            }
-
-            if (body_result.status == NodeStatus::Completed &&
-                body_result.output.contains("_callWorkflow")) {
-                runSubWorkflow(body_result, result.call_depth, result.execution_id);
-            }
+            // Store in main results with iteration suffix
+            const std::string result_key = body_node_id + "_iter_" + std::to_string(i);
 
-            rejectPauseInLoopBody(body_result);
+            const MarkerOutcome body_outcome = applyNodeMarkers(
+                body_result, body_node_id, result_key, result,
+                body_overlay, body_overlay_ignored, static_cast<int>(i));
 
             // Only a node that actually ran counts. A node on a branch that was
             // not taken is skipped, and a skipped loop-back edge is exactly the
@@ -2601,39 +2632,9 @@ bool WorkflowEngine::executeLoopBody(
 
             iteration_results[body_node_id] = body_result;
 
-            // Store in main results with iteration suffix
-            std::string result_key = body_node_id + "_iter_" + std::to_string(i);
-            result.node_results[result_key] = body_result;
-
-            if (body_result.status == NodeStatus::Completed &&
-                body_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);
-                }
-                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);
+            // The flag travels out on the result because this walk cannot end
+            // the outer one itself.
+            if (body_outcome == MarkerOutcome::Stop) {
                 stopped_in_body = true;
                 break;
             }

+ 30 - 0
src/runner/workflow_engine.hpp

@@ -104,6 +104,12 @@ struct ExecutionResult {
     nlohmann::json final_output;
     nlohmann::json workflow_snapshot;  // Snapshot of workflow at execution time
     nlohmann::json webhook_response;   // Set by a respond-to-webhook node, if any
+    // What a Workflow Output node said this workflow returns. On the result
+    // rather than local to the walk, because there are two walks - the main one
+    // and a loop body - and a workflow-output node is legitimate in either.
+    nlohmann::json explicit_output;
+    bool has_explicit_output = false;
+
     std::string paused_node_id;        // Node that asked to pause, when Waiting
     std::string pause_token;           // Must be presented to resume
     int64_t pause_expires_at = 0;      // Milliseconds since the epoch, 0 for never
@@ -275,6 +281,30 @@ private:
     // target_node_id and cached_outputs are set only when a single node is
     // being tested from the editor. A body node tested that way must run once
     // with the data the editor pinned, not once per item in the loop.
+    // What a node's output asked the engine to do, once its markers have been
+    // acted on.
+    enum class MarkerOutcome {
+        Continue,  // nothing special, carry on
+        Stop,      // end the run, but successfully
+        Pause      // end this pass; a resume picks it up
+    };
+
+    // Everything a node can ask of the engine through its output: run another
+    // workflow, name the workflow's result, set a webhook response, stop, pause.
+    //
+    // It lives in one function because there are two walks over nodes - the main
+    // one and the loop body - and every marker has to be understood by both.
+    // Handling them separately is what let a workflow-output node inside a loop
+    // be silently ignored, and it is the shape behind most of the marker bugs
+    // this engine has had. A marker added here reaches both walks or neither.
+    MarkerOutcome applyNodeMarkers(NodeExecutionResult& node_result,
+                                   const std::string& node_id,
+                                   const std::string& result_key,
+                                   ExecutionResult& result,
+                                   const nlohmann::json& overlay,
+                                   const std::vector<std::string>& overlay_ignored,
+                                   int iteration);
+
     bool executeLoopBody(const std::string& loop_node_id,
                         const Workflow& workflow,
                         const LoopContext& loop_ctx,

+ 40 - 0
tests/nodes/workflow-output-in-loop.json

@@ -0,0 +1,40 @@
+{
+  "name": "verify-workflow-output-in-loop",
+  "helpers": [
+    {
+      "key": "sub",
+      "name": "verify sub: names its result from inside a loop",
+      "active": true,
+      "nodes": [
+        {"id": "in", "name": "Input", "type": "workflow-input", "position": {"x": 0, "y": 0},
+         "config": {"fields": [{"name": "n", "type": "number", "required": true}]}},
+        {"id": "items", "name": "Items", "type": "code", "position": {"x": 0, "y": 100},
+         "config": {"code": "return { items: ['only'] };"}},
+        {"id": "loop", "name": "Loop", "type": "loop", "position": {"x": 0, "y": 200},
+         "config": {"inputField": "data.result.items"}},
+        {"id": "out", "name": "Out", "type": "workflow-output", "position": {"x": 0, "y": 300},
+         "config": {"source": "fields",
+                    "fields": [{"name": "named", "value": "from inside the loop"},
+                               {"name": "item", "value": "{{data.loop.item}}"}]}}
+      ],
+      "connections": [
+        {"sourceNodeId": "in", "sourceOutput": "main", "targetNodeId": "items", "targetInput": "data"},
+        {"sourceNodeId": "items", "sourceOutput": "main", "targetNodeId": "loop", "targetInput": "data"},
+        {"sourceNodeId": "loop", "sourceOutput": "loop", "targetNodeId": "out", "targetInput": "data"}
+      ]
+    }
+  ],
+  "nodes": [
+    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
+    {"id": "call", "name": "Call", "type": "call-workflow", "position": {"x": 0, "y": 100},
+     "config": {"workflowId": "{{helper:sub}}", "inputSource": "fields",
+                "fields": [{"name": "n", "value": "1"}]}}
+  ],
+  "connections": [
+    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "call", "targetInput": "data"}
+  ],
+  "expect": {
+    "call": {"status": "completed",
+             "output": {"result": {"named": "from inside the loop"}}}
+  }
+}