diff --git a/bun.lock b/bun.lock index 0919c85d..d13ee69b 100644 --- a/bun.lock +++ b/bun.lock @@ -292,7 +292,7 @@ }, "packages/cli": { "name": "@wrnexus/cli", - "version": "0.8.51", + "version": "0.8.52", "bin": { "wrnexus": "src/index.ts", }, @@ -357,7 +357,7 @@ }, "packages/dev-server": { "name": "@wrnexus/dev-server", - "version": "0.8.46", + "version": "0.8.47", "dependencies": { "@wrnexus/authz": "workspace:*", "@wrnexus/cache": "workspace:*", @@ -546,7 +546,7 @@ }, "packages/queue": { "name": "@wrnexus/queue", - "version": "0.8.10", + "version": "0.8.11", "dependencies": { "@wrnexus/core": "workspace:*", "@wrnexus/rpc": "workspace:*", @@ -591,7 +591,7 @@ }, "packages/router": { "name": "@wrnexus/router", - "version": "0.8.9", + "version": "0.8.10", "dependencies": { "@wrnexus/compiler": "workspace:*", "@wrnexus/core": "workspace:*", diff --git a/packages/cli/package.json b/packages/cli/package.json index 6cb1bccb..46d94a54 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -1,6 +1,6 @@ { "name": "@wrnexus/cli", - "version": "0.8.51", + "version": "0.8.52", "type": "module", "main": "src/index.ts", "exports": { diff --git a/packages/cli/src/build.ts b/packages/cli/src/build.ts index ee81b9ca..670a97ae 100644 --- a/packages/cli/src/build.ts +++ b/packages/cli/src/build.ts @@ -763,6 +763,15 @@ applyAuthzManifestEarly([${authzSetupEntries}]); .join(", "); if (router.services.length) console.log(`✓ RPC services: ${router.services.length}`); + const queuesLit = (router.queues ?? []) + .map((queue) => { + const v = `q${counter++}`; + imports.push(`import * as ${v} from ${JSON.stringify(fwd(queue.file))};`); + return `{ name: ${JSON.stringify(queue.name)}, mod: ${v} }`; + }) + .join(", "); + if (router.queues?.length) console.log(`✓ Queues: ${router.queues.length}`); + // Authorization declarations again, this time for ProdOptions.authz — a // SEPARATE set of static imports of the exact same files (harmless; ES // modules are evaluated once and shared across every importer), statically @@ -798,6 +807,7 @@ await createProductionServer( components: [${componentsLit}], layouts: [${layoutsLit}], services: [${servicesLit}], + queues: [${queuesLit}], }, { reactivePath: join(import.meta.dir, "reactive.js"), diff --git a/packages/dev-server/package.json b/packages/dev-server/package.json index 6158626e..ea3f2d0f 100644 --- a/packages/dev-server/package.json +++ b/packages/dev-server/package.json @@ -1,6 +1,6 @@ { "name": "@wrnexus/dev-server", - "version": "0.8.46", + "version": "0.8.47", "type": "module", "main": "src/index.ts", "exports": { diff --git a/packages/dev-server/src/index.ts b/packages/dev-server/src/index.ts index 35fc6852..aa4a91fe 100644 --- a/packages/dev-server/src/index.ts +++ b/packages/dev-server/src/index.ts @@ -813,6 +813,7 @@ export async function startServer(opts: ServeOptions): Promise { watcher?.close(); unsubscribeDevToolbar?.(); server.stop(); + void handlers.shutdown(); void pluginRunner.hook("shutdown"); rmSync(cacheDir, { recursive: true, force: true, maxRetries: 3, retryDelay: 100 }); }, diff --git a/packages/dev-server/src/prod.ts b/packages/dev-server/src/prod.ts index 21746129..0aefde9c 100644 --- a/packages/dev-server/src/prod.ts +++ b/packages/dev-server/src/prod.ts @@ -79,6 +79,8 @@ export interface ProdManifest { layouts: { name: string; mod: RouteModule }[]; /** RPC service implementations (from app/services/*.ts). */ services?: { name: string; mod: RouteModule }[]; + /** Background queue modules (from app/queues/*.ts). */ + queues?: { name: string; mod: RouteModule }[]; } export interface ProductionPluginAsset { @@ -292,6 +294,7 @@ function buildProdRouter(manifest: ProdManifest): { for (const l of manifest.layouts) modules.set(`layout:${l.name}`, l.mod); for (const service of manifest.services ?? []) modules.set(`service:${service.name}`, service.mod); + for (const queue of manifest.queues ?? []) modules.set(`queue:${queue.name}`, queue.mod); const router: Router = { pages, @@ -307,6 +310,10 @@ function buildProdRouter(manifest: ProdManifest): { name: service.name, file: `service:${service.name}`, })), + queues: (manifest.queues ?? []).map((queue) => ({ + name: queue.name, + file: `queue:${queue.name}`, + })), matchPage: optimizedMatcher(pages), matchApi: optimizedMatcher(api), matchRealtime: optimizedMatcher(realtime), @@ -652,11 +659,12 @@ export async function createProductionServer(manifest: ProdManifest, opts: ProdO // Graceful shutdown: stop accepting connections, then exit. let shuttingDown = false; - const shutdown = () => { + const shutdown = async () => { if (shuttingDown) return; shuttingDown = true; console.log("WrNexus: shutting down…"); server.stop(); + await handlers.shutdown(); process.exit(0); }; process.on("SIGTERM", shutdown); diff --git a/packages/dev-server/src/runtime.ts b/packages/dev-server/src/runtime.ts index 05d2994e..02746c55 100644 --- a/packages/dev-server/src/runtime.ts +++ b/packages/dev-server/src/runtime.ts @@ -909,6 +909,7 @@ export interface Ws { export interface Handlers { fetch(req: Request, server: UpgradeServer): Promise; + shutdown(): Promise; websocket: { open(ws: Ws): void; message(ws: Ws, message: string | Uint8Array): void; @@ -985,6 +986,27 @@ export function createHandlers(deps: RuntimeDeps): Handlers { } return servicesPromise; }; + type QueueModule = { start(): void; shutdown(options?: { force?: boolean }): Promise }; + let queuesPromise: Promise | undefined; + const loadQueues = (): Promise => { + queuesPromise ??= (async () => { + const queues: QueueModule[] = []; + for (const entry of router.queues ?? []) { + const imported = await loadModule(entry.file); + const queue = imported.default as QueueModule | undefined; + if (!queue || typeof queue.start !== "function" || typeof queue.shutdown !== "function") { + throw new Error(`Queue ${entry.file} must default-export defineQueue(...)`); + } + queue.start(); + queues.push(queue); + } + return queues; + })().catch((error) => { + queuesPromise = undefined; + throw error; + }); + return queuesPromise; + }; // Server-side realtime room manager (shared by every `defineRoom` connection). const realtime = createRealtimeRegistry(); @@ -1041,6 +1063,7 @@ export function createHandlers(deps: RuntimeDeps): Handlers { }; async function fetchHandler(req: Request, server: UpgradeServer): Promise { + await loadQueues(); // Behind a trusted proxy (nginx / gateway), honor X-Forwarded-Proto/Host so // ctx.url reflects the external HTTPS scheme — makes CSRF/session cookies Secure. const url = resolveRequestUrl(req, deps.security?.trustProxy); @@ -2321,6 +2344,10 @@ export function createHandlers(deps: RuntimeDeps): Handlers { return { fetch: fetchHandler, + async shutdown() { + const queues = await queuesPromise?.catch(() => []); + await Promise.allSettled((queues ?? []).map((queue) => queue.shutdown())); + }, websocket: { open(ws) { if (ws.data?.kind === "hmr") hub?.add(ws); diff --git a/packages/queue/README.md b/packages/queue/README.md index e2917849..fc3cf0d7 100644 --- a/packages/queue/README.md +++ b/packages/queue/README.md @@ -118,6 +118,45 @@ await queue.add("email", { to: "a@b.com" }, { delayMs: 5000, maxAttempts: 3 }); queue.start(); // begin polling; queue.stop() to halt ``` +### Application queues with `defineQueue` + +WrNexus applications can place typed queue modules in `app/queues`. The dev and +production servers discover them, start them after runtime initialization, and +shut them down gracefully: + +```ts +import { defineQueue } from "@wrnexus/queue"; + +export default defineQueue({ + name: "email", + jobs: { + send: { + maxAttempts: 3, + idempotency: ({ messageId }: { messageId: number }) => `message:${messageId}`, + run: async ({ messageId }) => sendMessage(messageId), + }, + }, +}); +``` + +Callers get a typed producer and inspection helpers: + +```ts +const job = await email.send.add({ messageId: 42 }); +await email.send.status(job.id); +await email.send.cancel(job.id); +await email.send.retry(job.id); +``` + +`createDurableQueue` also exposes `start()`, `stop()`, `isRunning()`, `health()` +and `get()` so applications do not need to build their own polling middleware. + +### SQLite + +Use `sqliteQueueStore(db)` from `@wrnexus/queue/sqlite` with a WrNexus-compatible +SQLite client. `installSqliteQueueSchema(db)` installs the queue table and its +due-job and idempotency indexes. + Use `context.signal` in network/database calls so forced shutdown and active cancellation finish promptly. For process termination, prefer `await queue.shutdown()`; use `{ force: true }` only after your grace period. diff --git a/packages/queue/package.json b/packages/queue/package.json index 0b40dbd7..d0d10767 100644 --- a/packages/queue/package.json +++ b/packages/queue/package.json @@ -1,11 +1,12 @@ { "name": "@wrnexus/queue", - "version": "0.8.10", + "version": "0.8.11", "private": true, "type": "module", "main": "src/index.ts", "exports": { - ".": "./src/index.ts" + ".": "./src/index.ts", + "./sqlite": "./src/sqlite.ts" }, "dependencies": { "@wrnexus/core": "workspace:*", diff --git a/packages/queue/src/defined.ts b/packages/queue/src/defined.ts new file mode 100644 index 00000000..7d658740 --- /dev/null +++ b/packages/queue/src/defined.ts @@ -0,0 +1,139 @@ +import type { AddOptions, Job, JobContext, JobHandler } from "./index.ts"; +import { + createDurableQueue, + type DurableQueue, + type DurableQueueOptions, + type QueueHealth, +} from "./durable.ts"; + +export type JobStatus = "queued" | "completed" | "failed" | "cancelled" | "missing"; + +export interface DefinedJobOptions extends Omit { + run: (data: T, context: JobContext & { job: Job }) => void | Promise; + failed?: (data: T, error: unknown, job: Job) => void | Promise; + idempotency?: (data: T) => string | undefined; + validate?: (data: unknown) => data is T; +} + +export type QueueJobDefinitions = Record>; + +export interface DefineQueueOptions { + name: string; + jobs: TJobs; + queue?: DurableQueue; + options?: DurableQueueOptions; +} + +type DataOf = T extends DefinedJobOptions ? I : never; + +export interface DefinedJob { + readonly name: string; + add(data: T, options?: AddOptions): Promise>; + addMany(entries: readonly { data: T; options?: AddOptions }[]): Promise[]>; + get(id: string): Promise | null>; + status(id: string): Promise; + list(): Promise[]>; + cancel(id: string): Promise; + retry(id: string): Promise; +} + +export type DefinedJobs = { + [K in keyof TJobs]: DefinedJob>; +}; + +export interface DefinedQueue { + readonly name: string; + readonly jobs: DefinedJobs; + readonly raw: DurableQueue; + start(): void; + stop(): void; + shutdown(options?: { force?: boolean }): Promise; + health(): Promise; + drain(): Promise; +} + +function fullName(queue: string, job: string): string { + return `${queue}:${job}`; +} + +export function defineQueue( + definition: DefineQueueOptions, +): DefinedQueue & DefinedJobs { + const name = definition.name.trim(); + if (!name) throw new TypeError("queue name cannot be empty"); + const queue = + definition.queue ?? + createDurableQueue({ + ...definition.options, + async onDeadLetter(job, error) { + await definition.options?.onDeadLetter?.(job, error); + const prefix = `${name}:`; + if (!job.name.startsWith(prefix)) return; + const failed = definition.jobs[job.name.slice(prefix.length)]?.failed; + await failed?.(job.data, error, job); + }, + }); + const jobs: Record> = {}; + + for (const [shortName, jobDefinition] of Object.entries(definition.jobs)) { + const jobName = fullName(name, shortName); + const handler: JobHandler = async (job, context) => { + if (jobDefinition.validate && !jobDefinition.validate(job.data)) { + throw new TypeError(`WRN-QUEUE-PAYLOAD: invalid payload for '${jobName}'`); + } + await jobDefinition.run(job.data, { ...context, job }); + }; + queue.process(jobName, handler); + + jobs[shortName] = { + name: jobName, + add(data, options = {}) { + if (jobDefinition.validate && !jobDefinition.validate(data)) { + return Promise.reject( + new TypeError(`WRN-QUEUE-PAYLOAD: invalid payload for '${jobName}'`), + ); + } + return queue.add(jobName, data, { + maxAttempts: jobDefinition.maxAttempts, + delayMs: jobDefinition.delayMs, + repeat: jobDefinition.repeat, + priority: jobDefinition.priority, + idempotencyKey: jobDefinition.idempotency?.(data), + ...options, + }); + }, + addMany(entries) { + return Promise.all(entries.map((entry) => this.add(entry.data, entry.options))); + }, + async get(id) { + const job = await queue.get(id); + return job?.name === jobName ? (job as Job) : null; + }, + async status(id) { + return queue.status(id); + }, + async list() { + return (await queue.list(jobName)) as Job[]; + }, + cancel(id) { + return queue.cancel(id); + }, + retry(id) { + return queue.retry(id); + }, + }; + } + + const result = { + name, + jobs: jobs as DefinedJobs, + raw: queue, + start: () => queue.start(), + stop: () => queue.stop(), + shutdown: (options?: { force?: boolean }) => queue.shutdown(options), + health: () => queue.health(), + drain: () => queue.drain(), + ...(jobs as DefinedJobs), + }; + return result; +} diff --git a/packages/queue/src/durable.ts b/packages/queue/src/durable.ts index 99929512..9509ddf3 100644 --- a/packages/queue/src/durable.ts +++ b/packages/queue/src/durable.ts @@ -17,10 +17,54 @@ export interface QueueStore { findByIdempotencyKey?(name: string, key: string): Promise; size?(): Promise; claim?(id: string, worker: string, leaseUntil: number, now?: number): Promise; + release?(id: string, worker: string): Promise; + archive?(record: QueueJobRecord): Promise; + history?(id?: string): Promise; +} + +export type QueueJobState = "completed" | "failed" | "cancelled"; +export interface QueueJobRecord { + job: Job; + state: QueueJobState; + finishedAt: number; + error?: string; +} + +export type QueueEventType = + | "job.added" + | "job.started" + | "job.completed" + | "job.retrying" + | "job.failed" + | "job.cancelled" + | "worker.started" + | "worker.stopped" + | "worker.error"; + +export interface QueueEvent { + type: QueueEventType; + at: number; + workerId: string; + job?: Job; + error?: unknown; + durationMs?: number; +} + +export interface QueueHealth { + running: boolean; + accepting: boolean; + draining: boolean; + active: number; + pending: number; + failed: number; + workerId: string; + lastPollAt?: number; + lastError?: unknown; } export function memoryQueueStore(): QueueStore { const jobs = new Map(); + const records = new Map(); return { async put(job) { @@ -59,6 +103,14 @@ export function memoryQueueStore(): QueueStore { async size() { return jobs.size; }, + async archive(record) { + records.set(record.job.id, structuredClone(record)); + }, + async history(id) { + return [...records.values()] + .filter((record) => !id || record.job.id === id) + .map((record) => structuredClone(record)); + }, }; } @@ -69,16 +121,27 @@ export interface DurableQueueOptions { concurrency?: number; capacity?: number; leaseMs?: number; + pollMs?: number; backoff?: (attempt: number) => number; now?: () => number; onDeadLetter?: (job: Job, error: unknown) => void | Promise; context?: (job: Job, signal: AbortSignal) => ExecutionContext; + onEvent?: (event: QueueEvent) => void | Promise; + /** Maintenance hook run before each poll, for application-specific recovery. */ + beforeDrain?: (queue: DurableQueue) => void | Promise; } export interface DurableQueue { add(name: string, data: T, options?: AddOptions): Promise>; process(name: string, handler: JobHandler): void; drain(): Promise; + start(): void; + stop(): void; + isRunning(): boolean; + health(): Promise; + get(id: string): Promise; + status(id: string): Promise<"queued" | QueueJobState | "missing">; + history(id?: string): Promise; cancel(id: string): Promise; shutdown(options?: { force?: boolean }): Promise; list(name?: string): Promise; @@ -114,11 +177,19 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu const concurrency = positiveInteger(options.concurrency ?? 10, "queue concurrency"); const capacity = positiveInteger(options.capacity ?? 10_000, "queue capacity"); const leaseMs = positiveInteger(options.leaseMs ?? 30_000, "queue leaseMs"); + const pollMs = positiveInteger(options.pollMs ?? 250, "queue pollMs"); let sequence = 0; let draining = false; let accepting = true; + let timer: ReturnType | null = null; + let lastPollAt: number | undefined; + let lastError: unknown; const active = new Map }>(); + const emit = async (event: Omit) => { + await options.onEvent?.({ ...event, at: now(), workerId }); + }; + function beginJob(job: Job): Promise { const controller = new AbortController(); const promise = runJob(job, controller).finally(() => active.delete(job.id)); @@ -135,6 +206,8 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu } job.attempts += 1; + const startedAt = now(); + void emit({ type: "job.started", job: structuredClone(job) }); try { await handler(job, { @@ -153,8 +226,14 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu job.runAt = now() + job.repeat; await store.put(job); } else { + await store.archive?.({ job: structuredClone(job), state: "completed", finishedAt: now() }); await store.remove(job.id); } + await emit({ + type: "job.completed", + job: structuredClone(job), + durationMs: now() - startedAt, + }); } catch (error) { if (controller.signal.aborted) { await store.remove(job.id); @@ -167,17 +246,25 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu ); job.runAt = now() + delay; await store.put(job); + await emit({ type: "job.retrying", job: structuredClone(job), error }); } else { await store.remove(job.id); deadLetters.set(job.id, structuredClone(job)); + await store.archive?.({ + job: structuredClone(job), + state: "failed", + finishedAt: now(), + error: error instanceof Error ? error.message : String(error), + }); await options.onDeadLetter?.(structuredClone(job), error); + await emit({ type: "job.failed", job: structuredClone(job), error }); } } return true; } - return { + const durableQueue: DurableQueue = { async add(name: string, data: T, add: AddOptions = {}) { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); if (!name.trim()) throw new TypeError("queue job name cannot be empty"); @@ -216,6 +303,7 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu createdAt, }; await store.put(job); + await emit({ type: "job.added", job: structuredClone(job) }); return structuredClone(job); }, @@ -233,7 +321,9 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu async drain() { if (draining) return 0; draining = true; + lastPollAt = now(); try { + if (options.beforeDrain) await options.beforeDrain(durableQueue); const due = await store.due(now(), concurrency); const results = await Promise.all(due.map(beginJob)); return results.filter(Boolean).length; @@ -250,19 +340,80 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu return store.list(name); }, + get(id) { + return store.get(id); + }, + + async status(id) { + if (await store.get(id)) return "queued"; + const record = (await store.history?.(id))?.at(-1); + if (record) return record.state; + if (deadLetters.has(id)) return "failed"; + return "missing"; + }, + + async history(id) { + return (await store.history?.(id)) ?? []; + }, + + start() { + if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); + if (timer) return; + timer = setInterval(() => { + void this.drain().catch(async (error: unknown) => { + lastError = error; + await emit({ type: "worker.error", error }); + }); + }, pollMs); + timer.unref?.(); + void emit({ type: "worker.started" }); + }, + + stop() { + if (timer) clearInterval(timer); + timer = null; + void emit({ type: "worker.stopped" }); + }, + + isRunning() { + return timer !== null; + }, + + async health() { + return { + running: timer !== null, + accepting, + draining, + active: active.size, + pending: store.size ? await store.size() : (await store.list()).length, + failed: deadLetters.size, + workerId, + lastPollAt, + lastError, + }; + }, + async cancel(id) { const running = active.get(id); if (running) { running.controller.abort(); + const job = await store.get(id); + if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() }); + await emit({ type: "job.cancelled" }); return true; } if (!(await store.get(id))) return false; + const job = await store.get(id); + if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() }); await store.remove(id); + await emit({ type: "job.cancelled" }); return true; }, async shutdown(shutdownOptions = {}) { accepting = false; + if (timer) clearInterval(timer); + timer = null; if (shutdownOptions.force) { for (const { controller } of active.values()) controller.abort(); } @@ -270,7 +421,13 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu }, async retry(id) { - const job = deadLetters.get(id); + let job = deadLetters.get(id); + if (!job) { + const record = (await store.history?.(id))?.find( + (candidate) => candidate.state === "failed", + ); + job = record?.job; + } if (!job) return false; deadLetters.delete(id); job.attempts = 0; @@ -279,4 +436,5 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu return true; }, }; + return durableQueue; } diff --git a/packages/queue/src/index.ts b/packages/queue/src/index.ts index 38d3e8b8..af96a612 100644 --- a/packages/queue/src/index.ts +++ b/packages/queue/src/index.ts @@ -298,6 +298,25 @@ export function createQueue(options: QueueOptions = {}): Queue { } export { memoryQueueStore, createDurableQueue } from "./durable.ts"; export type { QueueStore, DurableQueue, DurableQueueOptions } from "./durable.ts"; +export type { + QueueEvent, + QueueEventType, + QueueHealth, + QueueJobRecord, + QueueJobState, +} from "./durable.ts"; +export { defineQueue } from "./defined.ts"; +export type { + DefinedJob, + DefinedJobOptions, + DefinedJobs, + DefinedQueue, + DefineQueueOptions, + JobStatus, + QueueJobDefinitions, +} from "./defined.ts"; +export { installSqliteQueueSchema, sqliteQueueSchema, sqliteQueueStore } from "./sqlite.ts"; +export type { SqliteQueueClient } from "./sqlite.ts"; export { redisQueueStore, postgresQueueStore, POSTGRES_QUEUE_SCHEMA } from "./stores.ts"; export type { RedisQueueClient, SqlQueueClient } from "./stores.ts"; export { diff --git a/packages/queue/src/scheduler.ts b/packages/queue/src/scheduler.ts index dd8f23a0..0143d934 100644 --- a/packages/queue/src/scheduler.ts +++ b/packages/queue/src/scheduler.ts @@ -85,10 +85,11 @@ export interface QueueDashboardSnapshot { export async function queueDashboardSnapshot(queue: DurableQueue): Promise { const jobs = await queue.list(); + const history = await queue.history(); return { generatedAt: Date.now(), pending: jobs.length, - failed: queue.failed().length, + failed: history.filter((record) => record.state === "failed").length || queue.failed().length, byName: jobs.reduce>((counts, job) => { counts[job.name] = (counts[job.name] ?? 0) + 1; return counts; diff --git a/packages/queue/src/sqlite.ts b/packages/queue/src/sqlite.ts new file mode 100644 index 00000000..9200d7fc --- /dev/null +++ b/packages/queue/src/sqlite.ts @@ -0,0 +1,167 @@ +import type { Job } from "./index.ts"; +import type { QueueJobRecord, QueueStore } from "./durable.ts"; + +export interface SqliteQueueClient { + exec(sql: string, parameters?: unknown[]): Promise<{ changes: number }>; + one(sql: string, parameters?: unknown[]): Promise; + all(sql: string, parameters?: unknown[]): Promise; +} + +interface QueueJobRow { + id: string; + name: string; + payload: string; + run_at: number; + priority: number; + idempotency_key: string | null; +} + +const SAFE_TABLE = /^[a-z_][a-z0-9_]*$/i; + +export function sqliteQueueSchema(table = "wrnexus_jobs"): string { + if (!SAFE_TABLE.test(table)) throw new Error("Invalid queue table name"); + return `CREATE TABLE IF NOT EXISTS ${table} ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL, + payload TEXT NOT NULL, + run_at INTEGER NOT NULL, + priority INTEGER NOT NULL DEFAULT 0, + idempotency_key TEXT, + lease_owner TEXT, + lease_until INTEGER + ); + CREATE INDEX IF NOT EXISTS ${table}_due ON ${table} (run_at, priority DESC); + CREATE UNIQUE INDEX IF NOT EXISTS ${table}_idempotency ON ${table} (name, idempotency_key) + WHERE idempotency_key IS NOT NULL; + CREATE TABLE IF NOT EXISTS ${table}_history ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL, + state TEXT NOT NULL, + payload TEXT NOT NULL, + error TEXT, + finished_at INTEGER NOT NULL + ); + CREATE INDEX IF NOT EXISTS ${table}_history_finished ON ${table}_history (finished_at DESC);`; +} + +export async function installSqliteQueueSchema( + db: SqliteQueueClient, + table = "wrnexus_jobs", +): Promise { + for (const statement of sqliteQueueSchema(table) + .split(";") + .map((value) => value.trim())) { + if (statement) await db.exec(statement); + } +} + +export function sqliteQueueStore(db: SqliteQueueClient, table = "wrnexus_jobs"): QueueStore { + if (!SAFE_TABLE.test(table)) throw new Error("Invalid queue table name"); + const decode = (row: QueueJobRow): Job => JSON.parse(row.payload) as Job; + return { + async put(job) { + await db.exec( + `INSERT INTO ${table} (id,name,payload,run_at,priority,idempotency_key,lease_owner,lease_until) + VALUES (?,?,?,?,?,?,NULL,NULL) + ON CONFLICT(id) DO UPDATE SET name=excluded.name,payload=excluded.payload, + run_at=excluded.run_at,priority=excluded.priority,idempotency_key=excluded.idempotency_key, + lease_owner=NULL,lease_until=NULL`, + [ + job.id, + job.name, + JSON.stringify(job), + job.runAt, + job.priority, + job.idempotencyKey ?? null, + ], + ); + }, + async get(id) { + const row = await db.one(`SELECT * FROM ${table} WHERE id = ?`, [id]); + return row ? decode(row) : null; + }, + async remove(id) { + await db.exec(`DELETE FROM ${table} WHERE id = ?`, [id]); + }, + async due(now, limit) { + const rows = await db.all( + `SELECT * FROM ${table} WHERE run_at <= ? AND (lease_until IS NULL OR lease_until < ?) + ORDER BY priority DESC,run_at ASC LIMIT ?`, + [now, now, limit], + ); + return rows.map(decode); + }, + async list(name) { + const rows = name + ? await db.all(`SELECT * FROM ${table} WHERE name = ? ORDER BY run_at`, [name]) + : await db.all(`SELECT * FROM ${table} ORDER BY run_at`); + return rows.map(decode); + }, + async findByIdempotencyKey(name, key) { + const row = await db.one( + `SELECT * FROM ${table} WHERE name = ? AND idempotency_key = ? LIMIT 1`, + [name, key], + ); + return row ? decode(row) : null; + }, + async size() { + const row = await db.one<{ total: number }>(`SELECT COUNT(*) AS total FROM ${table}`); + return Number(row?.total ?? 0); + }, + async claim(id, worker, leaseUntil, now = Date.now()) { + const result = await db.exec( + `UPDATE ${table} SET lease_owner = ?,lease_until = ? + WHERE id = ? AND (lease_until IS NULL OR lease_until < ?)`, + [worker, leaseUntil, id, now], + ); + return result.changes > 0; + }, + async release(id, worker) { + await db.exec( + `UPDATE ${table} SET lease_owner = NULL,lease_until = NULL WHERE id = ? AND lease_owner = ?`, + [id, worker], + ); + }, + async archive(record) { + await db.exec( + `INSERT INTO ${table}_history (id,name,state,payload,error,finished_at) VALUES (?,?,?,?,?,?) + ON CONFLICT(id) DO UPDATE SET state=excluded.state,payload=excluded.payload, + error=excluded.error,finished_at=excluded.finished_at`, + [ + record.job.id, + record.job.name, + record.state, + JSON.stringify(record.job), + record.error ?? null, + record.finishedAt, + ], + ); + }, + async history(id) { + const rows = id + ? await db.all<{ + state: QueueJobRecord["state"]; + payload: string; + error: string | null; + finished_at: number; + }>( + `SELECT state,payload,error,finished_at FROM ${table}_history WHERE id = ? ORDER BY finished_at DESC`, + [id], + ) + : await db.all<{ + state: QueueJobRecord["state"]; + payload: string; + error: string | null; + finished_at: number; + }>( + `SELECT state,payload,error,finished_at FROM ${table}_history ORDER BY finished_at DESC`, + ); + return rows.map((row) => ({ + job: JSON.parse(row.payload) as Job, + state: row.state, + finishedAt: row.finished_at, + error: row.error ?? undefined, + })); + }, + }; +} diff --git a/packages/queue/test/defined.test.ts b/packages/queue/test/defined.test.ts new file mode 100644 index 00000000..13179a9d --- /dev/null +++ b/packages/queue/test/defined.test.ts @@ -0,0 +1,62 @@ +import { expect, test } from "bun:test"; +import { createDurableQueue, defineQueue, memoryQueueStore } from "../src/index.ts"; + +test("defineQueue creates typed producers and registers workers", async () => { + const seen: number[] = []; + const email = defineQueue({ + name: "email", + options: { store: memoryQueueStore() }, + jobs: { + send: { + maxAttempts: 4, + idempotency: (data: { messageId: number }) => `message:${data.messageId}`, + validate: (data: unknown): data is { messageId: number } => + typeof data === "object" && + data !== null && + Number.isInteger((data as { messageId?: unknown }).messageId), + run: async ({ messageId }: { messageId: number }) => void seen.push(messageId), + }, + }, + }); + + const first = await email.send.add({ messageId: 7 }); + const duplicate = await email.jobs.send.add({ messageId: 7 }); + expect(duplicate.id).toBe(first.id); + expect(await email.send.status(first.id)).toBe("queued"); + expect(await email.drain()).toBe(1); + expect(seen).toEqual([7]); + expect(await email.send.status(first.id)).toBe("completed"); +}); + +test("durable queues expose lifecycle, health and events", async () => { + const events: string[] = []; + const queue = createDurableQueue({ + pollMs: 10, + onEvent: (event) => void events.push(event.type), + }); + queue.process("work", async () => {}); + queue.start(); + expect(queue.isRunning()).toBe(true); + await queue.add("work", {}); + await queue.drain(); + const health = await queue.health(); + expect(health.pending).toBe(0); + expect(events).toContain("job.added"); + expect(events).toContain("job.completed"); + queue.stop(); + expect(queue.isRunning()).toBe(false); + await queue.shutdown(); +}); + +test("defineQueue validates payloads before persistence", async () => { + const queue = defineQueue({ + name: "safe", + jobs: { + number: { + validate: (data: unknown): data is number => typeof data === "number", + run: async (_data: number) => {}, + }, + }, + }); + await expect(queue.number.add("no" as never)).rejects.toThrow("WRN-QUEUE-PAYLOAD"); +}); diff --git a/packages/router/package.json b/packages/router/package.json index 1e8b313c..80cfd32e 100644 --- a/packages/router/package.json +++ b/packages/router/package.json @@ -1,6 +1,6 @@ { "name": "@wrnexus/router", - "version": "0.8.9", + "version": "0.8.10", "type": "module", "main": "src/index.ts", "exports": { diff --git a/packages/router/src/index.ts b/packages/router/src/index.ts index cbf3f4ec..2f252d18 100644 --- a/packages/router/src/index.ts +++ b/packages/router/src/index.ts @@ -63,6 +63,8 @@ export interface Router { authz: ComponentRef[]; /** Service implementations (`app/services/.ts`) mounted for inter-app calls. */ services: ComponentRef[]; + /** Background queues (`app/queues/.ts`) loaded and started with the server. */ + queues?: ComponentRef[]; matchPage(pathname: string): RouteMatch | null; matchApi(pathname: string): RouteMatch | null; matchRealtime(pathname: string): RouteMatch | null; @@ -329,6 +331,25 @@ export function buildRouter(appDir: string, opts: RouterOptions = {}): Router { services.push({ name, file: f.file }); } + const queues: ComponentRef[] = []; + const queueFilesByName = new Map(); + for (const f of scanDir(join(appDir, "queues"), [".js"])) { + if (!/\.(ts|js)$/.test(f.file) || /[.]gen[.](ts|js)$/.test(f.file)) continue; + const name = f.rel.replace(/\.(ts|js)$/, "").replace(/\//g, ":"); + if (!/^[A-Za-z][A-Za-z0-9_:-]*$/.test(name)) { + console.warn(`[wrnexus] skipping queue with unsafe name: ${name}`); + continue; + } + const existing = queueFilesByName.get(name); + if (existing) { + throw new Error( + `WRN-QUEUE-COLLISION: two queues are named "${name}": ${existing} and ${f.file}`, + ); + } + queueFilesByName.set(name, f.file); + queues.push({ name, file: f.file }); + } + return { pages, api, @@ -340,6 +361,7 @@ export function buildRouter(appDir: string, opts: RouterOptions = {}): Router { schemas, authz, services, + queues, matchPage: (p) => matchRoute(pages, p), matchApi: (p) => matchRoute(api, p), matchRealtime: (p) => matchRoute(realtime, p),