Author SHA1 Message Date
nsrb a20b26b4c7 feat(plugins): add duplicate plugin functionality and related API endpoints
- Implemented a new `duplicatePlugin` function to allow copying of installed plugins with a new ID.
- Added API endpoint for duplicating plugins, including error handling for various edge cases.
- Introduced frontend hooks and UI components for duplicating plugins in the script management interface.
- Enhanced tests to validate the duplication process and ensure proper error handling.
2026-08-30 19:33:36 +07:00
nsrb f1cdb7ac68 chore(paths): refactor data directory structure and update references
- Updated paths for instance data, workflows, and logs to use a unified `data/` directory.
- Adjusted related documentation in AGENTS.md and README.md to reflect the new data structure.
- Refactored path handling in server files to ensure consistency and clarity in data management.
- Enhanced test scripts to align with the new directory structure for improved organization.
2026-08-30 15:51:33 +07:00
nsrb 3f0774813d 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.
2026-08-30 15:37:58 +07:00
37 changed files with 482 additions and 674 deletions
+3 -1
View File
@@ -1,6 +1,8 @@
node_modules/ node_modules/
.pnpm-store/ .pnpm-store/
# Instance data (SQLite, live workflows, control-state, backups)
data/ data/
# Process logs
logs/ logs/
*.db *.db
*.db-* *.db-*
@@ -10,7 +12,7 @@ packages/web/dist
# Personal/local scripts and workflows (not for the repo) # Personal/local scripts and workflows (not for the repo)
debug-*.js debug-*.js
# Legacy live workflow tree (migrated to packages/server/data/workflows/) # Legacy live workflow trees
packages/server/workflows/ packages/server/workflows/
# Plugin install staging and per-plugin deps # Plugin install staging and per-plugin deps
+3 -3
View File
@@ -4,14 +4,14 @@ This file tells agents how to add a **user plugin**. Do not put personal or site
## Workflows (instance data) ## Workflows (instance data)
Live workflows are **not** product source. They live under `packages/server/data/workflows/<owner>/` (gitignored; override with `JFLOW_WORKFLOWS_DIR`). Live workflows are **not** product source. They live under `data/workflows/<owner>/` (gitignored; override with `JFLOW_WORKFLOWS_DIR`).
| Kind | In git? | Path | | Kind | In git? | Path |
|---|---|---| |---|---|---|
| Live / personal YAML | No | `packages/server/data/workflows/<owner>/` | | Live / personal YAML | No | `data/workflows/<owner>/` |
| Example presets | Yes | `examples/workflows/*.yaml` (copy into editor only; runner does not load them) | | 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 `data/workflows/`. Prefer owner `local`. Example presets must use **core** scripts only (no `plugin/…` that requires install).
## Where things live ## Where things live
+9 -7
View File
@@ -55,16 +55,16 @@ pnpm --dir packages/server reset-admin -- --username admin --password 'your-pass
| Kind | Loaded by runner? | Location | | Kind | Loaded by runner? | Location |
|---|---|---| |---|---|---|
| **Live workflows** | Yes | `packages/server/data/workflows/<owner>/` (gitignored) | | **Live workflows** | Yes | `data/workflows/<owner>/` (gitignored) |
| **Example presets** | No | `examples/workflows/*.yaml` — offered when creating a new workflow | | **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). - 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). - 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`. - Override the live store in tests with `JFLOW_WORKFLOWS_DIR`.
```bash ```bash
# Smoke # Smoke (isolated under packages/server/data — not the live instance tree)
JFLOW_DATA_DIR=packages/server/data \
JFLOW_PLUGINS_DIR=packages/server/data/plugins-smoke-test \ JFLOW_PLUGINS_DIR=packages/server/data/plugins-smoke-test \
JFLOW_DB_PATH=packages/server/data/plugins-smoke.db \ JFLOW_DB_PATH=packages/server/data/plugins-smoke.db \
node packages/server/test/plugins-smoke.js node packages/server/test/plugins-smoke.js
@@ -109,7 +109,7 @@ topic: "{{ vars.ntfy_prefix }}/{{ data.channel }}"
| `context` | Run clipboard (nested) | | `context` | Run clipboard (nested) |
| `data` | This step’s input (nested; numeric segments index arrays) | | `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. YAML **SET** evaluates JSONata against the full `ctx`; the result is `output` (the next step’s data). `jsonata.js` does the same.
@@ -145,7 +145,7 @@ Admin UI route **Ops** (`/ops`) talks to the control process.
| Drain restart | Pause → wait active=0 → stop children → migrate → recreate → resume | | Drain restart | Pause → wait active=0 → stop children → migrate → recreate → resume |
| Force restart | Same without waiting (interrupts active runs; orphans marked `worker_lost`) | | Force restart | Same without waiting (interrupts active runs; orphans marked `worker_lost`) |
Desired state is stored in `packages/server/data/control-state.json` (generation, worker count, restart-needed). Plugin installs (later) bump generation and set restart-needed; you apply with Drain restart. Desired state is stored in `data/control-state.json` (generation, worker count, restart-needed). Plugin installs (later) bump generation and set restart-needed; you apply with Drain restart.
## Environment ## Environment
@@ -153,8 +153,10 @@ Desired state is stored in `packages/server/data/control-state.json` (generation
|---|---|---| |---|---|---|
| `JFLOW_JWT_SECRET` | `jflow-dev-secret` (dev only) | **Required in production**. | | `JFLOW_JWT_SECRET` | `jflow-dev-secret` (dev only) | **Required in production**. |
| `JFLOW_SECRETS_KEY` | `jflow-dev-secrets-key` (dev only) | Master key for named secrets. **Required in production**. Changing it makes existing secrets unreadable. 64 hex chars are used as a raw AES-256 key; any other string is derived with scrypt. | | `JFLOW_SECRETS_KEY` | `jflow-dev-secrets-key` (dev only) | Master key for named secrets. **Required in production**. Changing it makes existing secrets unreadable. 64 hex chars are used as a raw AES-256 key; any other string is derived with scrypt. |
| `JFLOW_DB_PATH` | `packages/server/data/jerapah-flow.db` | SQLite file. | | `JFLOW_DATA_DIR` | `data/` | Instance data root (SQLite, workflows, control-state, backups, trash). Falls back to `packages/server/data` if that tree still has the db or workflows. |
| `JFLOW_WORKFLOWS_DIR` | `packages/server/data/workflows` | Live workflow YAML (instance data). | | `JFLOW_DB_PATH` | `data/jerapah-flow.db` | SQLite file. |
| `JFLOW_WORKFLOWS_DIR` | `data/workflows` | Live workflow YAML (instance data). |
| `JFLOW_LOGS_DIR` | `logs/` | Rolling process logs. |
| `REDIS_URL` | `redis://127.0.0.1:6379` | Redis for BullMQ workflow queue. **Required** — the server will not start if Redis is unreachable. | | `REDIS_URL` | `redis://127.0.0.1:6379` | Redis for BullMQ workflow queue. **Required** — the server will not start if Redis is unreachable. |
| `REDIS_PASS` | — | Optional Redis AUTH password (sent via ioredis `password`). Prefer this over embedding credentials in `REDIS_URL` so logs stay clean. | | `REDIS_PASS` | — | Optional Redis AUTH password (sent via ioredis `password`). Prefer this over embedding credentials in `REDIS_URL` so logs stay clean. |
| `JFLOW_QUEUE_NAME` | `jerapah-workflows` | BullMQ queue name. | | `JFLOW_QUEUE_NAME` | `jerapah-workflows` | BullMQ queue name. |
+2 -2
View File
@@ -1,8 +1,8 @@
import fs from "fs"; import fs from "fs";
import path from "path"; import path from "path";
import { SERVER_ROOT } from "./paths.js"; import { REPO_ROOT } from "./paths.js";
const ROOT_PKG = path.resolve(SERVER_ROOT, "../../package.json"); const ROOT_PKG = path.join(REPO_ROOT, "package.json");
/** /**
* JerapahFlow app version from the monorepo root package.json. * JerapahFlow app version from the monorepo root package.json.
-17
View File
@@ -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 };
}
-32
View File
@@ -8,7 +8,6 @@ const MUSTACHE_TOKEN_RE =
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g; /\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g;
const WHOLE_MUSTACHE_RE = const WHOLE_MUSTACHE_RE =
/^\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}$/; /^\{\{\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 {{ * @typedef {{
@@ -19,30 +18,6 @@ const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
* }} ConfigRefCtx * }} 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. * Walk config (objects/arrays) and interpolate `{{ path }}` strings.
* Does not walk trigger data. * Does not walk trigger data.
@@ -84,13 +59,6 @@ export async function resolveConfigRefs(value, ctx, seen = new WeakSet()) {
* @returns {Promise<unknown>} * @returns {Promise<unknown>}
*/ */
async function resolveStringRef(value, ctx) { 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); const whole = WHOLE_MUSTACHE_RE.exec(value);
if (whole && whole[0] === value) { if (whole && whole[0] === value) {
return resolvePath(whole[1], ctx, { raw: value, allowObject: true }); 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.
}
@@ -17,6 +17,9 @@ export async function up(knex) {
t.text("output"); t.text("output");
t.text("error"); t.text("error");
t.text("parent_run_id").references("id").inTable("workflow_runs"); 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( await knex.schema.raw(
@@ -28,45 +31,11 @@ export async function up(knex) {
await knex.schema.raw( await knex.schema.raw(
"CREATE INDEX workflow_runs_status_started_at_idx ON workflow_runs (status, started_at DESC)", "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( 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( 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 * @param {import("knex").Knex} knex
*/ */
export async function down(knex) { export async function down(knex) {
await knex.schema.dropTableIfExists("logs");
await knex.schema.dropTableIfExists("step_runs");
await knex.schema.dropTableIfExists("workflow_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");
}
@@ -29,35 +29,35 @@ const DEFAULT_EMAIL_TEMPLATE = `<!DOCTYPE html>
* @param {import("knex").Knex} knex * @param {import("knex").Knex} knex
*/ */
export async function up(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.text("kind").notNullable().defaultTo("response");
t.integer("system").notNullable().defaultTo(0); t.integer("system").notNullable().defaultTo(0);
t.text("created_at").notNullable();
t.text("updated_at").notNullable();
}); });
const now = new Date().toISOString(); const now = new Date().toISOString();
const existing = await knex("http_pages").where({ name: "email-default" }).first(); await knex("http_pages").insert({
if (!existing) { id: "00000000-0000-4000-8000-000000000001",
await knex("http_pages").insert({ name: "email-default",
id: "00000000-0000-4000-8000-000000000001", content: DEFAULT_EMAIL_TEMPLATE,
name: "email-default", mime: "html",
content: DEFAULT_EMAIL_TEMPLATE, status: 200,
mime: "html", kind: "template",
status: 200, system: 1,
kind: "template", created_at: now,
system: 1, updated_at: now,
created_at: now, });
updated_at: now,
});
}
} }
/** /**
* @param {import("knex").Knex} knex * @param {import("knex").Knex} knex
*/ */
export async function down(knex) { export async function down(knex) {
await knex("http_pages").where({ name: "email-default", system: 1 }).del(); await knex.schema.dropTableIfExists("http_pages");
await knex.schema.alterTable("http_pages", (t) => {
t.dropColumn("kind");
t.dropColumn("system");
});
} }
@@ -2,21 +2,11 @@
* @param {import("knex").Knex} knex * @param {import("knex").Knex} knex
*/ */
export async function up(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) => { await knex.schema.createTable("http_auths", (t) => {
t.text("id").primary(); t.text("id").primary();
t.text("name").notNullable().unique(); t.text("name").notNullable().unique();
t.text("type").notNullable(); // bearer | basic | header t.text("type").notNullable();
t.text("config").notNullable(); // JSON t.text("config").notNullable();
t.integer("unauthorized_status").nullable(); t.integer("unauthorized_status").nullable();
t.text("unauthorized_response").nullable(); t.text("unauthorized_response").nullable();
t.text("created_at").notNullable(); t.text("created_at").notNullable();
@@ -29,5 +19,4 @@ export async function up(knex) {
*/ */
export async function down(knex) { export async function down(knex) {
await knex.schema.dropTableIfExists("http_auths"); await knex.schema.dropTableIfExists("http_auths");
await knex.schema.dropTableIfExists("http_pages");
} }
@@ -21,29 +21,11 @@ export async function up(knex) {
await knex.schema.raw( await knex.schema.raw(
"CREATE INDEX workflow_revisions_workflow_id_created_at_idx ON workflow_revisions (workflow_id, created_at DESC)", "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 * @param {import("knex").Knex} knex
*/ */
export async function down(knex) { export async function down(knex) {
await knex.schema.dropTableIfExists("workflow_trash");
await knex.schema.dropTableIfExists("workflow_revisions"); 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");
}
@@ -14,7 +14,9 @@ export async function up(knex) {
t.unique(["owner", "name"]); 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)",
);
} }
/** /**
-268
View File
@@ -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;
}
+3 -4
View File
@@ -13,13 +13,12 @@
"start:control": "node control.js", "start:control": "node control.js",
"start:web": "node web-server.js", "start:web": "node web-server.js",
"migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\"", "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_DATA_DIR=./data JFLOW_PLUGINS_DIR=./data/plugins-smoke-test JFLOW_DB_PATH=./data/plugins-smoke.db node test/plugins-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:workflow-history": "JFLOW_DATA_DIR=./data JFLOW_WORKFLOWS_DIR=./data/workflow-history-smoke JFLOW_DB_PATH=./data/workflow-history-smoke.db node test/workflow-history-smoke.js",
"reset-admin": "node reset-admin.js", "reset-admin": "node reset-admin.js",
"test:profiles": "node test/profiles-smoke.js", "test:profiles": "node test/profiles-smoke.js",
"test:set-dry-run": "node test/set-dry-run-smoke.js", "test:set-dry-run": "node test/set-dry-run-smoke.js",
"test:config-refs": "node test/config-refs-smoke.js", "test:config-refs": "node test/config-refs-smoke.js"
"test:owner-migrate": "node test/owner-migrate-smoke.js"
}, },
"dependencies": { "dependencies": {
"@jerapah-flow/shared": "workspace:*", "@jerapah-flow/shared": "workspace:*",
+39 -17
View File
@@ -1,27 +1,49 @@
import fs from "fs";
import path from "path"; import path from "path";
import { fileURLToPath } from "url"; import { fileURLToPath } from "url";
export const SERVER_ROOT = path.dirname(fileURLToPath(import.meta.url)); export const SERVER_ROOT = path.dirname(fileURLToPath(import.meta.url));
export const REPO_ROOT = path.resolve(SERVER_ROOT, "../..");
export const SCRIPTS_DIR = path.join(SERVER_ROOT, "scripts"); export const SCRIPTS_DIR = path.join(SERVER_ROOT, "scripts");
export const DATA_DIR = path.join(SERVER_ROOT, "data");
/** Prefer `preferred` unless only `legacy` already has files. */
function existingDir(preferred, legacy, probe) {
const has = (dir) =>
probe ? probe(dir) : fs.existsSync(dir);
if (has(preferred) || !has(legacy)) return preferred;
return legacy;
}
function hasInstanceData(dir) {
return (
fs.existsSync(path.join(dir, "jerapah-flow.db")) ||
fs.existsSync(path.join(dir, "workflows"))
);
}
/** Instance data (SQLite, live workflows, control-state). Not product source. */
export const DATA_DIR = path.resolve(
process.env.JFLOW_DATA_DIR ??
existingDir(
path.join(REPO_ROOT, "data"),
path.join(SERVER_ROOT, "data"),
hasInstanceData,
),
);
/** Live instance workflows (not shipped in git). Override for tests. */ /** Live instance workflows (not shipped in git). Override for tests. */
export const WORKFLOWS_DIR = export const WORKFLOWS_DIR = path.resolve(
process.env.JFLOW_WORKFLOWS_DIR ?? path.join(DATA_DIR, "workflows"); 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). */ /** User plugins (repo-root /plugins, outside the pnpm workspace). */
export const PLUGINS_DIR = export const PLUGINS_DIR = path.resolve(
process.env.JFLOW_PLUGINS_DIR ?? process.env.JFLOW_PLUGINS_DIR ?? path.join(REPO_ROOT, "plugins"),
path.resolve(SERVER_ROOT, "../../plugins"); );
/** Example plugin sources shipped with the repo. */ /** Example plugin sources shipped with the repo. */
export const EXAMPLE_PLUGINS_DIR = path.resolve( export const EXAMPLE_PLUGINS_DIR = path.join(REPO_ROOT, "examples/plugins");
SERVER_ROOT,
"../../examples/plugins",
);
/** Example workflow YAML presets (not loaded by the runner). */ /** Example workflow YAML presets (not loaded by the runner). */
export const EXAMPLE_WORKFLOWS_DIR = path.resolve( export const EXAMPLE_WORKFLOWS_DIR = path.join(REPO_ROOT, "examples/workflows");
SERVER_ROOT, export const LOGS_DIR = path.resolve(
"../../examples/workflows", process.env.JFLOW_LOGS_DIR ??
existingDir(path.join(REPO_ROOT, "logs"), path.join(SERVER_ROOT, "logs")),
); );
export const LOGS_DIR = path.join(SERVER_ROOT, "logs"); export const WEB_DIST = path.join(REPO_ROOT, "packages/web/dist");
export const WEB_DIST = path.resolve(SERVER_ROOT, "../web/dist");
+90
View File
@@ -380,6 +380,96 @@ export function forkCoreScript(coreName, newId, opts = {}) {
} }
} }
/**
* Copy an installed plugin to a new plugin id.
*
* @param {string} sourceId
* @param {string} newId
* @param {{ description?: string }} [opts]
*/
export function duplicatePlugin(sourceId, newId, opts = {}) {
const fromId = assertPluginId(sourceId);
const id = assertPluginId(newId);
if (id === fromId) {
const err = new Error("cannot duplicate onto itself");
err.statusCode = 400;
throw err;
}
if (coreBareNames().has(id)) {
const err = new Error(`plugin id collides with core script: ${id}`);
err.statusCode = 409;
throw err;
}
if (fs.existsSync(pluginDir(id))) {
const err = new Error(`plugin already exists: ${id}`);
err.statusCode = 409;
throw err;
}
const source = getInstalledPlugin(fromId);
if (!source) {
const err = new Error(`plugin not found: ${fromId}`);
err.statusCode = 404;
throw err;
}
if (!source.manifest) {
const err = new Error(
source.compatError || `plugin has no valid manifest: ${fromId}`,
);
err.statusCode = 400;
throw err;
}
const staging = path.join(PLUGINS_DIR, `.staging-dup-${id}-${Date.now()}`);
fs.mkdirSync(staging, { recursive: true });
try {
fs.cpSync(source.dir, staging, {
recursive: true,
filter: (src) => {
const base = path.basename(src);
return base !== "node_modules" && base !== ".disabled";
},
});
const manifest = buildManifest({
id,
name: id,
version: source.manifest.version,
jerapah: source.manifest.jerapah,
main: source.manifest.main,
description:
opts.description ?? source.manifest.description ?? null,
});
fs.writeFileSync(
path.join(staging, PLUGIN_MANIFEST),
`${JSON.stringify(manifest, null, 2)}\n`,
"utf8",
);
const pkgPath = path.join(staging, "package.json");
if (fs.existsSync(pkgPath)) {
let pkg = {};
try {
pkg = JSON.parse(fs.readFileSync(pkgPath, "utf8"));
} catch {
pkg = {};
}
if (pkg == null || typeof pkg !== "object" || Array.isArray(pkg)) {
pkg = {};
}
pkg.name = `jflow-plugin-${id}`;
fs.writeFileSync(pkgPath, `${JSON.stringify(pkg, null, 2)}\n`, "utf8");
}
return installPluginFromDirectory(staging, {
overwrite: false,
reason: `plugin:${id} duplicated from ${fromId}`,
});
} finally {
fs.rmSync(staging, { recursive: true, force: true });
}
}
/** /**
* @param {string} pluginDirectory * @param {string} pluginDirectory
* @returns {((id: string) => unknown) | null} * @returns {((id: string) => unknown) | null}
+1 -3
View File
@@ -1,8 +1,6 @@
import path from "path"; import path from "path";
import pm2 from "pm2"; import pm2 from "pm2";
import { SERVER_ROOT } from "./paths.js"; import { REPO_ROOT, SERVER_ROOT } from "./paths.js";
const REPO_ROOT = path.resolve(SERVER_ROOT, "../..");
export const PM2_HTTP_NAME = "jflow-http"; export const PM2_HTTP_NAME = "jflow-http";
export const PM2_WORKER_NAME = "jflow-worker"; export const PM2_WORKER_NAME = "jflow-worker";
+9
View File
@@ -15,6 +15,10 @@ function ntfyHeaders(ctx) {
headers.Title = ctx.data.title; headers.Title = ctx.data.title;
} }
if (ctx.config?.markdown === true) {
headers.md = "true";
}
return headers; return headers;
} }
@@ -149,6 +153,11 @@ ntfy.meta = {
default: "https://ntfy.sh/jerapah-flow", default: "https://ntfy.sh/jerapah-flow",
description: "ntfy topic URL", description: "ntfy topic URL",
}, },
markdown: {
type: "boolean",
default: false,
description: "Send as Markdown (ntfy md header)",
},
fingerprint: { fingerprint: {
type: "string", type: "string",
required: false, required: false,
+36 -1
View File
@@ -8,6 +8,7 @@ import {
import * as fsStore from "../../fs-store.js"; import * as fsStore from "../../fs-store.js";
import { import {
forkCoreScript, forkCoreScript,
duplicatePlugin,
listCoreScriptNames, listCoreScriptNames,
listInstalledPlugins, listInstalledPlugins,
resolveScriptRef, resolveScriptRef,
@@ -29,7 +30,7 @@ import { normalizeStepResult } from "../../step-result.js";
import { resolveConfigRefs } from "../../config-refs.js"; import { resolveConfigRefs } from "../../config-refs.js";
import { getAppVersion } from "../../app-version.js"; import { getAppVersion } from "../../app-version.js";
import { EXAMPLE_PLUGINS_DIR } from "../../paths.js"; import { EXAMPLE_PLUGINS_DIR } from "../../paths.js";
import { pluginScriptRef } from "../../plugin-manifest.js"; import { parsePluginScriptRef, pluginScriptRef } from "../../plugin-manifest.js";
import { evaluateJsonata, SET_STEP_SCRIPT } from "../../workflow-parse.js"; import { evaluateJsonata, SET_STEP_SCRIPT } from "../../workflow-parse.js";
import { DEFAULT_OWNER } from "@jerapah-flow/shared"; import { DEFAULT_OWNER } from "@jerapah-flow/shared";
@@ -251,6 +252,40 @@ export default function scriptsPluginFactory(registry) {
} }
}); });
fastify.post("/scripts/:name/duplicate", async (req, reply) => {
const rawName = decodeURIComponent(
/** @type {{ name: string }} */ (req.params).name,
);
const body = /** @type {{ id?: string, description?: string }} */ (
req.body ?? {}
);
if (typeof body.id !== "string" || !body.id.trim()) {
return reply.code(400).send({ error: "id is required" });
}
const parsed = parsePluginScriptRef(rawName);
if (!parsed) {
return reply
.code(400)
.send({ error: "name must be a plugin ref (plugin/<id>)" });
}
try {
const installed = duplicatePlugin(parsed.id, body.id.trim(), {
description: body.description,
});
clearScriptCache();
return reply.code(201).send({
...installed,
restartNeeded: true,
warning:
"Plugins run as the JerapahFlow process user. Review code before install.",
});
} catch (err) {
return reply
.code(/** @type {any} */ (err).statusCode ?? 500)
.send({ error: err instanceof Error ? err.message : String(err) });
}
});
fastify.post("/scripts/:name/dry-run", async (req, reply) => { fastify.post("/scripts/:name/dry-run", async (req, reply) => {
const rawName = decodeURIComponent( const rawName = decodeURIComponent(
/** @type {{ name: string }} */ (req.params).name, /** @type {{ name: string }} */ (req.params).name,
-14
View File
@@ -31,8 +31,6 @@ import {
getRedisUrlForLog, getRedisUrlForLog,
} from "./workflow-queue.js"; } from "./workflow-queue.js";
import { purgeExpiredTrash } from "./workflow-trash.js"; import { purgeExpiredTrash } from "./workflow-trash.js";
import { migrateLegacyWorkflowsIfNeeded } from "./workflow-migrate.js";
import { migrateDefaultOwnerIfNeeded } from "./owner-migrate.js";
import { import {
getConfigGeneration, getConfigGeneration,
startHeartbeatLoop, 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, { const registry = createRegistry(server, {
queue: workflowQueue, queue: workflowQueue,
// Cron + HTTP triggers enqueue jobs; only the API process may own them. // Cron + HTTP triggers enqueue jobs; only the API process may own them.
+2 -36
View File
@@ -2,8 +2,7 @@ import { migrate, db } from "../db.js";
import { upsertSecret, deleteSecret } from "../secrets-store.js"; import { upsertSecret, deleteSecret } from "../secrets-store.js";
import { deleteVariable, upsertVariable } from "../variables-store.js"; import { deleteVariable, upsertVariable } from "../variables-store.js";
import { Secret } from "../secret-value.js"; import { Secret } from "../secret-value.js";
import { parseLegacyConfigRef, resolveConfigRefs } from "../config-refs.js"; import { resolveConfigRefs } from "../config-refs.js";
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
await migrate(); await migrate();
@@ -27,50 +26,17 @@ async function assertRejects(fn, match) {
const owner = "config_refs_smoke_owner"; const owner = "config_refs_smoke_owner";
const ctx = { owner, workflowKey: `${owner}/config-refs-smoke.yaml`, context: {}, data: {} }; 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); const literal = await resolveConfigRefs("password123", ctx);
assert(literal === "password123", "literal passthrough"); assert(literal === "password123", "literal passthrough");
const unknown = await resolveConfigRefs("$FOO_bar", ctx); const unknown = await resolveConfigRefs("$FOO_bar", ctx);
assert(unknown === "$FOO_bar", "$FOO_bar stays literal"); assert(unknown === "$FOO_bar", "$FOO_bar stays literal");
const embedded = await resolveConfigRefs("Bearer $SECRET_x", ctx); 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); const number = await resolveConfigRefs(42, ctx);
assert(number === 42, "number passthrough"); 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 = []; const created = [];
try { try {
const secret = await upsertSecret({ 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,
});
+23
View File
@@ -2,6 +2,7 @@
* Smoke: core vs plugin scripts, fork, example install, resolve, run. * Smoke: core vs plugin scripts, fork, example install, resolve, run.
* *
* Run: * Run:
* JFLOW_DATA_DIR=packages/server/data \
* JFLOW_PLUGINS_DIR=packages/server/data/plugins-smoke-test \ * JFLOW_PLUGINS_DIR=packages/server/data/plugins-smoke-test \
* JFLOW_DB_PATH=packages/server/data/plugins-smoke.db \ * JFLOW_DB_PATH=packages/server/data/plugins-smoke.db \
* node packages/server/test/plugins-smoke.js * node packages/server/test/plugins-smoke.js
@@ -12,6 +13,7 @@ import { migrate, db } from "../db.js";
import { getAppVersion, satisfiesRange } from "../app-version.js"; import { getAppVersion, satisfiesRange } from "../app-version.js";
import { import {
forkCoreScript, forkCoreScript,
duplicatePlugin,
resolveScriptRef, resolveScriptRef,
uninstallPlugin, uninstallPlugin,
listInstalledPlugins, listInstalledPlugins,
@@ -82,6 +84,26 @@ async function main() {
); );
assert.equal(blankRun.output.ok, true); assert.equal(blankRun.output.ok, true);
const duplicated = duplicatePlugin("blank-smoke", "blank-smoke-copy");
assert.equal(duplicated.scriptRef, "plugin/blank-smoke-copy");
clearScriptCache();
assert.equal(resolveScriptRef("plugin/blank-smoke-copy").kind, "plugin");
const dupRun = await runScript(
"plugin/blank-smoke-copy",
{ data: 1, context: {}, config: null },
{ log: silent, workflowName: "smoke", owner: "default" },
);
assert.equal(dupRun.output.ok, true);
let hitDupSelf = false;
try {
duplicatePlugin("blank-smoke", "blank-smoke");
} catch (err) {
hitDupSelf = true;
assert.match(String(err.message), /itself/);
}
assert.equal(hitDupSelf, true);
let hit = false; let hit = false;
try { try {
forkCoreScript("ntfy.js", "ntfy"); forkCoreScript("ntfy.js", "ntfy");
@@ -101,6 +123,7 @@ async function main() {
uninstallPlugin("jsonata-smoke-fork"); uninstallPlugin("jsonata-smoke-fork");
uninstallPlugin("blank-smoke"); uninstallPlugin("blank-smoke");
uninstallPlugin("blank-smoke-copy");
uninstallPlugin("get-current-time"); uninstallPlugin("get-current-time");
console.log("plugins-smoke: ok"); console.log("plugins-smoke: ok");
-51
View File
@@ -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;
}
+17
View File
@@ -59,6 +59,23 @@ export function useForkScript() {
}); });
} }
export function useDuplicatePlugin() {
const qc = useQueryClient();
return useMutation({
mutationFn: async ({ name, id, description }) =>
(
await api.post(`/scripts/${encodeURIComponent(name)}/duplicate`, {
id,
description,
})
).data,
onSuccess: () => {
qc.invalidateQueries({ queryKey: ["scripts"] });
qc.invalidateQueries({ queryKey: ["ops-status"] });
},
});
}
export function useInstallPlugin() { export function useInstallPlugin() {
const qc = useQueryClient(); const qc = useQueryClient();
return useMutation({ return useMutation({
@@ -6,32 +6,19 @@ import { useVariables } from "../../api/hooks.js";
const VAR_PEEK_MAX = 48; const VAR_PEEK_MAX = 48;
const MUSTACHE_RE = const MUSTACHE_RE =
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g; /\{\{\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 * @param {unknown} value
* @returns {{ * @returns {{
* kind: "mustache" | "legacy", * kind: "mustache",
* path?: string, * path?: string,
* root?: string, * root?: string,
* name?: string, * name?: string,
* legacyKind?: string,
* legacyName?: string,
* } | null} * } | null}
*/ */
function describeConfigRef(value) { function describeConfigRef(value) {
if (typeof value !== "string") return null; if (typeof value !== "string") return null;
const trimmed = value.trim(); 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; MUSTACHE_RE.lastIndex = 0;
const match = MUSTACHE_RE.exec(trimmed); const match = MUSTACHE_RE.exec(trimmed);
if (!match) return null; if (!match) return null;
@@ -90,20 +77,6 @@ export function ConfigRefHint({ value, owner }) {
if (!ref) return null; 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) { if (!isVar) {
const label = const label =
ref.root === "secrets" ref.root === "secrets"
+39 -1
View File
@@ -4,6 +4,7 @@ import { LuArrowLeft, LuCopy, LuPlay, LuSave } from "react-icons/lu";
import { errorMessage } from "../api/client.js"; import { errorMessage } from "../api/client.js";
import { import {
useCreatePlugin, useCreatePlugin,
useDuplicatePlugin,
useForkScript, useForkScript,
useSaveScript, useSaveScript,
useScript, useScript,
@@ -119,9 +120,11 @@ export function ScriptEditPage() {
const existing = useScript(name); const existing = useScript(name);
const save = useSaveScript(); const save = useSaveScript();
const fork = useForkScript(); const fork = useForkScript();
const duplicate = useDuplicatePlugin();
const [content, setContent] = useState(""); const [content, setContent] = useState("");
const [contentReady, setContentReady] = useState(false); const [contentReady, setContentReady] = useState(false);
const [forkId, setForkId] = useState(""); const [forkId, setForkId] = useState("");
const [duplicateId, setDuplicateId] = useState("");
const isCore = existing.data?.kind === "core" || existing.data?.editable === false; const isCore = existing.data?.kind === "core" || existing.data?.editable === false;
@@ -163,6 +166,21 @@ export function ScriptEditPage() {
); );
} }
function onDuplicate(e) {
e.preventDefault();
const id = normalizePluginId(duplicateId);
if (!id) return;
duplicate.mutate(
{ name, id },
{
onSuccess: (data) => {
notify.success("Duplicated — drain-restart recommended");
navigate(`/scripts/${encodeURIComponent(data.scriptRef)}/edit`);
},
},
);
}
if (existing.isLoading) { if (existing.isLoading) {
return ( return (
<div className="flex min-h-[12rem] items-center justify-center"> <div className="flex min-h-[12rem] items-center justify-center">
@@ -238,7 +256,27 @@ export function ScriptEditPage() {
<span className="text-error text-sm">{errorMessage(fork.error)}</span> <span className="text-error text-sm">{errorMessage(fork.error)}</span>
) : null} ) : null}
</form> </form>
) : null} ) : (
<form onSubmit={onDuplicate} className="flex flex-wrap items-center gap-2">
<input
className="input input-sm font-mono w-56"
placeholder="duplicate id (e.g. my-plugin-copy)"
value={duplicateId}
onChange={(e) => setDuplicateId(e.target.value)}
/>
<button
type="submit"
className="btn btn-sm"
disabled={duplicate.isPending || !normalizePluginId(duplicateId)}
>
<LuCopy className="size-4" />
Duplicate plugin
</button>
{duplicate.isError ? (
<span className="text-error text-sm">{errorMessage(duplicate.error)}</span>
) : null}
</form>
)}
<form id="script-edit-form" onSubmit={onSave} className="min-h-0 flex-1"> <form id="script-edit-form" onSubmit={onSave} className="min-h-0 flex-1">
<CodeEditor <CodeEditor
+85 -8
View File
@@ -4,6 +4,7 @@ import { LuCopy, LuPencil, LuPlay, LuPlus, LuSearch, LuTrash2 } from "react-icon
import { errorMessage } from "../api/client.js"; import { errorMessage } from "../api/client.js";
import { import {
useDeleteScript, useDeleteScript,
useDuplicatePlugin,
useForkScript, useForkScript,
useInstallPlugin, useInstallPlugin,
useScripts, useScripts,
@@ -23,8 +24,11 @@ export function ScriptsPage() {
const [query, setQuery] = useState(""); const [query, setQuery] = useState("");
const [forkFor, setForkFor] = useState(null); const [forkFor, setForkFor] = useState(null);
const [forkId, setForkId] = useState(""); const [forkId, setForkId] = useState("");
const [duplicateFor, setDuplicateFor] = useState(null);
const [duplicateId, setDuplicateId] = useState("");
const del = useDeleteScript(); const del = useDeleteScript();
const fork = useForkScript(); const fork = useForkScript();
const duplicate = useDuplicatePlugin();
const install = useInstallPlugin(); const install = useInstallPlugin();
const { notify } = useNotifications(); const { notify } = useNotifications();
@@ -182,14 +186,33 @@ export function ScriptsPage() {
<LuCopy className="size-4" /> <LuCopy className="size-4" />
</button> </button>
) : ( ) : (
<Link <>
to={`/scripts/${encodeURIComponent(name)}/edit`} <Link
className="btn btn-ghost btn-xs" to={`/scripts/${encodeURIComponent(name)}/edit`}
title="Edit" className="btn btn-ghost btn-xs"
aria-label="Edit" title="Edit"
> aria-label="Edit"
<LuPencil className="size-4" /> >
</Link> <LuPencil className="size-4" />
</Link>
<button
type="button"
className="btn btn-ghost btn-xs"
title="Duplicate"
aria-label="Duplicate"
onClick={() => {
setDuplicateFor(name);
setDuplicateId(
String(name)
.replace(/^plugin\//i, "")
.replace(/\.js$/i, "")
.toLowerCase() + "-copy",
);
}}
>
<LuCopy className="size-4" />
</button>
</>
)} )}
{!isCore ? ( {!isCore ? (
<button <button
@@ -264,6 +287,60 @@ export function ScriptsPage() {
</dialog> </dialog>
) : null} ) : null}
{duplicateFor ? (
<dialog className="modal modal-open">
<div className="modal-box">
<h3 className="font-semibold">Duplicate {duplicateFor}</h3>
<p className="text-sm opacity-70 py-2">
Creates <code>plugin/&lt;id&gt;</code> from this plugin.
</p>
<input
className="input input-bordered input-sm w-full font-mono"
value={duplicateId}
onChange={(e) => setDuplicateId(e.target.value)}
/>
{duplicate.isError ? (
<p className="text-error text-sm mt-2">{errorMessage(duplicate.error)}</p>
) : null}
<div className="modal-action">
<button
type="button"
className="btn btn-ghost btn-sm"
onClick={() => setDuplicateFor(null)}
>
Cancel
</button>
<button
type="button"
className="btn btn-primary btn-sm"
disabled={duplicate.isPending || !duplicateId.trim()}
onClick={() =>
duplicate.mutate(
{ name: duplicateFor, id: duplicateId.trim() },
{
onSuccess: (data) => {
setDuplicateFor(null);
notify.success("Duplicated — drain-restart recommended");
navigate(
`/scripts/${encodeURIComponent(data.scriptRef)}/edit`,
);
},
},
)
}
>
Duplicate
</button>
</div>
</div>
<form method="dialog" className="modal-backdrop">
<button type="button" onClick={() => setDuplicateFor(null)}>
close
</button>
</form>
</dialog>
) : null}
<ConfirmDialog <ConfirmDialog
open={Boolean(confirmDelete)} open={Boolean(confirmDelete)}
title={confirmDelete ? `Delete ${confirmDelete}?` : ""} title={confirmDelete ? `Delete ${confirmDelete}?` : ""}