diff --git a/packages/server/registry.js b/packages/server/registry.js index 65c13ab..8f9cc5b 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -7,7 +7,10 @@ import { log } from "./logger.js"; import * as store from "./store.js"; import { clearScriptCache, runScript } from "./script-sandbox.js"; import { + SET_STEP_SCRIPT, compileWorkflowScripts, + evaluateJsonata, + isJsonataTruthy, mergeStepData, namespacedPath, parseScriptStep, @@ -209,7 +212,7 @@ export function createRegistry(server) { let next = ctx; for (const index of compiled.order) { const parsed = compiled.steps[index]; - next = await executeStep(parsed, next, index, runId, runLog, key); + next = await runCompiledStep(parsed, next, index, runId, runLog, key); } return next; } @@ -230,7 +233,7 @@ export function createRegistry(server) { for (const [orderIndex, stepIndex] of compiled.order.entries()) { const parsed = compiled.steps[stepIndex]; const data = mergeStepData(parsed, outputsById, triggerData); - last = await executeStep( + last = await runCompiledStep( parsed, { ...ctx, data }, orderIndex, @@ -253,8 +256,9 @@ export function createRegistry(server) { * @param {import("pino").Logger} runLog * @param {string} key */ - async function executeStep(parsed, ctx, index, runId, runLog, key) { - const { script, config } = parsed; + async function runCompiledStep(parsed, ctx, index, runId, runLog, key) { + const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; + const config = parsed.config; const step = await store.startStep({ runId, index, @@ -263,6 +267,23 @@ export function createRegistry(server) { }); const stepLog = runLog.child({ stepId: step.id, script }); try { + if (parsed.when) { + const whenResult = await evaluateJsonata(parsed.when, ctx); + if (!isJsonataTruthy(whenResult)) { + stepLog.debug({ when: parsed.when }, "step skipped"); + await store.finishStep(step.id, "skipped", ctx); + return ctx; + } + } + if (parsed.kind === "set") { + if (ctx == null || typeof ctx !== "object" || Array.isArray(ctx)) { + throw new Error("set requires an object context"); + } + const value = await evaluateJsonata(parsed.expression, ctx); + const result = { ...ctx, [parsed.as]: value }; + await store.finishStep(step.id, "success", result); + return result; + } const result = await runScript(script, { ...ctx, config }, { log: stepLog, workflowName: key, @@ -334,7 +355,8 @@ export function createRegistry(server) { for (const { workflow } of workflows.values()) { for (const raw of workflow.scripts ?? []) { try { - refs.add(parseScriptStep(raw).script); + const parsed = parseScriptStep(raw); + if (parsed.kind === "script") refs.add(parsed.script); } catch { // skip invalid steps } diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index 923140f..276e0d3 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -26,7 +26,8 @@ function scriptNames(workflow) { const names = []; for (const raw of workflow.scripts ?? []) { try { - names.push(parseScriptStep(raw).script); + const parsed = parseScriptStep(raw); + names.push(parsed.kind === "set" ? `set:${parsed.as}` : parsed.script); } catch { names.push(null); } diff --git a/packages/server/store.js b/packages/server/store.js index 6c76017..c50c584 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -126,7 +126,7 @@ export async function startStep({ runId, index, script, config = null }) { /** * @param {string} id - * @param {"success" | "failed"} status + * @param {"success" | "failed" | "skipped"} status * @param {unknown} [output] * @param {Error | unknown} [err] */ diff --git a/packages/server/workflow-parse.js b/packages/server/workflow-parse.js index 2c9e6d5..0aa9d5b 100644 --- a/packages/server/workflow-parse.js +++ b/packages/server/workflow-parse.js @@ -1,12 +1,35 @@ +import jsonata from "jsonata"; + +const AS_IDENT = /^[A-Za-z_][A-Za-z0-9_]*$/; +const RESERVED_AS = new Set(["data", "config"]); + +export const SET_STEP_SCRIPT = "set"; + /** * @typedef {{ alias: string, from: string }} NeedEdge * @typedef {{ + * kind: "script", * script: string, * config: unknown | null, + * expression?: undefined, + * as?: undefined, * id: string | null, * needsKind: "none" | "list" | "map", * needs: NeedEdge[], - * }} ParsedStep + * when: string | null, + * }} ParsedScriptStep + * @typedef {{ + * kind: "set", + * script: typeof SET_STEP_SCRIPT, + * config: { expression: string, as: string }, + * expression: string, + * as: string, + * id: string | null, + * needsKind: "none", + * needs: NeedEdge[], + * when: string | null, + * }} ParsedSetStep + * @typedef {ParsedScriptStep | ParsedSetStep} ParsedStep * @typedef {ParsedStep & { index: number }} CompiledStep * @typedef {{ * dagMode: boolean, @@ -16,30 +39,36 @@ */ /** - * @param {unknown} step - * @returns {ParsedStep} + * Compile a JSONata source so save/load fails on syntax errors. + * @param {string} source + * @param {string} label */ -export function parseScriptStep(step) { - if (typeof step === "string") { - return { - script: step, - config: null, - id: null, - needsKind: "none", - needs: [], - }; +export function compileJsonata(source, label) { + try { + return jsonata(source); + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + throw new Error(`Invalid ${label}: ${msg}`); } - 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)}`); +} + +/** + * @param {unknown} value + */ +export function isJsonataTruthy(value) { + if (value == null || value === false) return false; + if (Array.isArray(value) && value.length === 0) return false; + return true; +} + +/** + * Evaluate JSONata against the full ctx (jsonata 2 may return a Promise). + * @param {string} source + * @param {unknown} ctx + */ +export async function evaluateJsonata(source, ctx) { + const result = compileJsonata(source, "expression").evaluate(ctx); + return await result; } /** @@ -52,6 +81,56 @@ export function namespacedPath(owner, triggerPath) { return `/u/${owner}/${cleaned}`; } +/** + * @param {unknown} step + * @returns {ParsedStep} + */ +export function parseScriptStep(step) { + if (typeof step === "string") { + return { + kind: "script", + script: step, + config: null, + id: null, + needsKind: "none", + needs: [], + when: null, + }; + } + if (step == null || typeof step !== "object" || Array.isArray(step)) { + throw new Error(`Invalid script step: ${JSON.stringify(step)}`); + } + + const hasScript = step.script != null && step.script !== ""; + const hasSet = step.set != null; + + if (hasScript && hasSet) { + throw new Error("Step cannot have both script and set"); + } + + if (hasSet) { + return parseSetStep(step); + } + + if (hasScript) { + if (typeof step.script !== "string") { + throw new Error(`Invalid script step: ${JSON.stringify(step)}`); + } + const { needsKind, needs } = parseNeeds(step.needs); + return { + kind: "script", + script: step.script, + config: step.config ?? null, + id: parseOptionalId(step.id), + needsKind, + needs, + when: parseWhen(step.when), + }; + } + + throw new Error(`Invalid script step: ${JSON.stringify(step)}`); +} + /** * Parse scripts, detect linear vs DAG mode, and return a topological order. * Linear mode (no `needs` on any step) keeps array order. @@ -81,6 +160,18 @@ export function compileWorkflowScripts(scripts) { ids.set(step.id, step.index); } + if (dagMode) { + for (const step of steps) { + const who = step.id ?? String(step.index); + if (step.kind === "set") { + throw new Error(`set steps are not allowed in DAG workflows (step ${who})`); + } + if (step.when) { + throw new Error(`when is not allowed in DAG workflows (step ${who})`); + } + } + } + if (!dagMode) { return { dagMode: false, steps, order: steps.map((s) => s.index) }; } @@ -167,6 +258,56 @@ export function mergeStepData(step, outputsById, triggerData) { return data; } +/** + * @param {object} step + * @returns {ParsedSetStep} + */ +function parseSetStep(step) { + if (step.needs != null) { + throw new Error("needs is not allowed on set steps"); + } + const spec = step.set; + if (spec == null || typeof spec !== "object" || Array.isArray(spec)) { + throw new Error(`Invalid set step: ${JSON.stringify(step)}`); + } + const expression = spec.expression; + const as = spec.as; + if (typeof expression !== "string" || !expression.trim()) { + throw new Error("set.expression is required"); + } + if (typeof as !== "string" || !AS_IDENT.test(as)) { + throw new Error(`Invalid set.as: ${JSON.stringify(as)}`); + } + if (RESERVED_AS.has(as)) { + throw new Error(`set.as "${as}" is reserved`); + } + compileJsonata(expression, "set.expression"); + return { + kind: "set", + script: SET_STEP_SCRIPT, + config: { expression, as }, + expression, + as, + id: parseOptionalId(step.id), + needsKind: "none", + needs: [], + when: parseWhen(step.when), + }; +} + +/** + * @param {unknown} when + * @returns {string | null} + */ +function parseWhen(when) { + if (when == null || when === "") return null; + if (typeof when !== "string") { + throw new Error(`Invalid when: ${JSON.stringify(when)}`); + } + compileJsonata(when, "when"); + return when; +} + /** * @param {unknown} id * @returns {string | null} @@ -181,7 +322,7 @@ function parseOptionalId(id) { /** * @param {unknown} needs - * @returns {{ needsKind: ParsedStep["needsKind"], needs: NeedEdge[] }} + * @returns {{ needsKind: ParsedScriptStep["needsKind"], needs: NeedEdge[] }} */ function parseNeeds(needs) { if (needs == null) { diff --git a/packages/web/src/lib/format.jsx b/packages/web/src/lib/format.jsx index 1cd3b04..2462203 100644 --- a/packages/web/src/lib/format.jsx +++ b/packages/web/src/lib/format.jsx @@ -13,7 +13,9 @@ export function StatusBadge({ status }) { ? "badge-error" : status === "running" ? "badge-warning" - : "badge-ghost"; + : status === "skipped" + ? "badge-info" + : "badge-ghost"; return {status ?? "—"}; } diff --git a/packages/web/src/lib/workflow-mermaid.js b/packages/web/src/lib/workflow-mermaid.js index 9876a67..3248d50 100644 --- a/packages/web/src/lib/workflow-mermaid.js +++ b/packages/web/src/lib/workflow-mermaid.js @@ -1,12 +1,26 @@ function parseStep(step) { if (typeof step === "string") { - return { script: step, id: null, needs: null }; + return { kind: "script", script: step, id: null, needs: null, when: null, as: null }; + } + if (step?.set) { + const as = typeof step.set?.as === "string" && step.set.as ? step.set.as : "set"; + return { + kind: "set", + script: null, + as, + id: typeof step.id === "string" && step.id ? step.id : null, + needs: null, + when: typeof step.when === "string" && step.when ? step.when : null, + }; } if (step?.script) { return { + kind: "script", script: step.script, + as: null, id: typeof step.id === "string" && step.id ? step.id : null, needs: step.needs ?? null, + when: typeof step.when === "string" && step.when ? step.when : null, }; } return null; @@ -38,6 +52,15 @@ function triggerLabel(t) { return `${method} ${path}`.trim(); } +function stepLabel(s) { + if (s.kind === "set") { + const base = s.id ? `${s.id}: set ${s.as}` : `set ${s.as}`; + return s.when ? `${base} when: ${s.when}` : base; + } + const base = s.id ? `${s.id}: ${s.script}` : s.script; + return s.when ? `${base} when: ${s.when}` : base; +} + export function workflowToFlowchart(parsed) { /** @type {Record} */ const scriptIds = {}; @@ -72,9 +95,13 @@ export function workflowToFlowchart(parsed) { } for (const s of parsedSteps) { - scriptIds[s.mermaidId] = s.script; - const label = s.id ? `${s.id}: ${s.script}` : s.script; - lines.push(` ${s.mermaidId}["${mermaidLabel(label)}"]`); + const label = mermaidLabel(stepLabel(s)); + if (s.kind === "set") { + lines.push(` ${s.mermaidId}(["${label}"])`); + } else { + scriptIds[s.mermaidId] = s.script; + lines.push(` ${s.mermaidId}["${label}"]`); + } } if (dagMode) { diff --git a/packages/web/src/pages/EventDetailPage.jsx b/packages/web/src/pages/EventDetailPage.jsx index 5537753..9b3bd31 100644 --- a/packages/web/src/pages/EventDetailPage.jsx +++ b/packages/web/src/pages/EventDetailPage.jsx @@ -4,6 +4,17 @@ import { useRun } from "../api/hooks.js"; import { LogViewer } from "../components/LogViewer.jsx"; import { formatTime, StatusBadge } from "../lib/format.jsx"; +function stepLabel(s) { + if (s.script === "set") { + return s.config?.as ? `set:${s.config.as}` : "set"; + } + return s.script; +} + +function isEditableScript(s) { + return Boolean(s.script) && s.script !== "set"; +} + export function EventDetailPage() { const { id } = useParams(); const { data: run, isLoading, error } = useRun(id); @@ -67,9 +78,13 @@ export function EventDetailPage() { {s.step_index} - - {s.script} - + {isEditableScript(s) ? ( + + {stepLabel(s)} + + ) : ( + stepLabel(s) + )} diff --git a/packages/web/src/pages/EventsPage.jsx b/packages/web/src/pages/EventsPage.jsx index fc4e22e..2bb8b15 100644 --- a/packages/web/src/pages/EventsPage.jsx +++ b/packages/web/src/pages/EventsPage.jsx @@ -38,6 +38,7 @@ export function EventsPage() { + {isLoading ? (