10 Commits
Author SHA1 Message Date
nsrb 707e2c0f1e feat(plugins): enhance SMS sending functionality with raw HTTP handling
- Added support for using the `net` module in the `send-sms` plugin to establish raw HTTP connections.
- Refactored the `ztePost` function to `zteHttpPost`, allowing for custom HTTP request formatting and improved error handling.
- Introduced a new utility function, `httpBodyFromRaw`, to process raw HTTP responses effectively.
- Updated the `script-sandbox.js` to allow the use of the `net` and `node:net` modules in the sandbox environment.
2026-09-08 22:58:52 +07:00
nsrb 6c82ff20eb fix(plugins): enhance plugin installation and script handling
- Updated the `pnpm install` command in `plugin-install.js` to include the `--ignore-workspace` flag, ensuring proper installation of plugins within the app tree.
- Improved the `instantiateScriptSource` function in `script-sandbox.js` to resolve `pluginDir` more effectively, allowing for better package management.
- Added a new smoke test in `plugins-smoke.js` to validate the ability to require additional packages from plugin directories, enhancing testing coverage for plugin functionality.
- Updated documentation in `AGENTS.md` to reflect changes in the installation command.
2026-09-08 22:00:26 +07:00
nsrb 67ed3eecca feat(prayer-times): add Jakarta prayer time plugins and enhance ntfy integration
- Introduced two new plugins: `jakarta-prayer-time-parser-for-ntfy` for formatting prayer times as ntfy Markdown, and `jakarta-prayer-time-today-and-tomorrow` for fetching prayer times from Kemenag for today and tomorrow.
- Enhanced `ntfy.js` to support improved Markdown header handling and added a utility function for Markdown configuration validation.
- Updated plugin metadata and descriptions for clarity on functionality and usage.
- Implemented robust error handling and logging for better debugging and user feedback.
2026-08-30 21:25:09 +07:00
nsrb 80e9c91b88 refactor(scripts): remove admin role checks from plugin API endpoints
- Eliminated admin role checks from the create, install, and delete plugin API endpoints to simplify access control.
- Updated the ScriptsPage component to remove unnecessary user role checks, streamlining the installation process for all users.
- Improved UI by consolidating the zip installation button, enhancing user experience during plugin management.
2026-08-30 20:39:25 +07:00
nsrb 88d3692116 feat(scripts): add zip file installation functionality for plugins
- Introduced the ability to install plugins from zip files, enhancing the script management interface.
- Implemented file selection and base64 encoding for zip files before installation.
- Added confirmation dialog for installing plugins from zip, including warnings for existing plugins.
- Integrated user role checks to restrict zip installation to admin users only.
- Improved UI feedback during the installation process with loading indicators.
2026-08-30 20:27:26 +07:00
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
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
57 changed files with 1598 additions and 506 deletions
+3 -1
View File
@@ -1,6 +1,8 @@
node_modules/
.pnpm-store/
# Instance data (SQLite, live workflows, control-state, backups)
data/
# Process logs
logs/
*.db
*.db-*
@@ -10,7 +12,7 @@ packages/web/dist
# Personal/local scripts and workflows (not for the repo)
debug-*.js
# Legacy live workflow tree (migrated to packages/server/data/workflows/)
# Legacy live workflow trees
packages/server/workflows/
# Plugin install staging and per-plugin deps
+7 -7
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)
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 |
|---|---|---|
| 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) |
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
@@ -78,7 +78,7 @@ Copy the layout from `plugins/joplin-api`, `plugins/send-sms`, or `examples/plug
Add `dependencies` only if the script `require()`s extra npm packages. Then:
```bash
pnpm install --dir plugins/<id> --ignore-scripts --prefer-offline
pnpm install --dir plugins/<id> --ignore-scripts --prefer-offline --ignore-workspace
```
`plugins/*/node_modules/` is gitignored. Host-allowlisted modules (`axios`, `jsonata`, …) come from the server; extra deps resolve from the plugin directory.
@@ -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 `{`).
+30 -7
View File
@@ -55,16 +55,16 @@ pnpm --dir packages/server reset-admin -- --username admin --password 'your-pass
| 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 |
- Live YAML is **instance data**, same as SQLite and secrets — not product source. New resources use owner `local` (owner remains in storage/URLs for a possible future multi-tenant mode; the UI hides it).
- On first start, if the instance store is empty and a legacy `packages/server/workflows/` tree still exists, it is copied into `data/workflows/`.
- New workflow editor starts empty; optional presets copy example YAML into the editor (nothing is saved until Save).
- Override the live store in tests with `JFLOW_WORKFLOWS_DIR`.
```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_DB_PATH=packages/server/data/plugins-smoke.db \
node packages/server/test/plugins-smoke.js
@@ -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.… }}`. 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.
@@ -124,7 +145,7 @@ Admin UI route **Ops** (`/ops`) talks to the control process.
| Drain restart | Pause → wait active=0 → stop children → migrate → recreate → resume |
| 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
@@ -132,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_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_WORKFLOWS_DIR` | `packages/server/data/workflows` | Live workflow YAML (instance data). |
| `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_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_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. |
@@ -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
+2 -2
View File
@@ -1,8 +1,8 @@
import fs from "fs";
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.
+156 -72
View File
@@ -3,36 +3,24 @@ 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*\}\}$/;
/**
* @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.
* @param {unknown} value
* @returns {ConfigRef | null}
*/
export function parseConfigRef(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;
}
/**
* Walk config (objects/arrays) and replace whole-value `$SECRET_` / `$CONTEXT_` / `$VAR_`
* strings. Does not walk trigger data.
* Walk config (objects/arrays) and interpolate `{{ path }}` strings.
* Does not walk trigger data.
*
* @param {unknown} value
* @param {ConfigRefCtx} ctx
@@ -68,83 +56,179 @@ 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 whole = WHOLE_MUSTACHE_RE.exec(value);
if (whole && whole[0] === value) {
return resolvePath(whole[1], ctx, { raw: value, allowObject: true });
}
if (ref.kind === "secret") {
return resolveSecretRef(ref, ctx);
if (!value.includes("{{")) {
return value;
}
if (ref.kind === "var") {
return resolveVarRef(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;
}
return resolveContextRef(ref, ctx);
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`);
}
@@ -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");
});
}
@@ -17,6 +17,9 @@ export async function up(knex) {
t.text("output");
t.text("error");
t.text("parent_run_id").references("id").inTable("workflow_runs");
t.text("job_id");
t.text("queued_at");
t.integer("workflow_revision");
});
await knex.schema.raw(
@@ -28,45 +31,11 @@ export async function up(knex) {
await knex.schema.raw(
"CREATE INDEX workflow_runs_status_started_at_idx ON workflow_runs (status, started_at DESC)",
);
await knex.schema.createTable("step_runs", (t) => {
t.text("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.integer("step_index").notNullable();
t.text("script").notNullable();
t.text("config");
t.text("status").notNullable();
t.text("started_at").notNullable();
t.text("finished_at");
t.integer("duration_ms");
t.text("output");
t.text("error");
});
await knex.schema.raw(
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
);
await knex.schema.createTable("logs", (t) => {
t.increments("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.text("step_id");
t.text("ts").notNullable();
t.integer("level").notNullable();
t.text("msg");
t.text("payload");
});
await knex.schema.raw(
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
"CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)",
);
}
@@ -74,7 +43,5 @@ export async function up(knex) {
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("logs");
await knex.schema.dropTableIfExists("step_runs");
await knex.schema.dropTableIfExists("workflow_runs");
}
@@ -0,0 +1,33 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("step_runs", (t) => {
t.text("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.integer("step_index").notNullable();
t.text("script").notNullable();
t.text("config");
t.text("status").notNullable();
t.text("started_at").notNullable();
t.text("finished_at");
t.integer("duration_ms");
t.text("output");
t.text("error");
});
await knex.schema.raw(
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("step_runs");
}
@@ -0,0 +1,29 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("logs", (t) => {
t.increments("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.text("step_id");
t.text("ts").notNullable();
t.integer("level").notNullable();
t.text("msg");
t.text("payload");
});
await knex.schema.raw(
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("logs");
}
@@ -29,35 +29,35 @@ const DEFAULT_EMAIL_TEMPLATE = `<!DOCTYPE html>
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.alterTable("http_pages", (t) => {
await knex.schema.createTable("http_pages", (t) => {
t.text("id").primary();
t.text("name").notNullable().unique();
t.text("content").notNullable();
t.text("mime").notNullable();
t.integer("status").notNullable().defaultTo(200);
t.text("kind").notNullable().defaultTo("response");
t.integer("system").notNullable().defaultTo(0);
t.text("created_at").notNullable();
t.text("updated_at").notNullable();
});
const now = new Date().toISOString();
const existing = await knex("http_pages").where({ name: "email-default" }).first();
if (!existing) {
await knex("http_pages").insert({
id: "00000000-0000-4000-8000-000000000001",
name: "email-default",
content: DEFAULT_EMAIL_TEMPLATE,
mime: "html",
status: 200,
kind: "template",
system: 1,
created_at: now,
updated_at: now,
});
}
await knex("http_pages").insert({
id: "00000000-0000-4000-8000-000000000001",
name: "email-default",
content: DEFAULT_EMAIL_TEMPLATE,
mime: "html",
status: 200,
kind: "template",
system: 1,
created_at: now,
updated_at: now,
});
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex("http_pages").where({ name: "email-default", system: 1 }).del();
await knex.schema.alterTable("http_pages", (t) => {
t.dropColumn("kind");
t.dropColumn("system");
});
await knex.schema.dropTableIfExists("http_pages");
}
@@ -2,21 +2,11 @@
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("http_pages", (t) => {
t.text("id").primary();
t.text("name").notNullable().unique();
t.text("content").notNullable();
t.text("mime").notNullable(); // html | json
t.integer("status").notNullable().defaultTo(200);
t.text("created_at").notNullable();
t.text("updated_at").notNullable();
});
await knex.schema.createTable("http_auths", (t) => {
t.text("id").primary();
t.text("name").notNullable().unique();
t.text("type").notNullable(); // bearer | basic | header
t.text("config").notNullable(); // JSON
t.text("type").notNullable();
t.text("config").notNullable();
t.integer("unauthorized_status").nullable();
t.text("unauthorized_response").nullable();
t.text("created_at").notNullable();
@@ -29,5 +19,4 @@ export async function up(knex) {
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("http_auths");
await knex.schema.dropTableIfExists("http_pages");
}
@@ -21,29 +21,11 @@ export async function up(knex) {
await knex.schema.raw(
"CREATE INDEX workflow_revisions_workflow_id_created_at_idx ON workflow_revisions (workflow_id, created_at DESC)",
);
await knex.schema.createTable("workflow_trash", (t) => {
t.text("id").primary();
t.text("workflow_id").notNullable();
t.text("owner").notNullable();
t.text("file").notNullable();
t.text("name");
t.text("deleted_at").notNullable();
t.text("trash_path").notNullable();
});
await knex.schema.raw(
"CREATE INDEX workflow_trash_deleted_at_idx ON workflow_trash (deleted_at ASC)",
);
await knex.schema.raw(
"CREATE UNIQUE INDEX workflow_trash_owner_file_idx ON workflow_trash (owner, file)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("workflow_trash");
await knex.schema.dropTableIfExists("workflow_revisions");
}
@@ -0,0 +1,28 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("workflow_trash", (t) => {
t.text("id").primary();
t.text("workflow_id").notNullable();
t.text("owner").notNullable();
t.text("file").notNullable();
t.text("name");
t.text("deleted_at").notNullable();
t.text("trash_path").notNullable();
});
await knex.schema.raw(
"CREATE INDEX workflow_trash_deleted_at_idx ON workflow_trash (deleted_at ASC)",
);
await knex.schema.raw(
"CREATE UNIQUE INDEX workflow_trash_owner_file_idx ON workflow_trash (owner, file)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("workflow_trash");
}
@@ -14,7 +14,9 @@ export async function up(knex) {
t.unique(["owner", "name"]);
});
await knex.schema.raw("CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)");
await knex.schema.raw(
"CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)",
);
}
/**
+4 -3
View File
@@ -13,11 +13,12 @@
"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",
"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_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",
"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"
},
"dependencies": {
"@jerapah-flow/shared": "workspace:*",
+39 -17
View File
@@ -1,27 +1,49 @@
import fs from "fs";
import path from "path";
import { fileURLToPath } from "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 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. */
export const WORKFLOWS_DIR =
process.env.JFLOW_WORKFLOWS_DIR ?? path.join(DATA_DIR, "workflows");
/** Pre-0.1 layout; used only for one-shot migrate into WORKFLOWS_DIR. */
export const LEGACY_WORKFLOWS_DIR = path.join(SERVER_ROOT, "workflows");
export const WORKFLOWS_DIR = path.resolve(
process.env.JFLOW_WORKFLOWS_DIR ?? path.join(DATA_DIR, "workflows"),
);
/** User plugins (repo-root /plugins, outside the pnpm workspace). */
export const PLUGINS_DIR =
process.env.JFLOW_PLUGINS_DIR ??
path.resolve(SERVER_ROOT, "../../plugins");
export const PLUGINS_DIR = path.resolve(
process.env.JFLOW_PLUGINS_DIR ?? path.join(REPO_ROOT, "plugins"),
);
/** Example plugin sources shipped with the repo. */
export const EXAMPLE_PLUGINS_DIR = path.resolve(
SERVER_ROOT,
"../../examples/plugins",
);
export const EXAMPLE_PLUGINS_DIR = path.join(REPO_ROOT, "examples/plugins");
/** Example workflow YAML presets (not loaded by the runner). */
export const EXAMPLE_WORKFLOWS_DIR = path.resolve(
SERVER_ROOT,
"../../examples/workflows",
export const EXAMPLE_WORKFLOWS_DIR = path.join(REPO_ROOT, "examples/workflows");
export const LOGS_DIR = path.resolve(
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.resolve(SERVER_ROOT, "../web/dist");
export const WEB_DIST = path.join(REPO_ROOT, "packages/web/dist");
+10 -1
View File
@@ -51,9 +51,18 @@ export async function pnpmInstallPlugin(dir) {
}
if (!hasDeps) return { skipped: true };
// --ignore-workspace: plugins live under the app tree but are not workspace
// packages; without this, pnpm install --dir can no-op against the root monorepo.
await execFileAsync(
"pnpm",
["install", "--dir", dir, "--ignore-scripts", "--prefer-offline"],
[
"install",
"--dir",
dir,
"--ignore-scripts",
"--prefer-offline",
"--ignore-workspace",
],
{
cwd: dir,
env: { ...process.env, CI: "1" },
+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
* @returns {((id: string) => unknown) | null}
+1 -3
View File
@@ -1,8 +1,6 @@
import path from "path";
import pm2 from "pm2";
import { SERVER_ROOT } from "./paths.js";
const REPO_ROOT = path.resolve(SERVER_ROOT, "../..");
import { REPO_ROOT, SERVER_ROOT } from "./paths.js";
export const PM2_HTTP_NAME = "jflow-http";
export const PM2_WORKER_NAME = "jflow-worker";
+1
View File
@@ -511,6 +511,7 @@ export function createRegistry(server, opts = {}) {
owner,
workflowKey: key,
context: incomingContext,
data: ctx.data,
});
const stepCtx = {
data: ctx.data,
+13 -3
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);
@@ -21,7 +22,9 @@ const ALLOWED_MODULES = new Set([
"basic-ftp",
"jsonata",
"mustache",
"net",
"node-html-parser",
"node:net",
"node:stream",
"nodemailer",
"rss-parser",
@@ -443,7 +446,7 @@ function createScriptSandbox({
log,
script,
workflowName,
owner = "default",
owner = DEFAULT_OWNER,
$workflows = $workflowsStub,
pluginDir = null,
}) {
@@ -556,16 +559,23 @@ function instantiateCompiled(
* $workflows?: { trigger: (name: string, data?: unknown) => Promise<unknown> },
* pluginDir?: string | null,
* }} [opts]
*
* When `pluginDir` is omitted, it is resolved from `script` so inspect/dry-run
* of `plugin/<id>` can `require()` extra packages from that plugin's node_modules.
*/
export function instantiateScriptSource(script, source, opts = {}) {
const compiled = compileScriptSource(source, script);
const pluginDir =
"pluginDir" in opts
? opts.pluginDir ?? null
: (resolveScriptRef(script).pluginDir ?? null);
const fn = instantiateCompiled(compiled, {
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,
pluginDir,
});
return { fn, ...extractScriptMeta(fn) };
}
+21
View File
@@ -15,9 +15,24 @@ function ntfyHeaders(ctx) {
headers.Title = ctx.data.title;
}
if (isMarkdownEnabled(ctx.config?.markdown)) {
// ntfy accepts X-Markdown / Markdown / md with true | 1 | yes
headers.Markdown = "yes";
log.info("ntfy: Markdown enabled");
}
return headers;
}
function isMarkdownEnabled(value) {
if (value === true || value === 1) return true;
if (typeof value === "string") {
const v = value.trim().toLowerCase();
return v === "true" || v === "yes" || v === "1";
}
return false;
}
function resolveFingerprintKey(config) {
const fingerprint = config?.fingerprint;
if (fingerprint === true) return "fingerprint:ntfy";
@@ -149,6 +164,12 @@ ntfy.meta = {
default: "https://ntfy.sh/jerapah-flow",
description: "ntfy topic URL",
},
markdown: {
type: "boolean",
default: false,
description:
"Render body as Markdown (sets Markdown: yes; aliases X-Markdown / md also accepted by ntfy)",
},
fingerprint: {
type: "string",
required: false,
+40 -5
View File
@@ -8,6 +8,7 @@ import {
import * as fsStore from "../../fs-store.js";
import {
forkCoreScript,
duplicatePlugin,
listCoreScriptNames,
listInstalledPlugins,
resolveScriptRef,
@@ -29,8 +30,9 @@ import { normalizeStepResult } from "../../step-result.js";
import { resolveConfigRefs } from "../../config-refs.js";
import { getAppVersion } from "../../app-version.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 { DEFAULT_OWNER } from "@jerapah-flow/shared";
/**
* @param {{ referencedScripts: () => Set<string> }} registry
@@ -250,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) => {
const rawName = decodeURIComponent(
/** @type {{ name: string }} */ (req.params).name,
@@ -258,7 +294,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 +340,7 @@ export default function scriptsPluginFactory(registry) {
owner,
workflowKey: "dry-run",
context: incomingContext,
data: incomingData,
},
);
const ctx = {
@@ -365,6 +402,7 @@ export default function scriptsPluginFactory(registry) {
owner,
workflowKey: "dry-run",
context: incomingContext,
data: incomingData,
});
const ctx = {
data: incomingData,
@@ -428,7 +466,6 @@ export default function scriptsPluginFactory(registry) {
fastify.post(
"/plugins/create",
{ onRequest: [fastify.requireAdmin] },
async (req, reply) => {
const body = /** @type {{ id?: string, content?: string, description?: string }} */ (
req.body ?? {}
@@ -460,7 +497,6 @@ export default function scriptsPluginFactory(registry) {
fastify.post(
"/plugins/install",
{ onRequest: [fastify.requireAdmin] },
async (req, reply) => {
const body = /** @type {{
source?: string,
@@ -535,7 +571,6 @@ export default function scriptsPluginFactory(registry) {
fastify.delete(
"/plugins/:id",
{ onRequest: [fastify.requireAdmin] },
async (req, reply) => {
const { id } = /** @type {{ id: string }} */ (req.params);
try {
+7 -8
View File
@@ -706,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
@@ -31,7 +31,6 @@ import {
getRedisUrlForLog,
} from "./workflow-queue.js";
import { purgeExpiredTrash } from "./workflow-trash.js";
import { migrateLegacyWorkflowsIfNeeded } from "./workflow-migrate.js";
import {
getConfigGeneration,
startHeartbeatLoop,
@@ -129,12 +128,6 @@ export async function startApp(opts = {}) {
}
});
try {
migrateLegacyWorkflowsIfNeeded();
} catch (err) {
log.warn({ err }, "legacy workflow migrate failed");
}
const registry = createRegistry(server, {
queue: workflowQueue,
// Cron + HTTP triggers enqueue jobs; only the API process may own them.
+109 -109
View File
@@ -2,7 +2,7 @@ import { migrate, db } from "../db.js";
import { upsertSecret, deleteSecret } from "../secrets-store.js";
import { deleteVariable, upsertVariable } from "../variables-store.js";
import { Secret } from "../secret-value.js";
import { parseConfigRef, resolveConfigRefs } from "../config-refs.js";
import { resolveConfigRefs } from "../config-refs.js";
await migrate();
@@ -23,126 +23,138 @@ 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: {} };
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)}`);
}
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" });
const owner = "config_refs_smoke_owner";
const ctx = { owner, workflowKey: `${owner}/config-refs-smoke.yaml`, context: {}, data: {} };
{
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");
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,
});
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 +167,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();
+72
View File
@@ -2,20 +2,24 @@
* Smoke: core vs plugin scripts, fork, example install, resolve, run.
*
* Run:
* JFLOW_DATA_DIR=packages/server/data \
* JFLOW_PLUGINS_DIR=packages/server/data/plugins-smoke-test \
* JFLOW_DB_PATH=packages/server/data/plugins-smoke.db \
* node packages/server/test/plugins-smoke.js
*/
import assert from "node:assert/strict";
import fs from "fs";
import path from "node:path";
import { migrate, db } from "../db.js";
import { getAppVersion, satisfiesRange } from "../app-version.js";
import {
forkCoreScript,
duplicatePlugin,
resolveScriptRef,
uninstallPlugin,
listInstalledPlugins,
createBlankPlugin,
pluginDir,
} from "../plugin-store.js";
import { installExamplePlugin } from "../plugin-install.js";
import {
@@ -82,6 +86,26 @@ async function main() {
);
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;
try {
forkCoreScript("ntfy.js", "ntfy");
@@ -97,10 +121,58 @@ async function main() {
);
assert.ok(meta);
const extra = createBlankPlugin(
"extra-require-smoke",
`const extra = require("smoke-extra");
async function main() {
return { output: { n: extra.n } };
}
main.meta = {
description: "smoke extra require",
config: {},
input: {},
output: { n: { type: "number" } },
example: { data: {}, config: {} },
};
export default main;
`,
);
assert.equal(extra.scriptRef, "plugin/extra-require-smoke");
const extraPkg = path.join(
pluginDir("extra-require-smoke"),
"node_modules",
"smoke-extra",
);
fs.mkdirSync(extraPkg, { recursive: true });
fs.writeFileSync(
path.join(extraPkg, "package.json"),
`${JSON.stringify({ name: "smoke-extra", main: "index.js" })}\n`,
);
fs.writeFileSync(path.join(extraPkg, "index.js"), "module.exports = { n: 9 };\n");
clearScriptCache();
const extraSource = fs.readFileSync(
path.join(pluginDir("extra-require-smoke"), "script.js"),
"utf8",
);
const extraInspect = inspectScriptSource(
"plugin/extra-require-smoke",
extraSource,
);
assert.equal(extraInspect.metaError, null, extraInspect.metaError);
assert.equal(extraInspect.meta?.description, "smoke extra require");
const extraRun = await runScript(
"plugin/extra-require-smoke",
{ data: null, context: {}, config: null },
{ log: silent, workflowName: "smoke", owner: "default" },
);
assert.equal(extraRun.output.n, 9);
assert.ok(listInstalledPlugins().some((p) => p.id === "get-current-time"));
uninstallPlugin("jsonata-smoke-fork");
uninstallPlugin("blank-smoke");
uninstallPlugin("blank-smoke-copy");
uninstallPlugin("extra-require-smoke");
uninstallPlugin("get-current-time");
console.log("plugins-smoke: ok");
+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,
+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 (
-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;
}
+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";
+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() {
const qc = useQueryClient();
return useMutation({
@@ -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,27 @@ 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 CONFIG_REF_PREFIXES = [
{ prefix: "$SECRET_", label: "secret" },
{ prefix: "$CONTEXT_", label: "context" },
{ prefix: "$VAR_", label: "variable" },
];
/**
* @param {unknown} value
* @returns {{
* kind: "mustache",
* path?: string,
* root?: string,
* name?: 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 };
}
}
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 +38,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 +67,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);
@@ -73,9 +78,17 @@ export function ConfigRefHint({ value, owner }) {
if (!ref) return null;
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>
);
}
+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";
+39 -1
View File
@@ -4,6 +4,7 @@ import { LuArrowLeft, LuCopy, LuPlay, LuSave } from "react-icons/lu";
import { errorMessage } from "../api/client.js";
import {
useCreatePlugin,
useDuplicatePlugin,
useForkScript,
useSaveScript,
useScript,
@@ -119,9 +120,11 @@ export function ScriptEditPage() {
const existing = useScript(name);
const save = useSaveScript();
const fork = useForkScript();
const duplicate = useDuplicatePlugin();
const [content, setContent] = useState("");
const [contentReady, setContentReady] = useState(false);
const [forkId, setForkId] = useState("");
const [duplicateId, setDuplicateId] = useState("");
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) {
return (
<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>
) : null}
</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">
<CodeEditor
+184 -11
View File
@@ -1,9 +1,10 @@
import { useMemo, useState } from "react";
import { useMemo, useRef, useState } from "react";
import { Link, Navigate, useNavigate, useSearchParams } from "react-router-dom";
import { LuCopy, LuPencil, LuPlay, LuPlus, LuSearch, LuTrash2 } from "react-icons/lu";
import { LuCopy, LuPencil, LuPlay, LuPlus, LuSearch, LuTrash2, LuUpload } from "react-icons/lu";
import { errorMessage } from "../api/client.js";
import {
useDeleteScript,
useDuplicatePlugin,
useForkScript,
useInstallPlugin,
useScripts,
@@ -14,6 +15,21 @@ import { TagBadge } from "../components/TagBadge.jsx";
import { scriptTags } from "../lib/script.js";
import { useNotifications } from "../notifications.jsx";
const BASE64_CHUNK = 0x8000;
function uint8ArrayToBase64(bytes) {
let binary = "";
for (let i = 0; i < bytes.length; i += BASE64_CHUNK) {
binary += String.fromCharCode(...bytes.subarray(i, i + BASE64_CHUNK));
}
return btoa(binary);
}
async function fileToBase64(file) {
const bytes = new Uint8Array(await file.arrayBuffer());
return uint8ArrayToBase64(bytes);
}
export function ScriptsPage() {
const [params] = useSearchParams();
const navigate = useNavigate();
@@ -23,10 +39,17 @@ export function ScriptsPage() {
const [query, setQuery] = useState("");
const [forkFor, setForkFor] = useState(null);
const [forkId, setForkId] = useState("");
const [duplicateFor, setDuplicateFor] = useState(null);
const [duplicateId, setDuplicateId] = useState("");
const del = useDeleteScript();
const fork = useForkScript();
const duplicate = useDuplicatePlugin();
const install = useInstallPlugin();
const { notify } = useNotifications();
const [pendingZip, setPendingZip] = useState(null);
const [zipEncoding, setZipEncoding] = useState(false);
const zipInFlight = useRef(false);
const zipBusy = zipEncoding || install.isPending;
const visible = useMemo(() => {
const term = query.trim().toLowerCase();
@@ -45,6 +68,44 @@ export function ScriptsPage() {
});
}, [scripts, query]);
function onPickZip(e) {
const file = e.target.files?.[0];
e.target.value = "";
if (!file) return;
install.reset();
setPendingZip(file);
}
async function installPendingZip() {
if (!pendingZip || zipInFlight.current) return;
zipInFlight.current = true;
setZipEncoding(true);
try {
const zipBase64 = await fileToBase64(pendingZip);
install.mutate(
{ source: "zip", zipBase64, overwrite: true },
{
onSuccess: (data) => {
setPendingZip(null);
const ref =
data.scriptRef ?? (data.id ? `plugin/${data.id}` : "plugin");
notify.success(`Installed ${ref} — drain-restart to load workers`);
if (data.warning) notify.warning(data.warning);
},
onError: (err) => notify.error(errorMessage(err)),
onSettled: () => {
zipInFlight.current = false;
setZipEncoding(false);
},
},
);
} catch (err) {
zipInFlight.current = false;
setZipEncoding(false);
notify.error(errorMessage(err));
}
}
if (editName) {
return <Navigate to={`/scripts/${encodeURIComponent(editName)}/edit`} replace />;
}
@@ -67,7 +128,7 @@ export function ScriptsPage() {
<button
type="button"
className="btn btn-outline btn-sm"
disabled={install.isPending}
disabled={zipBusy}
onClick={() =>
install.mutate(
{ source: "example", exampleId: "get-current-time", overwrite: true },
@@ -80,6 +141,21 @@ export function ScriptsPage() {
>
Install example
</button>
<label className={`btn btn-outline btn-sm ${zipBusy ? "btn-disabled" : ""}`}>
{zipBusy ? (
<span className="loading loading-spinner loading-xs" />
) : (
<LuUpload className="size-4" />
)}
Install from zip
<input
type="file"
accept=".zip,application/zip"
className="hidden"
disabled={zipBusy}
onChange={onPickZip}
/>
</label>
<Link to="/scripts/new" className="btn btn-primary btn-sm">
<LuPlus className="size-4" />
Add plugin
@@ -182,14 +258,33 @@ export function ScriptsPage() {
<LuCopy className="size-4" />
</button>
) : (
<Link
to={`/scripts/${encodeURIComponent(name)}/edit`}
className="btn btn-ghost btn-xs"
title="Edit"
aria-label="Edit"
>
<LuPencil className="size-4" />
</Link>
<>
<Link
to={`/scripts/${encodeURIComponent(name)}/edit`}
className="btn btn-ghost btn-xs"
title="Edit"
aria-label="Edit"
>
<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 ? (
<button
@@ -264,6 +359,84 @@ export function ScriptsPage() {
</dialog>
) : 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
open={Boolean(pendingZip)}
title="Install plugin from zip?"
message={
pendingZip ? (
<>
Installing{" "}
<span className="font-mono">{pendingZip.name}</span> replaces an
existing plugin with the same id. Drain-restart workers after
install so new dependencies load.
</>
) : null
}
confirmLabel="Install"
confirmClass="btn-primary"
error={install.isError ? errorMessage(install.error) : null}
loading={zipBusy}
onCancel={() => {
if (zipBusy) return;
setPendingZip(null);
}}
onConfirm={installPendingZip}
/>
<ConfirmDialog
open={Boolean(confirmDelete)}
title={confirmDelete ? `Delete ${confirmDelete}?` : ""}
+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}
@@ -0,0 +1,8 @@
{
"id": "jakarta-prayer-time-parser-for-ntfy",
"name": "Jakarta prayer times → ntfy",
"version": "0.2.0",
"jerapah": ">=0.1.0 <1.0.0",
"main": "script.js",
"description": "Format today/tomorrow prayer times as ntfy Markdown (headings + tables)"
}
@@ -0,0 +1,7 @@
{
"name": "jflow-plugin-jakarta-prayer-time-parser-for-ntfy",
"version": "0.2.0",
"private": true,
"type": "module",
"description": "JerapahFlow plugin: format Jakarta prayer times as ntfy Markdown"
}
@@ -0,0 +1,165 @@
const PRAYERS = [
["imsak", "Imsak"],
["subuh", "Subuh"],
["terbit", "Terbit"],
["dhuha", "Dhuha"],
["dzuhur", "Dzuhur"],
["ashar", "Ashar"],
["maghrib", "Maghrib"],
["isya", "Isya"],
];
function passContext(ctx) {
if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) {
return { ...ctx.context };
}
return {};
}
function unwrapScalar(value) {
if (value == null) return "";
if (typeof value === "string") return value.trim();
if (typeof value === "number" || typeof value === "boolean") return String(value);
return "";
}
function requireDay(day, label) {
if (day == null || typeof day !== "object" || Array.isArray(day)) {
throw new Error(
`parser: missing ${label} prayer-time data (run plugin/jakarta-prayer-time-today-and-tomorrow first)`,
);
}
return day;
}
function dayCaption(day, fallbackDate) {
const tanggal = unwrapScalar(day.tanggal);
if (tanggal) return tanggal;
return unwrapScalar(fallbackDate);
}
function formatDayTable(heading, day) {
const lines = [
`### ${heading}`,
"",
"| Prayer | Time |",
"| --- | --- |",
];
for (const [key, label] of PRAYERS) {
const time = unwrapScalar(day[key]) || "—";
lines.push(`| ${label} | ${time} |`);
}
return lines.join("\n");
}
async function parsePrayerTimesForNtfy(ctx) {
const data = ctx?.data && typeof ctx.data === "object" ? ctx.data : {};
const today = requireDay(data.today, "today");
const tomorrow = requireDay(data.tomorrow, "tomorrow");
const location = unwrapScalar(ctx.config?.location) || "Jakarta";
const todayCaption = dayCaption(today, data.todayDate);
const tomorrowCaption = dayCaption(tomorrow, data.tomorrowDate);
const todayHeading = todayCaption ? `Today — ${todayCaption}` : "Today";
const tomorrowHeading = tomorrowCaption
? `Tomorrow — ${tomorrowCaption}`
: "Tomorrow";
const message = [
formatDayTable(todayHeading, today),
"",
"---",
"",
formatDayTable(tomorrowHeading, tomorrow),
].join("\n");
const title = todayCaption
? `Prayer Times (${location}): ${todayCaption}`
: `Prayer Times (${location})`;
log.info(
{ location, todayCaption, tomorrowCaption },
"jakarta-prayer-time-parser-for-ntfy: formatted",
);
const output = { ok: true, title, message };
return { output, context: { ...passContext(ctx), ...output } };
}
parsePrayerTimesForNtfy.meta = {
description:
"Format today/tomorrow prayer times as ntfy Markdown (headings + tables). Set ntfy.js config.markdown to true.",
previewConfigKey: "location",
tags: ["ntfy", "Markdown", "Prayer Times"],
config: {
location: {
type: "string",
default: "Jakarta",
description: "City name used in the ntfy title",
},
},
input: {
today: {
type: "any",
required: true,
description: "Today's prayer-time record from the fetcher",
},
tomorrow: {
type: "any",
required: true,
description: "Tomorrow's prayer-time record from the fetcher",
},
todayDate: {
type: "string",
required: false,
description: "YYYY-MM-DD fallback when today.tanggal is missing",
},
tomorrowDate: {
type: "string",
required: false,
description: "YYYY-MM-DD fallback when tomorrow.tanggal is missing",
},
},
output: {
ok: { type: "boolean" },
title: { type: "string", description: "ntfy Title header" },
message: { type: "string", description: "Markdown body (headings + tables)" },
},
context: {
ok: { type: "boolean" },
title: { type: "string" },
message: { type: "string" },
},
example: {
data: {
todayDate: "2026-08-30",
tomorrowDate: "2026-08-31",
today: {
tanggal: "Minggu, 30/08/2026",
imsak: "04:32",
subuh: "04:42",
terbit: "05:57",
dhuha: "06:26",
dzuhur: "11:58",
ashar: "15:18",
maghrib: "17:58",
isya: "19:08",
},
tomorrow: {
tanggal: "Senin, 31/08/2026",
imsak: "04:32",
subuh: "04:42",
terbit: "05:57",
dhuha: "06:26",
dzuhur: "11:58",
ashar: "15:17",
maghrib: "17:58",
isya: "19:08",
},
},
config: { location: "Jakarta" },
},
};
export default parsePrayerTimesForNtfy;
@@ -0,0 +1,8 @@
{
"id": "jakarta-prayer-time-today-and-tomorrow",
"name": "Jakarta prayer times (today & tomorrow)",
"version": "0.2.0",
"jerapah": ">=0.1.0 <1.0.0",
"main": "script.js",
"description": "Fetch Kemenag prayer times for today and tomorrow using dayjs (Asia/Jakarta)"
}
@@ -0,0 +1,10 @@
{
"name": "jflow-plugin-jakarta-prayer-time-today-and-tomorrow",
"version": "0.2.0",
"private": true,
"type": "module",
"description": "JerapahFlow plugin: fetch Jakarta prayer times for today and tomorrow",
"dependencies": {
"dayjs": "^1.11.13"
}
}
@@ -0,0 +1,22 @@
lockfileVersion: '9.0'
settings:
autoInstallPeers: true
excludeLinksFromLockfile: false
importers:
.:
dependencies:
dayjs:
specifier: ^1.11.13
version: 1.11.21
packages:
dayjs@1.11.21:
resolution: {integrity: sha512-98IT+HOahAisibz/yjKbzuOBwYcjJ7BCLPzARyHiyEBmRz4fatF+KPJszEHXsGYjUG234aH/cOjW1wwTbKUZlA==}
snapshots:
dayjs@1.11.21: {}
@@ -0,0 +1,255 @@
import dayjs from "dayjs";
import utc from "dayjs/plugin/utc";
import timezone from "dayjs/plugin/timezone";
dayjs.extend(utc);
dayjs.extend(timezone);
const PAGE_URL = "https://bimasislam.kemenag.go.id/jadwalshalat";
const API_URL = "https://bimasislam.kemenag.go.id/ajax/getShalatbln";
/** Default Kemenag hashes: DKI Jakarta / Kota Jakarta. */
const DEFAULT_X = "c51ce410c124a10e0db5e4b97fc2af39";
const DEFAULT_Y = "58a2fc6ed39fd083f55d4182bf88826d";
const DEFAULT_TIMEZONE = "Asia/Jakarta";
const BROWSER_HEADERS = {
accept: "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
"accept-language": "en-US,en;q=0.9,id;q=0.8",
"cache-control": "no-cache",
pragma: "no-cache",
};
function passContext(ctx) {
if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) {
return { ...ctx.context };
}
return {};
}
function unwrapScalar(value) {
if (value == null) return "";
if (typeof value === "string") return value.trim();
if (typeof value === "number" || typeof value === "boolean") return String(value);
return "";
}
function monthKey(d) {
return `${d.year()}-${d.month() + 1}`;
}
function parseCookies(setCookie) {
if (!setCookie) return "";
return setCookie
.split(/,(?=[^;,]+=)/)
.map((cookie) => cookie.split(";")[0])
.join("; ");
}
async function getSessionCookie() {
const cookieFetch = await fetch(PAGE_URL, {
method: "GET",
headers: {
...BROWSER_HEADERS,
"upgrade-insecure-requests": "1",
},
});
const cookies = parseCookies(cookieFetch.headers.get("set-cookie"));
if (!cookies) {
throw new Error(
`Prayer-time session cookie missing (HTTP ${cookieFetch.status})`,
);
}
return cookies;
}
async function fetchPrayerMonth(cookies, month, year, x, y) {
const resp = await fetch(API_URL, {
method: "POST",
headers: {
accept: "application/json, text/javascript, */*; q=0.01",
"accept-language": BROWSER_HEADERS["accept-language"],
"cache-control": "no-cache",
"content-type": "application/x-www-form-urlencoded; charset=UTF-8",
pragma: "no-cache",
"x-requested-with": "XMLHttpRequest",
cookie: cookies,
Referer: PAGE_URL,
},
body: new URLSearchParams({
x,
y,
bln: String(month),
thn: String(year),
}).toString(),
});
const responseText = await resp.text();
let responseData;
try {
responseData = JSON.parse(responseText);
} catch {
throw new Error(
`Prayer-time API returned invalid JSON (HTTP ${resp.status}) for ${year}-${month}`,
);
}
if (!resp.ok) {
throw new Error(
`Prayer-time API request failed with HTTP ${resp.status} for ${year}-${month}`,
);
}
return {
data: responseData,
status: resp.status,
month: Number(month),
year: Number(year),
};
}
/**
* Resolve one day's record from a monthly API payload.
* Accepts `{ data: { "YYYY-MM-DD": {...} } }` or a date-keyed object.
*/
function getDateData(monthResponse, dateKey) {
const candidates = [
monthResponse?.data?.data,
monthResponse?.data,
monthResponse,
];
for (const map of candidates) {
if (map && typeof map === "object" && !Array.isArray(map) && map[dateKey]) {
const day = map[dateKey];
if (typeof day === "object" && !Array.isArray(day)) {
return day;
}
}
}
throw new Error(`No prayer times for ${dateKey}`);
}
async function fetchPrayerTimes(ctx) {
const timezoneName = unwrapScalar(ctx.config?.timezone) || DEFAULT_TIMEZONE;
const x = unwrapScalar(ctx.config?.x) || DEFAULT_X;
const y = unwrapScalar(ctx.config?.y) || DEFAULT_Y;
const today = dayjs().tz(timezoneName);
if (!today.isValid()) {
throw new Error(`invalid timezone "${timezoneName}"`);
}
const tomorrow = today.add(1, "day");
const todayDate = today.format("YYYY-MM-DD");
const tomorrowDate = tomorrow.format("YYYY-MM-DD");
log.info(
{ today: todayDate, tomorrow: tomorrowDate, timezone: timezoneName },
"jakarta-prayer-time-today-and-tomorrow: fetching",
);
const cookies = await getSessionCookie();
const needed = [
{ key: monthKey(today), month: today.month() + 1, year: today.year() },
{ key: monthKey(tomorrow), month: tomorrow.month() + 1, year: tomorrow.year() },
];
const unique = [];
const seen = new Set();
for (const item of needed) {
if (seen.has(item.key)) continue;
seen.add(item.key);
unique.push(item);
}
const monthResponses = {};
await Promise.all(
unique.map(async (item) => {
monthResponses[item.key] = await fetchPrayerMonth(
cookies,
item.month,
item.year,
x,
y,
);
}),
);
const todayData = getDateData(monthResponses[monthKey(today)], todayDate);
const tomorrowData = getDateData(
monthResponses[monthKey(tomorrow)],
tomorrowDate,
);
const output = {
ok: true,
today: todayData,
tomorrow: tomorrowData,
todayDate,
tomorrowDate,
timezone: timezoneName,
};
log.info(
{
todayFound: Boolean(output.today),
tomorrowFound: Boolean(output.tomorrow),
monthsFetched: unique.length,
},
"jakarta-prayer-time-today-and-tomorrow: extracted",
);
return {
output,
context: { ...passContext(ctx), ...output },
};
}
fetchPrayerTimes.meta = {
description:
"Fetch Kemenag Bimas Islam prayer times for today and tomorrow (Asia/Jakarta by default).",
previewConfigKey: "timezone",
tags: ["HTTP", "Prayer Times"],
config: {
timezone: {
type: "string",
default: DEFAULT_TIMEZONE,
description: "IANA timezone used to pick today/tomorrow",
},
x: {
type: "string",
default: DEFAULT_X,
description: "Kemenag kabupaten hash (x). Default is DKI Jakarta",
},
y: {
type: "string",
default: DEFAULT_Y,
description: "Kemenag kota hash (y). Default is Kota Jakarta",
},
},
input: {},
output: {
ok: { type: "boolean" },
today: { type: "any", description: "Prayer-time record for today" },
tomorrow: { type: "any", description: "Prayer-time record for tomorrow" },
todayDate: { type: "string", description: "YYYY-MM-DD in the configured timezone" },
tomorrowDate: { type: "string", description: "YYYY-MM-DD in the configured timezone" },
timezone: { type: "string" },
},
context: {
ok: { type: "boolean" },
today: { type: "any" },
tomorrow: { type: "any" },
todayDate: { type: "string" },
tomorrowDate: { type: "string" },
timezone: { type: "string" },
},
example: {
data: {},
config: {
timezone: "Asia/Jakarta",
},
},
};
export default fetchPrayerTimes;
+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,
},
},
+67 -28
View File
@@ -1,3 +1,5 @@
import net from "node:net";
function passContext(ctx) {
if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) {
return { ...ctx.context };
@@ -65,38 +67,75 @@ function parseResult(text) {
}
}
function leftoverFromError(err) {
let current = err;
for (let i = 0; i < 6 && current; i++) {
if (typeof current.data === "string") return current.data;
current = current.cause;
function httpBodyFromRaw(raw) {
const normalized = String(raw).replace(/\r\n/g, "\n");
const split = normalized.indexOf("\n\n");
return split === -1 ? normalized : normalized.slice(split + 2);
}
/**
* ZTE Demo-Webs uses LF-only status lines (`HTTP/1.1 200 OK\n`).
* Node 24 fetch/http throw HPE_CR_EXPECTED before any body is read.
*/
function zteHttpPost(url, fields) {
const target = new URL(url);
if (target.protocol !== "http:") {
throw new Error(`send-sms: only http: is supported (got ${target.protocol})`);
}
return null;
const body = new URLSearchParams(fields).toString();
const request = [
`POST ${target.pathname}${target.search} HTTP/1.1`,
`Host: ${target.host}`,
"accept: application/json, text/javascript, */*; q=0.01",
"accept-language: en-US,en;q=0.9",
"content-type: application/x-www-form-urlencoded; charset=UTF-8",
"x-requested-with: XMLHttpRequest",
`Referer: ${refererFromUrl(url)}`,
`Content-Length: ${Buffer.byteLength(body)}`,
"Connection: close",
"",
body,
].join("\r\n");
return new Promise((resolve, reject) => {
const socket = net.connect(
{ host: target.hostname, port: Number(target.port) || 80 },
() => {
socket.write(request);
},
);
const chunks = [];
let settled = false;
const finish = (raw) => {
if (settled) return;
settled = true;
socket.removeAllListeners();
socket.destroy();
resolve(parseResult(httpBodyFromRaw(raw)));
};
socket.on("data", (chunk) => {
chunks.push(chunk);
const raw = Buffer.concat(chunks).toString("utf8");
const parsed = parseResult(httpBodyFromRaw(raw));
if (parsed.result !== undefined) finish(raw);
});
socket.on("end", () => finish(Buffer.concat(chunks).toString("utf8")));
socket.on("error", (err) => {
if (settled) return;
settled = true;
reject(err);
});
socket.setTimeout(15000, () => {
if (settled) return;
settled = true;
socket.destroy();
reject(new Error("ZTE HTTP timeout"));
});
});
}
async function ztePost(url, fields) {
const body = new URLSearchParams(fields).toString();
try {
const response = await fetch(url, {
method: "POST",
headers: {
accept: "application/json, text/javascript, */*; q=0.01",
"content-type": "application/x-www-form-urlencoded; charset=UTF-8",
"x-requested-with": "XMLHttpRequest",
Referer: refererFromUrl(url),
},
body,
});
return parseResult(await response.text());
} catch (err) {
const leftover = leftoverFromError(err);
if (leftover != null) {
const parsed = parseResult(leftover);
log.info({ result: parsed.result ?? parsed }, "send-sms: recovered malformed response");
return parsed;
}
throw err;
}
return zteHttpPost(url, fields);
}
async function sendZteSms(ctx) {