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/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/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/src/api/workflows.js b/packages/server/src/api/workflows.js index 2904c38..d95f151 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -14,9 +14,33 @@ 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, +} 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. @@ -62,6 +86,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,6 +184,74 @@ 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(); @@ -93,6 +271,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 +281,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,14 +294,13 @@ 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, @@ -137,6 +312,94 @@ 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 }); + } + 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 +430,7 @@ export default function workflowsPluginFactory(registry) { return { owner, file, + workflow_id: workflowIdFromFile(file), key, content, parsed, @@ -188,48 +452,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 +561,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 +606,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 +648,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 +666,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 +699,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 +782,5 @@ export default function workflowsPluginFactory(registry) { }); }; } + +export { newWorkflowFilename }; diff --git a/packages/server/start-app.js b/packages/server/start-app.js index 6f2ddb7..e38ed79 100644 --- a/packages/server/start-app.js +++ b/packages/server/start-app.js @@ -28,11 +28,7 @@ import { createWorkflowWorker, getRedisUrlForLog, } from "./workflow-queue.js"; -import { - getConfigGeneration, - startHeartbeatLoop, - subscribeReload, -} from "./control-bus.js"; +import { purgeExpiredTrash } from "./workflow-trash.js"; /** * @param {{ @@ -48,6 +44,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(); diff --git a/packages/server/test/workflow-history-smoke.js b/packages/server/test/workflow-history-smoke.js new file mode 100644 index 0000000..00978b2 --- /dev/null +++ b/packages/server/test/workflow-history-smoke.js @@ -0,0 +1,121 @@ +/** + * 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, +} 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 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..4953cbd 100644 --- a/packages/server/workflow-duplicate.js +++ b/packages/server/workflow-duplicate.js @@ -1,4 +1,10 @@ import yaml from "yaml"; +import { + newWorkflowFilename, + workflowFileStem, +} from "./workflow-normalize.js"; + +export { workflowFileStem }; export function ensureWorkflowFilename(file) { const trimmed = String(file ?? "").trim(); @@ -6,13 +12,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 +34,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 +69,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..9063aa1 --- /dev/null +++ b/packages/server/workflow-history.js @@ -0,0 +1,128 @@ +import { randomUUID } from "node:crypto"; +import { db } from "./db.js"; +import { workflowContentSha } 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), + }; +} + +/** + * 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-normalize.js b/packages/server/workflow-normalize.js new file mode 100644 index 0000000..3486e85 --- /dev/null +++ b/packages/server/workflow-normalize.js @@ -0,0 +1,62 @@ +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; +} + +/** + * SHA256 of normalized workflow content (YAML → object → canonical JSON). + * @param {string} content + */ +export function workflowContentSha(content) { + const parsed = 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, ""); +} + +/** + * 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/web/src/App.jsx b/packages/web/src/App.jsx index 579c487..ccc3713 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"; @@ -59,6 +60,7 @@ export function App() { } /> } /> } /> + } /> } /> } /> } /> diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index d7fed6d..302a704 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"] }); }, }); @@ -572,3 +577,124 @@ 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 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/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 - + +
+ + + + + + ); +} + +/** + * Extract validation warnings from a failed save mutation error. + * @param {unknown} error + */ +export function saveWarningsFromError(error) { + const warnings = error?.response?.data?.warnings; + return Array.isArray(warnings) ? warnings : null; +} + +export function isSaveWarningsError(error) { + return error?.response?.status === 422 && saveWarningsFromError(error); +} + +export function saveErrorMessage(error) { + if (isSaveWarningsError(error)) { + return "Workflow has validation warnings"; + } + return errorMessage(error); +} diff --git a/packages/web/src/components/workflow/WorkflowHistoryPanel.jsx b/packages/web/src/components/workflow/WorkflowHistoryPanel.jsx new file mode 100644 index 0000000..fcdae63 --- /dev/null +++ b/packages/web/src/components/workflow/WorkflowHistoryPanel.jsx @@ -0,0 +1,113 @@ +import { useState } from "react"; +import { LuHistory, LuRotateCcw } from "react-icons/lu"; +import { errorMessage } from "../../api/client.js"; +import { + useRevertWorkflowRevision, + useWorkflowRevisions, +} from "../../api/hooks.js"; +import { formatTime } from "../../lib/format.jsx"; +import { useNotifications } from "../../notifications.jsx"; +import { + SaveWorkflowWarningsDialog, + isSaveWarningsError, + saveWarningsFromError, +} from "./SaveWorkflowWarningsDialog.jsx"; + +function reasonLabel(reason, meta) { + if (reason === "duplicated" && meta?.from) return `duplicated from ${meta.from}`; + if (reason === "revert" && meta?.fromRevision != null) { + return `reverted from #${meta.fromRevision}`; + } + if (reason === "restored-from-trash") return "restored from trash"; + return reason ?? "save"; +} + +export function WorkflowHistoryPanel({ owner, file, onReverted }) { + const { notify } = useNotifications(); + const revisions = useWorkflowRevisions(owner, file); + const revert = useRevertWorkflowRevision(); + const [pendingRevision, setPendingRevision] = useState(null); + const [warnings, setWarnings] = useState(null); + + function doRevert(revision, saveAnyway = false) { + setPendingRevision(revision); + revert.mutate( + { owner, file, revision, saveAnyway }, + { + onSuccess: (data) => { + setWarnings(null); + setPendingRevision(null); + notify.success(`Restored revision #${revision}`); + onReverted?.(data); + }, + onError: (err) => { + setPendingRevision(null); + if (isSaveWarningsError(err)) { + setWarnings({ revision, items: saveWarningsFromError(err) }); + return; + } + notify.error(errorMessage(err)); + }, + }, + ); + } + + const items = revisions.data?.revisions ?? []; + + return ( +
+
+ + History + ({items.length} / 50) +
+
+ {revisions.isLoading ? ( +
+ +
+ ) : !items.length ? ( +

No revisions yet.

+ ) : ( +
    + {items.map((rev) => ( +
  • +
    +
    #{rev.revision}
    +
    {formatTime(rev.created_at)}
    +
    + {reasonLabel(rev.reason, rev.meta)} +
    +
    + +
  • + ))} +
+ )} +
+ {warnings ? ( + setWarnings(null)} + onSaveAnyway={() => doRevert(warnings.revision, true)} + /> + ) : null} +
+ ); +} diff --git a/packages/web/src/pages/WorkflowEditPage.jsx b/packages/web/src/pages/WorkflowEditPage.jsx index 3f8ae0d..74d4061 100644 --- a/packages/web/src/pages/WorkflowEditPage.jsx +++ b/packages/web/src/pages/WorkflowEditPage.jsx @@ -3,6 +3,7 @@ import { Link, useNavigate, useParams } from "react-router-dom"; import { LuArrowLeft, LuCopy, LuPause, LuPlay, LuSave } from "react-icons/lu"; import { errorMessage } from "../api/client.js"; import { + useCreateWorkflow, useOwners, useSaveWorkflow, useSetWorkflowEnabled, @@ -11,6 +12,13 @@ import { import { DuplicateWorkflowDialog } from "../components/DuplicateWorkflowDialog.jsx"; import { WorkflowFileIcon } from "../components/WorkflowFileIcon.jsx"; import { WorkflowVisualEditor } from "../components/workflow/WorkflowVisualEditor.jsx"; +import { WorkflowHistoryPanel } from "../components/workflow/WorkflowHistoryPanel.jsx"; +import { + SaveWorkflowWarningsDialog, + isSaveWarningsError, + saveErrorMessage, + saveWarningsFromError, +} from "../components/workflow/SaveWorkflowWarningsDialog.jsx"; import { NEW_WORKFLOW_YAML, parseWorkflowYaml } from "../lib/workflow-doc.js"; import { useNotifications } from "../notifications.jsx"; @@ -87,48 +95,63 @@ function WorkflowEditorLayout({ export function WorkflowNewPage() { const navigate = useNavigate(); + const { notify } = useNotifications(); const { data: owners = [] } = useOwners(); const [owner, setOwner] = useState(""); - const [file, setFile] = useState(""); const [content, setContent] = useState(NEW_WORKFLOW_YAML); const [savedYaml] = useState(NEW_WORKFLOW_YAML); - const save = useSaveWorkflow(); + const [saveWarnings, setSaveWarnings] = useState(null); + const create = useCreateWorkflow(); useEffect(() => { if (!owner && owners[0]) setOwner(owners[0]); }, [owner, owners]); - function onSave() { - const yamlFile = file.endsWith(".yaml") || file.endsWith(".yml") ? file : `${file}.yaml`; - save.mutate( - { owner, file: yamlFile, content }, + function onSave(saveAnyway = false) { + create.mutate( + { owner, content, saveAnyway }, { - onSuccess: () => + onSuccess: (data) => { + setSaveWarnings(null); + notify.success(`Created ${data.file}`); navigate( - `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(yamlFile)}/edit`, - ), + `/workflows/${encodeURIComponent(data.owner)}/${encodeURIComponent(data.file)}/edit`, + ); + }, + onError: (err) => { + if (isSaveWarningsError(err)) { + setSaveWarnings(saveWarningsFromError(err)); + return; + } + notify.error(errorMessage(err)); + }, }, ); } return ( - -
- + <> + onSave(false)} + savePending={create.isPending} + saveDisabled={!owner} + saveError={create.isError && !saveWarnings ? errorMessage(create.error) : null} + > +
+

+ A UUID filename is assigned on save (for example{" "} + a1b2c3d4-….yaml). Edit the{" "} + name: field for the display name. +

+ setOwner(e.target.value)} required /> - setFile(e.target.value)} - required - /> - - } + } + /> +
+
+ {saveWarnings ? ( + setSaveWarnings(null)} + onSaveAnyway={() => onSave(true)} /> -
-
+ ) : null} + ); } @@ -165,6 +189,7 @@ export function WorkflowEditPage() { const [contentReady, setContentReady] = useState(false); const [testOpen, setTestOpen] = useState(false); const [duplicateOpen, setDuplicateOpen] = useState(false); + const [saveWarnings, setSaveWarnings] = useState(null); const routeKey = `${owner}/${file}`; const [activeKey, setActiveKey] = useState(routeKey); @@ -186,14 +211,21 @@ export function WorkflowEditPage() { } }, [existing.data, existing.isLoading, contentReady]); - function onSave() { + function onSave(saveAnyway = false) { save.mutate( - { owner, file, content }, + { owner, file, content, saveAnyway }, { onSuccess: () => { setSavedYaml(content); + setSaveWarnings(null); notify.success("Workflow saved"); }, + onError: (err) => { + if (isSaveWarningsError(err)) { + setSaveWarnings(saveWarningsFromError(err)); + return; + } + }, }, ); } @@ -221,16 +253,21 @@ export function WorkflowEditPage() { const parsedDoc = parseWorkflowYaml(content); const yamlOk = !parsedDoc.parseError; const workflowName = parsedDoc.doc?.name?.trim(); - const pageTitle = workflowName ? `${workflowName} (${file})` : file; + const pageTitle = workflowName ? `${workflowName}` : file; return ( <> + {pageTitle} + {file} + + } savePending={save.isPending} saveDisabled={!contentReady} - saveError={save.isError ? errorMessage(save.error) : null} - onSave={onSave} + saveError={save.isError && !saveWarnings ? saveErrorMessage(save.error) : null} + onSave={() => onSave(false)} onTest={() => setTestOpen(true)} onDuplicate={() => setDuplicateOpen(true)} onToggleEnabled={() => @@ -245,17 +282,34 @@ export function WorkflowEditPage() { enableDisabled={!contentReady || Boolean(existing.data?.parseError) || !yamlOk} enableError={setEnabled.isError ? errorMessage(setEnabled.error) : null} > -
- setTestOpen(false)} - /> +
+
+ setTestOpen(false)} + /> +
+
+ { + existing.refetch().then((result) => { + const next = result.data?.content; + if (next != null) { + setContent(next); + setSavedYaml(next); + } + }); + }} + /> +
{duplicateOpen ? ( @@ -271,6 +325,14 @@ export function WorkflowEditPage() { }} /> ) : null} + {saveWarnings ? ( + setSaveWarnings(null)} + onSaveAnyway={() => onSave(true)} + /> + ) : null} ); } diff --git a/packages/web/src/pages/WorkflowTrashPage.jsx b/packages/web/src/pages/WorkflowTrashPage.jsx new file mode 100644 index 0000000..f56d1a4 --- /dev/null +++ b/packages/web/src/pages/WorkflowTrashPage.jsx @@ -0,0 +1,162 @@ +import { useState } from "react"; +import { Link, useNavigate } from "react-router-dom"; +import { LuArrowLeft, LuRotateCcw, LuTrash2 } from "react-icons/lu"; +import { errorMessage } from "../api/client.js"; +import { + usePurgeWorkflowTrash, + useRestoreWorkflowTrash, + useWorkflowTrash, +} from "../api/hooks.js"; +import { formatTime } from "../lib/format.jsx"; +import { useNotifications } from "../notifications.jsx"; + +function formatAge(ms) { + if (ms == null || ms < 0) return "—"; + const mins = Math.floor(ms / 60_000); + if (mins < 60) return `${mins}m ago`; + const hours = Math.floor(mins / 60); + if (hours < 48) return `${hours}h ago`; + const days = Math.floor(hours / 24); + return `${days}d ago`; +} + +export function WorkflowTrashPage() { + const navigate = useNavigate(); + const { notify } = useNotifications(); + const { data: items = [], isLoading } = useWorkflowTrash(); + const restore = useRestoreWorkflowTrash(); + const purge = usePurgeWorkflowTrash(); + const [confirmPurge, setConfirmPurge] = useState(null); + + return ( +
+
+ + + +

Workflow trash

+
+ +

+ Deleted workflows are kept for 7 days. Restore to bring them back, or delete permanently + to remove the file and revision history. +

+ + {restore.isError ? ( +

{errorMessage(restore.error)}

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

{errorMessage(purge.error)}

+ ) : null} + + {isLoading ? ( + + ) : !items.length ? ( +

Trash is empty.

+ ) : ( +
+ + + + + + + + + + + + + {items.map((item) => ( + + + + + + + + + + ))} + +
NameFileOwnerDeletedAgePurge in +
{item.name ?? "—"}{item.file}{item.owner}{formatTime(item.deleted_at)}{formatAge(item.age_ms)} + {item.days_until_purge <= 0 ? ( + soon + ) : ( + `${item.days_until_purge}d` + )} + + + +
+
+ )} + + {confirmPurge ? ( + +
+

Delete permanently?

+

+ {confirmPurge.name ?? confirmPurge.file} ({confirmPurge.owner}/{confirmPurge.file}) + will be removed forever, including revision history. +

+
+ + +
+
+
+ +
+
+ ) : null} +
+ ); +} diff --git a/packages/web/src/pages/WorkflowsPage.jsx b/packages/web/src/pages/WorkflowsPage.jsx index 14252e8..47dedf5 100644 --- a/packages/web/src/pages/WorkflowsPage.jsx +++ b/packages/web/src/pages/WorkflowsPage.jsx @@ -4,6 +4,8 @@ import { LuActivity, LuCopy, LuPencil, LuPlay, LuPlus, LuRefreshCw, LuTrash2, Lu import { errorMessage } from "../api/client.js"; import { useDeleteWorkflow, + useDownloadWorkflowBackup, + useRestoreWorkflowBackup, useReregisterWorkflows, useRunWorkflow, useSetWorkflowEnabled, @@ -12,6 +14,7 @@ import { import { DuplicateWorkflowDialog } from "../components/DuplicateWorkflowDialog.jsx"; import { WorkflowFileIcon } from "../components/WorkflowFileIcon.jsx"; import { formatTime, WorkflowStatusBadge } from "../lib/format.jsx"; +import { useNotifications } from "../notifications.jsx"; export function WorkflowsPage() { const navigate = useNavigate(); @@ -25,6 +28,10 @@ export function WorkflowsPage() { const run = useRunWorkflow(); const setEnabled = useSetWorkflowEnabled(); const reregister = useReregisterWorkflows(); + const backup = useDownloadWorkflowBackup(); + const restoreBackup = useRestoreWorkflowBackup(); + const { notify } = useNotifications(); + const [restoreMode, setRestoreMode] = useState("merge"); if (editParam) { const slash = editParam.indexOf("/"); @@ -72,6 +79,10 @@ export function WorkflowsPage() {

Workflows

+ + + Trash + + + +
+ {isLoading ? ( ) : ( @@ -128,6 +189,7 @@ export function WorkflowsPage() { > {w.name} + {w.file} {!w.registered ? ( setConfirmDelete(w)} > @@ -231,7 +293,11 @@ export function WorkflowsPage() { {confirmDelete ? (
-

Delete {confirmDelete.key}?

+

Move {confirmDelete.key} to trash?

+

+ The workflow is removed from the list but kept in trash for 7 days. Revision history + is preserved. +

{del.isError ? (

{errorMessage(del.error)}

) : null} @@ -250,7 +316,7 @@ export function WorkflowsPage() { ) } > - Delete + Move to trash