Переглянути джерело

Merge branch 'sub-workflows': call a workflow as a step

fszontagh 1 місяць тому
батько
коміт
892ff16f07

+ 93 - 0
nodes/core/call-workflow.js

@@ -0,0 +1,93 @@
+/**
+ * @node call-workflow
+ * @name Call Workflow
+ * @category flow-control
+ * @version 1.0.0
+ * @description Run another workflow as a step and carry on with what it returns
+ * @icon workflow
+ */
+
+const configSchema = {
+    type: 'object',
+    properties: {
+        workflowId: {
+            type: 'string',
+            title: 'Workflow',
+            description: 'The workflow to run. It should start with a Workflow Input node, which says what it expects to be passed',
+            dynamicOptions: { source: 'workflows' }
+        },
+        inputSource: {
+            type: 'string',
+            title: 'Pass',
+            enum: ['fields', 'input'],
+            default: 'fields',
+            description: 'fields builds exactly what to pass. input forwards everything that reached this node, which is quicker but sends more than the called workflow asked for'
+        },
+        fields: {
+            type: 'array',
+            title: 'Input',
+            description: 'The values to pass. Names must match what the called workflow expects',
+            showWhen: { field: 'inputSource', value: 'fields' },
+            items: {
+                type: 'object',
+                properties: {
+                    name: { type: 'string', title: 'Name' },
+                    value: { type: 'string', title: 'Value', description: 'A literal, or an expression such as {{data.result.city}}' }
+                }
+            }
+        }
+    },
+    required: ['workflowId']
+};
+
+const inputSchema = { type: 'object', properties: { data: { type: 'any' } } };
+
+const outputSchema = {
+    type: 'object',
+    properties: {
+        result: { type: 'any', description: 'What the called workflow returned' },
+        workflowId: { type: 'string' },
+        workflowName: { type: 'string' },
+        executionId: { type: 'string', description: 'The called workflow has its own execution, which can be opened and inspected on its own' },
+        status: { type: 'string' }
+    }
+};
+
+async function execute(config, input, context) {
+    const workflowId = String(config.workflowId || '').trim();
+    if (!workflowId) {
+        throw new Error('Call Workflow: choose a workflow to run');
+    }
+
+    let payload;
+    if ((config.inputSource || 'fields') === 'input') {
+        payload = (input && input.data !== undefined) ? input.data : input;
+        if (payload === undefined || payload === null) {
+            payload = {};
+        }
+    } else {
+        payload = {};
+        const fields = config.fields || [];
+        for (let i = 0; i < fields.length; i++) {
+            const field = fields[i] || {};
+            const name = String(field.name || '').trim();
+            if (name) {
+                payload[name] = field.value;
+            }
+        }
+    }
+
+    smartbotic.log.info('Call Workflow: running ' + workflowId);
+
+    // A node cannot reach the engine, so this asks the engine to make the call
+    // and put the answer here. The engine replaces this marker with the called
+    // workflow's result, its execution id and its status.
+    return {
+        _callWorkflow: {
+            workflowId: workflowId,
+            input: payload
+        }
+    };
+}
+
+module.exports = { configSchema, inputSchema, outputSchema, execute };

+ 70 - 0
nodes/core/workflow-output.js

@@ -0,0 +1,70 @@
+/**
+ * @node workflow-output
+ * @name Workflow Output
+ * @category flow-control
+ * @version 1.0.0
+ * @description Name what this workflow returns to whoever called it
+ * @icon log-out
+ */
+
+const configSchema = {
+    type: 'object',
+    properties: {
+        source: {
+            type: 'string',
+            title: 'Return',
+            enum: ['input', 'fields'],
+            default: 'input',
+            description: 'input returns whatever reached this node. fields builds an object from the values below'
+        },
+        fields: {
+            type: 'array',
+            title: 'Fields',
+            description: 'The object to return',
+            showWhen: { field: 'source', value: 'fields' },
+            items: {
+                type: 'object',
+                properties: {
+                    name: { type: 'string', title: 'Name' },
+                    value: { type: 'string', title: 'Value', description: 'A literal, or an expression such as {{data.result.total}}' }
+                }
+            }
+        }
+    }
+};
+
+const inputSchema = { type: 'object', properties: { data: { type: 'any' } } };
+
+const outputSchema = {
+    type: 'object',
+    properties: {
+        output: { type: 'object', description: 'What the caller receives' }
+    }
+};
+
+async function execute(config, input, context) {
+    let value;
+
+    if ((config.source || 'input') === 'fields') {
+        value = {};
+        const fields = config.fields || [];
+        for (let i = 0; i < fields.length; i++) {
+            const field = fields[i] || {};
+            const name = String(field.name || '').trim();
+            if (name) {
+                value[name] = field.value;
+            }
+        }
+    } else {
+        value = (input && input.data !== undefined) ? input.data : input;
+    }
+
+    smartbotic.log.info('Workflow Output: returning ' +
+        (value && typeof value === 'object' ? Object.keys(value).length + ' key(s)' : typeof value));
+
+    // The engine reads this and makes it the execution's result, so a caller
+    // gets it rather than whatever node happened to run last.
+    return { _workflowOutput: value };
+}
+
+module.exports = { configSchema, inputSchema, outputSchema, execute };

