diff --git a/.gitignore b/.gitignore index 2e4d176..943a64c 100644 --- a/.gitignore +++ b/.gitignore @@ -5,3 +5,8 @@ logs/ *.db-* .DS_Store packages/web/dist + +# Personal/local scripts and workflows (not for the repo) +packages/server/scripts/dev-* +packages/server/workflows/**/dev-* +debug-*.js diff --git a/README.md b/README.md index c9d420a..2786685 100644 --- a/README.md +++ b/README.md @@ -1,11 +1,13 @@ -# scrunner +# JerapahFlow + +![JerapahFlow](packages/web/src/theme/brand/wordmark.png) Workflow runner with a sandboxed script engine, SQLite run history, and an admin UI. ## Packages -- `packages/server` — Fastify runner, HTTP/cron triggers, admin REST API -- `packages/web` — React admin UI (Vite, DaisyUI, React Query) +- `@jerapah-flow/server` (`packages/server`) — Fastify runner, HTTP/cron triggers, admin REST API +- `@jerapah-flow/web` (`packages/web`) — React admin UI (Vite, DaisyUI, React Query) ## Setup @@ -19,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 | @@ -34,12 +61,12 @@ The first account created becomes **admin**. Later accounts are created from Use | Variable | Default | Notes | |---|---|---| -| `SCRUNNER_JWT_SECRET` | `scrunner-dev-secret` (dev only) | **Required in production** | -| `SCRUNNER_SECRETS_KEY` | `scrunner-dev-secrets-key` (dev only) | Master key for named secrets. **Required in production**. Changing it makes existing secrets unreadable. 64 hex chars are used as a raw AES-256 key; any other string is derived with scrypt. | -| `SCRUNNER_DB_PATH` | `packages/server/data/scrunner.db` | SQLite file | -| `SCRUNNER_LOG_LEVEL` | `debug` | Pino level | -| `SCRUNNER_RETENTION_DAYS` | `30` | Run history prune | -| `SCRUNNER_CORS_ORIGIN` | `http://localhost:5173` | Vite origin in dev | +| `JERAPAH_FLOW_JWT_SECRET` | `scrunner-dev-secret` (dev only) | **Required in production**. `SCRUNNER_JWT_SECRET` is still accepted. | +| `JERAPAH_FLOW_SECRETS_KEY` | `scrunner-dev-secrets-key` (dev only) | Master key for named secrets. **Required in production**. Changing it makes existing secrets unreadable. 64 hex chars are used as a raw AES-256 key; any other string is derived with scrypt. `SCRUNNER_SECRETS_KEY` is still accepted. | +| `JERAPAH_FLOW_DB_PATH` | `packages/server/data/jerapah-flow.db` (or existing `scrunner.db`) | SQLite file. `SCRUNNER_DB_PATH` is still accepted. | +| `JERAPAH_FLOW_LOG_LEVEL` | `debug` | Pino level | +| `JERAPAH_FLOW_RETENTION_DAYS` | `30` | Run history prune | +| `JERAPAH_FLOW_CORS_ORIGIN` | `http://localhost:5173` | Vite origin in dev | | `PORT` | `9000` | HTTP port | | `NODE_ENV` | — | Set `production` for secure cookies | @@ -48,7 +75,7 @@ The first account created becomes **admin**. Later accounts are created from Use ```bash pnpm install pnpm build -SCRUNNER_JWT_SECRET=... SCRUNNER_SECRETS_KEY=... NODE_ENV=production pnpm start +JERAPAH_FLOW_JWT_SECRET=... JERAPAH_FLOW_SECRETS_KEY=... NODE_ENV=production pnpm start ``` The server serves `packages/web/dist` when that folder exists. diff --git a/package.json b/package.json index 040fda0..d001c21 100644 --- a/package.json +++ b/package.json @@ -1,16 +1,16 @@ { - "name": "scrunner", + "name": "jerapah-flow", "version": "1.0.0", "private": true, "description": "Script/workflow runner with admin UI", "type": "module", "scripts": { - "dev": "pnpm --filter @scrunner/server --filter @scrunner/web --parallel dev", - "dev:server": "pnpm --filter @scrunner/server dev", - "dev:web": "pnpm --filter @scrunner/web dev", - "build": "pnpm --filter @scrunner/web build", - "start": "pnpm --filter @scrunner/server start", - "migrate": "pnpm --filter @scrunner/server migrate" + "dev": "pnpm --filter @jerapah-flow/server --filter @jerapah-flow/web --parallel dev", + "dev:server": "pnpm --filter @jerapah-flow/server dev", + "dev:web": "pnpm --filter @jerapah-flow/web dev", + "build": "pnpm --filter @jerapah-flow/web build", + "start": "pnpm --filter @jerapah-flow/server start", + "migrate": "pnpm --filter @jerapah-flow/server migrate" }, "packageManager": "pnpm@10.25.0", "pnpm": { diff --git a/packages/server/config-refs.js b/packages/server/config-refs.js new file mode 100644 index 0000000..aa3c81b --- /dev/null +++ b/packages/server/config-refs.js @@ -0,0 +1,150 @@ +import { coerceCredentialString } from "./http-trigger-auth.js"; +import { assertSecretName, getSecretPlaintext } from "./secrets-store.js"; +import { isSecret } from "./secret-value.js"; +import { assertVariableName, getVariablePlain } from "./variables-store.js"; + +const PREFIXES = [ + { kind: "context", prefix: "$CONTEXT_" }, + { kind: "secret", prefix: "$SECRET_" }, + { kind: "var", prefix: "$VAR_" }, +]; + +/** + * @typedef {{ kind: "secret" | "context" | "var", name: string, raw: string }} ConfigRef + * @typedef {{ owner: string, workflowKey: string, context?: unknown }} ConfigRefCtx + */ + +/** + * Parse a whole-value config placeholder. Returns null for literals. + * @param {unknown} value + * @returns {ConfigRef | null} + */ +export function parseConfigRef(value) { + if (typeof value !== "string") return null; + const trimmed = value.trim(); + for (const { kind, prefix } of PREFIXES) { + if (trimmed.startsWith(prefix)) { + return { kind, name: trimmed.slice(prefix.length), raw: trimmed }; + } + } + return null; +} + +/** + * Walk config (objects/arrays) and replace whole-value `$SECRET_` / `$CONTEXT_` / `$VAR_` + * strings. Does not walk trigger data. + * + * @param {unknown} value + * @param {ConfigRefCtx} ctx + * @param {WeakSet} [seen] + * @returns {Promise} + */ +export async function resolveConfigRefs(value, ctx, seen = new WeakSet()) { + if (typeof value === "string") { + return resolveStringRef(value, ctx); + } + if (value == null || typeof value !== "object") { + return value; + } + if (seen.has(value)) return value; + seen.add(value); + + if (Array.isArray(value)) { + const out = []; + for (const item of value) { + out.push(await resolveConfigRefs(item, ctx, seen)); + } + return out; + } + + /** @type {Record} */ + const out = {}; + for (const [key, child] of Object.entries(value)) { + out[key] = await resolveConfigRefs(child, ctx, seen); + } + return out; +} + +/** + * @param {string} value + * @param {ConfigRefCtx} ctx + * @returns {Promise} + */ +async function resolveStringRef(value, ctx) { + const ref = parseConfigRef(value); + if (!ref) return value; + + if (ref.kind === "secret") { + return resolveSecretRef(ref, ctx); + } + if (ref.kind === "var") { + return resolveVarRef(ref, ctx); + } + return resolveContextRef(ref, ctx); +} + +/** + * @param {ConfigRef} ref + * @param {ConfigRefCtx} ctx + * @returns {Promise} + */ +async function resolveSecretRef(ref, ctx) { + try { + assertSecretName(ref.name); + } catch { + throw new Error(`config ref ${ref.raw}: invalid secret name`); + } + const plaintext = await getSecretPlaintext(ctx.owner, ref.name); + if (plaintext == null) { + throw new Error(`config ref ${ref.raw}: secret "${ref.name}" not found`); + } + return plaintext; +} + +/** + * @param {ConfigRef} ref + * @param {ConfigRefCtx} ctx + * @returns {Promise} + */ +async function resolveVarRef(ref, ctx) { + if (ref.name.length === 0) { + throw new Error(`config ref ${ref.raw}: empty variable name`); + } + try { + assertVariableName(ref.name); + } catch { + throw new Error(`config ref ${ref.raw}: invalid variable name`); + } + const value = await getVariablePlain(ctx.owner, ref.name); + if (value == null) { + throw new Error(`config ref ${ref.raw}: variable "${ref.name}" not found`); + } + return value; +} + +/** + * @param {ConfigRef} ref + * @param {ConfigRefCtx} ctx + * @returns {string} + */ +function resolveContextRef(ref, ctx) { + if (ref.name.length === 0) { + throw new Error(`config ref ${ref.raw}: empty context key`); + } + const bag = + ctx.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context) + ? /** @type {Record} */ (ctx.context) + : {}; + if (!Object.prototype.hasOwnProperty.call(bag, ref.name)) { + throw new Error(`config ref ${ref.raw}: context "${ref.name}" not found`); + } + const raw = bag[ref.name]; + if (isSecret(raw)) { + return raw.reveal(); + } + const coerced = coerceCredentialString(raw); + if (coerced == null) { + throw new Error(`config ref ${ref.raw}: context "${ref.name}" is not a scalar`); + } + return coerced; +} diff --git a/packages/server/fs-store.js b/packages/server/fs-store.js index 23061d1..2f2c842 100644 --- a/packages/server/fs-store.js +++ b/packages/server/fs-store.js @@ -55,11 +55,42 @@ export function writeScript(name, content) { fs.writeFileSync(path.join(SCRIPTS_DIR, name), content, "utf8"); } +/** + * Icon next to the script: `fetch-html.js` → `fetch-html.png` or `.jpg`. + * @returns {{ filePath: string, contentType: string } | null} + */ +export function resolveScriptIcon(name) { + assertScriptName(name); + const base = name.slice(0, -3); + for (const { ext, contentType } of [ + { ext: "png", contentType: "image/png" }, + { ext: "jpg", contentType: "image/jpeg" }, + ]) { + const filePath = path.join(SCRIPTS_DIR, `${base}.${ext}`); + if (fs.existsSync(filePath) && fs.statSync(filePath).isFile()) { + return { filePath, contentType }; + } + } + return null; +} + +export function scriptHasIcon(name) { + return resolveScriptIcon(name) != null; +} + export function deleteScript(name) { assertScriptName(name); const filePath = path.join(SCRIPTS_DIR, name); if (!fs.existsSync(filePath)) return false; + const icon = resolveScriptIcon(name); fs.unlinkSync(filePath); + if (icon) { + try { + fs.unlinkSync(icon.filePath); + } catch { + // ignore missing icon + } + } return true; } diff --git a/packages/server/knexfile.js b/packages/server/knexfile.js index eed0814..495f345 100644 --- a/packages/server/knexfile.js +++ b/packages/server/knexfile.js @@ -2,7 +2,12 @@ import fs from "fs"; import path from "path"; import { DATA_DIR, SERVER_ROOT } from "./paths.js"; -const dbPath = process.env.SCRUNNER_DB_PATH ?? path.join(DATA_DIR, "scrunner.db"); +const defaultNew = path.join(DATA_DIR, "jerapah-flow.db"); +const defaultLegacy = path.join(DATA_DIR, "scrunner.db"); +const dbPath = + process.env.JERAPAH_FLOW_DB_PATH ?? + process.env.SCRUNNER_DB_PATH ?? + (fs.existsSync(defaultLegacy) && !fs.existsSync(defaultNew) ? defaultLegacy : defaultNew); fs.mkdirSync(path.dirname(dbPath), { recursive: true }); /** @type {import("knex").Knex.Config} */ diff --git a/packages/server/logger.js b/packages/server/logger.js index 0c19b24..2cf80b6 100644 --- a/packages/server/logger.js +++ b/packages/server/logger.js @@ -105,7 +105,7 @@ fs.mkdirSync(LOGS_DIR, { recursive: true }); const rollingFile = pino.transport({ target: "pino-roll", options: { - file: path.join(LOGS_DIR, "scrunner.log"), + file: path.join(LOGS_DIR, "jerapah-flow.log"), size: "10m", mkdir: true, limit: { count: 5 }, @@ -113,7 +113,7 @@ const rollingFile = pino.transport({ }); export const log = pino( - { level: process.env.SCRUNNER_LOG_LEVEL ?? "debug" }, + { level: process.env.JERAPAH_FLOW_LOG_LEVEL ?? process.env.SCRUNNER_LOG_LEVEL ?? "debug" }, pino.multistream([ { level: "debug", stream: redactStream(process.stdout) }, { level: "debug", stream: redactStream(rollingFile) }, diff --git a/packages/server/migrations/20260816220000_variables.js b/packages/server/migrations/20260816220000_variables.js new file mode 100644 index 0000000..4228db4 --- /dev/null +++ b/packages/server/migrations/20260816220000_variables.js @@ -0,0 +1,26 @@ +/** + * @param {import("knex").Knex} knex + */ +export async function up(knex) { + await knex.schema.createTable("variables", (t) => { + t.text("id").primary(); + t.text("owner").notNullable(); + t.text("name").notNullable(); + t.text("type").notNullable(); + t.text("value").notNullable(); + t.text("created_at").notNullable(); + t.text("updated_at").notNullable(); + t.unique(["owner", "name"]); + }); + + await knex.schema.raw( + "CREATE INDEX variables_owner_name_idx ON variables (owner, name)", + ); +} + +/** + * @param {import("knex").Knex} knex + */ +export async function down(knex) { + await knex.schema.dropTableIfExists("variables"); +} diff --git a/packages/server/package.json b/packages/server/package.json index e5bc384..c3c1eb5 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -1,5 +1,5 @@ { - "name": "@scrunner/server", + "name": "@jerapah-flow/server", "version": "1.0.0", "private": true, "type": "module", diff --git a/packages/server/registry.js b/packages/server/registry.js index d25bc4c..f4d56ff 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, @@ -23,6 +30,7 @@ import { sendHttpPageOrJson, sendSuccessPage, } from "./http-trigger-auth.js"; +import { resolveConfigRefs } from "./config-refs.js"; /** * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry @@ -130,9 +138,12 @@ export function createRegistry(server) { const onDisk = fsStore.listOwnerYamlFiles(owner); for (const file of onDisk) { - if (!workflowFiles.includes(file)) { - log.warn(`Workflow file not in registers.yaml: ${owner}/${file}`); + if (workflowFiles.includes(file)) continue; + if (file.startsWith("dev-")) { + workflowFiles.push(file); + continue; } + log.warn(`Workflow file not in registers.yaml: ${owner}/${file}`); } for (const file of workflowFiles) { @@ -312,7 +323,7 @@ export function createRegistry(server) { pruneTask.destroy(); pruneTask = null; } - const days = Number(process.env.SCRUNNER_RETENTION_DAYS ?? 30); + const days = Number(process.env.JERAPAH_FLOW_RETENTION_DAYS ?? process.env.SCRUNNER_RETENTION_DAYS ?? 30); pruneTask = cron.schedule( "0 0 * * *", async () => { @@ -327,24 +338,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 +356,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 +369,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 +389,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 +405,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 +421,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 +514,10 @@ export function createRegistry(server) { } break; } + + for (const parsed of wave) completed.add(parsed.index); } + return last; } @@ -483,13 +556,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, @@ -502,39 +576,57 @@ export function createRegistry(server) { depth, ) { const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; - const config = parsed.config; + const unresolvedConfig = parsed.config; + const incomingContext = normalizeContext(ctx.context); const step = await store.startStep({ runId, index, script, - config, + config: unresolvedConfig, }); const stepLog = runLog.child({ stepId: step.id, script }); try { + const config = await resolveConfigRefs(unresolvedConfig, { + owner, + workflowKey: key, + context: incomingContext, + }); + const stepCtx = { + data: ctx.data, + context: incomingContext, + config, + }; 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 +636,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 +672,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 +701,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/runner.js b/packages/server/runner.js index 452e380..c0a610b 100644 --- a/packages/server/runner.js +++ b/packages/server/runner.js @@ -17,6 +17,7 @@ import runsPlugin from "./src/api/runs.js"; import dashboardPluginFactory from "./src/api/dashboard.js"; import secretsPlugin from "./src/api/secrets.js"; import kvPlugin from "./src/api/kv.js"; +import variablesPlugin from "./src/api/variables.js"; import httpPagesPlugin from "./src/api/http-pages.js"; import httpAuthsPlugin from "./src/api/http-auths.js"; import { WEB_DIST } from "./paths.js"; @@ -26,11 +27,12 @@ await migrate(); enableLogPersistence(); const jwtSecret = + process.env.JERAPAH_FLOW_JWT_SECRET ?? process.env.SCRUNNER_JWT_SECRET ?? (process.env.NODE_ENV === "production" ? "" : "scrunner-dev-secret"); if (!jwtSecret) { - log.error("SCRUNNER_JWT_SECRET is required in production"); + log.error("JERAPAH_FLOW_JWT_SECRET is required in production"); process.exit(1); } @@ -52,7 +54,7 @@ await server.register(jwt, { }, }); await server.register(cors, { - origin: process.env.SCRUNNER_CORS_ORIGIN ?? "http://localhost:5173", + origin: process.env.JERAPAH_FLOW_CORS_ORIGIN ?? process.env.SCRUNNER_CORS_ORIGIN ?? "http://localhost:5173", credentials: true, }); @@ -91,6 +93,7 @@ await server.register( await api.register(authPlugin); await api.register(usersPlugin); await api.register(secretsPlugin); + await api.register(variablesPlugin); await api.register(kvPlugin); await api.register(httpPagesPlugin); await api.register(httpAuthsPlugin); diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index aa22fc6..4b7c7c7 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -10,6 +10,7 @@ import { SCRIPTS_DIR } from "./paths.js"; import { isSecret, Secret, unwrapSecretsDeep } from "./secret-value.js"; import { getHttpPageByName, getHttpTemplateByName } from "./http-pages-store.js"; import { getSecretPlaintext } from "./secrets-store.js"; +import { getVariablePlain } from "./variables-store.js"; const hostRequire = createRequire(import.meta.url); @@ -301,6 +302,24 @@ function createRestrictedRequire(screenedAxios) { }; } +/** + * @param {string} owner + */ +function createVarsApi(owner) { + return { + /** + * @param {string} name + */ + async get(name) { + const value = await getVariablePlain(owner, name); + if (value == null) { + throw new Error(`variable "${name}" not found`); + } + return value; + }, + }; +} + /** * @param {string} owner */ @@ -380,6 +399,7 @@ function createScriptSandbox({ const $kv = createKvApi(workflowName); const $fingerprint = createFingerprintApi($kv); const $secrets = createSecretsApi(owner); + const $vars = createVarsApi(owner); const $responses = createResponsesApi(); const sandbox = { ...pickBuiltins(), @@ -389,13 +409,14 @@ function createScriptSandbox({ $kv, $fingerprint, $secrets, + $vars, $responses, $workflows, require: createRestrictedRequire($axios), }; vm.createContext(sandbox, { - name: `scrunner:${workflowName}:${script}`, + name: `jerapah-flow:${workflowName}:${script}`, codeGeneration: { strings: false, wasm: false }, }); diff --git a/packages/server/scripts/fetch-binary.js b/packages/server/scripts/fetch-binary.js index 5f6209a..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,26 +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: { @@ -87,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 f86a0ed..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,20 +72,21 @@ 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 = { description: "Fetch HTML and optionally select elements or transform with JSONata", + previewConfigKey: "url", + tags: ["HTTP"], config: { url: { type: "string", required: true, description: "Page URL" }, selector: { type: "string", required: false, description: "CSS selector" }, @@ -98,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 new file mode 100644 index 0000000..e276dd3 --- /dev/null +++ b/packages/server/scripts/fetch-http.js @@ -0,0 +1,172 @@ +const HTTP_METHODS = ["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"]; + +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} + +function isPlainObject(value) { + return value != null && typeof value === "object" && !Array.isArray(value); +} + +function resolveMethod(ctx) { + const raw = ctx.config?.method ?? ctx.data?.method ?? "GET"; + if (typeof raw !== "string" || raw.length === 0) { + throw new Error("fetch-http: method must be a non-empty string"); + } + const method = raw.toUpperCase(); + if (!HTTP_METHODS.includes(method)) { + throw new Error(`fetch-http: unsupported method "${raw}"`); + } + return method; +} + +function normalizeHeaders(value, label) { + if (value == null) return {}; + if (!isPlainObject(value)) { + throw new Error(`fetch-http: ${label} must be an object`); + } + + /** @type {Record} */ + const headers = {}; + for (const [key, val] of Object.entries(value)) { + if (val == null) continue; + if (typeof val === "string") { + headers[key] = val; + continue; + } + if (typeof val === "number" || typeof val === "boolean") { + headers[key] = String(val); + continue; + } + if (typeof val === "object") { + headers[key] = val; + continue; + } + throw new Error(`fetch-http: header "${key}" must be a string`); + } + return headers; +} + +function resolveHeaders(ctx) { + return { + ...normalizeHeaders(ctx.config?.headers, "config.headers"), + ...normalizeHeaders(ctx.data?.headers, "data.headers"), + }; +} + +function resolveBody(ctx) { + if (ctx.config != null && typeof ctx.config === "object" && "body" in ctx.config) { + return { present: true, value: ctx.config.body }; + } + if (ctx.data != null && typeof ctx.data === "object" && "body" in ctx.data) { + return { present: true, value: ctx.data.body }; + } + return { present: false, value: undefined }; +} + +function responseSize(data) { + if (data == null) return 0; + if (typeof data === "string") return data.length; + if (typeof Buffer !== "undefined" && Buffer.isBuffer(data)) return data.length; + try { + return JSON.stringify(data).length; + } catch { + return null; + } +} + +async function fetchHttp(ctx) { + const url = ctx.config?.url; + if (typeof url !== "string" || url.length === 0) { + throw new Error("fetch-http: ctx.config.url is required"); + } + + const method = resolveMethod(ctx); + const headers = resolveHeaders(ctx); + const body = resolveBody(ctx); + + /** @type {Record} */ + const request = { url, method, headers }; + if (body.present && method !== "HEAD") { + request.data = body.value; + } + + log.info({ url, method }, "fetch-http: fetching"); + const response = await $axios.request(request); + log.info( + { status: response.status, length: responseSize(response.data) }, + "fetch-http: fetch complete", + ); + + const output = { + httpResponse: response.data, + httpStatus: response.status, + }; + log.info( + { status: output.httpStatus, length: responseSize(output.httpResponse) }, + "fetch-http: saved httpResponse", + ); + + return { output, context: { ...passContext(ctx), ...output } }; +} + +fetchHttp.meta = { + description: "Fetch a URL and store the response body.", + previewConfigKey: "url", + tags: ["HTTP"], + config: { + url: { type: "string", required: true, description: "Request URL" }, + method: { + type: "string", + default: "GET", + enum: ["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], + description: "HTTP method", + }, + headers: { + type: "object", + required: false, + description: "Request headers (string values; Secret values are unwrapped)", + }, + body: { + type: "any", + required: false, + description: "Request body for POST/PUT/PATCH (object, string, or buffer)", + }, + }, + input: { + method: { + type: "string", + required: false, + enum: ["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], + description: "Fallback when config.method is omitted", + }, + headers: { + type: "object", + required: false, + description: "Merged over config.headers (e.g. Authorization from a secret)", + }, + body: { type: "any", required: false, description: "Fallback when config.body is omitted" }, + }, + output: { + 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: { + url: "https://httpbin.org/post", + method: "POST", + headers: { "Content-Type": "application/json", Accept: "application/json" }, + body: { hello: "world" }, + }, + }, +}; + +export default fetchHttp; diff --git a/packages/server/scripts/fetch-rss-feed.js b/packages/server/scripts/fetch-rss-feed.js index 64af8ef..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,24 +56,26 @@ 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 = { description: "Fetch an RSS/Atom feed and optionally transform it with JSONata", + previewConfigKey: "url", + tags: ["RSS"], config: { url: { type: "string", required: false, description: "Feed URL (or pass data.url)" }, 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", @@ -87,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 b02424c..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,16 +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", @@ -87,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 6091c2f..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,17 @@ 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: { type: "string", @@ -31,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 9f51b00..b05cf24 100644 --- a/packages/server/scripts/jsonata.js +++ b/packages/server/scripts/jsonata.js @@ -1,25 +1,34 @@ import jsonata from "jsonata"; -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; +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); + 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 7555ce8..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,18 +79,21 @@ 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; } const headers = ntfyHeaders(ctx); - const ntfyUrl = ctx.config?.url || "https://ntfy.sh/scrunner"; + const ntfyUrl = ctx.config?.url || "https://ntfy.sh/jerapah-flow"; if (hasFile) { const filename = ctx.data?.filename || "attachment"; @@ -114,7 +124,7 @@ async function ntfy(ctx) { log.info("ntfy sending message to %s", ntfyUrl); log.info("ntfy messsage: %s", truncatedMessage); - await $axios.post(ntfyUrl, ctx.data?.message || "Hello from scrunner", { + await $axios.post(ntfyUrl, ctx.data?.message || "Hello from jerapah-flow", { headers, }); } @@ -126,15 +136,17 @@ async function ntfy(ctx) { sent.fingerprint = stored.hash; sent.fingerprintAt = stored.at; } - return sent; + return { output: sent, context: passContext(ctx) }; } ntfy.meta = { description: "Send a message or file to an ntfy topic", + previewConfigKey: "url", + tags: ["channel"], config: { url: { type: "string", - default: "https://ntfy.sh/scrunner", + default: "https://ntfy.sh/jerapah-flow", description: "ntfy topic URL", }, fingerprint: { @@ -145,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", @@ -161,12 +173,13 @@ 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 scrunner" }, - config: { url: "https://ntfy.sh/scrunner", fingerprint: true }, + data: { title: "Hello", message: "Hello from jerapah-flow" }, + config: { url: "https://ntfy.sh/jerapah-flow", fingerprint: true }, }, }; diff --git a/packages/server/scripts/render-template.js b/packages/server/scripts/render-template.js index 7870055..23b2e2d 100644 --- a/packages/server/scripts/render-template.js +++ b/packages/server/scripts/render-template.js @@ -70,15 +70,22 @@ 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 + : {}, }; } renderTemplate.meta = { description: "Render an HTML template from Responses (kind: template) with Mustache and return html + plain-text fallback", + previewConfigKey: "template", config: { template: { type: "string", diff --git a/packages/server/scripts/send-email.js b/packages/server/scripts/send-email.js index cf88b25..8a9c612 100644 --- a/packages/server/scripts/send-email.js +++ b/packages/server/scripts/send-email.js @@ -218,24 +218,32 @@ 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 + : {}, }; } sendEmail.meta = { description: "Send an email via SMTP (nodemailer); plain text, HTML, or both", + previewConfigKey: "from", + tags: ["channel"], config: { service: { type: "string", required: false, description: - "Well-known provider shortcut (e.g. gmail, outlook365, sendgrid). Alternative to host.", + "Nodemailer well-known service ID (e.g. Gmail, Outlook365, SendGrid). Sets host, port, and TLS. See https://nodemailer.com/smtp/well-known-services", }, host: { type: "string", @@ -426,13 +434,11 @@ sendEmail.meta = { to: ["recipient@example.com", "other@example.com"], cc: "manager@example.com", bcc: "audit@example.com", - subject: "Hello from scrunner", + subject: "Hello from JerapahFlow", text: "This is a plain-text test message.", }, config: { - host: "smtp.gmail.com", - port: 587, - secure: false, + service: "Gmail", user: "smtp-login@gmail.com", from: "notifications@example.com", passwordSecret: "gmail_app_password", diff --git a/packages/server/scripts/trigger-workflow.js b/packages/server/scripts/trigger-workflow.js index 5e78d6b..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,25 +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", @@ -38,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: { @@ -56,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/secrets.js b/packages/server/secrets.js index 86f0865..45da609 100644 --- a/packages/server/secrets.js +++ b/packages/server/secrets.js @@ -11,10 +11,11 @@ const AUTH_TAG_LEN = 16; */ export function resolveSecretsKeyMaterial() { const raw = + process.env.JERAPAH_FLOW_SECRETS_KEY ?? process.env.SCRUNNER_SECRETS_KEY ?? (process.env.NODE_ENV === "production" ? "" : DEV_DEFAULT); if (!raw) { - throw new Error("SCRUNNER_SECRETS_KEY is required in production"); + throw new Error("JERAPAH_FLOW_SECRETS_KEY is required in production"); } return raw; } diff --git a/packages/server/src/api/auth.js b/packages/server/src/api/auth.js index d877c3d..8683230 100644 --- a/packages/server/src/api/auth.js +++ b/packages/server/src/api/auth.js @@ -1,7 +1,7 @@ import bcrypt from "bcryptjs"; import * as store from "../../store.js"; -export const COOKIE = "scrunner_token"; +export const COOKIE = "jerapah_flow_token"; export const OPEN_API_ROUTES = new Set([ "GET /auth/bootstrap", "POST /auth/register", diff --git a/packages/server/src/api/scripts.js b/packages/server/src/api/scripts.js index 939ade7..6d175d0 100644 --- a/packages/server/src/api/scripts.js +++ b/packages/server/src/api/scripts.js @@ -1,3 +1,4 @@ +import fs from "fs"; import { clearScriptCache, inspectScriptSource, @@ -5,6 +6,8 @@ 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"; +import { resolveConfigRefs } from "../../config-refs.js"; /** * @param {{ referencedScripts: () => Set }} registry @@ -21,11 +24,27 @@ export default function scriptsPluginFactory(registry) { content == null ? { meta: null, metaError: "script not found" } : inspectScriptSource(name, content); - return { name, ...inspected }; + return { name, hasIcon: fsStore.scriptHasIcon(name), ...inspected }; }); return { scripts }; }); + fastify.get("/scripts/:name/icon", async (req, reply) => { + const { name } = /** @type {{ name: string }} */ (req.params); + try { + fsStore.assertScriptName(name); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + const icon = fsStore.resolveScriptIcon(name); + if (icon == null) return reply.code(404).send({ error: "icon not found" }); + const body = fs.readFileSync(icon.filePath); + return reply + .type(icon.contentType) + .header("Cache-Control", "private, max-age=60") + .send(body); + }); + fastify.get("/scripts/:name", async (req, reply) => { const { name } = /** @type {{ name: string }} */ (req.params); try { @@ -35,7 +54,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, ...inspectScriptSource(name, content) }; + return { name, content, hasIcon: fsStore.scriptHasIcon(name), ...inspectScriptSource(name, content) }; }); fastify.put("/scripts/:name", async (req, reply) => { @@ -85,7 +104,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") { @@ -101,24 +120,37 @@ export default function scriptsPluginFactory(registry) { } } - const ctx = { - data: body.data ?? null, - config: body.config ?? null, - }; + const incomingContext = + body.context != null && typeof body.context === "object" && !Array.isArray(body.context) + ? body.context + : {}; const { log, logs } = createDryRunLogger(); const started = Date.now(); try { + const config = await resolveConfigRefs(body.config ?? null, { + owner, + workflowKey: "dry-run", + context: incomingContext, + }); + const ctx = { + data: body.data ?? null, + context: incomingContext, + config, + }; const { fn, meta, metaError } = instantiateScriptSource(name, body.content, { log, 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, @@ -130,6 +162,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/variables.js b/packages/server/src/api/variables.js new file mode 100644 index 0000000..27f359a --- /dev/null +++ b/packages/server/src/api/variables.js @@ -0,0 +1,55 @@ +import * as fsStore from "../../fs-store.js"; +import { + assertVariableName, + assertVariableType, + deleteVariable, + getVariableById, + listVariables, + upsertVariable, +} from "../../variables-store.js"; + +/** + * @param {import("fastify").FastifyInstance} fastify + */ +export default async function variablesPlugin(fastify) { + fastify.get("/variables", async (req, reply) => { + const q = /** @type {{ owner?: string }} */ (req.query ?? {}); + try { + const owner = q.owner ? fsStore.assertOwner(q.owner) : undefined; + const variables = await listVariables({ owner }); + return { variables }; + } catch (err) { + return reply.code(err.statusCode ?? 500).send({ error: err.message }); + } + }); + + fastify.put("/variables", async (req, reply) => { + const body = /** @type {{ owner?: string, name?: string, type?: unknown, value?: unknown }} */ ( + req.body ?? {} + ); + try { + fsStore.assertOwner(String(body.owner ?? "")); + assertVariableName(String(body.name ?? "")); + assertVariableType(body.type); + const variable = await upsertVariable({ + owner: String(body.owner), + name: String(body.name), + type: body.type, + value: body.value, + }); + return reply.send({ variable }); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + }); + + fastify.delete("/variables/:id", async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); + const existing = await getVariableById(id); + if (!existing) { + return reply.code(404).send({ error: "variable not found" }); + } + await deleteVariable(id); + return { ok: true }; + }); +} diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index b4b8da4..abc4bea 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -10,6 +10,11 @@ import { authLabel, validateWorkflowHttpTriggers, } from "../../workflow-http-validate.js"; +import { + duplicateWorkflowYaml, + ensureWorkflowFilename, + suggestCopyFilename, +} from "../../workflow-duplicate.js"; function triggerSummary(owner, workflow) { if (!workflow || typeof workflow !== "object") return []; @@ -32,7 +37,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); } @@ -259,6 +264,101 @@ export default function workflowsPluginFactory(registry) { return { ok: true }; }); + fastify.post("/workflows/:owner/:file/duplicate", async (req, reply) => { + const { owner, file } = /** @type {{ owner: string, file: string }} */ ( + req.params + ); + try { + fsStore.assertOwner(owner); + fsStore.assertWorkflowFile(file); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + const source = fsStore.readWorkflowYaml(owner, file); + if (source == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + + const body = /** @type {{ file?: unknown, owner?: unknown }} */ (req.body ?? {}); + let destOwner = owner; + if (body.owner != null && body.owner !== "") { + if (typeof body.owner !== "string") { + return reply.code(400).send({ error: "owner must be a string" }); + } + try { + destOwner = fsStore.assertOwner(body.owner); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + } + + let destFile; + try { + if (body.file == null || body.file === "") { + destFile = suggestCopyFilename(file, fsStore.listOwnerYamlFiles(destOwner)); + } else if (typeof body.file !== "string") { + return reply.code(400).send({ error: "file must be a string" }); + } else { + destFile = ensureWorkflowFilename(body.file); + } + fsStore.assertWorkflowFile(destFile); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + + if (destOwner === owner && destFile === file) { + return reply.code(400).send({ error: "cannot duplicate onto itself" }); + } + if (fsStore.readWorkflowYaml(destOwner, destFile) != null) { + return reply.code(409).send({ error: "workflow already exists" }); + } + + let content; + try { + content = duplicateWorkflowYaml(source, { + sourceFile: file, + destFile, + rewriteHttpPaths: destOwner === owner, + }); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + + let parsed; + try { + parsed = yaml.parse(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), + }); + } + try { + await validateWorkflowHttpTriggers(parsed); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + + fsStore.writeWorkflowYaml(destOwner, destFile, content); + const registered = fsStore.readRegisters(destOwner); + if (!registered.includes(destFile)) { + registered.push(destFile); + fsStore.writeRegisters(destOwner, registered); + } + registry.reregister(); + return reply.code(201).send({ owner: destOwner, file: destFile }); + }); + fastify.post("/workflows/:owner/:file/run", async (req, reply) => { const { owner, file } = /** @type {{ owner: string, file: string }} */ ( req.params 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/config-refs-smoke.js b/packages/server/test/config-refs-smoke.js new file mode 100644 index 0000000..c1c0eaf --- /dev/null +++ b/packages/server/test/config-refs-smoke.js @@ -0,0 +1,213 @@ +import { migrate, db } from "../db.js"; +import { upsertSecret, deleteSecret } from "../secrets-store.js"; +import { deleteVariable, upsertVariable } from "../variables-store.js"; +import { Secret } from "../secret-value.js"; +import { parseConfigRef, resolveConfigRefs } from "../config-refs.js"; + +await migrate(); + +function assert(cond, msg) { + if (!cond) throw new Error(msg); +} + +async function assertRejects(fn, match) { + try { + await fn(); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + if (match && !message.includes(match)) { + throw new Error(`rejected with "${message}", expected to include "${match}"`); + } + return; + } + throw new Error(`expected to reject (${match ?? "any error"})`); +} + +const owner = "default"; +const workflowKey = "default/config-refs-smoke.yaml"; +const ctx = { owner, workflowKey, context: {} }; + +function assertParse(value, expected) { + const got = parseConfigRef(value); + if (expected == null) { + assert(got == null, `expected null parse for ${JSON.stringify(value)}, got ${JSON.stringify(got)}`); + return; + } + assert(got != null, `expected parse for ${JSON.stringify(value)}`); + assert(got.kind === expected.kind, `kind ${got.kind} !== ${expected.kind}`); + assert(got.name === expected.name, `name ${JSON.stringify(got.name)} !== ${JSON.stringify(expected.name)}`); +} + +assertParse("password123", null); +assertParse("$FOO_bar", null); +assertParse("$SECRET", null); +assertParse(" password123 ", null); +assertParse("$SECRET_zte_modem_password", { kind: "secret", name: "zte_modem_password" }); +assertParse(" $SECRET_zte_modem_password ", { kind: "secret", name: "zte_modem_password" }); +assertParse("Bearer $SECRET_x", null); +assertParse("$KV_modem password", null); +assertParse("$VAR_ntfy_url", { kind: "var", name: "ntfy_url" }); +assertParse("$CONTEXT_token", { kind: "context", name: "token" }); +assertParse("$SECRET_", { kind: "secret", name: "" }); +assertParse("$CONTEXT_SECRET_foo", { kind: "context", name: "SECRET_foo" }); + +{ + const literal = await resolveConfigRefs("password123", ctx); + assert(literal === "password123", "literal passthrough"); + const unknown = await resolveConfigRefs("$FOO_bar", ctx); + assert(unknown === "$FOO_bar", "$FOO_bar stays literal"); + const kvLiteral = await resolveConfigRefs("$KV_modem_password", ctx); + assert(kvLiteral === "$KV_modem_password", "$KV_ stays literal"); + const embedded = await resolveConfigRefs("Bearer $SECRET_x", ctx); + assert(embedded === "Bearer $SECRET_x", "mid-string stays literal"); + const number = await resolveConfigRefs(42, ctx); + assert(number === 42, "number passthrough"); +} + +const secret = await upsertSecret({ + owner, + name: "config_refs_smoke_token", + value: "s3cret-ok", +}); +const varUrl = await upsertVariable({ + owner, + name: "config_refs_smoke_url", + type: "string", + value: "https://example.test", +}); +const varRetry = await upsertVariable({ + owner, + name: "config_refs_smoke_retry", + type: "number", + value: 3, +}); +const varDebug = await upsertVariable({ + owner, + name: "config_refs_smoke_debug", + type: "boolean", + value: false, +}); + +try { + { + const resolved = await resolveConfigRefs("$SECRET_config_refs_smoke_token", ctx); + assert(resolved === "s3cret-ok", "secret resolve"); + } + { + const resolved = await resolveConfigRefs(" $SECRET_config_refs_smoke_token ", ctx); + assert(resolved === "s3cret-ok", "secret resolve trimmed"); + } + { + const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_url", ctx); + assert(resolved === "https://example.test", "var string"); + } + { + const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_retry", ctx); + assert(resolved === 3, "var number stays number"); + assert(typeof resolved === "number", "var number type"); + } + { + const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_debug", ctx); + assert(resolved === false, "var boolean stays false"); + assert(typeof resolved === "boolean", "var boolean type"); + } + { + const resolved = await resolveConfigRefs("$CONTEXT_token", { + ...ctx, + context: { token: "ctx-token-ok" }, + }); + assert(resolved === "ctx-token-ok", "context string"); + } + { + const resolved = await resolveConfigRefs("$CONTEXT_n", { + ...ctx, + context: { n: 7 }, + }); + assert(resolved === "7", "context number stringify"); + } + { + const wrapped = new Secret("wrapped-secret-ok"); + const resolved = await resolveConfigRefs("$CONTEXT_tok", { + ...ctx, + context: { tok: wrapped }, + }); + assert(resolved === "wrapped-secret-ok", "context Secret unwrap"); + } + + { + const nested = await resolveConfigRefs( + { + url: "$VAR_config_refs_smoke_url", + retry: "$VAR_config_refs_smoke_retry", + debug: "$VAR_config_refs_smoke_debug", + password: "$SECRET_config_refs_smoke_token", + headers: { Authorization: "$KV_modem_password" }, + extra: ["$CONTEXT_token", "plain"], + }, + { ...ctx, context: { token: "ctx-token-ok" } }, + ); + assert(nested.url === "https://example.test", "nested var string"); + assert(nested.retry === 3, "nested var number"); + assert(nested.debug === false, "nested var boolean"); + assert(nested.password === "s3cret-ok", "nested secret"); + assert(nested.headers.Authorization === "$KV_modem_password", "nested $KV_ stays literal"); + assert(nested.extra[0] === "ctx-token-ok", "nested array context"); + assert(nested.extra[1] === "plain", "nested array literal"); + } + + const data = { password: "$SECRET_config_refs_smoke_token" }; + const config = { password: "$SECRET_config_refs_smoke_token" }; + const resolvedConfig = await resolveConfigRefs(config, ctx); + assert(resolvedConfig.password === "s3cret-ok", "config resolved"); + assert(data.password === "$SECRET_config_refs_smoke_token", "data not walked"); + assert(config.password === "$SECRET_config_refs_smoke_token", "input config not mutated"); + + await assertRejects( + () => resolveConfigRefs("$SECRET_does_not_exist_xyz", ctx), + 'secret "does_not_exist_xyz" not found', + ); + await assertRejects( + () => resolveConfigRefs("$SECRET_not valid", ctx), + "invalid secret name", + ); + await assertRejects( + () => resolveConfigRefs("$SECRET_", ctx), + "invalid secret name", + ); + await assertRejects( + () => resolveConfigRefs("$CONTEXT_missing", ctx), + 'context "missing" not found', + ); + await assertRejects( + () => + resolveConfigRefs("$CONTEXT_obj", { + ...ctx, + context: { obj: { a: 1 } }, + }), + 'context "obj" is not a scalar', + ); + await assertRejects( + () => resolveConfigRefs("$CONTEXT_", ctx), + "empty context key", + ); + await assertRejects( + () => resolveConfigRefs("$VAR_does_not_exist_xyz", ctx), + 'variable "does_not_exist_xyz" not found', + ); + await assertRejects( + () => resolveConfigRefs("$VAR_not valid", ctx), + "invalid variable name", + ); + await assertRejects( + () => resolveConfigRefs("$VAR_", ctx), + "empty variable name", + ); +} finally { + await deleteSecret(secret.id); + await deleteVariable(varUrl.id); + await deleteVariable(varRetry.id); + await deleteVariable(varDebug.id); +} + +console.log("config-refs smoke test passed"); +await db.destroy(); 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/test/variables-smoke.js b/packages/server/test/variables-smoke.js new file mode 100644 index 0000000..18e5908 --- /dev/null +++ b/packages/server/test/variables-smoke.js @@ -0,0 +1,136 @@ +import { migrate, db } from "../db.js"; +import { runScriptSource } from "../script-sandbox.js"; +import { log } from "../logger.js"; +import { + deleteVariable, + encodeVariableValue, + getVariablePlain, + upsertVariable, +} from "../variables-store.js"; + +await migrate(); + +function assert(cond, msg) { + if (!cond) throw new Error(msg); +} + +async function assertThrows(fn, match) { + try { + await fn(); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + if (match && !message.includes(match)) { + throw new Error(`threw "${message}", expected to include "${match}"`); + } + return; + } + throw new Error(`expected to throw (${match ?? "any error"})`); +} + +const owner = "default"; + +await assertThrows( + () => encodeVariableValue("number", "3"), + "finite number", +); +await assertThrows( + () => encodeVariableValue("number", NaN), + "finite number", +); +await assertThrows( + () => encodeVariableValue("boolean", "true"), + "boolean", +); +await assertThrows( + () => encodeVariableValue("string", 3), + "string", +); + +const created = []; +try { + const str = await upsertVariable({ + owner, + name: "variables_smoke_url", + type: "string", + value: "https://n.0dev.web.id/system", + }); + created.push(str.id); + assert(str.value === "https://n.0dev.web.id/system", "string stored"); + assert( + (await getVariablePlain(owner, "variables_smoke_url")) === str.value, + "string plain", + ); + + const num = await upsertVariable({ + owner, + name: "variables_smoke_retry", + type: "number", + value: 3, + }); + created.push(num.id); + assert(num.value === 3 && typeof num.value === "number", "number stored"); + + const flag = await upsertVariable({ + owner, + name: "variables_smoke_debug", + type: "boolean", + value: false, + }); + created.push(flag.id); + assert(flag.value === false, "boolean false stored"); + + const updated = await upsertVariable({ + owner, + name: "variables_smoke_retry", + type: "number", + value: 9, + }); + assert(updated.id === num.id, "upsert same id"); + assert(updated.value === 9, "upsert number"); + + await assertThrows( + () => + upsertVariable({ + owner, + name: "variables_smoke_bad", + type: "number", + value: "abc", + }), + "finite number", + ); + + const scriptOut = await runScriptSource( + "variables-smoke.js", + `export default async function () { + const url = await $vars.get("variables_smoke_url"); + const retry = await $vars.get("variables_smoke_retry"); + const debug = await $vars.get("variables_smoke_debug"); + return { url, retry, debug, retryType: typeof retry, debugType: typeof debug }; + }`, + { data: {} }, + { log, workflowName: "default/variables-smoke.yaml", owner }, + ); + assert(scriptOut.url === "https://n.0dev.web.id/system", "script $vars string"); + assert(scriptOut.retry === 9 && scriptOut.retryType === "number", "script $vars number"); + assert(scriptOut.debug === false && scriptOut.debugType === "boolean", "script $vars boolean"); + + await assertThrows( + () => + runScriptSource( + "variables-smoke.js", + `export default async function () { + return $vars.get("does_not_exist_xyz"); + }`, + { data: {} }, + { log, workflowName: "default/variables-smoke.yaml", owner }, + ), + 'variable "does_not_exist_xyz" not found', + ); +} finally { + for (const id of created) { + await deleteVariable(id); + } +} + +console.log("variables smoke test passed"); +await db.destroy(); diff --git a/packages/server/variables-store.js b/packages/server/variables-store.js new file mode 100644 index 0000000..2ed9f64 --- /dev/null +++ b/packages/server/variables-store.js @@ -0,0 +1,190 @@ +import { randomUUID } from "node:crypto"; +import { db } from "./db.js"; +import { assertOwner } from "./fs-store.js"; + +const MAX_NAME_LENGTH = 128; +const MAX_STRING_BYTES = 64 * 1024; +const VARIABLE_NAME_RE = /^[A-Za-z0-9._-]+$/; +export const VARIABLE_TYPES = /** @type {const} */ (["string", "number", "boolean"]); + +function nowIso() { + return new Date().toISOString(); +} + +function httpError(message, statusCode = 400) { + const err = new Error(message); + err.statusCode = statusCode; + return err; +} + +/** + * @param {unknown} name + * @returns {string} + */ +export function assertVariableName(name) { + if (typeof name !== "string" || !VARIABLE_NAME_RE.test(name)) { + throw httpError("invalid variable name"); + } + if (name.length > MAX_NAME_LENGTH) { + throw httpError(`variable name must be at most ${MAX_NAME_LENGTH} characters`); + } + return name; +} + +/** + * @param {unknown} type + * @returns {"string" | "number" | "boolean"} + */ +export function assertVariableType(type) { + if (type !== "string" && type !== "number" && type !== "boolean") { + throw httpError("type must be string, number, or boolean"); + } + return type; +} + +/** + * @param {"string" | "number" | "boolean"} type + * @param {unknown} value + * @returns {string} + */ +export function encodeVariableValue(type, value) { + if (type === "string") { + if (typeof value !== "string") { + throw httpError("value must be a string"); + } + if (Buffer.byteLength(value, "utf8") > MAX_STRING_BYTES) { + throw httpError(`value exceeds ${MAX_STRING_BYTES} byte limit`); + } + return value; + } + if (type === "number") { + if (typeof value !== "number" || !Number.isFinite(value)) { + throw httpError("value must be a finite number"); + } + return String(value); + } + if (typeof value !== "boolean") { + throw httpError("value must be a boolean"); + } + return value ? "true" : "false"; +} + +/** + * @param {"string" | "number" | "boolean"} type + * @param {string} stored + * @returns {string | number | boolean} + */ +export function decodeVariableValue(type, stored) { + if (type === "string") return stored; + if (type === "number") { + const n = Number(stored); + if (!Number.isFinite(n)) { + throw new Error(`corrupt number variable: ${JSON.stringify(stored)}`); + } + return n; + } + if (stored === "true") return true; + if (stored === "false") return false; + throw new Error(`corrupt boolean variable: ${JSON.stringify(stored)}`); +} + +/** + * @param {Record} row + */ +function publicVariable(row) { + const type = assertVariableType(row.type); + return { + id: row.id, + owner: row.owner, + name: row.name, + type, + value: decodeVariableValue(type, String(row.value ?? "")), + created_at: row.created_at, + updated_at: row.updated_at, + }; +} + +/** + * @param {{ owner?: string }} [filters] + */ +export async function listVariables(filters = {}) { + let q = db("variables") + .select("id", "owner", "name", "type", "value", "created_at", "updated_at") + .orderBy("owner", "asc") + .orderBy("name", "asc"); + if (filters.owner) { + q = q.where("owner", assertOwner(filters.owner)); + } + const rows = await q; + return rows.map((row) => publicVariable(row)); +} + +/** + * @param {string} id + */ +export async function getVariableById(id) { + const row = await db("variables").where({ id }).first(); + return row ? publicVariable(row) : null; +} + +/** + * @param {{ owner: string, name: string, type: unknown, value: unknown }} opts + */ +export async function upsertVariable({ owner, name, type, value }) { + const ownerName = assertOwner(owner); + const variableName = assertVariableName(name); + const variableType = assertVariableType(type); + const encoded = encodeVariableValue(variableType, value); + const now = nowIso(); + const existing = await db("variables") + .where({ owner: ownerName, name: variableName }) + .first(); + + if (existing) { + await db("variables") + .where({ id: existing.id }) + .update({ + type: variableType, + value: encoded, + updated_at: now, + }); + return getVariableById(existing.id); + } + + const id = randomUUID(); + await db("variables").insert({ + id, + owner: ownerName, + name: variableName, + type: variableType, + value: encoded, + created_at: now, + updated_at: now, + }); + return getVariableById(id); +} + +/** + * @param {string} id + * @returns {Promise} + */ +export async function deleteVariable(id) { + const n = await db("variables").where({ id }).del(); + return n > 0; +} + +/** + * Typed primitive for an owner/name. Returns null if missing. + * @param {string} owner + * @param {string} name + * @returns {Promise} + */ +export async function getVariablePlain(owner, name) { + const ownerName = assertOwner(owner); + const variableName = assertVariableName(name); + const row = await db("variables") + .where({ owner: ownerName, name: variableName }) + .first(); + if (!row) return null; + return decodeVariableValue(assertVariableType(row.type), String(row.value ?? "")); +} diff --git a/packages/server/workflow-duplicate.js b/packages/server/workflow-duplicate.js new file mode 100644 index 0000000..6e41506 --- /dev/null +++ b/packages/server/workflow-duplicate.js @@ -0,0 +1,111 @@ +import yaml from "yaml"; + +export function ensureWorkflowFilename(file) { + const trimmed = String(file ?? "").trim(); + if (!trimmed) return ""; + return /\.ya?ml$/i.test(trimmed) ? trimmed : `${trimmed}.yaml`; +} + +export function workflowFileStem(file) { + return String(file).replace(/\.ya?ml$/i, ""); +} + +/** + * Next unused copy filename: `track.yaml` → `track-copy.yaml`, + * `track-copy.yaml` → `track-copy-2.yaml`. + * @param {string} file + * @param {string[]} existingFiles + */ +export function suggestCopyFilename(file, existingFiles = []) { + const name = ensureWorkflowFilename(file) || "workflow.yaml"; + const match = name.match(/^(.*?)(\.ya?ml)$/i); + const base = match ? match[1] : name; + const ext = match ? match[2] : ".yaml"; + const existing = new Set(existingFiles); + + const copyMatch = base.match(/^(.*)-copy(?:-(\d+))?$/); + const root = copyMatch ? copyMatch[1] : base; + const candidate = (i) => + i <= 1 ? `${root}-copy${ext}` : `${root}-copy-${i}${ext}`; + + let n = copyMatch ? Number(copyMatch[2] || 1) + 1 : 1; + while (existing.has(candidate(n))) n += 1; + return candidate(n); +} + +export function nextCopyName(name) { + const trimmed = String(name ?? "").trim(); + if (!trimmed) return "copy"; + const match = trimmed.match(/^(.*) \(copy(?: (\d+))?\)$/); + if (!match) return `${trimmed} (copy)`; + const n = match[2] ? Number(match[2]) + 1 : 2; + return `${match[1]} (copy ${n})`; +} + +export function httpPathCopySuffix(sourceFile, destFile) { + const src = workflowFileStem(sourceFile); + const dest = workflowFileStem(destFile); + if (dest.startsWith(`${src}-`) && dest.length > src.length + 1) { + return dest.slice(src.length + 1); + } + if (dest === src) return "copy"; + return dest || "copy"; +} + +function suffixHttpPath(path, suffix) { + const trimmed = String(path).replace(/\/+$/, ""); + const withSlash = trimmed.startsWith("/") ? trimmed : `/${trimmed}`; + const safe = String(suffix).replace(/[^A-Za-z0-9._-]+/g, "-").replace(/^-+|-+$/g, "") || "copy"; + return `${withSlash}-${safe}`; +} + +function rewriteHttpTriggerPaths(doc, suffix) { + const triggers = doc.get("triggers"); + if (!yaml.isSeq(triggers)) return; + for (const item of triggers.items) { + if (!yaml.isMap(item)) continue; + if (String(item.get("type") ?? "").toLowerCase() !== "http") continue; + const path = item.get("path"); + if (typeof path !== "string" || !path.trim()) continue; + item.set("path", suffixHttpPath(path, suffix)); + } +} + +/** + * Copy workflow YAML: append " (copy)" to name, disable, optionally rewrite HTTP paths. + * Preserves comments via YAML CST. + * @param {string} content + * @param {{ sourceFile: string, destFile: string, rewriteHttpPaths?: boolean }} opts + */ +export function duplicateWorkflowYaml(content, opts) { + const sourceFile = opts?.sourceFile ?? ""; + const destFile = opts?.destFile ?? ""; + const rewriteHttpPaths = opts?.rewriteHttpPaths !== false; + + const doc = yaml.parseDocument(content); + if (doc.errors?.length) { + const err = new Error(doc.errors[0]?.message ?? "invalid yaml"); + err.statusCode = 400; + throw err; + } + const parsed = doc.toJSON(); + if (parsed == null || typeof parsed !== "object" || Array.isArray(parsed)) { + const err = new Error("workflow yaml must be an object"); + err.statusCode = 400; + throw err; + } + + const currentName = doc.get("name"); + if (typeof currentName === "string" && currentName.trim()) { + doc.set("name", nextCopyName(currentName)); + } else { + doc.set("name", workflowFileStem(destFile) || "copy"); + } + doc.set("enabled", false); + + if (rewriteHttpPaths) { + rewriteHttpTriggerPaths(doc, httpPathCopySuffix(sourceFile, destFile)); + } + + return String(doc); +} 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 345725f..8e8690b 100644 --- a/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml +++ b/packages/server/workflows/default/comic-monkeyuser-to-ntfy.yaml @@ -2,20 +2,23 @@ name: Comic - monkeyuser to ntfy scripts: - script: fetch-html.js config: - url: "https://www.monkeyuser.com/" - outputVar: "httpResponse" - selector: ".comic img" + url: https://www.monkeyuser.com/ + outputVar: comic + selector: .comic img jsonata: | {"url": "https://www.monkeyuser.com" & [attributes.src][0], "title": [attributes.title][0]} - script: jsonata.js config: expression: | - {"data": {"title": httpResponse.title, "message": httpResponse.title, "attach": httpResponse.url}} - - fetch-binary.js + { + "title": data.comic.title, + "message": data.comic.title, + "attach": data.comic.url + } + - script: fetch-binary.js - script: ntfy.js config: - url: https://ntfy.sh/scrunner - + url: https://ntfy.sh/jerapah-flow triggers: - type: HTTP method: POST diff --git a/packages/server/workflows/default/cron-example.yaml b/packages/server/workflows/default/cron-example.yaml index 5d2e484..b829f36 100644 --- a/packages/server/workflows/default/cron-example.yaml +++ b/packages/server/workflows/default/cron-example.yaml @@ -2,9 +2,13 @@ name: cron example description: | this workflow triggered by cron scripts: + - script: get-current-time.js + - set: + expression: '{"message": context.datetime}' - script: ntfy.js config: - url: https://ntfy.sh/scrunner + url: $VAR_ntfy_channel + title: $DATA triggers: - type: cron schedule: "*/20 * * * *" diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index d994fb3..7273056 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -6,4 +6,7 @@ scripts: - comic-monkeyuser-to-ntfy.yaml - time-and-comic-to-ntfy.yaml - rss-selfhst-to-ntfy.yaml - - detect-example-changes.yaml + - send-gmail.yaml + - test-send-gmail.yaml + - track.yaml + - rss-devto-to-ntfy.yaml diff --git a/packages/server/workflows/default/rss-devto-to-ntfy.yaml b/packages/server/workflows/default/rss-devto-to-ntfy.yaml new file mode 100644 index 0000000..16b7cb3 --- /dev/null +++ b/packages/server/workflows/default/rss-devto-to-ntfy.yaml @@ -0,0 +1,27 @@ +name: RSS - selfh.st first item to ntfy (copy) +scripts: + - script: fetch-rss-feed.js + config: + url: https://selfh.st/rss/ + outputVar: item + jsonata: items[0] + - script: fingerprint.js + config: + key: selfhst-latest + jsonata: "data.item.guid ? data.item.guid : data.item.link" + - script: jsonata.js + config: + expression: | + { + "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: + url: https://n.0dev.web.id/system +triggers: + - type: HTTP + method: POST + path: /selfhst-rss-rss-devto-to-ntfy +enabled: false diff --git a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml index dba12c2..900bfe2 100644 --- a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml +++ b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml @@ -2,7 +2,7 @@ name: RSS - selfh.st first item to ntfy scripts: - script: fetch-rss-feed.js config: - url: "https://selfh.st/rss/" + url: https://selfh.st/rss/ outputVar: item jsonata: items[0] - script: fingerprint.js @@ -13,16 +13,13 @@ 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: - url: https://ntfy.sh/scrunner - + url: https://n.0dev.web.id/system triggers: - type: HTTP method: POST diff --git a/packages/server/workflows/default/send-gmail.yaml b/packages/server/workflows/default/send-gmail.yaml new file mode 100644 index 0000000..d6b42cf --- /dev/null +++ b/packages/server/workflows/default/send-gmail.yaml @@ -0,0 +1,39 @@ +name: send-gmail +description: | + Send email via Gmail SMTP. Configure user/from and the named secret, + then call from other workflows with trigger-workflow.js (name: send-gmail). + + Input (data): + to optional recipient(s); defaults to your Gmail below + subject required + text plain body (aliases: body, message) + html optional HTML body + from, cc, bcc, replyTo, priority, headers optional +scripts: + - script: jsonata.js + config: + expression: | + { + "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: + service: Gmail + user: elevent16th@gmail.com + from: elevent16th@gmail.com + passwordSecret: gmail_app_password +triggers: + - type: workflow + - type: HTTP + method: POST + path: /send-gmail + auth: basic-auth diff --git a/packages/server/workflows/default/test-send-gmail.yaml b/packages/server/workflows/default/test-send-gmail.yaml new file mode 100644 index 0000000..0e67a39 --- /dev/null +++ b/packages/server/workflows/default/test-send-gmail.yaml @@ -0,0 +1,18 @@ +name: test-send-gmail +description: | + Kick send-gmail with sample data. Edit the payload below, then Run + (or POST /u/default/test-send-gmail). +scripts: + - script: trigger-workflow.js + config: + name: send-gmail + expression: | + { + "subject": "JerapahFlow test", + "text": "Hello from test-send-gmail", + "html": "

