diff --git a/ecosystem.config.cjs b/ecosystem.config.cjs new file mode 100644 index 0000000..f8ef770 --- /dev/null +++ b/ecosystem.config.cjs @@ -0,0 +1,43 @@ +const fs = require("fs"); +const path = require("path"); + +function loadEnv(file) { + const out = {}; + if (!fs.existsSync(file)) return out; + for (const line of fs.readFileSync(file, "utf8").split("\n")) { + const t = line.trim(); + if (!t || t.startsWith("#")) continue; + const i = t.indexOf("="); + if (i === -1) continue; + const key = t.slice(0, i).trim(); + let val = t.slice(i + 1).trim(); + if ( + (val.startsWith('"') && val.endsWith('"')) || + (val.startsWith("'") && val.endsWith("'")) + ) { + val = val.slice(1, -1); + } + out[key] = val; + } + return out; +} + +const root = __dirname; + +module.exports = { + apps: [ + { + name: "jerapah-flow", + cwd: root, + script: "packages/server/runner.js", + interpreter: "node", + instances: 1, + autorestart: true, + max_restarts: 20, + env: { + NODE_ENV: "production", + ...loadEnv(path.join(root, ".env")), + }, + }, + ], +}; \ No newline at end of file diff --git a/packages/server/control.js b/packages/server/control.js index ac1fd94..f79c885 100644 --- a/packages/server/control.js +++ b/packages/server/control.js @@ -35,6 +35,7 @@ import { PM2_HTTP_NAME, PM2_WORKER_NAME, recreateChildren, + restartPm2Process, stopPm2App, } from "./pm2-bridge.js"; @@ -149,6 +150,43 @@ async function waitUntilIdle(timeoutMs, lockToken) { return { ok: active === 0, active }; } +/** @param {unknown} raw */ +function parsePmId(raw) { + if (raw == null || raw === "") return null; + const id = Math.floor(Number(raw)); + if (!Number.isFinite(id) || id < 0) return Number.NaN; + return id; +} + +async function handleProcessRestart(pmId, reply) { + const lock = tryAcquireOpsLock(`process-restart:${process.pid}`); + if (!lock.ok) { + return reply.code(409).send({ error: lock.error, holder: lock.holder }); + } + + try { + const restarted = await restartPm2Process(pmId); + return reply.send({ ok: true, ...restarted, status: await buildStatus() }); + } catch (err) { + const code = err && typeof err === "object" ? err.code : undefined; + if (code === "BAD_REQUEST") { + return reply.code(400).send({ error: "pmId is required" }); + } + if (code === "NOT_FOUND") { + return reply.code(404).send({ error: "process not found" }); + } + if (code === "FORBIDDEN") { + return reply.code(403).send({ error: "process is not a JerapahFlow child" }); + } + log.error({ err, pmId }, "process restart failed"); + return reply.code(500).send({ + error: err instanceof Error ? err.message : String(err), + }); + } finally { + releaseOpsLock(lock.token); + } +} + async function buildStatus() { const state = readControlState(); const children = await describeChildren(); @@ -258,9 +296,18 @@ server.post( "/ops/restart", { onRequest: [server.authenticate, server.requireAdmin] }, async (req, reply) => { - const body = /** @type {{ force?: boolean, timeoutMs?: number }} */ ( + const body = /** @type {{ force?: boolean, timeoutMs?: number, pmId?: number }} */ ( req.body ?? {} ); + + const pmId = parsePmId(body.pmId); + if (pmId != null) { + if (Number.isNaN(pmId)) { + return reply.code(400).send({ error: "pmId is required" }); + } + return handleProcessRestart(pmId, reply); + } + const force = Boolean(body.force); const timeoutMs = Math.min( Math.max(Number(body.timeoutMs) || DRAIN_DEFAULT_MS, 1_000), @@ -342,6 +389,19 @@ server.post( }, ); +server.post( + "/ops/process/restart", + { onRequest: [server.authenticate, server.requireAdmin] }, + async (req, reply) => { + const body = /** @type {{ pmId?: number }} */ (req.body ?? {}); + const pmId = parsePmId(body.pmId); + if (pmId == null || Number.isNaN(pmId)) { + return reply.code(400).send({ error: "pmId is required" }); + } + return handleProcessRestart(pmId, reply); + }, +); + server.post( "/ops/scale", { onRequest: [server.authenticate, server.requireAdmin] }, diff --git a/packages/server/pm2-bridge.js b/packages/server/pm2-bridge.js index 5d72a28..9bbec47 100644 --- a/packages/server/pm2-bridge.js +++ b/packages/server/pm2-bridge.js @@ -130,6 +130,40 @@ export function restartPm2App(name) { }); } +/** + * Restart a single PM2 process by id. Only jflow-http / jflow-worker. + * @param {number} pmId + * @returns {Promise<{ name: string, pmId: number }>} + */ +export async function restartPm2Process(pmId) { + const id = Math.floor(Number(pmId)); + if (!Number.isFinite(id) || id < 0) { + const err = new Error("invalid pmId"); + err.code = "BAD_REQUEST"; + throw err; + } + + const list = await listPm2(); + const proc = list.find((p) => Number(p.pm_id) === id); + if (!proc) { + const err = new Error("process not found"); + err.code = "NOT_FOUND"; + throw err; + } + if (proc.name !== PM2_HTTP_NAME && proc.name !== PM2_WORKER_NAME) { + const err = new Error("process is not a JerapahFlow child"); + err.code = "FORBIDDEN"; + throw err; + } + await new Promise((resolve, reject) => { + pm2.restart(id, (err) => { + if (err) reject(err); + else resolve(); + }); + }); + return { name: proc.name, pmId: id }; +} + /** * Shared env for child processes. * @param {{ generation: number }} opts @@ -251,19 +285,29 @@ export async function describeChildren() { const mapOne = (p) => ({ name: p.name, - pmId: p.pm_id, + pmId: Number(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, + memory: Number(p.monit?.memory) || 0, + cpu: Number(p.monit?.cpu) || 0, }); + const httpMapped = http.map(mapOne); + const workersMapped = workers.map(mapOne); + const all = [...httpMapped, ...workersMapped]; + return { - http: http.map(mapOne), - workers: workers.map(mapOne), + http: httpMapped, + workers: workersMapped, httpOnline: http.some((p) => p.pm2_env?.status === "online"), workerOnlineCount: workers.filter((p) => p.pm2_env?.status === "online").length, + totals: { + memory: all.reduce((sum, p) => sum + p.memory, 0), + cpu: all.reduce((sum, p) => sum + p.cpu, 0), + }, }; } diff --git a/packages/server/start-app.js b/packages/server/start-app.js index 5b6a65e..39f0299 100644 --- a/packages/server/start-app.js +++ b/packages/server/start-app.js @@ -126,7 +126,11 @@ export async function startApp(opts = {}) { } }); - const registry = createRegistry(server, { queue: workflowQueue }); + const registry = createRegistry(server, { + queue: workflowQueue, + // Cron + HTTP triggers enqueue jobs; only the API process may own them. + enableTriggers: runApi, + }); registry.registerWorkflows(); if (runApi) { registry.registerHttpTriggers(); diff --git a/packages/web/src/components/HeartbeatsAndPm2Card.jsx b/packages/web/src/components/HeartbeatsAndPm2Card.jsx new file mode 100644 index 0000000..4e3c220 --- /dev/null +++ b/packages/web/src/components/HeartbeatsAndPm2Card.jsx @@ -0,0 +1,165 @@ +import { useState } from "react"; +import { errorMessage } from "../api/client.js"; +import { HeartbeatsToolbar } from "./HeartbeatsToolbar.jsx"; +import { useNotifications } from "../notifications.jsx"; +import { formatBytes, formatCpu, formatDuration } from "../lib/format.js"; + +const MAX_WORKERS = 32; + +export function HeartbeatsAndPm2Card({ + data, + children, + restart, + processRestart, + workerCount, + desiredWorkers, + onDecreaseWorkers, + onIncreaseWorkers, + onScale, + scalePending = false, + onReload, + reloadPending = false, +}) { + const { notify } = useNotifications(); + const [forceRestart, setForceRestart] = useState(false); + const [pendingPmId, setPendingPmId] = useState(null); + + async function onRestartAll() { + const label = forceRestart ? "Force restart" : "Restart"; + if ( + forceRestart && + !window.confirm("Force restart will interrupt active jobs. Continue?") + ) { + return; + } + + try { + await restart.mutateAsync({ force: forceRestart }); + notify.success(`${label} ok`); + } catch (e) { + notify.error(errorMessage(e, `${label} failed`)); + } + } + + async function onProcessRestart(pmId) { + const id = Number(pmId); + setPendingPmId(id); + try { + await processRestart.mutateAsync({ pmId: id }); + notify.success("Process restart ok"); + } catch (e) { + notify.error(errorMessage(e, "Process restart failed")); + } finally { + setPendingPmId(null); + } + } + + const heartbeats = data?.heartbeats ?? []; + const httpChildren = children?.http ?? []; + const workerChildren = children?.workers ?? []; + const pm2Children = [...httpChildren, ...workerChildren]; + const heartbeatByPid = new Map(); + for (const h of heartbeats) { + if (!heartbeatByPid.has(h.pid)) heartbeatByPid.set(h.pid, h); + } + + return ( +
+
+

Heartbeats

+ + + +
+ + + + + + + + + + + + + + + + + {pm2Children.map((p) => { + const pid = p.pid ?? null; + const hb = pid != null ? heartbeatByPid.get(pid) : undefined; + const role = + hb?.role ?? + (p.name === "jflow-http" + ? "http" + : p.name === "jflow-worker" + ? "worker" + : p.name ?? "unknown"); + const generation = hb?.generation ?? p.generation ?? "—"; + const uptimeMs = + p.uptime != null + ? Math.max(0, Date.now() - Number(p.uptime)) + : hb + ? Math.max(0, Date.now() - hb.ts) + : null; + + return ( + + + + + + + + + + + + + + ); + })} + {!pm2Children.length ? ( + + + + ) : null} + +
RoleNamePM IDStatusPIDGenerationUptimeMemoryCPURestarts +
{role}{p.name}{p.pmId}{p.status}{pid ?? "—"}{generation} + {formatDuration(uptimeMs)} + {formatBytes(p.memory)}{formatCpu(p.cpu)}{p.restarts ?? 0} + +
+ No PM2 children +
+
+
+
+ ); +} diff --git a/packages/web/src/components/HeartbeatsToolbar.jsx b/packages/web/src/components/HeartbeatsToolbar.jsx new file mode 100644 index 0000000..81292f1 --- /dev/null +++ b/packages/web/src/components/HeartbeatsToolbar.jsx @@ -0,0 +1,87 @@ +import { LuMinus, LuPlus } from "react-icons/lu"; + +export function HeartbeatsToolbar({ + forceRestart, + onForceRestartChange, + onRestartAll, + restartPending = false, + workerCount, + desiredWorkers, + onDecreaseWorkers, + onIncreaseWorkers, + onScale, + scalePending = false, + onReload, + reloadPending = false, + minWorkers = 0, + maxWorkers = 32, +}) { + return ( +
+
+ + + +
+ + + + {workerCount} worker{workerCount === 1 ? "" : "s"} + + + + +
+ + +
+ +
+
+ ); +} diff --git a/packages/web/src/components/HttpProcessCard.jsx b/packages/web/src/components/HttpProcessCard.jsx new file mode 100644 index 0000000..fa1e592 --- /dev/null +++ b/packages/web/src/components/HttpProcessCard.jsx @@ -0,0 +1,46 @@ +import { LuCircleCheck, LuCircleX } from "react-icons/lu"; + +export function HttpProcessCard({ + httpOnline, + onStart, + onStop, + startPending = false, + stopPending = false, +}) { + return ( +
+
+

HTTP process

+
+ {httpOnline ? ( + + ) : ( + + )} + {httpOnline ? "Running" : "Stopped"} +
+
+ {httpOnline ? ( + + ) : ( + + )} +
+
+
+ ); +} diff --git a/packages/web/src/components/ProcessResourcesCard.jsx b/packages/web/src/components/ProcessResourcesCard.jsx new file mode 100644 index 0000000..6296c37 --- /dev/null +++ b/packages/web/src/components/ProcessResourcesCard.jsx @@ -0,0 +1,17 @@ +import { formatBytes, formatCpu } from "../lib/format.js"; + +export function ProcessResourcesCard({ children }) { + const totals = children?.totals ?? { memory: 0, cpu: 0 }; + + return ( +
+
+

Resources

+
+

{formatBytes(totals.memory)} memory total

+

{formatCpu(totals.cpu)} CPU total

+
+
+
+ ); +} diff --git a/packages/web/src/components/QueueCard.jsx b/packages/web/src/components/QueueCard.jsx new file mode 100644 index 0000000..32244a5 --- /dev/null +++ b/packages/web/src/components/QueueCard.jsx @@ -0,0 +1,61 @@ +import { LuCircleCheck, LuCirclePause } from "react-icons/lu"; + +export function QueueCard({ + paused = false, + workersOnline = 0, + queuedJobs = 0, + activeJobs = 0, + onPause, + onResume, + pausePending = false, + resumePending = false, +}) { + return ( +
+
+

Queue

+
+ {paused ? ( + + ) : ( + + )} + + {paused + ? "Paused" + : `Running (${workersOnline} worker${workersOnline === 1 ? "" : "s"})`} + +
+
+

+ {activeJobs} running job{activeJobs === 1 ? "" : "s"} +

+

+ {queuedJobs} queued job{queuedJobs === 1 ? "" : "s"} +

+
+
+ {paused ? ( + + ) : ( + + )} +
+
+
+ ); +} diff --git a/packages/web/src/index.css b/packages/web/src/index.css index de45d5a..0070d50 100644 --- a/packages/web/src/index.css +++ b/packages/web/src/index.css @@ -3,12 +3,28 @@ themes: light --default, dark --prefersdark; } +@property --value { + syntax: ""; + inherits: true; + initial-value: 0; +} + html, body, #root { min-height: 100%; } +.radial-progress { + transition: + --value 0.7s ease-out, + --radialprogress 0.7s ease-out; +} + +.radial-progress::after { + transition: transform 0.7s ease-out; +} + .mermaid-clickable [id^="s"] { cursor: pointer; } diff --git a/packages/web/src/lib/format.js b/packages/web/src/lib/format.js new file mode 100644 index 0000000..984cd89 --- /dev/null +++ b/packages/web/src/lib/format.js @@ -0,0 +1,31 @@ +export function formatBytes(bytes) { + const n = Number(bytes); + if (!Number.isFinite(n)) return "—"; + if (n === 0) return "0 KB"; + if (n < 1024 * 1024) return `${Math.round(n / 1024)} KB`; + return `${(n / (1024 * 1024)).toFixed(1)} MB`; +} + +export function formatCpu(cpu) { + const n = Number(cpu); + if (!Number.isFinite(n)) return "—"; + return `${n.toFixed(1)}%`; +} + +/** Compact duration from milliseconds (e.g. `850ms`, `12s`, `5m 3s`, `2h 15m`, `3d 4h`). */ +export function formatDuration(ms) { + const n = Number(ms); + if (!Number.isFinite(n) || n < 0) return "—"; + if (n < 1000) return `${Math.round(n)}ms`; + const sec = Math.floor(n / 1000); + if (sec < 60) return `${sec}s`; + const min = Math.floor(sec / 60); + const s = sec % 60; + if (min < 60) return s ? `${min}m ${s}s` : `${min}m`; + const hr = Math.floor(min / 60); + const m = min % 60; + if (hr < 48) return m ? `${hr}h ${m}m` : `${hr}h`; + const days = Math.floor(hr / 24); + const h = hr % 24; + return h ? `${days}d ${h}h` : `${days}d`; +} diff --git a/packages/web/src/pages/OpsPage.jsx b/packages/web/src/pages/OpsPage.jsx index 72ca3a6..477921c 100644 --- a/packages/web/src/pages/OpsPage.jsx +++ b/packages/web/src/pages/OpsPage.jsx @@ -1,9 +1,9 @@ -import { useState } from "react"; +import { useEffect, useState } from "react"; import { - useOpsBumpGeneration, useOpsHttpStart, useOpsHttpStop, useOpsPause, + useOpsProcessRestart, useOpsReload, useOpsRestart, useOpsResume, @@ -11,11 +11,20 @@ import { useOpsStatus, } from "../api/hooks.js"; import { errorMessage } from "../api/client.js"; +import { HttpProcessCard } from "../components/HttpProcessCard.jsx"; +import { QueueCard } from "../components/QueueCard.jsx"; +import { ProcessResourcesCard } from "../components/ProcessResourcesCard.jsx"; +import { HeartbeatsAndPm2Card } from "../components/HeartbeatsAndPm2Card.jsx"; +import { useNotifications } from "../notifications.jsx"; + +const MAX_WORKERS = 32; function Stat({ label, value }) { return (
- {label} + + {label} + {value}
); @@ -30,10 +39,10 @@ export function OpsPage() { const scale = useOpsScale(); const httpStart = useOpsHttpStart(); const httpStop = useOpsHttpStop(); - const bump = useOpsBumpGeneration(); - const [workerCount, setWorkerCount] = useState(""); - const [msg, setMsg] = useState(null); - const [err, setErr] = useState(null); + const processRestart = useOpsProcessRestart(); + const { notify } = useNotifications(); + const [workerCount, setWorkerCount] = useState(1); + const [workerInit, setWorkerInit] = useState(false); const unavailable = status.isError && @@ -41,15 +50,34 @@ export function OpsPage() { status.error?.response?.status === 404 || status.error?.message?.includes("Network")); + const data = status.data; + const desired = data?.desired; + const queue = data?.queue; + const children = data?.children; + const httpOnline = Boolean(children?.httpOnline); + + useEffect(() => { + if (!workerInit && desired?.workers != null) { + setWorkerCount(desired.workers); + setWorkerInit(true); + } + }, [desired?.workers, workerInit]); + function run(label, mutateAsync, args) { - setMsg(null); - setErr(null); - mutateAsync(args) - .then((data) => { - setMsg(`${label} ok`); - return data; + const call = args === undefined ? mutateAsync() : mutateAsync(args); + call + .then((result) => { + notify.success(`${label} ok`); + return result; }) - .catch((e) => setErr(errorMessage(e, `${label} failed`))); + .catch((e) => notify.error(errorMessage(e, `${label} failed`))); + } + + function scaleWorkers() { + run("Scale", scale.mutateAsync, { + workers: workerCount, + force: false, + }); } if (unavailable) { @@ -68,11 +96,6 @@ export function OpsPage() { ); } - const data = status.data; - const desired = data?.desired; - const queue = data?.queue; - const children = data?.children; - return (
@@ -98,226 +121,50 @@ export function OpsPage() {
) : null} - {msg ? ( -
- {msg} -
- ) : null} - {err ? ( -
- {err} -
- ) : null} - {status.isLoading && !data ? ( ) : ( <> -
- - - - + run("HTTP start", httpStart.mutateAsync)} + onStop={() => run("HTTP stop", httpStop.mutateAsync)} /> - run("Pause", pause.mutateAsync)} + onResume={() => run("Resume", resume.mutateAsync)} /> - - - -
+ +
-
-

Queue

-
- - - -
-
- -
-

HTTP process

-
- - -
-
- -
-

Workers

-
- setWorkerCount(e.target.value)} - /> - - -
-
- -
-

Restart HTTP + workers

-

- Pauses the queue, waits until active jobs are 0, migrates DB, then - recreates PM2 processes. Force skips the wait (may interrupt runs). -

-
- - - -
-
- -
-

Heartbeats

-
- - - - - - - - - - - {(data?.heartbeats ?? []).map((h) => ( - - - - - - - ))} - {!data?.heartbeats?.length ? ( - - - - ) : null} - -
RolePIDGenerationAge (ms)
{h.role}{h.pid}{h.generation} - {Math.max(0, Date.now() - h.ts)} -
- No live heartbeats yet -
-
-
- -
-

PM2 children

-
-              {JSON.stringify(children ?? {}, null, 2)}
-            
-
+ setWorkerCount((n) => Math.max(0, n - 1))} + onIncreaseWorkers={() => + setWorkerCount((n) => Math.min(MAX_WORKERS, n + 1)) + } + onScale={scaleWorkers} + scalePending={scale.isPending} + onReload={() => run("Reload", reload.mutateAsync)} + reloadPending={reload.isPending} + /> )}