|
|
@@ -275,6 +275,14 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
resume_seed = actual_trigger_data["_resumeSeed"];
|
|
|
actual_trigger_data.erase("_resumeSeed");
|
|
|
}
|
|
|
+ // How deep in a chain of workflow calls this run is. Carried on the trigger
|
|
|
+ // data like the other internal keys, and stripped before anything sees it.
|
|
|
+ int call_depth = 0;
|
|
|
+ if (actual_trigger_data.contains("_callDepth")) {
|
|
|
+ call_depth = actual_trigger_data.value("_callDepth", 0);
|
|
|
+ actual_trigger_data.erase("_callDepth");
|
|
|
+ }
|
|
|
+
|
|
|
if (actual_trigger_data.contains("_resumeExecutionId")) {
|
|
|
resume_execution_id = actual_trigger_data["_resumeExecutionId"].get<std::string>();
|
|
|
actual_trigger_data.erase("_resumeExecutionId");
|
|
|
@@ -290,6 +298,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
result.trigger_type = trigger_type;
|
|
|
result.trigger_data = actual_trigger_data;
|
|
|
result.started_at = TimeUtils::nowMs();
|
|
|
+ result.call_depth = call_depth;
|
|
|
|
|
|
if (!resume_execution_id.empty()) {
|
|
|
result.execution_id = resume_execution_id;
|
|
|
@@ -374,6 +383,10 @@ 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;
|
|
|
@@ -718,6 +731,14 @@ 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.
|
|
|
@@ -757,6 +778,19 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
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
|
|
|
@@ -882,7 +916,11 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
if (result.status == ExecutionStatus::Running) {
|
|
|
result.status = ExecutionStatus::Completed;
|
|
|
|
|
|
- if (result.stop_requested) {
|
|
|
+ if (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;
|
|
|
+ } 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
|
|
|
// it, not an entry that will not be found.
|
|
|
@@ -1620,6 +1658,87 @@ static std::vector<std::string> applyConfigOverlay(nlohmann::json& config,
|
|
|
return ignored;
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+// The deepest a chain of workflow calls may go. A workflow that calls itself,
|
|
|
+// directly or through others, would otherwise run until the process died - and
|
|
|
+// the failure would look like a hang rather than a mistake in the workflow.
|
|
|
+static constexpr int kMaxCallDepth = 5;
|
|
|
+
|
|
|
+bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_depth,
|
|
|
+ const std::string& parent_execution_id) {
|
|
|
+ const auto call = node_result.output["_callWorkflow"];
|
|
|
+ const std::string workflow_id = call.value("workflowId", "");
|
|
|
+ node_result.output.erase("_callWorkflow");
|
|
|
+
|
|
|
+ if (workflow_id.empty()) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Call Workflow: no workflow was chosen";
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+
|
|
|
+ if (call_depth + 1 > kMaxCallDepth) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Call Workflow: workflows are nested more than " +
|
|
|
+ std::to_string(kMaxCallDepth) + " deep, which usually means one of "
|
|
|
+ "them calls itself. Stopping here rather than running out of memory";
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto stored = storage_.get("workflows", workflow_id);
|
|
|
+ if (stored.failed()) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Call Workflow: no workflow with id " + workflow_id;
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto sub_workflow = Workflow::fromJson(stored.value());
|
|
|
+
|
|
|
+ nlohmann::json sub_trigger = call.value("input", nlohmann::json::object());
|
|
|
+ if (!sub_trigger.is_object()) {
|
|
|
+ sub_trigger = nlohmann::json{{"data", sub_trigger}};
|
|
|
+ }
|
|
|
+ sub_trigger["_callDepth"] = call_depth + 1;
|
|
|
+ sub_trigger["_calledBy"] = parent_execution_id;
|
|
|
+
|
|
|
+ LOG_INFO("Execution {} calls workflow {} at depth {}", parent_execution_id, workflow_id,
|
|
|
+ call_depth + 1);
|
|
|
+
|
|
|
+ // No callback: the sub-workflow reports its own progress against its own
|
|
|
+ // execution, and forwarding its node events to the parent's subscribers
|
|
|
+ // would make the parent's canvas light up nodes it does not have.
|
|
|
+ auto outcome = execute(sub_workflow, "workflow", sub_trigger, nullptr);
|
|
|
+ if (outcome.failed()) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Call Workflow: " + outcome.error().message();
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+
|
|
|
+ const auto& sub = outcome.value();
|
|
|
+ node_result.output["workflowId"] = workflow_id;
|
|
|
+ node_result.output["workflowName"] = sub.workflow_name;
|
|
|
+ node_result.output["executionId"] = sub.execution_id;
|
|
|
+ node_result.output["status"] = executionStatusToString(sub.status);
|
|
|
+ node_result.output["result"] = sub.final_output;
|
|
|
+
|
|
|
+ // A sub-workflow that failed fails the node that called it. Continuing with
|
|
|
+ // an empty result would hide the failure one level up, where nobody is
|
|
|
+ // looking at that execution.
|
|
|
+ if (sub.status == ExecutionStatus::Failed) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Called workflow \"" + sub.workflow_name + "\" failed: " + sub.error;
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+ if (sub.status == ExecutionStatus::Waiting) {
|
|
|
+ node_result.status = NodeStatus::Failed;
|
|
|
+ node_result.error = "Called workflow \"" + sub.workflow_name + "\" paused for an "
|
|
|
+ "approval. A workflow that waits for a person cannot be called as a "
|
|
|
+ "step - collect the answer in the calling workflow instead";
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+
|
|
|
+ return true;
|
|
|
+}
|
|
|
+
|
|
|
nlohmann::json WorkflowEngine::collectConfigOverlay(
|
|
|
const std::string& node_id,
|
|
|
const Workflow& workflow,
|
|
|
@@ -2396,6 +2515,11 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
}
|
|
|
}
|
|
|
|
|
|
+ if (body_result.status == NodeStatus::Completed &&
|
|
|
+ body_result.output.contains("_callWorkflow")) {
|
|
|
+ runSubWorkflow(body_result, result.call_depth, result.execution_id);
|
|
|
+ }
|
|
|
+
|
|
|
rejectPauseInLoopBody(body_result);
|
|
|
|
|
|
// Only a node that actually ran counts. A node on a branch that was
|