diff --git a/README.md b/README.md index 633e9ef..2786685 100644 --- a/README.md +++ b/README.md @@ -21,6 +21,31 @@ pnpm dev The first account created becomes **admin**. Later accounts are created from Users. +## Script contract + +Each script is `async function main(ctx)` and **must** return: + +```js +{ output, context?, skipRemaining? } +``` + +| Field | Meaning | +|---|---| +| `ctx.data` | This step’s input (trigger payload, previous `output`, or DAG `needs`) | +| `ctx.context` | Run clipboard (plain object, default `{}`) | +| `ctx.config` | This step’s YAML config | +| `output` | Becomes the **next** step’s `data` | +| `context` | Next snapshot of the bag. Omitted → keep incoming | +| `skipRemaining` | Stop later steps. Sibling of `output`/`context`, not inside `output` | + +Returning the full `ctx` is an error. Mutating `ctx.data` or `ctx.context` does not persist unless returned. + +YAML **SET** evaluates JSONata against the full `ctx`; the result is `output` (the next step’s data). `jsonata.js` does the same. + +DAG `needs` assemble this step’s `data` from upstream **outputs**. Independent steps in the same wave share a context snapshot; sibling writes to the same context key fail the run. + +Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / `meta.output` / `meta.context` describe `data`, the return pipe, and clipboard keys. + ## Scripts | Command | Description | diff --git a/packages/server/registry.js b/packages/server/registry.js index 474da1d..3a00336 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -15,6 +15,13 @@ import { namespacedPath, parseScriptStep, } from "./workflow-parse.js"; +import { + chainCtx, + mergeContextWave, + normalizeContext, + normalizeStepResult, + storedEnvelope, +} from "./step-result.js"; import * as fsStore from "./fs-store.js"; import { checkHttpAuth, @@ -327,24 +334,15 @@ export function createRegistry(server) { ); } - function isSkipRemaining(result) { - return ( - result != null && - typeof result === "object" && - !Array.isArray(result) && - /** @type {{ skipRemaining?: unknown }} */ (result).skipRemaining === true - ); - } - /** * @param {import("./workflow-parse.js").CompiledStep} parsed * @param {number} index - * @param {unknown} ctx + * @param {import("./step-result.js").StepResult} last * @param {string} runId * @param {import("pino").Logger} runLog * @param {string} reason */ - async function markStepSkipped(parsed, index, ctx, runId, runLog, reason) { + async function markStepSkipped(parsed, index, last, runId, runLog, reason) { const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; const step = await store.startStep({ runId, @@ -354,12 +352,12 @@ export function createRegistry(server) { }); const stepLog = runLog.child({ stepId: step.id, script }); stepLog.debug({ reason }, "step skipped"); - await store.finishStep(step.id, "skipped", ctx, reason); + await store.finishStep(step.id, "skipped", storedEnvelope(last), reason); } /** * @param {import("./workflow-parse.js").CompiledScripts} compiled - * @param {{ data?: unknown }} ctx + * @param {{ data?: unknown, context?: Record }} ctx * @param {string} runId * @param {import("pino").Logger} runLog * @param {string} key @@ -367,11 +365,17 @@ export function createRegistry(server) { * @param {number} depth */ async function runLinearSteps(compiled, ctx, runId, runLog, key, owner, depth) { - let next = ctx; + /** @type {import("./step-result.js").StepResult} */ + let last = { + output: ctx.data, + context: normalizeContext(ctx.context), + skipRemaining: false, + }; + let next = { data: ctx.data, context: last.context }; for (let i = 0; i < compiled.order.length; i++) { const index = compiled.order[i]; const parsed = compiled.steps[index]; - next = await runCompiledStep( + last = await runCompiledStep( parsed, next, index, @@ -381,13 +385,14 @@ export function createRegistry(server) { owner, depth, ); - if (isSkipRemaining(next)) { + next = chainCtx(last); + if (last.skipRemaining) { for (let j = i + 1; j < compiled.order.length; j++) { const laterIndex = compiled.order[j]; await markStepSkipped( compiled.steps[laterIndex], laterIndex, - next, + last, runId, runLog, "skipRemaining", @@ -396,12 +401,12 @@ export function createRegistry(server) { break; } } - return next; + return last; } /** * @param {import("./workflow-parse.js").CompiledScripts} compiled - * @param {{ data?: unknown }} ctx + * @param {{ data?: unknown, context?: Record }} ctx * @param {string} runId * @param {import("pino").Logger} runLog * @param {string} key @@ -412,30 +417,91 @@ export function createRegistry(server) { const triggerData = ctx.data; /** @type {Map} */ const outputsById = new Map(); - let last = ctx; + /** @type {Map} */ + const idToIndex = new Map(); + for (const step of compiled.steps) { + if (step.id) idToIndex.set(step.id, step.index); + } - for (const [orderIndex, stepIndex] of compiled.order.entries()) { - const parsed = compiled.steps[stepIndex]; - const data = mergeStepData(parsed, outputsById, triggerData); - last = await runCompiledStep( - parsed, - { ...ctx, data }, - orderIndex, - runId, - runLog, - key, - owner, - depth, - ); - if (parsed.id) { - outputsById.set(parsed.id, last); + const remaining = new Set(compiled.steps.map((s) => s.index)); + const completed = new Set(); + let context = normalizeContext(ctx.context); + /** @type {import("./step-result.js").StepResult} */ + let last = { + output: ctx.data, + context, + skipRemaining: false, + }; + let execIndex = 0; + + /** + * @param {import("./workflow-parse.js").CompiledStep} step + */ + function needsMet(step) { + if (step.needsKind === "none" || step.needs.length === 0) return true; + return step.needs.every((edge) => completed.has(idToIndex.get(edge.from))); + } + + while (remaining.size) { + const wave = compiled.steps + .filter((s) => remaining.has(s.index) && needsMet(s)) + .sort((a, b) => a.index - b.index); + if (wave.length === 0) { + throw new Error("Workflow has a cycle"); } - if (isSkipRemaining(last)) { - for (let j = orderIndex + 1; j < compiled.order.length; j++) { - const laterParsed = compiled.steps[compiled.order[j]]; + + const snapshot = { ...context }; + /** @type {Array<{ id: string, context: unknown }>} */ + const patches = []; + let skipRest = false; + + for (const parsed of wave) { + remaining.delete(parsed.index); + const index = execIndex; + execIndex += 1; + if (skipRest) { await markStepSkipped( - laterParsed, - j, + parsed, + index, + last, + runId, + runLog, + "skipRemaining", + ); + continue; + } + const data = mergeStepData(parsed, outputsById, triggerData); + last = await runCompiledStep( + parsed, + { data, context: { ...snapshot } }, + index, + runId, + runLog, + key, + owner, + depth, + ); + if (parsed.id) { + outputsById.set(parsed.id, last.output); + } + patches.push({ + id: parsed.id ?? String(parsed.index), + context: last.context, + }); + if (last.skipRemaining) skipRest = true; + } + + context = mergeContextWave(snapshot, patches); + last = { ...last, context }; + + if (skipRest) { + for (const later of compiled.steps.filter((s) => remaining.has(s.index))) { + remaining.delete(later.index); + const index = execIndex; + execIndex += 1; + await markStepSkipped( + later, + index, last, runId, runLog, @@ -444,7 +510,10 @@ export function createRegistry(server) { } break; } + + for (const parsed of wave) completed.add(parsed.index); } + return last; } @@ -483,13 +552,14 @@ export function createRegistry(server) { /** * @param {import("./workflow-parse.js").CompiledStep} parsed - * @param {{ data?: unknown, config?: unknown }} ctx + * @param {{ data?: unknown, context?: unknown, config?: unknown }} ctx * @param {number} index * @param {string} runId * @param {import("pino").Logger} runLog * @param {string} key * @param {string} owner * @param {number} depth + * @returns {Promise} */ async function runCompiledStep( parsed, @@ -503,6 +573,12 @@ export function createRegistry(server) { ) { const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; const config = parsed.config; + const incomingContext = normalizeContext(ctx.context); + const stepCtx = { + data: ctx.data, + context: incomingContext, + config, + }; const step = await store.startStep({ runId, index, @@ -512,29 +588,36 @@ export function createRegistry(server) { const stepLog = runLog.child({ stepId: step.id, script }); try { if (parsed.when) { - const whenResult = await evaluateJsonata(parsed.when, ctx); + const whenResult = await evaluateJsonata(parsed.when, stepCtx); if (!isJsonataTruthy(whenResult)) { stepLog.debug({ when: parsed.when }, "step skipped"); - await store.finishStep(step.id, "skipped", ctx, "when condition"); - return ctx; + const skipped = { + output: stepCtx.data, + context: incomingContext, + skipRemaining: false, + }; + await store.finishStep(step.id, "skipped", storedEnvelope(skipped), "when condition"); + return skipped; } } 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); + const value = await evaluateJsonata(parsed.expression, stepCtx); + const result = { + output: value, + context: incomingContext, + skipRemaining: false, + }; + await store.finishStep(step.id, "success", storedEnvelope(result)); return result; } - const result = await runScript(script, { ...ctx, config }, { + const raw = await runScript(script, stepCtx, { log: stepLog, workflowName: key, owner, $workflows: createWorkflowsApi(owner, key, runId, depth), }); - await store.finishStep(step.id, "success", result); + const result = normalizeStepResult(raw, incomingContext, script); + await store.finishStep(step.id, "success", storedEnvelope(result)); return result; } catch (err) { await store.finishStep(step.id, "failed", null, err); @@ -544,7 +627,7 @@ export function createRegistry(server) { /** * @param {string} key - * @param {{ data?: unknown }} context + * @param {{ data?: unknown, context?: unknown }} context * @param {{ type: string, detail?: string | null }} trigger * @param {{ * parentRunId?: string | null, @@ -580,8 +663,8 @@ export function createRegistry(server) { runLog.debug(detach ? "running workflow (detached)" : "running workflow"); const initialCtx = { - ...context, data: context.data ?? workflow.data ?? null, + context: normalizeContext(context.context), }; const execute = async () => { @@ -609,8 +692,8 @@ export function createRegistry(server) { depth, ); } - await store.finishRun(run.id, "success", ctx); - return { runId: run.id, status: "success", result: ctx }; + await store.finishRun(run.id, "success", storedEnvelope(ctx)); + return { runId: run.id, status: "success", result: storedEnvelope(ctx) }; } catch (err) { runLog.error({ err }, "workflow failed"); await store.finishRun(run.id, "failed", null, err); diff --git a/packages/server/scripts/fetch-binary.js b/packages/server/scripts/fetch-binary.js index 8712dda..abeb8a8 100644 --- a/packages/server/scripts/fetch-binary.js +++ b/packages/server/scripts/fetch-binary.js @@ -1,7 +1,15 @@ -function ensureDataObject(ctx) { - if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { - ctx.data = {}; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; } + return {}; +} + +function mergeData(data) { + if (data != null && typeof data === "object" && !Array.isArray(data)) { + return { ...data }; + } + return {}; } function filenameFromUrl(url) { @@ -38,8 +46,6 @@ function resolveUrl(ctx) { } async function fetchBinary(ctx) { - ensureDataObject(ctx); - const url = resolveUrl(ctx); if (!url) { throw new Error( @@ -55,27 +61,32 @@ async function fetchBinary(ctx) { const response = await $axios.get(url, { responseType: "arraybuffer" }); const file = Buffer.from(response.data ?? []); const contentType = String(response.headers?.["content-type"] ?? "application/octet-stream"); - const filename = ctx.config?.filename || ctx.data.filename || filenameFromUrl(url); + const filename = ctx.config?.filename || ctx.data?.filename || filenameFromUrl(url); - ctx.data[outputVar] = file; - ctx.data.filename = filename; - ctx.data.contentType = contentType; + const extra = { + [outputVar]: file, + filename, + contentType, + }; log.info( { outputVar, filename, contentType, length: file.length }, "fetch-binary: saved", ); - return ctx; + return { + output: { ...mergeData(ctx.data), ...extra }, + context: { ...passContext(ctx), ...extra }, + }; } fetchBinary.meta = { - description: "Download a binary URL into ctx.data as a Buffer", + description: "Download a binary URL and add the Buffer to output (keeps previous data fields)", previewConfigKey: "url", 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" }, + outputVar: { type: "string", default: "file", description: "output key for the Buffer" }, filename: { type: "string", required: false, description: "Override saved filename" }, }, input: { @@ -88,6 +99,11 @@ fetchBinary.meta = { filename: { type: "string" }, contentType: { type: "string" }, }, + context: { + 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" }, diff --git a/packages/server/scripts/fetch-html.js b/packages/server/scripts/fetch-html.js index 3b7be2e..93ac59e 100644 --- a/packages/server/scripts/fetch-html.js +++ b/packages/server/scripts/fetch-html.js @@ -1,10 +1,11 @@ import { parse } from "node-html-parser"; import jsonata from "jsonata"; -function ensureDataObject(ctx) { - if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { - ctx.data = {}; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; } + return {}; } function serializeElement(el) { @@ -43,8 +44,6 @@ async function fetchHtml(ctx) { throw new Error("fetch-html: ctx.config.outputVar is required when selector or jsonata is set"); } - ensureDataObject(ctx); - log.info({ url }, "fetch-html: fetching page"); const response = await $axios.get(url); log.info( @@ -52,19 +51,20 @@ async function fetchHtml(ctx) { "fetch-html: fetch complete", ); - ctx.data.httpResponse = String(response.data ?? ""); + /** @type {Record} */ + const output = { httpResponse: String(response.data ?? "") }; log.info( - { length: ctx.data.httpResponse.length }, + { length: String(output.httpResponse).length }, "fetch-html: saved httpResponse", ); if (selector) { log.info({ selector, outputVar }, "fetch-html: selecting elements"); - const root = parse(ctx.data.httpResponse); + const root = parse(String(output.httpResponse)); const elements = root.querySelectorAll(selector); - ctx.data[outputVar] = elements.map(serializeElement); + output[outputVar] = elements.map(serializeElement); log.info( - { outputVar, count: ctx.data[outputVar].length }, + { outputVar, count: output[outputVar].length }, "fetch-html: saved selector matches", ); } @@ -72,16 +72,15 @@ async function fetchHtml(ctx) { if (jsonataExpr) { log.info({ outputVar, jsonata: jsonataExpr }, "fetch-html: evaluating jsonata"); const expression = jsonata(jsonataExpr); - const input = ctx.data[outputVar]; - const result = await expression.evaluate(input); - ctx.data[outputVar] = result; + const result = await expression.evaluate(output[outputVar]); + output[outputVar] = result; log.info( { outputVar, result: previewValue(result) }, "fetch-html: saved jsonata result", ); } - return ctx; + return { output, context: { ...passContext(ctx), ...output } }; } fetchHtml.meta = { @@ -100,7 +99,10 @@ fetchHtml.meta = { }, input: {}, output: { - httpResponse: { type: "string", description: "Raw HTML" }, + httpResponse: { type: "string", description: "Raw HTML (overwritten when outputVar is httpResponse)" }, + }, + context: { + httpResponse: { type: "any", description: "Same keys as output" }, }, example: { data: {}, diff --git a/packages/server/scripts/fetch-http.js b/packages/server/scripts/fetch-http.js index 1246036..e276dd3 100644 --- a/packages/server/scripts/fetch-http.js +++ b/packages/server/scripts/fetch-http.js @@ -1,9 +1,10 @@ const HTTP_METHODS = ["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"]; -function ensureDataObject(ctx) { - if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { - ctx.data = {}; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; } + return {}; } function isPlainObject(value) { @@ -40,7 +41,6 @@ function normalizeHeaders(value, label) { headers[key] = String(val); continue; } - // Secret wrappers are unwrapped by $axios; do not String() them. if (typeof val === "object") { headers[key] = val; continue; @@ -88,8 +88,6 @@ async function fetchHttp(ctx) { const headers = resolveHeaders(ctx); const body = resolveBody(ctx); - ensureDataObject(ctx); - /** @type {Record} */ const request = { url, method, headers }; if (body.present && method !== "HEAD") { @@ -103,14 +101,16 @@ async function fetchHttp(ctx) { "fetch-http: fetch complete", ); - ctx.data.httpResponse = response.data; - ctx.data.httpStatus = response.status; + const output = { + httpResponse: response.data, + httpStatus: response.status, + }; log.info( - { status: ctx.data.httpStatus, length: responseSize(ctx.data.httpResponse) }, + { status: output.httpStatus, length: responseSize(output.httpResponse) }, "fetch-http: saved httpResponse", ); - return ctx; + return { output, context: { ...passContext(ctx), ...output } }; } fetchHttp.meta = { @@ -154,6 +154,10 @@ fetchHttp.meta = { httpResponse: { type: "any", description: "Response body as returned by the server" }, httpStatus: { type: "number", description: "HTTP status code" }, }, + context: { + httpResponse: { type: "any", description: "Same as output.httpResponse" }, + httpStatus: { type: "number", description: "Same as output.httpStatus" }, + }, example: { data: {}, config: { diff --git a/packages/server/scripts/fetch-rss-feed.js b/packages/server/scripts/fetch-rss-feed.js index afcb6a5..7e8c170 100644 --- a/packages/server/scripts/fetch-rss-feed.js +++ b/packages/server/scripts/fetch-rss-feed.js @@ -1,10 +1,11 @@ import rssParser from "rss-parser"; import jsonata from "jsonata"; -function ensureDataObject(ctx) { - if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { - ctx.data = {}; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; } + return {}; } function previewValue(value) { @@ -28,8 +29,6 @@ async function fetchRssFeed(ctx) { throw new Error("fetch-rss-feed: ctx.config.outputVar is required when jsonata is set"); } - ensureDataObject(ctx); - log.info({ url }, "fetch-rss-feed: fetching feed"); const parser = new rssParser({ customFields: { @@ -46,7 +45,8 @@ async function fetchRssFeed(ctx) { "fetch-rss-feed: fetch complete", ); - ctx.data.rssFeed = feed; + /** @type {Record} */ + const output = { rssFeed: feed }; log.info( { title: feed.title, itemCount }, "fetch-rss-feed: saved rssFeed", @@ -56,14 +56,14 @@ async function fetchRssFeed(ctx) { log.info({ outputVar, jsonata: jsonataExpr }, "fetch-rss-feed: evaluating jsonata"); const expression = jsonata(jsonataExpr); const result = await expression.evaluate(feed); - ctx.data[outputVar] = result; + output[outputVar] = result; log.info( { outputVar, result: previewValue(result) }, "fetch-rss-feed: saved jsonata result", ); } - return ctx; + return { output, context: { ...passContext(ctx), ...output } }; } fetchRssFeed.meta = { @@ -75,7 +75,7 @@ fetchRssFeed.meta = { outputVar: { type: "string", required: false, - description: "Required when jsonata is set; destination on ctx.data", + description: "Required when jsonata is set; key on output and context", }, jsonata: { type: "string", @@ -89,6 +89,9 @@ fetchRssFeed.meta = { output: { rssFeed: { type: "object", description: "Parsed RSS/Atom feed from rss-parser" }, }, + context: { + rssFeed: { type: "object", description: "Same as output.rssFeed (plus outputVar when set)" }, + }, example: { data: {}, config: { diff --git a/packages/server/scripts/fingerprint.js b/packages/server/scripts/fingerprint.js index 1baf860..522ed16 100644 --- a/packages/server/scripts/fingerprint.js +++ b/packages/server/scripts/fingerprint.js @@ -1,9 +1,17 @@ import jsonata from "jsonata"; -function ensureDataObject(ctx) { - if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { - ctx.data = {}; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; } + return {}; +} + +function mergeData(data) { + if (data != null && typeof data === "object" && !Array.isArray(data)) { + return { ...data }; + } + return {}; } async function fingerprint(ctx) { @@ -28,13 +36,14 @@ async function fingerprint(ctx) { maxAge: ctx.config?.maxAge, }); - ensureDataObject(ctx); - ctx.data.fingerprint = result.hash; - ctx.data.fingerprintChanged = result.changed; - ctx.data.fingerprintPrevious = result.previous; - ctx.data.fingerprintAt = result.changed ? result.at : result.previousAt; - ctx.data.fingerprintAge = result.ageMs; - ctx.data.fingerprintExpired = result.expired; + const extra = { + fingerprint: result.hash, + fingerprintChanged: result.changed, + fingerprintPrevious: result.previous, + fingerprintAt: result.changed ? result.at : result.previousAt, + fingerprintAge: result.ageMs, + fingerprintExpired: result.expired, + }; log.info( { @@ -46,17 +55,22 @@ async function fingerprint(ctx) { "fingerprint: result", ); + /** @type {{ output: Record, context: Record, skipRemaining?: true }} */ + const envelope = { + output: { ...mergeData(ctx.data), ...extra }, + context: { ...passContext(ctx), ...extra }, + }; if (!result.changed && skipRemaining) { - ctx.skipRemaining = true; + envelope.skipRemaining = true; } - - return ctx; + return envelope; } fingerprint.meta = { description: "Hash a value, compare it to the last stored fingerprint, and skip remaining steps when unchanged", previewConfigKey: "key", + reads: "ctx", config: { key: { type: "string", @@ -88,6 +102,14 @@ fingerprint.meta = { fingerprintAge: { type: "number", required: false, description: "Age in milliseconds" }, fingerprintExpired: { type: "boolean" }, }, + context: { + fingerprint: { type: "string" }, + fingerprintChanged: { type: "boolean" }, + fingerprintPrevious: { type: "string", required: false }, + fingerprintAt: { type: "string", required: false }, + fingerprintAge: { type: "number", required: false }, + fingerprintExpired: { type: "boolean" }, + }, example: { data: { item: { guid: "https://example.com/post-1" } }, config: { diff --git a/packages/server/scripts/get-current-time.js b/packages/server/scripts/get-current-time.js index ab3231c..ad34931 100644 --- a/packages/server/scripts/get-current-time.js +++ b/packages/server/scripts/get-current-time.js @@ -1,22 +1,30 @@ -// this script will get current time +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} -function getCurrentTime() { - return { - data: { - datetime: new Date().toISOString(), - processId: "1234" - } - } +function getCurrentTime(ctx) { + const output = { + datetime: new Date().toISOString(), + processId: "1234", + }; + return { output, context: { ...passContext(ctx), ...output } }; } getCurrentTime.meta = { - description: "Return the current time as ctx.data.datetime", + description: "Return the current time as output.datetime", config: {}, input: {}, output: { datetime: { type: "string", description: "ISO timestamp" }, processId: { type: "string" }, }, + context: { + datetime: { type: "string", description: "ISO timestamp" }, + processId: { type: "string" }, + }, example: { data: {}, config: {}, diff --git a/packages/server/scripts/get-secret.js b/packages/server/scripts/get-secret.js index 2212955..123a660 100644 --- a/packages/server/scripts/get-secret.js +++ b/packages/server/scripts/get-secret.js @@ -1,3 +1,17 @@ +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} + +function mergeData(data) { + if (data != null && typeof data === "object" && !Array.isArray(data)) { + return { ...data }; + } + return {}; +} + async function getSecret(ctx) { const name = ctx.config?.name; if (typeof name !== "string" || name.length === 0) { @@ -9,19 +23,16 @@ async function getSecret(ctx) { : name; const value = await $secrets.get(name); - const base = - ctx != null && typeof ctx === "object" && !Array.isArray(ctx) ? { ...ctx } : {}; - const data = - base.data != null && typeof base.data === "object" && !Array.isArray(base.data) - ? { ...base.data } - : {}; - data[as] = value; - return { ...base, data }; + const extra = { [as]: value }; + return { + output: { ...mergeData(ctx.data), ...extra }, + context: { ...passContext(ctx), ...extra }, + }; } getSecret.meta = { description: - "Load a named secret for this workflow owner into ctx.data. The value is wrapped and redacted in logs.", + "Load a named secret for this workflow owner onto output and context. The value is wrapped and redacted in logs.", previewConfigKey: "name", config: { name: { @@ -32,16 +43,12 @@ getSecret.meta = { as: { type: "string", required: false, - description: "ctx.data field to write (defaults to name)", + description: "Field name on output and context (defaults to name)", }, }, input: {}, - output: { - data: { - type: "object", - description: "Previous ctx.data plus the retrieved Secret at [as]", - }, - }, + output: {}, + context: {}, example: { data: {}, config: { name: "ntfy_token", as: "ntfyToken" }, diff --git a/packages/server/scripts/jsonata.js b/packages/server/scripts/jsonata.js index 5948104..b05cf24 100644 --- a/packages/server/scripts/jsonata.js +++ b/packages/server/scripts/jsonata.js @@ -1,26 +1,34 @@ import jsonata from "jsonata"; async 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 = await expression.evaluate(ctx.data); - log.info({result}, "jsonata: expression result"); - return result; + log.info({ ctx }, "jsonata: context"); + log.info("jsonata: evaluating expression %s", ctx.config.expression); + const expression = jsonata(ctx.config.expression); + const result = await expression.evaluate(ctx); + log.info({ result }, "jsonata: expression result"); + return { + output: result, + context: + ctx.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context) + ? ctx.context + : {}, + }; } jsonataFn.meta = { - description: "Evaluate a JSONata expression against ctx.data and return the result as the next context", + description: + "Evaluate a JSONata expression against the full ctx (data, context, config); the result is the next data", previewConfigKey: "expression", + reads: "ctx", config: { - expression: { type: "string", required: true, description: "JSONata expression" }, + expression: { type: "string", required: true, description: "JSONata expression against ctx" }, }, input: {}, output: {}, example: { data: { title: "Hello", url: "https://example.com" }, config: { - expression: '{"data": {"title": title, "message": title, "attach": url}}', + expression: '{"title": data.title, "message": data.title, "attach": data.url}', }, }, }; diff --git a/packages/server/scripts/ntfy.js b/packages/server/scripts/ntfy.js index 408c600..3c589f5 100644 --- a/packages/server/scripts/ntfy.js +++ b/packages/server/scripts/ntfy.js @@ -1,5 +1,12 @@ import jsonata from "jsonata"; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} + function ntfyHeaders(ctx) { const headers = {}; @@ -72,11 +79,14 @@ async function ntfy(ctx) { "ntfy: skipped, fingerprint unchanged", ); return { - sent: "false", - skipped: true, - fingerprint: checked.hash, - fingerprintAt: checked.previousAt, - fingerprintAge: checked.ageMs, + output: { + sent: "false", + skipped: true, + fingerprint: checked.hash, + fingerprintAt: checked.previousAt, + fingerprintAge: checked.ageMs, + }, + context: passContext(ctx), }; } fp = checked; @@ -126,7 +136,7 @@ async function ntfy(ctx) { sent.fingerprint = stored.hash; sent.fingerprintAt = stored.at; } - return sent; + return { output: sent, context: passContext(ctx) }; } ntfy.meta = { @@ -147,7 +157,7 @@ ntfy.meta = { fingerprintJsonata: { type: "string", required: false, - description: "JSONata against ctx; default hashes title, message, attach, filename, contentType", + description: "JSONata against full ctx; default hashes data title, message, attach, filename, contentType", }, fingerprintMaxAge: { type: "string", @@ -163,8 +173,9 @@ ntfy.meta = { filename: { type: "string", required: false }, contentType: { type: "string", required: false }, }, + reads: "ctx", output: { - sent: { type: "string", description: "Replaces the workflow context with { sent: \"true\" }" }, + sent: { type: "string", description: '"true" when sent, "false" when fingerprint skipped' }, }, example: { data: { title: "Hello", message: "Hello from jerapah-flow" }, diff --git a/packages/server/scripts/render-template.js b/packages/server/scripts/render-template.js index 2c066a5..23b2e2d 100644 --- a/packages/server/scripts/render-template.js +++ b/packages/server/scripts/render-template.js @@ -70,9 +70,15 @@ async function renderTemplate(ctx) { const text = htmlToText(html); return { - html, - text, - template: templateName, + output: { + html, + text, + template: templateName, + }, + context: + ctx.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context) + ? ctx.context + : {}, }; } diff --git a/packages/server/scripts/send-email.js b/packages/server/scripts/send-email.js index 9d090e5..8a9c612 100644 --- a/packages/server/scripts/send-email.js +++ b/packages/server/scripts/send-email.js @@ -218,13 +218,19 @@ async function sendEmail(ctx) { log.info({ messageId: info.messageId }, "send-email: message sent"); return { - sent: true, - messageId: info.messageId ?? null, - from, - to, - cc: cc ?? null, - bcc: bcc ?? null, - subject, + output: { + sent: true, + messageId: info.messageId ?? null, + from, + to, + cc: cc ?? null, + bcc: bcc ?? null, + subject, + }, + context: + ctx.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context) + ? ctx.context + : {}, }; } diff --git a/packages/server/scripts/trigger-workflow.js b/packages/server/scripts/trigger-workflow.js index 2f478e2..0639472 100644 --- a/packages/server/scripts/trigger-workflow.js +++ b/packages/server/scripts/trigger-workflow.js @@ -1,5 +1,12 @@ import jsonata from "jsonata"; +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} + async function triggerWorkflow(ctx) { const name = ctx.config?.name; if (typeof name !== "string" || name.length === 0) { @@ -9,27 +16,23 @@ async function triggerWorkflow(ctx) { let data = ctx.data; const expression = ctx.config?.expression; if (typeof expression === "string" && expression.length > 0) { - const result = jsonata(expression).evaluate(ctx.data); + const result = jsonata(expression).evaluate(ctx); 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 }; + return { + output: { name, runId: started?.runId ?? null }, + context: passContext(ctx), + }; } 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.", + "Fire-and-forget another workflow by YAML name (same owner). Destination must declare triggers: [{ type: workflow }]. Optionally reshape the destination input with JSONata against full ctx.", previewConfigKey: "name", tags: ["trigger"], + reads: "ctx", config: { name: { type: "string", @@ -40,15 +43,13 @@ triggerWorkflow.meta = { type: "string", required: false, description: - "Optional JSONata expression evaluated against ctx.data; result becomes the destination run input", + "Optional JSONata expression evaluated against ctx; 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", - }, + name: { type: "string", description: "Destination workflow name" }, + runId: { type: "string", description: "Started run id (null if the destination failed to start)" }, }, example: { data: { @@ -58,7 +59,7 @@ triggerWorkflow.meta = { config: { name: "notify-comic", expression: - '{ "title": title, "message": title, "attach": url }', + '{ "title": data.title, "message": data.title, "attach": data.url }', }, }, }; diff --git a/packages/server/src/api/scripts.js b/packages/server/src/api/scripts.js index 43772b4..af9c5ef 100644 --- a/packages/server/src/api/scripts.js +++ b/packages/server/src/api/scripts.js @@ -6,6 +6,7 @@ import { } from "../../script-sandbox.js"; import * as fsStore from "../../fs-store.js"; import { createDryRunLogger, safeSerialize } from "./dry-run-logger.js"; +import { normalizeStepResult } from "../../step-result.js"; /** * @param {{ referencedScripts: () => Set }} registry @@ -102,7 +103,7 @@ export default function scriptsPluginFactory(registry) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } - const body = /** @type {{ content?: string, data?: unknown, config?: unknown, owner?: string }} */ ( + const body = /** @type {{ content?: string, data?: unknown, context?: unknown, config?: unknown, owner?: string }} */ ( req.body ?? {} ); if (typeof body.content !== "string") { @@ -118,8 +119,13 @@ export default function scriptsPluginFactory(registry) { } } + const incomingContext = + body.context != null && typeof body.context === "object" && !Array.isArray(body.context) + ? body.context + : {}; const ctx = { data: body.data ?? null, + context: incomingContext, config: body.config ?? null, }; @@ -132,10 +138,13 @@ export default function scriptsPluginFactory(registry) { workflowName: "dry-run", owner, }); - const output = await fn(ctx); + const raw = await fn(ctx); + const result = normalizeStepResult(raw, incomingContext, name); return { status: "success", - output: safeSerialize(output), + output: safeSerialize(result.output), + context: safeSerialize(result.context), + skipRemaining: result.skipRemaining, error: null, logs, durationMs: Date.now() - started, @@ -147,6 +156,8 @@ export default function scriptsPluginFactory(registry) { return { status: "failed", output: null, + context: null, + skipRemaining: false, error: err instanceof Error ? err.message : String(err), logs, durationMs: Date.now() - started, diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index b4b8da4..1be33bf 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -32,7 +32,7 @@ function scriptNames(workflow) { for (const raw of workflow.scripts ?? []) { try { const parsed = parseScriptStep(raw); - names.push(parsed.kind === "set" ? `set:${parsed.as}` : parsed.script); + names.push(parsed.kind === "set" ? "set" : parsed.script); } catch { names.push(null); } diff --git a/packages/server/step-result.js b/packages/server/step-result.js new file mode 100644 index 0000000..af7084f --- /dev/null +++ b/packages/server/step-result.js @@ -0,0 +1,141 @@ +/** + * Step return contract: `{ output, context?, skipRemaining? }`. + */ + +/** + * @param {unknown} value + */ +export function isPlainObject(value) { + return value != null && typeof value === "object" && !Array.isArray(value); +} + +/** + * @param {unknown} value + * @returns {Record} + */ +export function normalizeContext(value) { + return isPlainObject(value) ? /** @type {Record} */ (value) : {}; +} + +/** + * @typedef {{ + * output: unknown, + * context: Record, + * skipRemaining: boolean, + * }} StepResult + */ + +/** + * @param {unknown} raw + * @param {unknown} incomingContext + * @param {string} [label] + * @returns {StepResult} + */ +export function normalizeStepResult(raw, incomingContext, label = "script") { + const incoming = normalizeContext(incomingContext); + if (!isPlainObject(raw)) { + const got = raw == null ? String(raw) : typeof raw; + throw new Error(`${label} must return { output, context }. Got ${got}`); + } + if (!("output" in raw) && ("data" in raw || "config" in raw)) { + throw new Error( + `${label} must return { output, context }. Returning the full ctx is no longer valid.`, + ); + } + + let context = incoming; + if ("context" in raw) { + if (raw.context == null) { + context = incoming; + } else if (!isPlainObject(raw.context)) { + throw new Error(`${label} returned a non-object context`); + } else { + context = /** @type {Record} */ (raw.context); + } + } + + return { + output: "output" in raw ? raw.output : null, + context, + skipRemaining: raw.skipRemaining === true, + }; +} + +/** + * Persistable envelope (omit skipRemaining unless set). + * @param {StepResult} result + */ +export function storedEnvelope(result) { + /** @type {{ output: unknown, context: Record, skipRemaining?: true }} */ + const out = { output: result.output, context: result.context }; + if (result.skipRemaining) out.skipRemaining = true; + return out; +} + +/** + * Next step ctx (without config). + * @param {StepResult} result + */ +export function chainCtx(result) { + return { data: result.output, context: result.context }; +} + +/** + * Shallow-merge sibling context diffs vs a shared snapshot. + * Fails if two siblings change the same key. + * + * @param {unknown} snapshot + * @param {Array<{ id: string, context: unknown }>} patches + * @returns {Record} + */ +export function mergeContextWave(snapshot, patches) { + const base = normalizeContext(snapshot); + /** @type {Map} */ + const writers = new Map(); + /** @type {Map} */ + const changes = new Map(); + + for (const patch of patches) { + const who = patch.id; + const next = normalizeContext(patch.context); + const keys = new Set([...Object.keys(base), ...Object.keys(next)]); + for (const key of keys) { + const inBase = Object.prototype.hasOwnProperty.call(base, key); + const inNext = Object.prototype.hasOwnProperty.call(next, key); + if (inBase && inNext && Object.is(base[key], next[key])) continue; + if (!inBase && !inNext) continue; + if (inBase && !inNext) { + rememberChange(writers, changes, key, who, { deleted: true }); + continue; + } + if (!inBase || !Object.is(base[key], next[key])) { + rememberChange(writers, changes, key, who, { value: next[key] }); + } + } + } + + const merged = { ...base }; + for (const [key, spec] of changes) { + if ("deleted" in spec) delete merged[key]; + else merged[key] = spec.value; + } + return merged; +} + +/** + * @param {Map} writers + * @param {Map} changes + * @param {string} key + * @param {string} who + * @param {{ deleted: true } | { value: unknown }} spec + */ +function rememberChange(writers, changes, key, who, spec) { + const previous = writers.get(key); + if (previous != null && previous !== who) { + throw new Error( + `DAG context key conflict: "${key}" written by "${previous}" and "${who}"`, + ); + } + writers.set(key, who); + changes.set(key, spec); +} diff --git a/packages/server/test/fingerprint-smoke.js b/packages/server/test/fingerprint-smoke.js index 3331cd3..47e5b60 100644 --- a/packages/server/test/fingerprint-smoke.js +++ b/packages/server/test/fingerprint-smoke.js @@ -92,7 +92,7 @@ export default async function (ctx) { const scriptOut = await runScriptSource( "fingerprint-smoke.js", script, - { data: { n: 1 } }, + { data: { n: 1 }, context: {} }, { log, workflowName: ns }, ); if (!scriptOut.changed) throw new Error("sandbox $fingerprint.claim should be new"); diff --git a/packages/server/workflow-parse.js b/packages/server/workflow-parse.js index 0aa9d5b..059b707 100644 --- a/packages/server/workflow-parse.js +++ b/packages/server/workflow-parse.js @@ -1,8 +1,5 @@ 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"; /** @@ -12,7 +9,6 @@ export const SET_STEP_SCRIPT = "set"; * script: string, * config: unknown | null, * expression?: undefined, - * as?: undefined, * id: string | null, * needsKind: "none" | "list" | "map", * needs: NeedEdge[], @@ -21,11 +17,10 @@ export const SET_STEP_SCRIPT = "set"; * @typedef {{ * kind: "set", * script: typeof SET_STEP_SCRIPT, - * config: { expression: string, as: string }, + * config: { expression: string }, * expression: string, - * as: string, * id: string | null, - * needsKind: "none", + * needsKind: "none" | "list" | "map", * needs: NeedEdge[], * when: string | null, * }} ParsedSetStep @@ -163,9 +158,6 @@ export function compileWorkflowScripts(scripts) { 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})`); } @@ -227,13 +219,10 @@ export function compileWorkflowScripts(scripts) { } /** - * Prefer a script return's `.data`; otherwise treat the whole return as data. + * Upstream DAG output stored in outputsById (already the pipe value). * @param {unknown} result */ -export function extractStepData(result) { - if (result && typeof result === "object" && !Array.isArray(result) && "data" in result) { - return result.data; - } +export function extractStepOutput(result) { return result ?? null; } @@ -248,12 +237,12 @@ export function mergeStepData(step, outputsById, triggerData) { return triggerData; } if (step.needsKind === "list" && step.needs.length === 1) { - return extractStepData(outputsById.get(step.needs[0].from)); + return extractStepOutput(outputsById.get(step.needs[0].from)); } /** @type {Record} */ const data = {}; for (const { alias, from } of step.needs) { - data[alias] = extractStepData(outputsById.get(from)); + data[alias] = extractStepOutput(outputsById.get(from)); } return data; } @@ -263,34 +252,24 @@ export function mergeStepData(step, outputsById, triggerData) { * @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"); + const { needsKind, needs } = parseNeeds(step.needs); return { kind: "set", script: SET_STEP_SCRIPT, - config: { expression, as }, + config: { expression }, expression, - as, id: parseOptionalId(step.id), - needsKind: "none", - needs: [], + needsKind, + needs, when: parseWhen(step.when), }; } diff --git a/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml b/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml index c0161a2..dbc554e 100644 --- a/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml +++ b/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml @@ -10,7 +10,11 @@ scripts: - script: jsonata.js config: expression: | - {"data": {"title": httpResponse.title, "message": httpResponse.title, "attach": httpResponse.url}} + { + "title": data.httpResponse.title, + "message": data.httpResponse.title, + "attach": data.httpResponse.url + } - fetch-binary.js - script: ntfy.js config: diff --git a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml index 32b05b4..b305130 100644 --- a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml +++ b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml @@ -13,11 +13,9 @@ scripts: config: expression: | { - "data": { - "title": item.title, - "message": item.contentSnippet & "\n" & item.link, - "attach": item.mediaContent.url ? item.mediaContent.url : (item.mediaContent.$ ? item.mediaContent.$.url : undefined) - } + "title": data.item.title, + "message": data.item.contentSnippet & "\n" & data.item.link, + "attach": data.item.mediaContent.url ? data.item.mediaContent.url : (data.item.mediaContent.$ ? data.item.mediaContent.$.url : undefined) } - script: ntfy.js config: diff --git a/packages/server/workflows/default/send-gmail.yaml b/packages/server/workflows/default/send-gmail.yaml index 0bddf4d..3673194 100644 --- a/packages/server/workflows/default/send-gmail.yaml +++ b/packages/server/workflows/default/send-gmail.yaml @@ -14,18 +14,16 @@ scripts: config: expression: | { - "data": { - "to": $exists(to) ? to : "nasyarobby@gmail.com", - "from": from, - "subject": subject, - "text": $exists(text) ? text : ($exists(body) ? body : message), - "html": html, - "cc": cc, - "bcc": bcc, - "replyTo": replyTo, - "priority": priority, - "headers": headers - } + "to": $exists(data.to) ? data.to : "nasyarobby@gmail.com", + "from": data.from, + "subject": data.subject, + "text": $exists(data.text) ? data.text : ($exists(data.body) ? data.body : data.message), + "html": data.html, + "cc": data.cc, + "bcc": data.bcc, + "replyTo": data.replyTo, + "priority": data.priority, + "headers": data.headers } - script: send-email.js config: diff --git a/packages/server/workflows/default/test.yaml b/packages/server/workflows/default/test.yaml index 871e1b7..4c1cb48 100644 --- a/packages/server/workflows/default/test.yaml +++ b/packages/server/workflows/default/test.yaml @@ -3,7 +3,7 @@ scripts: - get-current-time.js - script: jsonata.js config: - expression: '{"data": {"message": datetime & " " & processId}}' + expression: '{"message": data.datetime & " " & data.processId}' - ntfy.js triggers: - type: HTTP diff --git a/packages/server/workflows/default/time-and-comic-to-ntfy.yaml b/packages/server/workflows/default/time-and-comic-to-ntfy.yaml index 640e683..4268411 100644 --- a/packages/server/workflows/default/time-and-comic-to-ntfy.yaml +++ b/packages/server/workflows/default/time-and-comic-to-ntfy.yaml @@ -19,11 +19,11 @@ scripts: needs: [time, comic] config: expression: | - {"data": { - "title": comic.httpResponse.title, - "message": time.datetime & " " & comic.httpResponse.title, - "attach": comic.httpResponse.url - }} + { + "title": data.comic.httpResponse.title, + "message": data.time.datetime & " " & data.comic.httpResponse.title, + "attach": data.comic.httpResponse.url + } - id: notify script: ntfy.js diff --git a/packages/server/workflows/default/time-to-ntfy-example.yaml b/packages/server/workflows/default/time-to-ntfy-example.yaml index 830c4e5..683d505 100644 --- a/packages/server/workflows/default/time-to-ntfy-example.yaml +++ b/packages/server/workflows/default/time-to-ntfy-example.yaml @@ -2,10 +2,10 @@ name: time to ntfy example description: | this workflow will send a message to ntfy with the current time scripts: - - get-current-time.js + - script: get-current-time.js - script: ntfy.js config: - url: https://ntfy.sh/scrunner + url: https://n.0dev.web.id/system triggers: - type: HTTP method: POST diff --git a/packages/server/workflows/default/track.yaml b/packages/server/workflows/default/track.yaml index e2f21a7..5593b67 100644 --- a/packages/server/workflows/default/track.yaml +++ b/packages/server/workflows/default/track.yaml @@ -16,10 +16,8 @@ scripts: config: expression: | { - "data": { - "title": httpResponse.customer_name, - "message": httpResponse.customer_name & " : " & httpResponse.status - } + "title": data.httpResponse.customer_name, + "message": data.httpResponse.customer_name & " : " & data.httpResponse.status } - script: ntfy.js config: diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index ec87e50..72f052c 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -107,11 +107,12 @@ export function useDeleteScript() { export function useDryRunScript() { return useMutation({ - mutationFn: async ({ name, content, data, config, owner }) => + mutationFn: async ({ name, content, data, context, config, owner }) => ( await api.post(`/scripts/${encodeURIComponent(name)}/dry-run`, { content, data, + context, config, owner, }) diff --git a/packages/web/src/components/ScriptMetaPanel.jsx b/packages/web/src/components/ScriptMetaPanel.jsx index 8bdea86..2d915b2 100644 --- a/packages/web/src/components/ScriptMetaPanel.jsx +++ b/packages/web/src/components/ScriptMetaPanel.jsx @@ -1,4 +1,4 @@ -import { scriptTags } from "../lib/script.js"; +import { scriptTags, scriptReadsCtx } from "../lib/script.js"; import { TagBadge } from "./TagBadge.jsx"; function FieldTable({ title, fields }) { @@ -64,17 +64,21 @@ export function ScriptMetaPanel({ meta, metaError, className = "" }) { return (
{meta.description ?

{meta.description}

: null} - {tags.length ? ( + {(tags.length || scriptReadsCtx(meta)) ? (
{tags.map((tag) => ( ))} + {scriptReadsCtx(meta) ? ( + reads: ctx + ) : null}
) : null} -
+
+
); diff --git a/packages/web/src/components/workflow/AddScriptDialog.jsx b/packages/web/src/components/workflow/AddScriptDialog.jsx index 8ed261c..8718f97 100644 --- a/packages/web/src/components/workflow/AddScriptDialog.jsx +++ b/packages/web/src/components/workflow/AddScriptDialog.jsx @@ -49,7 +49,9 @@ export function AddScriptDialog({ open, onClose, onPick }) { }} > Set (JSONata) - Assign a JSONata result to a variable + + JSONata over ctx; result becomes the next step’s data + {isLoading ? ( diff --git a/packages/web/src/components/workflow/ScriptCard.jsx b/packages/web/src/components/workflow/ScriptCard.jsx index 9aae03b..7148e36 100644 --- a/packages/web/src/components/workflow/ScriptCard.jsx +++ b/packages/web/src/components/workflow/ScriptCard.jsx @@ -42,10 +42,10 @@ export function ScriptCard({ step.id && otherSteps.some((s) => s.id === step.id && s.uiId !== step.uiId); const preview = previewConfigValue(step.config, meta?.previewConfigKey); const previewFull = configValueText(step.config, meta?.previewConfigKey); - const baseName = step.kind === "set" ? `set:${step.as || "…"}` : step.script || "untitled"; + const baseName = step.kind === "set" ? "set" : step.script || "untitled"; const titleFull = previewFull ? `${baseName} (${previewFull})` : baseName; const setConfig = - step.kind === "set" ? { as: step.as ?? "", expression: step.expression ?? "" } : null; + step.kind === "set" ? { expression: step.expression ?? "" } : null; return (
{meta.description}

) : null} {expanded && step.kind === "set" ? ( -

Assign a JSONata result onto the context

+

+ JSONata over ctx; result becomes the next step’s data +

) : null} {step.kind === "script" ? ( @@ -142,14 +144,11 @@ export function ScriptCard({ <> {step.kind === "set" ? (
- - onChange({ ...step, as: e.target.value })} + -