refactor(config): remove legacy config reference handling and streamline YAML management
- Deleted legacy config reference handling functions and related migration scripts to simplify the codebase. - Updated documentation to reflect the removal of legacy YAML paths and configurations. - Adjusted existing YAML management processes to ensure compatibility with the new mustache-style syntax. - Removed tests related to legacy config references, focusing on current functionality and ensuring clarity in the testing suite.
This commit is contained in:
@@ -11,7 +11,7 @@ Live workflows are **not** product source. They live under `packages/server/data
|
||||
| Live / personal YAML | No | `packages/server/data/workflows/<owner>/` |
|
||||
| Example presets | Yes | `examples/workflows/*.yaml` (copy into editor only; runner does not load them) |
|
||||
|
||||
Do **not** add personal YAML under `packages/server/`, `examples/workflows/`, or the legacy `packages/server/workflows/` tree. Prefer owner `local`. Example presets must use **core** scripts only (no `plugin/…` that requires install).
|
||||
Do **not** add personal YAML under `packages/server/`, `examples/workflows/`, or `packages/server/data/workflows/`. Prefer owner `local`. Example presets must use **core** scripts only (no `plugin/…` that requires install).
|
||||
|
||||
## Where things live
|
||||
|
||||
|
||||
@@ -59,7 +59,6 @@ pnpm --dir packages/server reset-admin -- --username admin --password 'your-pass
|
||||
| **Example presets** | No | `examples/workflows/*.yaml` — offered when creating a new workflow |
|
||||
|
||||
- Live YAML is **instance data**, same as SQLite and secrets — not product source. New resources use owner `local` (owner remains in storage/URLs for a possible future multi-tenant mode; the UI hides it).
|
||||
- On first start, if the instance store is empty and a legacy `packages/server/workflows/` tree still exists, it is copied into `data/workflows/`.
|
||||
- New workflow editor starts empty; optional presets copy example YAML into the editor (nothing is saved until Save).
|
||||
- Override the live store in tests with `JFLOW_WORKFLOWS_DIR`.
|
||||
|
||||
@@ -109,7 +108,7 @@ topic: "{{ vars.ntfy_prefix }}/{{ data.channel }}"
|
||||
| `context` | Run clipboard (nested) |
|
||||
| `data` | This step’s input (nested; numeric segments index arrays) |
|
||||
|
||||
A string that is exactly one `{{ path }}` keeps the native type (object/array/number/boolean). Mixed strings concatenate as text. Bare name fields such as `passwordSecret: gmail_app_password` stay names for `$secrets.get` — do not wrap them in `{{ secrets.… }}`. Legacy `$VAR_` / `$SECRET_` / `$CONTEXT_` whole-value refs throw; use mustache instead. Script APIs `$vars.get` / `$secrets.get` are unchanged.
|
||||
A string that is exactly one `{{ path }}` keeps the native type (object/array/number/boolean). Mixed strings concatenate as text. Bare name fields such as `passwordSecret: gmail_app_password` stay names for `$secrets.get` — do not wrap them in `{{ secrets.… }}`. Script APIs `$vars.get` / `$secrets.get` are unchanged.
|
||||
|
||||
YAML **SET** evaluates JSONata against the full `ctx`; the result is `output` (the next step’s data). `jsonata.js` does the same.
|
||||
|
||||
|
||||
@@ -1,17 +0,0 @@
|
||||
/**
|
||||
* Rewrite legacy `$VAR_` / `$SECRET_` / `$CONTEXT_` placeholders to mustache.
|
||||
* Safe for passwordSecret-style fields (those store bare names, not $SECRET_ prefixes).
|
||||
*
|
||||
* @param {string} text
|
||||
* @returns {{ text: string, changed: boolean }}
|
||||
*/
|
||||
export function rewriteLegacyConfigRefsInText(text) {
|
||||
if (typeof text !== "string" || text.length === 0) {
|
||||
return { text: text ?? "", changed: false };
|
||||
}
|
||||
const next = text
|
||||
.replace(/\$VAR_([A-Za-z0-9._-]+)/g, "{{ vars.$1 }}")
|
||||
.replace(/\$SECRET_([A-Za-z0-9._-]+)/g, "{{ secrets.$1 }}")
|
||||
.replace(/\$CONTEXT_([A-Za-z0-9._-]+)/g, "{{ context.$1 }}");
|
||||
return { text: next, changed: next !== text };
|
||||
}
|
||||
@@ -8,7 +8,6 @@ const MUSTACHE_TOKEN_RE =
|
||||
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g;
|
||||
const WHOLE_MUSTACHE_RE =
|
||||
/^\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}$/;
|
||||
const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
|
||||
/**
|
||||
* @typedef {{
|
||||
@@ -19,30 +18,6 @@ const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
* }} ConfigRefCtx
|
||||
*/
|
||||
|
||||
/**
|
||||
* Detect leftover `$VAR_` / `$SECRET_` / `$CONTEXT_` whole-value refs.
|
||||
* @param {unknown} value
|
||||
* @returns {{ kind: "var" | "secret" | "context", name: string, raw: string } | null}
|
||||
*/
|
||||
export function parseLegacyConfigRef(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const match = LEGACY_PREFIX_RE.exec(value);
|
||||
if (!match) return null;
|
||||
const kind =
|
||||
match[1] === "VAR" ? "var" : match[1] === "SECRET" ? "secret" : "context";
|
||||
return { kind, name: match[2] ?? "", raw: value.trim() };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {"var" | "secret" | "context"} kind
|
||||
* @param {string} name
|
||||
*/
|
||||
function legacyRenameHint(kind, name) {
|
||||
if (kind === "var") return `use {{ vars.${name || "name"} }}`;
|
||||
if (kind === "secret") return `use {{ secrets.${name || "name"} }}`;
|
||||
return `use {{ context.${name || "name"} }}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Walk config (objects/arrays) and interpolate `{{ path }}` strings.
|
||||
* Does not walk trigger data.
|
||||
@@ -84,13 +59,6 @@ export async function resolveConfigRefs(value, ctx, seen = new WeakSet()) {
|
||||
* @returns {Promise<unknown>}
|
||||
*/
|
||||
async function resolveStringRef(value, ctx) {
|
||||
const legacy = parseLegacyConfigRef(value);
|
||||
if (legacy) {
|
||||
throw new Error(
|
||||
`config ref ${legacy.raw}: removed; ${legacyRenameHint(legacy.kind, legacy.name)}`,
|
||||
);
|
||||
}
|
||||
|
||||
const whole = WHOLE_MUSTACHE_RE.exec(value);
|
||||
if (whole && whole[0] === value) {
|
||||
return resolvePath(whole[1], ctx, { raw: value, allowObject: true });
|
||||
|
||||
@@ -1,24 +0,0 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.text("job_id");
|
||||
t.text("queued_at");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_job_id_idx");
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.dropColumn("job_id");
|
||||
t.dropColumn("queued_at");
|
||||
});
|
||||
}
|
||||
@@ -1,21 +0,0 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.integer("workflow_revision");
|
||||
});
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_workflow_revision_idx");
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.dropColumn("workflow_revision");
|
||||
});
|
||||
}
|
||||
@@ -1,16 +0,0 @@
|
||||
/**
|
||||
* No-op Knex marker: owner + config-ref migrate runs in owner-migrate.js at startup
|
||||
* (needs filesystem + DB together). Kept so deploy tooling sees a versioned step.
|
||||
*
|
||||
* @param {import("knex").Knex} _knex
|
||||
*/
|
||||
export async function up(_knex) {
|
||||
// Intentionally empty — see migrateDefaultOwnerIfNeeded().
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} _knex
|
||||
*/
|
||||
export async function down(_knex) {
|
||||
// Irreversible data migrate.
|
||||
}
|
||||
+5
-38
@@ -17,6 +17,9 @@ export async function up(knex) {
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
t.text("parent_run_id").references("id").inTable("workflow_runs");
|
||||
t.text("job_id");
|
||||
t.text("queued_at");
|
||||
t.integer("workflow_revision");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
@@ -28,45 +31,11 @@ export async function up(knex) {
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_status_started_at_idx ON workflow_runs (status, started_at DESC)",
|
||||
);
|
||||
|
||||
await knex.schema.createTable("step_runs", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.integer("step_index").notNullable();
|
||||
t.text("script").notNullable();
|
||||
t.text("config");
|
||||
t.text("status").notNullable();
|
||||
t.text("started_at").notNullable();
|
||||
t.text("finished_at");
|
||||
t.integer("duration_ms");
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
|
||||
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
|
||||
);
|
||||
|
||||
await knex.schema.createTable("logs", (t) => {
|
||||
t.increments("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.text("step_id");
|
||||
t.text("ts").notNullable();
|
||||
t.integer("level").notNullable();
|
||||
t.text("msg");
|
||||
t.text("payload");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
|
||||
"CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)",
|
||||
);
|
||||
}
|
||||
|
||||
@@ -74,7 +43,5 @@ export async function up(knex) {
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("logs");
|
||||
await knex.schema.dropTableIfExists("step_runs");
|
||||
await knex.schema.dropTableIfExists("workflow_runs");
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("step_runs", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.integer("step_index").notNullable();
|
||||
t.text("script").notNullable();
|
||||
t.text("config");
|
||||
t.text("status").notNullable();
|
||||
t.text("started_at").notNullable();
|
||||
t.text("finished_at");
|
||||
t.integer("duration_ms");
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("step_runs");
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("logs", (t) => {
|
||||
t.increments("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.text("step_id");
|
||||
t.text("ts").notNullable();
|
||||
t.integer("level").notNullable();
|
||||
t.text("msg");
|
||||
t.text("payload");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("logs");
|
||||
}
|
||||
+20
-20
@@ -29,35 +29,35 @@ const DEFAULT_EMAIL_TEMPLATE = `<!DOCTYPE html>
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("http_pages", (t) => {
|
||||
await knex.schema.createTable("http_pages", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("content").notNullable();
|
||||
t.text("mime").notNullable();
|
||||
t.integer("status").notNullable().defaultTo(200);
|
||||
t.text("kind").notNullable().defaultTo("response");
|
||||
t.integer("system").notNullable().defaultTo(0);
|
||||
t.text("created_at").notNullable();
|
||||
t.text("updated_at").notNullable();
|
||||
});
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const existing = await knex("http_pages").where({ name: "email-default" }).first();
|
||||
if (!existing) {
|
||||
await knex("http_pages").insert({
|
||||
id: "00000000-0000-4000-8000-000000000001",
|
||||
name: "email-default",
|
||||
content: DEFAULT_EMAIL_TEMPLATE,
|
||||
mime: "html",
|
||||
status: 200,
|
||||
kind: "template",
|
||||
system: 1,
|
||||
created_at: now,
|
||||
updated_at: now,
|
||||
});
|
||||
}
|
||||
await knex("http_pages").insert({
|
||||
id: "00000000-0000-4000-8000-000000000001",
|
||||
name: "email-default",
|
||||
content: DEFAULT_EMAIL_TEMPLATE,
|
||||
mime: "html",
|
||||
status: 200,
|
||||
kind: "template",
|
||||
system: 1,
|
||||
created_at: now,
|
||||
updated_at: now,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex("http_pages").where({ name: "email-default", system: 1 }).del();
|
||||
await knex.schema.alterTable("http_pages", (t) => {
|
||||
t.dropColumn("kind");
|
||||
t.dropColumn("system");
|
||||
});
|
||||
await knex.schema.dropTableIfExists("http_pages");
|
||||
}
|
||||
+2
-13
@@ -2,21 +2,11 @@
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("http_pages", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("content").notNullable();
|
||||
t.text("mime").notNullable(); // html | json
|
||||
t.integer("status").notNullable().defaultTo(200);
|
||||
t.text("created_at").notNullable();
|
||||
t.text("updated_at").notNullable();
|
||||
});
|
||||
|
||||
await knex.schema.createTable("http_auths", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("type").notNullable(); // bearer | basic | header
|
||||
t.text("config").notNullable(); // JSON
|
||||
t.text("type").notNullable();
|
||||
t.text("config").notNullable();
|
||||
t.integer("unauthorized_status").nullable();
|
||||
t.text("unauthorized_response").nullable();
|
||||
t.text("created_at").notNullable();
|
||||
@@ -29,5 +19,4 @@ export async function up(knex) {
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("http_auths");
|
||||
await knex.schema.dropTableIfExists("http_pages");
|
||||
}
|
||||
-18
@@ -21,29 +21,11 @@ export async function up(knex) {
|
||||
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");
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
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");
|
||||
}
|
||||
+3
-1
@@ -14,7 +14,9 @@ export async function up(knex) {
|
||||
t.unique(["owner", "name"]);
|
||||
});
|
||||
|
||||
await knex.schema.raw("CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)");
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1,268 +0,0 @@
|
||||
import fs from "fs";
|
||||
import path from "path";
|
||||
import { DEFAULT_OWNER } from "@jerapah-flow/shared";
|
||||
import { db } from "./db.js";
|
||||
import { log } from "./logger.js";
|
||||
import { EXAMPLE_WORKFLOWS_DIR, WORKFLOWS_DIR } from "./paths.js";
|
||||
import { TRASH_WORKFLOWS_DIR } from "./workflow-trash.js";
|
||||
import { rewriteLegacyConfigRefsInText } from "./config-ref-rewrite.js";
|
||||
|
||||
const LEGACY_OWNER = "default";
|
||||
|
||||
/**
|
||||
* @param {string} dir
|
||||
* @returns {string[]}
|
||||
*/
|
||||
function listFilesRecursive(dir) {
|
||||
if (!fs.existsSync(dir)) return [];
|
||||
/** @type {string[]} */
|
||||
const out = [];
|
||||
for (const entry of fs.readdirSync(dir, { withFileTypes: true })) {
|
||||
const full = path.join(dir, entry.name);
|
||||
if (entry.isDirectory()) out.push(...listFilesRecursive(full));
|
||||
else out.push(full);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Move files from legacy owner dir into DEFAULT_OWNER. Skip name collisions.
|
||||
* @param {string} rootDir
|
||||
* @returns {{ moved: number, skipped: number }}
|
||||
*/
|
||||
function mergeOwnerDir(rootDir) {
|
||||
const fromDir = path.join(rootDir, LEGACY_OWNER);
|
||||
const toDir = path.join(rootDir, DEFAULT_OWNER);
|
||||
if (!fs.existsSync(fromDir)) return { moved: 0, skipped: 0 };
|
||||
|
||||
fs.mkdirSync(toDir, { recursive: true });
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
|
||||
for (const name of fs.readdirSync(fromDir)) {
|
||||
const from = path.join(fromDir, name);
|
||||
const to = path.join(toDir, name);
|
||||
const st = fs.statSync(from);
|
||||
if (!st.isFile()) {
|
||||
skipped += 1;
|
||||
log.warn({ from }, "owner migrate: skip non-file under legacy owner dir");
|
||||
continue;
|
||||
}
|
||||
if (fs.existsSync(to)) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ from, to },
|
||||
"owner migrate: skip YAML collision (local already has file)",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
fs.renameSync(from, to);
|
||||
moved += 1;
|
||||
}
|
||||
|
||||
const remaining = fs.existsSync(fromDir) ? fs.readdirSync(fromDir) : [];
|
||||
if (remaining.length === 0 && fs.existsSync(fromDir)) {
|
||||
fs.rmdirSync(fromDir);
|
||||
}
|
||||
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} filePath
|
||||
* @returns {boolean}
|
||||
*/
|
||||
function rewriteFileInPlace(filePath) {
|
||||
const raw = fs.readFileSync(filePath, "utf8");
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(raw);
|
||||
if (!changed) return false;
|
||||
fs.writeFileSync(filePath, text, "utf8");
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} rootDir
|
||||
* @returns {number}
|
||||
*/
|
||||
function rewriteYamlTree(rootDir) {
|
||||
let n = 0;
|
||||
for (const file of listFilesRecursive(rootDir)) {
|
||||
if (!/\.ya?ml$/i.test(file)) continue;
|
||||
if (rewriteFileInPlace(file)) n += 1;
|
||||
}
|
||||
return n;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
* @param {string} table
|
||||
* @param {"name" | "file" | null} uniqueCol
|
||||
*/
|
||||
async function migrateOwnerColumn(knex, table, uniqueCol) {
|
||||
const legacyRows = await knex(table).where({ owner: LEGACY_OWNER }).select("*");
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
for (const row of legacyRows) {
|
||||
if (uniqueCol != null) {
|
||||
const conflict = await knex(table)
|
||||
.where({ owner: DEFAULT_OWNER, [uniqueCol]: row[uniqueCol] })
|
||||
.first();
|
||||
if (conflict) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ table, id: row.id, [uniqueCol]: row[uniqueCol] },
|
||||
"owner migrate: skip row collision",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
await knex(table).where({ id: row.id }).update({ owner: DEFAULT_OWNER });
|
||||
moved += 1;
|
||||
}
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateWorkflowRuns(knex) {
|
||||
const n = await knex("workflow_runs")
|
||||
.where({ owner: LEGACY_OWNER })
|
||||
.update({ owner: DEFAULT_OWNER });
|
||||
return { moved: Number(n) || 0, skipped: 0 };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateWorkflowRevisionsOwner(knex) {
|
||||
const n = await knex("workflow_revisions")
|
||||
.where({ owner: LEGACY_OWNER })
|
||||
.update({ owner: DEFAULT_OWNER });
|
||||
return { moved: Number(n) || 0, skipped: 0 };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateScriptStateNamespaces(knex) {
|
||||
const rows = await knex("script_state")
|
||||
.where("namespace", "like", `${LEGACY_OWNER}/%`)
|
||||
.select("namespace", "key");
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
for (const row of rows) {
|
||||
const nextNs = `${DEFAULT_OWNER}${row.namespace.slice(LEGACY_OWNER.length)}`;
|
||||
const conflict = await knex("script_state")
|
||||
.where({ namespace: nextNs, key: row.key })
|
||||
.first();
|
||||
if (conflict) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ namespace: row.namespace, key: row.key, nextNs },
|
||||
"owner migrate: skip script_state collision",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
await knex("script_state")
|
||||
.where({ namespace: row.namespace, key: row.key })
|
||||
.update({ namespace: nextNs });
|
||||
moved += 1;
|
||||
}
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function rewriteDbConfigStrings(knex) {
|
||||
let profiles = 0;
|
||||
let revisions = 0;
|
||||
|
||||
const profileRows = await knex("profiles").select("id", "config");
|
||||
for (const row of profileRows) {
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(String(row.config ?? ""));
|
||||
if (!changed) continue;
|
||||
await knex("profiles").where({ id: row.id }).update({ config: text });
|
||||
profiles += 1;
|
||||
}
|
||||
|
||||
const revisionRows = await knex("workflow_revisions").select("id", "content");
|
||||
for (const row of revisionRows) {
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(String(row.content ?? ""));
|
||||
if (!changed) continue;
|
||||
await knex("workflow_revisions").where({ id: row.id }).update({ content: text });
|
||||
revisions += 1;
|
||||
}
|
||||
|
||||
return { profiles, revisions };
|
||||
}
|
||||
|
||||
/**
|
||||
* One-shot: move owner `default` → `local`, rewrite prefix refs to mustache.
|
||||
* Idempotent when there is no remaining `default` data / prefix refs.
|
||||
*/
|
||||
export async function migrateDefaultOwnerIfNeeded() {
|
||||
const knex = db;
|
||||
|
||||
const yamlLive = mergeOwnerDir(WORKFLOWS_DIR);
|
||||
const yamlTrash = mergeOwnerDir(TRASH_WORKFLOWS_DIR);
|
||||
|
||||
const variables = await migrateOwnerColumn(knex, "variables", "name");
|
||||
const secrets = await migrateOwnerColumn(knex, "secrets", "name");
|
||||
const profiles = await migrateOwnerColumn(knex, "profiles", "name");
|
||||
const trash = await migrateOwnerColumn(knex, "workflow_trash", "file");
|
||||
const runs = await migrateWorkflowRuns(knex);
|
||||
const revisionsOwner = await migrateWorkflowRevisionsOwner(knex);
|
||||
const scriptState = await migrateScriptStateNamespaces(knex);
|
||||
|
||||
const yamlRewritten =
|
||||
rewriteYamlTree(path.join(WORKFLOWS_DIR, DEFAULT_OWNER)) +
|
||||
rewriteYamlTree(path.join(TRASH_WORKFLOWS_DIR, DEFAULT_OWNER)) +
|
||||
rewriteYamlTree(path.join(WORKFLOWS_DIR, LEGACY_OWNER)) +
|
||||
rewriteYamlTree(path.join(TRASH_WORKFLOWS_DIR, LEGACY_OWNER));
|
||||
|
||||
let examplesRewritten = 0;
|
||||
if (fs.existsSync(EXAMPLE_WORKFLOWS_DIR)) {
|
||||
examplesRewritten = rewriteYamlTree(EXAMPLE_WORKFLOWS_DIR);
|
||||
}
|
||||
|
||||
const dbStrings = await rewriteDbConfigStrings(knex);
|
||||
|
||||
const summary = {
|
||||
yamlLive,
|
||||
yamlTrash,
|
||||
variables,
|
||||
secrets,
|
||||
profiles,
|
||||
trash,
|
||||
runs,
|
||||
revisionsOwner,
|
||||
scriptState,
|
||||
yamlRewritten,
|
||||
examplesRewritten,
|
||||
dbStrings,
|
||||
};
|
||||
|
||||
const touched =
|
||||
yamlLive.moved +
|
||||
yamlTrash.moved +
|
||||
variables.moved +
|
||||
secrets.moved +
|
||||
profiles.moved +
|
||||
trash.moved +
|
||||
runs.moved +
|
||||
revisionsOwner.moved +
|
||||
scriptState.moved +
|
||||
yamlRewritten +
|
||||
examplesRewritten +
|
||||
dbStrings.profiles +
|
||||
dbStrings.revisions >
|
||||
0;
|
||||
|
||||
if (touched) {
|
||||
log.info(summary, "migrated owner default → local and rewrote config refs");
|
||||
}
|
||||
|
||||
return summary;
|
||||
}
|
||||
@@ -18,8 +18,7 @@
|
||||
"reset-admin": "node reset-admin.js",
|
||||
"test:profiles": "node test/profiles-smoke.js",
|
||||
"test:set-dry-run": "node test/set-dry-run-smoke.js",
|
||||
"test:config-refs": "node test/config-refs-smoke.js",
|
||||
"test:owner-migrate": "node test/owner-migrate-smoke.js"
|
||||
"test:config-refs": "node test/config-refs-smoke.js"
|
||||
},
|
||||
"dependencies": {
|
||||
"@jerapah-flow/shared": "workspace:*",
|
||||
|
||||
@@ -7,8 +7,6 @@ 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 ??
|
||||
|
||||
@@ -31,8 +31,6 @@ import {
|
||||
getRedisUrlForLog,
|
||||
} from "./workflow-queue.js";
|
||||
import { purgeExpiredTrash } from "./workflow-trash.js";
|
||||
import { migrateLegacyWorkflowsIfNeeded } from "./workflow-migrate.js";
|
||||
import { migrateDefaultOwnerIfNeeded } from "./owner-migrate.js";
|
||||
import {
|
||||
getConfigGeneration,
|
||||
startHeartbeatLoop,
|
||||
@@ -130,18 +128,6 @@ export async function startApp(opts = {}) {
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
migrateLegacyWorkflowsIfNeeded();
|
||||
} catch (err) {
|
||||
log.warn({ err }, "legacy workflow migrate failed");
|
||||
}
|
||||
|
||||
try {
|
||||
await migrateDefaultOwnerIfNeeded();
|
||||
} catch (err) {
|
||||
log.warn({ err }, "owner default→local migrate failed");
|
||||
}
|
||||
|
||||
const registry = createRegistry(server, {
|
||||
queue: workflowQueue,
|
||||
// Cron + HTTP triggers enqueue jobs; only the API process may own them.
|
||||
|
||||
@@ -2,8 +2,7 @@ import { migrate, db } from "../db.js";
|
||||
import { upsertSecret, deleteSecret } from "../secrets-store.js";
|
||||
import { deleteVariable, upsertVariable } from "../variables-store.js";
|
||||
import { Secret } from "../secret-value.js";
|
||||
import { parseLegacyConfigRef, resolveConfigRefs } from "../config-refs.js";
|
||||
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
|
||||
import { resolveConfigRefs } from "../config-refs.js";
|
||||
|
||||
await migrate();
|
||||
|
||||
@@ -27,50 +26,17 @@ async function assertRejects(fn, match) {
|
||||
const owner = "config_refs_smoke_owner";
|
||||
const ctx = { owner, workflowKey: `${owner}/config-refs-smoke.yaml`, context: {}, data: {} };
|
||||
|
||||
{
|
||||
const rewritten = rewriteLegacyConfigRefsInText(
|
||||
'url: $VAR_ntfy\ntoken: $SECRET_tok\nid: $CONTEXT_user',
|
||||
);
|
||||
assert(rewritten.changed, "rewrite detects legacy refs");
|
||||
assert(
|
||||
rewritten.text.includes("{{ vars.ntfy }}") &&
|
||||
rewritten.text.includes("{{ secrets.tok }}") &&
|
||||
rewritten.text.includes("{{ context.user }}"),
|
||||
"rewrite maps prefixes",
|
||||
);
|
||||
assert(!rewriteLegacyConfigRefsInText("passwordSecret: gmail_app").changed, "bare names untouched");
|
||||
}
|
||||
|
||||
assert(parseLegacyConfigRef("$VAR_ntfy_url")?.kind === "var", "legacy var parse");
|
||||
assert(parseLegacyConfigRef("$SECRET_x")?.kind === "secret", "legacy secret parse");
|
||||
assert(parseLegacyConfigRef("$CONTEXT_token")?.kind === "context", "legacy context parse");
|
||||
assert(parseLegacyConfigRef("{{ vars.x }}") == null, "mustache is not legacy");
|
||||
assert(parseLegacyConfigRef("Bearer $SECRET_x") == null, "mid-string not whole-value legacy");
|
||||
|
||||
{
|
||||
const literal = await resolveConfigRefs("password123", ctx);
|
||||
assert(literal === "password123", "literal passthrough");
|
||||
const unknown = await resolveConfigRefs("$FOO_bar", ctx);
|
||||
assert(unknown === "$FOO_bar", "$FOO_bar stays literal");
|
||||
const embedded = await resolveConfigRefs("Bearer $SECRET_x", ctx);
|
||||
assert(embedded === "Bearer $SECRET_x", "mid-string legacy stays literal");
|
||||
assert(embedded === "Bearer $SECRET_x", "mid-string stays literal");
|
||||
const number = await resolveConfigRefs(42, ctx);
|
||||
assert(number === 42, "number passthrough");
|
||||
}
|
||||
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$VAR_ntfy_channel", ctx),
|
||||
"removed",
|
||||
);
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$SECRET_tok", ctx),
|
||||
"use {{ secrets.tok }}",
|
||||
);
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$CONTEXT_token", ctx),
|
||||
"use {{ context.token }}",
|
||||
);
|
||||
|
||||
const created = [];
|
||||
try {
|
||||
const secret = await upsertSecret({
|
||||
|
||||
@@ -1,30 +0,0 @@
|
||||
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
|
||||
import { migrateDefaultOwnerIfNeeded } from "../owner-migrate.js";
|
||||
import { migrate, db } from "../db.js";
|
||||
|
||||
function assert(cond, msg) {
|
||||
if (!cond) throw new Error(msg);
|
||||
}
|
||||
|
||||
{
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(
|
||||
"a: $VAR_x\nb: $SECRET_y\nc: $CONTEXT_z\nd: passwordSecret: gmail_app",
|
||||
);
|
||||
assert(changed, "detects legacy");
|
||||
assert(text.includes("{{ vars.x }}"), "var rewrite");
|
||||
assert(text.includes("{{ secrets.y }}"), "secret rewrite");
|
||||
assert(text.includes("{{ context.z }}"), "context rewrite");
|
||||
assert(text.includes("passwordSecret: gmail_app"), "bare secret name untouched");
|
||||
}
|
||||
|
||||
await migrate();
|
||||
const first = await migrateDefaultOwnerIfNeeded();
|
||||
const second = await migrateDefaultOwnerIfNeeded();
|
||||
assert(second.secrets.moved === 0, "second pass moves no secrets");
|
||||
assert(second.runs.moved === 0, "second pass moves no runs");
|
||||
await db.destroy();
|
||||
|
||||
console.log("owner-migrate smoke passed", {
|
||||
firstSecretsMoved: first.secrets.moved,
|
||||
secondSecretsMoved: second.secrets.moved,
|
||||
});
|
||||
@@ -1,51 +0,0 @@
|
||||
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;
|
||||
}
|
||||
@@ -6,32 +6,19 @@ import { useVariables } from "../../api/hooks.js";
|
||||
const VAR_PEEK_MAX = 48;
|
||||
const MUSTACHE_RE =
|
||||
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g;
|
||||
const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
|
||||
/**
|
||||
* @param {unknown} value
|
||||
* @returns {{
|
||||
* kind: "mustache" | "legacy",
|
||||
* kind: "mustache",
|
||||
* path?: string,
|
||||
* root?: string,
|
||||
* name?: string,
|
||||
* legacyKind?: string,
|
||||
* legacyName?: string,
|
||||
* } | null}
|
||||
*/
|
||||
function describeConfigRef(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const trimmed = value.trim();
|
||||
const legacy = LEGACY_PREFIX_RE.exec(trimmed);
|
||||
if (legacy) {
|
||||
const legacyKind =
|
||||
legacy[1] === "VAR" ? "variable" : legacy[1] === "SECRET" ? "secret" : "context";
|
||||
return {
|
||||
kind: "legacy",
|
||||
legacyKind,
|
||||
legacyName: legacy[2] ?? "",
|
||||
};
|
||||
}
|
||||
MUSTACHE_RE.lastIndex = 0;
|
||||
const match = MUSTACHE_RE.exec(trimmed);
|
||||
if (!match) return null;
|
||||
@@ -90,20 +77,6 @@ export function ConfigRefHint({ value, owner }) {
|
||||
|
||||
if (!ref) return null;
|
||||
|
||||
if (ref.kind === "legacy") {
|
||||
const hint =
|
||||
ref.legacyKind === "variable"
|
||||
? `{{ vars.${ref.legacyName || "name"} }}`
|
||||
: ref.legacyKind === "secret"
|
||||
? `{{ secrets.${ref.legacyName || "name"} }}`
|
||||
: `{{ context.${ref.legacyName || "name"} }}`;
|
||||
return (
|
||||
<p className="text-xs text-error">
|
||||
legacy ${ref.legacyKind} ref — use <span className="font-mono">{hint}</span>
|
||||
</p>
|
||||
);
|
||||
}
|
||||
|
||||
if (!isVar) {
|
||||
const label =
|
||||
ref.root === "secrets"
|
||||
|
||||
Reference in New Issue
Block a user