13 Commits
Author SHA1 Message Date
nsrb 5615d989fa feat(migration): migrate legacy owner data and update config reference handling
- Implemented a migration process to move resources from the legacy owner "default" to the new default owner "local", ensuring no data loss during the transition.
- Updated configuration reference handling to replace legacy `$VAR_`, `$SECRET_`, and `$CONTEXT_` prefixes with mustache-style `{{ vars.name }}`, `{{ secrets.name }}`, and `{{ context.name }}`.
- Enhanced YAML configuration files and scripts to reflect the new mustache syntax, improving consistency across the application.
- Added tests to validate the migration process and ensure proper handling of legacy references.
- Updated documentation to guide users on the new configuration reference format.
2026-08-30 05:32:48 +07:00
nsrbandCursor e2e70a5aad fix(workflows): unregister missing yaml on trash
Ghost register entries returned 404 on delete and stayed in the list. Unregister first so move-to-trash still removes them when the file is gone.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-23 09:17:10 +07:00
nsrb dd98421b08 fix(server): refine error handling and PM2 connection logic
Deploy to Raspberry Pi / deploy (push) Canceled after 1m16s
- Reduced the maximum restarts for PM2 processes from 20 to 8 and added a restart delay of 2000ms for better stability.
- Improved error handling during server startup and PM2 connection, ensuring clearer logging and graceful exits on failure.
- Introduced a utility function to filter out PM2 metadata from environment variables, preventing conflicts during process management.
2026-08-23 08:45:26 +07:00
nsrb 69312c9080 feat(pm2): update PM2 configuration and introduce new start script
Deploy to Raspberry Pi / deploy (push) Canceled after 37s
- Updated PM2 version to 6.0.14 in package.json and pnpm-lock.yaml for improved stability.
- Changed interpreter for PM2 processes to use `process.execPath` for consistency.
- Added a new `start:pm2` script in package.json to streamline PM2 process management.
- Updated README to reflect changes in PM2 usage and instructions for starting the application.
2026-08-23 07:25:25 +07:00
nsrb c8697532f9 feat(ecosystem): add exec_mode to PM2 configuration for control and web server
Deploy to Raspberry Pi / deploy (push) Failing after 10s
- Introduced `exec_mode: "fork"` for both the control plane and web server processes to enhance performance and resource management.
2026-08-22 23:11:48 +07:00
nsrb be8122f9e5 feat(ecosystem): update PM2 configuration for control plane and web server
Deploy to Raspberry Pi / deploy (push) Successful in 52s
- Renamed the control process to `jflow-control` and updated the script path to `control.js`.
- Introduced a new `jflow-web` process for the production UI, serving on port 8500.
- Enhanced environment variable handling for both processes, ensuring proper port assignments.
- Updated README to reflect new process architecture and usage instructions for starting the application.
2026-08-22 22:32:04 +07:00
nsrb 0a8463d8f4 Merge branch 'main' into dev
Deploy to Raspberry Pi / deploy (push) Successful in 4m3s
2026-08-22 21:48:51 +07:00
nsrb 721f15e900 Merge branch 'dev' of https://git.home.0dev.web.id/nsrb/jerapah-flow into dev
Deploy to Raspberry Pi / deploy (push) Successful in 1m44s
2026-08-21 10:34:34 +07:00
nsrb 877ea9e3f8 Merge branch 'main' into dev 2026-08-21 09:31:01 +07:00
nsrb 93b51dcaaa Update .gitea/workflows/deploy.yaml
Deploy to Raspberry Pi / deploy (push) Successful in 1m43s
2026-08-20 18:56:34 -04:00
nsrb eb7c318c13 Update .gitea/workflows/deploy.yaml
Deploy to Raspberry Pi / deploy (push) Failing after 4s
2026-08-20 18:40:36 -04:00
nsrb 864427b45c feat(workflows): add auto-disable feature for workflows on consecutive failures
Deploy to Raspberry Pi / deploy (push) Canceled after 0s
- Introduced `disableOnConsecutiveFailures` option in workflow triggers to automatically disable workflows after reaching a specified failure threshold.
- Updated related functions to handle the new feature, including persistence of the disabled state and reloading of registries.
- Enhanced UI components to support the new option, allowing users to toggle the auto-disable feature in the workflow configuration.
2026-08-21 05:37:24 +07:00
nsrb 46aa8ca327 add deploy.yml for branch dev pushes
Deploy to Raspberry Pi / deploy (push) Canceled after 0s
2026-08-20 18:09:14 -04:00
39 changed files with 1364 additions and 272 deletions
+33
View File
@@ -0,0 +1,33 @@
name: Deploy to Raspberry Pi
on:
push:
branches: [dev]
jobs:
deploy:
runs-on: home # must match a label on your act_runner
steps:
- name: Deploy
run: |
set -euo pipefail
# act_runner uses bash --noprofile --norc; load nvm/pnpm explicitly
export NVM_DIR="/home/nsrb/.nvm"
export PNPM_HOME="/home/nsrb/.local/share/pnpm"
# shellcheck disable=SC1091
[ -s "$NVM_DIR/nvm.sh" ] && . "$NVM_DIR/nvm.sh"
export PATH="$PNPM_HOME/bin:$PATH"
APP_DIR=/home/nsrb/apps/jerapah-flow
cd "$APP_DIR"
git fetch origin dev
git checkout dev
git pull --ff-only origin dev
pnpm install --frozen-lockfile
pnpm build
pm2 startOrReload ecosystem.config.cjs --update-env
pm2 save
+3 -3
View File
@@ -120,7 +120,7 @@ Return `{ output, context?, skipRemaining? }`. Do not return `ctx`. Mutations of
|---|---|
| `ctx.data` | Step input (trigger payload, previous `output`, or DAG `needs`) |
| `ctx.context` | Run clipboard (plain object) |
| `ctx.config` | YAML `config` (secrets already unwrapped from `$SECRET_name`) |
| `ctx.config` | YAML `config` (mustache refs like `{{ secrets.name }}` already resolved) |
| `output` | Next step’s `data` |
| `context` | Next clipboard. Omit to keep incoming |
@@ -142,7 +142,7 @@ scripts:
script: plugin/my-plugin
config:
url: http://10.8.0.6:3030/notes
token: $SECRET_joplin_api_token
token: "{{ secrets.joplin_api_token }}"
```
Optional `name` is the display title in the editor and graph (falls back to the script filename). Canonical ref is `plugin/<id>` (`.js` suffix is optional).
@@ -153,4 +153,4 @@ Optional `name` is the display title in the editor and graph (falls back to the
- Use an id that matches a core script file (`ntfy`, `jsonata`, …).
- Mismatch folder name and `jerapah-plugin.json` `id` (plugin is disabled).
- Commit `plugins/.staging-*` or `plugins/*/node_modules/`.
- Put secrets in `script.js`; use YAML `$SECRET_name` / `$VAR_name`.
- Put secrets in `script.js`; use YAML `{{ secrets.name }}` / `{{ vars.name }}` (quote if the value starts with `{`).
+53 -10
View File
@@ -83,13 +83,34 @@ Each script is `async function main(ctx)` and **must** return:
|---|---|
| `ctx.data` | This step’s input (trigger payload, previous `output`, or DAG `needs`) |
| `ctx.context` | Run clipboard (plain object, default `{}`) |
| `ctx.config` | This step’s YAML config |
| `ctx.config` | This step’s YAML config (mustache refs already resolved) |
| `output` | Becomes the **next** step’s `data` |
| `context` | Next snapshot of the bag. Omitted → keep incoming |
| `skipRemaining` | Stop later steps. Sibling of `output`/`context`, not inside `output` |
Returning the full `ctx` is an error. Mutating `ctx.data` or `ctx.context` does not persist unless returned.
### Config interpolation
YAML `config` strings may use mustache paths. Quote values that start with `{`.
```yaml
url: "{{ vars.ntfy_channel }}"
token: "{{ secrets.joplin_api_token }}"
id: "{{ context.user.id }}"
title: "{{ data.httpResponse.data.date }}"
topic: "{{ vars.ntfy_prefix }}/{{ data.channel }}"
```
| Root | Meaning |
|---|---|
| `vars` | Owner variable; remaining segments are the flat name (`{{ vars.foo.bar }}` → variable `foo.bar`) |
| `secrets` | Same for secrets |
| `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.
YAML **SET** evaluates JSONata against the full `ctx`; the result is `output` (the next step’s data). `jsonata.js` does the same.
DAG `needs` assemble this step’s `data` from upstream **outputs**. Independent steps in the same wave share a context snapshot; sibling writes to the same context key fail the run.
@@ -105,8 +126,9 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` /
| `pnpm dev:server` | Monolith API/runner only |
| `pnpm dev:web` | UI only (proxies `/api` → :8700, `/ops` → :8600) |
| `pnpm build` | Production UI build |
| `pnpm start` | Monolith: API + worker + built UI |
| `pnpm start` | Monolith: API + worker + built UI (serves `dist` on :8700) |
| `pnpm start:control` | Control plane only (migrates, manages PM2 children) |
| `pnpm start:web` | Production UI on :8500 (`dist` + proxies to control/HTTP) |
| `pnpm start:api` | HTTP API + cron enqueue (`JFLOW_ROLE=api`) |
| `pnpm start:worker` | BullMQ worker only |
| `pnpm migrate` | Apply SQLite migrations |
@@ -140,24 +162,45 @@ Desired state is stored in `packages/server/data/control-state.json` (generation
| `JFLOW_ROLE` | `all` | `all` (HTTP + cron + worker), `api`, or `worker`. Prefer `pnpm start:api` / `start:worker` under control. |
| `JFLOW_CONFIG_GENERATION` | `1` | Set by control/PM2 so children report config generation in heartbeats. |
| `JFLOW_CONTROL_PORT` | `8600` | Control ops API port. |
| `JFLOW_UI_PORT` | `8500` | Production UI server (`web-server.js`) port. |
| `JFLOW_HTTP_PORT` | `8700` | HTTP API port (PM2 children / UI proxy target). |
| `JFLOW_LOG_LEVEL` | `debug` | Pino level |
| `JFLOW_RETENTION_DAYS` | `30` | Run history prune |
| `JFLOW_CORS_ORIGIN` | `http://localhost:8500` | Vite origin in dev |
| `PORT` | `8700` | HTTP API port |
| `NODE_ENV` | — | Set `production` for secure cookies |
| `JFLOW_CORS_ORIGIN` | `http://localhost:8500` | Browser origin (Vite in dev, UI server in prod) |
| `PORT` | `8700` | HTTP API port (alias; prefer `JFLOW_HTTP_PORT` under control) |
| `NODE_ENV` | — | Set `production` for secure cookies (unless overridden) |
| `COOKIE_SECURE` | (from `NODE_ENV`) | `true`/`false` — force Secure cookie flag. Use `false` for plain HTTP LAN access (`http://192.168.x.x`) |
Workflow runs are **queued** via BullMQ. HTTP and manual triggers return `202 { runId, status: "queued" }` immediately; poll `GET /api/runs/:id` for progress (`queued` → `running` → `success` \| `failed`). Cron remains an in-process producer that enqueues jobs on each tick.
## Production
Control-plane topology (same ports as `pnpm dev:pm2`):
| Process | Port | Role |
|---|---|---|
| `jflow-web` | **8500** | Built UI + proxies `/api` → :8700, `/ops` + `/api/auth` → :8600 |
| `jflow-control` | **8600** | Migrations, Ops API, starts/stops PM2 HTTP + workers |
| `jflow-http` | **8700** | API + cron enqueue (managed by control) |
| `jflow-worker` | — | BullMQ workers (managed by control) |
```bash
pnpm install
pnpm build
# Redis must be reachable at REDIS_URL (set REDIS_PASS if Redis requires AUTH)
# Recommended: run control (migrates + manages PM2 HTTP/workers)
JFLOW_JWT_SECRET=... JFLOW_SECRETS_KEY=... REDIS_URL=redis://127.0.0.1:6379 REDIS_PASS=... NODE_ENV=production pnpm start:control
# Or monolith (dev-style):
# ... pnpm start
# Put secrets in .env (JFLOW_JWT_SECRET, JFLOW_SECRETS_KEY, REDIS_URL, …)
# Use in-tree PM2 6.x (same module control.js requires). A global `pm2` 7.x
# against a 6.x daemon pegs CPU even when ls shows only 2 fork instances.
pnpm start:pm2
# UI: http://localhost:8500
# If you already mixed versions: pnpm pm2 -- kill && pnpm start:pm2
```
With control, serve the built UI from Vite preview, a reverse proxy, or set `JFLOW_SERVE_UI=1` on the HTTP process.
Or without the ecosystem file:
```bash
NODE_ENV=production pnpm start:control # :8600 + PM2 children
NODE_ENV=production pnpm start:web # :8500
```
Monolith (no Ops stop/scale): `pnpm build && pnpm start` serves the UI from the API process on :8700. Optional `JFLOW_SERVE_UI=1` on `start:api` does the same when you run HTTP alone — do **not** use that under control-plane mode (stopping HTTP would take down the UI).
+36 -9
View File
@@ -1,7 +1,12 @@
/**
* Production PM2 ecosystem (monolith runner).
* Use for deployed/single-process starts. For local multi-process Ops UI, use
* `pnpm dev:pm2` → packages/server/ecosystem.dev.cjs / control.js instead.
* Production PM2 ecosystem (control plane + UI).
* Starts always-on processes only; HTTP (:8700) and workers are owned by
* control.js via PM2 (same as `pnpm dev:pm2`).
*
* Use `pnpm start:pm2` (in-tree PM2 6.x). Do not use a global `pm2` 7.x —
* a CLI/daemon version mismatch pegs CPU even with instances: 1.
*
* Prerequisites: `pnpm build` (packages/web/dist), Redis, .env secrets.
*/
const fs = require("fs");
const path = require("path");
@@ -28,21 +33,43 @@ function loadEnv(file) {
}
const root = __dirname;
const env = {
NODE_ENV: "production",
...loadEnv(path.join(root, ".env")),
};
module.exports = {
apps: [
{
name: "jerapah-flow",
name: "jflow-control",
cwd: root,
script: "packages/server/runner.js",
interpreter: "node",
script: "packages/server/control.js",
interpreter: process.execPath,
instances: 1,
exec_mode: "fork",
autorestart: true,
max_restarts: 8,
restart_delay: 2000,
env: {
...env,
JFLOW_CONTROL_PORT: env.JFLOW_CONTROL_PORT ?? "8600",
},
},
{
name: "jflow-web",
cwd: root,
script: "packages/server/web-server.js",
interpreter: process.execPath,
instances: 1,
exec_mode: "fork",
autorestart: true,
max_restarts: 20,
env: {
NODE_ENV: "production",
...loadEnv(path.join(root, ".env")),
...env,
JFLOW_UI_PORT: env.JFLOW_UI_PORT ?? "8500",
JFLOW_CONTROL_PORT: env.JFLOW_CONTROL_PORT ?? "8600",
JFLOW_HTTP_PORT: env.JFLOW_HTTP_PORT ?? "8700",
},
},
],
};
};
@@ -1,6 +1,6 @@
name: Comic - monkeyuser to ntfy
description: |
Scrape the latest MonkeyUser comic and send it to ntfy (requires $VAR_ntfy_channel).
Scrape the latest MonkeyUser comic and send it to ntfy (requires {{ vars.ntfy_channel }}).
scripts:
- script: fetch-html.js
config:
@@ -20,7 +20,7 @@ scripts:
- script: fetch-binary.js
- script: ntfy.js
config:
url: $VAR_ntfy_channel
url: "{{ vars.ntfy_channel }}"
triggers:
- type: HTTP
method: POST
+3
View File
@@ -14,6 +14,9 @@
"start:api": "pnpm --filter @jerapah-flow/server start:api",
"start:worker": "pnpm --filter @jerapah-flow/server start:worker",
"start:control": "pnpm --filter @jerapah-flow/server start:control",
"start:web": "pnpm --filter @jerapah-flow/server start:web",
"pm2": "node scripts/pm2.mjs",
"start:pm2": "node scripts/pm2.mjs start ecosystem.config.cjs",
"migrate": "pnpm --filter @jerapah-flow/server migrate",
"test": "pnpm --filter @jerapah-flow/shared test && pnpm --filter @jerapah-flow/web test",
"lint": "pnpm --filter @jerapah-flow/web lint"
+17
View File
@@ -0,0 +1,17 @@
/**
* 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 };
}
+182 -66
View File
@@ -3,36 +3,49 @@ import { assertSecretName, getSecretPlaintext } from "./secrets-store.js";
import { isSecret } from "./secret-value.js";
import { assertVariableName, getVariablePlain } from "./variables-store.js";
const PREFIXES = [
{ kind: "context", prefix: "$CONTEXT_" },
{ kind: "secret", prefix: "$SECRET_" },
{ kind: "var", prefix: "$VAR_" },
];
const FORBIDDEN_SEGMENTS = new Set(["__proto__", "constructor", "prototype"]);
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 {{ kind: "secret" | "context" | "var", name: string, raw: string }} ConfigRef
* @typedef {{ owner: string, workflowKey: string, context?: unknown }} ConfigRefCtx
* @typedef {{
* owner: string,
* workflowKey?: string,
* context?: unknown,
* data?: unknown,
* }} ConfigRefCtx
*/
/**
* Parse a whole-value config placeholder. Returns null for literals.
* Detect leftover `$VAR_` / `$SECRET_` / `$CONTEXT_` whole-value refs.
* @param {unknown} value
* @returns {ConfigRef | null}
* @returns {{ kind: "var" | "secret" | "context", name: string, raw: string } | null}
*/
export function parseConfigRef(value) {
export function parseLegacyConfigRef(value) {
if (typeof value !== "string") return null;
const trimmed = value.trim();
for (const { kind, prefix } of PREFIXES) {
if (trimmed.startsWith(prefix)) {
return { kind, name: trimmed.slice(prefix.length), raw: trimmed };
}
}
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() };
}
/**
* Walk config (objects/arrays) and replace whole-value `$SECRET_` / `$CONTEXT_` / `$VAR_`
* strings. Does not walk trigger data.
* @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.
*
* @param {unknown} value
* @param {ConfigRefCtx} ctx
@@ -68,83 +81,186 @@ export async function resolveConfigRefs(value, ctx, seen = new WeakSet()) {
/**
* @param {string} value
* @param {ConfigRefCtx} ctx
* @returns {Promise<string | number | boolean>}
* @returns {Promise<unknown>}
*/
async function resolveStringRef(value, ctx) {
const ref = parseConfigRef(value);
if (!ref) return value;
const legacy = parseLegacyConfigRef(value);
if (legacy) {
throw new Error(
`config ref ${legacy.raw}: removed; ${legacyRenameHint(legacy.kind, legacy.name)}`,
);
}
if (ref.kind === "secret") {
return resolveSecretRef(ref, ctx);
const whole = WHOLE_MUSTACHE_RE.exec(value);
if (whole && whole[0] === value) {
return resolvePath(whole[1], ctx, { raw: value, allowObject: true });
}
if (ref.kind === "var") {
return resolveVarRef(ref, ctx);
if (!value.includes("{{")) {
return value;
}
return resolveContextRef(ref, ctx);
MUSTACHE_TOKEN_RE.lastIndex = 0;
let out = "";
let lastIndex = 0;
let match;
while ((match = MUSTACHE_TOKEN_RE.exec(value)) != null) {
out += value.slice(lastIndex, match.index);
const resolved = await resolvePath(match[1], ctx, {
raw: match[0],
allowObject: false,
});
out += stringifyScalar(resolved, match[0]);
lastIndex = match.index + match[0].length;
}
out += value.slice(lastIndex);
return out;
}
/**
* @param {ConfigRef} ref
* @param {string} pathExpr
* @param {ConfigRefCtx} ctx
* @returns {Promise<string>}
* @param {{ raw: string, allowObject: boolean }} opts
*/
async function resolveSecretRef(ref, ctx) {
try {
assertSecretName(ref.name);
} catch {
throw new Error(`config ref ${ref.raw}: invalid secret name`);
async function resolvePath(pathExpr, ctx, opts) {
const segments = pathExpr.split(".");
if (segments.length === 0 || segments.some((s) => !s)) {
throw new Error(`config ref ${opts.raw}: empty path`);
}
const plaintext = await getSecretPlaintext(ctx.owner, ref.name);
if (plaintext == null) {
throw new Error(`config ref ${ref.raw}: secret "${ref.name}" not found`);
for (const seg of segments) {
if (FORBIDDEN_SEGMENTS.has(seg)) {
throw new Error(`config ref ${opts.raw}: forbidden path segment "${seg}"`);
}
}
return plaintext;
const root = segments[0];
const rest = segments.slice(1);
if (root === "vars") {
return resolveNamedStore("var", rest, ctx, opts);
}
if (root === "secrets") {
return resolveNamedStore("secret", rest, ctx, opts);
}
if (root === "context") {
return walkObject(ctx.context, rest, opts);
}
if (root === "data") {
return walkObject(ctx.data, rest, opts);
}
throw new Error(
`config ref ${opts.raw}: unknown root "${root}" (use vars, secrets, context, or data)`,
);
}
/**
* @param {ConfigRef} ref
* @param {"var" | "secret"} kind
* @param {string[]} rest
* @param {ConfigRefCtx} ctx
* @returns {Promise<string | number | boolean>}
* @param {{ raw: string, allowObject: boolean }} opts
*/
async function resolveVarRef(ref, ctx) {
if (ref.name.length === 0) {
throw new Error(`config ref ${ref.raw}: empty variable name`);
async function resolveNamedStore(kind, rest, ctx, opts) {
if (rest.length === 0) {
throw new Error(`config ref ${opts.raw}: empty ${kind} name`);
}
const name = rest.join(".");
try {
assertVariableName(ref.name);
if (kind === "secret") assertSecretName(name);
else assertVariableName(name);
} catch {
throw new Error(`config ref ${ref.raw}: invalid variable name`);
throw new Error(`config ref ${opts.raw}: invalid ${kind} name`);
}
const value = await getVariablePlain(ctx.owner, ref.name);
if (kind === "secret") {
const plaintext = await getSecretPlaintext(ctx.owner, name);
if (plaintext == null) {
throw new Error(`config ref ${opts.raw}: secret "${name}" not found`);
}
return plaintext;
}
const value = await getVariablePlain(ctx.owner, name);
if (value == null) {
throw new Error(`config ref ${ref.raw}: variable "${ref.name}" not found`);
throw new Error(`config ref ${opts.raw}: variable "${name}" not found`);
}
return value;
}
/**
* @param {ConfigRef} ref
* @param {ConfigRefCtx} ctx
* @returns {string}
* @param {unknown} root
* @param {string[]} rest
* @param {{ raw: string, allowObject: boolean }} opts
*/
function resolveContextRef(ref, ctx) {
if (ref.name.length === 0) {
throw new Error(`config ref ${ref.raw}: empty context key`);
function walkObject(root, rest, opts) {
if (rest.length === 0) {
return unwrapValue(root, opts);
}
const bag =
ctx.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)
? /** @type {Record<string, unknown>} */ (ctx.context)
: {};
if (!Object.prototype.hasOwnProperty.call(bag, ref.name)) {
throw new Error(`config ref ${ref.raw}: context "${ref.name}" not found`);
let cur = root;
for (const seg of rest) {
if (cur == null || typeof cur !== "object") {
throw new Error(`config ref ${opts.raw}: path not found`);
}
if (Array.isArray(cur)) {
if (!/^\d+$/.test(seg)) {
throw new Error(`config ref ${opts.raw}: path not found`);
}
const idx = Number(seg);
if (!Number.isInteger(idx) || idx < 0 || idx >= cur.length) {
throw new Error(`config ref ${opts.raw}: path not found`);
}
cur = cur[idx];
continue;
}
const bag = /** @type {Record<string, unknown>} */ (cur);
if (!Object.prototype.hasOwnProperty.call(bag, seg)) {
throw new Error(`config ref ${opts.raw}: path not found`);
}
cur = bag[seg];
}
const raw = bag[ref.name];
if (isSecret(raw)) {
return raw.reveal();
return unwrapValue(cur, opts);
}
/**
* @param {unknown} value
* @param {{ raw: string, allowObject: boolean }} opts
*/
function unwrapValue(value, opts) {
if (isSecret(value)) {
return value.reveal();
}
const coerced = coerceCredentialString(raw);
if (coerced == null) {
throw new Error(`config ref ${ref.raw}: context "${ref.name}" is not a scalar`);
if (!opts.allowObject && value != null && typeof value === "object") {
throw new Error(`config ref ${opts.raw}: value is not a scalar`);
}
return coerced;
// Whole-value context/data may be any JSON type; mixed strings need scalars only.
if (opts.allowObject) {
if (value != null && typeof value === "object") return value;
if (
typeof value === "string" ||
typeof value === "number" ||
typeof value === "boolean"
) {
return value;
}
// Prefer credential coercion for odd primitives (e.g. bigint) when whole-value.
const coerced = coerceCredentialString(value);
if (coerced != null) return coerced;
return value;
}
return value;
}
/**
* @param {unknown} value
* @param {string} raw
*/
function stringifyScalar(value, raw) {
if (value == null) {
throw new Error(`config ref ${raw}: value is null`);
}
if (typeof value === "string") return value;
if (typeof value === "number" || typeof value === "boolean") {
return String(value);
}
throw new Error(`config ref ${raw}: value is not a scalar`);
}
+21 -18
View File
@@ -64,13 +64,6 @@ try {
process.exit(1);
}
try {
await connectPm2();
} catch (err) {
log.error({ err }, "failed to connect to PM2 — is pm2 installed?");
process.exit(1);
}
async function applyDesiredState() {
const state = readControlState();
await ensureHttp({
@@ -98,8 +91,6 @@ async function applyDesiredState() {
);
}
await applyDesiredState();
const server = fastify({ loggerInstance: log });
await server.register(cookie);
await server.register(jwt, {
@@ -484,12 +475,24 @@ process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);
const port = Number(process.env.JFLOW_CONTROL_PORT ?? process.env.PORT ?? 8600);
server
.listen({ host: "0.0.0.0", port })
.then(() => {
log.info(`Control is running on port ${port}`);
})
.catch((err) => {
log.error({ err }, "failed to start control");
process.exit(1);
});
try {
await server.listen({ host: "0.0.0.0", port });
log.info(`Control is running on port ${port}`);
} catch (err) {
log.error({ err }, "failed to start control");
process.exit(1);
}
try {
await connectPm2();
} catch (err) {
log.error({ err }, "failed to connect to PM2 — is pm2 installed?");
process.exit(1);
}
try {
await applyDesiredState();
} catch (err) {
log.error({ err }, "failed to apply desired state");
process.exit(1);
}
@@ -0,0 +1,16 @@
/**
* 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.
}
+97
View File
@@ -0,0 +1,97 @@
const HOP_BY_HOP = new Set([
"connection",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"te",
"trailer",
"transfer-encoding",
"upgrade",
"host",
"content-length",
]);
const REPLY_SKIP = new Set([
"connection",
"keep-alive",
"transfer-encoding",
"content-encoding",
"content-length",
]);
/**
* Control origin used when the UI or HTTP process proxies to control.
* @returns {string}
*/
export function controlOrigin() {
const explicit = process.env.JFLOW_CONTROL_URL?.trim();
if (explicit) return explicit.replace(/\/$/, "");
const port = Number(process.env.JFLOW_CONTROL_PORT ?? 8600);
return `http://127.0.0.1:${port}`;
}
/**
* HTTP API origin (workflow triggers + REST).
* @returns {string}
*/
export function httpOrigin() {
const explicit = process.env.JFLOW_HTTP_URL?.trim();
if (explicit) return explicit.replace(/\/$/, "");
const port = Number(process.env.JFLOW_HTTP_PORT ?? process.env.PORT ?? 8700);
return `http://127.0.0.1:${port}`;
}
/**
* Forward the incoming request to `origin`, preserving path + query.
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
* @param {string} origin
* @param {{ unreachableMessage?: string }} [opts]
*/
export async function proxyToOrigin(req, reply, origin, opts = {}) {
const target = `${origin.replace(/\/$/, "")}${req.raw.url ?? "/"}`;
const headers = {};
for (const [key, value] of Object.entries(req.headers)) {
if (value == null || HOP_BY_HOP.has(key.toLowerCase())) continue;
headers[key] = Array.isArray(value) ? value.join(", ") : String(value);
}
const method = req.method.toUpperCase();
const hasBody = method !== "GET" && method !== "HEAD";
let body;
if (hasBody) {
if (Buffer.isBuffer(req.body)) body = req.body;
else if (typeof req.body === "string") body = req.body;
else if (req.body != null) {
body = JSON.stringify(req.body);
if (!headers["content-type"]) headers["content-type"] = "application/json";
}
}
let res;
try {
res = await fetch(target, { method, headers, body });
} catch (err) {
const message = opts.unreachableMessage ?? "upstream unreachable";
req.log.warn({ err, target }, `proxy: ${message}`);
return reply.code(502).send({ error: message });
}
reply.code(res.status);
res.headers.forEach((value, key) => {
if (REPLY_SKIP.has(key.toLowerCase())) return;
reply.header(key, value);
});
return reply.send(Buffer.from(await res.arrayBuffer()));
}
/**
* Forward `/ops/*` to the control plane (same-origin UI in production).
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
*/
export async function proxyOpsToControl(req, reply) {
return proxyToOrigin(req, reply, controlOrigin(), {
unreachableMessage: "control plane unreachable",
});
}
+268
View File
@@ -0,0 +1,268 @@
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;
}
+5 -2
View File
@@ -11,12 +11,15 @@
"start:api": "node server.js",
"start:worker": "node worker.js",
"start:control": "node control.js",
"start:web": "node web-server.js",
"migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\"",
"test:plugins": "JFLOW_PLUGINS_DIR=./data/plugins-smoke-test JFLOW_DB_PATH=./data/plugins-smoke.db node test/plugins-smoke.js",
"test:workflow-history": "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",
"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:owner-migrate": "node test/owner-migrate-smoke.js"
},
"dependencies": {
"@jerapah-flow/shared": "workspace:*",
@@ -41,7 +44,7 @@
"nodemailer": "^9.0.5",
"pino": "^10.3.1",
"pino-roll": "^4.0.0",
"pm2": "^6.0.13",
"pm2": "6.0.14",
"rss-parser": "^3.13.0",
"ssh2-sftp-client": "^12.1.1",
"webdav": "^5.10.0",
+45 -1
View File
@@ -164,13 +164,57 @@ export async function restartPm2Process(pmId) {
return { name: proc.name, pmId: id };
}
/**
* PM2 injects these into process.env of a managed app. Spreading them into
* `pm2.start({ env })` overwrites `name` / `pm_exec_path` so God restarts
* jflow-control instead of launching http/worker (EADDRINUSE :8600 loop).
*/
const PM2_META_KEYS = new Set([
"name",
"namespace",
"exec_mode",
"exec_interpreter",
"instances",
"instance_var",
"node_app_instance",
"unique_id",
"status",
"username",
"windowsHide",
"merge_logs",
"vizion",
"vizion_running",
"autostart",
"autorestart",
"automation",
"km_link",
]);
/**
* @param {NodeJS.ProcessEnv} env
* @returns {NodeJS.ProcessEnv}
*/
export function withoutPm2Meta(env) {
/** @type {NodeJS.ProcessEnv} */
const out = {};
for (const [key, val] of Object.entries(env)) {
if (val == null) continue;
if (PM2_META_KEYS.has(key)) continue;
if (key.startsWith("pm_") || key.startsWith("axm_") || key.startsWith("PM2_")) {
continue;
}
out[key] = val;
}
return out;
}
/**
* Shared env for child processes.
* @param {{ generation: number }} opts
*/
export function childEnv(opts) {
return {
...process.env,
...withoutPm2Meta(process.env),
JFLOW_CONFIG_GENERATION: String(opts.generation),
JFLOW_CORS_ORIGIN: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:8500",
PORT: process.env.JFLOW_HTTP_PORT ?? "8700",
+54 -1
View File
@@ -36,7 +36,9 @@ import {
resolveFailureTriggerConfig,
} from "./trigger-failure.js";
import { enqueueWorkflowJob } from "./workflow-queue.js";
import { ensureInitialRevision } from "./workflow-history.js";
import { ensureInitialRevision, recordRevision } from "./workflow-history.js";
import { workflowIdFromFile } from "./workflow-normalize.js";
import { publishReload } from "./control-bus.js";
/**
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
@@ -509,6 +511,7 @@ export function createRegistry(server, opts = {}) {
owner,
workflowKey: key,
context: incomingContext,
data: ctx.data,
});
const stepCtx = {
data: ctx.data,
@@ -553,6 +556,45 @@ export function createRegistry(server, opts = {}) {
}
}
/**
* Persist `enabled: false` for a workflow and reload registries across processes.
* @param {string} owner
* @param {string} file
* @param {string} key
*/
async function disableWorkflowForConsecutiveFailures(owner, file, key) {
const content = fsStore.readWorkflowYaml(owner, file);
if (content == null) {
throw new Error(`workflow file missing for ${key}`);
}
const doc = yaml.parseDocument(content);
if (doc.errors?.length) {
throw new Error(doc.errors[0]?.message ?? "invalid yaml");
}
const parsed = doc.toJSON();
if (parsed?.enabled === false) {
log.debug({ workflow: key }, "workflow already disabled");
return;
}
doc.set("enabled", false);
const nextContent = String(doc);
fsStore.writeWorkflowYaml(owner, file, nextContent);
await recordRevision({
workflowId: workflowIdFromFile(file),
owner,
file,
content: nextContent,
reason: "disable-on-consecutive-failures",
});
reregister();
try {
await publishReload({ type: "workflows" });
} catch {
// Redis may be briefly unavailable; local reload already applied.
}
log.warn({ workflow: key }, "disabled workflow after consecutive failures");
}
/**
* @param {{
* key: string,
@@ -589,6 +631,17 @@ export function createRegistry(server, opts = {}) {
return;
}
if (failureConfig.disableOnConsecutiveFailures) {
const entry = workflows.get(opts.key);
if (entry) {
await disableWorkflowForConsecutiveFailures(entry.owner, entry.file, opts.key);
} else {
log.warn({ workflow: opts.key }, "cannot disable missing workflow entry");
}
}
if (!failureConfig.workflowName) return;
const destKey = resolveWorkflowTriggerKey(opts.owner, failureConfig.workflowName);
const alertData = buildFailureAlertData({
sourceKey: opts.key,
+3 -2
View File
@@ -11,6 +11,7 @@ import { isSecret, Secret, unwrapSecretsDeep } from "./secret-value.js";
import { getHttpPageByName, getHttpTemplateByName } from "./http-pages-store.js";
import { getSecretPlaintext } from "./secrets-store.js";
import { getVariablePlain } from "./variables-store.js";
import { DEFAULT_OWNER } from "@jerapah-flow/shared";
const hostRequire = createRequire(import.meta.url);
@@ -443,7 +444,7 @@ function createScriptSandbox({
log,
script,
workflowName,
owner = "default",
owner = DEFAULT_OWNER,
$workflows = $workflowsStub,
pluginDir = null,
}) {
@@ -563,7 +564,7 @@ export function instantiateScriptSource(script, source, opts = {}) {
log: opts.log ?? inspectLog,
script,
workflowName: opts.workflowName ?? "inspect",
owner: opts.owner ?? "default",
owner: opts.owner ?? DEFAULT_OWNER,
$workflows: opts.$workflows,
pluginDir: opts.pluginDir ?? null,
});
+4 -1
View File
@@ -31,6 +31,7 @@ import { getAppVersion } from "../../app-version.js";
import { EXAMPLE_PLUGINS_DIR } from "../../paths.js";
import { pluginScriptRef } from "../../plugin-manifest.js";
import { evaluateJsonata, SET_STEP_SCRIPT } from "../../workflow-parse.js";
import { DEFAULT_OWNER } from "@jerapah-flow/shared";
/**
* @param {{ referencedScripts: () => Set<string> }} registry
@@ -258,7 +259,7 @@ export default function scriptsPluginFactory(registry) {
req.body ?? {}
);
let owner = "local";
let owner = DEFAULT_OWNER;
if (body.owner != null && body.owner !== "") {
try {
owner = fsStore.assertOwner(String(body.owner));
@@ -304,6 +305,7 @@ export default function scriptsPluginFactory(registry) {
owner,
workflowKey: "dry-run",
context: incomingContext,
data: incomingData,
},
);
const ctx = {
@@ -365,6 +367,7 @@ export default function scriptsPluginFactory(registry) {
owner,
workflowKey: "dry-run",
context: incomingContext,
data: incomingData,
});
const ctx = {
data: incomingData,
+8 -8
View File
@@ -76,6 +76,7 @@ function triggerSummary(owner, workflow, nameById) {
schedule: t?.schedule ?? null,
onConsecutiveFailures: t?.onConsecutiveFailures ?? null,
onFailureWorkflow: t?.onFailureWorkflow ?? null,
disableOnConsecutiveFailures: t?.disableOnConsecutiveFailures === true,
auth: isHttp ? authLabel(t?.auth, nameById) : null,
};
});
@@ -705,15 +706,14 @@ export default function workflowsPluginFactory(registry) {
return reply.code(err.statusCode ?? 400).send({ error: err.message });
}
const raw = fsStore.readWorkflowYaml(owner, file);
if (raw == null) {
return reply.code(404).send({ error: "workflow not found" });
}
let name = null;
try {
const parsed = yaml.parse(raw);
name = parsed?.name ?? null;
} catch {
// ignore
if (raw != null) {
try {
const parsed = yaml.parse(raw);
name = parsed?.name ?? null;
} catch {
// ignore
}
}
try {
const item = await moveWorkflowToTrash({
+7
View File
@@ -32,6 +32,7 @@ import {
} 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,
@@ -135,6 +136,12 @@ export async function startApp(opts = {}) {
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.
+140 -106
View File
@@ -2,7 +2,8 @@ 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 { parseConfigRef, resolveConfigRefs } from "../config-refs.js";
import { parseLegacyConfigRef, resolveConfigRefs } from "../config-refs.js";
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
await migrate();
@@ -23,126 +24,171 @@ async function assertRejects(fn, match) {
throw new Error(`expected to reject (${match ?? "any error"})`);
}
const owner = "default";
const workflowKey = "default/config-refs-smoke.yaml";
const ctx = { owner, workflowKey, context: {} };
const owner = "config_refs_smoke_owner";
const ctx = { owner, workflowKey: `${owner}/config-refs-smoke.yaml`, context: {}, data: {} };
function assertParse(value, expected) {
const got = parseConfigRef(value);
if (expected == null) {
assert(got == null, `expected null parse for ${JSON.stringify(value)}, got ${JSON.stringify(got)}`);
return;
}
assert(got != null, `expected parse for ${JSON.stringify(value)}`);
assert(got.kind === expected.kind, `kind ${got.kind} !== ${expected.kind}`);
assert(got.name === expected.name, `name ${JSON.stringify(got.name)} !== ${JSON.stringify(expected.name)}`);
{
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");
}
assertParse("password123", null);
assertParse("$FOO_bar", null);
assertParse("$SECRET", null);
assertParse(" password123 ", null);
assertParse("$SECRET_zte_modem_password", { kind: "secret", name: "zte_modem_password" });
assertParse(" $SECRET_zte_modem_password ", { kind: "secret", name: "zte_modem_password" });
assertParse("Bearer $SECRET_x", null);
assertParse("$KV_modem password", null);
assertParse("$VAR_ntfy_url", { kind: "var", name: "ntfy_url" });
assertParse("$CONTEXT_token", { kind: "context", name: "token" });
assertParse("$SECRET_", { kind: "secret", name: "" });
assertParse("$CONTEXT_SECRET_foo", { kind: "context", name: "SECRET_foo" });
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 kvLiteral = await resolveConfigRefs("$KV_modem_password", ctx);
assert(kvLiteral === "$KV_modem_password", "$KV_ stays literal");
const embedded = await resolveConfigRefs("Bearer $SECRET_x", ctx);
assert(embedded === "Bearer $SECRET_x", "mid-string stays literal");
assert(embedded === "Bearer $SECRET_x", "mid-string legacy stays literal");
const number = await resolveConfigRefs(42, ctx);
assert(number === 42, "number passthrough");
}
const secret = await upsertSecret({
owner,
name: "config_refs_smoke_token",
value: "s3cret-ok",
});
const varUrl = await upsertVariable({
owner,
name: "config_refs_smoke_url",
type: "string",
value: "https://example.test",
});
const varRetry = await upsertVariable({
owner,
name: "config_refs_smoke_retry",
type: "number",
value: 3,
});
const varDebug = await upsertVariable({
owner,
name: "config_refs_smoke_debug",
type: "boolean",
value: false,
});
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({
owner,
name: "config_refs_smoke_token",
value: "s3cret-ok",
});
created.push(["secret", secret.id]);
const varUrl = await upsertVariable({
owner,
name: "config_refs_smoke_url",
type: "string",
value: "https://example.test",
});
created.push(["var", varUrl.id]);
const varRetry = await upsertVariable({
owner,
name: "config_refs_smoke_retry",
type: "number",
value: 3,
});
created.push(["var", varRetry.id]);
const varDebug = await upsertVariable({
owner,
name: "config_refs_smoke_debug",
type: "boolean",
value: false,
});
created.push(["var", varDebug.id]);
{
const resolved = await resolveConfigRefs("$SECRET_config_refs_smoke_token", ctx);
assert(resolved === "s3cret-ok", "secret resolve");
const resolved = await resolveConfigRefs("{{ secrets.config_refs_smoke_token }}", ctx);
assert(resolved === "s3cret-ok", "secret whole-value");
}
{
const resolved = await resolveConfigRefs(" $SECRET_config_refs_smoke_token ", ctx);
assert(resolved === "s3cret-ok", "secret resolve trimmed");
}
{
const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_url", ctx);
const resolved = await resolveConfigRefs("{{ vars.config_refs_smoke_url }}", ctx);
assert(resolved === "https://example.test", "var string");
}
{
const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_retry", ctx);
assert(resolved === 3, "var number stays number");
assert(typeof resolved === "number", "var number type");
const resolved = await resolveConfigRefs("{{ vars.config_refs_smoke_retry }}", ctx);
assert(resolved === 3, "var number keeps type");
}
{
const resolved = await resolveConfigRefs("$VAR_config_refs_smoke_debug", ctx);
assert(resolved === false, "var boolean stays false");
assert(typeof resolved === "boolean", "var boolean type");
const resolved = await resolveConfigRefs("{{ vars.config_refs_smoke_debug }}", ctx);
assert(resolved === false, "var boolean keeps type");
}
{
const resolved = await resolveConfigRefs("$CONTEXT_token", {
const resolved = await resolveConfigRefs("{{ context.token }}", {
...ctx,
context: { token: "ctx-token-ok" },
});
assert(resolved === "ctx-token-ok", "context string");
}
{
const resolved = await resolveConfigRefs("$CONTEXT_n", {
const resolved = await resolveConfigRefs("{{ context.n }}", {
...ctx,
context: { n: 7 },
});
assert(resolved === "7", "context number stringify");
assert(resolved === 7, "context number keeps type");
}
{
const wrapped = new Secret("wrapped-secret-ok");
const resolved = await resolveConfigRefs("$CONTEXT_tok", {
const resolved = await resolveConfigRefs("{{ context.tok }}", {
...ctx,
context: { tok: wrapped },
});
assert(resolved === "wrapped-secret-ok", "context Secret unwrap");
}
{
const resolved = await resolveConfigRefs("{{ context.user }}", {
...ctx,
context: { user: { id: "u1", role: "admin" } },
});
assert(
resolved && typeof resolved === "object" && resolved.id === "u1",
"whole-value object pass-through",
);
}
{
const resolved = await resolveConfigRefs("{{ context.user.id }}", {
...ctx,
context: { user: { id: "nested-id" } },
});
assert(resolved === "nested-id", "nested context path");
}
{
const resolved = await resolveConfigRefs("{{ data.items.0.id }}", {
...ctx,
data: { items: [{ id: "row-0" }] },
});
assert(resolved === "row-0", "array index path");
}
{
const resolved = await resolveConfigRefs(
"{{ vars.config_refs_smoke_url }}/{{ data.channel }}",
{ ...ctx, data: { channel: "alerts" } },
);
assert(resolved === "https://example.test/alerts", "concatenation");
}
{
const resolved = await resolveConfigRefs("Bearer {{ context.token }}", {
...ctx,
context: { token: "abc" },
});
assert(resolved === "Bearer abc", "mixed string");
}
{
const nested = await resolveConfigRefs(
{
url: "$VAR_config_refs_smoke_url",
retry: "$VAR_config_refs_smoke_retry",
debug: "$VAR_config_refs_smoke_debug",
password: "$SECRET_config_refs_smoke_token",
url: "{{ vars.config_refs_smoke_url }}",
retry: "{{ vars.config_refs_smoke_retry }}",
debug: "{{ vars.config_refs_smoke_debug }}",
password: "{{ secrets.config_refs_smoke_token }}",
headers: { Authorization: "$KV_modem_password" },
extra: ["$CONTEXT_token", "plain"],
extra: ["{{ context.token }}", "plain"],
},
{ ...ctx, context: { token: "ctx-token-ok" } },
);
@@ -155,59 +201,47 @@ try {
assert(nested.extra[1] === "plain", "nested array literal");
}
const data = { password: "$SECRET_config_refs_smoke_token" };
const config = { password: "$SECRET_config_refs_smoke_token" };
const data = { password: "{{ secrets.config_refs_smoke_token }}" };
const config = { password: "{{ secrets.config_refs_smoke_token }}" };
const resolvedConfig = await resolveConfigRefs(config, ctx);
assert(resolvedConfig.password === "s3cret-ok", "config resolved");
assert(data.password === "$SECRET_config_refs_smoke_token", "data not walked");
assert(config.password === "$SECRET_config_refs_smoke_token", "input config not mutated");
assert(data.password === "{{ secrets.config_refs_smoke_token }}", "data not walked");
assert(config.password === "{{ secrets.config_refs_smoke_token }}", "input config not mutated");
await assertRejects(
() => resolveConfigRefs("$SECRET_does_not_exist_xyz", ctx),
() => resolveConfigRefs("{{ secrets.does_not_exist_xyz }}", ctx),
'secret "does_not_exist_xyz" not found',
);
await assertRejects(
() => resolveConfigRefs("$SECRET_not valid", ctx),
"invalid secret name",
);
await assertRejects(
() => resolveConfigRefs("$SECRET_", ctx),
"invalid secret name",
);
await assertRejects(
() => resolveConfigRefs("$CONTEXT_missing", ctx),
'context "missing" not found',
() => resolveConfigRefs("{{ context.missing }}", ctx),
"path not found",
);
await assertRejects(
() =>
resolveConfigRefs("$CONTEXT_obj", {
resolveConfigRefs("Bearer {{ context.obj }}", {
...ctx,
context: { obj: { a: 1 } },
}),
'context "obj" is not a scalar',
"not a scalar",
);
await assertRejects(
() => resolveConfigRefs("$CONTEXT_", ctx),
"empty context key",
() => resolveConfigRefs("{{ vars }}", ctx),
"empty var name",
);
await assertRejects(
() => resolveConfigRefs("$VAR_does_not_exist_xyz", ctx),
'variable "does_not_exist_xyz" not found',
() => resolveConfigRefs("{{ title }}", ctx),
"unknown root",
);
await assertRejects(
() => resolveConfigRefs("$VAR_not valid", ctx),
"invalid variable name",
);
await assertRejects(
() => resolveConfigRefs("$VAR_", ctx),
"empty variable name",
() => resolveConfigRefs("{{ context.__proto__.x }}", ctx),
"forbidden path segment",
);
} finally {
await deleteSecret(secret.id);
await deleteVariable(varUrl.id);
await deleteVariable(varRetry.id);
await deleteVariable(varDebug.id);
for (const [kind, id] of created.reverse()) {
if (kind === "secret") await deleteSecret(id);
else await deleteVariable(id);
}
await db.destroy();
}
console.log("config-refs smoke test passed");
await db.destroy();
@@ -0,0 +1,30 @@
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,
});
+2 -2
View File
@@ -71,11 +71,11 @@ const created = await upsertProfile({
owner,
name,
script: "ntfy.js",
config: { url: "$VAR_ntfy_channel" },
config: { url: "{{ vars.ntfy_channel }}" },
description: "smoke",
});
assert(created.name === name, "created");
assert(created.config.url === "$VAR_ntfy_channel", "config roundtrip");
assert(created.config.url === "{{ vars.ntfy_channel }}", "config roundtrip");
assert(created.script === "ntfy.js", "script locked on profile");
const fetched = await getProfilePlain(owner, name);
@@ -136,6 +136,18 @@ const warnings = collectWorkflowWarnings(
);
assert.ok(warnings.warnings.some((w) => w.code === "unknown_script"));
const ghostFile = newWorkflowFilename();
fsStore.writeRegisters(owner, [file, ghostFile]);
const ghostTrash = await moveWorkflowToTrash({
workflowId: workflowIdFromFile(ghostFile),
owner,
file: ghostFile,
name: null,
});
assert.equal(ghostTrash, null);
assert.ok(!fsStore.readRegisters(owner).includes(ghostFile));
assert.ok(fsStore.readRegisters(owner).includes(file));
const trashed = await moveWorkflowToTrash({
workflowId,
owner,
+19 -5
View File
@@ -47,12 +47,18 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) {
const threshold = Number(spec.onConsecutiveFailures);
const workflowName = onFailureWorkflowName(spec);
if (!Number.isFinite(threshold) || threshold < 1 || workflowName.length === 0) {
const disableOnConsecutiveFailures = isDisableOnConsecutiveFailures(spec);
if (
!Number.isFinite(threshold) ||
threshold < 1 ||
(workflowName.length === 0 && !disableOnConsecutiveFailures)
) {
return null;
}
return {
threshold: Math.floor(threshold),
workflowName,
workflowName: workflowName.length > 0 ? workflowName : null,
disableOnConsecutiveFailures,
};
}
return null;
@@ -66,6 +72,13 @@ function onFailureWorkflowName(trigger) {
return typeof value === "string" ? value.trim() : "";
}
/**
* @param {Record<string, unknown>} trigger
*/
export function isDisableOnConsecutiveFailures(trigger) {
return trigger?.disableOnConsecutiveFailures === true;
}
/**
* @param {unknown} workflow
*/
@@ -81,12 +94,13 @@ export async function validateWorkflowFailureTriggers(workflow) {
const hasThreshold =
trigger.onConsecutiveFailures != null && trigger.onConsecutiveFailures !== "";
const hasWorkflow = onFailureWorkflowName(trigger).length > 0;
const hasDisable = isDisableOnConsecutiveFailures(trigger);
if (!hasThreshold && !hasWorkflow) continue;
if (!hasThreshold && !hasWorkflow && !hasDisable) continue;
if (!hasThreshold || !hasWorkflow) {
if (!hasThreshold || (!hasWorkflow && !hasDisable)) {
const err = new Error(
"onConsecutiveFailures and onFailureWorkflow must both be set on a trigger",
"onConsecutiveFailures requires onFailureWorkflow and/or disableOnConsecutiveFailures",
);
err.statusCode = 400;
throw err;
+116
View File
@@ -0,0 +1,116 @@
/**
* Production UI server (:8500).
* Serves packages/web/dist and proxies /api, /ops, /admin, /u like Vite in dev.
* Always-on — survives Ops stop of jflow-http.
*/
import fs from "fs";
import fastify from "fastify";
import fastifyStatic from "@fastify/static";
import { log } from "./logger.js";
import { WEB_DIST } from "./paths.js";
import {
controlOrigin,
httpOrigin,
proxyToOrigin,
} from "./ops-proxy.js";
if (!fs.existsSync(WEB_DIST)) {
log.error(
{ WEB_DIST },
"web dist missing — run `pnpm build` before starting the UI server",
);
process.exit(1);
}
const port = Number(process.env.JFLOW_UI_PORT ?? 8500);
const control = controlOrigin();
const http = httpOrigin();
const server = fastify({ loggerInstance: log });
/**
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
*/
async function proxyApi(req, reply) {
const url = req.raw.url ?? "";
// Match Vite: /api/auth → control (login works when HTTP is stopped).
if (url === "/api/auth" || url.startsWith("/api/auth/") || url.startsWith("/api/auth?")) {
return proxyToOrigin(req, reply, control, {
unreachableMessage: "control plane unreachable",
});
}
return proxyToOrigin(req, reply, http, {
unreachableMessage: "HTTP API unreachable",
});
}
/**
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
*/
async function proxyOps(req, reply) {
return proxyToOrigin(req, reply, control, {
unreachableMessage: "control plane unreachable",
});
}
/**
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
*/
async function proxyHttp(req, reply) {
return proxyToOrigin(req, reply, http, {
unreachableMessage: "HTTP API unreachable",
});
}
server.all("/api", proxyApi);
server.all("/api/*", proxyApi);
server.all("/ops", proxyOps);
server.all("/ops/*", proxyOps);
server.all("/admin", proxyHttp);
server.all("/admin/*", proxyHttp);
server.all("/u", proxyHttp);
server.all("/u/*", proxyHttp);
await server.register(fastifyStatic, {
root: WEB_DIST,
wildcard: false,
});
server.setNotFoundHandler((req, reply) => {
const url = req.raw.url ?? "";
if (
url.startsWith("/api") ||
url.startsWith("/u/") ||
url.startsWith("/admin") ||
url.startsWith("/ops")
) {
return reply.code(404).send({ error: "not found" });
}
return reply.sendFile("index.html");
});
async function shutdown() {
try {
await server.close();
} catch (err) {
log.error({ err }, "web-server shutdown error");
}
process.exit(0);
}
process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);
try {
await server.listen({ host: "0.0.0.0", port });
log.info(
{ port, WEB_DIST, control, http },
"UI server listening (static + proxy)",
);
} catch (err) {
log.error({ err }, "failed to start UI server");
process.exit(1);
}
+9 -3
View File
@@ -1,4 +1,4 @@
import { HTTP_METHODS } from "@jerapah-flow/shared";
import { HTTP_METHODS, DEFAULT_OWNER } from "@jerapah-flow/shared";
import {
checkAnyHttpAuth,
resolveAuthMechanisms,
@@ -85,8 +85,14 @@ export function createHttpTriggerHandler({
return async function dispatchHttpTrigger(req, reply) {
const wildcard = /** @type {{ "*": string }} */ (req.params)["*"] ?? "";
const url = `/u/${String(wildcard).replace(/^\/+/, "")}`;
// Compat: leftover webhooks still hitting /u/default/... after owner rename.
const compatUrl = url.startsWith("/u/default/")
? `/u/${DEFAULT_OWNER}/${url.slice("/u/default/".length)}`
: url === "/u/default"
? `/u/${DEFAULT_OWNER}`
: url;
const method = String(req.method ?? "GET").toUpperCase();
const routeKey = `${method} ${url}`;
const routeKey = `${method} ${compatUrl}`;
const mapped = httpRoutes.get(routeKey);
if (!mapped) {
@@ -104,7 +110,7 @@ export function createHttpTriggerHandler({
if (t?.type !== "HTTP") return false;
const m = String(t.method ?? "POST").toUpperCase();
const p = namespacedPath(entry.owner, t.path);
return m === method && p === url;
return m === method && p === compatUrl;
}) ?? mapped.trigger;
if (
+26 -6
View File
@@ -62,29 +62,49 @@ export async function isInTrash(owner, file) {
return Boolean(row);
}
function unregisterWorkflow(owner, file) {
const registered = fsStore.readRegisters(owner).filter((f) => f !== file);
fsStore.writeRegisters(owner, registered);
}
/**
* Soft-delete: move YAML to trash dir, unregister, keep revision history.
* Soft-delete: unregister, move YAML to trash when present, keep revision history.
* Missing YAML (ghost register entries) is still unregistered so it leaves the list.
* @param {{
* workflowId: string,
* owner: string,
* file: string,
* name?: string | null,
* }} opts
* @returns {Promise<ReturnType<typeof rowToItem> | null>} trash item, or null if there was no file to keep
*/
export async function moveWorkflowToTrash(opts) {
unregisterWorkflow(opts.owner, opts.file);
const sourcePath = path.join(WORKFLOWS_DIR, opts.owner, opts.file);
if (!fs.existsSync(sourcePath)) {
const err = new Error("workflow not found");
err.statusCode = 404;
throw err;
const existing = await db("workflow_trash")
.where({ owner: opts.owner, file: opts.file })
.first();
return existing ? rowToItem(existing) : null;
}
const trashPath = trashFilePath(opts.owner, opts.file);
fs.mkdirSync(path.dirname(trashPath), { recursive: true });
fs.renameSync(sourcePath, trashPath);
const registered = fsStore.readRegisters(opts.owner).filter((f) => f !== opts.file);
fsStore.writeRegisters(opts.owner, registered);
const existing = await db("workflow_trash")
.where({ owner: opts.owner, file: opts.file })
.first();
if (existing) {
await db("workflow_trash").where({ id: existing.id }).update({
workflow_id: opts.workflowId,
name: opts.name ?? existing.name ?? null,
deleted_at: nowIso(),
trash_path: trashPath,
});
return rowToItem(await db("workflow_trash").where({ id: existing.id }).first());
}
const id = randomUUID();
const deleted_at = nowIso();
+1
View File
@@ -2,3 +2,4 @@ export { isPlainObject } from "./is-plain-object.js";
export { mergeProfileConfig, overlayFromMerged, configHasOverlay } from "./profile-config.js";
export { ensureWorkflowFilename, suggestCopyFilename } from "./workflow-filename.js";
export { HTTP_METHODS, namespacedPath, hasWorkflowTrigger } from "./workflow-path.js";
export { DEFAULT_OWNER } from "./tenant.js";
+2
View File
@@ -0,0 +1,2 @@
/** Default owner folder for new resources (latent tenant id). */
export const DEFAULT_OWNER = "local";
@@ -66,7 +66,8 @@ export function SecretEditorModal({ mode, initial, onClose, onSaved }) {
</label>
<p className="text-xs opacity-60">
Values are encrypted at rest and never shown again after save. Values shorter than 8
characters are not redacted from logs.
characters are not redacted from logs. In workflows use{" "}
<span className="font-mono">{"{{ secrets.name }}"}</span> (quote in YAML).
</p>
{upsert.isError ? (
<p className="text-error text-sm">{errorMessage(upsert.error)}</p>
@@ -129,7 +129,9 @@ export function VariableEditorModal({ mode, initial, onClose, onSaved }) {
</label>
<p className="text-xs opacity-60">
Stored in plaintext. Use Secrets for credentials. In workflows use{" "}
<span className="font-mono">$VAR_name</span> as a whole field.
<span className="font-mono">{"{{ vars.name }}"}</span> (quote strings that
start with <span className="font-mono">{"{"}</span>). Nested:{" "}
<span className="font-mono">{"{{ context.user.id }}"}</span>.
</p>
{formError ? <p className="text-error text-sm">{formError}</p> : null}
{upsert.isError ? (
@@ -4,22 +4,40 @@ import { LuEye, LuEyeOff, LuExternalLink } from "react-icons/lu";
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*$/;
const CONFIG_REF_PREFIXES = [
{ prefix: "$SECRET_", label: "secret" },
{ prefix: "$CONTEXT_", label: "context" },
{ prefix: "$VAR_", label: "variable" },
];
/**
* @param {unknown} value
* @returns {{
* kind: "mustache" | "legacy",
* path?: string,
* root?: string,
* name?: string,
* legacyKind?: string,
* legacyName?: string,
* } | null}
*/
function describeConfigRef(value) {
if (typeof value !== "string") return null;
const trimmed = value.trim();
for (const { prefix, label } of CONFIG_REF_PREFIXES) {
if (trimmed.startsWith(prefix) && trimmed.length > prefix.length) {
return { label, name: trimmed.slice(prefix.length), kind: prefix };
}
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] ?? "",
};
}
return null;
MUSTACHE_RE.lastIndex = 0;
const match = MUSTACHE_RE.exec(trimmed);
if (!match) return null;
const path = match[1];
const root = path.split(".")[0];
return { kind: "mustache", path, root, name: path.slice(root.length + 1) };
}
function formatVarDisplay(value) {
@@ -33,7 +51,7 @@ function truncatePeek(text, maxLen = VAR_PEEK_MAX) {
return `${text.slice(0, Math.max(0, maxLen - 1))}…`;
}
/** Edit-time lookup for `$VAR_` against workflow owner (not a runtime guarantee). */
/** Edit-time lookup for `{{ vars.name }}` against workflow owner (not a runtime guarantee). */
function lookupVariable(variables, owner, name) {
const list = Array.isArray(variables) ? variables : [];
const match = list.find((v) => v.owner === owner && v.name === name);
@@ -62,7 +80,7 @@ function variablesDeepLink({ owner, name, missing }) {
export function ConfigRefHint({ value, owner }) {
const ref = describeConfigRef(value);
const isVar = ref?.kind === "$VAR_";
const isVar = ref?.kind === "mustache" && ref.root === "vars" && Boolean(ref.name);
const { data: variables = [], isPending } = useVariables(undefined, { enabled: isVar });
const [revealed, setRevealed] = useState(false);
@@ -72,10 +90,32 @@ 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"
? "secret"
: ref.root === "context"
? "context"
: ref.root === "data"
? "data"
: "expression";
return (
<p className="text-xs opacity-60">
from {ref.label} <span className="font-mono">{ref.name}</span>
from {label} <span className="font-mono">{ref.path}</span>
</p>
);
}
@@ -143,10 +143,15 @@ function TriggerCardView({
function triggerSummary(trigger, owner) {
const type = trigger?.type;
const failure =
trigger.onConsecutiveFailures && trigger.onFailureWorkflow
? ` · onFailure@${trigger.onFailureWorkflow}`
: "";
/** @type {string[]} */
const failureParts = [];
if (trigger.onConsecutiveFailures && trigger.onFailureWorkflow) {
failureParts.push(`onFailure@${trigger.onFailureWorkflow}`);
}
if (trigger.disableOnConsecutiveFailures) {
failureParts.push("auto-disable");
}
const failure = failureParts.length ? ` · ${failureParts.join(", ")}` : "";
if (type === "HTTP") {
return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}${failure}`;
}
@@ -383,6 +388,20 @@ function FailureAlertFields({ trigger, disabled, onChange, alertDestinations })
) : null}
</FormSelect>
</Field>
<label className="label cursor-pointer justify-start gap-3 py-0">
<input
type="checkbox"
className="checkbox checkbox-sm"
checked={Boolean(trigger.disableOnConsecutiveFailures)}
disabled={disabled}
onChange={(e) =>
onChange({ ...trigger, disableOnConsecutiveFailures: e.target.checked })
}
/>
<span className="label-text">
Disable this workflow when the consecutive failure threshold is reached
</span>
</label>
</div>
);
}
+1 -2
View File
@@ -1,2 +1 @@
/** Default owner folder for new resources (latent tenant id). */
export const DEFAULT_OWNER = "local";
export { DEFAULT_OWNER } from "@jerapah-flow/shared";
+15 -1
View File
@@ -171,6 +171,7 @@ export function newHttpTrigger() {
unauthorized: null,
onConsecutiveFailures: "",
onFailureWorkflow: "",
disableOnConsecutiveFailures: false,
};
}
@@ -186,6 +187,7 @@ export function newCronTrigger() {
unauthorized: null,
onConsecutiveFailures: "",
onFailureWorkflow: "",
disableOnConsecutiveFailures: false,
};
}
@@ -201,6 +203,7 @@ export function newWorkflowTrigger() {
unauthorized: null,
onConsecutiveFailures: "",
onFailureWorkflow: "",
disableOnConsecutiveFailures: false,
};
}
@@ -303,9 +306,16 @@ function normalizeTrigger(raw) {
"unauthorized",
"onConsecutiveFailures",
"onFailureWorkflow",
"disableOnConsecutiveFailures",
])
: type === "cron"
? new Set(["type", "schedule", "onConsecutiveFailures", "onFailureWorkflow"])
? new Set([
"type",
"schedule",
"onConsecutiveFailures",
"onFailureWorkflow",
"disableOnConsecutiveFailures",
])
: new Set(["type"]);
/** @type {Record<string, unknown>} */
const extra = {};
@@ -323,6 +333,7 @@ function normalizeTrigger(raw) {
? ""
: String(raw.onConsecutiveFailures),
onFailureWorkflow: readOnFailureWorkflow(raw),
disableOnConsecutiveFailures: raw.disableOnConsecutiveFailures === true,
auth: Array.isArray(raw.auth) ? raw.auth : null,
response: typeof raw.response === "string" ? raw.response : "",
unauthorized: raw.unauthorized ?? null,
@@ -377,6 +388,9 @@ function dumpFailureTriggerFields(t, out) {
if (typeof t.onFailureWorkflow === "string" && t.onFailureWorkflow.trim()) {
out.onFailureWorkflow = t.onFailureWorkflow.trim();
}
if (t.disableOnConsecutiveFailures === true) {
out.disableOnConsecutiveFailures = true;
}
}
function dumpTrigger(t) {
+1 -1
View File
@@ -179,7 +179,7 @@ export function VariablesPage() {
title={confirmDelete ? `Delete ${confirmDelete.name}?` : ""}
message={
confirmDelete
? `This cannot be undone. Workflows that reference $VAR_${confirmDelete.name} will fail.`
? `This cannot be undone. Workflows that reference {{ vars.${confirmDelete.name} }} will fail.`
: ""
}
error={del.isError ? errorMessage(del.error) : null}
+2 -2
View File
@@ -198,7 +198,7 @@ joplinHttp.meta = {
token: {
type: "string",
required: true,
description: "Bearer token ($SECRET_); Bearer prefix is added if missing",
description: "Bearer token ({{ secrets.* }}); Bearer prefix is added if missing",
},
timeoutMs: {
type: "number",
@@ -237,7 +237,7 @@ joplinHttp.meta = {
config: {
url: "http://10.8.0.6:3030/notes",
method: "GET",
token: "$SECRET_joplin_api_token",
token: "{{ secrets.joplin_api_token }}",
timeoutMs: 60000,
},
},
+1 -1
View File
@@ -77,7 +77,7 @@ importers:
specifier: ^4.0.0
version: 4.0.0
pm2:
specifier: ^6.0.13
specifier: 6.0.14
version: 6.0.14(supports-color@7.2.0)
rss-parser:
specifier: ^3.13.0
+48
View File
@@ -0,0 +1,48 @@
#!/usr/bin/env node
/**
* Run the same PM2 binary control.js `require("pm2")` uses.
* A global `pm2` 7.x talking to an in-memory 6.x daemon pegs CPU on start.
*/
import { spawn, spawnSync } from "node:child_process";
import { createRequire } from "node:module";
import path from "node:path";
import { fileURLToPath } from "node:url";
const root = path.resolve(path.dirname(fileURLToPath(import.meta.url)), "..");
const require = createRequire(
path.join(root, "packages/server/package.json"),
);
const pm2Root = path.dirname(require.resolve("pm2/package.json"));
const pm2Bin = path.join(pm2Root, "bin/pm2");
const args = process.argv.slice(2);
function run(pm2Args, opts = {}) {
return spawnSync(process.execPath, [pm2Bin, ...pm2Args], {
cwd: root,
encoding: "utf8",
...opts,
});
}
const cmd = args[0];
if (cmd === "start" || cmd === "restart" || cmd === "reload") {
const probe = run(["ls"], { stdio: ["ignore", "pipe", "pipe"] });
const text = `${probe.stdout ?? ""}${probe.stderr ?? ""}`;
const mem = text.match(/In memory PM2 version:\s*(\S+)/);
const loc = text.match(/Local PM2 version:\s*(\S+)/);
if (mem && loc && mem[1] !== loc[1]) {
console.error(
`[jflow] PM2 daemon ${mem[1]} != CLI ${loc[1]}. Killing the daemon so control.js and the CLI share one version.`,
);
run(["kill"], { stdio: "inherit" });
}
}
const child = spawn(process.execPath, [pm2Bin, ...args], {
cwd: root,
stdio: "inherit",
});
child.on("exit", (code, signal) => {
if (signal) process.kill(process.pid, signal);
process.exit(code ?? 1);
});