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.
This commit is contained in:
+158
-9
@@ -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<string>} */
|
||||
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<string, unknown>} */
|
||||
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,21 +440,40 @@ 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,
|
||||
};
|
||||
|
||||
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);
|
||||
ctx = await runDagSteps(
|
||||
compiled,
|
||||
ctx,
|
||||
run.id,
|
||||
runLog,
|
||||
key,
|
||||
owner,
|
||||
depth,
|
||||
);
|
||||
} else {
|
||||
ctx = await runLinearSteps(compiled, ctx, run.id, runLog, key, owner);
|
||||
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 };
|
||||
@@ -347,6 +486,16 @@ export function createRegistry(server) {
|
||||
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() {
|
||||
|
||||
@@ -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<unknown> },
|
||||
* }} 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<unknown> },
|
||||
* }} 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<unknown> },
|
||||
* }} [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<unknown> },
|
||||
* }} 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<unknown> },
|
||||
* }} 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);
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
@@ -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();
|
||||
|
||||
@@ -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 ?? "—";
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user