feat(server): enhance script execution and workflow management

- Introduced compileWorkflowScripts function to parse and validate workflow scripts, supporting both linear and DAG execution modes.
- Added runLinearSteps and runDagSteps functions for executing compiled workflows based on their structure.
- Updated createRegistry function to utilize the new workflow compilation and execution logic.
- Enhanced script sandbox with meta extraction and instantiation capabilities for better script management.
- Improved API endpoints to return script metadata and handle errors more effectively.
- Removed deprecated scripts and updated workflows to reflect new structure and functionality.
This commit is contained in:
2026-08-14 19:14:44 +07:00
parent ddff33cf49
commit 049f5c56e3
21 changed files with 809 additions and 87 deletions
+89 -20
View File
@@ -6,7 +6,12 @@ import { WORKFLOWS_DIR } from "./paths.js";
import { log } from "./logger.js";
import * as store from "./store.js";
import { clearScriptCache, runScript } from "./script-sandbox.js";
import { namespacedPath, parseScriptStep } from "./workflow-parse.js";
import {
compileWorkflowScripts,
mergeStepData,
namespacedPath,
parseScriptStep,
} from "./workflow-parse.js";
import * as fsStore from "./fs-store.js";
/**
@@ -68,6 +73,7 @@ export function createRegistry(server) {
try {
const workflowData = fs.readFileSync(filePath, "utf8");
const workflow = yaml.parse(workflowData);
compileWorkflowScripts(workflow?.scripts);
workflows.set(key, { owner, file, workflow });
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
@@ -192,6 +198,83 @@ export function createRegistry(server) {
);
}
/**
* @param {import("./workflow-parse.js").CompiledScripts} compiled
* @param {{ data?: unknown }} ctx
* @param {string} runId
* @param {import("pino").Logger} runLog
* @param {string} key
*/
async function runLinearSteps(compiled, ctx, runId, runLog, key) {
let next = ctx;
for (const index of compiled.order) {
const parsed = compiled.steps[index];
next = await executeStep(parsed, next, index, runId, runLog, key);
}
return next;
}
/**
* @param {import("./workflow-parse.js").CompiledScripts} compiled
* @param {{ data?: unknown }} ctx
* @param {string} runId
* @param {import("pino").Logger} runLog
* @param {string} key
*/
async function runDagSteps(compiled, ctx, runId, runLog, key) {
const triggerData = ctx.data;
/** @type {Map<string, unknown>} */
const outputsById = new Map();
let last = ctx;
for (const [orderIndex, stepIndex] of compiled.order.entries()) {
const parsed = compiled.steps[stepIndex];
const data = mergeStepData(parsed, outputsById, triggerData);
last = await executeStep(
parsed,
{ ...ctx, data },
orderIndex,
runId,
runLog,
key,
);
if (parsed.id) {
outputsById.set(parsed.id, last);
}
}
return last;
}
/**
* @param {import("./workflow-parse.js").CompiledStep} parsed
* @param {{ data?: unknown, config?: unknown }} ctx
* @param {number} index
* @param {string} runId
* @param {import("pino").Logger} runLog
* @param {string} key
*/
async function executeStep(parsed, ctx, index, runId, runLog, key) {
const { script, config } = parsed;
const step = await store.startStep({
runId,
index,
script,
config,
});
const stepLog = runLog.child({ stepId: step.id, script });
try {
const result = await runScript(script, { ...ctx, config }, {
log: stepLog,
workflowName: key,
});
await store.finishStep(step.id, "success", result);
return result;
} catch (err) {
await store.finishStep(step.id, "failed", null, err);
throw err;
}
}
/**
* @param {string} key
* @param {{ data?: unknown }} context
@@ -221,25 +304,11 @@ export function createRegistry(server) {
};
try {
for (const [index, rawStep] of (workflow.scripts ?? []).entries()) {
const { script, config } = parseScriptStep(rawStep);
const step = await store.startStep({
runId: run.id,
index,
script,
config,
});
const stepLog = runLog.child({ stepId: step.id, script });
try {
ctx = await runScript(script, { ...ctx, config }, {
log: stepLog,
workflowName: key,
});
await store.finishStep(step.id, "success", ctx);
} catch (err) {
await store.finishStep(step.id, "failed", null, err);
throw err;
}
const compiled = compileWorkflowScripts(workflow.scripts);
if (compiled.dagMode) {
ctx = await runDagSteps(compiled, ctx, run.id, runLog, key);
} else {
ctx = await runLinearSteps(compiled, ctx, run.id, runLog, key);
}
await store.finishRun(run.id, "success", ctx);
return { runId: run.id, status: "success", result: ctx };
+75 -5
View File
@@ -3,6 +3,7 @@ import path from "path";
import vm from "node:vm";
import { createRequire } from "node:module";
import axios from "axios";
import pino from "pino";
import { createKvApi } from "./kv-store.js";
import { SCRIPTS_DIR } from "./paths.js";
@@ -298,6 +299,78 @@ function compileScriptSource(source, filename) {
return new vm.Script(wrapScriptSource(source, filename), { filename });
}
const inspectLog = pino({ level: "silent" });
/**
* @param {unknown} fn
* @returns {{ meta: Record<string, unknown> | null, metaError: string | null }}
*/
export function extractScriptMeta(fn) {
if (typeof fn !== "function") {
return { meta: null, metaError: "default export must be a function" };
}
if (!("meta" in fn) || fn.meta == null) {
return { meta: null, metaError: null };
}
try {
const serialized = JSON.parse(JSON.stringify(fn.meta));
if (serialized == null || typeof serialized !== "object" || Array.isArray(serialized)) {
return { meta: null, metaError: "script.meta must be a plain object" };
}
return { meta: serialized, metaError: null };
} catch (err) {
return {
meta: null,
metaError: err instanceof Error ? err.message : "script.meta could not be serialized",
};
}
}
/**
* @param {import("node:vm").Script} compiled
* @param {{ log: import("pino").Logger, script: string, workflowName: string }} opts
*/
function instantiateCompiled(compiled, { log, script, workflowName }) {
const sandbox = createScriptSandbox({ log, script, workflowName });
return compiled.runInContext(sandbox);
}
/**
* Compile source and return the default-export function plus extracted meta.
* Does not call the script.
*
* @param {string} script
* @param {string} source
* @param {{ log?: import("pino").Logger, workflowName?: string }} [opts]
*/
export function instantiateScriptSource(script, source, opts = {}) {
const compiled = compileScriptSource(source, script);
const fn = instantiateCompiled(compiled, {
log: opts.log ?? inspectLog,
script,
workflowName: opts.workflowName ?? "inspect",
});
return { fn, ...extractScriptMeta(fn) };
}
/**
* Evaluate module-level code and read `defaultExport.meta` without calling the script.
*
* @param {string} script
* @param {string} source
*/
export function inspectScriptSource(script, source) {
try {
const { meta, metaError } = instantiateScriptSource(script, source);
return { meta, metaError };
} catch (err) {
return {
meta: null,
metaError: err instanceof Error ? err.message : String(err),
};
}
}
function loadCompiledScript(script) {
const filePath = path.join(SCRIPTS_DIR, script);
const { mtimeMs } = fs.statSync(filePath);
@@ -321,8 +394,7 @@ function loadCompiledScript(script) {
*/
export async function runScript(script, ctx, { log, workflowName }) {
const compiled = loadCompiledScript(script);
const sandbox = createScriptSandbox({ log, script, workflowName });
const fn = compiled.runInContext(sandbox);
const fn = instantiateCompiled(compiled, { log, script, workflowName });
return await fn(ctx);
}
@@ -335,8 +407,6 @@ export async function runScript(script, ctx, { log, workflowName }) {
* @param {{ log: import("pino").Logger, workflowName: string }} opts
*/
export async function runScriptSource(script, source, ctx, { log, workflowName }) {
const compiled = compileScriptSource(source, script);
const sandbox = createScriptSandbox({ log, script, workflowName });
const fn = compiled.runInContext(sandbox);
const { fn } = instantiateScriptSource(script, source, { log, workflowName });
return await fn(ctx);
}
-4
View File
@@ -1,4 +0,0 @@
export default function main(context) {
log.info({ data: context.data }, "add-one");
return { data: (context.data || 0) + 1 };
}
-4
View File
@@ -1,4 +0,0 @@
export default function main(context) {
log.info({ input: context.input }, "add-ten");
return context.input + 10;
}
+27 -1
View File
@@ -37,7 +37,7 @@ function resolveUrl(ctx) {
return null;
}
export default async function fetchBinary(ctx) {
async function fetchBinary(ctx) {
ensureDataObject(ctx);
const url = resolveUrl(ctx);
@@ -68,3 +68,29 @@ export default async function fetchBinary(ctx) {
return ctx;
}
fetchBinary.meta = {
description: "Download a binary URL into ctx.data as a Buffer",
config: {
url: { type: "string", required: false, description: "Direct download URL" },
urlVar: { type: "string", required: false, description: "Key in ctx.data that holds the URL" },
outputVar: { type: "string", default: "file", description: "ctx.data key for the Buffer" },
filename: { type: "string", required: false, description: "Override saved filename" },
},
input: {
attach: { type: "string", required: false, description: "Fallback URL" },
url: { type: "string", required: false, description: "Fallback URL" },
filename: { type: "string", required: false },
},
output: {
file: { type: "buffer", description: "Downloaded bytes (or ctx.config.outputVar)" },
filename: { type: "string" },
contentType: { type: "string" },
},
example: {
data: { attach: "https://example.com/image.png" },
config: { outputVar: "file" },
},
};
export default fetchBinary;
+29 -1
View File
@@ -29,7 +29,7 @@ function previewValue(value) {
return { preview: `${json.slice(0, 500)}...`, truncated: true };
}
export default async function fetchHtml(ctx) {
async function fetchHtml(ctx) {
const url = ctx.config?.url;
if (typeof url !== "string" || url.length === 0) {
throw new Error("fetch-html: ctx.config.url is required");
@@ -83,3 +83,31 @@ export default async function fetchHtml(ctx) {
return ctx;
}
fetchHtml.meta = {
description: "Fetch HTML and optionally select elements or transform with JSONata",
config: {
url: { type: "string", required: true, description: "Page URL" },
selector: { type: "string", required: false, description: "CSS selector" },
outputVar: {
type: "string",
required: false,
description: "Required when selector or jsonata is set",
},
jsonata: { type: "string", required: false, description: "JSONata applied to selector matches" },
},
input: {},
output: {
httpResponse: { type: "string", description: "Raw HTML" },
},
example: {
data: {},
config: {
url: "https://example.com/",
selector: "h1",
outputVar: "headings",
},
},
};
export default fetchHtml;
+18 -2
View File
@@ -1,10 +1,26 @@
// this script will get current time
export default function getCurrentTime() {
function getCurrentTime() {
return {
data: {
datetime: new Date().toISOString(),
processId: "1234"
}
}
}
}
getCurrentTime.meta = {
description: "Return the current time as ctx.data.datetime",
config: {},
input: {},
output: {
datetime: { type: "string", description: "ISO timestamp" },
processId: { type: "string" },
},
example: {
data: {},
config: {},
},
};
export default getCurrentTime;
+19 -2
View File
@@ -1,10 +1,27 @@
import jsonata from "jsonata";
export default function jsonataFn(ctx) {
function jsonataFn(ctx) {
log.info({ctx}, "jsonata: context")
log.info("jsonata: evaluating expression %s", ctx.config.expression);
const expression = jsonata(ctx.config.expression);
const result = expression.evaluate(ctx.data);
log.info({result}, "jsonata: expression result");
return result;
}
}
jsonataFn.meta = {
description: "Evaluate a JSONata expression against ctx.data and return the result as the next context",
config: {
expression: { type: "string", required: true, description: "JSONata expression" },
},
input: {},
output: {},
example: {
data: { title: "Hello", url: "https://example.com" },
config: {
expression: '{"data": {"title": title, "message": title, "attach": url}}',
},
},
};
export default jsonataFn;
+29 -1
View File
@@ -9,7 +9,7 @@ function ntfyHeaders(ctx) {
return headers;
}
export default async function ntfy(ctx) {
async function ntfy(ctx) {
const file = ctx.data?.file;
const hasFile = Buffer.isBuffer(file) || file instanceof Uint8Array;
@@ -64,3 +64,31 @@ export default async function ntfy(ctx) {
});
return { sent: "true" };
}
ntfy.meta = {
description: "Send a message or file to an ntfy topic",
config: {
url: {
type: "string",
default: "https://ntfy.sh/scrunner",
description: "ntfy topic URL",
},
},
input: {
title: { type: "string", required: false },
message: { type: "string", required: false },
attach: { type: "string", required: false, description: "Remote attachment URL" },
file: { type: "buffer", required: false, description: "Binary body to PUT" },
filename: { type: "string", required: false },
contentType: { type: "string", required: false },
},
output: {
sent: { type: "string", description: "Replaces the workflow context with { sent: \"true\" }" },
},
example: {
data: { title: "Hello", message: "Hello from scrunner" },
config: { url: "https://ntfy.sh/scrunner" },
},
};
export default ntfy;
+25 -5
View File
@@ -1,4 +1,8 @@
import { clearScriptCache, runScriptSource } from "../../script-sandbox.js";
import {
clearScriptCache,
inspectScriptSource,
instantiateScriptSource,
} from "../../script-sandbox.js";
import * as fsStore from "../../fs-store.js";
import { createDryRunLogger, safeSerialize } from "./dry-run-logger.js";
@@ -11,7 +15,15 @@ export default function scriptsPluginFactory(registry) {
*/
return async function scriptsPlugin(fastify) {
fastify.get("/scripts", async () => {
return { scripts: fsStore.listScriptFiles() };
const scripts = fsStore.listScriptFiles().map((name) => {
const content = fsStore.readScript(name);
const inspected =
content == null
? { meta: null, metaError: "script not found" }
: inspectScriptSource(name, content);
return { name, ...inspected };
});
return { scripts };
});
fastify.get("/scripts/:name", async (req, reply) => {
@@ -23,7 +35,7 @@ export default function scriptsPluginFactory(registry) {
}
const content = fsStore.readScript(name);
if (content == null) return reply.code(404).send({ error: "script not found" });
return { name, content };
return { name, content, ...inspectScriptSource(name, content) };
});
fastify.put("/scripts/:name", async (req, reply) => {
@@ -40,7 +52,10 @@ export default function scriptsPluginFactory(registry) {
const existed = fsStore.readScript(name) != null;
fsStore.writeScript(name, body.content);
clearScriptCache();
return reply.code(existed ? 200 : 201).send({ name });
return reply.code(existed ? 200 : 201).send({
name,
...inspectScriptSource(name, body.content),
});
});
fastify.delete("/scripts/:name", async (req, reply) => {
@@ -86,24 +101,29 @@ export default function scriptsPluginFactory(registry) {
const started = Date.now();
try {
const output = await runScriptSource(name, body.content, ctx, {
const { fn, meta, metaError } = instantiateScriptSource(name, body.content, {
log,
workflowName: "dry-run",
});
const output = await fn(ctx);
return {
status: "success",
output: safeSerialize(output),
error: null,
logs,
durationMs: Date.now() - started,
meta,
metaError,
};
} catch (err) {
const inspected = inspectScriptSource(name, body.content);
return {
status: "failed",
output: null,
error: err instanceof Error ? err.message : String(err),
logs,
durationMs: Date.now() - started,
...inspected,
};
}
});
+14 -2
View File
@@ -1,7 +1,11 @@
import yaml from "yaml";
import * as store from "../../store.js";
import * as fsStore from "../../fs-store.js";
import { namespacedPath, parseScriptStep } from "../../workflow-parse.js";
import {
compileWorkflowScripts,
namespacedPath,
parseScriptStep,
} from "../../workflow-parse.js";
function triggerSummary(owner, workflow) {
if (!workflow || typeof workflow !== "object") return [];
@@ -160,13 +164,21 @@ export default function workflowsPluginFactory(registry) {
if (typeof body.content !== "string") {
return reply.code(400).send({ error: "content is required" });
}
let parsed;
try {
yaml.parse(body.content);
parsed = yaml.parse(body.content);
} catch (err) {
return reply.code(400).send({
error: `invalid yaml: ${err instanceof Error ? err.message : String(err)}`,
});
}
try {
compileWorkflowScripts(parsed?.scripts);
} catch (err) {
return reply.code(400).send({
error: err instanceof Error ? err.message : String(err),
});
}
const existed = fsStore.readWorkflowYaml(owner, file) != null;
fsStore.writeWorkflowYaml(owner, file, body.content);
const registered = fsStore.readRegisters(owner);
+219 -2
View File
@@ -1,6 +1,44 @@
/**
* @typedef {{ alias: string, from: string }} NeedEdge
* @typedef {{
* script: string,
* config: unknown | null,
* id: string | null,
* needsKind: "none" | "list" | "map",
* needs: NeedEdge[],
* }} ParsedStep
* @typedef {ParsedStep & { index: number }} CompiledStep
* @typedef {{
* dagMode: boolean,
* steps: CompiledStep[],
* order: number[],
* }} CompiledScripts
*/
/**
* @param {unknown} step
* @returns {ParsedStep}
*/
export function parseScriptStep(step) {
if (typeof step === "string") return { script: step, config: null };
if (step?.script) return { script: step.script, config: step.config ?? null };
if (typeof step === "string") {
return {
script: step,
config: null,
id: null,
needsKind: "none",
needs: [],
};
}
if (step?.script) {
const { needsKind, needs } = parseNeeds(step.needs);
return {
script: step.script,
config: step.config ?? null,
id: parseOptionalId(step.id),
needsKind,
needs,
};
}
throw new Error(`Invalid script step: ${JSON.stringify(step)}`);
}
@@ -13,3 +51,182 @@ export function namespacedPath(owner, triggerPath) {
const cleaned = String(triggerPath).replace(/^\/+/, "");
return `/u/${owner}/${cleaned}`;
}
/**
* Parse scripts, detect linear vs DAG mode, and return a topological order.
* Linear mode (no `needs` on any step) keeps array order.
* @param {unknown} scripts
* @returns {CompiledScripts}
*/
export function compileWorkflowScripts(scripts) {
if (scripts == null) {
return { dagMode: false, steps: [], order: [] };
}
if (!Array.isArray(scripts)) {
throw new Error("workflow scripts must be an array");
}
const steps = scripts.map((raw, index) => ({
...parseScriptStep(raw),
index,
}));
const dagMode = steps.some((s) => s.needsKind !== "none");
const ids = new Map();
for (const step of steps) {
if (!step.id) continue;
if (ids.has(step.id)) {
throw new Error(`Duplicate step id: ${step.id}`);
}
ids.set(step.id, step.index);
}
if (!dagMode) {
return { dagMode: false, steps, order: steps.map((s) => s.index) };
}
for (const step of steps) {
for (const edge of step.needs) {
if (!ids.has(edge.from)) {
const who = step.id ?? String(step.index);
throw new Error(`Unknown needs id "${edge.from}" (referenced by step ${who})`);
}
if (ids.get(edge.from) === step.index) {
throw new Error(`Step "${step.id ?? step.index}" cannot need itself`);
}
}
}
const n = steps.length;
const indegree = Array(n).fill(0);
const children = Array.from({ length: n }, () => []);
for (const step of steps) {
const parents = new Set(step.needs.map((e) => ids.get(e.from)));
indegree[step.index] = parents.size;
for (const p of parents) {
children[p].push(step.index);
}
}
/** @type {number[]} */
const ready = [];
for (let i = 0; i < n; i++) {
if (indegree[i] === 0) ready.push(i);
}
const order = [];
while (ready.length) {
const i = ready.shift();
order.push(i);
const nexts = children[i].slice().sort((a, b) => a - b);
for (const c of nexts) {
indegree[c] -= 1;
if (indegree[c] === 0) {
insertSorted(ready, c);
}
}
}
if (order.length !== n) {
throw new Error("Workflow has a cycle");
}
return { dagMode: true, steps, order };
}
/**
* Prefer a script return's `.data`; otherwise treat the whole return as data.
* @param {unknown} result
*/
export function extractStepData(result) {
if (result && typeof result === "object" && !Array.isArray(result) && "data" in result) {
return result.data;
}
return result ?? null;
}
/**
* Build ctx.data for a step from upstream outputs (DAG mode).
* @param {CompiledStep} step
* @param {Map<string, unknown>} outputsById
* @param {unknown} triggerData
*/
export function mergeStepData(step, outputsById, triggerData) {
if (step.needsKind === "none" || step.needs.length === 0) {
return triggerData;
}
if (step.needsKind === "list" && step.needs.length === 1) {
return extractStepData(outputsById.get(step.needs[0].from));
}
/** @type {Record<string, unknown>} */
const data = {};
for (const { alias, from } of step.needs) {
data[alias] = extractStepData(outputsById.get(from));
}
return data;
}
/**
* @param {unknown} id
* @returns {string | null}
*/
function parseOptionalId(id) {
if (id == null || id === "") return null;
if (typeof id !== "string") {
throw new Error(`Invalid step id: ${JSON.stringify(id)}`);
}
return id;
}
/**
* @param {unknown} needs
* @returns {{ needsKind: ParsedStep["needsKind"], needs: NeedEdge[] }}
*/
function parseNeeds(needs) {
if (needs == null) {
return { needsKind: "none", needs: [] };
}
if (Array.isArray(needs)) {
/** @type {NeedEdge[]} */
const list = [];
for (const item of needs) {
if (typeof item !== "string" || item.length === 0) {
throw new Error(`Invalid needs entry: ${JSON.stringify(item)}`);
}
list.push({ alias: item, from: item });
}
return { needsKind: "list", needs: list };
}
if (typeof needs === "object") {
/** @type {NeedEdge[]} */
const list = [];
for (const [alias, from] of Object.entries(needs)) {
if (
typeof alias !== "string" ||
alias.length === 0 ||
typeof from !== "string" ||
from.length === 0
) {
throw new Error(`Invalid needs map: ${JSON.stringify(needs)}`);
}
list.push({ alias, from });
}
return { needsKind: "map", needs: list };
}
throw new Error(`Invalid needs: ${JSON.stringify(needs)}`);
}
/**
* @param {number[]} arr
* @param {number} value
*/
function insertSorted(arr, value) {
for (let k = 0; k < arr.length; k++) {
if (value < arr[k]) {
arr.splice(k, 0, value);
return;
}
}
arr.push(value);
}
@@ -1,10 +0,0 @@
name: manual trigger
scripts:
- add-one.js
- add-ten.js
triggers:
- type: HTTP
method: POST
path: /mt
- type: manual
@@ -1,7 +1,7 @@
scripts:
- time-to-ntfy-example.yaml
- test.yaml
- manual-trigger.yaml
- cron-example.yaml
- fetch-devto.yaml
- comic-monkeyuser-to-ntfy.yaml
- time-and-comic-to-ntfy.yaml
@@ -0,0 +1,37 @@
name: time and comic to ntfy
description: >
Fan-in from two scripts (current time + monkeyuser comic), then send to ntfy.
scripts:
- id: time
script: get-current-time.js
- id: comic
script: fetch-html.js
config:
url: "https://www.monkeyuser.com/"
outputVar: "httpResponse"
selector: ".comic img"
jsonata: |
{"url": "https://www.monkeyuser.com" & [attributes.src][0], "title": [attributes.title][0]}
- id: compose
script: jsonata.js
needs: [time, comic]
config:
expression: |
{"data": {
"title": comic.httpResponse.title,
"message": time.datetime & " " & comic.httpResponse.title,
"attach": comic.httpResponse.url
}}
- id: notify
script: ntfy.js
needs: [compose]
config:
url: https://ntfy.sh/scrunner
triggers:
- type: HTTP
method: POST
path: /time-and-comic