Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
707e2c0f1e | ||
|
|
6c82ff20eb | ||
|
|
67ed3eecca | ||
|
|
80e9c91b88 | ||
|
|
88d3692116 | ||
|
|
a20b26b4c7 | ||
|
|
f1cdb7ac68 | ||
|
|
3f0774813d |
+3
-1
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
@@ -109,7 +109,7 @@ topic: "{{ vars.ntfy_prefix }}/{{ data.channel }}"
|
||||
| `context` | Run clipboard (nested) |
|
||||
| `data` | This step’s input (nested; numeric segments index arrays) |
|
||||
|
||||
A string that is exactly one `{{ path }}` keeps the native type (object/array/number/boolean). Mixed strings concatenate as text. Bare name fields such as `passwordSecret: gmail_app_password` stay names for `$secrets.get` — do not wrap them in `{{ secrets.… }}`. Legacy `$VAR_` / `$SECRET_` / `$CONTEXT_` whole-value refs throw; use mustache instead. Script APIs `$vars.get` / `$secrets.get` are unchanged.
|
||||
A string that is exactly one `{{ path }}` keeps the native type (object/array/number/boolean). Mixed strings concatenate as text. Bare name fields such as `passwordSecret: gmail_app_password` stay names for `$secrets.get` — do not wrap them in `{{ secrets.… }}`. Script APIs `$vars.get` / `$secrets.get` are unchanged.
|
||||
|
||||
YAML **SET** evaluates JSONata against the full `ctx`; the result is `output` (the next step’s data). `jsonata.js` does the same.
|
||||
|
||||
@@ -145,7 +145,7 @@ Admin UI route **Ops** (`/ops`) talks to the control process.
|
||||
| Drain restart | Pause → wait active=0 → stop children → migrate → recreate → resume |
|
||||
| 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
|
||||
|
||||
@@ -153,8 +153,10 @@ Desired state is stored in `packages/server/data/control-state.json` (generation
|
||||
|---|---|---|
|
||||
| `JFLOW_JWT_SECRET` | `jflow-dev-secret` (dev only) | **Required in production**. |
|
||||
| `JFLOW_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,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.
|
||||
|
||||
@@ -1,17 +0,0 @@
|
||||
/**
|
||||
* Rewrite legacy `$VAR_` / `$SECRET_` / `$CONTEXT_` placeholders to mustache.
|
||||
* Safe for passwordSecret-style fields (those store bare names, not $SECRET_ prefixes).
|
||||
*
|
||||
* @param {string} text
|
||||
* @returns {{ text: string, changed: boolean }}
|
||||
*/
|
||||
export function rewriteLegacyConfigRefsInText(text) {
|
||||
if (typeof text !== "string" || text.length === 0) {
|
||||
return { text: text ?? "", changed: false };
|
||||
}
|
||||
const next = text
|
||||
.replace(/\$VAR_([A-Za-z0-9._-]+)/g, "{{ vars.$1 }}")
|
||||
.replace(/\$SECRET_([A-Za-z0-9._-]+)/g, "{{ secrets.$1 }}")
|
||||
.replace(/\$CONTEXT_([A-Za-z0-9._-]+)/g, "{{ context.$1 }}");
|
||||
return { text: next, changed: next !== text };
|
||||
}
|
||||
@@ -8,7 +8,6 @@ const MUSTACHE_TOKEN_RE =
|
||||
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g;
|
||||
const WHOLE_MUSTACHE_RE =
|
||||
/^\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}$/;
|
||||
const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
|
||||
/**
|
||||
* @typedef {{
|
||||
@@ -19,30 +18,6 @@ const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
* }} ConfigRefCtx
|
||||
*/
|
||||
|
||||
/**
|
||||
* Detect leftover `$VAR_` / `$SECRET_` / `$CONTEXT_` whole-value refs.
|
||||
* @param {unknown} value
|
||||
* @returns {{ kind: "var" | "secret" | "context", name: string, raw: string } | null}
|
||||
*/
|
||||
export function parseLegacyConfigRef(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const match = LEGACY_PREFIX_RE.exec(value);
|
||||
if (!match) return null;
|
||||
const kind =
|
||||
match[1] === "VAR" ? "var" : match[1] === "SECRET" ? "secret" : "context";
|
||||
return { kind, name: match[2] ?? "", raw: value.trim() };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {"var" | "secret" | "context"} kind
|
||||
* @param {string} name
|
||||
*/
|
||||
function legacyRenameHint(kind, name) {
|
||||
if (kind === "var") return `use {{ vars.${name || "name"} }}`;
|
||||
if (kind === "secret") return `use {{ secrets.${name || "name"} }}`;
|
||||
return `use {{ context.${name || "name"} }}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Walk config (objects/arrays) and interpolate `{{ path }}` strings.
|
||||
* Does not walk trigger data.
|
||||
@@ -84,13 +59,6 @@ export async function resolveConfigRefs(value, ctx, seen = new WeakSet()) {
|
||||
* @returns {Promise<unknown>}
|
||||
*/
|
||||
async function resolveStringRef(value, ctx) {
|
||||
const legacy = parseLegacyConfigRef(value);
|
||||
if (legacy) {
|
||||
throw new Error(
|
||||
`config ref ${legacy.raw}: removed; ${legacyRenameHint(legacy.kind, legacy.name)}`,
|
||||
);
|
||||
}
|
||||
|
||||
const whole = WHOLE_MUSTACHE_RE.exec(value);
|
||||
if (whole && whole[0] === value) {
|
||||
return resolvePath(whole[1], ctx, { raw: value, allowObject: true });
|
||||
|
||||
@@ -1,24 +0,0 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.text("job_id");
|
||||
t.text("queued_at");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_job_id_idx");
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.dropColumn("job_id");
|
||||
t.dropColumn("queued_at");
|
||||
});
|
||||
}
|
||||
@@ -1,21 +0,0 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.integer("workflow_revision");
|
||||
});
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_workflow_revision_idx");
|
||||
await knex.schema.alterTable("workflow_runs", (t) => {
|
||||
t.dropColumn("workflow_revision");
|
||||
});
|
||||
}
|
||||
@@ -1,16 +0,0 @@
|
||||
/**
|
||||
* No-op Knex marker: owner + config-ref migrate runs in owner-migrate.js at startup
|
||||
* (needs filesystem + DB together). Kept so deploy tooling sees a versioned step.
|
||||
*
|
||||
* @param {import("knex").Knex} _knex
|
||||
*/
|
||||
export async function up(_knex) {
|
||||
// Intentionally empty — see migrateDefaultOwnerIfNeeded().
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} _knex
|
||||
*/
|
||||
export async function down(_knex) {
|
||||
// Irreversible data migrate.
|
||||
}
|
||||
+5
-38
@@ -17,6 +17,9 @@ export async function up(knex) {
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
t.text("parent_run_id").references("id").inTable("workflow_runs");
|
||||
t.text("job_id");
|
||||
t.text("queued_at");
|
||||
t.integer("workflow_revision");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
@@ -28,45 +31,11 @@ export async function up(knex) {
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_runs_status_started_at_idx ON workflow_runs (status, started_at DESC)",
|
||||
);
|
||||
|
||||
await knex.schema.createTable("step_runs", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.integer("step_index").notNullable();
|
||||
t.text("script").notNullable();
|
||||
t.text("config");
|
||||
t.text("status").notNullable();
|
||||
t.text("started_at").notNullable();
|
||||
t.text("finished_at");
|
||||
t.integer("duration_ms");
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
|
||||
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
|
||||
);
|
||||
|
||||
await knex.schema.createTable("logs", (t) => {
|
||||
t.increments("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.text("step_id");
|
||||
t.text("ts").notNullable();
|
||||
t.integer("level").notNullable();
|
||||
t.text("msg");
|
||||
t.text("payload");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
|
||||
"CREATE INDEX workflow_runs_workflow_revision_idx ON workflow_runs (workflow, workflow_revision)",
|
||||
);
|
||||
}
|
||||
|
||||
@@ -74,7 +43,5 @@ export async function up(knex) {
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("logs");
|
||||
await knex.schema.dropTableIfExists("step_runs");
|
||||
await knex.schema.dropTableIfExists("workflow_runs");
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("step_runs", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.integer("step_index").notNullable();
|
||||
t.text("script").notNullable();
|
||||
t.text("config");
|
||||
t.text("status").notNullable();
|
||||
t.text("started_at").notNullable();
|
||||
t.text("finished_at");
|
||||
t.integer("duration_ms");
|
||||
t.text("output");
|
||||
t.text("error");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("step_runs");
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("logs", (t) => {
|
||||
t.increments("id").primary();
|
||||
t.text("run_id")
|
||||
.notNullable()
|
||||
.references("id")
|
||||
.inTable("workflow_runs")
|
||||
.onDelete("CASCADE");
|
||||
t.text("step_id");
|
||||
t.text("ts").notNullable();
|
||||
t.integer("level").notNullable();
|
||||
t.text("msg");
|
||||
t.text("payload");
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("logs");
|
||||
}
|
||||
+20
-20
@@ -29,35 +29,35 @@ const DEFAULT_EMAIL_TEMPLATE = `<!DOCTYPE html>
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.alterTable("http_pages", (t) => {
|
||||
await knex.schema.createTable("http_pages", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("content").notNullable();
|
||||
t.text("mime").notNullable();
|
||||
t.integer("status").notNullable().defaultTo(200);
|
||||
t.text("kind").notNullable().defaultTo("response");
|
||||
t.integer("system").notNullable().defaultTo(0);
|
||||
t.text("created_at").notNullable();
|
||||
t.text("updated_at").notNullable();
|
||||
});
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const existing = await knex("http_pages").where({ name: "email-default" }).first();
|
||||
if (!existing) {
|
||||
await knex("http_pages").insert({
|
||||
id: "00000000-0000-4000-8000-000000000001",
|
||||
name: "email-default",
|
||||
content: DEFAULT_EMAIL_TEMPLATE,
|
||||
mime: "html",
|
||||
status: 200,
|
||||
kind: "template",
|
||||
system: 1,
|
||||
created_at: now,
|
||||
updated_at: now,
|
||||
});
|
||||
}
|
||||
await knex("http_pages").insert({
|
||||
id: "00000000-0000-4000-8000-000000000001",
|
||||
name: "email-default",
|
||||
content: DEFAULT_EMAIL_TEMPLATE,
|
||||
mime: "html",
|
||||
status: 200,
|
||||
kind: "template",
|
||||
system: 1,
|
||||
created_at: now,
|
||||
updated_at: now,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex("http_pages").where({ name: "email-default", system: 1 }).del();
|
||||
await knex.schema.alterTable("http_pages", (t) => {
|
||||
t.dropColumn("kind");
|
||||
t.dropColumn("system");
|
||||
});
|
||||
await knex.schema.dropTableIfExists("http_pages");
|
||||
}
|
||||
+2
-13
@@ -2,21 +2,11 @@
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("http_pages", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("content").notNullable();
|
||||
t.text("mime").notNullable(); // html | json
|
||||
t.integer("status").notNullable().defaultTo(200);
|
||||
t.text("created_at").notNullable();
|
||||
t.text("updated_at").notNullable();
|
||||
});
|
||||
|
||||
await knex.schema.createTable("http_auths", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("name").notNullable().unique();
|
||||
t.text("type").notNullable(); // bearer | basic | header
|
||||
t.text("config").notNullable(); // JSON
|
||||
t.text("type").notNullable();
|
||||
t.text("config").notNullable();
|
||||
t.integer("unauthorized_status").nullable();
|
||||
t.text("unauthorized_response").nullable();
|
||||
t.text("created_at").notNullable();
|
||||
@@ -29,5 +19,4 @@ export async function up(knex) {
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("http_auths");
|
||||
await knex.schema.dropTableIfExists("http_pages");
|
||||
}
|
||||
-18
@@ -21,29 +21,11 @@ export async function up(knex) {
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_revisions_workflow_id_created_at_idx ON workflow_revisions (workflow_id, created_at DESC)",
|
||||
);
|
||||
|
||||
await knex.schema.createTable("workflow_trash", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("workflow_id").notNullable();
|
||||
t.text("owner").notNullable();
|
||||
t.text("file").notNullable();
|
||||
t.text("name");
|
||||
t.text("deleted_at").notNullable();
|
||||
t.text("trash_path").notNullable();
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_trash_deleted_at_idx ON workflow_trash (deleted_at ASC)",
|
||||
);
|
||||
await knex.schema.raw(
|
||||
"CREATE UNIQUE INDEX workflow_trash_owner_file_idx ON workflow_trash (owner, file)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("workflow_trash");
|
||||
await knex.schema.dropTableIfExists("workflow_revisions");
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function up(knex) {
|
||||
await knex.schema.createTable("workflow_trash", (t) => {
|
||||
t.text("id").primary();
|
||||
t.text("workflow_id").notNullable();
|
||||
t.text("owner").notNullable();
|
||||
t.text("file").notNullable();
|
||||
t.text("name");
|
||||
t.text("deleted_at").notNullable();
|
||||
t.text("trash_path").notNullable();
|
||||
});
|
||||
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX workflow_trash_deleted_at_idx ON workflow_trash (deleted_at ASC)",
|
||||
);
|
||||
await knex.schema.raw(
|
||||
"CREATE UNIQUE INDEX workflow_trash_owner_file_idx ON workflow_trash (owner, file)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
export async function down(knex) {
|
||||
await knex.schema.dropTableIfExists("workflow_trash");
|
||||
}
|
||||
+3
-1
@@ -14,7 +14,9 @@ export async function up(knex) {
|
||||
t.unique(["owner", "name"]);
|
||||
});
|
||||
|
||||
await knex.schema.raw("CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)");
|
||||
await knex.schema.raw(
|
||||
"CREATE INDEX profiles_owner_name_idx ON profiles (owner, name)",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1,268 +0,0 @@
|
||||
import fs from "fs";
|
||||
import path from "path";
|
||||
import { DEFAULT_OWNER } from "@jerapah-flow/shared";
|
||||
import { db } from "./db.js";
|
||||
import { log } from "./logger.js";
|
||||
import { EXAMPLE_WORKFLOWS_DIR, WORKFLOWS_DIR } from "./paths.js";
|
||||
import { TRASH_WORKFLOWS_DIR } from "./workflow-trash.js";
|
||||
import { rewriteLegacyConfigRefsInText } from "./config-ref-rewrite.js";
|
||||
|
||||
const LEGACY_OWNER = "default";
|
||||
|
||||
/**
|
||||
* @param {string} dir
|
||||
* @returns {string[]}
|
||||
*/
|
||||
function listFilesRecursive(dir) {
|
||||
if (!fs.existsSync(dir)) return [];
|
||||
/** @type {string[]} */
|
||||
const out = [];
|
||||
for (const entry of fs.readdirSync(dir, { withFileTypes: true })) {
|
||||
const full = path.join(dir, entry.name);
|
||||
if (entry.isDirectory()) out.push(...listFilesRecursive(full));
|
||||
else out.push(full);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Move files from legacy owner dir into DEFAULT_OWNER. Skip name collisions.
|
||||
* @param {string} rootDir
|
||||
* @returns {{ moved: number, skipped: number }}
|
||||
*/
|
||||
function mergeOwnerDir(rootDir) {
|
||||
const fromDir = path.join(rootDir, LEGACY_OWNER);
|
||||
const toDir = path.join(rootDir, DEFAULT_OWNER);
|
||||
if (!fs.existsSync(fromDir)) return { moved: 0, skipped: 0 };
|
||||
|
||||
fs.mkdirSync(toDir, { recursive: true });
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
|
||||
for (const name of fs.readdirSync(fromDir)) {
|
||||
const from = path.join(fromDir, name);
|
||||
const to = path.join(toDir, name);
|
||||
const st = fs.statSync(from);
|
||||
if (!st.isFile()) {
|
||||
skipped += 1;
|
||||
log.warn({ from }, "owner migrate: skip non-file under legacy owner dir");
|
||||
continue;
|
||||
}
|
||||
if (fs.existsSync(to)) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ from, to },
|
||||
"owner migrate: skip YAML collision (local already has file)",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
fs.renameSync(from, to);
|
||||
moved += 1;
|
||||
}
|
||||
|
||||
const remaining = fs.existsSync(fromDir) ? fs.readdirSync(fromDir) : [];
|
||||
if (remaining.length === 0 && fs.existsSync(fromDir)) {
|
||||
fs.rmdirSync(fromDir);
|
||||
}
|
||||
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} filePath
|
||||
* @returns {boolean}
|
||||
*/
|
||||
function rewriteFileInPlace(filePath) {
|
||||
const raw = fs.readFileSync(filePath, "utf8");
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(raw);
|
||||
if (!changed) return false;
|
||||
fs.writeFileSync(filePath, text, "utf8");
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string} rootDir
|
||||
* @returns {number}
|
||||
*/
|
||||
function rewriteYamlTree(rootDir) {
|
||||
let n = 0;
|
||||
for (const file of listFilesRecursive(rootDir)) {
|
||||
if (!/\.ya?ml$/i.test(file)) continue;
|
||||
if (rewriteFileInPlace(file)) n += 1;
|
||||
}
|
||||
return n;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
* @param {string} table
|
||||
* @param {"name" | "file" | null} uniqueCol
|
||||
*/
|
||||
async function migrateOwnerColumn(knex, table, uniqueCol) {
|
||||
const legacyRows = await knex(table).where({ owner: LEGACY_OWNER }).select("*");
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
for (const row of legacyRows) {
|
||||
if (uniqueCol != null) {
|
||||
const conflict = await knex(table)
|
||||
.where({ owner: DEFAULT_OWNER, [uniqueCol]: row[uniqueCol] })
|
||||
.first();
|
||||
if (conflict) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ table, id: row.id, [uniqueCol]: row[uniqueCol] },
|
||||
"owner migrate: skip row collision",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
await knex(table).where({ id: row.id }).update({ owner: DEFAULT_OWNER });
|
||||
moved += 1;
|
||||
}
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateWorkflowRuns(knex) {
|
||||
const n = await knex("workflow_runs")
|
||||
.where({ owner: LEGACY_OWNER })
|
||||
.update({ owner: DEFAULT_OWNER });
|
||||
return { moved: Number(n) || 0, skipped: 0 };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateWorkflowRevisionsOwner(knex) {
|
||||
const n = await knex("workflow_revisions")
|
||||
.where({ owner: LEGACY_OWNER })
|
||||
.update({ owner: DEFAULT_OWNER });
|
||||
return { moved: Number(n) || 0, skipped: 0 };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function migrateScriptStateNamespaces(knex) {
|
||||
const rows = await knex("script_state")
|
||||
.where("namespace", "like", `${LEGACY_OWNER}/%`)
|
||||
.select("namespace", "key");
|
||||
let moved = 0;
|
||||
let skipped = 0;
|
||||
for (const row of rows) {
|
||||
const nextNs = `${DEFAULT_OWNER}${row.namespace.slice(LEGACY_OWNER.length)}`;
|
||||
const conflict = await knex("script_state")
|
||||
.where({ namespace: nextNs, key: row.key })
|
||||
.first();
|
||||
if (conflict) {
|
||||
skipped += 1;
|
||||
log.warn(
|
||||
{ namespace: row.namespace, key: row.key, nextNs },
|
||||
"owner migrate: skip script_state collision",
|
||||
);
|
||||
continue;
|
||||
}
|
||||
await knex("script_state")
|
||||
.where({ namespace: row.namespace, key: row.key })
|
||||
.update({ namespace: nextNs });
|
||||
moved += 1;
|
||||
}
|
||||
return { moved, skipped };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {import("knex").Knex} knex
|
||||
*/
|
||||
async function rewriteDbConfigStrings(knex) {
|
||||
let profiles = 0;
|
||||
let revisions = 0;
|
||||
|
||||
const profileRows = await knex("profiles").select("id", "config");
|
||||
for (const row of profileRows) {
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(String(row.config ?? ""));
|
||||
if (!changed) continue;
|
||||
await knex("profiles").where({ id: row.id }).update({ config: text });
|
||||
profiles += 1;
|
||||
}
|
||||
|
||||
const revisionRows = await knex("workflow_revisions").select("id", "content");
|
||||
for (const row of revisionRows) {
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(String(row.content ?? ""));
|
||||
if (!changed) continue;
|
||||
await knex("workflow_revisions").where({ id: row.id }).update({ content: text });
|
||||
revisions += 1;
|
||||
}
|
||||
|
||||
return { profiles, revisions };
|
||||
}
|
||||
|
||||
/**
|
||||
* One-shot: move owner `default` → `local`, rewrite prefix refs to mustache.
|
||||
* Idempotent when there is no remaining `default` data / prefix refs.
|
||||
*/
|
||||
export async function migrateDefaultOwnerIfNeeded() {
|
||||
const knex = db;
|
||||
|
||||
const yamlLive = mergeOwnerDir(WORKFLOWS_DIR);
|
||||
const yamlTrash = mergeOwnerDir(TRASH_WORKFLOWS_DIR);
|
||||
|
||||
const variables = await migrateOwnerColumn(knex, "variables", "name");
|
||||
const secrets = await migrateOwnerColumn(knex, "secrets", "name");
|
||||
const profiles = await migrateOwnerColumn(knex, "profiles", "name");
|
||||
const trash = await migrateOwnerColumn(knex, "workflow_trash", "file");
|
||||
const runs = await migrateWorkflowRuns(knex);
|
||||
const revisionsOwner = await migrateWorkflowRevisionsOwner(knex);
|
||||
const scriptState = await migrateScriptStateNamespaces(knex);
|
||||
|
||||
const yamlRewritten =
|
||||
rewriteYamlTree(path.join(WORKFLOWS_DIR, DEFAULT_OWNER)) +
|
||||
rewriteYamlTree(path.join(TRASH_WORKFLOWS_DIR, DEFAULT_OWNER)) +
|
||||
rewriteYamlTree(path.join(WORKFLOWS_DIR, LEGACY_OWNER)) +
|
||||
rewriteYamlTree(path.join(TRASH_WORKFLOWS_DIR, LEGACY_OWNER));
|
||||
|
||||
let examplesRewritten = 0;
|
||||
if (fs.existsSync(EXAMPLE_WORKFLOWS_DIR)) {
|
||||
examplesRewritten = rewriteYamlTree(EXAMPLE_WORKFLOWS_DIR);
|
||||
}
|
||||
|
||||
const dbStrings = await rewriteDbConfigStrings(knex);
|
||||
|
||||
const summary = {
|
||||
yamlLive,
|
||||
yamlTrash,
|
||||
variables,
|
||||
secrets,
|
||||
profiles,
|
||||
trash,
|
||||
runs,
|
||||
revisionsOwner,
|
||||
scriptState,
|
||||
yamlRewritten,
|
||||
examplesRewritten,
|
||||
dbStrings,
|
||||
};
|
||||
|
||||
const touched =
|
||||
yamlLive.moved +
|
||||
yamlTrash.moved +
|
||||
variables.moved +
|
||||
secrets.moved +
|
||||
profiles.moved +
|
||||
trash.moved +
|
||||
runs.moved +
|
||||
revisionsOwner.moved +
|
||||
scriptState.moved +
|
||||
yamlRewritten +
|
||||
examplesRewritten +
|
||||
dbStrings.profiles +
|
||||
dbStrings.revisions >
|
||||
0;
|
||||
|
||||
if (touched) {
|
||||
log.info(summary, "migrated owner default → local and rewrote config refs");
|
||||
}
|
||||
|
||||
return summary;
|
||||
}
|
||||
@@ -13,13 +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:config-refs": "node test/config-refs-smoke.js",
|
||||
"test:owner-migrate": "node test/owner-migrate-smoke.js"
|
||||
"test:config-refs": "node test/config-refs-smoke.js"
|
||||
},
|
||||
"dependencies": {
|
||||
"@jerapah-flow/shared": "workspace:*",
|
||||
|
||||
+39
-17
@@ -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");
|
||||
|
||||
@@ -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" },
|
||||
|
||||
@@ -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,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";
|
||||
|
||||
@@ -22,7 +22,9 @@ const ALLOWED_MODULES = new Set([
|
||||
"basic-ftp",
|
||||
"jsonata",
|
||||
"mustache",
|
||||
"net",
|
||||
"node-html-parser",
|
||||
"node:net",
|
||||
"node:stream",
|
||||
"nodemailer",
|
||||
"rss-parser",
|
||||
@@ -557,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,
|
||||
$workflows: opts.$workflows,
|
||||
pluginDir: opts.pluginDir ?? null,
|
||||
pluginDir,
|
||||
});
|
||||
return { fn, ...extractScriptMeta(fn) };
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -8,6 +8,7 @@ import {
|
||||
import * as fsStore from "../../fs-store.js";
|
||||
import {
|
||||
forkCoreScript,
|
||||
duplicatePlugin,
|
||||
listCoreScriptNames,
|
||||
listInstalledPlugins,
|
||||
resolveScriptRef,
|
||||
@@ -29,7 +30,7 @@ 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";
|
||||
|
||||
@@ -251,6 +252,40 @@ export default function scriptsPluginFactory(registry) {
|
||||
}
|
||||
});
|
||||
|
||||
fastify.post("/scripts/:name/duplicate", async (req, reply) => {
|
||||
const rawName = decodeURIComponent(
|
||||
/** @type {{ name: string }} */ (req.params).name,
|
||||
);
|
||||
const body = /** @type {{ id?: string, description?: string }} */ (
|
||||
req.body ?? {}
|
||||
);
|
||||
if (typeof body.id !== "string" || !body.id.trim()) {
|
||||
return reply.code(400).send({ error: "id is required" });
|
||||
}
|
||||
const parsed = parsePluginScriptRef(rawName);
|
||||
if (!parsed) {
|
||||
return reply
|
||||
.code(400)
|
||||
.send({ error: "name must be a plugin ref (plugin/<id>)" });
|
||||
}
|
||||
try {
|
||||
const installed = duplicatePlugin(parsed.id, body.id.trim(), {
|
||||
description: body.description,
|
||||
});
|
||||
clearScriptCache();
|
||||
return reply.code(201).send({
|
||||
...installed,
|
||||
restartNeeded: true,
|
||||
warning:
|
||||
"Plugins run as the JerapahFlow process user. Review code before install.",
|
||||
});
|
||||
} catch (err) {
|
||||
return reply
|
||||
.code(/** @type {any} */ (err).statusCode ?? 500)
|
||||
.send({ error: err instanceof Error ? err.message : String(err) });
|
||||
}
|
||||
});
|
||||
|
||||
fastify.post("/scripts/:name/dry-run", async (req, reply) => {
|
||||
const rawName = decodeURIComponent(
|
||||
/** @type {{ name: string }} */ (req.params).name,
|
||||
@@ -431,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 ?? {}
|
||||
@@ -463,7 +497,6 @@ export default function scriptsPluginFactory(registry) {
|
||||
|
||||
fastify.post(
|
||||
"/plugins/install",
|
||||
{ onRequest: [fastify.requireAdmin] },
|
||||
async (req, reply) => {
|
||||
const body = /** @type {{
|
||||
source?: string,
|
||||
@@ -538,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 {
|
||||
|
||||
@@ -31,8 +31,6 @@ import {
|
||||
getRedisUrlForLog,
|
||||
} from "./workflow-queue.js";
|
||||
import { purgeExpiredTrash } from "./workflow-trash.js";
|
||||
import { migrateLegacyWorkflowsIfNeeded } from "./workflow-migrate.js";
|
||||
import { migrateDefaultOwnerIfNeeded } from "./owner-migrate.js";
|
||||
import {
|
||||
getConfigGeneration,
|
||||
startHeartbeatLoop,
|
||||
@@ -130,18 +128,6 @@ export async function startApp(opts = {}) {
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
migrateLegacyWorkflowsIfNeeded();
|
||||
} catch (err) {
|
||||
log.warn({ err }, "legacy workflow migrate failed");
|
||||
}
|
||||
|
||||
try {
|
||||
await migrateDefaultOwnerIfNeeded();
|
||||
} catch (err) {
|
||||
log.warn({ err }, "owner default→local migrate failed");
|
||||
}
|
||||
|
||||
const registry = createRegistry(server, {
|
||||
queue: workflowQueue,
|
||||
// Cron + HTTP triggers enqueue jobs; only the API process may own them.
|
||||
|
||||
@@ -2,8 +2,7 @@ import { migrate, db } from "../db.js";
|
||||
import { upsertSecret, deleteSecret } from "../secrets-store.js";
|
||||
import { deleteVariable, upsertVariable } from "../variables-store.js";
|
||||
import { Secret } from "../secret-value.js";
|
||||
import { parseLegacyConfigRef, resolveConfigRefs } from "../config-refs.js";
|
||||
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
|
||||
import { resolveConfigRefs } from "../config-refs.js";
|
||||
|
||||
await migrate();
|
||||
|
||||
@@ -27,50 +26,17 @@ async function assertRejects(fn, match) {
|
||||
const owner = "config_refs_smoke_owner";
|
||||
const ctx = { owner, workflowKey: `${owner}/config-refs-smoke.yaml`, context: {}, data: {} };
|
||||
|
||||
{
|
||||
const rewritten = rewriteLegacyConfigRefsInText(
|
||||
'url: $VAR_ntfy\ntoken: $SECRET_tok\nid: $CONTEXT_user',
|
||||
);
|
||||
assert(rewritten.changed, "rewrite detects legacy refs");
|
||||
assert(
|
||||
rewritten.text.includes("{{ vars.ntfy }}") &&
|
||||
rewritten.text.includes("{{ secrets.tok }}") &&
|
||||
rewritten.text.includes("{{ context.user }}"),
|
||||
"rewrite maps prefixes",
|
||||
);
|
||||
assert(!rewriteLegacyConfigRefsInText("passwordSecret: gmail_app").changed, "bare names untouched");
|
||||
}
|
||||
|
||||
assert(parseLegacyConfigRef("$VAR_ntfy_url")?.kind === "var", "legacy var parse");
|
||||
assert(parseLegacyConfigRef("$SECRET_x")?.kind === "secret", "legacy secret parse");
|
||||
assert(parseLegacyConfigRef("$CONTEXT_token")?.kind === "context", "legacy context parse");
|
||||
assert(parseLegacyConfigRef("{{ vars.x }}") == null, "mustache is not legacy");
|
||||
assert(parseLegacyConfigRef("Bearer $SECRET_x") == null, "mid-string not whole-value legacy");
|
||||
|
||||
{
|
||||
const literal = await resolveConfigRefs("password123", ctx);
|
||||
assert(literal === "password123", "literal passthrough");
|
||||
const unknown = await resolveConfigRefs("$FOO_bar", ctx);
|
||||
assert(unknown === "$FOO_bar", "$FOO_bar stays literal");
|
||||
const embedded = await resolveConfigRefs("Bearer $SECRET_x", ctx);
|
||||
assert(embedded === "Bearer $SECRET_x", "mid-string legacy stays literal");
|
||||
assert(embedded === "Bearer $SECRET_x", "mid-string stays literal");
|
||||
const number = await resolveConfigRefs(42, ctx);
|
||||
assert(number === 42, "number passthrough");
|
||||
}
|
||||
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$VAR_ntfy_channel", ctx),
|
||||
"removed",
|
||||
);
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$SECRET_tok", ctx),
|
||||
"use {{ secrets.tok }}",
|
||||
);
|
||||
await assertRejects(
|
||||
() => resolveConfigRefs("$CONTEXT_token", ctx),
|
||||
"use {{ context.token }}",
|
||||
);
|
||||
|
||||
const created = [];
|
||||
try {
|
||||
const secret = await upsertSecret({
|
||||
|
||||
@@ -1,30 +0,0 @@
|
||||
import { rewriteLegacyConfigRefsInText } from "../config-ref-rewrite.js";
|
||||
import { migrateDefaultOwnerIfNeeded } from "../owner-migrate.js";
|
||||
import { migrate, db } from "../db.js";
|
||||
|
||||
function assert(cond, msg) {
|
||||
if (!cond) throw new Error(msg);
|
||||
}
|
||||
|
||||
{
|
||||
const { text, changed } = rewriteLegacyConfigRefsInText(
|
||||
"a: $VAR_x\nb: $SECRET_y\nc: $CONTEXT_z\nd: passwordSecret: gmail_app",
|
||||
);
|
||||
assert(changed, "detects legacy");
|
||||
assert(text.includes("{{ vars.x }}"), "var rewrite");
|
||||
assert(text.includes("{{ secrets.y }}"), "secret rewrite");
|
||||
assert(text.includes("{{ context.z }}"), "context rewrite");
|
||||
assert(text.includes("passwordSecret: gmail_app"), "bare secret name untouched");
|
||||
}
|
||||
|
||||
await migrate();
|
||||
const first = await migrateDefaultOwnerIfNeeded();
|
||||
const second = await migrateDefaultOwnerIfNeeded();
|
||||
assert(second.secrets.moved === 0, "second pass moves no secrets");
|
||||
assert(second.runs.moved === 0, "second pass moves no runs");
|
||||
await db.destroy();
|
||||
|
||||
console.log("owner-migrate smoke passed", {
|
||||
firstSecretsMoved: first.secrets.moved,
|
||||
secondSecretsMoved: second.secrets.moved,
|
||||
});
|
||||
@@ -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");
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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({
|
||||
|
||||
@@ -6,32 +6,19 @@ import { useVariables } from "../../api/hooks.js";
|
||||
const VAR_PEEK_MAX = 48;
|
||||
const MUSTACHE_RE =
|
||||
/\{\{\s*([A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z0-9_]+)*)\s*\}\}/g;
|
||||
const LEGACY_PREFIX_RE = /^\s*\$(VAR|SECRET|CONTEXT)_([A-Za-z0-9._-]*)\s*$/;
|
||||
|
||||
/**
|
||||
* @param {unknown} value
|
||||
* @returns {{
|
||||
* kind: "mustache" | "legacy",
|
||||
* kind: "mustache",
|
||||
* path?: string,
|
||||
* root?: string,
|
||||
* name?: string,
|
||||
* legacyKind?: string,
|
||||
* legacyName?: string,
|
||||
* } | null}
|
||||
*/
|
||||
function describeConfigRef(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const trimmed = value.trim();
|
||||
const legacy = LEGACY_PREFIX_RE.exec(trimmed);
|
||||
if (legacy) {
|
||||
const legacyKind =
|
||||
legacy[1] === "VAR" ? "variable" : legacy[1] === "SECRET" ? "secret" : "context";
|
||||
return {
|
||||
kind: "legacy",
|
||||
legacyKind,
|
||||
legacyName: legacy[2] ?? "",
|
||||
};
|
||||
}
|
||||
MUSTACHE_RE.lastIndex = 0;
|
||||
const match = MUSTACHE_RE.exec(trimmed);
|
||||
if (!match) return null;
|
||||
@@ -90,20 +77,6 @@ export function ConfigRefHint({ value, owner }) {
|
||||
|
||||
if (!ref) return null;
|
||||
|
||||
if (ref.kind === "legacy") {
|
||||
const hint =
|
||||
ref.legacyKind === "variable"
|
||||
? `{{ vars.${ref.legacyName || "name"} }}`
|
||||
: ref.legacyKind === "secret"
|
||||
? `{{ secrets.${ref.legacyName || "name"} }}`
|
||||
: `{{ context.${ref.legacyName || "name"} }}`;
|
||||
return (
|
||||
<p className="text-xs text-error">
|
||||
legacy ${ref.legacyKind} ref — use <span className="font-mono">{hint}</span>
|
||||
</p>
|
||||
);
|
||||
}
|
||||
|
||||
if (!isVar) {
|
||||
const label =
|
||||
ref.root === "secrets"
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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/<id></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}?` : ""}
|
||||
|
||||
@@ -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;
|
||||
+67
-28
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user