From a95e6d7bc86b9ab8de5e6b4a6f78cfb160fce985 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 19 Aug 2026 04:25:21 +0000 Subject: [PATCH 1/2] feat(runner): queue workflow runs with BullMQ (phase A) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- README.md | 9 +- .../20260819040000_workflow_run_queue.js | 24 ++ packages/server/package.json | 2 + packages/server/registry.js | 211 +++++++++------ packages/server/runner.js | 245 +++++++++++------- packages/server/src/api/dashboard.js | 6 +- packages/server/src/api/workflows.js | 9 +- packages/server/store.js | 45 +++- packages/server/workflow-queue.js | 103 ++++++++ packages/web/src/api/hooks.js | 6 +- .../components/workflow/WorkflowTestPanel.jsx | 65 ++++- packages/web/src/lib/format.jsx | 12 +- packages/web/src/pages/HomePage.jsx | 4 +- pnpm-lock.yaml | 176 +++++++++++++ 14 files changed, 716 insertions(+), 201 deletions(-) create mode 100644 packages/server/migrations/20260819040000_workflow_run_queue.js create mode 100644 packages/server/workflow-queue.js diff --git a/README.md b/README.md index 3a850ce..4f6da45 100644 --- a/README.md +++ b/README.md @@ -64,18 +64,25 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / | `JFLOW_JWT_SECRET` | `jflow-dev-secret` (dev only) | **Required in production**. | | `JFLOW_SECRETS_KEY` | `jflow-dev-secrets-key` (dev only) | Master key for named secrets. **Required in production**. Changing it makes existing secrets unreadable. 64 hex chars are used as a raw AES-256 key; any other string is derived with scrypt. | | `JFLOW_DB_PATH` | `packages/server/data/jerapah-flow.db` | SQLite file. | +| `REDIS_URL` | `redis://127.0.0.1:6379` | Redis for BullMQ workflow queue. **Required** — the server will not start if Redis is unreachable. | +| `JFLOW_QUEUE_NAME` | `jerapah-workflows` | BullMQ queue name. | +| `JFLOW_WORKER_CONCURRENCY` | `5` | Max parallel workflow jobs per worker process. | +| `JFLOW_ROLE` | `all` | `all` (API + cron producer + worker), `api` (HTTP/admin/cron enqueue only), or `worker` (consume queue only). | | `JFLOW_LOG_LEVEL` | `debug` | Pino level | | `JFLOW_RETENTION_DAYS` | `30` | Run history prune | | `JFLOW_CORS_ORIGIN` | `http://localhost:5173` | Vite origin in dev | | `PORT` | `9000` | HTTP port | | `NODE_ENV` | — | Set `production` for secure cookies | +Workflow runs are **queued** via BullMQ. HTTP and manual triggers return `202 { runId, status: "queued" }` immediately; poll `GET /api/runs/:id` for progress (`queued` → `running` → `success` \| `failed`). Cron remains an in-process producer that enqueues jobs on each tick. + ## Production ```bash pnpm install pnpm build -JFLOW_JWT_SECRET=... JFLOW_SECRETS_KEY=... NODE_ENV=production pnpm start +# Redis must be reachable at REDIS_URL +JFLOW_JWT_SECRET=... JFLOW_SECRETS_KEY=... REDIS_URL=redis://127.0.0.1:6379 NODE_ENV=production pnpm start ``` The server serves `packages/web/dist` when that folder exists. diff --git a/packages/server/migrations/20260819040000_workflow_run_queue.js b/packages/server/migrations/20260819040000_workflow_run_queue.js new file mode 100644 index 0000000..4695278 --- /dev/null +++ b/packages/server/migrations/20260819040000_workflow_run_queue.js @@ -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"); + }); +} diff --git a/packages/server/package.json b/packages/server/package.json index d3ac22a..a36bca8 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -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", diff --git a/packages/server/registry.js b/packages/server/registry.js index 39a4182..2c9b472 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -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} */ const workflows = new Map(); /** @type {Map} */ @@ -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, }; } diff --git a/packages/server/runner.js b/packages/server/runner.js index 87ca3c3..71f1ba8 100644 --- a/packages/server/runner.js +++ b/packages/server/runner.js @@ -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} */ (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} */ (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"); +} \ No newline at end of file diff --git a/packages/server/src/api/dashboard.js b/packages/server/src/api/dashboard.js index fe3831e..9ea8845 100644 --- a/packages/server/src/api/dashboard.js +++ b/packages/server/src/api/dashboard.js @@ -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, diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index f017f58..18ed3db 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -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 () => { diff --git a/packages/server/store.js b/packages/server/store.js index d8510dd..fc37b9e 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -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) => ({ diff --git a/packages/server/workflow-queue.js b/packages/server/workflow-queue.js new file mode 100644 index 0000000..d969d62 --- /dev/null +++ b/packages/server/workflow-queue.js @@ -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} 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()); +} diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index e4387f2..3d3f12f 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -271,8 +271,10 @@ export function useRun(id) { queryKey: ["runs", id], queryFn: async () => (await api.get(`/runs/${encodeURIComponent(id)}`)).data, enabled: Boolean(id), - refetchInterval: (query) => - query.state.data?.status === "running" ? 2000 : false, + refetchInterval: (query) => { + const status = query.state.data?.status; + return status === "running" || status === "queued" ? 1500 : false; + }, }); } diff --git a/packages/web/src/components/workflow/WorkflowTestPanel.jsx b/packages/web/src/components/workflow/WorkflowTestPanel.jsx index 4578750..4ac649c 100644 --- a/packages/web/src/components/workflow/WorkflowTestPanel.jsx +++ b/packages/web/src/components/workflow/WorkflowTestPanel.jsx @@ -1,7 +1,7 @@ import { useEffect, useRef, useState } from "react"; import { Link } from "react-router-dom"; import { LuMaximize2, LuMinimize2, LuPlay } from "react-icons/lu"; -import { errorMessage } from "../../api/client.js"; +import { api, errorMessage } from "../../api/client.js"; import { useRunWorkflow } from "../../api/hooks.js"; import { CodeEditor } from "../CodeEditor.jsx"; import { StatusBadge } from "../../lib/format.jsx"; @@ -14,11 +14,25 @@ import { import { SchemaTooltip } from "./FieldHelp.jsx"; const SAVE_DEBOUNCE_MS = 300; +const POLL_MS = 1500; function initialDataJson(owner, file, defaultData) { return prettyJson(overlayWorkflowTestData(defaultData, readWorkflowTestData(owner, file))); } +function isTerminalStatus(status) { + return status === "success" || status === "failed"; +} + +async function waitForRun(runId, { signal } = {}) { + while (!signal?.aborted) { + const run = (await api.get(`/runs/${encodeURIComponent(runId)}`)).data; + if (isTerminalStatus(run.status)) return run; + await new Promise((resolve) => setTimeout(resolve, POLL_MS)); + } + throw new Error("polling aborted"); +} + export function WorkflowTestPanel({ owner, file, @@ -34,8 +48,10 @@ export function WorkflowTestPanel({ const [inputTouched, setInputTouched] = useState(() => readWorkflowTestData(owner, file) != null); const [parseError, setParseError] = useState(null); const [last, setLast] = useState(null); + const [polling, setPolling] = useState(false); const [expanded, setExpanded] = useState(false); const saveTimer = useRef(null); + const pollAbort = useRef(null); const seedJson = prettyJson( overlayWorkflowTestData(defaultData, readWorkflowTestData(owner, file)), @@ -44,6 +60,7 @@ export function WorkflowTestPanel({ useEffect(() => { return () => { if (saveTimer.current) clearTimeout(saveTimer.current); + pollAbort.current?.abort(); }; }, []); @@ -79,16 +96,42 @@ export function WorkflowTestPanel({ return; } persist(data); + pollAbort.current?.abort(); + const controller = new AbortController(); + pollAbort.current = controller; + run.mutate( { owner, file, data }, { - onSuccess: (res) => { + onSuccess: async (res) => { + const runId = res.runId ?? null; setLast({ - status: res.status ?? "success", - runId: res.runId ?? null, - result: res.result, + status: res.status ?? "queued", + runId, + result: null, error: null, }); + if (!runId) return; + setPolling(true); + try { + const finished = await waitForRun(runId, { signal: controller.signal }); + setLast({ + status: finished.status, + runId, + result: finished.output, + error: finished.error ?? null, + }); + } catch (err) { + if (controller.signal.aborted) return; + setLast({ + status: "failed", + runId, + result: null, + error: errorMessage(err), + }); + } finally { + if (!controller.signal.aborted) setPolling(false); + } }, onError: (err) => { const runId = err?.response?.data?.runId; @@ -107,11 +150,15 @@ export function WorkflowTestPanel({ ? prettyJson(last.error ? { error: last.error } : last.result) : ""; + const busy = run.isPending || polling; + function close() { if (saveTimer.current) { clearTimeout(saveTimer.current); saveTimer.current = null; } + pollAbort.current?.abort(); + setPolling(false); try { persist(JSON.parse(dataJson || "null")); } catch { @@ -134,8 +181,8 @@ export function WorkflowTestPanel({

Test workflow

- Runs the saved YAML and writes an Event. Manual run works even when - the workflow is disabled. Save before testing unsaved edits. + Enqueues the saved YAML and polls the Event until it finishes. + Manual run works even when the workflow is disabled. Save before testing unsaved edits.