Ver Fonte

refactor: one implementation of the rules both graph walks share

executeLoopBody was a partial re-implementation of the main walk, and the two
had drifted four times. Three rules that existed twice now exist once:

  collectNodeInput   - the loop body collected its own input
  makeDisabledResult - the disabled-node record, kept in step by hand
  cachedResultFor    - replay-from-cache eligibility and its record

One of those drifts was a live bug, not just duplication. Inside a loop a
named port handed the target the whole source output instead of that port's
payload, so an if-condition's true branch delivered {result, matchedConditions,
true:{...}} where the same three nodes outside a loop delivered the payload
itself. The loop copy also treated an edge with an empty source_output as a
branch and dropped it, let a Skipped upstream contribute, and delivered a
Configurator's config edge as ordinary data. What a node received depended on
whether it happened to sit inside a loop.

Checked before changing it: 15 edges across 4 live workflows are fed by a
named port inside a loop, and none read their input - 14 go through $node[...]
or fixed config, and the one set-fields among them is in only-set mode, so it
does not pass its input through either. The change is invisible to all of them.

Two walk functions remain, on purpose. The top level goes through a topological
order once and marks a starved node "_skip"; the loop goes through a subset per
item, where a starved node may instead start from the current item, and stops
on a back edge. Those are different policies over the same per-node steps, and
merging them would give one walk threaded with "am I in a loop" at every
decision - harder to read than two callers sharing named helpers. So this is
not "one walk", and the duplication that remains is control flow rather than
rules.

Verified with a characterisation harness over eight loop shapes - nested loops,
branch gating, disabled transparency, back edges, both continueOnError
settings, an empty item list - captured before and after each of the three
steps. Exactly one case changed, the branch payload above; the rest are
byte-identical. Node suite 92 passed 0 failed after each step. The single-node
re-run path was re-checked separately across four cases, cachedResultFor
sitting directly on it.
fszontagh há 1 mês atrás
pai
commit
8e1ea33dcf
2 ficheiros alterados com 184 adições e 128 exclusões
  1. 131 127
      src/runner/workflow_engine.cpp
  2. 53 1
      src/runner/workflow_engine.hpp

+ 131 - 127
src/runner/workflow_engine.cpp

