feat(workflows): add revision history, trash, UUID naming, and backup

- Store up to 50 revisions per workflow in SQLite with normalized SHA dedup
- Soft-delete workflows to trash (7-day retention) with restore and permanent purge
- Assign UUID filenames for new and duplicated workflows; show name + file in UI
- Warn on invalid YAML, unknown scripts, and plaintext secrets with save-anyway option
- Add workflow backup zip (workflows + plugins) and merge/replace restore
- Trash page, history panel on editor, and smoke test

Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
Cursor Agent
2026-08-20 02:14:50 +00:00
co-authored by nsrb
parent 611511ae60
commit 7d51ff433f
20 changed files with 1943 additions and 206 deletions
+1
View File
@@ -0,0 +1 @@
dump.rdb
@@ -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");
}
+2 -1
View File
@@ -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",
+427 -86
View File
@@ -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<string, unknown> | 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<string, any>, loadErrors: Map<string, string>, 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 };
+9 -5
View File
@@ -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();
@@ -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();
+134
View File
@@ -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 });
}
}
+24 -7
View File
@@ -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}`;
}
+128
View File
@@ -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<string, unknown> | 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();
}
+62
View File
@@ -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));
}
+183
View File
@@ -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<number>}
*/
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;
}
@@ -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 };
}