diff --git a/.cursor/rules/background-work.mdc b/.cursor/rules/background-work.mdc index f3d3b7c..4a867ad 100644 --- a/.cursor/rules/background-work.mdc +++ b/.cursor/rules/background-work.mdc @@ -1,5 +1,5 @@ --- -description: Background work on the backend, cron jobs in the scheduler vs tasks in the Postgres queue, transactional enqueue, retries, idempotent handlers, PROCESS_ROLE for running a separate worker +description: Background work on the backend, cron jobs in the scheduler vs tasks in the Postgres queue, transactional enqueue, retries, idempotent handlers globs: apps/backend/src/**/* alwaysApply: false --- @@ -64,14 +64,6 @@ The task commits or rolls back with the row. A row can never sit in `queued` wit Workers are off under `NODE_ENV=test`, so nothing runs behind a test's back. `setupIntegrationTest()` starts the queue, which means code under test can enqueue. To test a handler, call the queue's `run` directly with a task object. To test the queue itself, see `src/libs/queue/queue.spec.ts`. -## Running a separate worker +## Where tasks run -`PROCESS_ROLE` picks what a process does. Same image, same start command. - -| Role | API routes | Cron jobs | Task workers | -| --- | --- | --- | --- | -| `all` (default) | yes | yes | yes | -| `web` | yes | no | no | -| `worker` | `/health` only | yes | yes | - -Start with `all`. When background work starts slowing requests down, deploy the backend a second time with `PROCESS_ROLE=worker` (same database, bucket, and provider keys, no public domain, no migration command) and set the public one to `web`. The two never talk to each other. Postgres is the hand-off, so you can run as many workers as you need. +Inside the backend process, next to the API and the scheduler. `main.ts` starts it with `startQueues({ queues, postgres })`, and it runs on that connection, so it opens no pool of its own. The queue is free until used: with nothing in `src/queues/index.ts` it opens no connections and creates no schema. There is no separate worker to deploy. Every instance works tasks, each up to the queue's `concurrency`, so adding a backend replica adds capacity. Keep `concurrency` low for anything memory-heavy: it shares the process with your requests. diff --git a/README.md b/README.md index 1a2be2e..a1cc89b 100644 --- a/README.md +++ b/README.md @@ -76,7 +76,7 @@ Features arrive when you ask for them. The CLI writes whole features into your r - 📡 **HyperFetch SDK**: typed HTTP and WebSocket client generated from your Hono routes - 🗄️ **Postgres + Prisma**: migrations, Kysely queries, pgvector ready - ⏰ **Cron scheduler**: in-process jobs with exactly-once runs and catch-up after downtime -- 📬 **Task queue**: Postgres-backed (pg-boss), retries with backoff, runs in-process or as a separate worker with one env var +- 📬 **Task queue**: Postgres-backed (pg-boss), retries with backoff, runs inside the API process - 🌍 **i18n ready**: Paraglide wired into every app, backend responses included - 🎨 **UI ready**: Tailwind v4 + shadcn/ui on web, NativeWind on mobile. CSR or SSR with a script switch - ⚡ **Rust-powered DX**: OXC lint/format, React Compiler through oxc, Vite 8 HMR in milliseconds diff --git a/apps/backend/.env.example b/apps/backend/.env.example index 9763eb1..e6366d0 100644 --- a/apps/backend/.env.example +++ b/apps/backend/.env.example @@ -26,14 +26,6 @@ RUSTFS_ENDPOINT=http://localhost:9000 RUSTFS_ACCESS_KEY=rustfsadmin RUSTFS_SECRET_KEY=rustfsadmin -# ----------------- -# Process role -# ----------------- - -# all (default): API, cron, and task workers in one process. -# To move background work off the API, run this image twice: web on one, worker on the other. -# PROCESS_ROLE=all - # ----------------- # External Services # ----------------- diff --git a/apps/backend/src/config/env/env.config.ts b/apps/backend/src/config/env/env.config.ts index e120cd8..8516916 100644 --- a/apps/backend/src/config/env/env.config.ts +++ b/apps/backend/src/config/env/env.config.ts @@ -21,8 +21,6 @@ export const envSchema = z.object({ RESEND_API_KEY: z.string(), // # Sentry SENTRY_DSN: z.string(), - // # Process role: "all" runs the API, cron, and task workers together; split them with "web" + "worker" - PROCESS_ROLE: z.enum(["all", "web", "worker"]).optional(), }); export type Env = z.infer; diff --git a/apps/backend/src/libs/queue/process-role.spec.ts b/apps/backend/src/libs/queue/process-role.spec.ts deleted file mode 100644 index 13c9406..0000000 --- a/apps/backend/src/libs/queue/process-role.spec.ts +++ /dev/null @@ -1,34 +0,0 @@ -import { getProcessRole, runsWorkers, servesApi } from "./process-role"; - -describe("process role", () => { - const original = process.env.PROCESS_ROLE; - - afterEach(() => { - if (original === undefined) delete process.env.PROCESS_ROLE; - else process.env.PROCESS_ROLE = original; - }); - - it("defaults to running everything in one process", () => { - delete process.env.PROCESS_ROLE; - expect(getProcessRole()).toBe("all"); - expect(servesApi()).toBe(true); - expect(runsWorkers()).toBe(true); - }); - - it("web serves the API and leaves background work to someone else", () => { - process.env.PROCESS_ROLE = "web"; - expect(servesApi()).toBe(true); - expect(runsWorkers()).toBe(false); - }); - - it("worker runs background work and serves no API", () => { - process.env.PROCESS_ROLE = "worker"; - expect(servesApi()).toBe(false); - expect(runsWorkers()).toBe(true); - }); - - it("fails loudly on a typo instead of silently running as something else", () => { - Object.assign(process.env, { PROCESS_ROLE: "workers" }); - expect(() => getProcessRole()).toThrow(/PROCESS_ROLE must be one of/); - }); -}); diff --git a/apps/backend/src/libs/queue/process-role.ts b/apps/backend/src/libs/queue/process-role.ts deleted file mode 100644 index 99b9c82..0000000 --- a/apps/backend/src/libs/queue/process-role.ts +++ /dev/null @@ -1,30 +0,0 @@ -/** - * What this process does, picked with PROCESS_ROLE. One image, one entrypoint: - * - * - `all` (default): API, scheduler, and task workers in one process. Right until background work - * starts competing with requests. - * - `web`: API only. Still enqueues tasks, never runs them, never ticks cron. - * - `worker`: scheduler and task workers. Serves /health and nothing else, so platform health - * checks keep passing without a second Dockerfile or start command. - * - * To scale out, run the same image twice: PROCESS_ROLE=web on the public service, - * PROCESS_ROLE=worker on a copy with the same database, bucket, and provider keys, no public - * domain, and no migration command. The two never talk to each other; Postgres is the hand-off. - */ -export const PROCESS_ROLES = ["all", "web", "worker"] as const; - -export type ProcessRole = (typeof PROCESS_ROLES)[number]; - -export const getProcessRole = (): ProcessRole => { - const role = process.env.PROCESS_ROLE ?? "all"; - if (!PROCESS_ROLES.includes(role as ProcessRole)) { - throw new Error(`PROCESS_ROLE must be one of ${PROCESS_ROLES.join(", ")}, got "${role}"`); - } - return role as ProcessRole; -}; - -/** True when this process serves the API routes and sockets. */ -export const servesApi = (): boolean => getProcessRole() !== "worker"; - -/** True when this process runs cron jobs and queued tasks. */ -export const runsWorkers = (): boolean => getProcessRole() !== "web"; diff --git a/apps/backend/src/libs/queue/queue.spec.ts b/apps/backend/src/libs/queue/queue.spec.ts index dd870fd..d8717b7 100644 --- a/apps/backend/src/libs/queue/queue.spec.ts +++ b/apps/backend/src/libs/queue/queue.spec.ts @@ -84,7 +84,9 @@ describe("task queue", () => { // is a no-op while an instance is up, so swap it for this spec's queues in a throwaway schema. const start = async (workers: boolean) => { await stopQueues(); - await startQueues([plain, failing, capped, exclusive], { + await startQueues({ + queues: [plain, failing, capped, exclusive], + postgres: env.postgres, workers, pollingIntervalSeconds: 0.5, boss: { schema: SCHEMA }, @@ -103,6 +105,18 @@ describe("task queue", () => { await sql`drop schema if exists ${sql.id(SCHEMA)} cascade`.execute(env.db); }); + it("costs nothing until a queue is registered: no schema, and a clear error on enqueue", async () => { + await stopQueues(); + + await startQueues({ queues: [], postgres: env.postgres, boss: { schema: SCHEMA } }); + + const schemas = await sql<{ count: string }>` + select count(*) as count from information_schema.schemata where schema_name = ${SCHEMA} + `.execute(env.db); + expect(Number(schemas.rows[0]?.count)).toBe(0); + await expect(enqueue(plain, { value: "nowhere to go" })).rejects.toThrow(/task queue is not running/); + }); + it("runs an enqueued task with its payload and attempt number", async () => { await start(true); await enqueue(plain, { value: "hello" }); diff --git a/apps/backend/src/libs/queue/queue.ts b/apps/backend/src/libs/queue/queue.ts index 86072e3..80ce525 100644 --- a/apps/backend/src/libs/queue/queue.ts +++ b/apps/backend/src/libs/queue/queue.ts @@ -1,11 +1,10 @@ import { captureException } from "@sentry/node"; -import type { Transaction } from "kysely"; -import { fromKysely, PgBoss, type ConstructorOptions, type JobWithMetadata } from "pg-boss"; +import { CompiledQuery, type Kysely, type Transaction } from "kysely"; +import { fromKysely, PgBoss, type ConstructorOptions, type Db, type JobWithMetadata } from "pg-boss"; import type { z } from "zod"; import type { DB } from "../../db/postgres/types/types"; import { logger } from "../logger/logger"; -import { runsWorkers } from "./process-role"; /** * Task queue on Postgres (pg-boss). @@ -14,12 +13,15 @@ import { runsWorkers } from "./process-role"; * tick that should return fast; anything slow or heavy (parsing an upload, calling a model, * sending a batch) belongs here as a task: * - * - Tasks live in Postgres, in their own `pgboss` schema, so there is nothing extra to deploy. + * - Tasks live in Postgres, in their own `pgboss` schema, so there is nothing extra to deploy. The + * queue runs on the connection it is handed (`startQueues({ queues, postgres })`), so it opens no + * pool of its own, and the schema only appears once a queue is registered; until then this + * module does nothing at all. * `enqueue` accepts the caller's transaction: the task commits or rolls back together with the * row it is about, so a row can never sit "queued" with no task behind it. - * - Every backend instance works tasks, with a per-queue concurrency cap per instance. Adding a - * replica adds capacity. To split the load off the API, run the same image a second time with - * PROCESS_ROLE=worker and set the API to PROCESS_ROLE=web. See `process-role.ts`. + * - Tasks run inside the backend process, with a per-queue concurrency cap. Every instance works + * them, so adding a replica adds capacity. Nothing here assumes the server is the only place + * that can work them: anything that has a database connection can call `startQueues`. * - A failed run is retried with backoff up to `retryLimit`, then `onExhausted` fires once so the * feature can mark its own row as failed. A handler that hits an error no retry can fix (bad * file, unsupported format) should record that itself and return normally. @@ -80,7 +82,11 @@ export type EnqueueOptions = { }; export type StartQueuesOptions = { - /** Overrides the role and test-environment default. Specs turn workers on to exercise retries. */ + /** The queues to register, normally the array from src/queues/index.ts. */ + queues: AnyQueueDefinition[]; + /** The backend's database handle. The queue runs on its connections instead of opening a pool. */ + postgres: { qb: Kysely }; + /** Work tasks in this process. Defaults to true, and to false under NODE_ENV=test. */ workers?: boolean; /** Seconds between polls of an idle queue. Defaults to 2. */ pollingIntervalSeconds?: number; @@ -154,19 +160,34 @@ const registerQueue = async (instance: PgBoss, queue: AnyQueueDefinition): Promi }; /** - * Connect, install or migrate the `pgboss` schema, register every queue, and (unless this process - * is web-only or under test) start working them. Every role calls this, because every role may - * enqueue. Safe to call once per process; call `stopQueues` on shutdown. + * pg-boss speaks plain SQL through one method, so it can run on the backend's own connections. The + * plugins come off for it: CamelCasePlugin would rename the columns pg-boss reads back. */ -export const startQueues = async (queues: AnyQueueDefinition[], options: StartQueuesOptions = {}): Promise => { +const toBossDb = (postgres: StartQueuesOptions["postgres"]): Db => { + const raw = postgres.qb.withoutPlugins(); + return { + executeSql: async (text, values = []) => { + const result = await raw.executeQuery(CompiledQuery.raw(text, values)); + return { rows: result.rows }; + }, + }; +}; + +/** + * Install or migrate the `pgboss` schema, register every queue, and (unless under test) start + * working them. Safe to call once per process; call `stopQueues` on shutdown. + */ +export const startQueues = async (options: StartQueuesOptions): Promise => { + const { queues } = options; if (boss) return; + // Free until used: a project with no queue gets no `pgboss` schema and no polling. The registry + // ships empty, and the first pack that adds a queue turns this on. + if (queues.length === 0) return; assertValidQueues(queues); const instance = new PgBoss({ - connectionString: process.env.DATABASE_URL, + db: toBossDb(options.postgres), schema: SCHEMA, - // Small on purpose: this pool sits next to the Prisma and Kysely pools on the same database. - max: 4, // Cron stays with the scheduler module; one clock is enough. schedule: false, ...options.boss, @@ -182,7 +203,7 @@ export const startQueues = async (queues: AnyQueueDefinition[], options: StartQu boss = instance; registered = new Set(queues.map((queue) => queue.name)); - const working = options.workers ?? (runsWorkers() && process.env.NODE_ENV !== "test"); + const working = options.workers ?? process.env.NODE_ENV !== "test"; if (working) { await Promise.all( queues.map((queue) => @@ -216,7 +237,12 @@ export const enqueue = async ( data: T, options: EnqueueOptions = {}, ): Promise => { - if (!boss) throw new Error("Queue: enqueue called before startQueues"); + if (!boss) { + throw new Error( + `Queue: cannot enqueue "${queue.name}", the task queue is not running. ` + + "Register the queue in src/queues/index.ts (it only starts when at least one is registered).", + ); + } if (!registered.has(queue.name)) { throw new Error(`Queue: "${queue.name}" is not registered in src/queues/index.ts`); } @@ -230,7 +256,7 @@ export const enqueue = async ( }); }; -/** Lets in-flight tasks finish (up to 30s), then closes the pool. */ +/** Lets in-flight tasks finish (up to 30s). The database connection belongs to the caller and stays open. */ export const stopQueues = async (): Promise => { if (!boss) return; const instance = boss; diff --git a/apps/backend/src/main.ts b/apps/backend/src/main.ts index 9380a66..49d1652 100644 --- a/apps/backend/src/main.ts +++ b/apps/backend/src/main.ts @@ -6,14 +6,13 @@ import { Hono } from "hono"; import { cors } from "hono/cors"; import { Env, validateEnv } from "./config/env/env.config"; -import { setupContext } from "./context"; +import { postgres, setupContext } from "./context"; import { createAppSwapper, createBootApp } from "./libs/boot/boot-app"; import { BootState, BootStage, getBootState, setBootError } from "./libs/boot/boot-state"; import { logger } from "./libs/logger/logger"; import { ApplicationError, AuthorizationError, DatabaseError, ValidationError } from "./middleware/error"; import { AuthError } from "./middleware/error/auth-error/types"; import { errorMiddleware, onError } from "./middleware/error/error-middleware"; -import { getProcessRole, runsWorkers, servesApi } from "./libs/queue/process-role"; import { startQueues, stopQueues } from "./libs/queue/queue"; import { startScheduler, stopScheduler } from "./libs/scheduler/scheduler"; import { m } from "./paraglide/messages.js"; @@ -72,20 +71,16 @@ const initialize = async (server: ReturnType): Promise => { app.get("/ping/*", (c) => { return c.json<{ message: string; success: boolean }>({ message: m.pong(), success: true }); }); - // PROCESS_ROLE=worker keeps /health for the platform's checks and serves nothing else. - if (servesApi()) { - registerSockets(app); - registerRoutes(app); - } + registerSockets(app); + registerRoutes(app); - // Every role connects to the task queue, because every role may enqueue; only roles that run - // workers pick tasks up. Queues are registered in src/queues/index.ts. - await startQueues(queues); + // The task queue lives in Postgres, so it starts after context setup; queues are registered in + // src/queues/index.ts and worked by this same process. + await startQueues({ queues, postgres }); // The scheduler needs the database (advisory locks, job_run bookkeeping), so it starts after // context setup; jobs are registered in src/jobs/index.ts. - if (runsWorkers()) startScheduler(jobs); - logger.info(`Process role: ${getProcessRole()}`); + startScheduler(jobs); /* ------------------------------------------------------------------------------------------------- * Handlers diff --git a/apps/backend/src/queues/index.ts b/apps/backend/src/queues/index.ts index 387e86a..fe6e7e5 100644 --- a/apps/backend/src/queues/index.ts +++ b/apps/backend/src/queues/index.ts @@ -3,7 +3,7 @@ import type { AnyQueueDefinition } from "../libs/queue/queue"; /** * Registered task queues. Ships empty on purpose: feature packs append entries at install time * through the CLI's queues codemod, and app code can add its own the same way. Every entry gets - * workers on each instance that runs them (see PROCESS_ROLE); see the queue module for retries, - * transactional enqueue, and why handlers must be safe to run twice. + * workers in this process; see the queue module for retries, transactional enqueue, and why + * handlers must be safe to run twice. */ export const queues: AnyQueueDefinition[] = []; diff --git a/apps/backend/src/testing/setup/test-utils.ts b/apps/backend/src/testing/setup/test-utils.ts index 3e336a4..781be6d 100644 --- a/apps/backend/src/testing/setup/test-utils.ts +++ b/apps/backend/src/testing/setup/test-utils.ts @@ -82,7 +82,7 @@ export function setupIntegrationTest(): TestEnv { await setupContext(testHonoApp); // Code under test may enqueue. Workers stay off in tests, so nothing runs behind a test's // back; call a queue's `run` directly to exercise a handler. - await startQueues(queues); + await startQueues({ queues, postgres: context.postgres }); sharedTestEnv = new TestEnv(context.postgres, context.valkey); }