diff --git a/.gitignore b/.gitignore index 943a64c..8312f02 100644 --- a/.gitignore +++ b/.gitignore @@ -7,6 +7,8 @@ logs/ packages/web/dist # Personal/local scripts and workflows (not for the repo) -packages/server/scripts/dev-* -packages/server/workflows/**/dev-* debug-*.js + +# Plugin install staging and per-plugin deps +plugins/.staging-* +plugins/*/node_modules/ diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..62d0cb9 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,144 @@ +# JerapahFlow plugins + +This file tells agents how to add a **user plugin**. Do not put personal or site-specific scripts in `packages/server/scripts/` (core, read-only). Do not put them in `examples/plugins/` (shipped examples only). + +## Where things live + +| Kind | YAML `script` | Editable | Path | +|---|---|---|---| +| Core | `fetch-http.js` | No | `packages/server/scripts/` | +| User plugin | `plugin/` | Yes | `plugins//` | +| Example source | install → `plugin/` | After install | `examples/plugins//` | + +Runtime load path is repo-root `plugins/` (`PLUGINS_DIR` in `packages/server/paths.js`). Override with `JFLOW_PLUGINS_DIR` only in tests. + +Plugins are **outside** the pnpm workspace (`packages/*`). Do not add `plugins/*` to `pnpm-workspace.yaml`. + +## When to create a plugin + +Create a user plugin when the script is: + +- Site-specific (LAN IPs, personal modems, private APIs) +- A fork of a core script the user wants to edit +- Anything that should stay in git but not ship as core + +Use native `fetch` (not `$axios`) when the URL is RFC1918 / WG (`10.x`, `192.168.x`, …). `$axios` is screened and **blocks** those hosts. + +## Create a new plugin + +1. Pick an id: lowercase letters, numbers, hyphens; max 64 chars; `^[a-z0-9]+(?:-[a-z0-9]+)*$`. +2. Folder name **must equal** manifest `id`. Example: `plugins/joplin-api`. +3. Id **must not** collide with a core script basename (`fetch-http`, `ntfy`, …). +4. Create three files (see below). Prefer `main: "script.js"`. +5. Point workflows at `plugin/`. +6. Restart the runner (`pnpm dev` / drain-restart under `pnpm dev:pm2`) so the plugin is picked up. + +Copy the layout from `plugins/joplin-api`, `plugins/send-sms`, or `examples/plugins/get-current-time`. + +### `plugins//jerapah-plugin.json` + +```json +{ + "id": "my-plugin", + "name": "My plugin", + "version": "0.1.0", + "jerapah": ">=0.1.0 <1.0.0", + "main": "script.js", + "description": "One-line description" +} +``` + +- `version` must be semver (`0.1.0`). +- `jerapah` must match the app (`0.1.0` in root `package.json`). Use `">=0.1.0 <1.0.0"` unless you know otherwise. +- `main` is a relative path; no `..`, not absolute. + +### `plugins//package.json` + +```json +{ + "name": "jflow-plugin-my-plugin", + "version": "0.1.0", + "private": true, + "type": "module", + "description": "JerapahFlow plugin: …" +} +``` + +Add `dependencies` only if the script `require()`s extra npm packages. Then: + +```bash +pnpm install --dir plugins/ --ignore-scripts --prefer-offline +``` + +`plugins/*/node_modules/` is gitignored. Host-allowlisted modules (`axios`, `jsonata`, …) come from the server; extra deps resolve from the plugin directory. + +### `plugins//script.js` + +Must `export default` a function. The sandbox rewrites ESM `import`/`export default` to CJS. + +```js +function passContext(ctx) { + if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) { + return { ...ctx.context }; + } + return {}; +} + +async function myPlugin(ctx) { + const output = { ok: true }; + return { output, context: { ...passContext(ctx), ...output } }; +} + +myPlugin.meta = { + description: "What this step does", + previewConfigKey: "url", + tags: ["HTTP"], + config: {}, + input: {}, + output: { ok: { type: "boolean" } }, + context: { ok: { type: "boolean" } }, + example: { data: {}, config: {} }, +}; + +export default myPlugin; +``` + +Return `{ output, context?, skipRemaining? }`. Do not return `ctx`. Mutations of `ctx.data` / `ctx.context` are discarded unless returned. + +| Field | Meaning | +|---|---| +| `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`) | +| `output` | Next step’s `data` | +| `context` | Next clipboard. Omit to keep incoming | + +`fn.meta` must be JSON-serializable (UI + dry-run). Include `config` / `input` / `output` field schemas and an `example`. + +## Sandbox globals (do not import these) + +Injected: `log` (pino), `console`, `fetch`, `require`, `$axios`, `$kv`, `$fingerprint`, `$secrets`, `$vars`, `$responses`, `$workflows`. + +- Use `log.info({ … }, "my-plugin: …")` — `log` is not an import. +- `require("axios")` is the screened `$axios` (RFC1918 blocked). Prefer `fetch` for LAN. +- Plugin `require("some-npm-dep")` uses the plugin’s `node_modules`. + +## Workflow YAML + +```yaml +scripts: + - script: plugin/my-plugin + config: + url: http://10.8.0.6:3030/notes + token: $SECRET_joplin_api_token +``` + +Canonical ref is `plugin/` (`.js` suffix is optional). + +## Do not + +- Add user plugins under `packages/server/scripts/` or `examples/plugins/`. +- 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`. diff --git a/README.md b/README.md index eff3864..d12a67d 100644 --- a/README.md +++ b/README.md @@ -13,14 +13,45 @@ Workflow runner with a sandboxed script engine, SQLite run history, and an admin ```bash pnpm install +# Redis required for the workflow queue pnpm dev ``` +- UI (dev): http://localhost:8500 - API: http://localhost:8700 -- UI (dev): http://localhost:5173 The first account created becomes **admin**. Later accounts are created from Users. +### Process modes + +| Command | Processes | Ports | +|---|---|---| +| `pnpm dev` | Monolith (`runner.js` = API + worker) + Vite | UI **8500**, API **8700** | +| `pnpm dev:pm2` | Control + PM2 HTTP + PM2 workers + Vite | UI **8500**, control **8600**, API **8700** | + +`pnpm dev:pm2` is the mode for Ops (start/stop HTTP, scale workers, drain restart). Control owns SQLite migrations; HTTP/workers do not migrate. + +## Scripts (core vs plugins) + +| Kind | Name in YAML | Editable | Location | +|---|---|---|---| +| **Core** | `fetch-http.js`, `s3.js`, … | No (fork only) | `packages/server/scripts/` | +| **Plugin** | `plugin/` | Yes | `plugins//` | + +- App version is **`0.1.0`** (root `package.json`). Plugin manifests declare `jerapah: ">=0.1.0 <1.0.0"`. +- Install plugins via admin API: zip (base64), HTTPS git URL, example, or fork a core script. +- Install/update/uninstall sets **restart-needed** — drain-restart HTTP + workers under `pnpm dev:pm2`. +- Example plugin: `examples/plugins/get-current-time` → `plugin/get-current-time`. +- User plugins in this repo: `plugins/joplin-api` → `plugin/joplin-api`, `plugins/send-sms` → `plugin/send-sms`. + +```bash +# Smoke +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 +``` + + ## Script contract Each script is `async function main(ctx)` and **must** return: @@ -50,13 +81,31 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / | Command | Description | |---|---| -| `pnpm dev` | Server + Vite together | -| `pnpm dev:server` | API/runner only | -| `pnpm dev:web` | UI only (proxies `/api` to :8700) | +| `pnpm dev` | Monolith server + Vite (no PM2) | +| `pnpm dev:pm2` | Control + PM2 HTTP/workers + Vite (Ops UI) | +| `pnpm dev:server` | Monolith API/runner only | +| `pnpm dev:web` | UI only (proxies `/api` → :8700, `/ops` → :8600) | | `pnpm build` | Production UI build | -| `pnpm start` | Serve API and built UI from :8700 | +| `pnpm start` | Monolith: API + worker + built UI | +| `pnpm start:control` | Control plane only (migrates, manages PM2 children) | +| `pnpm start:api` | HTTP API + cron enqueue (`JFLOW_ROLE=api`) | +| `pnpm start:worker` | BullMQ worker only | | `pnpm migrate` | Apply SQLite migrations | +## Ops (control plane) + +Admin UI route **Ops** (`/ops`) talks to the control process. + +| Action | Behavior | +|---|---| +| Pause / resume | BullMQ `queue.pause()` / `resume()` — cron/HTTP still enqueue | +| Reload workflows | Redis pub/sub → all live HTTP/worker processes re-read YAML | +| Scale workers | PM2 scale; scale-down drains active jobs unless `force` | +| 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. + ## Environment | Variable | Default | Notes | @@ -68,11 +117,13 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / | `REDIS_PASS` | — | Optional Redis AUTH password (sent via ioredis `password`). Prefer this over embedding credentials in `REDIS_URL` so logs stay clean. | | `JFLOW_QUEUE_NAME` | `jerapah-workflows` | BullMQ queue name. | | `JFLOW_WORKER_CONCURRENCY` | `5` | Max parallel workflow jobs per worker process. | -| `JFLOW_ROLE` | `all` | Which duties this process performs: `all` (HTTP/admin + cron producer + worker), `api` (HTTP/admin + cron enqueue only), or `worker` (consume queue only). Use separate processes in production when you want to scale workers independently. | +| `JFLOW_ROLE` | `all` | `all` (HTTP + cron + worker), `api`, or `worker`. Prefer `pnpm start:api` / `start:worker` under control. | +| `JFLOW_CONFIG_GENERATION` | `1` | Set by control/PM2 so children report config generation in heartbeats. | +| `JFLOW_CONTROL_PORT` | `8600` | Control ops API port. | | `JFLOW_LOG_LEVEL` | `debug` | Pino level | | `JFLOW_RETENTION_DAYS` | `30` | Run history prune | -| `JFLOW_CORS_ORIGIN` | `http://localhost:5173` | Vite origin in dev | -| `PORT` | `8700` | HTTP port | +| `JFLOW_CORS_ORIGIN` | `http://localhost:8500` | Vite origin in dev | +| `PORT` | `8700` | HTTP API port | | `NODE_ENV` | — | Set `production` for secure cookies | Workflow runs are **queued** via BullMQ. HTTP and manual triggers return `202 { runId, status: "queued" }` immediately; poll `GET /api/runs/:id` for progress (`queued` → `running` → `success` \| `failed`). Cron remains an in-process producer that enqueues jobs on each tick. @@ -83,7 +134,10 @@ Workflow runs are **queued** via BullMQ. HTTP and manual triggers return `202 { pnpm install pnpm build # Redis must be reachable at REDIS_URL (set REDIS_PASS if Redis requires AUTH) -JFLOW_JWT_SECRET=... JFLOW_SECRETS_KEY=... REDIS_URL=redis://127.0.0.1:6379 REDIS_PASS=... NODE_ENV=production pnpm start +# Recommended: run control (migrates + manages PM2 HTTP/workers) +JFLOW_JWT_SECRET=... JFLOW_SECRETS_KEY=... REDIS_URL=redis://127.0.0.1:6379 REDIS_PASS=... NODE_ENV=production pnpm start:control +# Or monolith (dev-style): +# ... pnpm start ``` -The server serves `packages/web/dist` when that folder exists. +With control, serve the built UI from Vite preview, a reverse proxy, or set `JFLOW_SERVE_UI=1` on the HTTP process. diff --git a/examples/plugins/get-current-time/jerapah-plugin.json b/examples/plugins/get-current-time/jerapah-plugin.json new file mode 100644 index 0000000..35294d2 --- /dev/null +++ b/examples/plugins/get-current-time/jerapah-plugin.json @@ -0,0 +1,8 @@ +{ + "id": "get-current-time", + "name": "Get current time", + "version": "0.1.0", + "jerapah": ">=0.1.0 <1.0.0", + "main": "script.js", + "description": "Example user plugin: return the current time as output.datetime" +} diff --git a/examples/plugins/get-current-time/package.json b/examples/plugins/get-current-time/package.json new file mode 100644 index 0000000..2ca2114 --- /dev/null +++ b/examples/plugins/get-current-time/package.json @@ -0,0 +1,7 @@ +{ + "name": "jflow-plugin-get-current-time", + "version": "0.1.0", + "private": true, + "type": "module", + "description": "Example JerapahFlow plugin (no npm dependencies)" +} diff --git a/packages/server/scripts/get-current-time.js b/examples/plugins/get-current-time/script.js similarity index 100% rename from packages/server/scripts/get-current-time.js rename to examples/plugins/get-current-time/script.js diff --git a/package.json b/package.json index d001c21..8cc5674 100644 --- a/package.json +++ b/package.json @@ -1,15 +1,19 @@ { "name": "jerapah-flow", - "version": "1.0.0", + "version": "0.1.0", "private": true, "description": "Script/workflow runner with admin UI", "type": "module", "scripts": { "dev": "pnpm --filter @jerapah-flow/server --filter @jerapah-flow/web --parallel dev", + "dev:pm2": "pnpm --filter @jerapah-flow/server dev:pm2", "dev:server": "pnpm --filter @jerapah-flow/server dev", "dev:web": "pnpm --filter @jerapah-flow/web dev", "build": "pnpm --filter @jerapah-flow/web build", "start": "pnpm --filter @jerapah-flow/server start", + "start:api": "pnpm --filter @jerapah-flow/server start:api", + "start:worker": "pnpm --filter @jerapah-flow/server start:worker", + "start:control": "pnpm --filter @jerapah-flow/server start:control", "migrate": "pnpm --filter @jerapah-flow/server migrate" }, "packageManager": "pnpm@10.25.0", diff --git a/packages/server/app-version.js b/packages/server/app-version.js new file mode 100644 index 0000000..c81629f --- /dev/null +++ b/packages/server/app-version.js @@ -0,0 +1,79 @@ +import fs from "fs"; +import path from "path"; +import { fileURLToPath } from "url"; +import { SERVER_ROOT } from "./paths.js"; + +const ROOT_PKG = path.resolve(SERVER_ROOT, "../../package.json"); + +/** + * JerapahFlow app version from the monorepo root package.json. + * @returns {string} + */ +export function getAppVersion() { + try { + const raw = JSON.parse(fs.readFileSync(ROOT_PKG, "utf8")); + if (typeof raw.version === "string" && raw.version.trim()) { + return raw.version.trim(); + } + } catch { + // fall through + } + return "0.1.0"; +} + +/** + * @param {string} version + * @returns {[number, number, number]} + */ +function parseSemver(version) { + const cleaned = String(version).trim().replace(/^v/i, ""); + const core = cleaned.split("-")[0].split("+")[0]; + const parts = core.split(".").map((p) => Number(p)); + return [parts[0] || 0, parts[1] || 0, parts[2] || 0]; +} + +/** + * @param {[number, number, number]} a + * @param {[number, number, number]} b + */ +function cmp(a, b) { + for (let i = 0; i < 3; i++) { + if (a[i] !== b[i]) return a[i] < b[i] ? -1 : 1; + } + return 0; +} + +/** + * Minimal semver range check for patterns used in manifests: + * `1.2.3`, `>=0.1.0`, `<1.0.0`, `>=0.1.0 <1.0.0` + * + * @param {string} version + * @param {string} range + * @returns {boolean} + */ +export function satisfiesRange(version, range) { + if (typeof range !== "string" || !range.trim()) return false; + const ver = parseSemver(version); + const tokens = range.trim().split(/\s+/); + for (const token of tokens) { + if (/^\d+\.\d+\.\d+/.test(token) && !token.startsWith(">") && !token.startsWith("<")) { + if (cmp(ver, parseSemver(token)) !== 0) return false; + continue; + } + const m = /^(>=|<=|>|<|=)?\s*v?(\d+\.\d+\.\d+(?:[-+][0-9A-Za-z.-]+)?)$/.exec( + token, + ); + if (!m) return false; + const op = m[1] || "="; + const bound = parseSemver(m[2]); + const c = cmp(ver, bound); + if (op === ">=" && c < 0) return false; + if (op === "<=" && c > 0) return false; + if (op === ">" && c <= 0) return false; + if (op === "<" && c >= 0) return false; + if (op === "=" && c !== 0) return false; + } + return true; +} + +void fileURLToPath; diff --git a/packages/server/control-bus.js b/packages/server/control-bus.js new file mode 100644 index 0000000..136fd0b --- /dev/null +++ b/packages/server/control-bus.js @@ -0,0 +1,142 @@ +import IORedis from "ioredis"; +import { log } from "./logger.js"; +import { + getRedisPassword, + getRedisUrl, + getSharedConnection, +} from "./workflow-queue.js"; + +export const CHANNEL_RELOAD = "jflow:reload"; +export const HEARTBEAT_KEY = "jflow:heartbeats"; +export const HEARTBEAT_TTL_SEC = 20; +export const HEARTBEAT_INTERVAL_MS = 5_000; + +/** + * @returns {number} + */ +export function getConfigGeneration() { + const raw = Number(process.env.JFLOW_CONFIG_GENERATION ?? 1); + if (!Number.isFinite(raw) || raw < 1) return 1; + return Math.floor(raw); +} + +/** + * Dedicated Redis connection for pub/sub (ioredis cannot mix pub/sub with other commands). + * @returns {IORedis} + */ +export function createPubSubConnection() { + /** @type {import("ioredis").RedisOptions} */ + const options = { maxRetriesPerRequest: null, enableReadyCheck: true }; + const password = getRedisPassword(); + if (password) options.password = password; + const conn = new IORedis(getRedisUrl(), options); + conn.on("error", (err) => { + log.error({ err }, "redis pub/sub connection error"); + }); + return conn; +} + +/** + * @param {string} role + * @param {{ pid?: number, hostname?: string }} [extra] + */ +export async function writeHeartbeat(role, extra = {}) { + const conn = getSharedConnection(); + const payload = JSON.stringify({ + role, + pid: extra.pid ?? process.pid, + hostname: extra.hostname ?? process.env.HOSTNAME ?? "local", + generation: getConfigGeneration(), + ts: Date.now(), + }); + const field = `${role}:${process.pid}`; + await conn.hset(HEARTBEAT_KEY, field, payload); + await conn.expire(HEARTBEAT_KEY, HEARTBEAT_TTL_SEC * 3); +} + +/** + * @returns {Promise>} + */ +export async function readHeartbeats() { + const conn = getSharedConnection(); + const all = await conn.hgetall(HEARTBEAT_KEY); + const now = Date.now(); + /** @type {Array} */ + const out = []; + for (const [field, raw] of Object.entries(all)) { + try { + const parsed = JSON.parse(raw); + const ts = Number(parsed.ts) || 0; + out.push({ + field, + role: String(parsed.role ?? "unknown"), + pid: Number(parsed.pid) || 0, + hostname: String(parsed.hostname ?? ""), + generation: Number(parsed.generation) || 0, + ts, + stale: now - ts > HEARTBEAT_TTL_SEC * 1000, + }); + } catch { + // skip bad rows + } + } + return out; +} + +/** + * @param {{ type?: string }} [payload] + */ +export async function publishReload(payload = { type: "workflows" }) { + const conn = getSharedConnection(); + await conn.publish(CHANNEL_RELOAD, JSON.stringify(payload)); +} + +/** + * @param {(msg: { type?: string }) => void | Promise} handler + * @returns {Promise<{ stop: () => Promise }>} + */ +export async function subscribeReload(handler) { + const sub = createPubSubConnection(); + await sub.subscribe(CHANNEL_RELOAD); + sub.on("message", (_channel, message) => { + let parsed = { type: "workflows" }; + try { + parsed = JSON.parse(message); + } catch { + // use default + } + void Promise.resolve(handler(parsed)).catch((err) => { + log.error({ err }, "reload handler failed"); + }); + }); + return { + async stop() { + await sub.unsubscribe(CHANNEL_RELOAD).catch(() => {}); + await sub.quit().catch(() => sub.disconnect()); + }, + }; +} + +/** + * Start periodic heartbeats. Returns a stop function. + * @param {string} role + */ +export function startHeartbeatLoop(role) { + const tick = () => { + void writeHeartbeat(role).catch((err) => { + log.error({ err, role }, "heartbeat failed"); + }); + }; + tick(); + const timer = setInterval(tick, HEARTBEAT_INTERVAL_MS); + timer.unref?.(); + return () => clearInterval(timer); +} diff --git a/packages/server/control-state.js b/packages/server/control-state.js new file mode 100644 index 0000000..52ab475 --- /dev/null +++ b/packages/server/control-state.js @@ -0,0 +1,164 @@ +import fs from "fs"; +import path from "path"; +import { DATA_DIR } from "./paths.js"; + +const STATE_PATH = path.join(DATA_DIR, "control-state.json"); +const LOCK_PATH = path.join(DATA_DIR, "ops.lock"); + +/** + * @typedef {{ + * http: "running" | "stopped", + * workers: number, + * queuePaused: boolean, + * generation: number, + * restartNeeded: boolean, + * restartReason: string | null, + * }} ControlState + */ + +/** @returns {ControlState} */ +export function defaultControlState() { + return { + http: "running", + workers: 1, + queuePaused: false, + generation: 1, + restartNeeded: false, + restartReason: null, + }; +} + +/** + * @returns {ControlState} + */ +export function readControlState() { + fs.mkdirSync(DATA_DIR, { recursive: true }); + if (!fs.existsSync(STATE_PATH)) { + const initial = defaultControlState(); + writeControlState(initial); + return initial; + } + try { + const raw = JSON.parse(fs.readFileSync(STATE_PATH, "utf8")); + const base = defaultControlState(); + return { + http: raw.http === "stopped" ? "stopped" : "running", + workers: Math.max(0, Math.min(32, Number(raw.workers) || 1)), + queuePaused: Boolean(raw.queuePaused), + generation: Math.max(1, Math.floor(Number(raw.generation) || 1)), + restartNeeded: Boolean(raw.restartNeeded), + restartReason: + typeof raw.restartReason === "string" ? raw.restartReason : null, + }; + } catch { + return defaultControlState(); + } +} + +/** + * @param {ControlState} state + */ +export function writeControlState(state) { + fs.mkdirSync(DATA_DIR, { recursive: true }); + const tmp = `${STATE_PATH}.tmp`; + fs.writeFileSync(tmp, `${JSON.stringify(state, null, 2)}\n`, "utf8"); + fs.renameSync(tmp, STATE_PATH); +} + +/** + * @param {Partial} patch + * @returns {ControlState} + */ +export function patchControlState(patch) { + const next = { ...readControlState(), ...patch }; + writeControlState(next); + return next; +} + +/** + * Bump config generation and mark restart needed. + * @param {string} reason + * @returns {ControlState} + */ +export function bumpGeneration(reason) { + const cur = readControlState(); + return patchControlState({ + generation: cur.generation + 1, + restartNeeded: true, + restartReason: reason, + }); +} + +/** + * Clear restart-needed after processes match generation. + * @returns {ControlState} + */ +export function clearRestartNeeded() { + return patchControlState({ + restartNeeded: false, + restartReason: null, + }); +} + +/** + * @param {string} owner + * @param {number} [ttlMs] + * @returns {{ ok: true, token: string } | { ok: false, error: string, holder?: string }} + */ +export function tryAcquireOpsLock(owner, ttlMs = 120_000) { + fs.mkdirSync(DATA_DIR, { recursive: true }); + const now = Date.now(); + if (fs.existsSync(LOCK_PATH)) { + try { + const existing = JSON.parse(fs.readFileSync(LOCK_PATH, "utf8")); + if (existing.expiresAt > now) { + return { + ok: false, + error: "ops lock held", + holder: String(existing.owner ?? "unknown"), + }; + } + } catch { + // stale/corrupt lock — overwrite + } + } + const token = `${owner}:${now}:${Math.random().toString(36).slice(2)}`; + const payload = { + owner, + token, + expiresAt: now + ttlMs, + }; + fs.writeFileSync(LOCK_PATH, `${JSON.stringify(payload)}\n`, "utf8"); + return { ok: true, token }; +} + +/** + * @param {string} token + */ +export function releaseOpsLock(token) { + if (!fs.existsSync(LOCK_PATH)) return; + try { + const existing = JSON.parse(fs.readFileSync(LOCK_PATH, "utf8")); + if (existing.token !== token) return; + } catch { + // ignore + } + fs.unlinkSync(LOCK_PATH); +} + +/** + * @param {string} token + * @param {number} [ttlMs] + */ +export function refreshOpsLock(token, ttlMs = 120_000) { + if (!fs.existsSync(LOCK_PATH)) return false; + try { + const existing = JSON.parse(fs.readFileSync(LOCK_PATH, "utf8")); + if (existing.token !== token) return false; + existing.expiresAt = Date.now() + ttlMs; + fs.writeFileSync(LOCK_PATH, `${JSON.stringify(existing)}\n`, "utf8"); + return true; + } catch { + return false; + } +} diff --git a/packages/server/control.js b/packages/server/control.js new file mode 100644 index 0000000..b89c38a --- /dev/null +++ b/packages/server/control.js @@ -0,0 +1,427 @@ +import fastify from "fastify"; +import cookie from "@fastify/cookie"; +import cors from "@fastify/cors"; +import jwt from "@fastify/jwt"; +import { migrate, db } from "./db.js"; +import { log, enableLogPersistence, flushLogs } from "./logger.js"; +import { COOKIE } from "./src/api/auth.js"; +import { + clearRestartNeeded, + bumpGeneration, + patchControlState, + readControlState, + refreshOpsLock, + releaseOpsLock, + tryAcquireOpsLock, +} from "./control-state.js"; +import { + getConfigGeneration, + publishReload, + readHeartbeats, +} from "./control-bus.js"; +import { + closeRedis, + createWorkflowQueue, + getRedisUrlForLog, +} from "./workflow-queue.js"; +import { reconcileOrphanRuns } from "./orphan-runs.js"; +import { + connectPm2, + deletePm2App, + describeChildren, + disconnectPm2, + ensureHttp, + ensureWorkers, + PM2_HTTP_NAME, + PM2_WORKER_NAME, + recreateChildren, + stopPm2App, +} from "./pm2-bridge.js"; + +const DRAIN_DEFAULT_MS = 60_000; +const DRAIN_POLL_MS = 500; + +const jwtSecret = + process.env.JFLOW_JWT_SECRET ?? + (process.env.NODE_ENV === "production" ? "" : "jflow-dev-secret"); + +if (!jwtSecret) { + log.error("JFLOW_JWT_SECRET is required in production"); + process.exit(1); +} + +await migrate(); +enableLogPersistence(); + +log.info({ redis: getRedisUrlForLog() }, "starting jerapah-flow control"); + +const workflowQueue = createWorkflowQueue(); +try { + await workflowQueue.waitUntilReady(); +} catch (err) { + log.error({ err, redis: getRedisUrlForLog() }, "failed to connect to Redis"); + process.exit(1); +} + +try { + await connectPm2(); +} catch (err) { + log.error({ err }, "failed to connect to PM2 — is pm2 installed?"); + process.exit(1); +} + +async function applyDesiredState() { + const state = readControlState(); + await ensureHttp({ + generation: state.generation, + running: state.http === "running", + }); + await ensureWorkers({ + generation: state.generation, + count: state.workers, + }); + if (state.queuePaused) { + await workflowQueue.pause(); + } else { + const paused = await workflowQueue.isPaused(); + if (paused) await workflowQueue.resume(); + } + log.info( + { + http: state.http, + workers: state.workers, + generation: state.generation, + queuePaused: state.queuePaused, + }, + "applied desired state", + ); +} + +await applyDesiredState(); + +const server = fastify({ loggerInstance: log }); +await server.register(cookie); +await server.register(jwt, { + secret: jwtSecret, + cookie: { cookieName: COOKIE, signed: false }, +}); +await server.register(cors, { + origin: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:8500", + credentials: true, +}); + +server.decorate("authenticate", async function authenticate(req, reply) { + try { + await req.jwtVerify(); + } catch { + return reply.code(401).send({ error: "unauthorized" }); + } +}); + +server.decorate("requireAdmin", async function requireAdmin(req, reply) { + if (req.user?.role !== "admin") { + return reply.code(403).send({ error: "forbidden" }); + } +}); + +/** + * @param {number} timeoutMs + * @param {string} lockToken + */ +async function waitUntilIdle(timeoutMs, lockToken) { + const started = Date.now(); + while (Date.now() - started < timeoutMs) { + refreshOpsLock(lockToken); + await reconcileOrphanRuns(workflowQueue); + const active = await workflowQueue.getActiveCount(); + if (active === 0) return { ok: true, active: 0 }; + await new Promise((r) => setTimeout(r, DRAIN_POLL_MS)); + } + const active = await workflowQueue.getActiveCount(); + return { ok: active === 0, active }; +} + +async function buildStatus() { + const state = readControlState(); + const children = await describeChildren(); + const heartbeats = await readHeartbeats(); + const live = heartbeats.filter((h) => !h.stale); + const queuePaused = await workflowQueue.isPaused(); + const counts = await workflowQueue.getJobCounts( + "active", + "waiting", + "delayed", + "paused", + "failed", + "completed", + ); + const orphans = await reconcileOrphanRuns(workflowQueue); + + const expectedGen = state.generation; + const mismatched = live.filter((h) => h.generation !== expectedGen); + + return { + control: { + pid: process.pid, + generation: getConfigGeneration(), + }, + desired: state, + children, + heartbeats: live, + generationMismatch: mismatched.length > 0 || state.restartNeeded, + mismatched, + queue: { + paused: queuePaused, + counts, + }, + orphans, + }; +} + +server.get( + "/ops/status", + { onRequest: [server.authenticate] }, + async () => buildStatus(), +); + +server.get("/ops/health", async () => ({ ok: true, role: "control" })); + +server.post( + "/ops/pause", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (_req, reply) => { + await workflowQueue.pause(); + patchControlState({ queuePaused: true }); + return reply.send({ ok: true, queuePaused: true }); + }, +); + +server.post( + "/ops/resume", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (_req, reply) => { + await workflowQueue.resume(); + patchControlState({ queuePaused: false }); + return reply.send({ ok: true, queuePaused: false }); + }, +); + +server.post( + "/ops/reload", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (_req, reply) => { + await publishReload({ type: "workflows" }); + return reply.send({ ok: true, published: true }); + }, +); + +server.post( + "/ops/generation/bump", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ reason?: string }} */ (req.body ?? {}); + const reason = String(body.reason ?? "manual bump").slice(0, 200); + const state = bumpGeneration(reason); + return reply.send({ ok: true, desired: state }); + }, +); + +server.post( + "/ops/http/start", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (_req, reply) => { + const state = patchControlState({ http: "running" }); + await ensureHttp({ generation: state.generation, running: true }); + return reply.send({ ok: true, desired: state }); + }, +); + +server.post( + "/ops/http/stop", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (_req, reply) => { + const state = patchControlState({ http: "stopped" }); + await stopPm2App(PM2_HTTP_NAME); + return reply.send({ ok: true, desired: state }); + }, +); + +server.post( + "/ops/restart", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ force?: boolean, timeoutMs?: number }} */ ( + req.body ?? {} + ); + const force = Boolean(body.force); + const timeoutMs = Math.min( + Math.max(Number(body.timeoutMs) || DRAIN_DEFAULT_MS, 1_000), + 10 * 60_000, + ); + + const lock = tryAcquireOpsLock(`restart:${process.pid}`); + if (!lock.ok) { + return reply.code(409).send({ error: lock.error, holder: lock.holder }); + } + + const stateBefore = readControlState(); + const wantPaused = stateBefore.queuePaused; + try { + await workflowQueue.pause(); + + if (!force) { + const drained = await waitUntilIdle(timeoutMs, lock.token); + if (!drained.ok) { + if (!wantPaused) { + await workflowQueue.resume(); + } + return reply.code(409).send({ + error: "active jobs still running", + active: drained.active, + hint: "wait or retry with force=true", + }); + } + } + + const desired = readControlState(); + await stopPm2App(PM2_HTTP_NAME); + await deletePm2App(PM2_HTTP_NAME); + await stopPm2App(PM2_WORKER_NAME); + await deletePm2App(PM2_WORKER_NAME); + + await reconcileOrphanRuns(workflowQueue); + // Drop stale heartbeats from killed PIDs + try { + const { HEARTBEAT_KEY } = await import("./control-bus.js"); + const { getSharedConnection } = await import("./workflow-queue.js"); + await getSharedConnection().del(HEARTBEAT_KEY); + } catch { + // ignore + } + await migrate(); + + await recreateChildren({ + generation: desired.generation, + http: desired.http === "running", + workers: desired.workers, + }); + + if (wantPaused) { + await workflowQueue.pause(); + patchControlState({ queuePaused: true }); + } else { + await workflowQueue.resume(); + patchControlState({ queuePaused: false }); + } + + await new Promise((r) => setTimeout(r, 1500)); + clearRestartNeeded(); + + return reply.send({ + ok: true, + forced: force, + desired: readControlState(), + status: await buildStatus(), + }); + } catch (err) { + log.error({ err }, "restart failed"); + return reply.code(500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } finally { + releaseOpsLock(lock.token); + } + }, +); + +server.post( + "/ops/scale", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ workers?: number, force?: boolean, timeoutMs?: number }} */ ( + req.body ?? {} + ); + const workers = Math.floor(Number(body.workers)); + if (!Number.isFinite(workers) || workers < 0 || workers > 32) { + return reply.code(400).send({ error: "workers must be 0..32" }); + } + const force = Boolean(body.force); + const timeoutMs = Math.min( + Math.max(Number(body.timeoutMs) || DRAIN_DEFAULT_MS, 1_000), + 10 * 60_000, + ); + + const current = readControlState(); + const scalingDown = workers < current.workers; + + const lock = tryAcquireOpsLock(`scale:${process.pid}`); + if (!lock.ok) { + return reply.code(409).send({ error: lock.error, holder: lock.holder }); + } + + const wasPaused = await workflowQueue.isPaused(); + try { + if (scalingDown) { + await workflowQueue.pause(); + if (!force) { + const drained = await waitUntilIdle(timeoutMs, lock.token); + if (!drained.ok) { + if (!wasPaused) { + await workflowQueue.resume(); + patchControlState({ queuePaused: false }); + } + return reply.code(409).send({ + error: "active jobs still running", + active: drained.active, + hint: "wait or retry with force=true", + }); + } + } + } + + const state = patchControlState({ workers }); + await ensureWorkers({ generation: state.generation, count: workers }); + + if (scalingDown && !wasPaused && !state.queuePaused) { + await workflowQueue.resume(); + patchControlState({ queuePaused: false }); + } + + return reply.send({ ok: true, desired: readControlState() }); + } catch (err) { + log.error({ err }, "scale failed"); + return reply.code(500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } finally { + releaseOpsLock(lock.token); + } + }, +); + +async function shutdown() { + try { + await workflowQueue.close(); + await closeRedis(); + await flushLogs(); + await db.destroy(); + disconnectPm2(); + } catch (err) { + log.error({ err }, "control shutdown error"); + } + process.exit(0); +} + +process.on("SIGINT", shutdown); +process.on("SIGTERM", shutdown); + +const port = Number(process.env.JFLOW_CONTROL_PORT ?? process.env.PORT ?? 8600); +server + .listen({ host: "0.0.0.0", port }) + .then(() => { + log.info(`Control is running on port ${port}`); + }) + .catch((err) => { + log.error({ err }, "failed to start control"); + process.exit(1); + }); diff --git a/packages/server/dev-pm2.mjs b/packages/server/dev-pm2.mjs new file mode 100644 index 0000000..92e71d9 --- /dev/null +++ b/packages/server/dev-pm2.mjs @@ -0,0 +1,114 @@ +#!/usr/bin/env node +/** + * Start Vite (:8500) + control (:8600). Control connects to PM2 and starts HTTP/workers. + * + * Usage: pnpm dev:pm2 + */ +import { spawn } from "node:child_process"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import net from "node:net"; + +const root = path.resolve(path.dirname(fileURLToPath(import.meta.url)), ".."); +const children = []; + +function run(command, args, opts = {}) { + const child = spawn(command, args, { + cwd: root, + stdio: "inherit", + env: { + ...process.env, + JFLOW_CORS_ORIGIN: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:8500", + JFLOW_CONTROL_PORT: process.env.JFLOW_CONTROL_PORT ?? "8600", + PORT: process.env.JFLOW_HTTP_PORT ?? "8700", + ...opts.env, + }, + shell: false, + }); + children.push(child); + child.on("exit", (code, signal) => { + if (shuttingDown) return; + console.error(`[dev:pm2] ${command} ${args.join(" ")} exited (${code ?? signal})`); + shutdown(code ?? 1); + }); + return child; +} + +let shuttingDown = false; + +async function redisReachable() { + const url = process.env.REDIS_URL || "redis://127.0.0.1:6379"; + let host = "127.0.0.1"; + let port = 6379; + try { + const u = new URL(url); + host = u.hostname || host; + port = Number(u.port || 6379); + } catch { + // keep defaults + } + return new Promise((resolve) => { + const socket = net.connect({ host, port }); + socket.setTimeout(1500); + socket.on("connect", () => { + socket.destroy(); + resolve(true); + }); + socket.on("timeout", () => { + socket.destroy(); + resolve(false); + }); + socket.on("error", () => resolve(false)); + }); +} + +async function shutdown(code = 0) { + if (shuttingDown) return; + shuttingDown = true; + for (const child of children) { + if (!child.killed) child.kill("SIGTERM"); + } + // Best-effort: stop jflow apps so the next run is clean + try { + const stop = spawn( + process.platform === "win32" ? "npx.cmd" : "npx", + ["pm2", "delete", "jflow-http", "jflow-worker"], + { cwd: root, stdio: "ignore", shell: false }, + ); + await new Promise((r) => stop.on("exit", r)); + } catch { + // ignore + } + process.exit(code); +} + +process.on("SIGINT", () => void shutdown(0)); +process.on("SIGTERM", () => void shutdown(0)); + +if (!(await redisReachable())) { + console.error( + "[dev:pm2] Redis is not reachable. Start Redis (default redis://127.0.0.1:6379) and retry.", + ); + process.exit(1); +} + +console.log("[dev:pm2] starting control on :8600 (PM2 will start HTTP :8700 + workers)"); +run(process.execPath, [path.join(root, "server/control.js")], { + env: { + JFLOW_CONTROL_PORT: "8600", + }, +}); + +// Give control a moment to migrate + spawn before Vite opens +await new Promise((r) => setTimeout(r, 1500)); + +console.log("[dev:pm2] starting Vite on :8500"); +run( + process.platform === "win32" ? "pnpm.cmd" : "pnpm", + ["--filter", "@jerapah-flow/web", "dev"], + { + env: { + // vite.config reads nothing; port is set in vite.config.js + }, + }, +); diff --git a/packages/server/ecosystem.dev.cjs b/packages/server/ecosystem.dev.cjs new file mode 100644 index 0000000..c832870 --- /dev/null +++ b/packages/server/ecosystem.dev.cjs @@ -0,0 +1,38 @@ +/** + * Dev ecosystem for `pnpm dev:pm2`. + * Prefer starting children via control.js (desired state) rather than this file. + * Kept as a reference / fallback: `pm2 start packages/server/ecosystem.dev.cjs` + */ +const path = require("path"); + +const root = path.resolve(__dirname, "../.."); + +module.exports = { + apps: [ + { + name: "jflow-http", + script: path.join(__dirname, "server.js"), + cwd: root, + instances: 1, + exec_mode: "fork", + env: { + JFLOW_ROLE: "api", + PORT: "8700", + JFLOW_CORS_ORIGIN: "http://localhost:8500", + JFLOW_CONFIG_GENERATION: "1", + }, + }, + { + name: "jflow-worker", + script: path.join(__dirname, "worker.js"), + cwd: root, + instances: 1, + exec_mode: "fork", + env: { + JFLOW_ROLE: "worker", + JFLOW_CORS_ORIGIN: "http://localhost:8500", + JFLOW_CONFIG_GENERATION: "1", + }, + }, + ], +}; diff --git a/packages/server/orphan-runs.js b/packages/server/orphan-runs.js new file mode 100644 index 0000000..912f5db --- /dev/null +++ b/packages/server/orphan-runs.js @@ -0,0 +1,58 @@ +import * as store from "./store.js"; +import { log } from "./logger.js"; + +/** + * Mark SQLite runs stuck in `running` when BullMQ no longer has them active. + * + * @param {import("bullmq").Queue} queue + * @returns {Promise<{ repaired: number, ids: string[] }>} + */ +export async function reconcileOrphanRuns(queue) { + const runs = await store.listRuns({ status: "running", limit: 200 }); + if (!runs.length) return { repaired: 0, ids: [] }; + + /** @type {Set} */ + const activeIds = new Set(); + try { + const active = await queue.getJobs(["active"]); + for (const job of active) { + if (job?.id != null) activeIds.add(String(job.id)); + } + } catch (err) { + log.error({ err }, "failed to list active jobs for orphan reconcile"); + return { repaired: 0, ids: [] }; + } + + /** @type {string[]} */ + const ids = []; + for (const run of runs) { + const jobId = run.job_id ? String(run.job_id) : run.id; + if (activeIds.has(jobId)) continue; + + let stillAlive = false; + try { + const job = await queue.getJob(jobId); + if (job) { + const state = await job.getState(); + // waiting/delayed means not started as running worker — still orphan for SQLite running + stillAlive = state === "active"; + } + } catch { + stillAlive = false; + } + if (stillAlive) continue; + + await store.finishRun( + run.id, + "failed", + null, + new Error("worker_lost: run interrupted by process stop or crash"), + ); + ids.push(run.id); + } + + if (ids.length) { + log.warn({ count: ids.length, ids }, "reconciled orphan running runs"); + } + return { repaired: ids.length, ids }; +} diff --git a/packages/server/package.json b/packages/server/package.json index 95d29cd..39d66a7 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -1,13 +1,18 @@ { "name": "@jerapah-flow/server", - "version": "1.0.0", + "version": "0.1.0", "private": true, "type": "module", "main": "runner.js", "scripts": { "dev": "node --watch runner.js", + "dev:pm2": "node dev-pm2.mjs", "start": "node runner.js", - "migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\"" + "start:api": "node server.js", + "start:worker": "node worker.js", + "start:control": "node control.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" }, "dependencies": { "@aws-sdk/client-s3": "^3.1111.0", @@ -31,6 +36,7 @@ "nodemailer": "^9.0.5", "pino": "^10.3.1", "pino-roll": "^4.0.0", + "pm2": "^6.0.13", "ssh2-sftp-client": "^12.1.1", "webdav": "^5.10.0", "yaml": "^2.9.0" diff --git a/packages/server/paths.js b/packages/server/paths.js index f84f6a5..f591164 100644 --- a/packages/server/paths.js +++ b/packages/server/paths.js @@ -5,5 +5,14 @@ export const SERVER_ROOT = path.dirname(fileURLToPath(import.meta.url)); export const SCRIPTS_DIR = path.join(SERVER_ROOT, "scripts"); export const WORKFLOWS_DIR = path.join(SERVER_ROOT, "workflows"); export const DATA_DIR = path.join(SERVER_ROOT, "data"); +/** User plugins (repo-root /plugins, outside the pnpm workspace). */ +export const PLUGINS_DIR = + process.env.JFLOW_PLUGINS_DIR ?? + path.resolve(SERVER_ROOT, "../../plugins"); +/** Example plugin sources shipped with the repo. */ +export const EXAMPLE_PLUGINS_DIR = path.resolve( + SERVER_ROOT, + "../../examples/plugins", +); export const LOGS_DIR = path.join(SERVER_ROOT, "logs"); export const WEB_DIST = path.resolve(SERVER_ROOT, "../web/dist"); diff --git a/packages/server/plugin-install.js b/packages/server/plugin-install.js new file mode 100644 index 0000000..06c1f05 --- /dev/null +++ b/packages/server/plugin-install.js @@ -0,0 +1,218 @@ +import fs from "fs"; +import os from "os"; +import path from "path"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; +import { EXAMPLE_PLUGINS_DIR, PLUGINS_DIR } from "./paths.js"; +import { + installPluginFromDirectory, + pluginDir, +} from "./plugin-store.js"; +import { readManifestFile } from "./plugin-manifest.js"; + +const execFileAsync = promisify(execFile); + +/** + * @param {string} url + */ +export function assertHttpsGitUrl(url) { + let parsed; + try { + parsed = new URL(url); + } catch { + const err = new Error("invalid git URL"); + err.statusCode = 400; + throw err; + } + if (parsed.protocol !== "https:") { + const err = new Error("git URL must use https://"); + err.statusCode = 400; + throw err; + } + return parsed.toString(); +} + +/** + * Run pnpm install in a plugin directory (ignore lifecycle scripts). + * @param {string} dir + */ +export async function pnpmInstallPlugin(dir) { + const pkg = path.join(dir, "package.json"); + if (!fs.existsSync(pkg)) return { skipped: true }; + let hasDeps = false; + try { + const raw = JSON.parse(fs.readFileSync(pkg, "utf8")); + hasDeps = Boolean( + (raw.dependencies && Object.keys(raw.dependencies).length) || + (raw.optionalDependencies && + Object.keys(raw.optionalDependencies).length), + ); + } catch { + hasDeps = true; + } + if (!hasDeps) return { skipped: true }; + + await execFileAsync( + "pnpm", + ["install", "--dir", dir, "--ignore-scripts", "--prefer-offline"], + { + cwd: dir, + env: { ...process.env, CI: "1" }, + timeout: 5 * 60_000, + maxBuffer: 10 * 1024 * 1024, + }, + ); + return { skipped: false }; +} + +/** + * @param {string} url + * @param {{ ref?: string, overwrite?: boolean }} [opts] + */ +export async function installPluginFromGit(url, opts = {}) { + const httpsUrl = assertHttpsGitUrl(url); + const staging = path.join( + PLUGINS_DIR, + `.staging-git-${Date.now()}-${Math.random().toString(36).slice(2)}`, + ); + fs.mkdirSync(PLUGINS_DIR, { recursive: true }); + fs.mkdirSync(staging, { recursive: true }); + try { + const args = ["clone", "--depth", "1"]; + if (opts.ref) { + args.push("--branch", String(opts.ref)); + } + args.push(httpsUrl, staging); + await execFileAsync("git", args, { + timeout: 5 * 60_000, + maxBuffer: 5 * 1024 * 1024, + }); + // Remove .git to keep plugins lean + fs.rmSync(path.join(staging, ".git"), { recursive: true, force: true }); + const installed = installPluginFromDirectory(staging, { + overwrite: Boolean(opts.overwrite), + markRestart: false, + }); + await pnpmInstallPlugin(installed.dir); + const { bumpGeneration } = await import("./control-state.js"); + bumpGeneration(`plugin:${installed.id} installed from git`); + return installed; + } finally { + fs.rmSync(staging, { recursive: true, force: true }); + } +} + +/** + * @param {string} zipPath + * @param {{ overwrite?: boolean }} [opts] + */ +export async function installPluginFromZipFile(zipPath, opts = {}) { + const abs = path.resolve(zipPath); + if (!fs.existsSync(abs)) { + const err = new Error("zip file not found"); + err.statusCode = 400; + throw err; + } + const staging = path.join( + PLUGINS_DIR, + `.staging-zip-${Date.now()}-${Math.random().toString(36).slice(2)}`, + ); + const extractDir = path.join(staging, "extract"); + fs.mkdirSync(extractDir, { recursive: true }); + try { + await execFileAsync("unzip", ["-q", abs, "-d", extractDir], { + timeout: 120_000, + }); + const root = findPluginRoot(extractDir); + const installed = installPluginFromDirectory(root, { + overwrite: Boolean(opts.overwrite), + markRestart: false, + }); + await pnpmInstallPlugin(installed.dir); + const { bumpGeneration } = await import("./control-state.js"); + bumpGeneration(`plugin:${installed.id} installed from zip`); + return installed; + } finally { + fs.rmSync(staging, { recursive: true, force: true }); + } +} + +/** + * @param {Buffer} buffer + * @param {{ overwrite?: boolean }} [opts] + */ +export async function installPluginFromZipBuffer(buffer, opts = {}) { + const tmp = path.join( + os.tmpdir(), + `jflow-plugin-${Date.now()}-${Math.random().toString(36).slice(2)}.zip`, + ); + fs.writeFileSync(tmp, buffer); + try { + return await installPluginFromZipFile(tmp, opts); + } finally { + fs.unlinkSync(tmp); + } +} + +/** + * Find directory containing jerapah-plugin.json (zip may have a single top folder). + * @param {string} extractDir + */ +function findPluginRoot(extractDir) { + const direct = path.join(extractDir, "jerapah-plugin.json"); + if (fs.existsSync(direct)) return extractDir; + const entries = fs.readdirSync(extractDir, { withFileTypes: true }); + const dirs = entries.filter((e) => e.isDirectory() && !e.name.startsWith(".")); + if (dirs.length === 1) { + const nested = path.join(extractDir, dirs[0].name); + if (fs.existsSync(path.join(nested, "jerapah-plugin.json"))) return nested; + } + for (const e of dirs) { + const nested = path.join(extractDir, e.name); + if (fs.existsSync(path.join(nested, "jerapah-plugin.json"))) return nested; + } + const err = new Error("zip missing jerapah-plugin.json"); + err.statusCode = 400; + throw err; +} + +/** + * Install a shipped example plugin by id (from examples/plugins/). + * @param {string} exampleId + * @param {{ overwrite?: boolean }} [opts] + */ +export async function installExamplePlugin(exampleId, opts = {}) { + const src = path.join(EXAMPLE_PLUGINS_DIR, exampleId); + if (!fs.existsSync(src)) { + const err = new Error(`example plugin not found: ${exampleId}`); + err.statusCode = 404; + throw err; + } + const installed = installPluginFromDirectory(src, { + overwrite: Boolean(opts.overwrite), + markRestart: false, + }); + await pnpmInstallPlugin(installed.dir); + const { bumpGeneration } = await import("./control-state.js"); + bumpGeneration(`plugin:${installed.id} installed from example`); + return installed; +} + +/** + * Copy example into plugins if missing (used by smoke / first-run helpers). + * @param {string} exampleId + */ +export function ensureExampleInstalled(exampleId) { + const dest = pluginDir(exampleId); + if (fs.existsSync(dest)) { + return { id: exampleId, already: true, dir: dest }; + } + const src = path.join(EXAMPLE_PLUGINS_DIR, exampleId); + const installed = installPluginFromDirectory(src, { + overwrite: false, + markRestart: false, + }); + return { ...installed, already: false }; +} + +void readManifestFile; diff --git a/packages/server/plugin-manifest.js b/packages/server/plugin-manifest.js new file mode 100644 index 0000000..e03ea69 --- /dev/null +++ b/packages/server/plugin-manifest.js @@ -0,0 +1,159 @@ +import fs from "fs"; +import path from "path"; +import { getAppVersion, satisfiesRange } from "./app-version.js"; + +export const PLUGIN_MANIFEST = "jerapah-plugin.json"; +export const PLUGIN_PREFIX = "plugin/"; + +/** + * @param {string} id + * @returns {string} + */ +export function assertPluginId(id) { + if (typeof id !== "string" || !/^[a-z0-9]+(?:-[a-z0-9]+)*$/.test(id)) { + const err = new Error( + "invalid plugin id (use lowercase letters, numbers, hyphens)", + ); + err.statusCode = 400; + throw err; + } + if (id.length > 64) { + const err = new Error("plugin id too long"); + err.statusCode = 400; + throw err; + } + return id; +} + +/** + * @param {string} scriptRef e.g. plugin/foo or plugin/foo.js + * @returns {{ id: string, scriptRef: string } | null} + */ +export function parsePluginScriptRef(scriptRef) { + if (typeof scriptRef !== "string") return null; + let rest = scriptRef; + if (rest.startsWith(PLUGIN_PREFIX)) { + rest = rest.slice(PLUGIN_PREFIX.length); + } else { + return null; + } + if (rest.endsWith(".js")) rest = rest.slice(0, -3); + if (!rest || rest.includes("/") || rest.includes("\\")) return null; + try { + const id = assertPluginId(rest); + return { id, scriptRef: `${PLUGIN_PREFIX}${id}` }; + } catch { + return null; + } +} + +/** + * @param {string} id + * @returns {string} + */ +export function pluginScriptRef(id) { + return `${PLUGIN_PREFIX}${assertPluginId(id)}`; +} + +/** + * @param {unknown} raw + * @returns {{ + * id: string, + * name: string, + * version: string, + * jerapah: string, + * main: string, + * description: string | null, + * }} + */ +export function validateManifest(raw) { + if (raw == null || typeof raw !== "object" || Array.isArray(raw)) { + const err = new Error("manifest must be an object"); + err.statusCode = 400; + throw err; + } + const id = assertPluginId(String(/** @type {any} */ (raw).id ?? "")); + const version = String(/** @type {any} */ (raw).version ?? "").trim(); + if (!/^\d+\.\d+\.\d+/.test(version)) { + const err = new Error("manifest.version must be semver (e.g. 0.1.0)"); + err.statusCode = 400; + throw err; + } + const jerapah = String(/** @type {any} */ (raw).jerapah ?? "").trim(); + if (!jerapah) { + const err = new Error("manifest.jerapah range is required"); + err.statusCode = 400; + throw err; + } + const main = String(/** @type {any} */ (raw).main ?? "script.js").trim(); + if (!main || main.includes("..") || path.isAbsolute(main)) { + const err = new Error("manifest.main must be a relative file path"); + err.statusCode = 400; + throw err; + } + const name = + String(/** @type {any} */ (raw).name ?? id).trim() || id; + const descriptionRaw = /** @type {any} */ (raw).description; + const description = + typeof descriptionRaw === "string" && descriptionRaw.trim() + ? descriptionRaw.trim() + : null; + return { id, name, version, jerapah, main, description }; +} + +/** + * @param {string} pluginDir + */ +export function readManifestFile(pluginDir) { + const filePath = path.join(pluginDir, PLUGIN_MANIFEST); + if (!fs.existsSync(filePath)) { + const err = new Error(`missing ${PLUGIN_MANIFEST}`); + err.statusCode = 400; + throw err; + } + let raw; + try { + raw = JSON.parse(fs.readFileSync(filePath, "utf8")); + } catch { + const err = new Error(`invalid ${PLUGIN_MANIFEST} JSON`); + err.statusCode = 400; + throw err; + } + return validateManifest(raw); +} + +/** + * @param {ReturnType} manifest + * @returns {{ ok: true } | { ok: false, error: string }} + */ +export function checkJerapahCompat(manifest) { + const appVersion = getAppVersion(); + if (!satisfiesRange(appVersion, manifest.jerapah)) { + return { + ok: false, + error: `plugin requires JerapahFlow ${manifest.jerapah}; app is ${appVersion}`, + }; + } + return { ok: true }; +} + +/** + * @param {{ + * id: string, + * name?: string, + * version?: string, + * jerapah?: string, + * main?: string, + * description?: string | null, + * }} opts + */ +export function buildManifest(opts) { + return validateManifest({ + id: opts.id, + name: opts.name ?? opts.id, + version: opts.version ?? "0.1.0", + jerapah: opts.jerapah ?? ">=0.1.0 <1.0.0", + main: opts.main ?? "script.js", + description: opts.description ?? null, + }); +} diff --git a/packages/server/plugin-store.js b/packages/server/plugin-store.js new file mode 100644 index 0000000..e10201c --- /dev/null +++ b/packages/server/plugin-store.js @@ -0,0 +1,456 @@ +import fs from "fs"; +import path from "path"; +import { createRequire } from "node:module"; +import { PLUGINS_DIR, SCRIPTS_DIR } from "./paths.js"; +import { + PLUGIN_MANIFEST, + assertPluginId, + buildManifest, + checkJerapahCompat, + parsePluginScriptRef, + pluginScriptRef, + readManifestFile, + validateManifest, +} from "./plugin-manifest.js"; +import { listScriptFiles } from "./fs-store.js"; +import { bumpGeneration } from "./control-state.js"; + +/** + * @param {string} id + */ +export function pluginDir(id) { + return path.join(PLUGINS_DIR, assertPluginId(id)); +} + +export function ensurePluginsDir() { + fs.mkdirSync(PLUGINS_DIR, { recursive: true }); +} + +/** + * Core script file names (*.js) currently shipped under SCRIPTS_DIR. + * @returns {string[]} + */ +export function listCoreScriptNames() { + return listScriptFiles(); +} + +/** + * Bare core names without .js (for collision checks). + * @returns {Set} + */ +export function coreBareNames() { + return new Set( + listCoreScriptNames().map((n) => (n.endsWith(".js") ? n.slice(0, -3) : n)), + ); +} + +/** + * @returns {Array<{ + * id: string, + * scriptRef: string, + * dir: string, + * manifest: ReturnType, + * compatible: boolean, + * compatError: string | null, + * disabled: boolean, + * }>} + */ +export function listInstalledPlugins() { + ensurePluginsDir(); + if (!fs.existsSync(PLUGINS_DIR)) return []; + /** @type {Array} */ + const out = []; + for (const entry of fs.readdirSync(PLUGINS_DIR, { withFileTypes: true })) { + if (!entry.isDirectory()) continue; + let id; + try { + id = assertPluginId(entry.name); + } catch { + continue; + } + const dir = pluginDir(id); + try { + const manifest = readManifestFile(dir); + if (manifest.id !== id) { + out.push({ + id, + scriptRef: pluginScriptRef(id), + dir, + manifest, + compatible: false, + compatError: `manifest id "${manifest.id}" does not match folder "${id}"`, + disabled: true, + }); + continue; + } + const compat = checkJerapahCompat(manifest); + const disabledFlag = fs.existsSync(path.join(dir, ".disabled")); + out.push({ + id, + scriptRef: pluginScriptRef(id), + dir, + manifest, + compatible: compat.ok, + compatError: compat.ok ? null : compat.error, + disabled: disabledFlag || !compat.ok, + }); + } catch (err) { + out.push({ + id, + scriptRef: pluginScriptRef(id), + dir, + manifest: null, + compatible: false, + compatError: err instanceof Error ? err.message : String(err), + disabled: true, + }); + } + } + return out.sort((a, b) => a.id.localeCompare(b.id)); +} + +/** + * @param {string} id + */ +export function getInstalledPlugin(id) { + const needle = assertPluginId(id); + return listInstalledPlugins().find((p) => p.id === needle) ?? null; +} + +/** + * Resolve a workflow script ref to a filesystem path + kind. + * + * @param {string} scriptRef + * @returns {{ + * kind: "core" | "plugin", + * scriptRef: string, + * filePath: string, + * pluginId?: string, + * pluginDir?: string, + * disabled?: boolean, + * error?: string, + * }} + */ +export function resolveScriptRef(scriptRef) { + if (typeof scriptRef !== "string" || !scriptRef.trim()) { + return { + kind: "core", + scriptRef: String(scriptRef), + filePath: "", + error: "invalid script ref", + }; + } + + const plugin = parsePluginScriptRef(scriptRef); + if (plugin) { + const installed = getInstalledPlugin(plugin.id); + if (!installed) { + return { + kind: "plugin", + scriptRef: plugin.scriptRef, + pluginId: plugin.id, + filePath: "", + error: `plugin not installed: ${plugin.scriptRef}`, + }; + } + if (installed.disabled) { + return { + kind: "plugin", + scriptRef: plugin.scriptRef, + pluginId: plugin.id, + pluginDir: installed.dir, + filePath: "", + disabled: true, + error: + installed.compatError || + `plugin disabled: ${plugin.scriptRef}`, + }; + } + const mainPath = path.join(installed.dir, installed.manifest.main); + if (!fs.existsSync(mainPath)) { + return { + kind: "plugin", + scriptRef: plugin.scriptRef, + pluginId: plugin.id, + pluginDir: installed.dir, + filePath: "", + error: `plugin main missing: ${installed.manifest.main}`, + }; + } + return { + kind: "plugin", + scriptRef: plugin.scriptRef, + pluginId: plugin.id, + pluginDir: installed.dir, + filePath: mainPath, + }; + } + + // Core: must be a plain *.js filename + if ( + scriptRef.includes("/") || + scriptRef.includes("\\") || + scriptRef.includes("..") + ) { + return { + kind: "core", + scriptRef, + filePath: "", + error: "invalid core script name", + }; + } + const name = scriptRef.endsWith(".js") ? scriptRef : `${scriptRef}.js`; + const filePath = path.join(SCRIPTS_DIR, name); + if (!fs.existsSync(filePath)) { + return { + kind: "core", + scriptRef: name, + filePath: "", + error: `core script not found: ${name}`, + }; + } + return { kind: "core", scriptRef: name, filePath }; +} + +/** + * Copy a prepared plugin directory into PLUGINS_DIR. + * + * @param {string} sourceDir directory containing jerapah-plugin.json + * @param {{ overwrite?: boolean, markRestart?: boolean, reason?: string }} [opts] + */ +export function installPluginFromDirectory(sourceDir, opts = {}) { + const abs = path.resolve(sourceDir); + if (!fs.existsSync(abs) || !fs.statSync(abs).isDirectory()) { + const err = new Error("plugin source directory not found"); + err.statusCode = 400; + throw err; + } + const manifest = readManifestFile(abs); + const compat = checkJerapahCompat(manifest); + if (!compat.ok) { + const err = new Error(compat.error); + err.statusCode = 409; + throw err; + } + if (coreBareNames().has(manifest.id)) { + const err = new Error( + `plugin id "${manifest.id}" collides with a core script name`, + ); + err.statusCode = 409; + throw err; + } + + const mainPath = path.join(abs, manifest.main); + if (!fs.existsSync(mainPath)) { + const err = new Error(`manifest.main not found: ${manifest.main}`); + err.statusCode = 400; + throw err; + } + + ensurePluginsDir(); + const dest = pluginDir(manifest.id); + if (fs.existsSync(dest)) { + if (!opts.overwrite) { + const err = new Error(`plugin already installed: ${manifest.id}`); + err.statusCode = 409; + throw err; + } + fs.rmSync(dest, { recursive: true, force: true }); + } + + fs.cpSync(abs, dest, { recursive: true }); + + // Ensure package.json exists (fork / thin plugins). + const pkgPath = path.join(dest, "package.json"); + if (!fs.existsSync(pkgPath)) { + fs.writeFileSync( + pkgPath, + `${JSON.stringify( + { + name: `jflow-plugin-${manifest.id}`, + version: manifest.version, + private: true, + type: "module", + }, + null, + 2, + )}\n`, + "utf8", + ); + } + + if (opts.markRestart !== false) { + bumpGeneration(opts.reason ?? `plugin:${manifest.id} installed`); + } + + return { + id: manifest.id, + scriptRef: pluginScriptRef(manifest.id), + dir: dest, + manifest, + }; +} + +/** + * @param {string} id + * @param {{ markRestart?: boolean }} [opts] + */ +export function uninstallPlugin(id, opts = {}) { + const pluginId = assertPluginId(id); + const dir = pluginDir(pluginId); + if (!fs.existsSync(dir)) { + const err = new Error("plugin not found"); + err.statusCode = 404; + throw err; + } + fs.rmSync(dir, { recursive: true, force: true }); + if (opts.markRestart !== false) { + bumpGeneration(`plugin:${pluginId} uninstalled`); + } + return { ok: true, id: pluginId }; +} + +/** + * Fork a core script into a new plugin. + * + * @param {string} coreName e.g. fetch-http.js + * @param {string} newId + * @param {{ description?: string }} [opts] + */ +export function forkCoreScript(coreName, newId, opts = {}) { + const id = assertPluginId(newId); + 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 coreFile = coreName.endsWith(".js") ? coreName : `${coreName}.js`; + const src = path.join(SCRIPTS_DIR, coreFile); + if (!fs.existsSync(src)) { + const err = new Error(`core script not found: ${coreFile}`); + err.statusCode = 404; + throw err; + } + + const staging = path.join(PLUGINS_DIR, `.staging-fork-${id}-${Date.now()}`); + fs.mkdirSync(staging, { recursive: true }); + try { + const main = "script.js"; + fs.copyFileSync(src, path.join(staging, main)); + const manifest = buildManifest({ + id, + name: id, + version: "0.1.0", + jerapah: ">=0.1.0 <1.0.0", + main, + description: + opts.description ?? `Fork of core script ${coreFile}`, + }); + fs.writeFileSync( + path.join(staging, PLUGIN_MANIFEST), + `${JSON.stringify(manifest, null, 2)}\n`, + "utf8", + ); + fs.writeFileSync( + path.join(staging, "package.json"), + `${JSON.stringify( + { + name: `jflow-plugin-${id}`, + version: "0.1.0", + private: true, + type: "module", + }, + null, + 2, + )}\n`, + "utf8", + ); + return installPluginFromDirectory(staging, { + overwrite: false, + reason: `plugin:${id} forked from ${coreFile}`, + }); + } finally { + fs.rmSync(staging, { recursive: true, force: true }); + } +} + +/** + * @param {string} pluginDirectory + * @returns {((id: string) => unknown) | null} + */ +export function createPluginRequire(pluginDirectory) { + const pkg = path.resolve(pluginDirectory, "package.json"); + if (!fs.existsSync(pkg)) return null; + return createRequire(pkg); +} + +/** + * Create an empty plugin from the new-script template. + * @param {string} newId + * @param {string} source + * @param {{ description?: string }} [opts] + */ +export function createBlankPlugin(newId, source, opts = {}) { + const id = assertPluginId(newId); + 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; + } + if (typeof source !== "string") { + const err = new Error("source content is required"); + err.statusCode = 400; + throw err; + } + + const staging = path.join(PLUGINS_DIR, `.staging-new-${id}-${Date.now()}`); + fs.mkdirSync(staging, { recursive: true }); + try { + const main = "script.js"; + fs.writeFileSync(path.join(staging, main), source, "utf8"); + const manifest = buildManifest({ + id, + name: id, + version: "0.1.0", + jerapah: ">=0.1.0 <1.0.0", + main, + description: opts.description ?? null, + }); + fs.writeFileSync( + path.join(staging, PLUGIN_MANIFEST), + `${JSON.stringify(manifest, null, 2)}\n`, + "utf8", + ); + fs.writeFileSync( + path.join(staging, "package.json"), + `${JSON.stringify( + { + name: `jflow-plugin-${id}`, + version: "0.1.0", + private: true, + type: "module", + }, + null, + 2, + )}\n`, + "utf8", + ); + return installPluginFromDirectory(staging, { + overwrite: false, + reason: `plugin:${id} created`, + }); + } finally { + fs.rmSync(staging, { recursive: true, force: true }); + } +} diff --git a/packages/server/pm2-bridge.js b/packages/server/pm2-bridge.js new file mode 100644 index 0000000..5d72a28 --- /dev/null +++ b/packages/server/pm2-bridge.js @@ -0,0 +1,272 @@ +import path from "path"; +import pm2 from "pm2"; +import { SERVER_ROOT } from "./paths.js"; + +const REPO_ROOT = path.resolve(SERVER_ROOT, "../.."); + +export const PM2_HTTP_NAME = "jflow-http"; +export const PM2_WORKER_NAME = "jflow-worker"; + +/** + * @returns {Promise} + */ +export function connectPm2() { + return new Promise((resolve, reject) => { + pm2.connect((err) => { + if (err) reject(err); + else resolve(); + }); + }); +} + +export function disconnectPm2() { + try { + pm2.disconnect(); + } catch { + // ignore + } +} + +/** + * @returns {Promise} + */ +export function listPm2() { + return new Promise((resolve, reject) => { + pm2.list((err, list) => { + if (err) reject(err); + else resolve(list ?? []); + }); + }); +} + +/** + * @param {string} name + * @returns {Promise} + */ +export async function listPm2ByName(name) { + const list = await listPm2(); + return list.filter((p) => p.name === name); +} + +/** + * @param {object} app + * @returns {Promise} + */ +function startPm2App(app) { + return new Promise((resolve, reject) => { + pm2.start(app, (err) => { + if (err) reject(err); + else resolve(); + }); + }); +} + +/** + * @param {string} name + * @returns {Promise} + */ +export function stopPm2App(name) { + return new Promise((resolve, reject) => { + pm2.stop(name, (err) => { + if (err) { + const msg = err instanceof Error ? err.message : String(err); + if (/not found|doesn't exist|process or namespace/i.test(msg)) { + resolve(); + return; + } + reject(err); + return; + } + resolve(); + }); + }); +} + +/** + * @param {string} name + * @returns {Promise} + */ +export function deletePm2App(name) { + return new Promise((resolve, reject) => { + pm2.delete(name, (err) => { + if (err) { + const msg = err instanceof Error ? err.message : String(err); + if (/not found|doesn't exist|process or namespace/i.test(msg)) { + resolve(); + return; + } + reject(err); + return; + } + resolve(); + }); + }); +} + +/** + * @param {string} name + * @param {number} instances + * @returns {Promise} + */ +export function scalePm2App(name, instances) { + return new Promise((resolve, reject) => { + pm2.scale(name, instances, (err) => { + if (err) reject(err); + else resolve(); + }); + }); +} + +/** + * @param {string} name + * @returns {Promise} + */ +export function restartPm2App(name) { + return new Promise((resolve, reject) => { + pm2.restart(name, (err) => { + if (err) reject(err); + else resolve(); + }); + }); +} + +/** + * Shared env for child processes. + * @param {{ generation: number }} opts + */ +export function childEnv(opts) { + return { + ...process.env, + JFLOW_CONFIG_GENERATION: String(opts.generation), + JFLOW_CORS_ORIGIN: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:8500", + PORT: process.env.JFLOW_HTTP_PORT ?? "8700", + }; +} + +/** + * Ensure HTTP app exists and matches desired running/stopped state. + * @param {{ generation: number, running: boolean }} opts + */ +export async function ensureHttp(opts) { + const existing = await listPm2ByName(PM2_HTTP_NAME); + if (!opts.running) { + if (existing.length) await stopPm2App(PM2_HTTP_NAME); + return; + } + + if (existing.length === 0) { + await startPm2App({ + name: PM2_HTTP_NAME, + script: path.join(SERVER_ROOT, "server.js"), + cwd: REPO_ROOT, + instances: 1, + exec_mode: "fork", + autorestart: true, + max_restarts: 20, + env: { + ...childEnv(opts), + JFLOW_ROLE: "api", + }, + }); + return; + } + + const online = existing.some((p) => p.pm2_env?.status === "online"); + if (!online) { + await restartPm2App(PM2_HTTP_NAME); + } +} + +/** + * Ensure worker app has `count` online forks (0 = stopped/deleted). + * @param {{ generation: number, count: number }} opts + */ +export async function ensureWorkers(opts) { + const count = Math.max(0, Math.min(32, Math.floor(opts.count))); + const existing = await listPm2ByName(PM2_WORKER_NAME); + + if (count === 0) { + if (existing.length) { + await stopPm2App(PM2_WORKER_NAME); + await deletePm2App(PM2_WORKER_NAME); + } + return; + } + + if (existing.length === 0) { + await startPm2App({ + name: PM2_WORKER_NAME, + script: path.join(SERVER_ROOT, "worker.js"), + cwd: REPO_ROOT, + instances: count, + exec_mode: "fork", + autorestart: true, + max_restarts: 50, + env: { + ...childEnv(opts), + JFLOW_ROLE: "worker", + }, + }); + return; + } + + const current = existing.length; + if (current !== count) { + await scalePm2App(PM2_WORKER_NAME, count); + } + + const stopped = existing.filter((p) => p.pm2_env?.status !== "online"); + if (stopped.length) { + await restartPm2App(PM2_WORKER_NAME); + } +} + +/** + * Hard recycle HTTP + workers with updated generation env. + * Children must be stopped first for a clean migrate window. + * + * @param {{ generation: number, http: boolean, workers: number }} opts + */ +export async function recreateChildren(opts) { + await stopPm2App(PM2_HTTP_NAME); + await deletePm2App(PM2_HTTP_NAME); + await stopPm2App(PM2_WORKER_NAME); + await deletePm2App(PM2_WORKER_NAME); + + if (opts.http) { + await ensureHttp({ generation: opts.generation, running: true }); + } + if (opts.workers > 0) { + await ensureWorkers({ generation: opts.generation, count: opts.workers }); + } +} + +/** + * Summarize PM2 process status for the ops UI. + */ +export async function describeChildren() { + const list = await listPm2(); + const http = list.filter((p) => p.name === PM2_HTTP_NAME); + const workers = list.filter((p) => p.name === PM2_WORKER_NAME); + + const mapOne = (p) => ({ + name: p.name, + pmId: p.pm_id, + status: p.pm2_env?.status ?? "unknown", + pid: p.pid ?? null, + restarts: p.pm2_env?.restart_time ?? 0, + uptime: p.pm2_env?.pm_uptime ?? null, + generation: Number(p.pm2_env?.JFLOW_CONFIG_GENERATION ?? 0) || null, + }); + + return { + http: http.map(mapOne), + workers: workers.map(mapOne), + httpOnline: http.some((p) => p.pm2_env?.status === "online"), + workerOnlineCount: workers.filter((p) => p.pm2_env?.status === "online").length, + }; +} + +export function getRepoRoot() { + return REPO_ROOT; +} diff --git a/packages/server/runner.js b/packages/server/runner.js index 8394fe2..d040858 100644 --- a/packages/server/runner.js +++ b/packages/server/runner.js @@ -1,245 +1,8 @@ -import fs from "fs"; -import fastify from "fastify"; -import cookie from "@fastify/cookie"; -import cors from "@fastify/cors"; -import jwt from "@fastify/jwt"; -import fastifyStatic from "@fastify/static"; -import { migrate, db } from "./db.js"; -import { log, enableLogPersistence, flushLogs } from "./logger.js"; -import * as store from "./store.js"; -import { createRegistry } from "./registry.js"; -import { COOKIE, OPEN_API_ROUTES } from "./src/api/auth.js"; -import authPlugin from "./src/api/auth.js"; -import usersPlugin from "./src/api/users.js"; -import scriptsPluginFactory from "./src/api/scripts.js"; -import workflowsPluginFactory from "./src/api/workflows.js"; -import runsPlugin from "./src/api/runs.js"; -import dashboardPluginFactory from "./src/api/dashboard.js"; -import secretsPlugin from "./src/api/secrets.js"; -import kvPlugin from "./src/api/kv.js"; -import variablesPlugin from "./src/api/variables.js"; -import httpPagesPlugin from "./src/api/http-pages.js"; -import httpAuthsPlugin from "./src/api/http-auths.js"; -import { WEB_DIST } from "./paths.js"; -import { resolveSecretsKeyMaterial } from "./secrets.js"; -import { - closeRedis, - createWorkflowQueue, - createWorkflowWorker, - getRedisUrlForLog, -} from "./workflow-queue.js"; +import { startApp } from "./start-app.js"; -await migrate(); -enableLogPersistence(); - -const jwtSecret = - process.env.JFLOW_JWT_SECRET ?? - (process.env.NODE_ENV === "production" ? "" : "jflow-dev-secret"); - -if (!jwtSecret) { - log.error("JFLOW_JWT_SECRET is required in production"); - process.exit(1); -} - -try { - resolveSecretsKeyMaterial(); -} catch (err) { - log.error(err instanceof Error ? err.message : String(err)); - process.exit(1); -} - -const role = (process.env.JFLOW_ROLE || "all").toLowerCase(); -const runApi = role === "all" || role === "api"; -const runWorker = role === "all" || role === "worker"; - -log.info({ redis: getRedisUrlForLog(), role }, "starting jerapah-flow"); - -const workflowQueue = createWorkflowQueue(); -try { - await workflowQueue.waitUntilReady(); -} catch (err) { - log.error({ err, redis: getRedisUrlForLog() }, "failed to connect to Redis"); - process.exit(1); -} - -const server = fastify({ loggerInstance: log }); - -await server.register(cookie); -await server.register(jwt, { - secret: jwtSecret, - cookie: { - cookieName: COOKIE, - signed: false, - }, +// Monolith / local `pnpm dev`: API + worker + migrate in one process. +await startApp({ + role: process.env.JFLOW_ROLE || "all", + migrate: true, + serveStaticUi: true, }); -await server.register(cors, { - origin: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:5173", - credentials: true, -}); - -server.decorate("authenticate", async function authenticate(req, reply) { - try { - await req.jwtVerify(); - } catch { - return reply.code(401).send({ error: "unauthorized" }); - } -}); - -server.decorate("requireAdmin", async function requireAdmin(req, reply) { - if (req.user?.role !== "admin") { - return reply.code(403).send({ error: "forbidden" }); - } -}); - -const registry = createRegistry(server, { queue: workflowQueue }); -registry.registerWorkflows(); -if (runApi) { - registry.registerHttpTriggers(); - registry.registerCronTriggers(); - registry.registerPruneJob(); -} - -/** @type {import("bullmq").Worker | null} */ -let workflowWorker = null; -if (runWorker) { - workflowWorker = createWorkflowWorker(async (job) => { - const data = /** @type {{ runId?: string, key?: string, depth?: number }} */ ( - job.data ?? {} - ); - if (typeof data.runId !== "string" || typeof data.key !== "string") { - throw new Error("invalid workflow job payload"); - } - const result = await registry.executeQueuedRun({ - runId: data.runId, - key: data.key, - depth: data.depth ?? 0, - }); - if (result.status === "failed") { - throw new Error(result.error || "workflow failed"); - } - return result; - }); -} - -if (runApi) { - await server.register( - async (api) => { - api.addHook("onRequest", async (req, reply) => { - const raw = (req.url || "").split("?")[0]; - const stripped = raw.replace(/^\/api/, "") || "/"; - const routeUrl = req.routeOptions?.url || stripped; - const open = - OPEN_API_ROUTES.has(`${req.method} ${routeUrl}`) || - OPEN_API_ROUTES.has(`${req.method} ${stripped}`); - if (open) return; - await server.authenticate(req, reply); - }); - await api.register(authPlugin); - await api.register(usersPlugin); - await api.register(secretsPlugin); - await api.register(variablesPlugin); - await api.register(kvPlugin); - await api.register(httpPagesPlugin); - await api.register(httpAuthsPlugin); - await api.register(scriptsPluginFactory(registry)); - await api.register(workflowsPluginFactory(registry)); - await api.register(runsPlugin); - await api.register(dashboardPluginFactory(registry)); - }, - { prefix: "/api" }, - ); - - server.post( - "/admin/workflows/reregister", - { onRequest: [server.authenticate] }, - async (_req, reply) => { - registry.reregister(); - return reply.send({ message: "Workflows refreshed" }); - }, - ); - - server.get( - "/admin/runs", - { onRequest: [server.authenticate] }, - async (req, reply) => { - const q = /** @type {Record} */ (req.query); - const limit = q.limit ? Number(q.limit) : undefined; - const runs = await store.listRuns({ - owner: q.owner, - workflow: q.workflow, - status: q.status, - limit: Number.isFinite(limit) ? limit : undefined, - before: q.before, - }); - return reply.send({ runs }); - }, - ); - - server.get( - "/admin/runs/:id", - { onRequest: [server.authenticate] }, - async (req, reply) => { - const { id } = /** @type {{ id: string }} */ (req.params); - const run = await store.getRun(id); - if (!run) { - return reply.code(404).send({ error: "run not found" }); - } - return reply.send(run); - }, - ); - - if (fs.existsSync(WEB_DIST)) { - await server.register(fastifyStatic, { - root: WEB_DIST, - wildcard: false, - }); - server.setNotFoundHandler((req, reply) => { - const url = req.raw.url ?? ""; - if ( - url.startsWith("/api") || - url.startsWith("/u/") || - url.startsWith("/admin") - ) { - return reply.code(404).send({ error: "not found" }); - } - return reply.sendFile("index.html"); - }); - } -} - -async function shutdown() { - try { - if (workflowWorker) { - await workflowWorker.close(); - } - await workflowQueue.close(); - await closeRedis(); - await flushLogs(); - await db.destroy(); - } catch (err) { - log.error({ err }, "shutdown error"); - } - process.exit(0); -} - -process.on("SIGINT", shutdown); -process.on("SIGTERM", shutdown); - -const port = Number(process.env.PORT ?? 8700); - -if (runApi) { - server - .listen({ - host: "0.0.0.0", - port, - }) - .then(() => { - log.info(`Server is running on port ${port}`); - }) - .catch((err) => { - log.error({ err }, "failed to start server"); - process.exit(1); - }); -} else { - log.info("worker-only mode; HTTP server not started"); -} \ No newline at end of file diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index 067ad02..7996f7a 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -6,7 +6,7 @@ import axios from "axios"; import pino from "pino"; import { createKvApi } from "./kv-store.js"; import { createFingerprintApi } from "./script-fingerprint.js"; -import { SCRIPTS_DIR } from "./paths.js"; +import { resolveScriptRef, createPluginRequire } from "./plugin-store.js"; import { isSecret, Secret, unwrapSecretsDeep } from "./secret-value.js"; import { getHttpPageByName, getHttpTemplateByName } from "./http-pages-store.js"; import { getSecretPlaintext } from "./secrets-store.js"; @@ -296,15 +296,60 @@ function createScreenedAxios(log) { }); } -function createRestrictedRequire(screenedAxios) { +const BLOCKED_PLUGIN_MODULES = new Set([ + "child_process", + "node:child_process", + "cluster", + "node:cluster", + "fs", + "node:fs", + "fs/promises", + "node:fs/promises", + "module", + "node:module", + "vm", + "node:vm", + "worker_threads", + "node:worker_threads", + "v8", + "node:v8", + "inspector", + "node:inspector", + "sqlite", + "node:sqlite", +]); + +/** + * @param {import("axios").AxiosInstance} screenedAxios + * @param {string | null} [pluginDirectory] + */ +function createRestrictedRequire(screenedAxios, pluginDirectory = null) { + const pluginRequire = pluginDirectory + ? createPluginRequire(pluginDirectory) + : null; + return function restrictedRequire(id) { - if (typeof id !== "string" || !ALLOWED_MODULES.has(id)) { + if (typeof id !== "string") { throw new Error(`require(${JSON.stringify(id)}) is not allowed`); } - if (id === "axios") { - return screenedAxios; + if (BLOCKED_PLUGIN_MODULES.has(id)) { + throw new Error(`require(${JSON.stringify(id)}) is not allowed`); } - return hostRequire(id); + if (ALLOWED_MODULES.has(id)) { + if (id === "axios") return screenedAxios; + return hostRequire(id); + } + if (pluginRequire) { + try { + return pluginRequire(id); + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + throw new Error( + `require(${JSON.stringify(id)}) failed in plugin: ${msg}`, + ); + } + } + throw new Error(`require(${JSON.stringify(id)}) is not allowed`); }; } @@ -391,6 +436,7 @@ const $workflowsStub = { * workflowName: string, * owner?: string, * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * pluginDir?: string | null, * }} opts */ function createScriptSandbox({ @@ -399,6 +445,7 @@ function createScriptSandbox({ workflowName, owner = "default", $workflows = $workflowsStub, + pluginDir = null, }) { const scriptLog = log.child({ workflow: workflowName, script }); const $axios = createScreenedAxios(scriptLog); @@ -418,7 +465,7 @@ function createScriptSandbox({ $vars, $responses, $workflows, - require: createRestrictedRequire($axios), + require: createRestrictedRequire($axios, pluginDir), }; vm.createContext(sandbox, { @@ -470,8 +517,29 @@ export function extractScriptMeta(fn) { * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, * }} opts */ -function instantiateCompiled(compiled, { log, script, workflowName, owner, $workflows }) { - const sandbox = createScriptSandbox({ log, script, workflowName, owner, $workflows }); +/** + * @param {import("vm").Script} compiled + * @param {{ + * log: import("pino").Logger, + * script: string, + * workflowName: string, + * owner?: string, + * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * pluginDir?: string | null, + * }} opts + */ +function instantiateCompiled( + compiled, + { log, script, workflowName, owner, $workflows, pluginDir = null }, +) { + const sandbox = createScriptSandbox({ + log, + script, + workflowName, + owner, + $workflows, + pluginDir, + }); return compiled.runInContext(sandbox); } @@ -486,6 +554,7 @@ function instantiateCompiled(compiled, { log, script, workflowName, owner, $work * workflowName?: string, * owner?: string, * $workflows?: { trigger: (name: string, data?: unknown) => Promise }, + * pluginDir?: string | null, * }} [opts] */ export function instantiateScriptSource(script, source, opts = {}) { @@ -496,6 +565,7 @@ export function instantiateScriptSource(script, source, opts = {}) { workflowName: opts.workflowName ?? "inspect", owner: opts.owner ?? "default", $workflows: opts.$workflows, + pluginDir: opts.pluginDir ?? null, }); return { fn, ...extractScriptMeta(fn) }; } @@ -519,17 +589,28 @@ export function inspectScriptSource(script, source) { } function loadCompiledScript(script) { - const filePath = path.join(SCRIPTS_DIR, script); + const resolved = resolveScriptRef(script); + if (resolved.error || !resolved.filePath) { + throw new Error(resolved.error || `script not found: ${script}`); + } + const filePath = resolved.filePath; const { mtimeMs } = fs.statSync(filePath); - const cached = scriptCache.get(script); + const cacheKey = `${resolved.kind}:${resolved.scriptRef}:${filePath}`; + const cached = scriptCache.get(cacheKey); if (cached && cached.mtimeMs === mtimeMs) { - return cached.compiled; + return cached; } const source = fs.readFileSync(filePath, "utf8"); const compiled = compileScriptSource(source, filePath); - scriptCache.set(script, { compiled, mtimeMs }); - return compiled; + const entry = { + compiled, + mtimeMs, + pluginDir: resolved.pluginDir ?? null, + scriptRef: resolved.scriptRef, + }; + scriptCache.set(cacheKey, entry); + return entry; } /** @@ -545,13 +626,14 @@ function loadCompiledScript(script) { * }} opts */ export async function runScript(script, ctx, { log, workflowName, owner, $workflows }) { - const compiled = loadCompiledScript(script); - const fn = instantiateCompiled(compiled, { + const loaded = loadCompiledScript(script); + const fn = instantiateCompiled(loaded.compiled, { log, - script, + script: loaded.scriptRef, workflowName, owner, $workflows, + pluginDir: loaded.pluginDir, }); return await fn(ctx); } diff --git a/packages/server/server.js b/packages/server/server.js new file mode 100644 index 0000000..4f3c1a7 --- /dev/null +++ b/packages/server/server.js @@ -0,0 +1,11 @@ +process.env.JFLOW_ROLE = "api"; + +import { startApp } from "./start-app.js"; + +// HTTP API + cron enqueue only. Schema migrations are owned by control. +await startApp({ + role: "api", + migrate: false, + // In PM2/split mode the Vite/control process serves the SPA. + serveStaticUi: process.env.JFLOW_SERVE_UI === "1", +}); diff --git a/packages/server/src/api/scripts.js b/packages/server/src/api/scripts.js index 6d175d0..989f0b1 100644 --- a/packages/server/src/api/scripts.js +++ b/packages/server/src/api/scripts.js @@ -1,13 +1,32 @@ import fs from "fs"; +import path from "path"; import { clearScriptCache, inspectScriptSource, instantiateScriptSource, } from "../../script-sandbox.js"; import * as fsStore from "../../fs-store.js"; +import { + forkCoreScript, + getInstalledPlugin, + listCoreScriptNames, + listInstalledPlugins, + resolveScriptRef, + uninstallPlugin, + createBlankPlugin, +} from "../../plugin-store.js"; +import { + installExamplePlugin, + installPluginFromDirectory, + installPluginFromGit, + installPluginFromZipBuffer, +} from "../../plugin-install.js"; import { createDryRunLogger, safeSerialize } from "./dry-run-logger.js"; 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"; /** * @param {{ referencedScripts: () => Set }} registry @@ -18,19 +37,58 @@ export default function scriptsPluginFactory(registry) { */ return async function scriptsPlugin(fastify) { fastify.get("/scripts", async () => { - const scripts = fsStore.listScriptFiles().map((name) => { + const core = listCoreScriptNames().map((name) => { const content = fsStore.readScript(name); const inspected = content == null ? { meta: null, metaError: "script not found" } : inspectScriptSource(name, content); - return { name, hasIcon: fsStore.scriptHasIcon(name), ...inspected }; + return { + name, + kind: "core", + editable: false, + hasIcon: fsStore.scriptHasIcon(name), + ...inspected, + }; }); - return { scripts }; + + const plugins = listInstalledPlugins().map((p) => { + let inspected = { meta: null, metaError: null }; + if (!p.disabled && p.manifest) { + try { + const mainPath = path.join(p.dir, p.manifest.main); + const content = fs.readFileSync(mainPath, "utf8"); + inspected = inspectScriptSource(p.scriptRef, content); + } catch (err) { + inspected = { + meta: null, + metaError: err instanceof Error ? err.message : String(err), + }; + } + } else if (p.compatError) { + inspected = { meta: null, metaError: p.compatError }; + } + return { + name: p.scriptRef, + kind: "plugin", + editable: true, + pluginId: p.id, + disabled: p.disabled, + version: p.manifest?.version ?? null, + hasIcon: false, + ...inspected, + }; + }); + + return { + scripts: [...core, ...plugins], + appVersion: getAppVersion(), + }; }); fastify.get("/scripts/:name/icon", async (req, reply) => { const { name } = /** @type {{ name: string }} */ (req.params); + // Icons only for core scripts today try { fsStore.assertScriptName(name); } catch (err) { @@ -46,64 +104,152 @@ export default function scriptsPluginFactory(registry) { }); fastify.get("/scripts/:name", async (req, reply) => { - const { name } = /** @type {{ name: string }} */ (req.params); - try { - fsStore.assertScriptName(name); - } catch (err) { - return reply.code(err.statusCode ?? 400).send({ error: err.message }); + const rawName = decodeURIComponent( + /** @type {{ name: string }} */ (req.params).name, + ); + const resolved = resolveScriptRef(rawName); + if (resolved.error || !resolved.filePath) { + return reply + .code(404) + .send({ error: resolved.error || "script not found" }); } - const content = fsStore.readScript(name); - if (content == null) return reply.code(404).send({ error: "script not found" }); - return { name, content, hasIcon: fsStore.scriptHasIcon(name), ...inspectScriptSource(name, content) }; + const content = fs.readFileSync(resolved.filePath, "utf8"); + const inspected = inspectScriptSource(resolved.scriptRef, content); + return { + name: resolved.scriptRef, + kind: resolved.kind, + editable: resolved.kind === "plugin", + pluginId: resolved.pluginId ?? null, + content, + hasIcon: + resolved.kind === "core" + ? fsStore.scriptHasIcon(resolved.scriptRef) + : false, + ...inspected, + }; }); + // Core scripts are read-only. Creating/editing bare *.js writes is disabled. + // New user scripts must be plugins (fork / zip / git). fastify.put("/scripts/:name", async (req, reply) => { - const { name } = /** @type {{ name: string }} */ (req.params); - try { - fsStore.assertScriptName(name); - } catch (err) { - return reply.code(err.statusCode ?? 400).send({ error: err.message }); + const rawName = decodeURIComponent( + /** @type {{ name: string }} */ (req.params).name, + ); + const pluginRef = resolveScriptRef(rawName); + if (pluginRef.kind === "core" || !rawName.startsWith("plugin/")) { + // Attempt to treat as core name + try { + fsStore.assertScriptName( + rawName.endsWith(".js") ? rawName : `${rawName}.js`, + ); + } catch { + // continue + } + if (!rawName.startsWith("plugin/")) { + return reply.code(403).send({ + error: + "core scripts are read-only; fork to a plugin or install a plugin", + }); + } } + const body = /** @type {{ content?: string }} */ (req.body ?? {}); if (typeof body.content !== "string") { return reply.code(400).send({ error: "content is required" }); } - const existed = fsStore.readScript(name) != null; - fsStore.writeScript(name, body.content); + + const resolved = resolveScriptRef(rawName); + if (resolved.kind !== "plugin" || !resolved.filePath || !resolved.pluginDir) { + return reply.code(404).send({ + error: resolved.error || "plugin not found (install or fork first)", + }); + } + if (resolved.disabled) { + return reply.code(409).send({ error: resolved.error || "plugin disabled" }); + } + + fs.writeFileSync(resolved.filePath, body.content, "utf8"); clearScriptCache(); - return reply.code(existed ? 200 : 201).send({ - name, - ...inspectScriptSource(name, body.content), + return reply.send({ + name: resolved.scriptRef, + kind: "plugin", + editable: true, + ...inspectScriptSource(resolved.scriptRef, body.content), }); }); fastify.delete("/scripts/:name", async (req, reply) => { - const { name } = /** @type {{ name: string }} */ (req.params); - try { - fsStore.assertScriptName(name); - } catch (err) { - return reply.code(err.statusCode ?? 400).send({ error: err.message }); + const rawName = decodeURIComponent( + /** @type {{ name: string }} */ (req.params).name, + ); + if (!rawName.startsWith("plugin/")) { + return reply.code(403).send({ + error: "core scripts cannot be deleted", + }); } - if (registry.referencedScripts().has(name)) { + const resolved = resolveScriptRef(rawName); + const id = resolved.pluginId; + if (!id) { + // may be installed but disabled — still allow uninstall via plugin id parse + const installed = listInstalledPlugins().find( + (p) => p.scriptRef === rawName || `plugin/${p.id}` === rawName, + ); + if (!installed) { + return reply.code(404).send({ error: "plugin not found" }); + } + if (registry.referencedScripts().has(installed.scriptRef)) { + return reply + .code(409) + .send({ error: "plugin is referenced by a workflow" }); + } + uninstallPlugin(installed.id); + clearScriptCache(); + return { ok: true, restartNeeded: true }; + } + if (registry.referencedScripts().has(resolved.scriptRef)) { return reply .code(409) - .send({ error: "script is referenced by a workflow" }); - } - if (!fsStore.deleteScript(name)) { - return reply.code(404).send({ error: "script not found" }); + .send({ error: "plugin is referenced by a workflow" }); } + uninstallPlugin(id); clearScriptCache(); - return { ok: true }; + return { ok: true, restartNeeded: true }; + }); + + fastify.post("/scripts/:name/fork", 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" }); + } + try { + const coreName = rawName.endsWith(".js") ? rawName : `${rawName}.js`; + fsStore.assertScriptName(coreName); + const installed = forkCoreScript(coreName, 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 { name } = /** @type {{ name: string }} */ (req.params); - try { - fsStore.assertScriptName(name); - } catch (err) { - return reply.code(err.statusCode ?? 400).send({ error: err.message }); - } - + const rawName = decodeURIComponent( + /** @type {{ name: string }} */ (req.params).name, + ); const body = /** @type {{ content?: string, data?: unknown, context?: unknown, config?: unknown, owner?: string }} */ ( req.body ?? {} ); @@ -121,10 +267,18 @@ export default function scriptsPluginFactory(registry) { } const incomingContext = - body.context != null && typeof body.context === "object" && !Array.isArray(body.context) + body.context != null && + typeof body.context === "object" && + !Array.isArray(body.context) ? body.context : {}; + const resolved = resolveScriptRef(rawName); + const pluginDir = + resolved.kind === "plugin" && !resolved.error + ? resolved.pluginDir ?? null + : null; + const { log, logs } = createDryRunLogger(); const started = Date.now(); @@ -139,13 +293,22 @@ export default function scriptsPluginFactory(registry) { context: incomingContext, config, }; - const { fn, meta, metaError } = instantiateScriptSource(name, body.content, { - log, - workflowName: "dry-run", - owner, - }); + const { fn, meta, metaError } = instantiateScriptSource( + resolved.scriptRef || rawName, + body.content, + { + log, + workflowName: "dry-run", + owner, + pluginDir, + }, + ); const raw = await fn(ctx); - const result = normalizeStepResult(raw, incomingContext, name); + const result = normalizeStepResult( + raw, + incomingContext, + resolved.scriptRef || rawName, + ); return { status: "success", output: safeSerialize(result.output), @@ -158,7 +321,7 @@ export default function scriptsPluginFactory(registry) { metaError, }; } catch (err) { - const inspected = inspectScriptSource(name, body.content); + const inspected = inspectScriptSource(rawName, body.content); return { status: "failed", output: null, @@ -171,5 +334,148 @@ export default function scriptsPluginFactory(registry) { }; } }); + + // --- Plugins --- + + fastify.get("/plugins", async () => { + return { + appVersion: getAppVersion(), + plugins: listInstalledPlugins(), + warning: + "Installing plugins runs third-party code as the JerapahFlow OS user.", + }; + }); + + fastify.post( + "/plugins/create", + { onRequest: [fastify.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ id?: string, content?: string, description?: string }} */ ( + req.body ?? {} + ); + if (typeof body.id !== "string" || !body.id.trim()) { + return reply.code(400).send({ error: "id is required" }); + } + if (typeof body.content !== "string") { + return reply.code(400).send({ error: "content is required" }); + } + try { + const installed = createBlankPlugin(body.id.trim(), body.content, { + 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( + "/plugins/install", + { onRequest: [fastify.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ + source?: string, + url?: string, + ref?: string, + path?: string, + exampleId?: string, + zipBase64?: string, + overwrite?: boolean, + }} */ (req.body ?? {}); + + try { + let installed; + if (body.source === "git") { + if (!body.url) { + return reply.code(400).send({ error: "url is required" }); + } + installed = await installPluginFromGit(body.url, { + ref: body.ref, + overwrite: Boolean(body.overwrite), + }); + } else if (body.source === "example") { + const id = body.exampleId || "get-current-time"; + installed = await installExamplePlugin(id, { + overwrite: Boolean(body.overwrite), + }); + } else if (body.source === "dir") { + if (!body.path) { + return reply.code(400).send({ error: "path is required" }); + } + // Only allow examples/ or existing staging under plugins for safety + const abs = path.resolve(body.path); + const allowed = + abs.startsWith(EXAMPLE_PLUGINS_DIR + path.sep) || + abs.startsWith(EXAMPLE_PLUGINS_DIR); + if (!allowed) { + return reply.code(403).send({ + error: "dir install only allowed under examples/plugins", + }); + } + installed = installPluginFromDirectory(abs, { + overwrite: Boolean(body.overwrite), + }); + } else if (body.source === "zip") { + if (!body.zipBase64) { + return reply.code(400).send({ error: "zipBase64 is required" }); + } + const buf = Buffer.from(body.zipBase64, "base64"); + installed = await installPluginFromZipBuffer(buf, { + overwrite: Boolean(body.overwrite), + }); + } else { + return reply.code(400).send({ + error: "source must be git | zip | example | dir", + }); + } + clearScriptCache(); + return reply.code(201).send({ + ...installed, + scriptRef: pluginScriptRef(installed.id), + restartNeeded: true, + warning: + "Plugins run as the JerapahFlow process user. Review code before install. Drain-restart workers to load new dependencies.", + }); + } catch (err) { + return reply + .code(/** @type {any} */ (err).statusCode ?? 500) + .send({ error: err instanceof Error ? err.message : String(err) }); + } + }, + ); + + fastify.delete( + "/plugins/:id", + { onRequest: [fastify.requireAdmin] }, + async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); + try { + const scriptRef = pluginScriptRef(id); + if (registry.referencedScripts().has(scriptRef)) { + return reply + .code(409) + .send({ error: "plugin is referenced by a workflow" }); + } + uninstallPlugin(id); + clearScriptCache(); + return { ok: true, restartNeeded: true }; + } catch (err) { + return reply + .code(/** @type {any} */ (err).statusCode ?? 500) + .send({ error: err instanceof Error ? err.message : String(err) }); + } + }, + ); + + void getInstalledPlugin; }; } diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index dd73fc1..2904c38 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -16,6 +16,20 @@ import { ensureWorkflowFilename, suggestCopyFilename, } from "../../workflow-duplicate.js"; +import { publishReload } from "../../control-bus.js"; + +/** + * Reload this process and notify other HTTP/worker processes via Redis. + * @param {{ reregister: () => void }} registry + */ +async function reregisterAll(registry) { + registry.reregister(); + try { + await publishReload({ type: "workflows" }); + } catch { + // Redis may be briefly unavailable; local reload already applied. + } +} function triggerSummary(owner, workflow) { if (!workflow || typeof workflow !== "object") return []; @@ -214,7 +228,7 @@ export default function workflowsPluginFactory(registry) { registered.push(file); fsStore.writeRegisters(owner, registered); } - registry.reregister(); + await reregisterAll(registry); return reply.code(existed ? 200 : 201).send({ owner, file }); }); @@ -251,7 +265,7 @@ export default function workflowsPluginFactory(registry) { doc.set("enabled", false); } fsStore.writeWorkflowYaml(owner, file, String(doc)); - registry.reregister(); + await reregisterAll(registry); return { owner, file, enabled: body.enabled }; }); @@ -270,7 +284,7 @@ export default function workflowsPluginFactory(registry) { } const registered = fsStore.readRegisters(owner).filter((f) => f !== file); fsStore.writeRegisters(owner, registered); - registry.reregister(); + await reregisterAll(registry); return { ok: true }; }); @@ -372,7 +386,7 @@ export default function workflowsPluginFactory(registry) { registered.push(destFile); fsStore.writeRegisters(destOwner, registered); } - registry.reregister(); + await reregisterAll(registry); return reply.code(201).send({ owner: destOwner, file: destFile }); }); @@ -394,9 +408,9 @@ export default function workflowsPluginFactory(registry) { if (!registered.includes(file)) { registered.push(file); fsStore.writeRegisters(owner, registered); - registry.reregister(); + await reregisterAll(registry); } else if (!registry.workflows.has(key) && !registry.loadErrors.has(key)) { - registry.reregister(); + await reregisterAll(registry); } if (!registry.workflows.has(key)) { return reply.code(404).send({ @@ -424,7 +438,7 @@ export default function workflowsPluginFactory(registry) { }); fastify.post("/workflows/reregister", async () => { - registry.reregister(); + await reregisterAll(registry); return { message: "Workflows refreshed" }; }); }; diff --git a/packages/server/start-app.js b/packages/server/start-app.js new file mode 100644 index 0000000..6f2ddb7 --- /dev/null +++ b/packages/server/start-app.js @@ -0,0 +1,292 @@ +import fs from "fs"; +import fastify from "fastify"; +import cookie from "@fastify/cookie"; +import cors from "@fastify/cors"; +import jwt from "@fastify/jwt"; +import fastifyStatic from "@fastify/static"; +import { migrate, db } from "./db.js"; +import { log, enableLogPersistence, flushLogs } from "./logger.js"; +import * as store from "./store.js"; +import { createRegistry } from "./registry.js"; +import { COOKIE, OPEN_API_ROUTES } from "./src/api/auth.js"; +import authPlugin from "./src/api/auth.js"; +import usersPlugin from "./src/api/users.js"; +import scriptsPluginFactory from "./src/api/scripts.js"; +import workflowsPluginFactory from "./src/api/workflows.js"; +import runsPlugin from "./src/api/runs.js"; +import dashboardPluginFactory from "./src/api/dashboard.js"; +import secretsPlugin from "./src/api/secrets.js"; +import kvPlugin from "./src/api/kv.js"; +import variablesPlugin from "./src/api/variables.js"; +import httpPagesPlugin from "./src/api/http-pages.js"; +import httpAuthsPlugin from "./src/api/http-auths.js"; +import { WEB_DIST } from "./paths.js"; +import { resolveSecretsKeyMaterial } from "./secrets.js"; +import { + closeRedis, + createWorkflowQueue, + createWorkflowWorker, + getRedisUrlForLog, +} from "./workflow-queue.js"; +import { + getConfigGeneration, + startHeartbeatLoop, + subscribeReload, +} from "./control-bus.js"; + +/** + * @param {{ + * role?: string, + * migrate?: boolean, + * serveStaticUi?: boolean, + * }} [opts] + */ +export async function startApp(opts = {}) { + const role = (opts.role || process.env.JFLOW_ROLE || "all").toLowerCase(); + const shouldMigrate = opts.migrate ?? role === "all"; + const serveStaticUi = opts.serveStaticUi ?? (role === "all" || role === "api"); + + if (shouldMigrate) { + await migrate(); + } + enableLogPersistence(); + + const jwtSecret = + process.env.JFLOW_JWT_SECRET ?? + (process.env.NODE_ENV === "production" ? "" : "jflow-dev-secret"); + + if (!jwtSecret) { + log.error("JFLOW_JWT_SECRET is required in production"); + process.exit(1); + } + + try { + resolveSecretsKeyMaterial(); + } catch (err) { + log.error(err instanceof Error ? err.message : String(err)); + process.exit(1); + } + + const runApi = role === "all" || role === "api"; + const runWorker = role === "all" || role === "worker"; + + log.info( + { + redis: getRedisUrlForLog(), + role, + generation: getConfigGeneration(), + migrate: shouldMigrate, + }, + "starting jerapah-flow", + ); + + const workflowQueue = createWorkflowQueue(); + try { + await workflowQueue.waitUntilReady(); + } catch (err) { + log.error({ err, redis: getRedisUrlForLog() }, "failed to connect to Redis"); + process.exit(1); + } + + const server = fastify({ loggerInstance: log }); + + await server.register(cookie); + await server.register(jwt, { + secret: jwtSecret, + cookie: { + cookieName: COOKIE, + signed: false, + }, + }); + await server.register(cors, { + origin: process.env.JFLOW_CORS_ORIGIN ?? "http://localhost:8500", + credentials: true, + }); + + server.decorate("authenticate", async function authenticate(req, reply) { + try { + await req.jwtVerify(); + } catch { + return reply.code(401).send({ error: "unauthorized" }); + } + }); + + server.decorate("requireAdmin", async function requireAdmin(req, reply) { + if (req.user?.role !== "admin") { + return reply.code(403).send({ error: "forbidden" }); + } + }); + + const registry = createRegistry(server, { queue: workflowQueue }); + registry.registerWorkflows(); + if (runApi) { + registry.registerHttpTriggers(); + registry.registerCronTriggers(); + registry.registerPruneJob(); + } + + /** @type {import("bullmq").Worker | null} */ + let workflowWorker = null; + if (runWorker) { + workflowWorker = createWorkflowWorker(async (job) => { + const data = /** @type {{ runId?: string, key?: string, depth?: number }} */ ( + job.data ?? {} + ); + if (typeof data.runId !== "string" || typeof data.key !== "string") { + throw new Error("invalid workflow job payload"); + } + const result = await registry.executeQueuedRun({ + runId: data.runId, + key: data.key, + depth: data.depth ?? 0, + }); + if (result.status === "failed") { + throw new Error(result.error || "workflow failed"); + } + return result; + }); + } + + if (runApi) { + await server.register( + async (api) => { + api.addHook("onRequest", async (req, reply) => { + const raw = (req.url || "").split("?")[0]; + const stripped = raw.replace(/^\/api/, "") || "/"; + const routeUrl = req.routeOptions?.url || stripped; + const open = + OPEN_API_ROUTES.has(`${req.method} ${routeUrl}`) || + OPEN_API_ROUTES.has(`${req.method} ${stripped}`); + if (open) return; + await server.authenticate(req, reply); + }); + await api.register(authPlugin); + await api.register(usersPlugin); + await api.register(secretsPlugin); + await api.register(variablesPlugin); + await api.register(kvPlugin); + await api.register(httpPagesPlugin); + await api.register(httpAuthsPlugin); + await api.register(scriptsPluginFactory(registry)); + await api.register(workflowsPluginFactory(registry)); + await api.register(runsPlugin); + await api.register(dashboardPluginFactory(registry)); + }, + { prefix: "/api" }, + ); + + server.post( + "/admin/workflows/reregister", + { onRequest: [server.authenticate] }, + async (_req, reply) => { + registry.reregister(); + try { + const { publishReload } = await import("./control-bus.js"); + await publishReload({ type: "workflows" }); + } catch { + // ignore + } + return reply.send({ message: "Workflows refreshed" }); + }, + ); + + server.get( + "/admin/runs", + { onRequest: [server.authenticate] }, + async (req, reply) => { + const q = /** @type {Record} */ (req.query); + const limit = q.limit ? Number(q.limit) : undefined; + const runs = await store.listRuns({ + owner: q.owner, + workflow: q.workflow, + status: q.status, + limit: Number.isFinite(limit) ? limit : undefined, + before: q.before, + }); + return reply.send({ runs }); + }, + ); + + server.get( + "/admin/runs/:id", + { onRequest: [server.authenticate] }, + async (req, reply) => { + const { id } = /** @type {{ id: string }} */ (req.params); + const run = await store.getRun(id); + if (!run) { + return reply.code(404).send({ error: "run not found" }); + } + return reply.send(run); + }, + ); + + if (serveStaticUi && fs.existsSync(WEB_DIST)) { + await server.register(fastifyStatic, { + root: WEB_DIST, + wildcard: false, + }); + server.setNotFoundHandler((req, reply) => { + const url = req.raw.url ?? ""; + if ( + url.startsWith("/api") || + url.startsWith("/u/") || + url.startsWith("/admin") || + url.startsWith("/ops") + ) { + return reply.code(404).send({ error: "not found" }); + } + return reply.sendFile("index.html"); + }); + } + } + + const stopHeartbeat = startHeartbeatLoop(runApi && runWorker ? "all" : runApi ? "api" : "worker"); + + const reloadSub = await subscribeReload(async () => { + log.info("received reload signal"); + registry.reregister(); + }); + + async function shutdown() { + try { + stopHeartbeat(); + await reloadSub.stop(); + if (workflowWorker) { + await workflowWorker.close(); + } + await workflowQueue.close(); + await closeRedis(); + await flushLogs(); + await db.destroy(); + } catch (err) { + log.error({ err }, "shutdown error"); + } + process.exit(0); + } + + process.on("SIGINT", shutdown); + process.on("SIGTERM", shutdown); + + const port = Number(process.env.PORT ?? 8700); + + if (runApi) { + try { + await server.listen({ host: "0.0.0.0", port }); + log.info(`Server is running on port ${port}`); + } catch (err) { + log.error({ err }, "failed to start server"); + process.exit(1); + } + } else { + log.info("worker-only mode; HTTP server not started"); + } + + return { + role, + server, + workflowQueue, + workflowWorker, + registry, + shutdown, + }; +} diff --git a/packages/server/test/plugins-smoke.js b/packages/server/test/plugins-smoke.js new file mode 100644 index 0000000..29776e0 --- /dev/null +++ b/packages/server/test/plugins-smoke.js @@ -0,0 +1,118 @@ +/** + * Smoke: core vs plugin scripts, fork, example install, resolve, run. + * + * Run: + * 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 { migrate, db } from "../db.js"; +import { getAppVersion, satisfiesRange } from "../app-version.js"; +import { + forkCoreScript, + resolveScriptRef, + uninstallPlugin, + listInstalledPlugins, + createBlankPlugin, +} from "../plugin-store.js"; +import { installExamplePlugin } from "../plugin-install.js"; +import { + runScript, + clearScriptCache, + inspectScriptSource, +} from "../script-sandbox.js"; +import { PLUGINS_DIR } from "../paths.js"; +import pino from "pino"; + +const silent = pino({ level: "silent" }); + +async function main() { + assert.equal(getAppVersion(), "0.1.0"); + assert.equal(satisfiesRange("0.1.0", ">=0.1.0 <1.0.0"), true); + assert.equal(satisfiesRange("1.0.0", ">=0.1.0 <1.0.0"), false); + assert.equal(satisfiesRange("0.2.0", ">=0.1.0 <1.0.0"), true); + + await migrate(); + + if (fs.existsSync(PLUGINS_DIR)) { + fs.rmSync(PLUGINS_DIR, { recursive: true, force: true }); + } + fs.mkdirSync(PLUGINS_DIR, { recursive: true }); + + const core = resolveScriptRef("fetch-http.js"); + assert.equal(core.kind, "core"); + assert.ok(core.filePath && fs.existsSync(core.filePath)); + assert.ok(resolveScriptRef("nope.js").error?.includes("not found")); + + const example = await installExamplePlugin("get-current-time", { + overwrite: true, + }); + assert.equal(example.id, "get-current-time"); + assert.equal(example.scriptRef, "plugin/get-current-time"); + + clearScriptCache(); + const pluginResolved = resolveScriptRef("plugin/get-current-time"); + assert.equal(pluginResolved.kind, "plugin"); + assert.ok(pluginResolved.filePath); + + const result = await runScript( + "plugin/get-current-time", + { data: null, context: {}, config: null }, + { log: silent, workflowName: "smoke", owner: "default" }, + ); + assert.ok(result?.output?.datetime); + + const forked = forkCoreScript("jsonata.js", "jsonata-smoke-fork"); + assert.equal(forked.scriptRef, "plugin/jsonata-smoke-fork"); + clearScriptCache(); + assert.equal(resolveScriptRef("plugin/jsonata-smoke-fork").kind, "plugin"); + + const blank = createBlankPlugin( + "blank-smoke", + `export default async function main(ctx) { return { output: { ok: true }, context: ctx.context ?? {} }; }`, + ); + assert.equal(blank.scriptRef, "plugin/blank-smoke"); + clearScriptCache(); + const blankRun = await runScript( + "plugin/blank-smoke", + { data: 1, context: {}, config: null }, + { log: silent, workflowName: "smoke", owner: "default" }, + ); + assert.equal(blankRun.output.ok, true); + + let hit = false; + try { + forkCoreScript("ntfy.js", "ntfy"); + } catch (err) { + hit = true; + assert.match(String(err.message), /collides/); + } + assert.equal(hit, true); + + const meta = inspectScriptSource( + "fetch-http.js", + fs.readFileSync(core.filePath, "utf8"), + ); + assert.ok(meta); + + assert.ok(listInstalledPlugins().some((p) => p.id === "get-current-time")); + + uninstallPlugin("jsonata-smoke-fork"); + uninstallPlugin("blank-smoke"); + uninstallPlugin("get-current-time"); + + console.log("plugins-smoke: ok"); + await db.destroy(); +} + +main().catch(async (err) => { + console.error(err); + try { + await db.destroy(); + } catch { + // ignore + } + process.exit(1); +}); diff --git a/packages/server/worker.js b/packages/server/worker.js new file mode 100644 index 0000000..8da34de --- /dev/null +++ b/packages/server/worker.js @@ -0,0 +1,10 @@ +process.env.JFLOW_ROLE = "worker"; + +import { startApp } from "./start-app.js"; + +// BullMQ worker only. Schema migrations are owned by control. +await startApp({ + role: "worker", + migrate: false, + serveStaticUi: false, +}); diff --git a/packages/server/workflows/default/cron-example.yaml b/packages/server/workflows/default/cron-example.yaml index b829f36..f29501a 100644 --- a/packages/server/workflows/default/cron-example.yaml +++ b/packages/server/workflows/default/cron-example.yaml @@ -2,7 +2,7 @@ name: cron example description: | this workflow triggered by cron scripts: - - script: get-current-time.js + - script: plugin/get-current-time - set: expression: '{"message": context.datetime}' - script: ntfy.js diff --git a/packages/server/workflows/default/dev-joplin-daily.yaml b/packages/server/workflows/default/dev-joplin-daily.yaml new file mode 100644 index 0000000..4010125 --- /dev/null +++ b/packages/server/workflows/default/dev-joplin-daily.yaml @@ -0,0 +1,29 @@ +name: Joplin nightly daily log +description: | + POST joplin-auto /api/logs/run at 04:00 (sync + yearly/monthly/today), then ntfy. + Secret joplin_setup_token (SETUP_API_TOKEN). Variable ntfy_channel. +scripts: + - script: plugin/joplin-api + config: + url: http://10.8.0.6:3040/api/logs/run + method: POST + token: $SECRET_joplin_setup_token + timeoutMs: 600000 + - set: + expression: | + { + "title": data.ok ? "Bullet journal" : "Bullet journal failed", + "message": data.ok + ? "The bullet journal for " & data.httpResponse.date & " has been created." + : data.message + } + - script: ntfy.js + config: + url: $VAR_ntfy_channel +triggers: + - type: cron + schedule: "0 4 * * *" + - type: HTTP + method: POST + path: /dev-joplin-daily +enabled: false diff --git a/packages/server/workflows/default/dev-joplin-get-note.yaml b/packages/server/workflows/default/dev-joplin-get-note.yaml new file mode 100644 index 0000000..a3d22f1 --- /dev/null +++ b/packages/server/workflows/default/dev-joplin-get-note.yaml @@ -0,0 +1,17 @@ +name: Joplin get note +description: | + GET a Joplin note by id from joplin-api. Input (data): id (32-char hex). + Secret joplin_api_token (JOPLIN_API_TOKEN / API_KEYS). Returns the full note. +scripts: + - script: plugin/joplin-api + config: + url: http://10.8.0.6:3030/notes + method: GET + token: $SECRET_joplin_api_token + timeoutMs: 60000 +triggers: + - type: workflow + - type: HTTP + method: POST + path: /dev-joplin-get-note +enabled: false diff --git a/packages/server/workflows/default/dev-joplin-sync.yaml b/packages/server/workflows/default/dev-joplin-sync.yaml new file mode 100644 index 0000000..f87b5f0 --- /dev/null +++ b/packages/server/workflows/default/dev-joplin-sync.yaml @@ -0,0 +1,27 @@ +name: Joplin nightly sync +description: | + POST joplin-auto /api/sync at 03:00, then ntfy success or failure. + Secret joplin_setup_token (SETUP_API_TOKEN). Variable ntfy_channel. +scripts: + - script: plugin/joplin-api + config: + url: http://10.8.0.6:3040/api/sync + method: POST + token: $SECRET_joplin_setup_token + timeoutMs: 600000 + - set: + expression: | + { + "title": data.ok ? "Joplin sync ok" : "Joplin sync failed", + "message": data.message + } + - script: ntfy.js + config: + url: $VAR_ntfy_channel +triggers: + - type: cron + schedule: "0 3 * * *" + - type: HTTP + method: POST + path: /dev-joplin-sync +enabled: false diff --git a/packages/server/workflows/default/dev-zte-sms.yaml b/packages/server/workflows/default/dev-zte-sms.yaml new file mode 100644 index 0000000..e351541 --- /dev/null +++ b/packages/server/workflows/default/dev-zte-sms.yaml @@ -0,0 +1,16 @@ +name: dev-zte-sms +description: | + Send an SMS via a ZTE modem web UI. + Input (data): to, message + Password is the named secret zte_modem_password ($SECRET_). +scripts: + - script: plugin/send-sms + config: + url: http://192.168.5.1/reqproc/proc_post + password: $SECRET_sms_secret +triggers: + - type: workflow + - type: HTTP + method: POST + path: /dev-zte-sms + auth: basic-auth diff --git a/packages/server/workflows/default/test.yaml b/packages/server/workflows/default/test.yaml index 4c1cb48..e8ec686 100644 --- a/packages/server/workflows/default/test.yaml +++ b/packages/server/workflows/default/test.yaml @@ -1,6 +1,6 @@ name: jsonata scripts: - - get-current-time.js + - plugin/get-current-time - script: jsonata.js config: expression: '{"message": data.datetime & " " & data.processId}' diff --git a/packages/server/workflows/default/time-and-comic-to-ntfy.yaml b/packages/server/workflows/default/time-and-comic-to-ntfy.yaml index 4268411..ab1156d 100644 --- a/packages/server/workflows/default/time-and-comic-to-ntfy.yaml +++ b/packages/server/workflows/default/time-and-comic-to-ntfy.yaml @@ -3,7 +3,7 @@ description: > Fan-in from two scripts (current time + monkeyuser comic), then send to ntfy. scripts: - id: time - script: get-current-time.js + script: plugin/get-current-time - id: comic script: fetch-html.js diff --git a/packages/server/workflows/default/time-to-ntfy-example.yaml b/packages/server/workflows/default/time-to-ntfy-example.yaml index df6144c..34404f9 100644 --- a/packages/server/workflows/default/time-to-ntfy-example.yaml +++ b/packages/server/workflows/default/time-to-ntfy-example.yaml @@ -2,7 +2,7 @@ name: time to ntfy example description: | this workflow will send a message to ntfy with the current time scripts: - - script: get-current-time.js + - script: plugin/get-current-time config: key: "" - script: ntfy.js diff --git a/packages/web/package.json b/packages/web/package.json index 0c484f9..64e7136 100644 --- a/packages/web/package.json +++ b/packages/web/package.json @@ -1,7 +1,7 @@ { "name": "@jerapah-flow/web", "private": true, - "version": "1.0.0", + "version": "0.1.0", "type": "module", "scripts": { "dev": "vite", diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index e90e00f..579c487 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -19,6 +19,7 @@ import { ResponsesPage } from "./pages/ResponsesPage.jsx"; import { UsersPage } from "./pages/UsersPage.jsx"; import { SecretsPage } from "./pages/SecretsPage.jsx"; import { VariablesPage } from "./pages/VariablesPage.jsx"; +import { OpsPage } from "./pages/OpsPage.jsx"; export function App() { const qc = useQueryClient(); @@ -71,11 +72,13 @@ export function App() { <> } /> } /> + } /> ) : ( <> } /> } /> + } /> )} } /> diff --git a/packages/web/src/api/client.js b/packages/web/src/api/client.js index 5d9a7aa..d4f5c75 100644 --- a/packages/web/src/api/client.js +++ b/packages/web/src/api/client.js @@ -5,11 +5,21 @@ export const api = axios.create({ withCredentials: true, }); +/** Control-plane ops API (proxied to :8600 in `pnpm dev:pm2`). */ +export const opsApi = axios.create({ + baseURL: "/ops", + withCredentials: true, +}); + api.interceptors.response.use( (res) => res, (err) => { const url = err.config?.url ?? ""; - const isAuthCall = url.includes("/auth/login") || url.includes("/auth/register") || url.includes("/auth/me") || url.includes("/auth/bootstrap"); + const isAuthCall = + url.includes("/auth/login") || + url.includes("/auth/register") || + url.includes("/auth/me") || + url.includes("/auth/bootstrap"); if (err.response?.status === 401 && !isAuthCall) { window.dispatchEvent(new Event("jerapah-flow:unauthorized")); } @@ -17,6 +27,16 @@ api.interceptors.response.use( }, ); +opsApi.interceptors.response.use( + (res) => res, + (err) => { + if (err.response?.status === 401) { + window.dispatchEvent(new Event("jerapah-flow:unauthorized")); + } + return Promise.reject(err); + }, +); + export function errorMessage(err, fallback = "request failed") { return err?.response?.data?.error ?? err?.message ?? fallback; } diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index abda51e..d7fed6d 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -1,5 +1,5 @@ import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; -import { api } from "./client.js"; +import { api, opsApi } from "./client.js"; export function useBootstrap() { return useQuery({ @@ -89,6 +89,47 @@ export function useSaveScript() { qc.invalidateQueries({ queryKey: ["scripts"] }); qc.invalidateQueries({ queryKey: ["scripts", vars.name] }); qc.invalidateQueries({ queryKey: ["dashboard"] }); + qc.invalidateQueries({ queryKey: ["ops-status"] }); + }, + }); +} + +export function useCreatePlugin() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ id, content, description }) => + (await api.post("/plugins/create", { id, content, description })).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["scripts"] }); + qc.invalidateQueries({ queryKey: ["ops-status"] }); + }, + }); +} + +export function useForkScript() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ name, id, description }) => + ( + await api.post(`/scripts/${encodeURIComponent(name)}/fork`, { + id, + description, + }) + ).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["scripts"] }); + qc.invalidateQueries({ queryKey: ["ops-status"] }); + }, + }); +} + +export function useInstallPlugin() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (body) => (await api.post("/plugins/install", body)).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["scripts"] }); + qc.invalidateQueries({ queryKey: ["ops-status"] }); }, }); } @@ -450,3 +491,84 @@ export function useDeleteHttpAuth() { export async function fetchHttpAuthLiterals(name) { return (await api.get(`/http-auths/${encodeURIComponent(name)}/reveal`)).data; } + +export function useOpsStatus(enabled = true) { + return useQuery({ + queryKey: ["ops-status"], + queryFn: async () => (await opsApi.get("/status")).data, + enabled, + retry: false, + refetchInterval: enabled ? 3000 : false, + }); +} + +export function useOpsPause() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async () => (await opsApi.post("/pause")).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsResume() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async () => (await opsApi.post("/resume")).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsReload() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async () => (await opsApi.post("/reload")).data, + onSuccess: () => { + qc.invalidateQueries({ queryKey: ["ops-status"] }); + qc.invalidateQueries({ queryKey: ["workflows"] }); + }, + }); +} + +export function useOpsRestart() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ force = false } = {}) => + (await opsApi.post("/restart", { force })).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsScale() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ workers, force = false }) => + (await opsApi.post("/scale", { workers, force })).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsHttpStart() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async () => (await opsApi.post("/http/start")).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsHttpStop() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async () => (await opsApi.post("/http/stop")).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + +export function useOpsBumpGeneration() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async (reason) => + (await opsApi.post("/generation/bump", { reason })).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + diff --git a/packages/web/src/components/Layout.jsx b/packages/web/src/components/Layout.jsx index 501b400..737fded 100644 --- a/packages/web/src/components/Layout.jsx +++ b/packages/web/src/components/Layout.jsx @@ -1,4 +1,4 @@ -import { NavLink, useNavigate } from "react-router-dom"; +import { Link, NavLink, useNavigate } from "react-router-dom"; import { LuActivity, LuCode, @@ -10,12 +10,13 @@ import { LuLogOut, LuMenu, LuMoon, + LuServer, LuShield, LuSun, LuTags, LuUsers, } from "react-icons/lu"; -import { useLogout } from "../api/hooks.js"; +import { useLogout, useOpsStatus } from "../api/hooks.js"; import { brandMark } from "../theme/brand.js"; import { useTheme } from "../theme.jsx"; @@ -34,6 +35,8 @@ export function Layout({ user, children }) { const { theme, toggle } = useTheme(); const logout = useLogout(); const navigate = useNavigate(); + const ops = useOpsStatus(user.role === "admin"); + const restartNeeded = Boolean(ops.data?.desired?.restartNeeded); const navItems = [ ...links, @@ -41,6 +44,7 @@ export function Layout({ user, children }) { ? [ { to: "/secrets", label: "Secrets", icon: LuKey }, { to: "/users", label: "Users", icon: LuUsers }, + { to: "/ops", label: "Ops", icon: LuServer }, ] : []), ]; @@ -88,7 +92,24 @@ export function Layout({ user, children }) { -
{children}
+
+ {restartNeeded ? ( +
+ + Config generation changed + {ops.data?.desired?.restartReason + ? ` (${ops.data.desired.restartReason})` + : ""} + .{" "} + + Drain restart from Ops + {" "} + to apply on HTTP + workers. + +
+ ) : null} + {children} +
@@ -92,8 +118,12 @@ export function ScriptEditPage() { const { notify } = useNotifications(); const existing = useScript(name); const save = useSaveScript(); + const fork = useForkScript(); const [content, setContent] = useState(""); const [contentReady, setContentReady] = useState(false); + const [forkId, setForkId] = useState(""); + + const isCore = existing.data?.kind === "core" || existing.data?.editable === false; useEffect(() => { if (existing.isLoading) return; @@ -105,9 +135,10 @@ export function ScriptEditPage() { function onSave(e) { e.preventDefault(); + if (isCore) return; save.mutate( { name, content }, - { onSuccess: () => notify.success("Script saved") }, + { onSuccess: () => notify.success("Plugin saved") }, ); } @@ -117,6 +148,21 @@ export function ScriptEditPage() { }); } + function onFork(e) { + e.preventDefault(); + const id = normalizePluginId(forkId); + if (!id) return; + fork.mutate( + { name, id }, + { + onSuccess: (data) => { + notify.success("Forked to plugin — drain-restart recommended"); + navigate(`/scripts/${encodeURIComponent(data.scriptRef)}/edit`); + }, + }, + ); + } + if (existing.isLoading) { return (
@@ -145,24 +191,63 @@ export function ScriptEditPage() {

{name}

+ + {isCore ? "core" : "plugin"} +
- + {!isCore ? ( + + ) : null}
+ {isCore ? ( +
+ Core scripts are read-only. Fork to create an editable plugin copy. +
+ ) : null} + + {isCore ? ( +
+ setForkId(e.target.value)} + /> + + {fork.isError ? ( + {errorMessage(fork.error)} + ) : null} +
+ ) : null} +
- + {} : setContent} + height="100%" + readOnly={isCore} + />
diff --git a/packages/web/src/pages/ScriptsPage.jsx b/packages/web/src/pages/ScriptsPage.jsx index 83aa616..a53c358 100644 --- a/packages/web/src/pages/ScriptsPage.jsx +++ b/packages/web/src/pages/ScriptsPage.jsx @@ -1,11 +1,17 @@ import { useMemo, useState } from "react"; import { Link, Navigate, useNavigate, useSearchParams } from "react-router-dom"; -import { LuPencil, LuPlay, LuPlus, LuSearch, LuTrash2 } from "react-icons/lu"; +import { LuCopy, LuPencil, LuPlay, LuPlus, LuSearch, LuTrash2 } from "react-icons/lu"; import { errorMessage } from "../api/client.js"; -import { useDeleteScript, useScripts } from "../api/hooks.js"; +import { + useDeleteScript, + useForkScript, + useInstallPlugin, + useScripts, +} from "../api/hooks.js"; import { ScriptIcon } from "../components/ScriptIcon.jsx"; import { TagBadge } from "../components/TagBadge.jsx"; import { scriptTags } from "../lib/script.js"; +import { useNotifications } from "../notifications.jsx"; export function ScriptsPage() { const [params] = useSearchParams(); @@ -14,7 +20,12 @@ export function ScriptsPage() { const { data: scripts = [], isLoading } = useScripts(); const [confirmDelete, setConfirmDelete] = useState(null); const [query, setQuery] = useState(""); + const [forkFor, setForkFor] = useState(null); + const [forkId, setForkId] = useState(""); const del = useDeleteScript(); + const fork = useForkScript(); + const install = useInstallPlugin(); + const { notify } = useNotifications(); const visible = useMemo(() => { const term = query.trim().toLowerCase(); @@ -23,9 +34,11 @@ export function ScriptsPage() { const name = typeof s === "string" ? s : s.name ?? ""; const description = typeof s === "string" ? "" : s.meta?.description ?? ""; const tags = typeof s === "string" ? [] : scriptTags(s.meta); + const kind = typeof s === "string" ? "" : s.kind ?? ""; return ( name.toLowerCase().includes(term) || description.toLowerCase().includes(term) || + kind.toLowerCase().includes(term) || tags.some((t) => t.toLowerCase().includes(term)) ); }); @@ -50,13 +63,33 @@ export function ScriptsPage() { onChange={(e) => setQuery(e.target.value)} /> + - Add + Add plugin
+ {install.isError ? ( +

{errorMessage(install.error)}

+ ) : null} + {isLoading ? ( ) : scripts.length === 0 ? ( @@ -71,6 +104,8 @@ export function ScriptsPage() { const metaError = typeof s === "string" ? null : s.metaError; const hasIcon = typeof s === "string" ? undefined : s.hasIcon; const tags = typeof s === "string" ? [] : scriptTags(s.meta); + const kind = typeof s === "string" ? "core" : s.kind ?? "core"; + const isCore = kind === "core"; return (
+
+ + {kind} + +

{metaError ? ( {metaError} @@ -114,25 +156,47 @@ export function ScriptsPage() { type="button" className="btn btn-ghost btn-xs" title="Dry run" - onClick={() => navigate(`/scripts/${encodeURIComponent(name)}/dry-run`)} + onClick={() => + navigate(`/scripts/${encodeURIComponent(name)}/dry-run`) + } > - - - - + {isCore ? ( + + ) : ( + + + + )} + {!isCore ? ( + + ) : null}

@@ -141,20 +205,79 @@ export function ScriptsPage() { )} - {confirmDelete ? ( + {forkFor ? (
-

Delete {confirmDelete}?

- {del.isError ? ( -

{errorMessage(del.error)}

+

Fork {forkFor}

+

+ Creates plugin/<id> from this core script. +

+ setForkId(e.target.value)} + /> + {fork.isError ? ( +

{errorMessage(fork.error)}

) : null}
- +
+
+
+ +
+
+ ) : null} + + {confirmDelete ? ( + +
+

Delete {confirmDelete}?

+

This uninstalls the plugin.

+ {del.isError ? ( +

{errorMessage(del.error)}

+ ) : null} +
+ +