Profiles store script + default config per owner. Workflow steps reference them with profile:, overlay keys win, and the UI marks overrides. Profile name is immutable after create so YAML refs stay stable. Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
292 lines
8.1 KiB
JavaScript
292 lines
8.1 KiB
JavaScript
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 { addApiAuthGuard, COOKIE } 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 { queryRunsFromRequest } from "./src/api/run-query.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 profilesPlugin from "./src/api/profiles.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 { purgeExpiredTrash } from "./workflow-trash.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();
|
|
try {
|
|
const purged = await purgeExpiredTrash();
|
|
if (purged > 0) {
|
|
log.info({ purged }, "purged expired workflow trash");
|
|
}
|
|
} catch (err) {
|
|
log.warn({ err }, "workflow trash purge failed");
|
|
}
|
|
}
|
|
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,
|
|
// Cron + HTTP triggers enqueue jobs; only the API process may own them.
|
|
enableTriggers: runApi,
|
|
});
|
|
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) => {
|
|
addApiAuthGuard(api, server);
|
|
await api.register(authPlugin);
|
|
await api.register(usersPlugin);
|
|
await api.register(secretsPlugin);
|
|
await api.register(variablesPlugin);
|
|
await api.register(profilesPlugin);
|
|
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 ?? {});
|
|
return reply.send(await queryRunsFromRequest(q));
|
|
},
|
|
);
|
|
|
|
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,
|
|
};
|
|
}
|