Hello from test-send-gmail

" + } +triggers: + - type: HTTP + method: POST + path: /test-send-gmail 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 52e1077..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,17 +19,17 @@ 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 needs: [compose] config: - url: https://ntfy.sh/scrunner + url: https://ntfy.sh/jerapah-flow triggers: - type: HTTP diff --git a/packages/server/workflows/default/time-to-ntfy-example.yaml b/packages/server/workflows/default/time-to-ntfy-example.yaml index 830c4e5..df6144c 100644 --- a/packages/server/workflows/default/time-to-ntfy-example.yaml +++ b/packages/server/workflows/default/time-to-ntfy-example.yaml @@ -2,10 +2,12 @@ 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 + config: + key: "" - 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 new file mode 100644 index 0000000..5593b67 --- /dev/null +++ b/packages/server/workflows/default/track.yaml @@ -0,0 +1,31 @@ +name: Track Gojek Trip +scripts: + - script: fetch-http.js + config: + url: https://api.gojekapi.com/live/track/491c3132c2cc7043 + method: GET + headers: + Content-Type: application/json + Accept: application/json + Origin: https://track.gojek.com + - script: fingerprint.js + config: + key: gojek-track-status + jsonata: data.httpResponse.status + - script: jsonata.js + config: + expression: | + { + "title": data.httpResponse.customer_name, + "message": data.httpResponse.customer_name & " : " & data.httpResponse.status + } + - script: ntfy.js + config: + url: https://n.0dev.web.id/system +triggers: + - type: cron + schedule: "* * * * *" + - type: HTTP + method: GET + path: /new +enabled: false diff --git a/packages/web/index.html b/packages/web/index.html index 41012af..f0550d8 100644 --- a/packages/web/index.html +++ b/packages/web/index.html @@ -3,7 +3,9 @@ - scrunner + + + JerapahFlow
diff --git a/packages/web/package.json b/packages/web/package.json index b22974e..0c484f9 100644 --- a/packages/web/package.json +++ b/packages/web/package.json @@ -1,5 +1,5 @@ { - "name": "@scrunner/web", + "name": "@jerapah-flow/web", "private": true, "version": "1.0.0", "type": "module", @@ -9,9 +9,14 @@ "preview": "vite preview" }, "dependencies": { + "@dnd-kit/core": "^6.3.1", + "@dnd-kit/sortable": "^10.0.0", + "@dnd-kit/utilities": "^3.2.2", "@monaco-editor/react": "^4.7.0", "@tanstack/react-query": "^5.84.0", + "@uiw/react-json-view": "2.0.0-alpha.43", "axios": "^1.11.0", + "cronstrue": "^3.24.0", "mermaid": "^11.9.0", "react": "^19.1.1", "react-dom": "^19.1.1", diff --git a/packages/web/public/apple-touch-icon.png b/packages/web/public/apple-touch-icon.png new file mode 100644 index 0000000..0072eac Binary files /dev/null and b/packages/web/public/apple-touch-icon.png differ diff --git a/packages/web/public/favicon.png b/packages/web/public/favicon.png new file mode 100644 index 0000000..0072eac Binary files /dev/null and b/packages/web/public/favicon.png differ diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index e067c3f..6bf8baf 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -17,6 +17,7 @@ import { AuthProfilesPage } from "./pages/AuthProfilesPage.jsx"; import { ResponsesPage } from "./pages/ResponsesPage.jsx"; import { UsersPage } from "./pages/UsersPage.jsx"; import { SecretsPage } from "./pages/SecretsPage.jsx"; +import { VariablesPage } from "./pages/VariablesPage.jsx"; export function App() { const qc = useQueryClient(); @@ -27,8 +28,8 @@ export function App() { function onUnauthorized() { qc.setQueryData(["me"], null); } - window.addEventListener("scrunner:unauthorized", onUnauthorized); - return () => window.removeEventListener("scrunner:unauthorized", onUnauthorized); + window.addEventListener("jerapah-flow:unauthorized", onUnauthorized); + return () => window.removeEventListener("jerapah-flow:unauthorized", onUnauthorized); }, [qc]); const user = me.data?.user; @@ -61,6 +62,7 @@ export function App() { } /> } /> } /> + } /> } /> } /> {user.role === "admin" ? ( diff --git a/packages/web/src/api/client.js b/packages/web/src/api/client.js index 14ab21e..5d9a7aa 100644 --- a/packages/web/src/api/client.js +++ b/packages/web/src/api/client.js @@ -11,7 +11,7 @@ api.interceptors.response.use( const url = err.config?.url ?? ""; const isAuthCall = url.includes("/auth/login") || url.includes("/auth/register") || url.includes("/auth/me") || url.includes("/auth/bootstrap"); if (err.response?.status === 401 && !isAuthCall) { - window.dispatchEvent(new Event("scrunner:unauthorized")); + window.dispatchEvent(new Event("jerapah-flow:unauthorized")); } return Promise.reject(err); }, diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index ec87e50..e4387f2 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, }) @@ -198,6 +199,27 @@ export function useDeleteWorkflow() { }); } +export function useDuplicateWorkflow() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ owner, file, destOwner, destFile }) => + ( + await api.post( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}/duplicate`, + { + ...(destOwner ? { owner: destOwner } : {}), + ...(destFile ? { file: destFile } : {}), + }, + ) + ).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["owners"] }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + }, + }); +} + export function useRunWorkflow() { const qc = useQueryClient(); return useMutation({ @@ -314,6 +336,33 @@ export function useDeleteSecret() { }); } +export function useVariables(owner) { + return useQuery({ + queryKey: ["variables", owner ?? "all"], + queryFn: async () => { + const params = owner ? { owner } : {}; + return (await api.get("/variables", { params })).data.variables; + }, + }); +} + +export function useUpsertVariable() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (body) => (await api.put("/variables", body)).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["variables"] }), + }); +} + +export function useDeleteVariable() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (id) => + (await api.delete(`/variables/${encodeURIComponent(id)}`)).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["variables"] }), + }); +} + export function useKvNamespaces() { return useQuery({ queryKey: ["kv", "namespaces"], diff --git a/packages/web/src/components/DuplicateWorkflowDialog.jsx b/packages/web/src/components/DuplicateWorkflowDialog.jsx new file mode 100644 index 0000000..d45cb58 --- /dev/null +++ b/packages/web/src/components/DuplicateWorkflowDialog.jsx @@ -0,0 +1,120 @@ +import { useEffect, useMemo, useRef, useState } from "react"; +import { errorMessage } from "../api/client.js"; +import { useDuplicateWorkflow, useOwners, useWorkflows } from "../api/hooks.js"; +import { ensureWorkflowFilename, suggestCopyFilename } from "../lib/workflow-doc.js"; +import { useNotifications } from "../notifications.jsx"; + +const EMPTY_WORKFLOWS = []; + +export function DuplicateWorkflowDialog({ source, warnUnsaved, onClose, onDuplicated }) { + const { notify } = useNotifications(); + const { data: owners = [] } = useOwners(); + const { data: workflows = EMPTY_WORKFLOWS } = useWorkflows(); + const duplicate = useDuplicateWorkflow(); + const [destOwner, setDestOwner] = useState(source.owner); + const [destFile, setDestFile] = useState(() => suggestCopyFilename(source.file)); + const fileTouched = useRef(false); + + const existingFiles = useMemo( + () => workflows.filter((w) => w.owner === destOwner).map((w) => w.file), + [workflows, destOwner], + ); + + useEffect(() => { + if (fileTouched.current) return; + setDestFile(suggestCopyFilename(source.file, existingFiles)); + }, [source.file, existingFiles]); + + const yamlFile = ensureWorkflowFilename(destFile); + const sameAsSource = destOwner === source.owner && yamlFile === source.file; + const exists = existingFiles.includes(yamlFile); + const canSubmit = Boolean(destOwner && yamlFile) && !sameAsSource && !exists && !duplicate.isPending; + + function onSubmit(e) { + e.preventDefault(); + if (!canSubmit) return; + duplicate.mutate( + { + owner: source.owner, + file: source.file, + destOwner, + destFile: yamlFile, + }, + { + onSuccess: (data) => { + notify.success(`Duplicated to ${data.owner}/${data.file}`); + onDuplicated?.(data); + }, + }, + ); + } + + return ( + +
+

Duplicate {source.key}?

+

+ The copy starts disabled. HTTP paths are rewritten when staying under the same owner so + triggers do not collide. +

+ {warnUnsaved ? ( +

The copy uses the last saved YAML, not unsaved edits.

+ ) : null} +
+ + + {sameAsSource ? ( +

Choose a different owner or filename.

+ ) : exists ? ( +

{destOwner}/{yamlFile} already exists.

+ ) : null} + {duplicate.isError ? ( +

{errorMessage(duplicate.error)}

+ ) : null} +
+ + +
+
+
+
+ +
+
+ ); +} diff --git a/packages/web/src/components/JsonViewBlock.jsx b/packages/web/src/components/JsonViewBlock.jsx new file mode 100644 index 0000000..8860e04 --- /dev/null +++ b/packages/web/src/components/JsonViewBlock.jsx @@ -0,0 +1,41 @@ +import JsonView from "@uiw/react-json-view"; +import { darkTheme } from "@uiw/react-json-view/dark"; +import { lightTheme } from "@uiw/react-json-view/light"; +import { useTheme } from "../theme.jsx"; + +/** + * Collapsible JSON tree using @uiw/react-json-view. + * Docs: value (JSON), keyName (root label), collapsed (depth), + * enableClipboard, displayObjectSize, displayDataTypes, style theme vars. + */ +export function JsonViewBlock({ title, value }) { + const { theme } = useTheme() ?? { theme: "light" }; + if (value === undefined) return null; + + const isTree = value != null && typeof value === "object"; + + return ( +
+ {title} +
+
+ +
+
+
+ ); +} diff --git a/packages/web/src/components/Layout.jsx b/packages/web/src/components/Layout.jsx index 5ec3262..501b400 100644 --- a/packages/web/src/components/Layout.jsx +++ b/packages/web/src/components/Layout.jsx @@ -12,9 +12,11 @@ import { LuMoon, LuShield, LuSun, + LuTags, LuUsers, } from "react-icons/lu"; import { useLogout } from "../api/hooks.js"; +import { brandMark } from "../theme/brand.js"; import { useTheme } from "../theme.jsx"; const links = [ @@ -23,6 +25,7 @@ const links = [ { to: "/workflows", label: "Workflows", icon: LuGitBranch }, { to: "/events", label: "Events", icon: LuActivity }, { to: "/kv", label: "KV", icon: LuDatabase }, + { to: "/variables", label: "Variables", icon: LuTags }, { to: "/auth", label: "Auth", icon: LuShield }, { to: "/responses", label: "Responses", icon: LuFileText }, ]; @@ -57,7 +60,10 @@ export function Layout({ user, children }) { -
scrunner
+
+ + JerapahFlow +
+ + {isLoading ? ( +
  • + +
  • + ) : ( + filtered.map((s) => ( +
  • + +
  • + )) + )} + +
    + +
    +
    +
    + +
    + + ); +} diff --git a/packages/web/src/components/workflow/ConfigFields.jsx b/packages/web/src/components/workflow/ConfigFields.jsx new file mode 100644 index 0000000..365a4f4 --- /dev/null +++ b/packages/web/src/components/workflow/ConfigFields.jsx @@ -0,0 +1,473 @@ +import { useEffect, useState } from "react"; +import { LuPlus, LuTrash2 } from "react-icons/lu"; +import { triggerDestinations } from "../../lib/workflow-doc.js"; +import { prettyJson } from "../../lib/script.js"; +import { FieldLabel } from "./FieldHelp.jsx"; + +const MULTILINE_KEYS = new Set(["expression", "jsonata"]); + +const CONFIG_REF_PREFIXES = [ + { prefix: "$SECRET_", label: "secret" }, + { prefix: "$CONTEXT_", label: "context" }, + { prefix: "$VAR_", label: "variable" }, +]; + +function describeConfigRef(value) { + if (typeof value !== "string") return null; + const trimmed = value.trim(); + for (const { prefix, label } of CONFIG_REF_PREFIXES) { + if (trimmed.startsWith(prefix) && trimmed.length > prefix.length) { + return { label, name: trimmed.slice(prefix.length) }; + } + } + return null; +} + +function ConfigRefHint({ value }) { + const ref = describeConfigRef(value); + if (!ref) return null; + return ( +

    + from {ref.label} {ref.name} +

    + ); +} + +function fieldSpec(meta, key) { + const spec = meta?.config?.[key]; + if (spec && typeof spec === "object") return spec; + return {}; +} + +/** Support meta `enum` or `options`: string[], or { value, label }[]. */ +function enumOptions(spec) { + const raw = spec?.enum ?? spec?.options; + if (!Array.isArray(raw) || raw.length === 0) return null; + return raw + .map((item) => { + if (typeof item === "string" || typeof item === "number" || typeof item === "boolean") { + return { value: item, label: String(item) }; + } + if (item && typeof item === "object" && item.value != null) { + return { value: item.value, label: item.label != null ? String(item.label) : String(item.value) }; + } + return null; + }) + .filter(Boolean); +} + +function EditableText({ value, onCommit, className = "", placeholder = "…", disabled }) { + const [editing, setEditing] = useState(false); + const [draft, setDraft] = useState(value ?? ""); + + useEffect(() => { + setDraft(value ?? ""); + }, [value]); + + if (!editing) { + return ( + + ); + } + + return ( + setDraft(e.target.value)} + onBlur={() => { + onCommit(draft); + setEditing(false); + }} + onKeyDown={(e) => { + if (e.key === "Enter") e.currentTarget.blur(); + if (e.key === "Escape") { + setDraft(value ?? ""); + setEditing(false); + } + }} + /> + ); +} + +function ValueEditor({ value, onChange, spec, script, fieldKey, workflows, owner, excludeFile, disabled }) { + const type = spec?.type ?? (value != null && typeof value === "object" ? "object" : "string"); + + if (script === "trigger-workflow.js" && fieldKey === "name") { + return ( + + ); + } + + const options = enumOptions(spec); + if (options) { + const fallback = spec?.default ?? (spec?.required ? options[0].value : ""); + const current = value === undefined || value === null ? fallback : value; + const known = options.some((o) => o.value === current); + return ( + + ); + } + + if (type === "boolean") { + return ( + onChange(e.target.checked)} + /> + ); + } + + if (type === "number") { + return ( + { + const raw = e.target.value; + if (raw === "") { + onChange(undefined); + return; + } + const n = Number(raw); + onChange(Number.isNaN(n) ? raw : n); + }} + /> + ); + } + + if (type === "object" || (value != null && typeof value === "object" && !Array.isArray(value) && type !== "any")) { + const obj = value && typeof value === "object" && !Array.isArray(value) ? value : {}; + return ; + } + + if (type === "any") { + const text = + typeof value === "string" ? value : value === undefined ? "" : prettyJson(value) || ""; + return ( + { + if (next.trim() === "") { + onChange(undefined); + return; + } + try { + onChange(JSON.parse(next)); + } catch { + onChange(next); + } + }} + disabled={disabled} + /> + ); + } + + const str = value == null ? "" : String(value); + const multiline = MULTILINE_KEYS.has(fieldKey) || str.includes("\n"); + if (multiline) { + return ( +