|
|
@@ -239,8 +239,38 @@ nlohmann::json ExecutionResult::toJson() const {
|
|
|
// would silently reintroduce truncated resume data.
|
|
|
const bool keep_full_outputs = (status == ExecutionStatus::Waiting);
|
|
|
|
|
|
- j["nodeExecutions"] = nlohmann::json::array();
|
|
|
+ // node_results is an unordered_map, keyed by node id (or "<id>_iter_<n>"
|
|
|
+ // inside a loop) purely so lookups are cheap - its iteration order is
|
|
|
+ // hash-bucket order and has never had anything to do with when nodes
|
|
|
+ // ran. Building the array straight from that iteration is why it used to
|
|
|
+ // come back in an order the caller could not rely on. Sorting by
|
|
|
+ // (startedAt, seq) here fixes that: seq is a monotonic counter stamped
|
|
|
+ // on every result as it is produced, so two nodes that started in the
|
|
|
+ // same millisecond still come out in the order they actually ran, not in
|
|
|
+ // whatever order the hash table happened to visit them.
|
|
|
+ std::vector<const NodeExecutionResult*> ordered;
|
|
|
+ ordered.reserve(node_results.size());
|
|
|
for (const auto& [id, result] : node_results) {
|
|
|
+ // A loop body node's plain-id entry is a lookup mirror of its last
|
|
|
+ // iteration, kept only so an expression outside the loop can name
|
|
|
+ // that node directly. It is not a second execution and must not be
|
|
|
+ // counted as one.
|
|
|
+ if (result.is_loop_mirror) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ ordered.push_back(&result);
|
|
|
+ }
|
|
|
+ std::sort(ordered.begin(), ordered.end(),
|
|
|
+ [](const NodeExecutionResult* a, const NodeExecutionResult* b) {
|
|
|
+ if (a->started_at != b->started_at) {
|
|
|
+ return a->started_at < b->started_at;
|
|
|
+ }
|
|
|
+ return a->seq < b->seq;
|
|
|
+ });
|
|
|
+
|
|
|
+ j["nodeExecutions"] = nlohmann::json::array();
|
|
|
+ for (const auto* result_ptr : ordered) {
|
|
|
+ const auto& result = *result_ptr;
|
|
|
nlohmann::json nr;
|
|
|
nr["nodeId"] = result.node_id;
|
|
|
nr["status"] = nodeStatusToString(result.status);
|
|
|
@@ -251,6 +281,20 @@ nlohmann::json ExecutionResult::toJson() const {
|
|
|
nr["output"] = keep_full_outputs ? result.output : truncateLargeValues(result.output);
|
|
|
nr["error"] = result.error;
|
|
|
nr["retryCount"] = result.retry_count;
|
|
|
+ // Absent for a node that did not run inside a loop. Present for one
|
|
|
+ // that did - even a disabled one, recorded once for the whole loop
|
|
|
+ // run rather than once per pass - naming the loop that owns it and,
|
|
|
+ // for a node that actually executed, which 0-based pass produced
|
|
|
+ // this record. A nested loop's body is tagged with the innermost
|
|
|
+ // loop; the outer loop's own pass is recoverable from its own
|
|
|
+ // record in this same array, which is tagged with whatever loop
|
|
|
+ // contains it, in turn.
|
|
|
+ if (result.loop_node_id) {
|
|
|
+ nr["loopNodeId"] = *result.loop_node_id;
|
|
|
+ }
|
|
|
+ if (result.loop_iteration) {
|
|
|
+ nr["loopIteration"] = *result.loop_iteration;
|
|
|
+ }
|
|
|
j["nodeExecutions"].push_back(nr);
|
|
|
}
|
|
|
|
|
|
@@ -387,6 +431,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
seeded.status = nodeStatusFromString(it.value().value("status", "completed"));
|
|
|
seeded.output = it.value().value("output", nlohmann::json::object());
|
|
|
seeded.started_at = TimeUtils::nowMs();
|
|
|
+ seeded.seq = nextSeq();
|
|
|
seeded.finished_at = seeded.started_at;
|
|
|
result.node_results[it.key()] = seeded;
|
|
|
}
|
|
|
@@ -616,6 +661,23 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
continue;
|
|
|
}
|
|
|
|
|
|
+ // Skip nodes that were already executed (e.g., loop body nodes, or -
|
|
|
+ // on a resume - nodes seeded from the stored execution that ran
|
|
|
+ // before the pause). This is the mechanism that keeps a resume from
|
|
|
+ // re-running work: every seeded node already has a result here.
|
|
|
+ //
|
|
|
+ // This has to come before the disabled check below. A disabled
|
|
|
+ // node inside a loop body is recorded once by the loop itself,
|
|
|
+ // tagged with the loop it belongs to; the topological order still
|
|
|
+ // visits that node's own position afterwards; checking "disabled"
|
|
|
+ // first would have unconditionally overwritten that record with
|
|
|
+ // an untagged one carrying whatever time the outer walk happened
|
|
|
+ // to reach it, which is not when the node was actually reached.
|
|
|
+ if (result.node_results.contains(node_id)) {
|
|
|
+ LOG_DEBUG("Node {} already has results, skipping in main loop", node_id);
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+
|
|
|
// A disabled node does not run, but what that means downstream
|
|
|
// depends on what the node does. One that passes work along is
|
|
|
// simply transparent - the flow carries on without it. One that
|
|
|
@@ -636,6 +698,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
disabled_result.output["_activeBranch"] = "";
|
|
|
}
|
|
|
disabled_result.started_at = TimeUtils::nowMs();
|
|
|
+ disabled_result.seq = nextSeq();
|
|
|
disabled_result.finished_at = disabled_result.started_at;
|
|
|
result.node_results[node_id] = disabled_result;
|
|
|
|
|
|
@@ -653,15 +716,6 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
continue;
|
|
|
}
|
|
|
|
|
|
- // Skip nodes that were already executed (e.g., loop body nodes, or -
|
|
|
- // on a resume - nodes seeded from the stored execution that ran
|
|
|
- // before the pause). This is the mechanism that keeps a resume from
|
|
|
- // re-running work: every seeded node already has a result here.
|
|
|
- if (result.node_results.contains(node_id)) {
|
|
|
- LOG_DEBUG("Node {} already has results, skipping in main loop", node_id);
|
|
|
- continue;
|
|
|
- }
|
|
|
-
|
|
|
// Skip other trigger nodes when a specific trigger is selected
|
|
|
auto node_def = registry_.getNode(node->type);
|
|
|
if (node_def && node_def->is_trigger && node_id != trigger_node_id) {
|
|
|
@@ -671,6 +725,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
skip_result.input = nlohmann::json::object();
|
|
|
skip_result.output = nlohmann::json::object();
|
|
|
skip_result.started_at = TimeUtils::nowMs();
|
|
|
+ skip_result.seq = nextSeq();
|
|
|
skip_result.finished_at = skip_result.started_at;
|
|
|
|
|
|
result.node_results[node_id] = skip_result;
|
|
|
@@ -704,6 +759,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
skip_result.input = input;
|
|
|
skip_result.output = nlohmann::json::object();
|
|
|
skip_result.started_at = TimeUtils::nowMs();
|
|
|
+ skip_result.seq = nextSeq();
|
|
|
skip_result.finished_at = skip_result.started_at;
|
|
|
|
|
|
result.node_results[node_id] = skip_result;
|
|
|
@@ -756,6 +812,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
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;
|
|
|
|
|
|
@@ -801,6 +858,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
node_result.output = nlohmann::json::object();
|
|
|
node_result.error = overlay_conflict;
|
|
|
node_result.started_at = TimeUtils::nowMs();
|
|
|
+ node_result.seq = nextSeq();
|
|
|
node_result.finished_at = node_result.started_at;
|
|
|
} else if (!t_disabled_reference_error.empty()) {
|
|
|
node_result.node_id = node_id;
|
|
|
@@ -809,6 +867,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
node_result.output = nlohmann::json::object();
|
|
|
node_result.error = t_disabled_reference_error;
|
|
|
node_result.started_at = TimeUtils::nowMs();
|
|
|
+ node_result.seq = nextSeq();
|
|
|
node_result.finished_at = node_result.started_at;
|
|
|
} else {
|
|
|
node_result = executeNode(evaluated_node, input, result.execution_id, workflow,
|
|
|
@@ -885,19 +944,18 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
// "Loop iteration failed" on its own says nothing, and it is
|
|
|
// what an error workflow forwards to whoever is watching -
|
|
|
// so the node that actually failed and what it said are
|
|
|
- // carried out with it. The results are keyed
|
|
|
- // "<node>_iter_<n>", and the first failure is the one that
|
|
|
- // matters; the rest are usually the same thing repeated.
|
|
|
+ // carried out with it. loop_node_id/loop_iteration (rather
|
|
|
+ // than parsing the "<node>_iter_<n>" map key) are what make
|
|
|
+ // this correct for a node inside a nested loop too, where
|
|
|
+ // the key itself carries the outer loop's own prefix as
|
|
|
+ // well and cannot be split apart unambiguously.
|
|
|
std::string detail;
|
|
|
for (const auto& [key, body_result] : result.node_results) {
|
|
|
- if (body_result.status != NodeStatus::Failed) {
|
|
|
+ if (body_result.status != NodeStatus::Failed ||
|
|
|
+ !body_result.loop_iteration) {
|
|
|
continue;
|
|
|
}
|
|
|
- const auto suffix = key.find("_iter_");
|
|
|
- if (suffix == std::string::npos) {
|
|
|
- continue;
|
|
|
- }
|
|
|
- const std::string body_node_id = key.substr(0, suffix);
|
|
|
+ const std::string& body_node_id = body_result.node_id;
|
|
|
std::string body_name = body_node_id;
|
|
|
for (const auto& n : workflow.nodes) {
|
|
|
if (n.id == body_node_id && !n.name.empty()) {
|
|
|
@@ -906,7 +964,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
|
}
|
|
|
detail = "Loop body node \"" + body_name + "\" failed on item " +
|
|
|
- key.substr(suffix + 6) + ": " + body_result.error;
|
|
|
+ std::to_string(*body_result.loop_iteration) + ": " + body_result.error;
|
|
|
break;
|
|
|
}
|
|
|
|
|
|
@@ -1353,6 +1411,7 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
result.node_id = node.id;
|
|
|
result.input = input;
|
|
|
result.started_at = TimeUtils::nowMs();
|
|
|
+ result.seq = nextSeq();
|
|
|
result.status = NodeStatus::Running;
|
|
|
|
|
|
// Get node definition - reuse the caller's lookup if it already made one
|
|
|
@@ -2481,7 +2540,8 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
ExecutionResult& result,
|
|
|
ExecutionCallback callback,
|
|
|
const std::string& target_node_id,
|
|
|
- const nlohmann::json& cached_outputs) {
|
|
|
+ const nlohmann::json& cached_outputs,
|
|
|
+ const std::string& key_prefix) {
|
|
|
|
|
|
// Find nodes connected to "loop" output
|
|
|
auto body_start_nodes = findLoopBodyNodes(loop_node_id, workflow);
|
|
|
@@ -2515,11 +2575,34 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
|
|
|
body_nodes_set.insert(node_id);
|
|
|
|
|
|
+ // A nested loop node belongs to this body - it has to, or it would
|
|
|
+ // never run and never recurse - but what it feeds through its "loop"
|
|
|
+ // output is its own body, discovered separately when the recursive
|
|
|
+ // call below reaches it. Following those edges here as well as there
|
|
|
+ // would fold the inner loop's body into the outer one, so both walks
|
|
|
+ // execute it: once correctly, through the recursion, and once more
|
|
|
+ // directly, on whatever stale input the outer body last gave it.
|
|
|
+ // Only "done" - what happens after the inner loop finishes - is
|
|
|
+ // genuinely a continuation of the outer body.
|
|
|
+ bool node_is_nested_loop = false;
|
|
|
+ if (node_id != loop_node_id) {
|
|
|
+ for (const auto& n : workflow.nodes) {
|
|
|
+ if (n.id == node_id && n.type == "loop") {
|
|
|
+ node_is_nested_loop = true;
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
// Find downstream nodes
|
|
|
for (const auto& conn : workflow.connections) {
|
|
|
- if (conn.source_node_id == node_id) {
|
|
|
- to_process.push(conn.target_node_id);
|
|
|
+ if (conn.source_node_id != node_id) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ if (node_is_nested_loop && conn.source_output != "done") {
|
|
|
+ continue;
|
|
|
}
|
|
|
+ to_process.push(conn.target_node_id);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
@@ -2740,10 +2823,25 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
disabled_result.output["_activeBranch"] = "";
|
|
|
}
|
|
|
disabled_result.started_at = TimeUtils::nowMs();
|
|
|
+ disabled_result.seq = nextSeq();
|
|
|
disabled_result.finished_at = disabled_result.started_at;
|
|
|
+ // 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;
|
|
|
|
|
|
iteration_results[body_node_id] = disabled_result;
|
|
|
- result.node_results[body_node_id + "_iter_" + std::to_string(i)] = disabled_result;
|
|
|
+
|
|
|
+ // A node turned off never runs, so it never produces a new
|
|
|
+ // record: one entry per iteration said "the busiest thing in
|
|
|
+ // the workflow" about a node that did nothing sixty times.
|
|
|
+ // But leaving it out entirely would erase a decision someone
|
|
|
+ // made on purpose, so it is recorded once for the whole loop -
|
|
|
+ // on the first iteration that reaches it - rather than once
|
|
|
+ // per pass or not at all.
|
|
|
+ const std::string disabled_key = key_prefix + body_node_id;
|
|
|
+ if (!result.node_results.contains(disabled_key)) {
|
|
|
+ result.node_results[disabled_key] = disabled_result;
|
|
|
+ }
|
|
|
|
|
|
// This returns early, before the back-edge check further down,
|
|
|
// so the same rule has to be applied here. Everywhere else a
|
|
|
@@ -2866,15 +2964,18 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
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;
|
|
|
+ cached_result.loop_node_id = loop_node_id;
|
|
|
+ cached_result.loop_iteration = static_cast<int>(i);
|
|
|
|
|
|
// 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,
|
|
|
+ key_prefix + body_node_id + "_iter_" + std::to_string(i), result,
|
|
|
nlohmann::json::object(), {}, static_cast<int>(i));
|
|
|
|
|
|
iteration_results[body_node_id] = cached_result;
|
|
|
@@ -2973,6 +3074,7 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
body_result.output = nlohmann::json::object();
|
|
|
body_result.error = body_overlay_conflict;
|
|
|
body_result.started_at = TimeUtils::nowMs();
|
|
|
+ body_result.seq = nextSeq();
|
|
|
body_result.finished_at = body_result.started_at;
|
|
|
} else if (!t_disabled_reference_error.empty()) {
|
|
|
body_result.node_id = body_node_id;
|
|
|
@@ -2981,19 +3083,79 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
body_result.output = nlohmann::json::object();
|
|
|
body_result.error = t_disabled_reference_error;
|
|
|
body_result.started_at = TimeUtils::nowMs();
|
|
|
+ body_result.seq = nextSeq();
|
|
|
body_result.finished_at = body_result.started_at;
|
|
|
} else {
|
|
|
body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow,
|
|
|
body_node_def_for_defaults);
|
|
|
}
|
|
|
|
|
|
+ // Named for the loop that runs it and which pass this is - the
|
|
|
+ // same tag a node outside a loop never carries at all.
|
|
|
+ body_result.loop_node_id = loop_node_id;
|
|
|
+ body_result.loop_iteration = static_cast<int>(i);
|
|
|
+
|
|
|
+ // A body node that is itself a loop node does not iterate on its
|
|
|
+ // own - it only ever produces the "_isLoop" marker once, same as
|
|
|
+ // at the top level, and something has to act on that marker or a
|
|
|
+ // loop nested inside another loop's body would silently run its
|
|
|
+ // "body" zero times. This is that recursion: the inner loop's
|
|
|
+ // own pass is this outer iteration's body_result, tagged above
|
|
|
+ // with the outer loop's id and iteration like any other body
|
|
|
+ // node; everything the inner loop runs is tagged with the inner
|
|
|
+ // loop's id and its own 0-based iteration instead, and the key
|
|
|
+ // prefix keeps those results from colliding with the same inner
|
|
|
+ // loop's run under a different outer iteration.
|
|
|
+ bool inner_loop_left_early = false;
|
|
|
+ if (body_result.status == NodeStatus::Completed &&
|
|
|
+ body_result.output.contains("_isLoop") &&
|
|
|
+ body_result.output["_isLoop"].get<bool>()) {
|
|
|
+
|
|
|
+ LoopContext inner_ctx;
|
|
|
+ for (const auto& item : body_result.output["_items"]) {
|
|
|
+ inner_ctx.items.push_back(item);
|
|
|
+ }
|
|
|
+ inner_ctx.output_field = body_result.output.value("_outputField", "results");
|
|
|
+ inner_ctx.item_variable = body_result.output.value("_itemVariable", "item");
|
|
|
+ inner_ctx.index_variable = body_result.output.value("_indexVariable", "index");
|
|
|
+ inner_ctx.continue_on_error = body_result.output.value("_continueOnError", true);
|
|
|
+
|
|
|
+ const std::string inner_key_prefix =
|
|
|
+ key_prefix + body_node_id + "_iter_" + std::to_string(i) + "_";
|
|
|
+ const bool inner_success = executeLoopBody(
|
|
|
+ body_node_id, workflow, inner_ctx, result, callback,
|
|
|
+ target_node_id, cached_outputs, inner_key_prefix);
|
|
|
+
|
|
|
+ body_result.output["_loopCompleted"] = true;
|
|
|
+ body_result.output[inner_ctx.output_field] = inner_ctx.results;
|
|
|
+ body_result.output["_activeBranch"] = "done";
|
|
|
+
|
|
|
+ nlohmann::json inner_done_data;
|
|
|
+ inner_done_data[inner_ctx.output_field] = inner_ctx.results;
|
|
|
+ inner_done_data["totalProcessed"] = inner_ctx.results.size();
|
|
|
+ inner_done_data["data"] = body_result.output.value("data", nlohmann::json::object());
|
|
|
+ body_result.output["done"] = inner_done_data;
|
|
|
+
|
|
|
+ if (!inner_success && !inner_ctx.continue_on_error) {
|
|
|
+ inner_loop_left_early = true;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
// Store in main results with iteration suffix
|
|
|
- const std::string result_key = body_node_id + "_iter_" + std::to_string(i);
|
|
|
+ const std::string result_key = key_prefix + body_node_id + "_iter_" + std::to_string(i);
|
|
|
|
|
|
const MarkerOutcome body_outcome = applyNodeMarkers(
|
|
|
body_result, body_node_id, result_key, result,
|
|
|
body_overlay, body_overlay_ignored, static_cast<int>(i));
|
|
|
|
|
|
+ if (inner_loop_left_early) {
|
|
|
+ iteration_failed = true;
|
|
|
+ all_succeeded = false;
|
|
|
+ if (!ctx.continue_on_error) {
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
// 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
|
|
|
// case that should end the iteration.
|
|
|
@@ -3101,14 +3263,38 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
// Copy results back (ctx was passed by const ref, but we used a mutable copy)
|
|
|
const_cast<LoopContext&>(loop_ctx).results = ctx.results;
|
|
|
|
|
|
- // Store base results for body nodes so main loop will skip them
|
|
|
- // (main loop checks result.node_results.contains(node_id))
|
|
|
+ // Store a lookup mirror for every body node under its own plain, never-
|
|
|
+ // prefixed id - never as one of the caller's own results (an expression
|
|
|
+ // outside the loop naming this node directly has to find something), and
|
|
|
+ // always under the same key an enclosing walk checks with
|
|
|
+ // "result.node_results.contains(node_id)" before deciding whether to run
|
|
|
+ // it. That check is what keeps a node from running twice; for a loop
|
|
|
+ // nested inside another loop's body the enclosing walk is itself
|
|
|
+ // recursive and has no idea how deep this one went, so the mirror has to
|
|
|
+ // be reachable by plain id regardless of nesting - key_prefix is only
|
|
|
+ // ever for the real per-iteration records, which do need to stay apart
|
|
|
+ // between one outer pass's run of this loop and the next's.
|
|
|
for (const auto& body_node_id : sorted_body) {
|
|
|
+ const std::string disabled_key = key_prefix + body_node_id;
|
|
|
+ auto disabled_it = result.node_results.find(disabled_key);
|
|
|
+ if (disabled_it != result.node_results.end() &&
|
|
|
+ disabled_it->second.status == NodeStatus::Disabled) {
|
|
|
+ // At the top level disabled_key already equals body_node_id -
|
|
|
+ // that one real record IS the plain-id entry, and must stay
|
|
|
+ // visible rather than being relabelled a mirror of itself.
|
|
|
+ if (!key_prefix.empty()) {
|
|
|
+ NodeExecutionResult mirror = disabled_it->second;
|
|
|
+ mirror.is_loop_mirror = true;
|
|
|
+ result.node_results[body_node_id] = mirror;
|
|
|
+ }
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+
|
|
|
// Find the most recent iteration result for this node
|
|
|
NodeExecutionResult base_result;
|
|
|
bool found = false;
|
|
|
for (size_t i = ctx.items.size(); i > 0; --i) {
|
|
|
- std::string iter_key = body_node_id + "_iter_" + std::to_string(i - 1);
|
|
|
+ std::string iter_key = key_prefix + body_node_id + "_iter_" + std::to_string(i - 1);
|
|
|
auto it = result.node_results.find(iter_key);
|
|
|
if (it != result.node_results.end()) {
|
|
|
base_result = it->second;
|
|
|
@@ -3119,6 +3305,7 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
if (found) {
|
|
|
// Mark as executed in loop so main loop skips it
|
|
|
base_result.output["_executedInLoop"] = true;
|
|
|
+ base_result.is_loop_mirror = true;
|
|
|
result.node_results[body_node_id] = base_result;
|
|
|
}
|
|
|
}
|