+ 99 - 0
nodes/triggers/workflow-input.js

@@ -0,0 +1,99 @@
+/**
+ * @node workflow-input
+ * @name Workflow Input
+ * @category triggers
+ * @version 1.0.0
+ * @description Start a workflow that another workflow calls, and say what it expects to be given
+ * @trigger
+ * @icon log-in
+ */
+
+const configSchema = {
+    type: 'object',
+    properties: {
+        fields: {
+            type: 'array',
+            title: 'Expected Input',
+            description: 'What the calling workflow must or may pass. A caller that leaves out a required one is refused here rather than failing somewhere deeper',
+            items: {
+                type: 'object',
+                properties: {
+                    name: { type: 'string', title: 'Name', description: 'Key the caller passes, such as city' },
+                    required: { type: 'boolean', title: 'Required', description: 'Refuse the call when this is missing', default: true },
+                    description: { type: 'string', title: 'Description', description: 'What the caller should put here' },
+                    defaultValue: { type: 'string', title: 'Default', description: 'Used when an optional key is not passed' }
+                }
+            }
+        },
+        allowExtra: {
+            type: 'boolean',
+            title: 'Allow Anything Else',
+            description: 'Pass through keys that are not listed above. Turn off to keep the contract exact and catch a caller sending the wrong thing',
+            default: true
+        }
+    }
+};
+
+const inputSchema = { type: 'object', properties: { data: { type: 'any' } } };
+
+const outputSchema = {
+    type: 'object',
+    properties: {
+        calledBy: { type: 'string', description: 'Execution id of the workflow that called this one, empty when run directly' }
+    }
+};
+
+async function execute(config, input, context) {
+    // The caller's input arrives as the trigger data, the same way a webhook
+    // body does. Internal keys the engine put there are not part of it.
+    const given = (input && typeof input === 'object') ? input : {};
+    const fields = config.fields || [];
+
+    const result = {};
+    const missing = [];
+
+    for (let i = 0; i < fields.length; i++) {
+        const field = fields[i] || {};
+        const name = String(field.name || '').trim();
+        if (!name) continue;
+
+        const has = Object.prototype.hasOwnProperty.call(given, name) &&
+                    given[name] !== undefined && given[name] !== null && given[name] !== '';
+
+        if (has) {
+            result[name] = given[name];
+        } else if (field.required !== false) {
+            missing.push(name);
+        } else if (field.defaultValue !== undefined && field.defaultValue !== '') {
+            result[name] = field.defaultValue;
+        }
+    }
+
+    // Refusing here, naming every missing key at once, beats failing three
+    // nodes later on an undefined that has to be traced back.
+    if (missing.length > 0) {
+        throw new Error('Workflow Input: the caller did not pass ' + missing.join(', ') +
+            '. ' + (missing.length === 1 ? 'It is' : 'They are') + ' marked required on this node');
+    }
+
+    if (config.allowExtra !== false) {
+        const keys = Object.keys(given);
+        for (let i = 0; i < keys.length; i++) {
+            const key = keys[i];
+            // Engine bookkeeping, not part of what the caller sent.
+            if (key.charAt(0) === '_') continue;
+            if (!Object.prototype.hasOwnProperty.call(result, key)) {
+                result[key] = given[key];
+            }
+        }
+    }
+
+    result.calledBy = (given && given._calledBy) || '';
+
+    smartbotic.log.info('Workflow Input: accepted ' +
+        Object.keys(result).filter(function (k) { return k !== 'calledBy'; }).length + ' value(s)');
+
+    return result;
+}
+
+module.exports = { configSchema, inputSchema, outputSchema, execute };

+ 125 - 1
src/runner/workflow_engine.cpp

@@ -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

+ 11 - 0
src/runner/workflow_engine.hpp

