feat(runner): queue workflow runs with BullMQ (phase A)

Enqueue HTTP, cron, and manual runs via Redis/BullMQ; workers execute
asynchronously with configurable concurrency. Triggers return 202 and
clients poll GET /api/runs/:id for progress (queued → running → done).

Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
Cursor Agent
2026-08-19 04:25:21 +00:00
co-authored by nsrb
parent 93362208a2
commit a95e6d7bc8
14 changed files with 716 additions and 201 deletions
@@ -0,0 +1,24 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.alterTable("workflow_runs", (t) => {
t.text("job_id");
t.text("queued_at");
});
await knex.schema.raw(
"CREATE INDEX workflow_runs_job_id_idx ON workflow_runs (job_id)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.raw("DROP INDEX IF EXISTS workflow_runs_job_id_idx");
await knex.schema.alterTable("workflow_runs", (t) => {
t.dropColumn("job_id");
t.dropColumn("queued_at");
});
}
+2
View File
@@ -17,7 +17,9 @@
"axios": "^1.19.0",
"bcryptjs": "^3.0.2",
"better-sqlite3": "^13.0.3",
"bullmq": "^6.1.2",
"fastify": "^5.12.0",
"ioredis": "^6.0.0",
"jsonata": "^2.2.2",
"knex": "^3.3.0",
"mustache": "^4.2.0",
+134 -77
View File
@@ -35,6 +35,7 @@ import {
buildFailureAlertData,
resolveFailureTriggerConfig,
} from "./trigger-failure.js";
import { enqueueWorkflowJob } from "./workflow-queue.js";
/**
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
@@ -56,8 +57,10 @@ function hasWorkflowTrigger(workflow) {
/**
* @param {import("fastify").FastifyInstance} server
* @param {{ queue?: import("bullmq").Queue | null }} [opts]
*/
export function createRegistry(server) {
export function createRegistry(server, opts = {}) {
const queue = opts.queue ?? null;
/** @type {Map<string, WorkflowEntry>} */
const workflows = new Map();
/** @type {Map<string, string>} */
@@ -261,26 +264,27 @@ export function createRegistry(server) {
}
}
const result = await runWorkflow(
const result = await enqueueWorkflow(
mapped.key,
{ data: req.body },
{ type: "http", detail: `${method} ${url}` },
);
if (result.status === "failed") {
return reply.code(500).send({
return reply.code(result.runId ? 500 : 404).send({
runId: result.runId,
status: result.status,
error: result.error,
});
}
const defaultBody = {
runId: result.runId,
result: result.result,
status: result.status,
};
if (typeof liveTrigger.response === "string" && liveTrigger.response) {
return sendSuccessPage(reply, liveTrigger.response, defaultBody);
}
return reply.send(defaultBody);
return reply.code(202).send(defaultBody);
}
function registerCronTriggers() {
@@ -308,7 +312,7 @@ export function createRegistry(server) {
schedule,
() => {
log.debug(`cron firing ${key} (${schedule})`);
return runWorkflow(
return enqueueWorkflow(
key,
{ data: workflow.data ?? null },
{ type: "cron", detail: schedule },
@@ -544,14 +548,13 @@ export function createRegistry(server) {
);
}
const destKey = resolveWorkflowTriggerKey(owner, name);
return runWorkflow(
return enqueueWorkflow(
destKey,
{ data },
{ type: "workflow", detail: parentKey },
{
parentRunId,
depth: depth + 1,
detach: true,
},
);
},
@@ -696,32 +699,36 @@ export function createRegistry(server) {
"triggering failure alert workflow",
);
await runWorkflow(
await enqueueWorkflow(
destKey,
{ data: alertData },
{ type: "workflow", detail: `failure:${opts.key}` },
{
parentRunId: opts.runId,
depth: opts.depth + 1,
detach: true,
},
);
}
/**
* Create a queued run and push a BullMQ job. Returns immediately.
*
* @param {string} key
* @param {{ data?: unknown, context?: unknown }} context
* @param {{ type: string, detail?: string | null }} trigger
* @param {{
* parentRunId?: string | null,
* depth?: number,
* detach?: boolean,
* }} [opts]
*/
async function runWorkflow(key, context, trigger, opts = {}) {
async function enqueueWorkflow(key, context, trigger, opts = {}) {
const parentRunId = opts.parentRunId ?? null;
const depth = opts.depth ?? 0;
const detach = opts.detach === true;
if (!queue) {
log.error({ workflow: key }, "workflow queue is not configured");
return { runId: null, status: "failed", error: "workflow queue is not configured" };
}
const entry = workflows.get(key);
if (!entry) {
@@ -734,82 +741,130 @@ export function createRegistry(server) {
log.debug({ workflow: key, trigger }, "skipping disabled workflow");
return { runId: null, status: "failed", error: "workflow disabled" };
}
const input = context.data ?? workflow.data ?? null;
const run = await store.startRun({
owner,
workflow: key,
workflowName: workflow?.name,
trigger,
input: context.data,
input,
parentRunId,
status: "queued",
});
const runLog = log.child({ runId: run.id, owner, workflow: key });
runLog.debug(detach ? "running workflow (detached)" : "running workflow");
const initialCtx = {
data: context.data ?? workflow.data ?? null,
context: normalizeContext(context.context),
};
const execute = async () => {
let ctx = initialCtx;
try {
const compiled = compileWorkflowScripts(workflow.scripts);
if (compiled.dagMode) {
ctx = await runDagSteps(
compiled,
ctx,
run.id,
runLog,
key,
owner,
depth,
);
} else {
ctx = await runLinearSteps(
compiled,
ctx,
run.id,
runLog,
key,
owner,
depth,
);
}
await store.finishRun(run.id, "success", storedEnvelope(ctx));
return { runId: run.id, status: "success", result: storedEnvelope(ctx) };
} catch (err) {
runLog.error({ err }, "workflow failed");
const error = err instanceof Error ? err.message : String(err);
await store.finishRun(run.id, "failed", null, err);
try {
await maybeTriggerFailureWorkflow({
key,
owner,
workflow,
runId: run.id,
trigger,
error,
depth,
});
} catch (alertErr) {
runLog.error({ err: alertErr }, "failed to trigger failure alert workflow");
}
return {
runId: run.id,
status: "failed",
error,
};
}
};
if (detach) {
execute().catch((err) => {
runLog.error({ err }, "detached workflow failed unexpectedly");
try {
const job = await enqueueWorkflowJob(queue, {
runId: run.id,
key,
depth,
});
return { runId: run.id, status: "started" };
await store.setRunJobId(run.id, String(job.id));
runLog.debug({ jobId: job.id }, "workflow queued");
return { runId: run.id, status: "queued", jobId: String(job.id) };
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
runLog.error({ err }, "failed to enqueue workflow");
await store.finishRun(run.id, "failed", null, err);
return { runId: run.id, status: "failed", error: message };
}
}
/**
* Execute a previously queued run (BullMQ worker entrypoint).
*
* @param {{ runId: string, key: string, depth?: number }} jobData
*/
async function executeQueuedRun(jobData) {
const runId = jobData.runId;
const key = jobData.key;
const depth = jobData.depth ?? 0;
const entry = workflows.get(key);
if (!entry) {
await store.finishRun(runId, "failed", null, new Error("workflow not found"));
return { runId, status: "failed", error: "workflow not found" };
}
return execute();
const existing = await store.getRun(runId);
if (!existing) {
return { runId, status: "failed", error: "run not found" };
}
if (existing.status === "success" || existing.status === "failed") {
return { runId, status: existing.status };
}
const marked = await store.markRunRunning(runId);
if (!marked.updated && existing.status !== "running") {
return { runId, status: existing.status };
}
const { owner, workflow } = entry;
const runLog = log.child({ runId, owner, workflow: key });
runLog.debug("running queued workflow");
const trigger = {
type: existing.trigger_type,
detail: existing.trigger_detail,
};
const initialCtx = {
data: existing.input ?? workflow.data ?? null,
context: normalizeContext(null),
};
try {
const compiled = compileWorkflowScripts(workflow.scripts);
let ctx;
if (compiled.dagMode) {
ctx = await runDagSteps(
compiled,
initialCtx,
runId,
runLog,
key,
owner,
depth,
);
} else {
ctx = await runLinearSteps(
compiled,
initialCtx,
runId,
runLog,
key,
owner,
depth,
);
}
await store.finishRun(runId, "success", storedEnvelope(ctx));
return { runId, status: "success", result: storedEnvelope(ctx) };
} catch (err) {
runLog.error({ err }, "workflow failed");
const error = err instanceof Error ? err.message : String(err);
await store.finishRun(runId, "failed", null, err);
try {
await maybeTriggerFailureWorkflow({
key,
owner,
workflow,
runId,
trigger,
error,
depth,
});
} catch (alertErr) {
runLog.error({ err: alertErr }, "failed to trigger failure alert workflow");
}
return { runId, status: "failed", error };
}
}
/**
* @deprecated Prefer enqueueWorkflow; kept as alias for callers.
*/
async function runWorkflow(key, context, trigger, opts = {}) {
return enqueueWorkflow(key, context, trigger, opts);
}
function reregister() {
@@ -842,6 +897,8 @@ export function createRegistry(server) {
registerPruneJob,
reregister,
runWorkflow,
enqueueWorkflow,
executeQueuedRun,
referencedScripts,
};
}
+150 -95
View File
@@ -22,6 +22,12 @@ 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,
getRedisUrl,
} from "./workflow-queue.js";
await migrate();
enableLogPersistence();
@@ -42,6 +48,20 @@ try {
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: getRedisUrl(), role }, "starting jerapah-flow");
const workflowQueue = createWorkflowQueue();
try {
await workflowQueue.waitUntilReady();
} catch (err) {
log.error({ err, redis: getRedisUrl() }, "failed to connect to Redis");
process.exit(1);
}
const server = fastify({ loggerInstance: log });
await server.register(cookie);
@@ -71,98 +91,129 @@ server.decorate("requireAdmin", async function requireAdmin(req, reply) {
}
});
const registry = createRegistry(server);
const registry = createRegistry(server, { queue: workflowQueue });
registry.registerWorkflows();
registry.registerHttpTriggers();
registry.registerCronTriggers();
registry.registerPruneJob();
if (runApi) {
registry.registerHttpTriggers();
registry.registerCronTriggers();
registry.registerPruneJob();
}
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" });
/** @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");
}
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" });
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 reply.sendFile("index.html");
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) {
@@ -176,15 +227,19 @@ process.on("SIGTERM", shutdown);
const port = Number(process.env.PORT ?? 9000);
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);
});
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");
}
+3 -3
View File
@@ -63,8 +63,8 @@ export default function dashboardPluginFactory(registry) {
}
}
const [running, failed, recent] = await Promise.all([
store.listRuns({ status: "running", limit: 10 }),
const [active, failed, recent] = await Promise.all([
store.listRuns({ status: ["queued", "running"], limit: 10 }),
store.listRuns({ status: "failed", limit: 20 }),
store.listRuns({ limit: 10 }),
]);
@@ -74,7 +74,7 @@ export default function dashboardPluginFactory(registry) {
scriptCount: fsStore.listScriptFiles().length,
enabledCount,
brokenCount,
running,
running: active,
needsAttention: {
failed,
brokenWorkflows,
+5 -4
View File
@@ -404,7 +404,7 @@ export default function workflowsPluginFactory(registry) {
});
}
const body = /** @type {{ data?: unknown }} */ (req.body ?? {});
const result = await registry.runWorkflow(
const result = await registry.enqueueWorkflow(
key,
{ data: body.data ?? null },
{ type: "manual", detail: "ui" },
@@ -412,14 +412,15 @@ export default function workflowsPluginFactory(registry) {
if (result.status === "failed") {
return reply.code(result.runId ? 500 : 404).send({
runId: result.runId,
status: result.status,
error: result.error,
});
}
return {
return reply.code(202).send({
runId: result.runId,
status: result.status,
result: result.result,
};
jobId: result.jobId ?? null,
});
});
fastify.post("/workflows/reregister", async () => {
+39 -6
View File
@@ -49,6 +49,7 @@ function nowIso() {
* trigger: { type: string, detail?: string | null },
* input?: unknown,
* parentRunId?: string | null,
* status?: "queued" | "running",
* }} opts
*/
export async function startRun({
@@ -58,9 +59,11 @@ export async function startRun({
trigger,
input = null,
parentRunId = null,
status = "queued",
}) {
const id = randomUUID();
const started_at = nowIso();
const now = nowIso();
const isQueued = status === "queued";
await db("workflow_runs").insert({
id,
owner,
@@ -68,12 +71,36 @@ export async function startRun({
workflow_name: workflowName ?? null,
trigger_type: trigger.type,
trigger_detail: trigger.detail ?? null,
status: "running",
started_at,
status,
started_at: now,
queued_at: isQueued ? now : null,
input: serialize(input),
parent_run_id: parentRunId ?? null,
});
return { id, started_at };
return { id, started_at: now, queued_at: isQueued ? now : null };
}
/**
* @param {string} id
* @param {string} jobId
*/
export async function setRunJobId(id, jobId) {
await db("workflow_runs").where({ id }).update({ job_id: jobId });
}
/**
* @param {string} id
*/
export async function markRunRunning(id) {
const started_at = nowIso();
const updated = await db("workflow_runs")
.where({ id })
.whereIn("status", ["queued", "running"])
.update({
status: "running",
started_at,
});
return { updated: Number(updated) > 0, started_at };
}
/**
@@ -181,7 +208,7 @@ export async function insertLogs(rows) {
* @param {{
* owner?: string,
* workflow?: string,
* status?: string,
* status?: string | string[],
* limit?: number,
* before?: string,
* }} [filters]
@@ -203,7 +230,13 @@ export async function listRuns(filters = {}) {
q = q.where("workflow", key);
}
}
if (filters.status) q = q.where("status", filters.status);
if (filters.status) {
if (Array.isArray(filters.status)) {
q = q.whereIn("status", filters.status);
} else {
q = q.where("status", filters.status);
}
}
if (filters.before) q = q.where("started_at", "<", filters.before);
const rows = await q.limit(limit);
return rows.map((row) => ({
+103
View File
@@ -0,0 +1,103 @@
import { Queue, Worker } from "bullmq";
import IORedis from "ioredis";
import { log } from "./logger.js";
const DEFAULT_REDIS_URL = "redis://127.0.0.1:6379";
const DEFAULT_QUEUE_NAME = "jerapah-workflows";
const DEFAULT_CONCURRENCY = 5;
/** @type {IORedis | null} */
let sharedConnection = null;
export function getRedisUrl() {
return process.env.REDIS_URL || DEFAULT_REDIS_URL;
}
export function getQueueName() {
return process.env.JFLOW_QUEUE_NAME || DEFAULT_QUEUE_NAME;
}
export function getWorkerConcurrency() {
const raw = Number(process.env.JFLOW_WORKER_CONCURRENCY ?? DEFAULT_CONCURRENCY);
if (!Number.isFinite(raw) || raw < 1) return DEFAULT_CONCURRENCY;
return Math.floor(raw);
}
/**
* BullMQ requires maxRetriesPerRequest: null for blocking commands.
* @returns {IORedis}
*/
export function getSharedConnection() {
if (sharedConnection) return sharedConnection;
sharedConnection = new IORedis(getRedisUrl(), {
maxRetriesPerRequest: null,
enableReadyCheck: true,
});
sharedConnection.on("error", (err) => {
log.error({ err }, "redis connection error");
});
return sharedConnection;
}
/**
* @returns {Queue}
*/
export function createWorkflowQueue() {
return new Queue(getQueueName(), {
connection: getSharedConnection(),
defaultJobOptions: {
removeOnComplete: { count: 1000 },
removeOnFail: { count: 5000 },
attempts: 1,
},
});
}
/**
* @param {(job: import("bullmq").Job) => Promise<unknown>} processor
* @returns {Worker}
*/
export function createWorkflowWorker(processor) {
const concurrency = getWorkerConcurrency();
const worker = new Worker(getQueueName(), processor, {
connection: getSharedConnection(),
concurrency,
});
worker.on("error", (err) => {
log.error({ err }, "workflow worker error");
});
log.info({ concurrency, queue: getQueueName() }, "workflow worker started");
return worker;
}
/**
* @param {Queue} queue
* @param {{
* runId: string,
* key: string,
* depth?: number,
* }} data
*/
export async function enqueueWorkflowJob(queue, data) {
const job = await queue.add(
"run",
{
runId: data.runId,
key: data.key,
depth: data.depth ?? 0,
},
{
jobId: data.runId,
},
);
return job;
}
/**
* @param {IORedis | null} [connection]
*/
export async function closeRedis(connection = sharedConnection) {
if (!connection) return;
if (connection === sharedConnection) sharedConnection = null;
await connection.quit().catch(() => connection.disconnect());
}