diff --git a/README.md b/README.md index eff3864..17191c3 100644 --- a/README.md +++ b/README.md @@ -13,14 +13,24 @@ 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. + ## Script contract Each script is `async function main(ctx)` and **must** return: @@ -50,13 +60,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 +96,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 +113,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/package.json b/package.json index d001c21..528db8e 100644 --- a/package.json +++ b/package.json @@ -6,10 +6,14 @@ "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/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/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..a6a795c 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -6,7 +6,11 @@ "main": "runner.js", "scripts": { "dev": "node --watch runner.js", + "dev:pm2": "node scripts/dev-pm2.mjs", "start": "node runner.js", + "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)))\"" }, "dependencies": { @@ -31,6 +35,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/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/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/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/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/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..2b80df5 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({ @@ -450,3 +450,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} +