@@ -116,6 +116,12 @@ struct ExecutionResult {
     // It lives on the result rather than being handled locally because a stop
     // can be raised inside a loop body, which runs in a separate walk and has
     // to hand the decision back to the caller.
+    // How deep in a chain of workflow calls this run is. Lives on the result so
+    // the loop-body walk can see it too - it takes the result but not the
+    // trigger data, and a sub-workflow called from inside a loop must count
+    // towards the same depth limit as one called outside it.
+    int call_depth = 0;
+
     bool stop_requested = false;
     std::string stopped_node_id;
     std::string stop_reason;
@@ -238,6 +244,11 @@ private:
 
     // Gather the settings any connected Configurator supplies for this node.
     // conflict_error is set when two live Configurators claim the same key.
+    // Run another workflow as a step and put its result on the calling node.
+    // Returns false and fills node_result.error when the call cannot be made.
+    bool runSubWorkflow(NodeExecutionResult& node_result, int call_depth,
+                        const std::string& parent_execution_id);
+
     nlohmann::json collectConfigOverlay(
         const std::string& node_id,
         const Workflow& workflow,

+ 12 - 0
tests/nodes/call-workflow-missing-input.json

@@ -0,0 +1,12 @@
+{
+  "name": "verify-call-workflow-missing-input",
+  "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": "wf_5cfe91eb-a7f0-4da8-aca9-2300ed40b399", "inputSource": "fields", "fields": [{"name": "label", "value": "no number given"}]}}
+  ],
+  "connections": [
+    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "call", "targetInput": "data"}
+  ],
+  "expect": {"call": {"status": "failed", "errorContains": "did not pass n"}}
+}

+ 19 - 0
tests/nodes/call-workflow.json

@@ -0,0 +1,19 @@
+{
+  "name": "verify-call-workflow",
+  "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": "wf_5cfe91eb-a7f0-4da8-aca9-2300ed40b399", "inputSource": "fields",
+                "fields": [{"name": "n", "value": "21"}]}},
+    {"id": "after", "name": "After", "type": "code", "position": {"x": 0, "y": 200},
+     "config": {"code": "const d=input.data||input; return { got: d.result, status: d.status };"}}
+  ],
+  "connections": [
+    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "call", "targetInput": "data"},
+    {"sourceNodeId": "call", "sourceOutput": "main", "targetNodeId": "after", "targetInput": "data"}
+  ],
+  "expectStatus": "completed",
+  "expect": {
+    "after": {"status": "completed", "output": {"result": {"status": "completed", "got": {"doubled": 42, "label": "unnamed"}}}}
+  }
+}

+ 31 - 1
webui/src/components/workflow/NodeConfigModal.tsx

@@ -6,7 +6,7 @@ import { useTheme } from '../../contexts/ThemeContext'
 import { ConfiguratorFields, ConfigTarget } from './ConfiguratorFields'
 import { NodeOptionsSelect } from './NodeOptionsSelect'
 import { groupFields, isCredentialKey } from './fieldGroups'
-import { NodeDefinition } from '../../api/workflows'
+import { NodeDefinition, workflowsApi } from '../../api/workflows'
 import { credentialsApi, CredentialInfo } from '../../api/credentials'
 import { ConditionBuilder, Condition, AvailableField } from '../ConditionBuilder'
 import { ExpressionInput } from '../ExpressionInput'
@@ -73,6 +73,19 @@ export function NodeConfigModal({
   })
 
   const credentials: CredentialInfo[] = credentialsData?.credentials || []
+
+  // Only fetched when the node being edited actually offers a workflow field.
+  const needsWorkflowList = Object.values(
+    ((nodeDefs.find((nd) => nd.id === selectedNodeData?.type)?.configSchema as any)?.properties) || {}
+  ).some((p: any) => p?.dynamicOptions?.source === 'workflows')
+
+  const { data: workflowsData } = useQuery({
+    queryKey: ['workflows', 'all-for-picker'],
+    queryFn: () => workflowsApi.list(1, 200),
+    enabled: needsWorkflowList,
+    staleTime: 60 * 1000,
+  })
+  const workflowChoices: { id: string; name: string }[] = workflowsData?.workflows || []
   const { resolvedTheme } = useTheme()
 
   // Editing a node is sometimes a one-line change and sometimes a hundred lines
@@ -484,6 +497,23 @@ export function NodeConfigModal({
                         placeholder={prop.description}
                         className="w-full px-3 py-2 border border-gray-200 dark:border-slate-600 rounded-lg bg-white dark:bg-slate-900 text-gray-900 dark:text-gray-100 focus:ring-2 focus:ring-primary-500 focus:border-primary-500"
                       />
+                    ) : prop.dynamicOptions?.source === 'workflows' ? (
+                      <select
+                        value={editingConfig[key] ?? ''}
+                        onChange={(e) => onConfigChange({ ...editingConfig, [key]: e.target.value })}
+                        className="w-full px-3 py-2 border border-gray-300 dark:border-slate-600 rounded-lg bg-white dark:bg-slate-700 text-gray-900 dark:text-gray-100"
+                      >
+                        <option value="">Choose a workflow</option>
+                        {(workflowChoices || [])
+                          // Calling itself is the one choice that is never
+                          // right, so it is not offered.
+                          .filter((w) => w.id !== workflowId)
+                          .map((w) => (
+                            <option key={w.id} value={w.id}>
+                              {w.name}
+                            </option>
+                          ))}
+                      </select>
                     ) : prop.dynamicOptions?.source === 'node' ? (
                       <NodeOptionsSelect
                         fieldKey={key}