|
@@ -18,6 +18,7 @@ std::string nodeStatusToString(NodeStatus status) {
|
|
|
case NodeStatus::Completed: return "completed";
|
|
case NodeStatus::Completed: return "completed";
|
|
|
case NodeStatus::Failed: return "failed";
|
|
case NodeStatus::Failed: return "failed";
|
|
|
case NodeStatus::Skipped: return "skipped";
|
|
case NodeStatus::Skipped: return "skipped";
|
|
|
|
|
+ case NodeStatus::Disabled: return "disabled";
|
|
|
default: return "unknown";
|
|
default: return "unknown";
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
@@ -73,6 +74,12 @@ namespace {
|
|
|
// the order they were inserted, while nlohmann writes them sorted. Comparing
|
|
// the order they were inserted, while nlohmann writes them sorted. Comparing
|
|
|
// the two as text calls two identical configs different and re-runs a node that
|
|
// the two as text calls two identical configs different and re-runs a node that
|
|
|
// did not change, so both sides are compared as parsed JSON instead.
|
|
// did not change, so both sides are compared as parsed JSON instead.
|
|
|
|
|
+// Reading the output of a node that was turned off is a mistake, not a missing
|
|
|
|
|
+// value, so it has to reach the node that made it rather than being logged and
|
|
|
|
|
+// replaced with an empty string. Expressions are resolved synchronously on the
|
|
|
|
|
+// thread running the node, so the message is handed back the same way.
|
|
|
|
|
+thread_local std::string t_disabled_reference_error;
|
|
|
|
|
+
|
|
|
static bool configMatchesHash(const std::string& hashed, const nlohmann::json& config) {
|
|
static bool configMatchesHash(const std::string& hashed, const nlohmann::json& config) {
|
|
|
try {
|
|
try {
|
|
|
return nlohmann::json::parse(hashed) == config;
|
|
return nlohmann::json::parse(hashed) == config;
|
|
@@ -357,6 +364,37 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A workflow whose trigger is turned off still runs when asked, and
|
|
|
|
|
+ // does nothing - that is what turning the trigger off is for. Reporting
|
|
|
|
|
+ // it as a failure would be wrong; nothing failed.
|
|
|
|
|
+ bool trigger_disabled = false;
|
|
|
|
|
+ for (const auto& node : workflow.nodes) {
|
|
|
|
|
+ if (node.id == trigger_node_id && node.disabled) {
|
|
|
|
|
+ trigger_disabled = true;
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (trigger_disabled) {
|
|
|
|
|
+ LOG_INFO("Trigger {} is disabled, so the workflow ran without doing anything",
|
|
|
|
|
+ trigger_node_id);
|
|
|
|
|
+ result.status = ExecutionStatus::Completed;
|
|
|
|
|
+ result.final_output = nlohmann::json::object();
|
|
|
|
|
+ result.final_output["triggerDisabled"] = true;
|
|
|
|
|
+ result.finished_at = TimeUtils::nowMs();
|
|
|
|
|
+ storeExecution(result);
|
|
|
|
|
+ if (callback) {
|
|
|
|
|
+ callback("execution." + executionStatusToString(result.status), {
|
|
|
|
|
+ {"executionId", result.execution_id},
|
|
|
|
|
+ {"workflowId", result.workflow_id},
|
|
|
|
|
+ {"status", executionStatusToString(result.status)},
|
|
|
|
|
+ {"error", result.error},
|
|
|
|
|
+ {"output", result.final_output}
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+ return result;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// Execute nodes in order
|
|
// Execute nodes in order
|
|
|
for (const auto& node_id : execution_order) {
|
|
for (const auto& node_id : execution_order) {
|
|
|
// Check for cancellation
|
|
// Check for cancellation
|
|
@@ -378,7 +416,43 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- if (!node || node->disabled) {
|
|
|
|
|
|
|
+ if (!node) {
|
|
|
|
|
+ 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
|
|
|
|
|
+ // 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;
|
|
|
|
|
+
|
|
|
|
|
+ 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.finished_at = disabled_result.started_at;
|
|
|
|
|
+ result.node_results[node_id] = disabled_result;
|
|
|
|
|
+
|
|
|
|
|
+ if (callback) {
|
|
|
|
|
+ callback("node.disabled", {
|
|
|
|
|
+ {"executionId", result.execution_id},
|
|
|
|
|
+ {"nodeId", node_id},
|
|
|
|
|
+ {"status", "disabled"},
|
|
|
|
|
+ {"stopsFlow", decides_branch}
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Node {} is disabled{}", node_id,
|
|
|
|
|
+ decides_branch ? ", so nothing after it will run" : "");
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -495,10 +569,20 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
} else {
|
|
} else {
|
|
|
// Evaluate expressions in node config
|
|
// Evaluate expressions in node config
|
|
|
WorkflowNode evaluated_node = *node;
|
|
WorkflowNode evaluated_node = *node;
|
|
|
|
|
+ t_disabled_reference_error.clear();
|
|
|
evaluated_node.config = evaluateExpressions(node->config, input, result.node_results, workflow);
|
|
evaluated_node.config = evaluateExpressions(node->config, input, result.node_results, workflow);
|
|
|
|
|
|
|
|
- // Execute node
|
|
|
|
|
- node_result = executeNode(evaluated_node, input, result.execution_id, workflow);
|
|
|
|
|
|
|
+ if (!t_disabled_reference_error.empty()) {
|
|
|
|
|
+ node_result.node_id = node_id;
|
|
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
|
|
+ node_result.input = input;
|
|
|
|
|
+ node_result.output = nlohmann::json::object();
|
|
|
|
|
+ node_result.error = t_disabled_reference_error;
|
|
|
|
|
+ node_result.started_at = TimeUtils::nowMs();
|
|
|
|
|
+ node_result.finished_at = node_result.started_at;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ node_result = executeNode(evaluated_node, input, result.execution_id, workflow);
|
|
|
|
|
+ }
|
|
|
result.node_results[node_id] = node_result;
|
|
result.node_results[node_id] = node_result;
|
|
|
|
|
|
|
|
if (callback) {
|
|
if (callback) {
|
|
@@ -1125,6 +1209,17 @@ nlohmann::json WorkflowEngine::collectNodeInput(
|
|
|
continue; // Skip input from skipped nodes
|
|
continue; // Skip input from skipped nodes
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ if (it->second.status == NodeStatus::Disabled) {
|
|
|
|
|
+ // A disabled node that decides a branch leaves every branch
|
|
|
|
|
+ // inactive, and the nodes after it are skipped with it. One
|
|
|
|
|
+ // that does not is transparent: it contributes no input, but
|
|
|
|
|
+ // it does not stop the node after it from running either.
|
|
|
|
|
+ if (!it->second.output.contains("_activeBranch")) {
|
|
|
|
|
+ all_inputs_from_branches = false;
|
|
|
|
|
+ }
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
if (it->second.status == NodeStatus::Completed) {
|
|
if (it->second.status == NodeStatus::Completed) {
|
|
|
const auto& output = it->second.output;
|
|
const auto& output = it->second.output;
|
|
|
|
|
|
|
@@ -1407,7 +1502,29 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- if (!body_node || body_node->disabled) continue;
|
|
|
|
|
|
|
+ if (!body_node) continue;
|
|
|
|
|
+
|
|
|
|
|
+ if (body_node->disabled) {
|
|
|
|
|
+ // Recorded rather than passed over, so the node after it can
|
|
|
|
|
+ // 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_result.output["_activeBranch"] = "";
|
|
|
|
|
+ }
|
|
|
|
|
+ disabled_result.started_at = TimeUtils::nowMs();
|
|
|
|
|
+ disabled_result.finished_at = disabled_result.started_at;
|
|
|
|
|
+
|
|
|
|
|
+ iteration_results[body_node_id] = disabled_result;
|
|
|
|
|
+ result.node_results[body_node_id + "_iter_" + std::to_string(i)] = disabled_result;
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
// Collect input - from iteration results or loop input
|
|
// Collect input - from iteration results or loop input
|
|
|
nlohmann::json node_input;
|
|
nlohmann::json node_input;
|
|
@@ -1428,8 +1545,19 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
has_active_upstream = true;
|
|
has_active_upstream = true;
|
|
|
} else {
|
|
} else {
|
|
|
// Input from previous body node
|
|
// Input from previous body node
|
|
|
- has_body_upstream = true;
|
|
|
|
|
auto it = iteration_results.find(conn.source_node_id);
|
|
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() &&
|
|
if (it != iteration_results.end() &&
|
|
|
it->second.status == NodeStatus::Completed) {
|
|
it->second.status == NodeStatus::Completed) {
|
|
|
|
|
|
|
@@ -1529,9 +1657,21 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
|
|
|
|
|
// Evaluate expressions in body node config before execution
|
|
// Evaluate expressions in body node config before execution
|
|
|
WorkflowNode evaluated_body_node = *body_node;
|
|
WorkflowNode evaluated_body_node = *body_node;
|
|
|
|
|
+ t_disabled_reference_error.clear();
|
|
|
evaluated_body_node.config = evaluateExpressions(body_node->config, node_input, merged_results, workflow);
|
|
evaluated_body_node.config = evaluateExpressions(body_node->config, node_input, merged_results, workflow);
|
|
|
|
|
|
|
|
- auto body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow);
|
|
|
|
|
|
|
+ NodeExecutionResult body_result;
|
|
|
|
|
+ if (!t_disabled_reference_error.empty()) {
|
|
|
|
|
+ body_result.node_id = body_node_id;
|
|
|
|
|
+ body_result.status = NodeStatus::Failed;
|
|
|
|
|
+ body_result.input = node_input;
|
|
|
|
|
+ body_result.output = nlohmann::json::object();
|
|
|
|
|
+ body_result.error = t_disabled_reference_error;
|
|
|
|
|
+ body_result.started_at = TimeUtils::nowMs();
|
|
|
|
|
+ body_result.finished_at = body_result.started_at;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow);
|
|
|
|
|
+ }
|
|
|
iteration_results[body_node_id] = body_result;
|
|
iteration_results[body_node_id] = body_result;
|
|
|
|
|
|
|
|
// Store in main results with iteration suffix
|
|
// Store in main results with iteration suffix
|
|
@@ -1896,10 +2036,15 @@ nlohmann::json WorkflowEngine::evaluateJavaScriptExpression(
|
|
|
|
|
|
|
|
// Build $node object with all node outputs indexed by name
|
|
// Build $node object with all node outputs indexed by name
|
|
|
nlohmann::json node_outputs = nlohmann::json::object();
|
|
nlohmann::json node_outputs = nlohmann::json::object();
|
|
|
|
|
+ nlohmann::json disabled_names = nlohmann::json::array();
|
|
|
for (const auto& node : workflow.nodes) {
|
|
for (const auto& node : workflow.nodes) {
|
|
|
auto it = results.find(node.id);
|
|
auto it = results.find(node.id);
|
|
|
- if (it != results.end() && it->second.status == NodeStatus::Completed) {
|
|
|
|
|
|
|
+ if (it == results.end()) continue;
|
|
|
|
|
+
|
|
|
|
|
+ if (it->second.status == NodeStatus::Completed) {
|
|
|
node_outputs[node.name] = it->second.output;
|
|
node_outputs[node.name] = it->second.output;
|
|
|
|
|
+ } else if (it->second.status == NodeStatus::Disabled) {
|
|
|
|
|
+ disabled_names.push_back(node.name);
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -1908,6 +2053,17 @@ nlohmann::json WorkflowEngine::evaluateJavaScriptExpression(
|
|
|
std::string js_code = R"(
|
|
std::string js_code = R"(
|
|
|
const $node = )" + node_outputs.dump() + R"(;
|
|
const $node = )" + node_outputs.dump() + R"(;
|
|
|
|
|
|
|
|
|
|
+ // Reading the output of a node that was turned off is a mistake worth
|
|
|
|
|
+ // hearing about: it yields nothing, and silently returning undefined
|
|
|
|
|
+ // would show up much later as an unexplained empty value.
|
|
|
|
|
+ )" + disabled_names.dump() + R"(.forEach(function (name) {
|
|
|
|
|
+ Object.defineProperty($node, name, {
|
|
|
|
|
+ get: function () {
|
|
|
|
|
+ throw new Error('Node "' + name + '" is disabled and produced no output');
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
async function execute(config, input, context) {
|
|
async function execute(config, input, context) {
|
|
|
const $trigger = input;
|
|
const $trigger = input;
|
|
|
// Handle data - unwrap double nesting if present (input.data.data)
|
|
// Handle data - unwrap double nesting if present (input.data.data)
|
|
@@ -1955,6 +2111,9 @@ nlohmann::json WorkflowEngine::evaluateJavaScriptExpression(
|
|
|
result.output.is_null() ? "null" : result.output.dump().substr(0, 100));
|
|
result.output.is_null() ? "null" : result.output.dump().substr(0, 100));
|
|
|
return result.output;
|
|
return result.output;
|
|
|
} else {
|
|
} else {
|
|
|
|
|
+ if (result.error.find("is disabled and produced no output") != std::string::npos) {
|
|
|
|
|
+ t_disabled_reference_error = result.error;
|
|
|
|
|
+ }
|
|
|
LOG_WARN("JavaScript expression evaluation failed: {} - Error: {}", expression, result.error);
|
|
LOG_WARN("JavaScript expression evaluation failed: {} - Error: {}", expression, result.error);
|
|
|
return nullptr;
|
|
return nullptr;
|
|
|
}
|
|
}
|