@@ -725,22 +725,9 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
             // decides where the flow goes has no answer to give, so nothing
             // after it can run.
             if (node->disabled) {
-                auto disabled_def = registry_.getNode(node->type);
-                const bool decides_branch = disabled_def &&
-                    (disabled_def->outputs.size() > 1 || !disabled_def->dynamic_outputs.is_null());
-
-                NodeExecutionResult disabled_result;
-                disabled_result.node_id = node_id;
-                disabled_result.status = NodeStatus::Disabled;
-                disabled_result.input = nlohmann::json::object();
-                disabled_result.output = nlohmann::json::object();
-                if (decides_branch) {
-                    // No branch is active, which is what stops the nodes after it.
-                    disabled_result.output["_activeBranch"] = "";
-                }
-                disabled_result.started_at = TimeUtils::nowMs();
-                disabled_result.seq = nextSeq();
-                disabled_result.finished_at = disabled_result.started_at;
+                bool decides_branch = false;
+                NodeExecutionResult disabled_result =
+                    makeDisabledResult(node_id, node->type, decides_branch);
                 result.node_results[node_id] = disabled_result;
 
                 if (callback) {
@@ -820,25 +807,19 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
 
             // Check if we can use cached output (single-node mode optimization)
             // Only use cache for upstream nodes (not the target node itself)
-            bool use_cached = false;
-            if (!target_node_id.empty() && node_id != target_node_id && cached_outputs.contains(node_id)) {
-                auto& cache_entry = cached_outputs[node_id];
-                if (cache_entry.contains("configHash") && cache_entry.contains("output")) {
-                    // Compare config hash - if unchanged, we can use cached output
-                    std::string cached_hash = cache_entry["configHash"].get<std::string>();
-                    if (configMatchesHash(cached_hash, node->config)) {
-                        use_cached = true;
-                        LOG_INFO("Using cached output for node {} (config unchanged)", node_id);
-
-                        if (node->type == "loop" && target_inside_loop) {
-                            use_cached = false;
-                            LOG_INFO("Node {} is the loop holding the node under test, so it runs", node_id);
-                        }
-                    } else {
-                        LOG_DEBUG("Node {} config changed, re-executing", node_id);
-                    }
+            std::optional<NodeExecutionResult> cached;
+            if (!target_node_id.empty() && node_id != target_node_id) {
+                cached = cachedResultFor(node_id, *node, input, cached_outputs);
+                if (cached && node->type == "loop" && target_inside_loop) {
+                    // The loop holding the node under test has to run, or the
+                    // node under test is never reached at all.
+                    cached.reset();
+                    LOG_INFO("Node {} is the loop holding the node under test, so it runs", node_id);
+                } else if (cached) {
+                    LOG_INFO("Using cached output for node {} (config unchanged)", node_id);
                 }
             }
+            const bool use_cached = cached.has_value();
 
             NodeExecutionResult node_result;
             // What this node's markers asked for. Set in both branches below,
@@ -847,15 +828,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
             MarkerOutcome marker_outcome = MarkerOutcome::Continue;
 
             if (use_cached) {
-                // Use cached output instead of executing
-                node_result.node_id = node_id;
-                node_result.status = NodeStatus::Completed;
-                node_result.input = input;
-                node_result.output = cached_outputs[node_id]["output"];
-                node_result.started_at = TimeUtils::nowMs();
-                node_result.seq = nextSeq();
-                node_result.finished_at = node_result.started_at;
-                node_result.from_cache = true;
+                node_result = *cached;
 
                 marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
                                                   nlohmann::json::object(), {}, -1);
@@ -2363,14 +2336,70 @@ nlohmann::json WorkflowEngine::collectConfigOverlay(
     return overlay;
 }
 
+std::optional<NodeExecutionResult> WorkflowEngine::cachedResultFor(
+    const std::string& node_id,
+    const WorkflowNode& node,
+    const nlohmann::json& input,
+    const nlohmann::json& cached_outputs) {
+
+    if (!cached_outputs.contains(node_id)) {
+        return std::nullopt;
+    }
+    const auto& cache_entry = cached_outputs[node_id];
+    if (!cache_entry.contains("configHash") || !cache_entry.contains("output")) {
+        return std::nullopt;
+    }
+    if (!configMatchesHash(cache_entry["configHash"].get<std::string>(), node.config)) {
+        LOG_DEBUG("Node {} config changed, re-executing", node_id);
+        return std::nullopt;
+    }
+
+    NodeExecutionResult cached;
+    cached.node_id = node_id;
+    cached.status = NodeStatus::Completed;
+    cached.input = input;
+    cached.output = cache_entry["output"];
+    cached.started_at = TimeUtils::nowMs();
+    cached.seq = nextSeq();
+    cached.finished_at = cached.started_at;
+    cached.from_cache = true;
+    return cached;
+}
+
+NodeExecutionResult WorkflowEngine::makeDisabledResult(const std::string& node_id,
+                                                       const std::string& node_type,
+                                                       bool& decides_branch) {
+    auto disabled_def = registry_.getNode(node_type);
+    decides_branch = disabled_def &&
+        (disabled_def->outputs.size() > 1 || !disabled_def->dynamic_outputs.is_null());
+
+    NodeExecutionResult disabled_result;
+    disabled_result.node_id = node_id;
+    disabled_result.status = NodeStatus::Disabled;
+    disabled_result.input = nlohmann::json::object();
+    disabled_result.output = nlohmann::json::object();
+    if (decides_branch) {
+        // No branch is active, which is what stops the nodes after it.
+        disabled_result.output["_activeBranch"] = "";
+    }
+    disabled_result.started_at = TimeUtils::nowMs();
+    disabled_result.seq = nextSeq();
+    disabled_result.finished_at = disabled_result.started_at;
+    return disabled_result;
+}
+
 nlohmann::json WorkflowEngine::collectNodeInput(
     const std::string& node_id,
     const Workflow& workflow,
-    const std::unordered_map<std::string, NodeExecutionResult>& results) {
+    const std::unordered_map<std::string, NodeExecutionResult>& results,
+    const std::string& loop_node_id,
+    const nlohmann::json* loop_input,
+    InputSources* sources) {
 
     nlohmann::json input;
     bool has_any_active_input = false;
     bool all_inputs_from_branches = true;
+    InputSources found;
 
     for (const auto& conn : workflow.connections) {
         if (conn.target_node_id == node_id) {
@@ -2382,7 +2411,32 @@ nlohmann::json WorkflowEngine::collectNodeInput(
                 continue;
             }
 
+            // The loop's own port hands the body its iteration, whole - it is
+            // already shaped as an input, so it replaces rather than merges.
+            if (loop_input && conn.source_node_id == loop_node_id) {
+                input = *loop_input;
+                found.from_loop_node = true;
+                found.any_active = true;
+                has_any_active_input = true;
+                all_inputs_from_branches = false;
+                continue;
+            }
+
             auto it = results.find(conn.source_node_id);
+
+            if (loop_input) {
+                // A disabled pass-through is not an upstream that has to
+                // produce something, or the node after it would be skipped
+                // along with it. Anything else counts, including a source that
+                // has not run this iteration.
+                const bool transparent =
+                    it != results.end() &&
+                    it->second.status == NodeStatus::Disabled &&
+                    !it->second.output.contains("_activeBranch");
+                if (!transparent) {
+                    found.any_body_upstream = true;
+                }
+            }
             if (it != results.end()) {
                 // Check if source node was skipped - propagate skip
                 if (it->second.status == NodeStatus::Skipped) {
@@ -2423,10 +2477,12 @@ nlohmann::json WorkflowEngine::collectNodeInput(
                                 input[conn.target_input] = output.value("data", output);
                             }
                             has_any_active_input = true;
+                            found.any_active = true;
                         } else {
                             // "main" output - pass through regardless of branch
                             input[conn.target_input] = output;
                             has_any_active_input = true;
+                            found.any_active = true;
                             all_inputs_from_branches = false;
                         }
                     } else {
@@ -2440,12 +2496,24 @@ nlohmann::json WorkflowEngine::collectNodeInput(
                             input[conn.target_input] = output;
                         }
                         has_any_active_input = true;
+                        found.any_active = true;
                     }
                 }
             }
         }
     }
 
+    if (sources) {
+        *sources = found;
+    }
+
+    // Inside a loop body this is the caller's decision: a node with no active
+    // upstream may need to start from the current item rather than be skipped,
+    // which "_skip" would prevent.
+    if (loop_input) {
+        return input;
+    }
+
     // If this node only receives input from branching nodes and none are active,
     // mark it for skipping
     if (all_inputs_from_branches && !has_any_active_input) {
@@ -2957,19 +3025,9 @@ bool WorkflowEngine::executeLoopBody(
                 // tell "turned off" from "never reached" - and so a disabled
                 // branching node stops the rest of the body, as it does outside
                 // a loop.
-                auto disabled_def = registry_.getNode(body_node->type);
-                NodeExecutionResult disabled_result;
-                disabled_result.node_id = body_node_id;
-                disabled_result.status = NodeStatus::Disabled;
-                disabled_result.input = nlohmann::json::object();
-                disabled_result.output = nlohmann::json::object();
-                if (disabled_def && (disabled_def->outputs.size() > 1 ||
-                                     !disabled_def->dynamic_outputs.is_null())) {
-                    disabled_result.output["_activeBranch"] = "";
-                }
-                disabled_result.started_at = TimeUtils::nowMs();
-                disabled_result.seq = nextSeq();
-                disabled_result.finished_at = disabled_result.started_at;
+                bool decides_branch = false;
+                NodeExecutionResult disabled_result =
+                    makeDisabledResult(body_node_id, body_node->type, decides_branch);
                 // Named for the loop it sits in, but not for a particular pass -
                 // it never runs, so no pass produced it.
                 disabled_result.loop_node_id = loop_node_id;
@@ -3005,60 +3063,17 @@ bool WorkflowEngine::executeLoopBody(
             }
 
             // Collect input - from iteration results or loop input
-            nlohmann::json node_input;
-            bool has_loop_input = false;
-            bool has_active_upstream = false;
-            bool has_body_upstream = false;   // wired to another body node, active or not
-
-            for (const auto& conn : workflow.connections) {
-                if (conn.target_node_id == body_node_id) {
-                    LOG_INFO("  Connection to {}: source={}, source_output={}, target_input={}",
-                             body_node_id, conn.source_node_id, conn.source_output, conn.target_input);
-                    if (conn.source_node_id == loop_node_id) {
-                        // Input from loop node - pass directly (loop_input already has data field)
-                        LOG_INFO("  -> Using loop_input (has currentItem: {})",
-                                 loop_input.contains("currentItem") ? "yes" : "no");
-                        node_input = loop_input;
-                        has_loop_input = true;
-                        has_active_upstream = true;
-                    } else {
-                        // Input from previous body node
-                        auto it = iteration_results.find(conn.source_node_id);
-
-                        // A disabled pass-through node does not count as an
-                        // upstream that has to produce something, or the node
-                        // after it would be skipped along with it.
-                        const bool transparent =
-                            it != iteration_results.end() &&
-                            it->second.status == NodeStatus::Disabled &&
-                            !it->second.output.contains("_activeBranch");
-                        if (!transparent) {
-                            has_body_upstream = true;
-                        }
-
-                        if (it != iteration_results.end() &&
-                            it->second.status == NodeStatus::Completed) {
-
-                            // Check for branch filtering (IF condition handling)
-                            const auto& source_output = it->second.output;
-                            if (source_output.contains("_activeBranch")) {
-                                std::string active_branch = source_output["_activeBranch"].get<std::string>();
-                                // Only use this input if connection matches active branch
-                                // or if source_output is "main" (default)
-                                if (conn.source_output != active_branch && conn.source_output != "main") {
-                                    // This connection is from an inactive branch - skip this input
-                                    LOG_DEBUG("Skipping input from {} via inactive branch {} (active: {})",
-                                             conn.source_node_id, conn.source_output, active_branch);
-                                    continue;
-                                }
-                            }
-
-                            node_input[conn.target_input] = source_output;
-                            has_active_upstream = true;
-                        }
-                    }
-                }
-            }
+            // The same collection the top-level walk uses. This block used
+            // to be a second implementation, and the two had drifted: a named
+            // port handed over the whole source output instead of that port's
+            // payload, an edge with an empty source_output was treated as a
+            // branch and dropped, a Skipped upstream still contributed, and a
+            // Configurator's edge arrived as ordinary data. A node's input
+            // depended on whether it happened to sit inside a loop.
+            WorkflowEngine::InputSources sources;
+            nlohmann::json node_input = collectNodeInput(body_node_id, workflow, iteration_results,
+                                                         loop_node_id, &loop_input, &sources);
+            const bool has_loop_input = sources.from_loop_node;
 
             // The loop item only used to reach the node wired directly to the
             // loop node: every later body node received nothing but its
@@ -3084,34 +3099,23 @@ bool WorkflowEngine::executeLoopBody(
             // was not taken, or the node before it was itself skipped - and the
             // skip has to carry down the chain. Falling back to the loop item
             // here would run the whole tail of an untaken branch on every item.
-            if (!has_loop_input && has_body_upstream && !has_active_upstream) {
+            if (!has_loop_input && sources.any_body_upstream && !sources.any_active) {
                 LOG_DEBUG("Skipping node {} - every upstream body node was skipped or on an inactive branch",
                           body_node_id);
                 continue;
             }
 
             // A node wired to nothing upstream starts from the loop item.
-            if (!has_loop_input && !has_body_upstream && node_input.empty()) {
+            if (!has_loop_input && !sources.any_body_upstream && node_input.empty()) {
                 node_input = loop_input;
             }
 
             // An upstream body node whose config has not changed since the
             // pinned execution is not run again - its recorded output stands in.
-            if (single_node_mode && body_node_id != target_node_id &&
-                cached_outputs.contains(body_node_id)) {
-                const auto& cache_entry = cached_outputs[body_node_id];
-                if (cache_entry.contains("configHash") && cache_entry.contains("output") &&
-                    configMatchesHash(cache_entry["configHash"].get<std::string>(), body_node->config)) {
-
-                    NodeExecutionResult cached_result;
-                    cached_result.node_id = body_node_id;
-                    cached_result.status = NodeStatus::Completed;
-                    cached_result.input = node_input;
-                    cached_result.output = cache_entry["output"];
-                    cached_result.started_at = TimeUtils::nowMs();
-                    cached_result.seq = nextSeq();
-                    cached_result.finished_at = cached_result.started_at;
-                    cached_result.from_cache = true;
+            if (single_node_mode && body_node_id != target_node_id) {
+                auto maybe_cached = cachedResultFor(body_node_id, *body_node, node_input, cached_outputs);
+                if (maybe_cached) {
+                    NodeExecutionResult cached_result = *maybe_cached;
                     cached_result.loop_node_id = loop_node_id;
                     cached_result.loop_iteration = static_cast<int>(i);
 

+ 53 - 1
src/runner/workflow_engine.hpp

@@ -353,9 +353,61 @@ private:
         const std::unordered_map<std::string, NodeExecutionResult>& results,
         std::string& conflict_error);
 
+    // What collecting a node's input revealed about its upstreams. Only the
+    // loop walk needs these: at the top level a node with no active upstream is
+    // marked with "_skip" here, while inside a loop body the same situation can
+    // mean "start from the current item instead", which is the caller's policy
+    // to apply rather than this function's.
+    struct InputSources {
+        bool from_loop_node = false;   // fed directly by the loop's own port
+        bool any_active = false;       // at least one upstream produced something
+        bool any_body_upstream = false;  // wired to a body node, active or not,
+                                         // ignoring a transparent disabled one
+    };
+
+    // The record a disabled node leaves behind, identical at the top level and
+    // inside a loop body - it was written out twice, and the two copies had to
+    // be kept in step by hand.
+    //
+    // A disabled node does not run, but what that means downstream depends on
+    // what the node does. One that passes work along is transparent: the flow
+    // carries on without it. One that decides where the flow goes has no answer
+    // to give, so nothing after it can run - which is what the empty
+    // "_activeBranch" says. decides_branch reports which of the two it was, for
+    // the caller's event and log.
+    NodeExecutionResult makeDisabledResult(const std::string& node_id, const std::string& node_type,
+                                           bool& decides_branch);
+
+    // The stand-in result for a node whose recorded output can be replayed
+    // instead of running it again, or nullopt when it cannot be.
+    //
+    // Usable means the caller pinned an execution, this is not the node under
+    // test, and the node's config still hashes to what it hashed to when that
+    // output was recorded. Both walks decided that and built the record the
+    // same way, in two places; what stays with each caller is what only it
+    // knows - the top level that a loop holding the node under test has to run
+    // anyway, the loop body which iteration to tag the record with.
+    std::optional<NodeExecutionResult> cachedResultFor(const std::string& node_id,
+                                                       const WorkflowNode& node,
+                                                       const nlohmann::json& input,
+                                                       const nlohmann::json& cached_outputs);
+
+    // One implementation for both walks.
+    //
+    // The loop body used to collect its own input, and the two had drifted
+    // apart in ways that changed what a node received depending only on whether
+    // it happened to sit inside a loop: a named port delivered the whole
+    // source output rather than that port's payload, an edge with an empty
+    // source_output was treated as a branch, a Skipped upstream still
+    // contributed, and a Configurator's edge arrived as data. Passing
+    // loop_input makes the loop-specific parts apply - an edge from the loop
+    // node carries the iteration, and "_skip" is left to the caller.
     nlohmann::json collectNodeInput(const std::string& node_id,
                                     const Workflow& workflow,
-                                    const std::unordered_map<std::string, NodeExecutionResult>& results);
+                                    const std::unordered_map<std::string, NodeExecutionResult>& results,
+                                    const std::string& loop_node_id = {},
+                                    const nlohmann::json* loop_input = nullptr,
+                                    InputSources* sources = nullptr);
 
     void storeExecution(const ExecutionResult& result);