From b8bbd560354b65d6b0cd8802501a6a92368ce023 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Fri, 14 Aug 2026 22:12:34 +0700 Subject: [PATCH] feat(server): implement workflow triggering and enhance script execution - Added support for triggering workflows via a new triggerWorkflow function, allowing for fire-and-forget execution of workflows with optional data transformation using JSONata. - Introduced a $workflows API in the script sandbox for accessing workflow triggers. - Enhanced existing script execution functions to accommodate the new workflow triggering capabilities. - Updated workflow-mermaid.js and WorkflowsPage to reflect the new workflow type in labels for better clarity. --- packages/server/registry.js | 195 +++++++++++++++++--- packages/server/script-sandbox.js | 76 ++++++-- packages/server/scripts/trigger-workflow.js | 64 +++++++ packages/web/src/lib/workflow-mermaid.js | 1 + packages/web/src/pages/WorkflowsPage.jsx | 2 + 5 files changed, 303 insertions(+), 35 deletions(-) create mode 100644 packages/server/scripts/trigger-workflow.js diff --git a/packages/server/registry.js b/packages/server/registry.js index a5d89ec..8313dc5 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -21,6 +21,18 @@ import * as fsStore from "./fs-store.js"; * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry */ +const MAX_WORKFLOW_TRIGGER_DEPTH = 8; + +/** + * @param {unknown} workflow + */ +function hasWorkflowTrigger(workflow) { + if (!workflow || typeof workflow !== "object") return false; + const triggers = /** @type {{ triggers?: Array<{ type?: string }> }} */ (workflow) + .triggers; + return (triggers ?? []).some((t) => t?.type === "workflow"); +} + /** * @param {import("fastify").FastifyInstance} server */ @@ -36,6 +48,49 @@ export function createRegistry(server) { /** @type {Set} */ const registeredHttpRoutes = new Set(); + /** + * Resolve a same-owner workflow that opts in with `type: workflow`. + * Prefers YAML `name`, then filename (`name` or `name.yaml`). + * + * @param {string} owner + * @param {string} name + */ + function resolveWorkflowTriggerKey(owner, name) { + if (typeof name !== "string" || name.length === 0) { + throw new Error("workflow name is required"); + } + + /** @type {string[]} */ + const byName = []; + for (const [key, entry] of workflows) { + if (entry.owner !== owner) continue; + if (entry.workflow?.enabled === false) continue; + if (!hasWorkflowTrigger(entry.workflow)) continue; + if (entry.workflow?.name === name) byName.push(key); + } + if (byName.length === 1) return byName[0]; + if (byName.length > 1) { + throw new Error(`ambiguous workflow name "${name}"`); + } + + const fileCandidates = + name.endsWith(".yaml") || name.endsWith(".yml") + ? [`${owner}/${name}`] + : [`${owner}/${name}`, `${owner}/${name}.yaml`]; + + for (const key of fileCandidates) { + const entry = workflows.get(key); + if (!entry) continue; + if (entry.workflow?.enabled === false) continue; + if (!hasWorkflowTrigger(entry.workflow)) continue; + return key; + } + + throw new Error( + `workflow "${name}" not found or has no workflow trigger (owner "${owner}")`, + ); + } + function registerWorkflows() { workflows.clear(); loadErrors.clear(); @@ -208,12 +263,22 @@ export function createRegistry(server) { * @param {import("pino").Logger} runLog * @param {string} key * @param {string} owner + * @param {number} depth */ - async function runLinearSteps(compiled, ctx, runId, runLog, key, owner) { + async function runLinearSteps(compiled, ctx, runId, runLog, key, owner, depth) { let next = ctx; for (const index of compiled.order) { const parsed = compiled.steps[index]; - next = await runCompiledStep(parsed, next, index, runId, runLog, key, owner); + next = await runCompiledStep( + parsed, + next, + index, + runId, + runLog, + key, + owner, + depth, + ); } return next; } @@ -225,8 +290,9 @@ export function createRegistry(server) { * @param {import("pino").Logger} runLog * @param {string} key * @param {string} owner + * @param {number} depth */ - async function runDagSteps(compiled, ctx, runId, runLog, key, owner) { + async function runDagSteps(compiled, ctx, runId, runLog, key, owner, depth) { const triggerData = ctx.data; /** @type {Map} */ const outputsById = new Map(); @@ -243,6 +309,7 @@ export function createRegistry(server) { runLog, key, owner, + depth, ); if (parsed.id) { outputsById.set(parsed.id, last); @@ -251,6 +318,39 @@ export function createRegistry(server) { return last; } + /** + * @param {string} owner + * @param {string} parentKey + * @param {string} parentRunId + * @param {number} depth + */ + function createWorkflowsApi(owner, parentKey, parentRunId, depth) { + return { + /** + * @param {string} name + * @param {unknown} [data] + */ + async trigger(name, data) { + if (depth >= MAX_WORKFLOW_TRIGGER_DEPTH) { + throw new Error( + `workflow trigger depth limit (${MAX_WORKFLOW_TRIGGER_DEPTH}) exceeded`, + ); + } + const destKey = resolveWorkflowTriggerKey(owner, name); + return runWorkflow( + destKey, + { data }, + { type: "workflow", detail: parentKey }, + { + parentRunId, + depth: depth + 1, + detach: true, + }, + ); + }, + }; + } + /** * @param {import("./workflow-parse.js").CompiledStep} parsed * @param {{ data?: unknown, config?: unknown }} ctx @@ -259,8 +359,18 @@ export function createRegistry(server) { * @param {import("pino").Logger} runLog * @param {string} key * @param {string} owner + * @param {number} depth */ - async function runCompiledStep(parsed, ctx, index, runId, runLog, key, owner) { + async function runCompiledStep( + parsed, + ctx, + index, + runId, + runLog, + key, + owner, + depth, + ) { const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; const config = parsed.config; const step = await store.startStep({ @@ -292,6 +402,7 @@ export function createRegistry(server) { log: stepLog, workflowName: key, owner, + $workflows: createWorkflowsApi(owner, key, runId, depth), }); await store.finishStep(step.id, "success", result); return result; @@ -305,8 +416,17 @@ export function createRegistry(server) { * @param {string} key * @param {{ data?: unknown }} context * @param {{ type: string, detail?: string | null }} trigger + * @param {{ + * parentRunId?: string | null, + * depth?: number, + * detach?: boolean, + * }} [opts] */ - async function runWorkflow(key, context, trigger) { + async function runWorkflow(key, context, trigger, opts = {}) { + const parentRunId = opts.parentRunId ?? null; + const depth = opts.depth ?? 0; + const detach = opts.detach === true; + const entry = workflows.get(key); if (!entry) { log.error({ workflow: key }, "workflow not found"); @@ -320,33 +440,62 @@ export function createRegistry(server) { workflowName: workflow?.name, trigger, input: context.data, + parentRunId, }); const runLog = log.child({ runId: run.id, owner, workflow: key }); - runLog.debug("running workflow"); + runLog.debug(detach ? "running workflow (detached)" : "running workflow"); - let ctx = { + const initialCtx = { ...context, data: context.data ?? workflow.data ?? null, }; - try { - const compiled = compileWorkflowScripts(workflow.scripts); - if (compiled.dagMode) { - ctx = await runDagSteps(compiled, ctx, run.id, runLog, key, owner); - } else { - ctx = await runLinearSteps(compiled, ctx, run.id, runLog, key, owner); + const execute = async () => { + let ctx = initialCtx; + try { + const compiled = compileWorkflowScripts(workflow.scripts); + if (compiled.dagMode) { + ctx = await runDagSteps( + compiled, + ctx, + run.id, + runLog, + key, + owner, + depth, + ); + } else { + ctx = await runLinearSteps( + compiled, + ctx, + run.id, + runLog, + key, + owner, + depth, + ); + } + await store.finishRun(run.id, "success", ctx); + return { runId: run.id, status: "success", result: ctx }; + } catch (err) { + runLog.error({ err }, "workflow failed"); + await store.finishRun(run.id, "failed", null, err); + return { + runId: run.id, + status: "failed", + error: err instanceof Error ? err.message : String(err), + }; } - await store.finishRun(run.id, "success", ctx); - return { runId: run.id, status: "success", result: ctx }; - } catch (err) { - runLog.error({ err }, "workflow failed"); - await store.finishRun(run.id, "failed", null, err); - return { - runId: run.id, - status: "failed", - error: err instanceof Error ? err.message : String(err), - }; + }; + + if (detach) { + execute().catch((err) => { + runLog.error({ err }, "detached workflow failed unexpectedly"); + }); + return { runId: run.id, status: "started" }; } + + return execute(); } function reregister() { diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index 557943e..c193163 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -317,10 +317,28 @@ function createSecretsApi(owner) { }; } +const $workflowsStub = { + async trigger() { + throw new Error("workflow runner is not available"); + }, +}; + /** - * @param {{ log: import("pino").Logger, script: string, workflowName: string, owner?: string }} opts + * @param {{ + * log: import("pino").Logger, + * script: string, + * workflowName: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * }} opts */ -function createScriptSandbox({ log, script, workflowName, owner = "default" }) { +function createScriptSandbox({ + log, + script, + workflowName, + owner = "default", + $workflows = $workflowsStub, +}) { const scriptLog = log.child({ workflow: workflowName, script }); const $axios = createScreenedAxios(scriptLog); const $kv = createKvApi(workflowName); @@ -332,6 +350,7 @@ function createScriptSandbox({ log, script, workflowName, owner = "default" }) { $axios, $kv, $secrets, + $workflows, require: createRestrictedRequire($axios), }; @@ -376,10 +395,16 @@ export function extractScriptMeta(fn) { /** * @param {import("node:vm").Script} compiled - * @param {{ log: import("pino").Logger, script: string, workflowName: string, owner?: string }} opts + * @param {{ + * log: import("pino").Logger, + * script: string, + * workflowName: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * }} opts */ -function instantiateCompiled(compiled, { log, script, workflowName, owner }) { - const sandbox = createScriptSandbox({ log, script, workflowName, owner }); +function instantiateCompiled(compiled, { log, script, workflowName, owner, $workflows }) { + const sandbox = createScriptSandbox({ log, script, workflowName, owner, $workflows }); return compiled.runInContext(sandbox); } @@ -389,7 +414,12 @@ function instantiateCompiled(compiled, { log, script, workflowName, owner }) { * * @param {string} script * @param {string} source - * @param {{ log?: import("pino").Logger, workflowName?: string, owner?: string }} [opts] + * @param {{ + * log?: import("pino").Logger, + * workflowName?: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * }} [opts] */ export function instantiateScriptSource(script, source, opts = {}) { const compiled = compileScriptSource(source, script); @@ -398,6 +428,7 @@ export function instantiateScriptSource(script, source, opts = {}) { script, workflowName: opts.workflowName ?? "inspect", owner: opts.owner ?? "default", + $workflows: opts.$workflows, }); return { fn, ...extractScriptMeta(fn) }; } @@ -439,11 +470,22 @@ function loadCompiledScript(script) { * * @param {string} script * @param {unknown} ctx - * @param {{ log: import("pino").Logger, workflowName: string, owner?: string }} opts + * @param {{ + * log: import("pino").Logger, + * workflowName: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * }} opts */ -export async function runScript(script, ctx, { log, workflowName, owner }) { +export async function runScript(script, ctx, { log, workflowName, owner, $workflows }) { const compiled = loadCompiledScript(script); - const fn = instantiateCompiled(compiled, { log, script, workflowName, owner }); + const fn = instantiateCompiled(compiled, { + log, + script, + workflowName, + owner, + $workflows, + }); return await fn(ctx); } @@ -453,9 +495,19 @@ export async function runScript(script, ctx, { log, workflowName, owner }) { * @param {string} script * @param {string} source * @param {unknown} ctx - * @param {{ log: import("pino").Logger, workflowName: string, owner?: string }} opts + * @param {{ + * log: import("pino").Logger, + * workflowName: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * }} opts */ -export async function runScriptSource(script, source, ctx, { log, workflowName, owner }) { - const { fn } = instantiateScriptSource(script, source, { log, workflowName, owner }); +export async function runScriptSource(script, source, ctx, { log, workflowName, owner, $workflows }) { + const { fn } = instantiateScriptSource(script, source, { + log, + workflowName, + owner, + $workflows, + }); return await fn(ctx); } diff --git a/packages/server/scripts/trigger-workflow.js b/packages/server/scripts/trigger-workflow.js new file mode 100644 index 0000000..5e78d6b --- /dev/null +++ b/packages/server/scripts/trigger-workflow.js @@ -0,0 +1,64 @@ +import jsonata from "jsonata"; + +async function triggerWorkflow(ctx) { + const name = ctx.config?.name; + if (typeof name !== "string" || name.length === 0) { + throw new Error("config.name is required"); + } + + let data = ctx.data; + const expression = ctx.config?.expression; + if (typeof expression === "string" && expression.length > 0) { + const result = jsonata(expression).evaluate(ctx.data); + data = await result; + } + + const started = await $workflows.trigger(name, data); + const triggered = { name, runId: started?.runId ?? null }; + + const base = + ctx != null && typeof ctx === "object" && !Array.isArray(ctx) ? { ...ctx } : {}; + + if (base.data != null && typeof base.data === "object" && !Array.isArray(base.data)) { + return { ...base, data: { ...base.data, triggered } }; + } + return { ...base, triggered }; +} + +triggerWorkflow.meta = { + description: + "Fire-and-forget another workflow by YAML name (same owner). Destination must declare triggers: [{ type: workflow }]. Optionally reshape ctx.data with JSONata before sending.", + config: { + name: { + type: "string", + required: true, + description: "Destination workflow YAML name (same owner)", + }, + expression: { + type: "string", + required: false, + description: + "Optional JSONata expression evaluated against ctx.data; result becomes the destination run input", + }, + }, + input: {}, + output: { + triggered: { + type: "object", + description: "Record of the kicked-off run ({ name, runId }); under data when data is an object", + }, + }, + example: { + data: { + title: "Hello", + url: "https://example.com/img.png", + }, + config: { + name: "notify-comic", + expression: + '{ "title": title, "message": title, "attach": url }', + }, + }, +}; + +export default triggerWorkflow; diff --git a/packages/web/src/lib/workflow-mermaid.js b/packages/web/src/lib/workflow-mermaid.js index 3248d50..5fc3b84 100644 --- a/packages/web/src/lib/workflow-mermaid.js +++ b/packages/web/src/lib/workflow-mermaid.js @@ -47,6 +47,7 @@ function mermaidLabel(text) { function triggerLabel(t) { const type = String(t?.type ?? "").toLowerCase(); if (type === "cron") return `cron ${t.schedule ?? ""}`.trim(); + if (type === "workflow") return "workflow"; const method = t?.method ?? "POST"; const path = t?.path ?? ""; return `${method} ${path}`.trim(); diff --git a/packages/web/src/pages/WorkflowsPage.jsx b/packages/web/src/pages/WorkflowsPage.jsx index 17dc895..1498d10 100644 --- a/packages/web/src/pages/WorkflowsPage.jsx +++ b/packages/web/src/pages/WorkflowsPage.jsx @@ -211,12 +211,14 @@ function triggerLabel(t) { const type = String(t?.type ?? "").toLowerCase(); if (type === "cron") return t.schedule || "cron"; if (type === "http") return t.path || "/"; + if (type === "workflow") return ""; return t?.type ?? "—"; } function triggerKind(t) { const type = String(t?.type ?? "").toLowerCase(); if (type === "http") return t.method || "POST"; + if (type === "workflow") return "workflow"; return t?.type ?? "—"; }