feat(ops): restart individual PM2 children and enrich manage UI

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-08-20 22:31:34 +07:00
co-authored by Cursor
parent 84f29057df
commit aa3c2a25b2
12 changed files with 658 additions and 237 deletions
+61 -1
View File
@@ -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] },
+47 -3
View File
@@ -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),
},
};
}
+5 -1
View File
@@ -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();