feat(workflows): migrate legacy workflows and enhance example handling
- Migrated legacy workflow files from `packages/server/workflows/` to `packages/server/data/workflows/` when the new store is empty. - Introduced new API endpoints to list and retrieve example workflows. - Added a dialog component for selecting example workflows when creating new workflows. - Updated workflow paths and configurations to reflect the new structure. - Enhanced the README and AGENTS.md documentation to clarify workflow management and usage. Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
@@ -13,7 +13,7 @@
|
||||
"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:workflow-history": "node test/workflow-history-smoke.js",
|
||||
"test:workflow-history": "JFLOW_WORKFLOWS_DIR=./data/workflow-history-smoke JFLOW_DB_PATH=./data/workflow-history-smoke.db node test/workflow-history-smoke.js",
|
||||
"test:profiles": "node test/profiles-smoke.js",
|
||||
"test:set-dry-run": "node test/set-dry-run-smoke.js"
|
||||
},
|
||||
|
||||
@@ -3,8 +3,12 @@ import { fileURLToPath } from "url";
|
||||
|
||||
export const SERVER_ROOT = path.dirname(fileURLToPath(import.meta.url));
|
||||
export const SCRIPTS_DIR = path.join(SERVER_ROOT, "scripts");
|
||||
export const WORKFLOWS_DIR = path.join(SERVER_ROOT, "workflows");
|
||||
export const DATA_DIR = path.join(SERVER_ROOT, "data");
|
||||
/** Live instance workflows (not shipped in git). Override for tests. */
|
||||
export const WORKFLOWS_DIR =
|
||||
process.env.JFLOW_WORKFLOWS_DIR ?? path.join(DATA_DIR, "workflows");
|
||||
/** Pre-0.1 layout; used only for one-shot migrate into WORKFLOWS_DIR. */
|
||||
export const LEGACY_WORKFLOWS_DIR = path.join(SERVER_ROOT, "workflows");
|
||||
/** User plugins (repo-root /plugins, outside the pnpm workspace). */
|
||||
export const PLUGINS_DIR =
|
||||
process.env.JFLOW_PLUGINS_DIR ??
|
||||
@@ -14,5 +18,10 @@ export const EXAMPLE_PLUGINS_DIR = path.resolve(
|
||||
SERVER_ROOT,
|
||||
"../../examples/plugins",
|
||||
);
|
||||
/** Example workflow YAML presets (not loaded by the runner). */
|
||||
export const EXAMPLE_WORKFLOWS_DIR = path.resolve(
|
||||
SERVER_ROOT,
|
||||
"../../examples/workflows",
|
||||
);
|
||||
export const LOGS_DIR = path.join(SERVER_ROOT, "logs");
|
||||
export const WEB_DIST = path.resolve(SERVER_ROOT, "../web/dist");
|
||||
|
||||
@@ -45,6 +45,11 @@ import {
|
||||
createWorkflowBackupBuffer,
|
||||
restoreWorkflowBackup,
|
||||
} from "../../workflow-backup.js";
|
||||
import {
|
||||
listExampleWorkflows,
|
||||
readExampleWorkflow,
|
||||
assertExampleWorkflowId,
|
||||
} from "../../workflow-examples.js";
|
||||
|
||||
/**
|
||||
* Reload this process and notify other HTTP/worker processes via Redis.
|
||||
@@ -246,6 +251,22 @@ export default function workflowsPluginFactory(registry) {
|
||||
return { owners: fsStore.listOwners() };
|
||||
});
|
||||
|
||||
fastify.get("/workflow-examples", async () => {
|
||||
return { examples: listExampleWorkflows() };
|
||||
});
|
||||
|
||||
fastify.get("/workflow-examples/:id", async (req, reply) => {
|
||||
const { id } = /** @type {{ id: string }} */ (req.params);
|
||||
if (!assertExampleWorkflowId(id)) {
|
||||
return reply.code(400).send({ error: "invalid example id" });
|
||||
}
|
||||
const example = readExampleWorkflow(id);
|
||||
if (!example) {
|
||||
return reply.code(404).send({ error: "example not found" });
|
||||
}
|
||||
return example;
|
||||
});
|
||||
|
||||
fastify.get("/workflows/trash", async () => {
|
||||
return { items: await listTrash() };
|
||||
});
|
||||
|
||||
@@ -31,6 +31,7 @@ import {
|
||||
getRedisUrlForLog,
|
||||
} from "./workflow-queue.js";
|
||||
import { purgeExpiredTrash } from "./workflow-trash.js";
|
||||
import { migrateLegacyWorkflowsIfNeeded } from "./workflow-migrate.js";
|
||||
import {
|
||||
getConfigGeneration,
|
||||
startHeartbeatLoop,
|
||||
@@ -128,6 +129,12 @@ export async function startApp(opts = {}) {
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
migrateLegacyWorkflowsIfNeeded();
|
||||
} catch (err) {
|
||||
log.warn({ err }, "legacy workflow migrate failed");
|
||||
}
|
||||
|
||||
const registry = createRegistry(server, {
|
||||
queue: workflowQueue,
|
||||
// Cron + HTTP triggers enqueue jobs; only the API process may own them.
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
import fs from "fs";
|
||||
import path from "path";
|
||||
import yaml from "yaml";
|
||||
import { EXAMPLE_WORKFLOWS_DIR } from "./paths.js";
|
||||
|
||||
/**
|
||||
* @param {string} id
|
||||
* @returns {string | null} safe basename without extension, or null if invalid
|
||||
*/
|
||||
export function assertExampleWorkflowId(id) {
|
||||
if (typeof id !== "string" || !/^[a-z0-9]+(?:-[a-z0-9]+)*$/i.test(id)) {
|
||||
return null;
|
||||
}
|
||||
return id;
|
||||
}
|
||||
|
||||
/**
|
||||
* Absolute path to an example YAML, or null if missing/unsafe.
|
||||
* @param {string} id
|
||||
*/
|
||||
export function exampleWorkflowPath(id) {
|
||||
const safe = assertExampleWorkflowId(id);
|
||||
if (!safe) return null;
|
||||
const filePath = path.join(EXAMPLE_WORKFLOWS_DIR, `${safe}.yaml`);
|
||||
const resolved = path.resolve(filePath);
|
||||
if (
|
||||
resolved !== EXAMPLE_WORKFLOWS_DIR &&
|
||||
!resolved.startsWith(EXAMPLE_WORKFLOWS_DIR + path.sep)
|
||||
) {
|
||||
return null;
|
||||
}
|
||||
if (!fs.existsSync(resolved) || !fs.statSync(resolved).isFile()) {
|
||||
return null;
|
||||
}
|
||||
return resolved;
|
||||
}
|
||||
|
||||
/**
|
||||
* @returns {{ id: string, name: string, description: string }[]}
|
||||
*/
|
||||
export function listExampleWorkflows() {
|
||||
if (!fs.existsSync(EXAMPLE_WORKFLOWS_DIR)) return [];
|
||||
return fs
|
||||
.readdirSync(EXAMPLE_WORKFLOWS_DIR)
|
||||
.filter((f) => f.endsWith(".yaml") || f.endsWith(".yml"))
|
||||
.sort()
|
||||
.map((f) => {
|
||||
const id = f.replace(/\.ya?ml$/i, "");
|
||||
const filePath = path.join(EXAMPLE_WORKFLOWS_DIR, f);
|
||||
let name = id;
|
||||
let description = "";
|
||||
try {
|
||||
const parsed = yaml.parse(fs.readFileSync(filePath, "utf8")) ?? {};
|
||||
if (typeof parsed.name === "string" && parsed.name.trim()) {
|
||||
name = parsed.name.trim();
|
||||
}
|
||||
if (parsed.description != null) {
|
||||
description = String(parsed.description).trim();
|
||||
}
|
||||
} catch {
|
||||
// keep id as name
|
||||
}
|
||||
return { id, name, description };
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} id
|
||||
* @returns {{ id: string, content: string } | null}
|
||||
*/
|
||||
export function readExampleWorkflow(id) {
|
||||
const filePath = exampleWorkflowPath(id);
|
||||
if (!filePath) return null;
|
||||
const safe = assertExampleWorkflowId(id);
|
||||
return {
|
||||
id: /** @type {string} */ (safe),
|
||||
content: fs.readFileSync(filePath, "utf8"),
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
import fs from "fs";
|
||||
import path from "path";
|
||||
import { LEGACY_WORKFLOWS_DIR, WORKFLOWS_DIR } from "./paths.js";
|
||||
import { log } from "./logger.js";
|
||||
|
||||
/**
|
||||
* Recursively copy a directory.
|
||||
* @param {string} src
|
||||
* @param {string} dest
|
||||
*/
|
||||
function copyDir(src, dest) {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True when WORKFLOWS_DIR has no owner subdirectories.
|
||||
* @param {string} dir
|
||||
*/
|
||||
function isEmptyWorkflowsDir(dir) {
|
||||
if (!fs.existsSync(dir)) return true;
|
||||
const entries = fs.readdirSync(dir, { withFileTypes: true });
|
||||
return !entries.some((e) => e.isDirectory());
|
||||
}
|
||||
|
||||
/**
|
||||
* One-shot: copy packages/server/workflows → data/workflows when the new
|
||||
* store is empty and the legacy tree still exists.
|
||||
* Does not copy from examples/workflows.
|
||||
*/
|
||||
export function migrateLegacyWorkflowsIfNeeded() {
|
||||
if (!isEmptyWorkflowsDir(WORKFLOWS_DIR)) return false;
|
||||
if (!fs.existsSync(LEGACY_WORKFLOWS_DIR)) return false;
|
||||
const legacyEntries = fs.readdirSync(LEGACY_WORKFLOWS_DIR, {
|
||||
withFileTypes: true,
|
||||
});
|
||||
if (!legacyEntries.some((e) => e.isDirectory())) return false;
|
||||
|
||||
fs.mkdirSync(WORKFLOWS_DIR, { recursive: true });
|
||||
copyDir(LEGACY_WORKFLOWS_DIR, WORKFLOWS_DIR);
|
||||
log.info(
|
||||
{ from: LEGACY_WORKFLOWS_DIR, to: WORKFLOWS_DIR },
|
||||
"migrated legacy workflows into instance store",
|
||||
);
|
||||
return true;
|
||||
}
|
||||
@@ -3,3 +3,4 @@ scripts:
|
||||
- detect-example-changes.yaml
|
||||
- time-to-ntfy-example.yaml
|
||||
- comic-monkeyuser-to-ntfy.yaml
|
||||
- afb272d4-b217-49ac-8c8b-755a6d8dac4a.yaml
|
||||
|
||||
Reference in New Issue
Block a user