Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 3 additions & 11 deletions .cursor/rules/background-work.mdc
Original file line number Diff line number Diff line change
@@ -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
---
Expand Down Expand Up @@ -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.
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 0 additions & 8 deletions apps/backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
# -----------------
Expand Down
2 changes: 0 additions & 2 deletions apps/backend/src/config/env/env.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof envSchema>;
Expand Down
34 changes: 0 additions & 34 deletions apps/backend/src/libs/queue/process-role.spec.ts

This file was deleted.

30 changes: 0 additions & 30 deletions apps/backend/src/libs/queue/process-role.ts

This file was deleted.

16 changes: 15 additions & 1 deletion apps/backend/src/libs/queue/queue.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand All @@ -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" });
Expand Down
62 changes: 44 additions & 18 deletions apps/backend/src/libs/queue/queue.ts
Original file line number Diff line number Diff line change
@@ -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).
Expand All @@ -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.
Expand Down Expand Up @@ -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<DB> };
/** 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;
Expand Down Expand Up @@ -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<void> => {
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<void> => {
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,
Expand All @@ -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) =>
Expand Down Expand Up @@ -216,7 +237,12 @@ export const enqueue = async <T>(
data: T,
options: EnqueueOptions = {},
): Promise<string | null> => {
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`);
}
Expand All @@ -230,7 +256,7 @@ export const enqueue = async <T>(
});
};

/** 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<void> => {
if (!boss) return;
const instance = boss;
Expand Down
19 changes: 7 additions & 12 deletions apps/backend/src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -72,20 +71,16 @@ const initialize = async (server: ReturnType<typeof serve>): Promise<Hono> => {
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
Expand Down
4 changes: 2 additions & 2 deletions apps/backend/src/queues/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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[] = [];
2 changes: 1 addition & 1 deletion apps/backend/src/testing/setup/test-utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Loading