diff --git a/ecosystem.config.cjs b/ecosystem.config.cjs new file mode 100644 index 0000000..f8ef770 --- /dev/null +++ b/ecosystem.config.cjs @@ -0,0 +1,43 @@ +const fs = require("fs"); +const path = require("path"); + +function loadEnv(file) { + const out = {}; + if (!fs.existsSync(file)) return out; + for (const line of fs.readFileSync(file, "utf8").split("\n")) { + const t = line.trim(); + if (!t || t.startsWith("#")) continue; + const i = t.indexOf("="); + if (i === -1) continue; + const key = t.slice(0, i).trim(); + let val = t.slice(i + 1).trim(); + if ( + (val.startsWith('"') && val.endsWith('"')) || + (val.startsWith("'") && val.endsWith("'")) + ) { + val = val.slice(1, -1); + } + out[key] = val; + } + return out; +} + +const root = __dirname; + +module.exports = { + apps: [ + { + name: "jerapah-flow", + cwd: root, + script: "packages/server/runner.js", + interpreter: "node", + instances: 1, + autorestart: true, + max_restarts: 20, + env: { + NODE_ENV: "production", + ...loadEnv(path.join(root, ".env")), + }, + }, + ], +}; \ No newline at end of file diff --git a/packages/server/.gitignore b/packages/server/.gitignore new file mode 100644 index 0000000..24ebcf5 --- /dev/null +++ b/packages/server/.gitignore @@ -0,0 +1 @@ +dump.rdb diff --git a/packages/server/control.js b/packages/server/control.js index b89c38a..f79c885 100644 --- a/packages/server/control.js +++ b/packages/server/control.js @@ -4,7 +4,7 @@ import cors from "@fastify/cors"; import jwt from "@fastify/jwt"; import { migrate, db } from "./db.js"; import { log, enableLogPersistence, flushLogs } from "./logger.js"; -import { COOKIE } from "./src/api/auth.js"; +import authPlugin, { addApiAuthGuard, COOKIE } from "./src/api/auth.js"; import { clearRestartNeeded, bumpGeneration, @@ -35,6 +35,7 @@ import { PM2_HTTP_NAME, PM2_WORKER_NAME, recreateChildren, + restartPm2Process, stopPm2App, } from "./pm2-bridge.js"; @@ -124,6 +125,14 @@ server.decorate("requireAdmin", async function requireAdmin(req, reply) { } }); +await server.register( + async (api) => { + addApiAuthGuard(api, server); + await api.register(authPlugin); + }, + { prefix: "/api" }, +); + /** * @param {number} timeoutMs * @param {string} lockToken @@ -141,6 +150,43 @@ async function waitUntilIdle(timeoutMs, lockToken) { return { ok: active === 0, active }; } +/** @param {unknown} raw */ +function parsePmId(raw) { + if (raw == null || raw === "") return null; + const id = Math.floor(Number(raw)); + if (!Number.isFinite(id) || id < 0) return Number.NaN; + return id; +} + +async function handleProcessRestart(pmId, reply) { + const lock = tryAcquireOpsLock(`process-restart:${process.pid}`); + if (!lock.ok) { + return reply.code(409).send({ error: lock.error, holder: lock.holder }); + } + + try { + const restarted = await restartPm2Process(pmId); + return reply.send({ ok: true, ...restarted, status: await buildStatus() }); + } catch (err) { + const code = err && typeof err === "object" ? err.code : undefined; + if (code === "BAD_REQUEST") { + return reply.code(400).send({ error: "pmId is required" }); + } + if (code === "NOT_FOUND") { + return reply.code(404).send({ error: "process not found" }); + } + if (code === "FORBIDDEN") { + return reply.code(403).send({ error: "process is not a JerapahFlow child" }); + } + log.error({ err, pmId }, "process restart failed"); + return reply.code(500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } finally { + releaseOpsLock(lock.token); + } +} + async function buildStatus() { const state = readControlState(); const children = await describeChildren(); @@ -250,9 +296,18 @@ server.post( "/ops/restart", { onRequest: [server.authenticate, server.requireAdmin] }, async (req, reply) => { - const body = /** @type {{ force?: boolean, timeoutMs?: number }} */ ( + const body = /** @type {{ force?: boolean, timeoutMs?: number, pmId?: number }} */ ( req.body ?? {} ); + + const pmId = parsePmId(body.pmId); + if (pmId != null) { + if (Number.isNaN(pmId)) { + return reply.code(400).send({ error: "pmId is required" }); + } + return handleProcessRestart(pmId, reply); + } + const force = Boolean(body.force); const timeoutMs = Math.min( Math.max(Number(body.timeoutMs) || DRAIN_DEFAULT_MS, 1_000), @@ -334,6 +389,19 @@ server.post( }, ); +server.post( + "/ops/process/restart", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ pmId?: number }} */ (req.body ?? {}); + const pmId = parsePmId(body.pmId); + if (pmId == null || Number.isNaN(pmId)) { + return reply.code(400).send({ error: "pmId is required" }); + } + return handleProcessRestart(pmId, reply); + }, +); + server.post( "/ops/scale", { onRequest: [server.authenticate, server.requireAdmin] }, diff --git a/packages/server/dev-pm2.mjs b/packages/server/dev-pm2.mjs index 92e71d9..fd9c1da 100644 --- a/packages/server/dev-pm2.mjs +++ b/packages/server/dev-pm2.mjs @@ -108,7 +108,7 @@ run( ["--filter", "@jerapah-flow/web", "dev"], { env: { - // vite.config reads nothing; port is set in vite.config.js + JFLOW_AUTH_PROXY: "http://127.0.0.1:8600", }, }, ); diff --git a/packages/server/http-auths-store.js b/packages/server/http-auths-store.js index afb805c..45e4661 100644 --- a/packages/server/http-auths-store.js +++ b/packages/server/http-auths-store.js @@ -4,12 +4,27 @@ import { assertHttpStatus } from "./http-pages-store.js"; const MAX_NAME_LENGTH = 128; const NAME_RE = /^[A-Za-z0-9._-]+$/; +const UUID_RE = + /^[0-9a-f]{8}-[0-9a-f]{4}-[1-8][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i; const ALLOWED_TYPES = new Set(["bearer", "basic", "header"]); function nowIso() { return new Date().toISOString(); } +/** + * @param {unknown} id + * @returns {string} + */ +export function assertAuthId(id) { + if (typeof id !== "string" || !UUID_RE.test(id)) { + const err = new Error("invalid auth id"); + err.statusCode = 400; + throw err; + } + return id.toLowerCase(); +} + /** * @param {unknown} name * @returns {string} @@ -218,10 +233,11 @@ function publicAuth(row, { includeConfig = true } = {}) { /** * Internal: full config including literals (for runtime auth checks). - * @param {string} name + * @param {string} id */ -export async function getHttpAuthInternal(name) { - const row = await db("http_auths").where({ name: assertAuthName(name) }).first(); +export async function getHttpAuthInternal(id) { + const authId = assertAuthId(id); + const row = await db("http_auths").where({ id: authId }).first(); if (!row) return null; return { id: row.id, @@ -235,11 +251,11 @@ export async function getHttpAuthInternal(name) { /** * Return only plaintext literal credential fields (not KV refs or encrypted secrets). - * @param {string} name - * @returns {Promise<{ name: string, type: string, literals: Record } | null>} + * @param {string} id + * @returns {Promise<{ id: string, name: string, type: string, literals: Record } | null>} */ -export async function revealHttpAuthLiterals(name) { - const internal = await getHttpAuthInternal(name); +export async function revealHttpAuthLiterals(id) { + const internal = await getHttpAuthInternal(id); if (!internal) return null; /** @type {Record} */ const literals = {}; @@ -248,7 +264,12 @@ export async function revealHttpAuthLiterals(name) { const v = cfg[key]; if (typeof v === "string") literals[key] = v; } - return { name: internal.name, type: internal.type, literals }; + return { + id: internal.id, + name: internal.name, + type: internal.type, + literals, + }; } export async function listHttpAuths() { @@ -256,24 +277,23 @@ export async function listHttpAuths() { return rows.map((r) => publicAuth(r)); } -/** - * @param {string} name - */ -export async function getHttpAuthByName(name) { - const row = await db("http_auths").where({ name: assertAuthName(name) }).first(); - return row ? publicAuth(row) : null; -} - /** * @param {string} id */ export async function getHttpAuthById(id) { - const row = await db("http_auths").where({ id }).first(); + let authId; + try { + authId = assertAuthId(id); + } catch { + return null; + } + const row = await db("http_auths").where({ id: authId }).first(); return row ? publicAuth(row) : null; } /** * @param {{ + * id?: string | null, * name: string, * type: string, * config?: unknown, @@ -282,6 +302,7 @@ export async function getHttpAuthById(id) { * }} opts */ export async function upsertHttpAuth({ + id, name, type, config, @@ -290,7 +311,26 @@ export async function upsertHttpAuth({ }) { const authName = assertAuthName(name); const authType = assertAuthType(type); - const existing = await db("http_auths").where({ name: authName }).first(); + + /** @type {Record | null} */ + let existing = null; + if (id != null && String(id).length > 0) { + const authId = assertAuthId(id); + existing = await db("http_auths").where({ id: authId }).first(); + if (!existing) { + const err = new Error("auth not found"); + err.statusCode = 404; + throw err; + } + } + + const nameClash = await db("http_auths").where({ name: authName }).first(); + if (nameClash && (!existing || nameClash.id !== existing.id)) { + const err = new Error(`auth name "${authName}" already exists`); + err.statusCode = 409; + throw err; + } + const prevConfig = existing ? parseConfig(existing.config) : {}; const normalized = normalizeAuthConfig(authType, config, { keepLiteralsFrom: prevConfig, @@ -315,18 +355,19 @@ export async function upsertHttpAuth({ await db("http_auths") .where({ id: existing.id }) .update({ + name: authName, type: authType, config: configJson, unauthorized_status: unauthStatus, unauthorized_response: unauthResponse, updated_at: now, }); - return getHttpAuthById(existing.id); + return getHttpAuthById(/** @type {string} */ (existing.id)); } - const id = randomUUID(); + const newId = randomUUID(); await db("http_auths").insert({ - id, + id: newId, name: authName, type: authType, config: configJson, @@ -335,7 +376,7 @@ export async function upsertHttpAuth({ created_at: now, updated_at: now, }); - return getHttpAuthById(id); + return getHttpAuthById(newId); } /** @@ -343,6 +384,7 @@ export async function upsertHttpAuth({ * @returns {Promise} */ export async function deleteHttpAuth(id) { - const n = await db("http_auths").where({ id }).del(); + const authId = assertAuthId(id); + const n = await db("http_auths").where({ id: authId }).del(); return n > 0; } diff --git a/packages/server/http-trigger-auth.js b/packages/server/http-trigger-auth.js index c5d1e86..24b0eff 100644 --- a/packages/server/http-trigger-auth.js +++ b/packages/server/http-trigger-auth.js @@ -77,38 +77,47 @@ export async function resolveCredentialValue(field, ctx) { } /** - * Normalize trigger.auth into an inline auth mechanism object. - * @param {unknown} authField - * @returns {Promise<{ + * @typedef {{ * type: string, * config: Record, * unauthorized_status?: number | null, * unauthorized_response?: string | null, * label: string, - * } | null>} + * }} AuthMechanism */ -export async function resolveAuthMechanism(authField) { - if (authField == null || authField === false) return null; - if (typeof authField === "string") { - const named = await getHttpAuthInternal(authField); - if (!named) { - log.warn({ name: authField }, "http auth: named profile not found"); +/** + * Resolve one auth entry (auth profile id UUID, or inline object). + * @param {unknown} entry + * @returns {Promise} + */ +export async function resolveAuthMechanism(entry) { + if (entry == null || entry === false) return null; + + if (typeof entry === "string") { + try { + const named = await getHttpAuthInternal(entry); + if (!named) { + log.warn({ id: entry }, "http auth: profile id not found"); + return null; + } + return { + type: named.type, + config: named.config, + unauthorized_status: named.unauthorized_status, + unauthorized_response: named.unauthorized_response, + label: named.name, + }; + } catch (err) { + log.warn({ err, id: entry }, "http auth: invalid profile id"); return null; } - return { - type: named.type, - config: named.config, - unauthorized_status: named.unauthorized_status, - unauthorized_response: named.unauthorized_response, - label: authField, - }; } - if (typeof authField === "object" && !Array.isArray(authField)) { - const obj = /** @type {Record} */ (authField); - if (typeof obj.name === "string" && obj.name.length > 0 && !obj.type) { - return resolveAuthMechanism(obj.name); + if (typeof entry === "object" && !Array.isArray(entry)) { + const obj = /** @type {Record} */ (entry); + if (typeof obj.id === "string" && obj.id.length > 0 && !obj.type) { + return resolveAuthMechanism(obj.id); } try { const type = assertAuthType(obj.type); @@ -116,6 +125,7 @@ export async function resolveAuthMechanism(authField) { const config = { ...obj }; delete config.type; delete config.name; + delete config.id; return { type, config, @@ -133,18 +143,64 @@ export async function resolveAuthMechanism(authField) { } /** - * Label for mermaid / summary (sync, no DB). + * Normalize trigger.auth (array of auth ids / inline objects) into mechanisms. + * Empty / null / false → no auth. Any entry that fails to resolve is skipped; + * if the field was non-empty but nothing resolves, returns [] (caller treats as unauthorized). * @param {unknown} authField + * @returns {Promise} */ -export function authLabel(authField) { - if (authField == null) return null; - if (typeof authField === "string") return authField; - if (typeof authField === "object" && !Array.isArray(authField)) { - const o = /** @type {Record} */ (authField); - if (typeof o.name === "string" && o.name) return o.name; - if (typeof o.type === "string" && o.type) return o.type; +export async function resolveAuthMechanisms(authField) { + if (authField == null || authField === false) return []; + if (!Array.isArray(authField) || authField.length === 0) return []; + + /** @type {AuthMechanism[]} */ + const out = []; + for (const entry of authField) { + const mech = await resolveAuthMechanism(entry); + if (mech) out.push(mech); } - return "auth"; + return out; +} + +/** + * True if any mechanism accepts the request (OR). + * @param {import("fastify").FastifyRequest} req + * @param {AuthMechanism[]} mechanisms + * @param {{ owner: string, workflowKey: string }} ctx + */ +export async function checkAnyHttpAuth(req, mechanisms, ctx) { + for (const mechanism of mechanisms) { + if (await checkHttpAuth(req, mechanism, ctx)) return true; + } + return false; +} + +/** + * Label for mermaid / summary (sync). Prefer resolved display names when provided. + * @param {unknown} authField + * @param {Map | Record} [nameById] + */ +export function authLabel(authField, nameById) { + if (authField == null || authField === false) return null; + if (!Array.isArray(authField) || authField.length === 0) return null; + const lookup = + nameById instanceof Map + ? (id) => nameById.get(id) + : nameById + ? (id) => nameById[id] + : () => undefined; + const parts = authField.map((entry) => { + if (typeof entry === "string") return lookup(entry) ?? entry; + if (entry && typeof entry === "object" && !Array.isArray(entry)) { + const o = /** @type {Record} */ (entry); + if (typeof o.id === "string" && o.id && !o.type) { + return lookup(o.id) ?? o.id; + } + if (typeof o.type === "string" && o.type) return o.type; + } + return "auth"; + }); + return parts.join("|"); } /** diff --git a/packages/server/migrations/20260820030000_workflow_history_trash.js b/packages/server/migrations/20260820030000_workflow_history_trash.js new file mode 100644 index 0000000..3401fbe --- /dev/null +++ b/packages/server/migrations/20260820030000_workflow_history_trash.js @@ -0,0 +1,49 @@ +/** + * @param {import("knex").Knex} knex + */ +export async function up(knex) { + await knex.schema.createTable("workflow_revisions", (t) => { + t.text("id").primary(); + t.text("workflow_id").notNullable(); + t.text("owner").notNullable(); + t.text("file").notNullable(); + t.integer("revision").notNullable(); + t.text("content_sha").notNullable(); + t.text("content").notNullable(); + t.text("reason"); + t.text("meta"); + t.text("created_at").notNullable(); + }); + + await knex.schema.raw( + "CREATE UNIQUE INDEX workflow_revisions_workflow_id_revision_idx ON workflow_revisions (workflow_id, revision)", + ); + await knex.schema.raw( + "CREATE INDEX workflow_revisions_workflow_id_created_at_idx ON workflow_revisions (workflow_id, created_at DESC)", + ); + + await knex.schema.createTable("workflow_trash", (t) => { + t.text("id").primary(); + t.text("workflow_id").notNullable(); + t.text("owner").notNullable(); + t.text("file").notNullable(); + t.text("name"); + t.text("deleted_at").notNullable(); + t.text("trash_path").notNullable(); + }); + + await knex.schema.raw( + "CREATE INDEX workflow_trash_deleted_at_idx ON workflow_trash (deleted_at ASC)", + ); + await knex.schema.raw( + "CREATE UNIQUE INDEX workflow_trash_owner_file_idx ON workflow_trash (owner, file)", + ); +} + +/** + * @param {import("knex").Knex} knex + */ +export async function down(knex) { + await knex.schema.dropTableIfExists("workflow_trash"); + await knex.schema.dropTableIfExists("workflow_revisions"); +} diff --git a/packages/server/migrations/20260820050000_workflow_run_revision.js b/packages/server/migrations/20260820050000_workflow_run_revision.js new file mode 100644 index 0000000..7312b3d --- /dev/null +++ b/packages/server/migrations/20260820050000_workflow_run_revision.js @@ -0,0 +1,21 @@ +/** + * @param {import("knex").Knex} knex + */ +export async function up(knex) { + await knex.schema.alterTable("workflow_runs", (t) => { + t.integer("workflow_revision"); + }); + await knex.schema.raw( + "CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)", + ); +} + +/** + * @param {import("knex").Knex} knex + */ +export async function down(knex) { + await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_workflow_revision_idx"); + await knex.schema.alterTable("workflow_runs", (t) => { + t.dropColumn("workflow_revision"); + }); +} diff --git a/packages/server/package.json b/packages/server/package.json index 39d66a7..345b80d 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -12,7 +12,8 @@ "start:worker": "node worker.js", "start:control": "node control.js", "migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\"", - "test:plugins": "JFLOW_PLUGINS_DIR=./data/plugins-smoke-test JFLOW_DB_PATH=./data/plugins-smoke.db node test/plugins-smoke.js" + "test:plugins": "JFLOW_PLUGINS_DIR=./data/plugins-smoke-test JFLOW_DB_PATH=./data/plugins-smoke.db node test/plugins-smoke.js", + "test:workflow-history": "node test/workflow-history-smoke.js" }, "dependencies": { "@aws-sdk/client-s3": "^3.1111.0", diff --git a/packages/server/pm2-bridge.js b/packages/server/pm2-bridge.js index 5d72a28..9bbec47 100644 --- a/packages/server/pm2-bridge.js +++ b/packages/server/pm2-bridge.js @@ -130,6 +130,40 @@ export function restartPm2App(name) { }); } +/** + * Restart a single PM2 process by id. Only jflow-http / jflow-worker. + * @param {number} pmId + * @returns {Promise<{ name: string, pmId: number }>} + */ +export async function restartPm2Process(pmId) { + const id = Math.floor(Number(pmId)); + if (!Number.isFinite(id) || id < 0) { + const err = new Error("invalid pmId"); + err.code = "BAD_REQUEST"; + throw err; + } + + const list = await listPm2(); + const proc = list.find((p) => Number(p.pm_id) === id); + if (!proc) { + const err = new Error("process not found"); + err.code = "NOT_FOUND"; + throw err; + } + if (proc.name !== PM2_HTTP_NAME && proc.name !== PM2_WORKER_NAME) { + const err = new Error("process is not a JerapahFlow child"); + err.code = "FORBIDDEN"; + throw err; + } + await new Promise((resolve, reject) => { + pm2.restart(id, (err) => { + if (err) reject(err); + else resolve(); + }); + }); + return { name: proc.name, pmId: id }; +} + /** * Shared env for child processes. * @param {{ generation: number }} opts @@ -251,19 +285,29 @@ export async function describeChildren() { const mapOne = (p) => ({ name: p.name, - pmId: p.pm_id, + pmId: Number(p.pm_id), status: p.pm2_env?.status ?? "unknown", pid: p.pid ?? null, restarts: p.pm2_env?.restart_time ?? 0, uptime: p.pm2_env?.pm_uptime ?? null, generation: Number(p.pm2_env?.JFLOW_CONFIG_GENERATION ?? 0) || null, + memory: Number(p.monit?.memory) || 0, + cpu: Number(p.monit?.cpu) || 0, }); + const httpMapped = http.map(mapOne); + const workersMapped = workers.map(mapOne); + const all = [...httpMapped, ...workersMapped]; + return { - http: http.map(mapOne), - workers: workers.map(mapOne), + http: httpMapped, + workers: workersMapped, httpOnline: http.some((p) => p.pm2_env?.status === "online"), workerOnlineCount: workers.filter((p) => p.pm2_env?.status === "online").length, + totals: { + memory: all.reduce((sum, p) => sum + p.memory, 0), + cpu: all.reduce((sum, p) => sum + p.cpu, 0), + }, }; } diff --git a/packages/server/registry.js b/packages/server/registry.js index 17d8874..b4bb4e2 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -24,8 +24,8 @@ import { } from "./step-result.js"; import * as fsStore from "./fs-store.js"; import { - checkHttpAuth, - resolveAuthMechanism, + checkAnyHttpAuth, + resolveAuthMechanisms, resolveUnauthorizedSpec, sendHttpPageOrJson, sendSuccessPage, @@ -36,6 +36,7 @@ import { resolveFailureTriggerConfig, } from "./trigger-failure.js"; import { enqueueWorkflowJob } from "./workflow-queue.js"; +import { ensureInitialRevision } from "./workflow-history.js"; /** * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry @@ -57,10 +58,16 @@ function hasWorkflowTrigger(workflow) { /** * @param {import("fastify").FastifyInstance} server - * @param {{ queue?: import("bullmq").Queue | null }} [opts] + * @param {{ + * queue?: import("bullmq").Queue | null, + * enableTriggers?: boolean, + * }} [opts] + * `enableTriggers` must be true only on the API (or all-in-one) process. + * Workers reload workflow defs but must not own cron/HTTP trigger registration. */ export function createRegistry(server, opts = {}) { const queue = opts.queue ?? null; + const enableTriggers = opts.enableTriggers ?? true; /** @type {Map} */ const workflows = new Map(); /** @type {Map} */ @@ -241,22 +248,26 @@ export function createRegistry(server, opts = {}) { return m === method && p === url; }) ?? mapped.trigger; - if (liveTrigger.auth != null && liveTrigger.auth !== false) { - const mechanism = await resolveAuthMechanism(liveTrigger.auth); - if (!mechanism) { + if ( + liveTrigger.auth != null && + liveTrigger.auth !== false && + !(Array.isArray(liveTrigger.auth) && liveTrigger.auth.length === 0) + ) { + const mechanisms = await resolveAuthMechanisms(liveTrigger.auth); + if (mechanisms.length === 0) { const { status, pageName } = resolveUnauthorizedSpec(liveTrigger, null); return sendHttpPageOrJson(reply, status, pageName, { error: "unauthorized", }); } - const ok = await checkHttpAuth(req, mechanism, { + const ok = await checkAnyHttpAuth(req, mechanisms, { owner: entry.owner, workflowKey: mapped.key, }); if (!ok) { const { status, pageName } = resolveUnauthorizedSpec( liveTrigger, - mechanism, + mechanisms[0], ); return sendHttpPageOrJson(reply, status, pageName, { error: "unauthorized", @@ -736,13 +747,14 @@ export function createRegistry(server, opts = {}) { return { runId: null, status: "failed", error: "workflow not found" }; } - const { owner, workflow } = entry; + const { owner, file, workflow } = entry; if (workflow?.enabled === false && trigger.type !== "manual") { log.debug({ workflow: key, trigger }, "skipping disabled workflow"); return { runId: null, status: "failed", error: "workflow disabled" }; } const input = context.data ?? workflow.data ?? null; + const ensured = await ensureInitialRevision({ owner, file }); const run = await store.startRun({ owner, workflow: key, @@ -751,6 +763,7 @@ export function createRegistry(server, opts = {}) { input, parentRunId, status: "queued", + workflowRevision: ensured?.revision ?? null, }); const runLog = log.child({ runId: run.id, owner, workflow: key }); @@ -808,9 +821,17 @@ export function createRegistry(server, opts = {}) { type: existing.trigger_type, detail: existing.trigger_detail, }; + // jobId is BullMQ's id; enqueue uses runId as jobId, so they match today. + const jobId = + existing.job_id != null && String(existing.job_id).length > 0 + ? String(existing.job_id) + : runId; const initialCtx = { data: existing.input ?? workflow.data ?? null, - context: normalizeContext(null), + context: { + runId, + jobId, + }, }; try { @@ -869,6 +890,10 @@ export function createRegistry(server, opts = {}) { function reregister() { registerWorkflows(); + // Cron/HTTP triggers are API-owned. Workers also subscribe to reload and + // must only refresh the in-memory workflow map — otherwise N workers each + // schedule the same cron and enqueue N duplicate jobs. + if (!enableTriggers) return; registerHttpTriggers(); registerCronTriggers(); } diff --git a/packages/server/src/api/auth.js b/packages/server/src/api/auth.js index 8683230..576d2c8 100644 --- a/packages/server/src/api/auth.js +++ b/packages/server/src/api/auth.js @@ -18,6 +18,24 @@ export function cookieOpts() { }; } +/** + * Require JWT for /api routes except bootstrap, login, and register. + * @param {import("fastify").FastifyInstance} api + * @param {import("fastify").FastifyInstance} root + */ +export function addApiAuthGuard(api, root) { + api.addHook("onRequest", async (req, reply) => { + const raw = (req.url || "").split("?")[0]; + const stripped = raw.replace(/^\/api/, "") || "/"; + const routeUrl = req.routeOptions?.url || stripped; + const open = + OPEN_API_ROUTES.has(`${req.method} ${routeUrl}`) || + OPEN_API_ROUTES.has(`${req.method} ${stripped}`); + if (open) return; + await root.authenticate(req, reply); + }); +} + export function validateCredentials(username, password) { if (!/^[A-Za-z0-9_]{3,32}$/.test(username)) { return "username must be 3-32 letters, numbers, or underscore"; diff --git a/packages/server/src/api/http-auths.js b/packages/server/src/api/http-auths.js index c7d85f0..b729953 100644 --- a/packages/server/src/api/http-auths.js +++ b/packages/server/src/api/http-auths.js @@ -1,9 +1,9 @@ import { + assertAuthId, assertAuthName, assertAuthType, listHttpAuths, getHttpAuthById, - getHttpAuthByName, upsertHttpAuth, deleteHttpAuth, revealHttpAuthLiterals, @@ -18,28 +18,28 @@ export default async function httpAuthsPlugin(fastify) { return { auths: await listHttpAuths() }; }); - fastify.get("/http-auths/:name/reveal", async (req, reply) => { - const { name } = /** @type {{ name: string }} */ (req.params); + fastify.get("/http-auths/:id/reveal", async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); try { - assertAuthName(name); + assertAuthId(id); } catch (err) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } - const revealed = await revealHttpAuthLiterals(name); + const revealed = await revealHttpAuthLiterals(id); if (!revealed) { return reply.code(404).send({ error: "auth not found" }); } return revealed; }); - fastify.get("/http-auths/:name", async (req, reply) => { - const { name } = /** @type {{ name: string }} */ (req.params); + fastify.get("/http-auths/:id", async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); try { - assertAuthName(name); + assertAuthId(id); } catch (err) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } - const auth = await getHttpAuthByName(name); + const auth = await getHttpAuthById(id); if (!auth) { return reply.code(404).send({ error: "auth not found" }); } @@ -48,6 +48,7 @@ export default async function httpAuthsPlugin(fastify) { fastify.put("/http-auths", async (req, reply) => { const body = /** @type {{ + id?: string | null, name?: string, type?: string, config?: unknown, @@ -57,6 +58,9 @@ export default async function httpAuthsPlugin(fastify) { try { assertAuthName(String(body.name ?? "")); assertAuthType(body.type); + if (body.id != null && String(body.id).length > 0) { + assertAuthId(String(body.id)); + } if ( body.unauthorized_response != null && String(body.unauthorized_response).length > 0 @@ -69,6 +73,7 @@ export default async function httpAuthsPlugin(fastify) { ); } const auth = await upsertHttpAuth({ + id: body.id != null && String(body.id).length > 0 ? String(body.id) : null, name: String(body.name), type: String(body.type), config: body.config, @@ -83,6 +88,11 @@ export default async function httpAuthsPlugin(fastify) { fastify.delete("/http-auths/:id", async (req, reply) => { const { id } = /** @type {{ id: string }} */ (req.params); + try { + assertAuthId(id); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } const existing = await getHttpAuthById(id); if (!existing) { return reply.code(404).send({ error: "auth not found" }); diff --git a/packages/server/src/api/scripts.js b/packages/server/src/api/scripts.js index 989f0b1..8858255 100644 --- a/packages/server/src/api/scripts.js +++ b/packages/server/src/api/scripts.js @@ -14,10 +14,10 @@ import { resolveScriptRef, uninstallPlugin, createBlankPlugin, + installPluginFromDirectory, } from "../../plugin-store.js"; import { installExamplePlugin, - installPluginFromDirectory, installPluginFromGit, installPluginFromZipBuffer, } from "../../plugin-install.js"; diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index 2904c38..eed98e7 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -10,13 +10,39 @@ import { authLabel, validateWorkflowHttpTriggers, } from "../../workflow-http-validate.js"; +import { listHttpAuths } from "../../http-auths-store.js"; import { validateWorkflowFailureTriggers } from "../../trigger-failure.js"; import { duplicateWorkflowYaml, ensureWorkflowFilename, - suggestCopyFilename, + suggestDuplicateFilename, } from "../../workflow-duplicate.js"; import { publishReload } from "../../control-bus.js"; +import { + workflowIdFromFile, + newWorkflowFilename, +} from "../../workflow-normalize.js"; +import { + collectWorkflowWarnings, + parseWorkflowDocument, +} from "../../workflow-validate-warnings.js"; +import { + recordRevision, + listRevisions, + getRevision, + ensureInitialRevision, +} from "../../workflow-history.js"; +import { + moveWorkflowToTrash, + listTrash, + restoreFromTrash, + purgeTrashItem, + isInTrash, +} from "../../workflow-trash.js"; +import { + createWorkflowBackupBuffer, + restoreWorkflowBackup, +} from "../../workflow-backup.js"; /** * Reload this process and notify other HTTP/worker processes via Redis. @@ -31,7 +57,7 @@ async function reregisterAll(registry) { } } -function triggerSummary(owner, workflow) { +function triggerSummary(owner, workflow, nameById) { if (!workflow || typeof workflow !== "object") return []; return (workflow.triggers ?? []).map((t) => { const type = t?.type ?? "unknown"; @@ -43,7 +69,7 @@ function triggerSummary(owner, workflow) { schedule: t?.schedule ?? null, onConsecutiveFailures: t?.onConsecutiveFailures ?? null, onFailureWorkflow: t?.onFailureWorkflow ?? null, - auth: isHttp ? authLabel(t?.auth) : null, + auth: isHttp ? authLabel(t?.auth, nameById) : null, }; }); } @@ -62,6 +88,92 @@ function scriptNames(workflow) { return names; } +/** + * @param {unknown} parsed + */ +async function validateStrictWorkflow(parsed) { + compileWorkflowScripts(parsed?.scripts); + await validateWorkflowHttpTriggers(parsed); + await validateWorkflowFailureTriggers(parsed); +} + +/** + * @param {{ + * owner: string, + * file: string, + * content: string, + * saveAnyway?: boolean, + * reason?: string | null, + * meta?: Record | null, + * forceRevision?: boolean, + * }} opts + */ +async function saveWorkflowContent(opts) { + const { warnings, parsed, parseError } = collectWorkflowWarnings(opts.content); + const saveAnyway = Boolean(opts.saveAnyway); + + if (!saveAnyway) { + if (parseError) { + const err = new Error("workflow has validation warnings"); + err.statusCode = 422; + err.warnings = warnings; + throw err; + } + try { + await validateStrictWorkflow(parsed); + } catch (validationErr) { + const err = new Error("workflow has validation warnings"); + err.statusCode = 422; + err.warnings = [ + ...warnings, + { + code: "validation_error", + message: + validationErr instanceof Error + ? validationErr.message + : String(validationErr), + }, + ]; + throw err; + } + if (warnings.length) { + const err = new Error("workflow has validation warnings"); + err.statusCode = 422; + err.warnings = warnings; + throw err; + } + } + + const workflowId = workflowIdFromFile(opts.file); + const existed = fsStore.readWorkflowYaml(opts.owner, opts.file) != null; + fsStore.writeWorkflowYaml(opts.owner, opts.file, opts.content); + + const registered = fsStore.readRegisters(opts.owner); + if (!registered.includes(opts.file)) { + registered.push(opts.file); + fsStore.writeRegisters(opts.owner, registered); + } + + const revision = await recordRevision({ + workflowId, + owner: opts.owner, + file: opts.file, + content: opts.content, + reason: opts.reason ?? "save", + meta: opts.meta ?? null, + force: opts.forceRevision, + }); + + return { + owner: opts.owner, + file: opts.file, + workflow_id: workflowId, + existed, + warnings, + revision, + }; +} + /** * @param {{ workflows: Map, loadErrors: Map, reregister: () => void }} registry */ @@ -74,12 +186,83 @@ export default function workflowsPluginFactory(registry) { return { owners: fsStore.listOwners() }; }); + fastify.get("/workflows/trash", async () => { + return { items: await listTrash() }; + }); + + fastify.post("/workflows/trash/:id/restore", async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); + try { + const restored = await restoreFromTrash(id); + await recordRevision({ + workflowId: restored.workflow_id, + owner: restored.owner, + file: restored.file, + content: restored.content, + reason: "restored-from-trash", + force: true, + }); + await reregisterAll(registry); + return { owner: restored.owner, file: restored.file }; + } catch (err) { + return reply.code(err.statusCode ?? 500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + }); + + fastify.delete("/workflows/trash/:id", async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); + try { + return await purgeTrashItem(id); + } catch (err) { + return reply.code(err.statusCode ?? 500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + }); + + fastify.get("/workflows/backup", async (_req, reply) => { + const buffer = await createWorkflowBackupBuffer(); + const stamp = new Date().toISOString().slice(0, 10); + return reply + .header("Content-Type", "application/zip") + .header( + "Content-Disposition", + `attachment; filename="jerapah-flow-backup-${stamp}.zip"`, + ) + .send(buffer); + }); + + fastify.post("/workflows/backup/restore", async (req, reply) => { + const body = /** @type {{ zipBase64?: string, mode?: string }} */ ( + req.body ?? {} + ); + if (typeof body.zipBase64 !== "string" || !body.zipBase64.trim()) { + return reply.code(400).send({ error: "zipBase64 is required" }); + } + const mode = body.mode === "replace" ? "replace" : "merge"; + try { + const buffer = Buffer.from(body.zipBase64, "base64"); + const result = await restoreWorkflowBackup(buffer, { mode }); + await reregisterAll(registry); + return result; + } catch (err) { + return reply.code(400).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + }); + fastify.get("/workflows", async (req) => { const q = /** @type {{ owner?: string }} */ (req.query ?? {}); const stats = await store.workflowStats(); const owners = q.owner ? [fsStore.assertOwner(q.owner)] : fsStore.listOwners(); + const authNameById = Object.fromEntries( + (await listHttpAuths()).map((a) => [a.id, a.name]), + ); const items = []; for (const owner of owners) { @@ -93,6 +276,7 @@ export default function workflowsPluginFactory(registry) { const files = [...new Set([...registered, ...onDisk])]; for (const file of files) { + if (await isInTrash(owner, file)) continue; const key = `${owner}/${file}`; const loaded = registry.workflows.get(key); const loadError = registry.loadErrors.get(key) ?? null; @@ -102,11 +286,8 @@ export default function workflowsPluginFactory(registry) { if (raw != null) { try { parsed = yaml.parse(raw); - } catch (err) { + } catch { // keep loadError - if (!loadError) { - // file on disk but unparseable and not in registers - } } } } @@ -118,18 +299,17 @@ export default function workflowsPluginFactory(registry) { items.push({ owner, file, + workflow_id: workflowIdFromFile(file), key, name: parsed?.name ?? file, description: parsed?.description ?? null, enabled: parsed ? parsed.enabled !== false : false, registered: registered.includes(file), - loadError: - loadError ?? - (parsed ? null : "unreadable"), + loadError: loadError ?? (parsed ? null : "unreadable"), lastInvokedAt: st.lastInvokedAt, lastStatus: st.lastStatus ?? null, invocationCount: st.invocationCount, - triggers: triggerSummary(owner, parsed), + triggers: triggerSummary(owner, parsed, authNameById), scripts: scriptNames(parsed), }); } @@ -137,6 +317,98 @@ export default function workflowsPluginFactory(registry) { return { workflows: items }; }); + fastify.get("/workflows/:owner/:file/revisions", 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 }); + } + if (fsStore.readWorkflowYaml(owner, file) == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + await ensureInitialRevision({ owner, file }); + const workflowId = workflowIdFromFile(file); + return { workflow_id: workflowId, revisions: await listRevisions(workflowId) }; + }); + + fastify.get( + "/workflows/:owner/:file/revisions/:revision", + async (req, reply) => { + const { owner, file, revision } = /** @type {{ owner: string, file: string, revision: string }} */ ( + req.params + ); + try { + fsStore.assertOwner(owner); + fsStore.assertWorkflowFile(file); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + const workflowId = workflowIdFromFile(file); + const rev = await getRevision(workflowId, Number(revision)); + if (!rev) { + return reply.code(404).send({ error: "revision not found" }); + } + return { + workflow_id: workflowId, + revision: rev.revision, + content: rev.content, + reason: rev.reason, + meta: rev.meta, + created_at: rev.created_at, + }; + }, + ); + + fastify.post( + "/workflows/:owner/:file/revisions/:revision/revert", + async (req, reply) => { + const { owner, file, revision } = /** @type {{ owner: string, file: string, revision: string }} */ ( + req.params + ); + try { + fsStore.assertOwner(owner); + fsStore.assertWorkflowFile(file); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + if (fsStore.readWorkflowYaml(owner, file) == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + const workflowId = workflowIdFromFile(file); + const rev = await getRevision(workflowId, Number(revision)); + if (!rev) { + return reply.code(404).send({ error: "revision not found" }); + } + const body = /** @type {{ saveAnyway?: boolean }} */ (req.body ?? {}); + try { + const saved = await saveWorkflowContent({ + owner, + file, + content: rev.content, + saveAnyway: body.saveAnyway, + reason: "revert", + meta: { fromRevision: rev.revision }, + }); + await reregisterAll(registry); + return saved; + } catch (err) { + if (err.statusCode === 422) { + return reply.code(422).send({ + error: err.message, + warnings: err.warnings ?? [], + }); + } + return reply.code(err.statusCode ?? 500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } + }, + ); + fastify.get("/workflows/:owner/:file", async (req, reply) => { const { owner, file } = /** @type {{ owner: string, file: string }} */ ( req.params @@ -167,6 +439,7 @@ export default function workflowsPluginFactory(registry) { return { owner, file, + workflow_id: workflowIdFromFile(file), key, content, parsed, @@ -188,48 +461,95 @@ export default function workflowsPluginFactory(registry) { } catch (err) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } - const body = /** @type {{ content?: string }} */ (req.body ?? {}); + const body = /** @type {{ content?: string, saveAnyway?: boolean }} */ ( + req.body ?? {} + ); if (typeof body.content !== "string") { return reply.code(400).send({ error: "content is required" }); } - let parsed; try { - parsed = yaml.parse(body.content); - } catch (err) { - return reply.code(400).send({ - error: `invalid yaml: ${err instanceof Error ? err.message : String(err)}`, + const saved = await saveWorkflowContent({ + owner, + file, + content: body.content, + saveAnyway: body.saveAnyway, + reason: "save", + }); + await reregisterAll(registry); + return reply.code(saved.existed ? 200 : 201).send({ + owner: saved.owner, + file: saved.file, + workflow_id: saved.workflow_id, + warnings: saved.warnings, + revision: saved.revision, }); - } - try { - compileWorkflowScripts(parsed?.scripts); } catch (err) { - return reply.code(400).send({ + if (err.statusCode === 422) { + return reply.code(422).send({ + error: err.message, + warnings: err.warnings ?? [], + }); + } + return reply.code(err.statusCode ?? 500).send({ error: err instanceof Error ? err.message : String(err), }); } + }); + + fastify.post("/workflows/:owner", async (req, reply) => { + const { owner } = /** @type {{ owner: string }} */ (req.params); try { - await validateWorkflowHttpTriggers(parsed); + fsStore.assertOwner(owner); } catch (err) { - return reply.code(err.statusCode ?? 400).send({ + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + const body = /** @type {{ content?: string, file?: string, saveAnyway?: boolean }} */ ( + req.body ?? {} + ); + if (typeof body.content !== "string") { + return reply.code(400).send({ error: "content is required" }); + } + let file = body.file?.trim() ? ensureWorkflowFilename(body.file) : ""; + if (!file) { + const existing = fsStore.listOwnerYamlFiles(owner); + file = suggestDuplicateFilename(existing); + } + try { + fsStore.assertWorkflowFile(file); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + if (fsStore.readWorkflowYaml(owner, file) != null) { + return reply.code(409).send({ error: "workflow already exists" }); + } + try { + const saved = await saveWorkflowContent({ + owner, + file, + content: body.content, + saveAnyway: body.saveAnyway, + reason: "create", + forceRevision: true, + }); + await reregisterAll(registry); + return reply.code(201).send({ + owner: saved.owner, + file: saved.file, + workflow_id: saved.workflow_id, + warnings: saved.warnings, + revision: saved.revision, + }); + } catch (err) { + if (err.statusCode === 422) { + return reply.code(422).send({ + error: err.message, + warnings: err.warnings ?? [], + }); + } + return reply.code(err.statusCode ?? 500).send({ error: err instanceof Error ? err.message : String(err), }); } - try { - await validateWorkflowFailureTriggers(parsed); - } catch (err) { - return reply.code(err.statusCode ?? 400).send({ - error: err instanceof Error ? err.message : String(err), - }); - } - const existed = fsStore.readWorkflowYaml(owner, file) != null; - fsStore.writeWorkflowYaml(owner, file, body.content); - const registered = fsStore.readRegisters(owner); - if (!registered.includes(file)) { - registered.push(file); - fsStore.writeRegisters(owner, registered); - } - await reregisterAll(registry); - return reply.code(existed ? 200 : 201).send({ owner, file }); }); fastify.patch("/workflows/:owner/:file", async (req, reply) => { @@ -250,23 +570,39 @@ export default function workflowsPluginFactory(registry) { if (content == null) { return reply.code(404).send({ error: "workflow not found" }); } - const doc = yaml.parseDocument(content); - if (doc.errors?.length) { - const msg = doc.errors[0]?.message ?? "invalid yaml"; - return reply.code(400).send({ error: msg }); - } - const parsed = doc.toJSON(); - if (parsed == null || typeof parsed !== "object" || Array.isArray(parsed)) { - return reply.code(400).send({ error: "workflow yaml must be an object" }); + let doc; + try { + ({ doc } = parseWorkflowDocument(content)); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ + error: err instanceof Error ? err.message : String(err), + }); } if (body.enabled) { doc.delete("enabled"); } else { doc.set("enabled", false); } - fsStore.writeWorkflowYaml(owner, file, String(doc)); - await reregisterAll(registry); - return { owner, file, enabled: body.enabled }; + const nextContent = String(doc); + try { + const saved = await saveWorkflowContent({ + owner, + file, + content: nextContent, + reason: body.enabled ? "enable" : "disable", + }); + await reregisterAll(registry); + return { + owner, + file, + enabled: body.enabled, + revision: saved.revision, + }; + } catch (err) { + return reply.code(err.statusCode ?? 500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } }); fastify.delete("/workflows/:owner/:file", async (req, reply) => { @@ -279,13 +615,31 @@ export default function workflowsPluginFactory(registry) { } catch (err) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } - if (!fsStore.deleteWorkflowYaml(owner, file)) { + const raw = fsStore.readWorkflowYaml(owner, file); + if (raw == null) { return reply.code(404).send({ error: "workflow not found" }); } - const registered = fsStore.readRegisters(owner).filter((f) => f !== file); - fsStore.writeRegisters(owner, registered); - await reregisterAll(registry); - return { ok: true }; + let name = null; + try { + const parsed = yaml.parse(raw); + name = parsed?.name ?? null; + } catch { + // ignore + } + try { + const item = await moveWorkflowToTrash({ + workflowId: workflowIdFromFile(file), + owner, + file, + name, + }); + await reregisterAll(registry); + return { ok: true, trash: item }; + } catch (err) { + return reply.code(err.statusCode ?? 500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } }); fastify.post("/workflows/:owner/:file/duplicate", async (req, reply) => { @@ -303,7 +657,9 @@ export default function workflowsPluginFactory(registry) { return reply.code(404).send({ error: "workflow not found" }); } - const body = /** @type {{ file?: unknown, owner?: unknown }} */ (req.body ?? {}); + const body = /** @type {{ file?: unknown, owner?: unknown, saveAnyway?: boolean }} */ ( + req.body ?? {} + ); let destOwner = owner; if (body.owner != null && body.owner !== "") { if (typeof body.owner !== "string") { @@ -319,7 +675,9 @@ export default function workflowsPluginFactory(registry) { let destFile; try { if (body.file == null || body.file === "") { - destFile = suggestCopyFilename(file, fsStore.listOwnerYamlFiles(destOwner)); + destFile = suggestDuplicateFilename( + fsStore.listOwnerYamlFiles(destOwner), + ); } else if (typeof body.file !== "string") { return reply.code(400).send({ error: "file must be a string" }); } else { @@ -350,44 +708,34 @@ export default function workflowsPluginFactory(registry) { }); } - let parsed; try { - parsed = yaml.parse(content); - } catch (err) { - return reply.code(400).send({ - error: `invalid yaml: ${err instanceof Error ? err.message : String(err)}`, + const saved = await saveWorkflowContent({ + owner: destOwner, + file: destFile, + content, + saveAnyway: body.saveAnyway, + reason: "duplicated", + meta: { from: `${owner}/${file}` }, + forceRevision: true, + }); + await reregisterAll(registry); + return reply.code(201).send({ + owner: destOwner, + file: destFile, + workflow_id: saved.workflow_id, + revision: saved.revision, }); - } - try { - compileWorkflowScripts(parsed?.scripts); } catch (err) { - return reply.code(400).send({ + if (err.statusCode === 422) { + return reply.code(422).send({ + error: err.message, + warnings: err.warnings ?? [], + }); + } + return reply.code(err.statusCode ?? 500).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), - }); - } - try { - await validateWorkflowFailureTriggers(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); - } - await reregisterAll(registry); - return reply.code(201).send({ owner: destOwner, file: destFile }); }); fastify.post("/workflows/:owner/:file/run", async (req, reply) => { @@ -443,3 +791,5 @@ export default function workflowsPluginFactory(registry) { }); }; } + +export { newWorkflowFilename }; diff --git a/packages/server/start-app.js b/packages/server/start-app.js index 6f2ddb7..39f0299 100644 --- a/packages/server/start-app.js +++ b/packages/server/start-app.js @@ -8,7 +8,7 @@ import { migrate, db } from "./db.js"; import { log, enableLogPersistence, flushLogs } from "./logger.js"; import * as store from "./store.js"; import { createRegistry } from "./registry.js"; -import { COOKIE, OPEN_API_ROUTES } from "./src/api/auth.js"; +import { addApiAuthGuard, COOKIE } from "./src/api/auth.js"; import authPlugin from "./src/api/auth.js"; import usersPlugin from "./src/api/users.js"; import scriptsPluginFactory from "./src/api/scripts.js"; @@ -28,6 +28,7 @@ import { createWorkflowWorker, getRedisUrlForLog, } from "./workflow-queue.js"; +import { purgeExpiredTrash } from "./workflow-trash.js"; import { getConfigGeneration, startHeartbeatLoop, @@ -48,6 +49,14 @@ export async function startApp(opts = {}) { if (shouldMigrate) { await migrate(); + try { + const purged = await purgeExpiredTrash(); + if (purged > 0) { + log.info({ purged }, "purged expired workflow trash"); + } + } catch (err) { + log.warn({ err }, "workflow trash purge failed"); + } } enableLogPersistence(); @@ -117,7 +126,11 @@ export async function startApp(opts = {}) { } }); - const registry = createRegistry(server, { queue: workflowQueue }); + const registry = createRegistry(server, { + queue: workflowQueue, + // Cron + HTTP triggers enqueue jobs; only the API process may own them. + enableTriggers: runApi, + }); registry.registerWorkflows(); if (runApi) { registry.registerHttpTriggers(); @@ -150,16 +163,7 @@ export async function startApp(opts = {}) { if (runApi) { await server.register( async (api) => { - api.addHook("onRequest", async (req, reply) => { - const raw = (req.url || "").split("?")[0]; - const stripped = raw.replace(/^\/api/, "") || "/"; - const routeUrl = req.routeOptions?.url || stripped; - const open = - OPEN_API_ROUTES.has(`${req.method} ${routeUrl}`) || - OPEN_API_ROUTES.has(`${req.method} ${stripped}`); - if (open) return; - await server.authenticate(req, reply); - }); + addApiAuthGuard(api, server); await api.register(authPlugin); await api.register(usersPlugin); await api.register(secretsPlugin); diff --git a/packages/server/store.js b/packages/server/store.js index b348d29..c680529 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -65,6 +65,7 @@ function nowIso() { * input?: unknown, * parentRunId?: string | null, * status?: "queued" | "running", + * workflowRevision?: number | null, * }} opts */ export async function startRun({ @@ -75,6 +76,7 @@ export async function startRun({ input = null, parentRunId = null, status = "queued", + workflowRevision = null, }) { const id = randomUUID(); const now = nowIso(); @@ -91,6 +93,7 @@ export async function startRun({ queued_at: isQueued ? now : null, input: serialize(input), parent_run_id: parentRunId ?? null, + workflow_revision: workflowRevision ?? null, }); return { id, started_at: now, queued_at: isQueued ? now : null }; } diff --git a/packages/server/test/http-trigger-auth-smoke.js b/packages/server/test/http-trigger-auth-smoke.js index 06e64ea..7d300e5 100644 --- a/packages/server/test/http-trigger-auth-smoke.js +++ b/packages/server/test/http-trigger-auth-smoke.js @@ -12,9 +12,12 @@ import { getHttpAuthInternal, } from "../http-auths-store.js"; import { + authLabel, + checkAnyHttpAuth, checkHttpAuth, coerceCredentialString, resolveAuthMechanism, + resolveAuthMechanisms, resolveCredentialValue, resolveUnauthorizedSpec, sendHttpPageOrJson, @@ -176,11 +179,11 @@ const secret = await upsertSecret({ config: { token: { secret: "does_not_exist_xyz" } }, }, ctx, - ); + ); assert(!missingSec, "missing secret fails closed"); } -// --- named profile --- +// --- named profile (by id) --- const profile = await upsertHttpAuth({ name: "webhook-smoke", type: "bearer", @@ -188,9 +191,10 @@ const profile = await upsertHttpAuth({ unauthorized_status: 403, unauthorized_response: "deny-smoke", }); +assert(typeof profile.id === "string" && profile.id.length > 0, "profile has id"); { - const mech = await resolveAuthMechanism("webhook-smoke"); - assert(mech?.label === "webhook-smoke", "named profile"); + const mech = await resolveAuthMechanism(profile.id); + assert(mech?.label === "webhook-smoke", "profile by id"); const ok = await checkHttpAuth( mockReq({ authorization: "Bearer named-token" }), mech, @@ -202,9 +206,68 @@ const profile = await upsertHttpAuth({ assert(pageName === "deny-smoke", "profile unauth page"); } +// --- rename keeps id --- +{ + const renamed = await upsertHttpAuth({ + id: profile.id, + name: "webhook-renamed", + type: "bearer", + config: { token: { keep: true } }, + unauthorized_status: 403, + unauthorized_response: "deny-smoke", + }); + assert(renamed.id === profile.id, "rename keeps id"); + assert(renamed.name === "webhook-renamed", "rename updates name"); + const mech = await resolveAuthMechanism(profile.id); + assert(mech?.label === "webhook-renamed", "resolve uses new name label"); + const ok = await checkHttpAuth( + mockReq({ authorization: "Bearer named-token" }), + mech, + ctx, + ); + assert(ok, "credentials survive rename"); +} + +// --- multi-auth OR --- +const basicProfile = await upsertHttpAuth({ + name: "basic-smoke", + type: "basic", + config: { user: "bob", password: "p@ss" }, +}); +{ + const mechs = await resolveAuthMechanisms([profile.id, basicProfile.id]); + assert(mechs.length === 2, "resolve two mechanisms"); + assert( + authLabel([profile.id, basicProfile.id], { + [profile.id]: "webhook-renamed", + [basicProfile.id]: "basic-smoke", + }) === "webhook-renamed|basic-smoke", + "authLabel", + ); + const viaBearer = await checkAnyHttpAuth( + mockReq({ authorization: "Bearer named-token" }), + mechs, + ctx, + ); + assert(viaBearer, "OR accepts bearer"); + const encoded = Buffer.from("bob:p@ss").toString("base64"); + const viaBasic = await checkAnyHttpAuth( + mockReq({ authorization: `Basic ${encoded}` }), + mechs, + ctx, + ); + assert(viaBasic, "OR accepts basic"); + const neither = await checkAnyHttpAuth( + mockReq({ authorization: "Bearer wrong" }), + mechs, + ctx, + ); + assert(!neither, "OR rejects when none match"); +} + // trigger-level override { - const mech = await getHttpAuthInternal("webhook-smoke"); + const mech = await getHttpAuthInternal(profile.id); const { status, pageName } = resolveUnauthorizedSpec( { unauthorized: { status: 401, response: "deny-smoke" } }, mech, @@ -226,7 +289,7 @@ await validateWorkflowHttpTriggers({ type: "HTTP", method: "POST", path: "/x", - auth: "webhook-smoke", + auth: [profile.id, basicProfile.id], response: "deny-smoke", }, ], @@ -235,12 +298,38 @@ await validateWorkflowHttpTriggers({ let threw = false; try { await validateWorkflowHttpTriggers({ - triggers: [{ type: "HTTP", path: "/x", auth: "no-such-profile" }], + triggers: [{ type: "HTTP", path: "/x", auth: profile.id }], }); } catch { threw = true; } -assert(threw, "unknown auth fails validation"); +assert(threw, "non-array auth fails validation"); + +threw = false; +try { + await validateWorkflowHttpTriggers({ + triggers: [{ type: "HTTP", path: "/x", auth: ["webhook-renamed"] }], + }); +} catch { + threw = true; +} +assert(threw, "name string fails validation"); + +threw = false; +try { + await validateWorkflowHttpTriggers({ + triggers: [ + { + type: "HTTP", + path: "/x", + auth: ["00000000-0000-4000-8000-000000000000"], + }, + ], + }); +} catch { + threw = true; +} +assert(threw, "unknown auth id fails validation"); threw = false; try { @@ -280,6 +369,7 @@ assert(threw, "unknown page fails validation"); // cleanup await deleteHttpAuth(profile.id); +await deleteHttpAuth(basicProfile.id); await deleteHttpPage(page.id); await deleteSecret(secret.id); const leftover = (await listSecrets({ owner })).find( diff --git a/packages/server/test/workflow-history-smoke.js b/packages/server/test/workflow-history-smoke.js new file mode 100644 index 0000000..3694988 --- /dev/null +++ b/packages/server/test/workflow-history-smoke.js @@ -0,0 +1,163 @@ +/** + * Smoke: workflow revisions, SHA dedup, trash, restore, purge. + * + * Run: pnpm --dir packages/server test:workflow-history + */ +import assert from "node:assert/strict"; +import fs from "fs"; +import path from "path"; +import { fileURLToPath } from "url"; +import { db, migrate } from "../db.js"; +import { WORKFLOWS_DIR } from "../paths.js"; +import * as fsStore from "../fs-store.js"; +import { + workflowContentSha, + workflowIdFromFile, + newWorkflowFilename, +} from "../workflow-normalize.js"; +import { + recordRevision, + listRevisions, + getLatestRevision, + deleteRevisionHistory, + ensureInitialRevision, +} from "../workflow-history.js"; +import { + moveWorkflowToTrash, + listTrash, + restoreFromTrash, + purgeTrashItem, + TRASH_WORKFLOWS_DIR, +} from "../workflow-trash.js"; +import { collectWorkflowWarnings } from "../workflow-validate-warnings.js"; + +const owner = "__workflow_history_smoke__"; +const file = newWorkflowFilename(); +const workflowId = workflowIdFromFile(file); +const ownerDir = path.join(WORKFLOWS_DIR, owner); +const trashPath = path.join(TRASH_WORKFLOWS_DIR, owner, file); + +function cleanup() { + if (fs.existsSync(trashPath)) fs.unlinkSync(trashPath); + if (fs.existsSync(ownerDir)) fs.rmSync(ownerDir, { recursive: true, force: true }); +} + +cleanup(); +await migrate(); + +const yamlV1 = `name: smoke test +scripts: + - plugin/get-current-time +triggers: + - type: HTTP + method: POST + path: /smoke +`; + +fsStore.writeWorkflowYaml(owner, file, yamlV1); +fsStore.writeRegisters(owner, [file]); + +assert.equal(workflowContentSha(yamlV1), workflowContentSha(`${yamlV1}\n\n`)); + +const rev1 = await recordRevision({ + workflowId, + owner, + file, + content: yamlV1, + reason: "create", + force: true, +}); +assert.equal(rev1.skipped, false); +assert.equal(rev1.revision, 1); + +const revDup = await recordRevision({ + workflowId, + owner, + file, + content: `${yamlV1}\n\n`, + reason: "save", +}); +assert.equal(revDup.skipped, true, "normalized SHA should dedupe blank lines"); + +const yamlV2 = yamlV1.replace("smoke test", "smoke test v2"); +const rev2 = await recordRevision({ + workflowId, + owner, + file, + content: yamlV2, + reason: "save", +}); +assert.equal(rev2.revision, 2); +assert.equal((await listRevisions(workflowId)).length, 2); + +const yamlDisabled = `${yamlV2}\nenabled: false\n`; +assert.equal( + workflowContentSha(yamlV2), + workflowContentSha(yamlDisabled), + "enabled-only change should not change content SHA", +); +const revDisable = await recordRevision({ + workflowId, + owner, + file, + content: yamlDisabled, + reason: "disable", +}); +assert.equal(revDisable.skipped, true, "enable/disable should not create a revision"); +assert.equal((await listRevisions(workflowId)).length, 2); + +const { startRun, getRun } = await import("../store.js"); +const latest = await getLatestRevision(workflowId); +assert.equal(latest?.revision, 2); +const run = await startRun({ + owner, + workflow: `${owner}/${file}`, + workflowName: "smoke test v2", + trigger: { type: "manual", detail: "smoke" }, + input: null, + workflowRevision: latest?.revision ?? null, +}); +const loaded = await getRun(run.id); +assert.equal(loaded.workflow_revision, 2); +await db("workflow_runs").where({ id: run.id }).del(); + +await deleteRevisionHistory(workflowId); +assert.equal(await getLatestRevision(workflowId), null); +const seeded = await ensureInitialRevision({ owner, file }); +assert.ok(seeded); +assert.equal(seeded.revision, 1); +assert.equal(seeded.seeded, true); +const again = await ensureInitialRevision({ owner, file }); +assert.equal(again?.revision, 1); +assert.equal(again?.seeded, false); + +const warnings = collectWorkflowWarnings( + `name: bad\nscripts:\n - unknown-script-xyz\n`, +); +assert.ok(warnings.warnings.some((w) => w.code === "unknown_script")); + +const trashed = await moveWorkflowToTrash({ + workflowId, + owner, + file, + name: "smoke test v2", +}); +assert.ok(trashed.id); +assert.equal(fsStore.readWorkflowYaml(owner, file), null); +assert.ok(fs.existsSync(trashPath)); + +const restored = await restoreFromTrash(trashed.id); +assert.equal(restored.file, file); +assert.ok(fsStore.readWorkflowYaml(owner, file)); + +await moveWorkflowToTrash({ workflowId, owner, file, name: "smoke test v2" }); +const trashAgain = (await listTrash()).find((t) => t.file === file); +assert.ok(trashAgain); +await purgeTrashItem(trashAgain.id); +assert.ok(!(await listTrash()).some((t) => t.file === file)); + +await deleteRevisionHistory(workflowId); +cleanup(); + +console.log("workflow-history-smoke: ok"); +await db.destroy(); diff --git a/packages/server/workflow-backup.js b/packages/server/workflow-backup.js new file mode 100644 index 0000000..f3d20d1 --- /dev/null +++ b/packages/server/workflow-backup.js @@ -0,0 +1,134 @@ +import fs from "fs"; +import os from "os"; +import path from "path"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; +import { randomUUID } from "node:crypto"; +import { DATA_DIR, PLUGINS_DIR, WORKFLOWS_DIR } from "./paths.js"; +import { getAppVersion } from "./app-version.js"; +import * as fsStore from "./fs-store.js"; +import { listInstalledPlugins } from "./plugin-store.js"; + +const execFileAsync = promisify(execFile); + +/** + * @param {string} dir + * @param {string} zipPath + */ +async function zipDirectory(dir, zipPath) { + await execFileAsync("zip", ["-r", zipPath, "."], { cwd: dir }); +} + +/** + * @param {string} zipPath + * @param {string} destDir + */ +async function unzipArchive(zipPath, destDir) { + fs.mkdirSync(destDir, { recursive: true }); + await execFileAsync("unzip", ["-o", zipPath, "-d", destDir]); +} + +/** + * Copy directory recursively. + * @param {string} src + * @param {string} dest + */ +function copyDir(src, dest) { + if (!fs.existsSync(src)) return; + fs.mkdirSync(dest, { recursive: true }); + for (const entry of fs.readdirSync(src, { withFileTypes: true })) { + const from = path.join(src, entry.name); + const to = path.join(dest, entry.name); + if (entry.isDirectory()) copyDir(from, to); + else fs.copyFileSync(from, to); + } +} + +/** + * Build a backup zip buffer (workflows + installed plugins + manifest). + */ +export async function createWorkflowBackupBuffer() { + const staging = path.join(DATA_DIR, `.backup-staging-${randomUUID()}`); + fs.mkdirSync(staging, { recursive: true }); + const zipPath = path.join(DATA_DIR, `.backup-${randomUUID()}.zip`); + + try { + const wfDest = path.join(staging, "workflows"); + copyDir(WORKFLOWS_DIR, wfDest); + + const pluginsDest = path.join(staging, "plugins"); + copyDir(PLUGINS_DIR, pluginsDest); + + const manifest = { + version: getAppVersion(), + created_at: new Date().toISOString(), + plugins: listInstalledPlugins().map((p) => p.id), + owners: fsStore.listOwners(), + }; + fs.writeFileSync( + path.join(staging, "manifest.json"), + JSON.stringify(manifest, null, 2), + "utf8", + ); + + await zipDirectory(staging, zipPath); + return fs.readFileSync(zipPath); + } finally { + fs.rmSync(staging, { recursive: true, force: true }); + if (fs.existsSync(zipPath)) fs.unlinkSync(zipPath); + } +} + +/** + * @param {Buffer} zipBuffer + * @param {{ mode?: "merge" | "replace" }} [opts] + */ +export async function restoreWorkflowBackup(zipBuffer, opts = {}) { + const mode = opts.mode === "replace" ? "replace" : "merge"; + const extractDir = fs.mkdtempSync(path.join(os.tmpdir(), "jflow-restore-")); + const zipPath = path.join(extractDir, "backup.zip"); + fs.writeFileSync(zipPath, zipBuffer); + + /** @type {string[]} */ + const warnings = []; + + try { + const contentDir = path.join(extractDir, "content"); + await unzipArchive(zipPath, contentDir); + + const manifestPath = path.join(contentDir, "manifest.json"); + if (fs.existsSync(manifestPath)) { + try { + const manifest = JSON.parse(fs.readFileSync(manifestPath, "utf8")); + for (const pluginId of manifest.plugins ?? []) { + const dir = path.join(PLUGINS_DIR, pluginId); + if (!fs.existsSync(dir)) { + warnings.push(`Plugin "${pluginId}" from backup is not installed`); + } + } + } catch { + warnings.push("Could not read backup manifest.json"); + } + } + + const wfSrc = path.join(contentDir, "workflows"); + if (fs.existsSync(wfSrc)) { + if (mode === "replace" && fs.existsSync(WORKFLOWS_DIR)) { + fs.rmSync(WORKFLOWS_DIR, { recursive: true, force: true }); + } + copyDir(wfSrc, WORKFLOWS_DIR); + } + + const pluginsSrc = path.join(contentDir, "plugins"); + if (fs.existsSync(pluginsSrc)) { + if (mode === "replace" && fs.existsSync(PLUGINS_DIR)) { + fs.rmSync(PLUGINS_DIR, { recursive: true, force: true }); + } + copyDir(pluginsSrc, PLUGINS_DIR); + } + + return { ok: true, mode, warnings }; + } finally { + fs.rmSync(extractDir, { recursive: true, force: true }); + } +} diff --git a/packages/server/workflow-duplicate.js b/packages/server/workflow-duplicate.js index 6e41506..d540b40 100644 --- a/packages/server/workflow-duplicate.js +++ b/packages/server/workflow-duplicate.js @@ -1,4 +1,8 @@ import yaml from "yaml"; +import { + newWorkflowFilename, + workflowFileStem, +} from "./workflow-normalize.js"; export function ensureWorkflowFilename(file) { const trimmed = String(file ?? "").trim(); @@ -6,13 +10,8 @@ export function ensureWorkflowFilename(file) { 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`. + * Legacy human-readable copy name (kept for UI hints). * @param {string} file * @param {string[]} existingFiles */ @@ -33,6 +32,19 @@ export function suggestCopyFilename(file, existingFiles = []) { return candidate(n); } +/** + * UUID-based duplicate filename (default for new duplicates). + * @param {string[]} existingFiles + */ +export function suggestDuplicateFilename(existingFiles = []) { + const existing = new Set(existingFiles); + let file = newWorkflowFilename(); + while (existing.has(file)) { + file = newWorkflowFilename(); + } + return file; +} + export function nextCopyName(name) { const trimmed = String(name ?? "").trim(); if (!trimmed) return "copy"; @@ -55,7 +67,10 @@ export function httpPathCopySuffix(sourceFile, destFile) { 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"; + const safe = + String(suffix) + .replace(/[^A-Za-z0-9._-]+/g, "-") + .replace(/^-+|-+$/g, "") || "copy"; return `${withSlash}-${safe}`; } diff --git a/packages/server/workflow-history.js b/packages/server/workflow-history.js new file mode 100644 index 0000000..5a04eff --- /dev/null +++ b/packages/server/workflow-history.js @@ -0,0 +1,173 @@ +import { randomUUID } from "node:crypto"; +import { db } from "./db.js"; +import * as fsStore from "./fs-store.js"; +import { workflowContentSha, workflowIdFromFile } from "./workflow-normalize.js"; + +const MAX_REVISIONS = 50; + +function nowIso() { + return new Date().toISOString(); +} + +/** + * @param {string | null} meta + */ +function parseMeta(meta) { + if (!meta) return null; + try { + return JSON.parse(meta); + } catch { + return null; + } +} + +/** + * @param {string} workflowId + */ +export async function getLatestRevision(workflowId) { + const row = await db("workflow_revisions") + .where({ workflow_id: workflowId }) + .orderBy("revision", "desc") + .first(); + if (!row) return null; + return { + ...row, + meta: parseMeta(row.meta), + }; +} + +/** + * @param {string} workflowId + */ +export async function listRevisions(workflowId) { + const rows = await db("workflow_revisions") + .where({ workflow_id: workflowId }) + .orderBy("revision", "desc"); + return rows.map((row) => ({ + id: row.id, + workflow_id: row.workflow_id, + owner: row.owner, + file: row.file, + revision: row.revision, + content_sha: row.content_sha, + reason: row.reason ?? null, + meta: parseMeta(row.meta), + created_at: row.created_at, + })); +} + +/** + * @param {string} workflowId + * @param {number} revision + */ +export async function getRevision(workflowId, revision) { + const row = await db("workflow_revisions") + .where({ workflow_id: workflowId, revision }) + .first(); + if (!row) return null; + return { + ...row, + meta: parseMeta(row.meta), + }; +} + +/** + * Ensure a workflow has at least revision #1 (seed from disk when history is empty). + * @param {{ owner: string, file: string }} opts + * @returns {Promise<{ revision: number, id: string, content_sha: string, created_at: string, seeded: boolean } | null>} + */ +export async function ensureInitialRevision(opts) { + const workflowId = workflowIdFromFile(opts.file); + const latest = await getLatestRevision(workflowId); + if (latest) { + return { + revision: latest.revision, + id: latest.id, + content_sha: latest.content_sha, + created_at: latest.created_at, + seeded: false, + }; + } + + const content = fsStore.readWorkflowYaml(opts.owner, opts.file); + if (content == null) return null; + + const recorded = await recordRevision({ + workflowId, + owner: opts.owner, + file: opts.file, + content, + reason: "seed", + force: true, + }); + + if (recorded.revision == null || recorded.id == null) return null; + + const row = await getLatestRevision(workflowId); + if (!row) return null; + + return { + revision: row.revision, + id: row.id, + content_sha: row.content_sha, + created_at: row.created_at, + seeded: true, + }; +} + +/** + * Insert a revision when content changed (SHA dedup skips identical saves). + * @param {{ + * workflowId: string, + * owner: string, + * file: string, + * content: string, + * reason?: string | null, + * meta?: Record | null, + * force?: boolean, + * }} opts + * @returns {Promise<{ skipped: boolean, revision: number | null, id: string | null }>} + */ +export async function recordRevision(opts) { + const sha = workflowContentSha(opts.content); + const latest = await getLatestRevision(opts.workflowId); + if (!opts.force && latest && latest.content_sha === sha) { + return { skipped: true, revision: latest.revision, id: latest.id }; + } + + const nextRevision = latest ? latest.revision + 1 : 1; + const id = randomUUID(); + const created_at = nowIso(); + + await db("workflow_revisions").insert({ + id, + workflow_id: opts.workflowId, + owner: opts.owner, + file: opts.file, + revision: nextRevision, + content_sha: sha, + content: opts.content, + reason: opts.reason ?? null, + meta: opts.meta ? JSON.stringify(opts.meta) : null, + created_at, + }); + + const overflow = await db("workflow_revisions") + .where({ workflow_id: opts.workflowId }) + .orderBy("revision", "desc") + .offset(MAX_REVISIONS) + .pluck("id"); + + if (overflow.length) { + await db("workflow_revisions").whereIn("id", overflow).del(); + } + + return { skipped: false, revision: nextRevision, id }; +} + +/** + * @param {string} workflowId + */ +export async function deleteRevisionHistory(workflowId) { + return db("workflow_revisions").where({ workflow_id: workflowId }).del(); +} diff --git a/packages/server/workflow-http-validate.js b/packages/server/workflow-http-validate.js index 58a0b06..2f44ce7 100644 --- a/packages/server/workflow-http-validate.js +++ b/packages/server/workflow-http-validate.js @@ -1,7 +1,11 @@ /** * Validate HTTP trigger auth / response fields on workflow save. */ -import { assertAuthType, getHttpAuthByName } from "./http-auths-store.js"; +import { + assertAuthId, + assertAuthType, + getHttpAuthById, +} from "./http-auths-store.js"; import { getHttpPageByName, assertHttpResponsePage } from "./http-pages-store.js"; import { authLabel } from "./http-trigger-auth.js"; @@ -24,51 +28,80 @@ function assertCredentialFieldShape(field, label) { } /** - * @param {unknown} auth + * @param {unknown} entry + * @param {string} path */ -async function validateAuthField(auth) { - if (auth == null || auth === false) return; - - if (typeof auth === "string") { - const named = await getHttpAuthByName(auth); +async function validateAuthEntry(entry, path) { + if (typeof entry === "string") { + try { + assertAuthId(entry); + } catch { + const err = new Error(`${path} must be an auth profile UUID`); + err.statusCode = 400; + throw err; + } + const named = await getHttpAuthById(entry); if (!named) { - const err = new Error(`unknown auth profile "${auth}"`); + const err = new Error(`unknown auth profile id "${entry}"`); err.statusCode = 400; throw err; } return; } - if (typeof auth === "object" && !Array.isArray(auth)) { - const obj = /** @type {Record} */ (auth); - if (typeof obj.name === "string" && obj.name.length > 0 && !obj.type) { - await validateAuthField(obj.name); + if (entry && typeof entry === "object" && !Array.isArray(entry)) { + const obj = /** @type {Record} */ (entry); + if (typeof obj.id === "string" && obj.id.length > 0 && !obj.type) { + await validateAuthEntry(obj.id, path); return; } const type = assertAuthType(obj.type); if (type === "bearer") { - assertCredentialFieldShape(obj.token, "auth.token"); + assertCredentialFieldShape(obj.token, `${path}.token`); } else if (type === "basic") { - assertCredentialFieldShape(obj.user, "auth.user"); + assertCredentialFieldShape(obj.user, `${path}.user`); if (obj.password != null && obj.password !== "") { - assertCredentialFieldShape(obj.password, "auth.password"); + assertCredentialFieldShape(obj.password, `${path}.password`); } } else if (type === "header") { if (typeof obj.header !== "string" || obj.header.length === 0) { - const err = new Error("auth.header must be a non-empty string"); + const err = new Error(`${path}.header must be a non-empty string`); err.statusCode = 400; throw err; } - assertCredentialFieldShape(obj.value, "auth.value"); + assertCredentialFieldShape(obj.value, `${path}.value`); } return; } - const err = new Error("auth must be a profile name or an auth object"); + const err = new Error( + `${path} must be an auth profile UUID or an inline auth object`, + ); err.statusCode = 400; throw err; } +/** + * auth is an array of auth profile UUIDs and/or inline auth objects (OR). + * null / false / [] = no auth. + * @param {unknown} auth + */ +async function validateAuthField(auth) { + if (auth == null || auth === false) return; + + if (!Array.isArray(auth)) { + const err = new Error( + "auth must be an array of auth profile UUIDs and/or inline auth objects", + ); + err.statusCode = 400; + throw err; + } + + for (let i = 0; i < auth.length; i++) { + await validateAuthEntry(auth[i], `auth[${i}]`); + } +} + /** * @param {unknown} pageName * @param {string} label diff --git a/packages/server/workflow-normalize.js b/packages/server/workflow-normalize.js new file mode 100644 index 0000000..6407005 --- /dev/null +++ b/packages/server/workflow-normalize.js @@ -0,0 +1,77 @@ +import { createHash, randomUUID } from "node:crypto"; +import yaml from "yaml"; + +/** + * Stable key order for canonical JSON (dedup ignores YAML formatting). + * @param {unknown} value + */ +export function canonicalize(value) { + if (value == null || typeof value !== "object") return value; + if (Array.isArray(value)) return value.map(canonicalize); + const out = {}; + for (const key of Object.keys(value).sort()) { + out[key] = canonicalize(value[key]); + } + return out; +} + +/** + * Parse YAML to a JS object (null when empty/invalid for callers that handle errors). + * @param {string} content + */ +export function parseWorkflowObject(content) { + if (typeof content !== "string" || !content.trim()) return null; + return yaml.parse(content) ?? null; +} + +/** + * Drop `enabled` before hashing so enable/disable does not create revision points. + * @param {unknown} parsed + */ +function stripEnabledForHash(parsed) { + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) return parsed; + const { enabled, ...rest } = parsed; + return rest; +} + +/** + * SHA256 of normalized workflow content (YAML → object → canonical JSON). + * @param {string} content + */ +export function workflowContentSha(content) { + const parsed = stripEnabledForHash(parseWorkflowObject(content)); + const canonical = canonicalize(parsed); + const json = JSON.stringify(canonical); + return createHash("sha256").update(json, "utf8").digest("hex"); +} + +/** + * @param {string} file + */ +export function workflowIdFromFile(file) { + return String(file).replace(/\.ya?ml$/i, ""); +} + +/** @param {string} file */ +export function workflowFileStem(file) { + return workflowIdFromFile(file); +} + +/** + * New on-disk workflow filename: `{uuid}.yaml`. + * @param {string} [uuid] + */ +export function newWorkflowFilename(uuid) { + const id = uuid ?? randomUUID(); + return `${id}.yaml`; +} + +const UUID_FILE_RE = + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\.ya?ml$/i; + +/** + * @param {string} file + */ +export function isUuidWorkflowFile(file) { + return UUID_FILE_RE.test(String(file)); +} diff --git a/packages/server/workflow-trash.js b/packages/server/workflow-trash.js new file mode 100644 index 0000000..fac75e1 --- /dev/null +++ b/packages/server/workflow-trash.js @@ -0,0 +1,183 @@ +import fs from "fs"; +import path from "path"; +import { randomUUID } from "node:crypto"; +import { db } from "./db.js"; +import { deleteRevisionHistory } from "./workflow-history.js"; +import { DATA_DIR, WORKFLOWS_DIR } from "./paths.js"; +import * as fsStore from "./fs-store.js"; + +export const TRASH_WORKFLOWS_DIR = path.join(DATA_DIR, "trash", "workflows"); +export const TRASH_RETENTION_DAYS = 7; + +function nowIso() { + return new Date().toISOString(); +} + +function trashFilePath(owner, file) { + return path.join(TRASH_WORKFLOWS_DIR, owner, file); +} + +/** + * @param {string} deletedAtIso + */ +export function trashAgeMs(deletedAtIso) { + return Date.now() - Date.parse(deletedAtIso); +} + +/** + * @param {string} deletedAtIso + */ +export function trashDaysRemaining(deletedAtIso) { + const purgeAt = + Date.parse(deletedAtIso) + TRASH_RETENTION_DAYS * 24 * 60 * 60 * 1000; + return Math.max(0, Math.ceil((purgeAt - Date.now()) / (24 * 60 * 60 * 1000))); +} + +function rowToItem(row) { + return { + id: row.id, + workflow_id: row.workflow_id, + owner: row.owner, + file: row.file, + name: row.name ?? null, + deleted_at: row.deleted_at, + trash_path: row.trash_path, + age_ms: trashAgeMs(row.deleted_at), + days_until_purge: trashDaysRemaining(row.deleted_at), + }; +} + +export async function listTrash() { + const rows = await db("workflow_trash").orderBy("deleted_at", "desc"); + return rows.map(rowToItem); +} + +export async function getTrashItem(id) { + const row = await db("workflow_trash").where({ id }).first(); + return row ? rowToItem(row) : null; +} + +export async function isInTrash(owner, file) { + const row = await db("workflow_trash").where({ owner, file }).first(); + return Boolean(row); +} + +/** + * Soft-delete: move YAML to trash dir, unregister, keep revision history. + * @param {{ + * workflowId: string, + * owner: string, + * file: string, + * name?: string | null, + * }} opts + */ +export async function moveWorkflowToTrash(opts) { + const sourcePath = path.join(WORKFLOWS_DIR, opts.owner, opts.file); + if (!fs.existsSync(sourcePath)) { + const err = new Error("workflow not found"); + err.statusCode = 404; + throw err; + } + + const trashPath = trashFilePath(opts.owner, opts.file); + fs.mkdirSync(path.dirname(trashPath), { recursive: true }); + fs.renameSync(sourcePath, trashPath); + + const registered = fsStore.readRegisters(opts.owner).filter((f) => f !== opts.file); + fsStore.writeRegisters(opts.owner, registered); + + const id = randomUUID(); + const deleted_at = nowIso(); + await db("workflow_trash").insert({ + id, + workflow_id: opts.workflowId, + owner: opts.owner, + file: opts.file, + name: opts.name ?? null, + deleted_at, + trash_path: trashPath, + }); + + return rowToItem(await db("workflow_trash").where({ id }).first()); +} + +/** + * Restore workflow from trash. + * @param {string} trashId + */ +export async function restoreFromTrash(trashId) { + const row = await db("workflow_trash").where({ id: trashId }).first(); + if (!row) { + const err = new Error("trash item not found"); + err.statusCode = 404; + throw err; + } + + const destPath = path.join(WORKFLOWS_DIR, row.owner, row.file); + if (fs.existsSync(destPath)) { + const err = new Error("workflow file already exists"); + err.statusCode = 409; + throw err; + } + if (!fs.existsSync(row.trash_path)) { + const err = new Error("trash file missing on disk"); + err.statusCode = 410; + throw err; + } + + fs.mkdirSync(path.dirname(destPath), { recursive: true }); + fs.renameSync(row.trash_path, destPath); + + const registered = fsStore.readRegisters(row.owner); + if (!registered.includes(row.file)) { + registered.push(row.file); + fsStore.writeRegisters(row.owner, registered); + } + + await db("workflow_trash").where({ id: trashId }).del(); + + return { + owner: row.owner, + file: row.file, + workflow_id: row.workflow_id, + content: fs.readFileSync(destPath, "utf8"), + }; +} + +/** + * Permanently delete a trash item and its revision history. + * @param {string} trashId + */ +export async function purgeTrashItem(trashId) { + const row = await db("workflow_trash").where({ id: trashId }).first(); + if (!row) { + const err = new Error("trash item not found"); + err.statusCode = 404; + throw err; + } + + if (fs.existsSync(row.trash_path)) { + fs.unlinkSync(row.trash_path); + } + + await deleteRevisionHistory(row.workflow_id); + await db("workflow_trash").where({ id: trashId }).del(); + return { ok: true }; +} + +/** + * Auto-purge trash older than retention window. + * @returns {Promise} + */ +export async function purgeExpiredTrash() { + const cutoff = new Date( + Date.now() - TRASH_RETENTION_DAYS * 24 * 60 * 60 * 1000, + ).toISOString(); + const rows = await db("workflow_trash") + .where("deleted_at", "<", cutoff) + .select("id"); + for (const row of rows) { + await purgeTrashItem(row.id); + } + return rows.length; +} diff --git a/packages/server/workflow-validate-warnings.js b/packages/server/workflow-validate-warnings.js new file mode 100644 index 0000000..91aad42 --- /dev/null +++ b/packages/server/workflow-validate-warnings.js @@ -0,0 +1,136 @@ +import yaml from "yaml"; +import { parseScriptStep } from "./workflow-parse.js"; +import { resolveScriptRef } from "./plugin-store.js"; +import { parseWorkflowObject } from "./workflow-normalize.js"; + +const SECRET_KEY_RE = + /(?:password|passwd|secret|token|api[_-]?key|auth(?:orization)?|credential|private[_-]?key)/i; + +const BEARER_RE = /Bearer\s+[A-Za-z0-9._~+/=-]{8,}/; + +/** + * Walk parsed YAML for suspicious secret-like string values. + * @param {unknown} value + * @param {string} pathKey + * @param {Array<{ code: string, message: string, path?: string }>} warnings + */ +function scanSecrets(value, pathKey, warnings) { + if (value == null) return; + if (typeof value === "string") { + if (BEARER_RE.test(value)) { + warnings.push({ + code: "plaintext_secret", + message: "Possible Bearer token in workflow YAML", + path: pathKey, + }); + } + return; + } + if (Array.isArray(value)) { + value.forEach((item, i) => scanSecrets(item, `${pathKey}[${i}]`, warnings)); + return; + } + if (typeof value === "object") { + for (const [k, v] of Object.entries(value)) { + const childPath = pathKey ? `${pathKey}.${k}` : k; + if (typeof v === "string" && v.trim() && SECRET_KEY_RE.test(k)) { + warnings.push({ + code: "plaintext_secret", + message: `Possible secret in field "${k}"`, + path: childPath, + }); + } + scanSecrets(v, childPath, warnings); + } + } +} + +/** + * Collect non-blocking save warnings for workflow YAML. + * @param {string} content + * @returns {{ warnings: Array<{ code: string, message: string, path?: string }>, parsed: unknown | null, parseError: string | null }} + */ +export function collectWorkflowWarnings(content) { + /** @type {Array<{ code: string, message: string, path?: string }>} */ + const warnings = []; + + let parsed = null; + let parseError = null; + try { + parsed = parseWorkflowObject(content); + if (parsed == null) { + warnings.push({ + code: "invalid_yaml", + message: "Workflow YAML is empty or not an object", + }); + } else if (typeof parsed !== "object" || Array.isArray(parsed)) { + warnings.push({ + code: "invalid_yaml", + message: "Workflow YAML must be a mapping/object", + }); + parsed = null; + } + } catch (err) { + parseError = err instanceof Error ? err.message : String(err); + warnings.push({ + code: "invalid_yaml", + message: `Invalid YAML: ${parseError}`, + }); + } + + if (parsed && typeof parsed === "object" && !Array.isArray(parsed)) { + scanSecrets(parsed, "", warnings); + + for (const [i, raw] of (parsed.scripts ?? []).entries()) { + try { + const step = parseScriptStep(raw); + if (step.kind === "set") continue; + const resolved = resolveScriptRef(step.script); + if (resolved.error) { + warnings.push({ + code: "unknown_script", + message: resolved.error, + path: `scripts[${i}]`, + }); + } + } catch (err) { + warnings.push({ + code: "invalid_script_step", + message: err instanceof Error ? err.message : String(err), + path: `scripts[${i}]`, + }); + } + } + } + + return { warnings, parsed, parseError }; +} + +/** + * Strict validation used when saveAnyway is false. + * @param {unknown} parsed + */ +export function assertStrictWorkflow(parsed) { + if (parsed == null || typeof parsed !== "object" || Array.isArray(parsed)) { + const err = new Error("workflow yaml must be an object"); + err.statusCode = 400; + throw err; + } + return parsed; +} + +/** + * Parse for PATCH/enable toggles (must be valid YAML document). + * @param {string} content + */ +export function parseWorkflowDocument(content) { + 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(); + assertStrictWorkflow(parsed); + return { doc, parsed }; +} diff --git a/packages/server/workflows/default/dev-zte-sms.yaml b/packages/server/workflows/default/dev-zte-sms.yaml index e351541..1a85a56 100644 --- a/packages/server/workflows/default/dev-zte-sms.yaml +++ b/packages/server/workflows/default/dev-zte-sms.yaml @@ -13,4 +13,5 @@ triggers: - type: HTTP method: POST path: /dev-zte-sms - auth: basic-auth + auth: + - 0f78d6d7-bd44-45d7-a826-f51c027b767f diff --git a/packages/server/workflows/default/send-gmail.yaml b/packages/server/workflows/default/send-gmail.yaml index d6b42cf..e0d5889 100644 --- a/packages/server/workflows/default/send-gmail.yaml +++ b/packages/server/workflows/default/send-gmail.yaml @@ -36,4 +36,5 @@ triggers: - type: HTTP method: POST path: /send-gmail - auth: basic-auth + auth: + - 0f78d6d7-bd44-45d7-a826-f51c027b767f diff --git a/packages/server/workflows/default/time-to-ntfy-example.yaml b/packages/server/workflows/default/time-to-ntfy-example.yaml index 34404f9..67dcdf3 100644 --- a/packages/server/workflows/default/time-to-ntfy-example.yaml +++ b/packages/server/workflows/default/time-to-ntfy-example.yaml @@ -1,6 +1,7 @@ name: time to ntfy example description: | this workflow will send a message to ntfy with the current time +enabled: false scripts: - script: plugin/get-current-time config: @@ -12,3 +13,5 @@ triggers: - type: HTTP method: POST path: /time-to-ntfy + auth: + - 0f78d6d7-bd44-45d7-a826-f51c027b767f diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index 579c487..351f59f 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -10,6 +10,7 @@ import { ScriptDryRunPage } from "./pages/ScriptDryRunPage.jsx"; import { ScriptEditPage, ScriptNewPage } from "./pages/ScriptEditPage.jsx"; import { WorkflowsPage } from "./pages/WorkflowsPage.jsx"; import { WorkflowEditPage, WorkflowNewPage } from "./pages/WorkflowEditPage.jsx"; +import { WorkflowTrashPage } from "./pages/WorkflowTrashPage.jsx"; import { EventsPage } from "./pages/EventsPage.jsx"; import { EventDetailPage } from "./pages/EventDetailPage.jsx"; import { FailuresPage } from "./pages/FailuresPage.jsx"; @@ -20,6 +21,7 @@ import { UsersPage } from "./pages/UsersPage.jsx"; import { SecretsPage } from "./pages/SecretsPage.jsx"; import { VariablesPage } from "./pages/VariablesPage.jsx"; import { OpsPage } from "./pages/OpsPage.jsx"; +import { BackupPage } from "./pages/BackupPage.jsx"; export function App() { const qc = useQueryClient(); @@ -59,6 +61,7 @@ export function App() { } /> } /> } /> + } /> } /> } /> } /> @@ -66,19 +69,33 @@ export function App() { } /> } /> } /> + } /> + } /> } /> + } /> + } /> } /> {user.role === "admin" ? ( <> } /> + } /> + } /> } /> - } /> + } /> + } /> + } /> + } /> + } /> ) : ( <> } /> + } /> } /> + } /> + } /> } /> + } /> )} } /> diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index d7fed6d..cd07e4b 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -191,17 +191,21 @@ export function useOwners() { export function useSaveWorkflow() { const qc = useQueryClient(); return useMutation({ - mutationFn: async ({ owner, file, content }) => + mutationFn: async ({ owner, file, content, saveAnyway }) => ( await api.put( `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}`, - { content }, + { content, ...(saveAnyway ? { saveAnyway: true } : {}) }, ) ).data, - onSuccess: () => { + onSuccess: (_data, vars) => { qc.invalidateQueries({ queryKey: ["workflows"] }); qc.invalidateQueries({ queryKey: ["owners"] }); qc.invalidateQueries({ queryKey: ["dashboard"] }); + qc.invalidateQueries({ + queryKey: ["workflows", vars.owner, vars.file, "revisions"], + }); + qc.invalidateQueries({ queryKey: ["workflows", vars.owner, vars.file] }); }, }); } @@ -235,6 +239,7 @@ export function useDeleteWorkflow() { ).data, onSuccess: () => { qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["workflows", "trash"] }); qc.invalidateQueries({ queryKey: ["dashboard"] }); }, }); @@ -390,13 +395,14 @@ export function useDeleteSecret() { }); } -export function useVariables(owner) { +export function useVariables(owner, options = {}) { return useQuery({ queryKey: ["variables", owner ?? "all"], queryFn: async () => { const params = owner ? { owner } : {}; return (await api.get("/variables", { params })).data.variables; }, + ...options, }); } @@ -488,8 +494,8 @@ export function useDeleteHttpAuth() { } /** Fetch plaintext literals only (not encrypted secrets). */ -export async function fetchHttpAuthLiterals(name) { - return (await api.get(`/http-auths/${encodeURIComponent(name)}/reveal`)).data; +export async function fetchHttpAuthLiterals(id) { + return (await api.get(`/http-auths/${encodeURIComponent(id)}/reveal`)).data; } export function useOpsStatus(enabled = true) { @@ -563,6 +569,15 @@ export function useOpsHttpStop() { }); } +export function useOpsProcessRestart() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ pmId }) => + (await opsApi.post("/restart", { pmId: Number(pmId) })).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + export function useOpsBumpGeneration() { const qc = useQueryClient(); return useMutation({ @@ -572,3 +587,137 @@ export function useOpsBumpGeneration() { }); } +export function useWorkflowTrash() { + return useQuery({ + queryKey: ["workflows", "trash"], + queryFn: async () => (await api.get("/workflows/trash")).data.items, + }); +} + +export function useRestoreWorkflowTrash() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (id) => + (await api.post(`/workflows/trash/${encodeURIComponent(id)}/restore`)).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["workflows", "trash"] }); + qc.invalidateQueries({ queryKey: ["owners"] }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + }, + }); +} + +export function usePurgeWorkflowTrash() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (id) => + (await api.delete(`/workflows/trash/${encodeURIComponent(id)}`)).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["workflows", "trash"] }); + }, + }); +} + +export function useWorkflowRevisions(owner, file, enabled = true) { + return useQuery({ + queryKey: ["workflows", owner, file, "revisions"], + queryFn: async () => + ( + await api.get( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}/revisions`, + ) + ).data, + enabled: Boolean(owner && file) && enabled, + }); +} + +export function useWorkflowRevision(owner, file, revision) { + return useQuery({ + queryKey: ["workflows", owner, file, "revisions", revision], + queryFn: async () => + ( + await api.get( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}/revisions/${revision}`, + ) + ).data, + enabled: Boolean(owner && file && revision != null), + }); +} + +export function useRevertWorkflowRevision() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ owner, file, revision, saveAnyway }) => + ( + await api.post( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}/revisions/${revision}/revert`, + saveAnyway ? { saveAnyway: true } : {}, + ) + ).data, + onSuccess: (_data, vars) => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["workflows", vars.owner, vars.file] }); + qc.invalidateQueries({ + queryKey: ["workflows", vars.owner, vars.file, "revisions"], + }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + }, + }); +} + +export function useCreateWorkflow() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ owner, content, file, saveAnyway }) => + ( + await api.post(`/workflows/${encodeURIComponent(owner)}`, { + content, + ...(file ? { file } : {}), + ...(saveAnyway ? { saveAnyway: true } : {}), + }) + ).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["owners"] }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + }, + }); +} + +export function useDownloadWorkflowBackup() { + return useMutation({ + mutationFn: async () => { + const res = await api.get("/workflows/backup", { responseType: "blob" }); + const disposition = res.headers["content-disposition"] ?? ""; + const match = disposition.match(/filename="([^"]+)"/); + const filename = match?.[1] ?? "jerapah-flow-backup.zip"; + const url = URL.createObjectURL(res.data); + const a = document.createElement("a"); + a.href = url; + a.download = filename; + a.click(); + URL.revokeObjectURL(url); + return { ok: true }; + }, + }); +} + +export function useRestoreWorkflowBackup() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ file, mode }) => { + const buffer = await file.arrayBuffer(); + const zipBase64 = btoa(String.fromCharCode(...new Uint8Array(buffer))); + return (await api.post("/workflows/backup/restore", { zipBase64, mode })).data; + }, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["owners"] }); + qc.invalidateQueries({ queryKey: ["scripts"] }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + qc.invalidateQueries({ queryKey: ["ops-status"] }); + }, + }); +} + diff --git a/packages/web/src/components/AuthEditorModal.jsx b/packages/web/src/components/AuthEditorModal.jsx new file mode 100644 index 0000000..8b01e7d --- /dev/null +++ b/packages/web/src/components/AuthEditorModal.jsx @@ -0,0 +1,422 @@ +import { useEffect, useState } from "react"; +import { LuEye, LuEyeOff } from "react-icons/lu"; +import { errorMessage } from "../api/client.js"; +import { + fetchHttpAuthLiterals, + useHttpPages, + useUpsertHttpAuth, +} from "../api/hooks.js"; + +function emptyCred(source = "literal") { + return { source, value: "", kv: "", namespace: "", secret: "" }; +} + +function credFromPublic(field, literalValue) { + if (!field || field.source === "missing") return emptyCred("literal"); + if (field.source === "kv") { + return { + source: "kv", + value: "", + kv: field.kv ?? "", + namespace: field.namespace ?? "", + secret: "", + }; + } + if (field.source === "secret") { + return { + source: "secret", + value: "", + kv: "", + namespace: "", + secret: field.secret ?? "", + }; + } + if (typeof literalValue === "string") { + return { + source: "literal", + value: literalValue, + kv: "", + namespace: "", + secret: "", + keep: true, + }; + } + return { + source: "literal", + value: "", + kv: "", + namespace: "", + secret: "", + keep: field.set === true, + }; +} + +function toApiField(cred, { required = true } = {}) { + if (cred.source === "kv") { + const out = { kv: cred.kv }; + if (cred.namespace) out.namespace = cred.namespace; + return out; + } + if (cred.source === "secret") { + return { secret: cred.secret }; + } + if (cred.value) return cred.value; + if (cred.keep) return { keep: true }; + if (!required) return ""; + return null; +} + +function emptyForm() { + return { + id: null, + name: "", + type: "bearer", + token: emptyCred(), + user: emptyCred(), + password: emptyCred(), + header: "", + value: emptyCred(), + unauthorized_status: "", + unauthorized_response: "", + }; +} + +function formFromAuth(auth, literals = {}) { + const cfg = auth.config ?? {}; + return { + id: auth.id, + name: auth.name, + type: auth.type, + token: credFromPublic(cfg.token, literals.token), + user: credFromPublic(cfg.user, literals.user), + password: credFromPublic(cfg.password, literals.password), + header: cfg.header ?? "", + value: credFromPublic(cfg.value, literals.value), + unauthorized_status: auth.unauthorized_status ?? "", + unauthorized_response: auth.unauthorized_response ?? "", + }; +} + +function Field({ label, children, hint }) { + return ( +
+
+ {label} +
+ {children} + {hint ? ( +
+ {hint} +
+ ) : null} +
+ ); +} + +function CredentialFields({ label, cred, onChange, allowEmpty, masked }) { + const [show, setShow] = useState(false); + + return ( +
+

{label}

+ + + + {cred.source === "literal" ? ( + +
+ onChange({ ...cred, value: e.target.value, keep: false })} + placeholder={ + cred.keep && !cred.value ? "(unchanged — leave blank to keep)" : "" + } + required={!allowEmpty && !cred.keep && !cred.value} + autoComplete="off" + /> + {masked ? ( + + ) : null} +
+
+ ) : null} + {cred.source === "kv" ? ( + <> + + onChange({ ...cred, namespace: e.target.value })} + /> + + + onChange({ ...cred, kv: e.target.value })} + required + /> + + + ) : null} + {cred.source === "secret" ? ( + + onChange({ ...cred, secret: e.target.value })} + required + pattern="[A-Za-z0-9._-]+" + /> + + ) : null} +
+ ); +} + +/** + * Add / edit an HTTP trigger auth profile. + * Reusable: mount when open; pass `auth` for edit (literals loaded inside). + * + * @param {"add" | "edit"} mode + * @param {object} [auth] Public auth row when mode is "edit" + * @param {() => void} onClose + * @param {(saved: unknown) => void} [onSaved] + */ +export function AuthEditorModal({ mode, auth, onClose, onSaved }) { + const { data: pages = [] } = useHttpPages(); + const upsert = useUpsertHttpAuth(); + const [form, setForm] = useState(emptyForm); + const [loading, setLoading] = useState(mode === "edit"); + + useEffect(() => { + if (mode !== "edit" || !auth?.id) { + setForm(emptyForm()); + setLoading(false); + return; + } + let cancelled = false; + setLoading(true); + (async () => { + /** @type {Record} */ + let literals = {}; + try { + const data = await fetchHttpAuthLiterals(auth.id); + literals = data.literals ?? {}; + } catch { + // Form still works with keep markers + } + if (cancelled) return; + setForm(formFromAuth(auth, literals)); + setLoading(false); + })(); + return () => { + cancelled = true; + }; + }, [mode, auth]); + + function onSubmit(e) { + e.preventDefault(); + /** @type {Record} */ + let config = {}; + if (form.type === "bearer") { + const token = toApiField(form.token); + if (token == null) return; + config = { token }; + } else if (form.type === "basic") { + const user = toApiField(form.user); + if (user == null) return; + const password = toApiField(form.password, { required: false }); + config = { user, password: password ?? "" }; + } else { + const value = toApiField(form.value); + if (value == null) return; + config = { header: form.header, value }; + } + + upsert.mutate( + { + id: form.id, + name: form.name, + type: form.type, + config, + unauthorized_status: + form.unauthorized_status === "" ? null : Number(form.unauthorized_status), + unauthorized_response: form.unauthorized_response || null, + }, + { + onSuccess: (data) => { + onSaved?.(data?.auth ?? data); + onClose(); + }, + }, + ); + } + + const title = mode === "add" ? "New auth profile" : `Edit ${form.name || auth?.name || ""}`; + + return ( + +
+

{title}

+ {loading ? ( +
+ +
+ ) : ( +
+ + setForm({ ...form, name: e.target.value })} + required + pattern="[A-Za-z0-9._-]+" + autoComplete="off" + /> + + + + + + + {form.type === "bearer" ? ( + setForm({ ...form, token })} + /> + ) : null} + {form.type === "basic" ? ( + <> + setForm({ ...form, user })} + /> + setForm({ ...form, password })} + allowEmpty + masked + /> + + ) : null} + {form.type === "header" ? ( + <> + + setForm({ ...form, header: e.target.value })} + required + placeholder="X-Webhook-Secret" + /> + + setForm({ ...form, value })} + /> + + ) : null} + + + setForm({ ...form, unauthorized_status: e.target.value })} + min={100} + max={599} + placeholder="401" + /> + + + + + + + {upsert.isError ? ( +

{errorMessage(upsert.error)}

+ ) : null} +
+ + +
+ + )} + {loading ? ( +
+ +
+ ) : null} +
+
+ +
+
+ ); +} diff --git a/packages/web/src/components/DuplicateWorkflowDialog.jsx b/packages/web/src/components/DuplicateWorkflowDialog.jsx index d45cb58..168d004 100644 --- a/packages/web/src/components/DuplicateWorkflowDialog.jsx +++ b/packages/web/src/components/DuplicateWorkflowDialog.jsx @@ -1,44 +1,22 @@ -import { useEffect, useMemo, useRef, useState } from "react"; +import { 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 { useDuplicateWorkflow, useOwners } from "../api/hooks.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; + if (destOwner === source.owner && duplicate.isPending) return; duplicate.mutate( { owner: source.owner, file: source.file, destOwner, - destFile: yamlFile, }, { onSuccess: (data) => { @@ -54,11 +32,13 @@ export function DuplicateWorkflowDialog({ source, warnUnsaved, onClose, onDuplic

Duplicate {source.key}?

- The copy starts disabled. HTTP paths are rewritten when staying under the same owner so - triggers do not collide. + A new UUID filename is assigned automatically. 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.

+

+ 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} @@ -101,7 +63,7 @@ export function DuplicateWorkflowDialog({ source, warnUnsaved, onClose, onDuplic - + +
+ + + + {workerCount} worker{workerCount === 1 ? "" : "s"} + + + + +
+ + +
+ +
+
+ ); +} diff --git a/packages/web/src/components/HttpProcessCard.jsx b/packages/web/src/components/HttpProcessCard.jsx new file mode 100644 index 0000000..fa1e592 --- /dev/null +++ b/packages/web/src/components/HttpProcessCard.jsx @@ -0,0 +1,46 @@ +import { LuCircleCheck, LuCircleX } from "react-icons/lu"; + +export function HttpProcessCard({ + httpOnline, + onStart, + onStop, + startPending = false, + stopPending = false, +}) { + return ( +
+
+

HTTP process

+
+ {httpOnline ? ( + + ) : ( + + )} + {httpOnline ? "Running" : "Stopped"} +
+
+ {httpOnline ? ( + + ) : ( + + )} +
+
+
+ ); +} diff --git a/packages/web/src/components/Layout.jsx b/packages/web/src/components/Layout.jsx index 737fded..45fd436 100644 --- a/packages/web/src/components/Layout.jsx +++ b/packages/web/src/components/Layout.jsx @@ -1,6 +1,8 @@ +import { Fragment } from "react"; import { Link, NavLink, useNavigate } from "react-router-dom"; import { LuActivity, + LuArchive, LuCode, LuDatabase, LuFileText, @@ -20,17 +22,50 @@ import { useLogout, useOpsStatus } from "../api/hooks.js"; import { brandMark } from "../theme/brand.js"; import { useTheme } from "../theme.jsx"; -const links = [ - { to: "/", label: "Home", icon: LuHouse, end: true }, - { to: "/scripts", label: "Scripts", icon: LuCode }, - { 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 }, +const navSections = [ + { + items: [{ to: "/", label: "Home", icon: LuHouse, end: true }], + }, + { + title: "Automate", + items: [ + { to: "/workflows", label: "Workflows", icon: LuGitBranch }, + { to: "/scripts", label: "Scripts", icon: LuCode }, + { to: "/events", label: "Events", icon: LuActivity }, + ], + }, + { + title: "Platform", + items: [ + { to: "/variables", label: "Variables", icon: LuTags }, + { to: "/kv", label: "KV", icon: LuDatabase }, + { to: "/auth", label: "Auth", icon: LuShield }, + { to: "/responses", label: "Responses", icon: LuFileText }, + { to: "/secrets", label: "Secrets", icon: LuKey, admin: true }, + ], + }, + { + title: "Admin", + admin: true, + items: [ + { to: "/manage", label: "Manage", icon: LuServer }, + { to: "/backup", label: "Backup", icon: LuArchive }, + { to: "/users", label: "Users", icon: LuUsers }, + ], + }, ]; +function visibleSections(role) { + const isAdmin = role === "admin"; + return navSections + .filter((s) => !s.admin || isAdmin) + .map((s) => ({ + ...s, + items: s.items.filter((item) => !item.admin || isAdmin), + })) + .filter((s) => s.items.length > 0); +} + export function Layout({ user, children }) { const { theme, toggle } = useTheme(); const logout = useLogout(); @@ -38,17 +73,6 @@ export function Layout({ user, children }) { const ops = useOpsStatus(user.role === "admin"); const restartNeeded = Boolean(ops.data?.desired?.restartNeeded); - const navItems = [ - ...links, - ...(user.role === "admin" - ? [ - { to: "/secrets", label: "Secrets", icon: LuKey }, - { to: "/users", label: "Users", icon: LuUsers }, - { to: "/ops", label: "Ops", icon: LuServer }, - ] - : []), - ]; - function closeDrawer() { const el = document.getElementById("nav-drawer"); if (el) el.checked = false; @@ -101,8 +125,8 @@ export function Layout({ user, children }) { ? ` (${ops.data.desired.restartReason})` : ""} .{" "} - - Drain restart from Ops + + Drain restart from Manage {" "} to apply on HTTP + workers. @@ -120,13 +144,20 @@ export function Layout({ user, children }) { JerapahFlow
- {navItems.map(({ to, label, icon: Icon, end }) => ( -
  • - - - {label} - -
  • + {visibleSections(user.role).map((section) => ( + + {section.title ? ( +
  • {section.title}
  • + ) : null} + {section.items.map(({ to, label, icon: Icon, end }) => ( +
  • + + + {label} + +
  • + ))} +
    ))}
    diff --git a/packages/web/src/components/ProcessResourcesCard.jsx b/packages/web/src/components/ProcessResourcesCard.jsx new file mode 100644 index 0000000..6296c37 --- /dev/null +++ b/packages/web/src/components/ProcessResourcesCard.jsx @@ -0,0 +1,17 @@ +import { formatBytes, formatCpu } from "../lib/format.js"; + +export function ProcessResourcesCard({ children }) { + const totals = children?.totals ?? { memory: 0, cpu: 0 }; + + return ( +
    +
    +

    Resources

    +
    +

    {formatBytes(totals.memory)} memory total

    +

    {formatCpu(totals.cpu)} CPU total

    +
    +
    +
    + ); +} diff --git a/packages/web/src/components/QueueCard.jsx b/packages/web/src/components/QueueCard.jsx new file mode 100644 index 0000000..32244a5 --- /dev/null +++ b/packages/web/src/components/QueueCard.jsx @@ -0,0 +1,61 @@ +import { LuCircleCheck, LuCirclePause } from "react-icons/lu"; + +export function QueueCard({ + paused = false, + workersOnline = 0, + queuedJobs = 0, + activeJobs = 0, + onPause, + onResume, + pausePending = false, + resumePending = false, +}) { + return ( +
    +
    +

    Queue

    +
    + {paused ? ( + + ) : ( + + )} + + {paused + ? "Paused" + : `Running (${workersOnline} worker${workersOnline === 1 ? "" : "s"})`} + +
    +
    +

    + {activeJobs} running job{activeJobs === 1 ? "" : "s"} +

    +

    + {queuedJobs} queued job{queuedJobs === 1 ? "" : "s"} +

    +
    +
    + {paused ? ( + + ) : ( + + )} +
    +
    +
    + ); +} diff --git a/packages/web/src/components/SecretEditorModal.jsx b/packages/web/src/components/SecretEditorModal.jsx new file mode 100644 index 0000000..3d1644e --- /dev/null +++ b/packages/web/src/components/SecretEditorModal.jsx @@ -0,0 +1,127 @@ +import { useState } from "react"; +import { errorMessage } from "../api/client.js"; +import { useOwners, useUpsertSecret } from "../api/hooks.js"; + +/** + * Add / replace an encrypted secret. + * Reusable: mount when open; parent supplies mode + initial fields. + * + * @param {"add" | "replace"} mode + * @param {{ owner: string, name?: string }} initial + * @param {() => void} onClose + * @param {(saved: unknown) => void} [onSaved] + * @param {boolean} [lockOwner] + */ +export function SecretEditorModal({ + mode, + initial, + onClose, + onSaved, + lockOwner = false, +}) { + const { data: owners = [] } = useOwners(); + const upsert = useUpsertSecret(); + const [form, setForm] = useState(() => ({ + owner: initial.owner || owners[0] || "default", + name: initial.name || "", + value: "", + })); + + function onSubmit(e) { + e.preventDefault(); + upsert.mutate( + { owner: form.owner, name: form.name, value: form.value }, + { + onSuccess: (data) => { + onSaved?.(data?.secret ?? data); + onClose(); + }, + }, + ); + } + + const title = mode === "add" ? "New secret" : `Replace ${form.owner}/${form.name}`; + + return ( + +
    +

    {title}

    + + {mode === "add" ? ( + <> + + + + ) : null} + +

    + Values are encrypted at rest and never shown again after save. Values shorter than 8 + characters are not redacted from logs. +

    + {upsert.isError ? ( +

    {errorMessage(upsert.error)}

    + ) : null} +
    + + +
    + +
    +
    + +
    +
    + ); +} diff --git a/packages/web/src/components/UserEditorModal.jsx b/packages/web/src/components/UserEditorModal.jsx new file mode 100644 index 0000000..6a1d828 --- /dev/null +++ b/packages/web/src/components/UserEditorModal.jsx @@ -0,0 +1,120 @@ +import { useState } from "react"; +import { errorMessage } from "../api/client.js"; +import { useCreateUser, useUpdateUser } from "../api/hooks.js"; + +/** + * Add / edit a local user. + * Reusable: mount when open; pass `user` for edit. + * + * @param {"add" | "edit"} mode + * @param {{ id: string, username: string, role: string }} [user] + * @param {() => void} onClose + * @param {(saved: unknown) => void} [onSaved] + */ +export function UserEditorModal({ mode, user, onClose, onSaved }) { + const create = useCreateUser(); + const update = useUpdateUser(); + const [form, setForm] = useState(() => ({ + username: user?.username || "", + password: "", + role: user?.role || "operator", + id: user?.id || null, + })); + + const mutation = mode === "add" ? create : update; + + function onSubmit(e) { + e.preventDefault(); + if (mode === "add") { + create.mutate( + { + username: form.username, + password: form.password, + role: form.role, + }, + { + onSuccess: (data) => { + onSaved?.(data?.user ?? data); + onClose(); + }, + }, + ); + return; + } + const body = { id: form.id, role: form.role }; + if (form.password) body.password = form.password; + update.mutate(body, { + onSuccess: (data) => { + onSaved?.(data?.user ?? data); + onClose(); + }, + }); + } + + const title = mode === "add" ? "New user" : form.username; + + return ( + +
    +

    {title}

    +
    + {mode === "add" ? ( + + ) : null} + + + {mutation.isError ? ( +

    {errorMessage(mutation.error)}

    + ) : null} +
    + + +
    +
    +
    +
    + +
    +
    + ); +} diff --git a/packages/web/src/components/VariableEditorModal.jsx b/packages/web/src/components/VariableEditorModal.jsx new file mode 100644 index 0000000..4f54ea3 --- /dev/null +++ b/packages/web/src/components/VariableEditorModal.jsx @@ -0,0 +1,191 @@ +import { useState } from "react"; +import { errorMessage } from "../api/client.js"; +import { useOwners, useUpsertVariable } from "../api/hooks.js"; + +const TYPES = ["string", "number", "boolean"]; + +function defaultValue(type) { + if (type === "boolean") return false; + if (type === "number") return ""; + return ""; +} + +/** + * Add / edit a plaintext workflow variable. + * Reusable: mount when open; parent supplies mode + initial fields. + * + * @param {"add" | "edit"} mode + * @param {{ owner: string, name?: string, type?: string, value?: string | number | boolean }} initial + * @param {() => void} onClose + * @param {(saved: unknown) => void} [onSaved] + * @param {boolean} [lockOwner] When true, owner cannot be changed (add mode). + */ +export function VariableEditorModal({ + mode, + initial, + onClose, + onSaved, + lockOwner = false, +}) { + const { data: owners = [] } = useOwners(); + const upsert = useUpsertVariable(); + const [form, setForm] = useState(() => ({ + owner: initial.owner || owners[0] || "default", + name: initial.name || "", + type: initial.type || "string", + value: + initial.type === "number" + ? String(initial.value ?? "") + : initial.value !== undefined + ? initial.value + : defaultValue(initial.type || "string"), + })); + const [formError, setFormError] = useState(null); + + function onTypeChange(type) { + setForm({ ...form, type, value: defaultValue(type) }); + } + + function onSubmit(e) { + e.preventDefault(); + let value = form.value; + if (form.type === "number") { + value = Number(form.value); + if (!Number.isFinite(value)) { + setFormError("value must be a finite number"); + return; + } + } + if (form.type === "boolean") { + value = form.value === true; + } + setFormError(null); + upsert.mutate( + { owner: form.owner, name: form.name, type: form.type, value }, + { + onSuccess: (data) => { + onSaved?.(data?.variable ?? data); + onClose(); + }, + }, + ); + } + + const title = mode === "add" ? "New variable" : `Edit ${form.owner}/${form.name}`; + + return ( + +
    +

    {title}

    +
    + {mode === "add" ? ( + <> + + + + ) : null} + +