feat(web): add Ops page and proxy control on :8600

Vite listens on 8500 and proxies /ops to the control plane. Admin Ops UI
covers pause/resume, scale, drain/force restart, and restart-needed banner.

Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
Cursor Agent
2026-08-19 10:55:02 +00:00
co-authored by nsrb
parent 756ba8a7ac
commit 4e6c336ae8
21 changed files with 2931 additions and 269 deletions
+142
View File
@@ -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<Array<{
* field: string,
* role: string,
* pid: number,
* hostname: string,
* generation: number,
* ts: number,
* stale: boolean,
* }>>}
*/
export async function readHeartbeats() {
const conn = getSharedConnection();
const all = await conn.hgetall(HEARTBEAT_KEY);
const now = Date.now();
/** @type {Array<any>} */
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<void>} handler
* @returns {Promise<{ stop: () => Promise<void> }>}
*/
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);
}
+164
View File
@@ -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<ControlState>} 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;
}
}
+427
View File
@@ -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);
});
+38
View File
@@ -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",
},
},
],
};
+58
View File
@@ -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<string>} */
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 };
}
+5
View File
@@ -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"
+272
View File
@@ -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<void>}
*/
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<import("pm2").ProcessDescription[]>}
*/
export function listPm2() {
return new Promise((resolve, reject) => {
pm2.list((err, list) => {
if (err) reject(err);
else resolve(list ?? []);
});
});
}
/**
* @param {string} name
* @returns {Promise<import("pm2").ProcessDescription[]>}
*/
export async function listPm2ByName(name) {
const list = await listPm2();
return list.filter((p) => p.name === name);
}
/**
* @param {object} app
* @returns {Promise<void>}
*/
function startPm2App(app) {
return new Promise((resolve, reject) => {
pm2.start(app, (err) => {
if (err) reject(err);
else resolve();
});
});
}
/**
* @param {string} name
* @returns {Promise<void>}
*/
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<void>}
*/
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<void>}
*/
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<void>}
*/
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;
}
+6 -243
View File
@@ -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<string, string | undefined>} */ (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");
}
+11
View File
@@ -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",
});
+21 -7
View File
@@ -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" };
});
};
+292
View File
@@ -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<string, string | undefined>} */ (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,
};
}
+10
View File
@@ -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,
});