diff --git a/.env.example b/.env.example index d6b8c521..1929c080 100644 --- a/.env.example +++ b/.env.example @@ -39,6 +39,13 @@ REDIS_URL=redis://localhost:6379 # Wallet-provisioning worker poll interval (ms) WORKER_POLL_INTERVAL_MS=5000 +SCHEDULER_INTERVAL_MS=15000 +SCHEDULER_LEASE_MS=60000 +SCHEDULER_SHUTDOWN_TIMEOUT_MS=30000 +SCHEDULER_QUEUES= +SCHEDULER_DISABLED_QUEUES= +SCHEDULER_IN_PROCESS=false + # Logging Configuration # LOG_LEVEL=info (options: error, warn, info, http, verbose, debug, silly) @@ -55,7 +62,6 @@ STELLAR_FUNDING_MAX_RETRIES=5 # Account Lifecycle Configuration DELETION_COOLING_OFF_DAYS=30 EXPORT_TTL_DAYS=7 -# Interval for the background lifecycle sweep (export generation, deletion finalization). 0 = disabled (lazy sweep on requests only) LIFECYCLE_SWEEP_INTERVAL_MS=0 # Phone OTP Configuration @@ -63,9 +69,9 @@ LIFECYCLE_SWEEP_INTERVAL_MS=0 SMS_PROVIDER=mock RATE_LIMIT_OTP_WINDOW_MS=900000 RATE_LIMIT_OTP_MAX=5 - -# Data Lifecycle / Audit Configuration — see docs/DATA_LIFECYCLE.md -# HMAC key for the source-IP hash on immutable audit events. Rotating it makes -# older hashes uncorrelatable with newer ones. If unset in production, audit -# events omit the IP hash rather than storing an unkeyed (reversible) digest. -# AUDIT_IP_HASH_SECRET=change-me-in-production + +# Data Lifecycle / Audit Configuration — see docs/DATA_LIFECYCLE.md +# HMAC key for the source-IP hash on immutable audit events. Rotating it makes +# older hashes uncorrelatable with newer ones. If unset in production, audit +# events omit the IP hash rather than storing an unkeyed (reversible) digest. +# AUDIT_IP_HASH_SECRET=change-me-in-production diff --git a/Dockerfile b/Dockerfile index db4dc82f..9541e982 100644 --- a/Dockerfile +++ b/Dockerfile @@ -85,8 +85,9 @@ COPY --from=build /app/prisma.config.ts ./ # Entrypoint scripts COPY docker/entrypoint-api.sh ./entrypoint-api.sh COPY docker/entrypoint-worker.sh ./entrypoint-worker.sh +COPY docker/entrypoint-scheduler.sh ./entrypoint-scheduler.sh -RUN chmod +x entrypoint-api.sh entrypoint-worker.sh +RUN chmod +x entrypoint-api.sh entrypoint-worker.sh entrypoint-scheduler.sh # Non-root user RUN groupadd --gid 1001 appgroup && \ diff --git a/docker-compose.yml b/docker-compose.yml index 8921f699..5341e6f9 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,12 +1,11 @@ # ============================================================================ # Learnault API — Local Development Stack # -# One command starts a healthy stack (API + wallet worker + PostgreSQL + Redis): +# One command starts a healthy stack (API + PostgreSQL + Redis): # docker compose up -d --build # # The API container applies migrations and seeds deterministic fixtures on -# boot (see docker/entrypoint-dev-api.sh). The worker drains the idempotent -# wallet-provisioning outbox (see src/workers/wallet-provisioning.worker.ts). +# boot (see docker/entrypoint-dev-api.sh). # # Useful commands: # docker compose ps → service status + health @@ -95,10 +94,7 @@ services: start_period: 20s stop_grace_period: 30s - # -------------------------------------------------------------------------- - # Worker — drains the wallet-provisioning outbox - # -------------------------------------------------------------------------- - worker: + scheduler: build: context: . dockerfile: docker/Dockerfile.dev @@ -110,12 +106,17 @@ services: NODE_ENV: development DATABASE_URL: postgresql://${POSTGRES_USER:-learnault}:${POSTGRES_PASSWORD:-learnault}@db:5432/${POSTGRES_DB:-learnault_dev}?schema=public RUN_MIGRATIONS: 'true' - WORKER_POLL_INTERVAL_MS: ${WORKER_POLL_INTERVAL_MS:-5000} + LOG_LEVEL: ${LOG_LEVEL:-info} + SCHEDULER_INTERVAL_MS: ${SCHEDULER_INTERVAL_MS:-15000} + SCHEDULER_LEASE_MS: ${SCHEDULER_LEASE_MS:-60000} + SCHEDULER_SHUTDOWN_TIMEOUT_MS: ${SCHEDULER_SHUTDOWN_TIMEOUT_MS:-30000} + SCHEDULER_QUEUES: ${SCHEDULER_QUEUES:-} + SCHEDULER_DISABLED_QUEUES: ${SCHEDULER_DISABLED_QUEUES:-} volumes: - .:/app - /app/node_modules - command: ['./docker/entrypoint-dev-worker.sh'] - stop_grace_period: 30s + command: ['./docker/entrypoint-dev-scheduler.sh'] + stop_grace_period: 40s volumes: pgdata: diff --git a/docker/entrypoint-dev-scheduler.sh b/docker/entrypoint-dev-scheduler.sh new file mode 100755 index 00000000..f88f4ea1 --- /dev/null +++ b/docker/entrypoint-dev-scheduler.sh @@ -0,0 +1,11 @@ +#!/bin/sh +set -e + +if [ "${RUN_MIGRATIONS:-true}" = "true" ]; then + echo "[entrypoint] Applying database migrations …" + npx prisma migrate deploy + echo "[entrypoint] Migrations applied." +fi + +echo "[entrypoint] Starting scheduled job runner …" +exec pnpm scheduler:dev diff --git a/docker/entrypoint-dev-worker.sh b/docker/entrypoint-dev-worker.sh deleted file mode 100755 index 3fb8fcdb..00000000 --- a/docker/entrypoint-dev-worker.sh +++ /dev/null @@ -1,18 +0,0 @@ -#!/bin/sh -set -e - -# --------------------------------------------------------------------------- -# Learnault Worker — Development Container Entrypoint -# -# Environment variables: -# RUN_MIGRATIONS = "true" (default) → apply pending migrations on boot -# --------------------------------------------------------------------------- - -if [ "${RUN_MIGRATIONS:-true}" = "true" ]; then - echo "[entrypoint] Applying database migrations …" - npx prisma migrate deploy - echo "[entrypoint] Migrations applied." -fi - -echo "[entrypoint] Starting wallet-provisioning worker …" -exec pnpm worker:dev diff --git a/docker/entrypoint-scheduler.sh b/docker/entrypoint-scheduler.sh new file mode 100755 index 00000000..6b7de1e1 --- /dev/null +++ b/docker/entrypoint-scheduler.sh @@ -0,0 +1,11 @@ +#!/bin/sh +set -e + +if [ "${RUN_MIGRATIONS}" = "true" ]; then + echo "[entrypoint] Running database migrations …" + npx prisma migrate deploy + echo "[entrypoint] Migrations applied." +fi + +echo "[entrypoint] Starting scheduled job runner …" +exec node dist/workers/scheduler.worker.js diff --git a/docs/DATA_LIFECYCLE.md b/docs/DATA_LIFECYCLE.md index 92982178..6e136725 100644 --- a/docs/DATA_LIFECYCLE.md +++ b/docs/DATA_LIFECYCLE.md @@ -154,12 +154,15 @@ rather than removing it. | `OutboxEvent` | MUTABLE | 30d (`createdAt`) | Retain | No | | `JobAttempt` | MUTABLE | 30d (`createdAt`) | Cascade | No | | `RolledBackRecord` | IMMUTABLE | 30d (`createdAt`) | Retain | No | +| `QueueLease` | MUTABLE | Indefinite | Retain | No | | `WalletProvisioningJob` | MUTABLE | 90d (`updatedAt`) | Cascade | No | `EmailDelivery` and `NotificationLog` hold rendered message bodies, which is personal data — hence the short window and hard deletion on erasure. `DeviceToken` is deleted rather than archived: an archived push token would still -be a live address. +be a live address. `QueueLease` holds one long-lived row per recurring queue +drain — a queue name, the current lease token, and the holder id — so there is +no user data to erase and nothing to age out. --- diff --git a/docs/DEVELOPMENT_STACK.md b/docs/DEVELOPMENT_STACK.md index 2f293475..9fb734f3 100644 --- a/docs/DEVELOPMENT_STACK.md +++ b/docs/DEVELOPMENT_STACK.md @@ -1,6 +1,6 @@ # Local Development Stack (Docker Compose) -A reproducible local stack for the Learnault API: **API**, **wallet worker**, **PostgreSQL**, and **Redis** — started with one command. +A reproducible local stack for the Learnault API: **API**, **wallet worker**, **scheduler**, **PostgreSQL**, and **Redis** — started with one command. ## Prerequisites @@ -23,6 +23,7 @@ docker compose ps # learnault-dev-db Up ... (healthy) # learnault-dev-redis Up ... (healthy) # learnault-dev-worker Up ... (healthy) +# learnault-dev-scheduler Up ... ``` The API is available at `http://localhost:5000` (Swagger UI at `http://localhost:5000/api-docs`). @@ -37,6 +38,37 @@ The `api` service entrypoint (`docker/entrypoint-dev-api.sh`) waits for PostgreS The `worker` service runs `src/workers/wallet-provisioning.worker.ts`, which polls the idempotent wallet-provisioning outbox and generates Stellar keys through the dev in-memory KMS adapter. In production, swap the KMS adapter for a real one (e.g. AWS KMS) behind the same `KmsSecretStore` interface. +The `scheduler` service runs `src/workers/scheduler.worker.ts`. See below. + +## Scheduled job runner + +Every recurring queue drain is owned by the `scheduler` service, not by the request that enqueued the work — so a delivery whose `nextAttemptAt` falls due is retried on time even when the API is receiving no traffic, and request latency never includes queue-drain work. + +Registered queues: `email`, `notification`, `webhook`, `stellar-funding`, `data-export`, `account-lifecycle`. + +Each tick takes a row lease on `queue_leases` via `JobLeaseService.acquireQueueLease()` before draining, so extra replicas are safe: + +```bash +docker compose up -d --scale scheduler=2 +``` + +A replica that loses the race logs a skipped tick and moves on; a replica that crashes mid-drain has its lease expire, and the next tick reclaims the queue. + +| Variable | Default | Purpose | +| --- | --- | --- | +| `SCHEDULER_INTERVAL_MS` | `15000` | Base tick interval for every queue | +| `SCHEDULER__INTERVAL_MS` | — | Per-queue override, e.g. `SCHEDULER_WEBHOOK_INTERVAL_MS` | +| `SCHEDULER_LEASE_MS` | `60000` | Lease held per tick (floored at 2× the interval) | +| `SCHEDULER_QUEUES` | all | Comma list restricting which queues this replica runs | +| `SCHEDULER_DISABLED_QUEUES` | — | Comma list of queues to skip | +| `SCHEDULER_SHUTDOWN_TIMEOUT_MS` | `30000` | How long `SIGTERM` waits for in-flight ticks | +| `SCHEDULER_IN_PROCESS` | `false` | Opt-in: run the runner inside the API process for single-process deployments | +| `LIFECYCLE_SWEEP_INTERVAL_MS` | `0` | When `> 0`, overrides the `account-lifecycle` queue interval | + +Every tick emits a structured log line carrying per-queue `depth`, `due`, `lagMs` (age of the oldest due row), `durationMs`, and cumulative `attempts` / `failures` / `skipped`. + +`pnpm scheduler:verify` runs both evidence scenarios against the stack: a due-but-failed delivery drained with no inbound HTTP traffic, then a batch drained by two replicas with no row processed twice. + ## Health checks & readiness | Endpoint | Meaning | @@ -54,7 +86,7 @@ The API container only reports **healthy** after `/health/live` responds; `depen pnpm stack:up # docker compose up -d --build pnpm stack:down # stop the stack (keeps data volumes) pnpm stack:reset # stop + delete data volumes (project-scoped reset) -pnpm stack:logs # follow API + worker logs +pnpm stack:logs # follow API + worker + scheduler logs pnpm stack:validate # docker compose config --quiet pnpm stack:smoke # validate + start + probe health endpoints ``` @@ -65,9 +97,10 @@ pnpm stack:smoke # validate + start + probe health endpoints docker compose logs -f # all services docker compose logs -f api # API only docker compose logs worker # worker only +docker compose logs -f scheduler # scheduled job runner only ``` -Both services have `stop_grace_period: 30s`, matching the app's graceful-shutdown handler (`SHUTDOWN_TIMEOUT_MS`): `docker compose down` sends `SIGTERM`, the server drains HTTP connections and closes the Prisma pool before exiting. +`api` and `worker` have `stop_grace_period: 30s` and `scheduler` has `40s`, matching each process's graceful-shutdown handler (`SHUTDOWN_TIMEOUT_MS` / `SCHEDULER_SHUTDOWN_TIMEOUT_MS`): `docker compose down` sends `SIGTERM`, the server drains HTTP connections and closes the Prisma pool before exiting, and the scheduler stops scheduling, waits for in-flight ticks, and releases their queue leases so no queue is left parked. ## Data persistence & reset diff --git a/docs/domains/REQUEST_AND_EVENT_FLOWS.md b/docs/domains/REQUEST_AND_EVENT_FLOWS.md index 1f1075ed..0a96b34a 100644 --- a/docs/domains/REQUEST_AND_EVENT_FLOWS.md +++ b/docs/domains/REQUEST_AND_EVENT_FLOWS.md @@ -14,6 +14,7 @@ This document maps the key request flows and domain event propagation patterns a 6. [Referral Application Flow](#referral-application-flow) 7. [Withdrawal Flow](#withdrawal-flow) 8. [Notification Delivery Flow](#notification-delivery-flow) +9. [Subscribing to Domain Events](#subscribing-to-domain-events) --- @@ -508,3 +509,121 @@ Target state: Event-driven communication via domain events | Credentials | Blockchain Infra | Service Call | On-chain storage | | Credentials | Notifications | Domain Event | Credential notification | | All Domains | Shared Kernel | Direct Import | Config, errors, middleware, utils | + +--- + +## Subscribing to Domain Events + +A new domain event needs a registered handler, not a new worker process. The +outbox relay leases every pending event and dispatches it by `eventType` to the +handlers registered for that type. + +### 1. Declare the event schema + +Add the payload schema in `src/lib/transactions/event-schema.ts`. A handler for +an event type with no schema is rejected at startup. + +```ts +registry.register({ + version: 1, + eventType: 'ModuleCompleted', + validate: async (payload) => { + await z.object({ + completionId: z.string().uuid(), + userId: z.string().uuid(), + }).parseAsync(payload) + }, +}) +``` + +### 2. Emit the event in the same transaction as the domain write + +```ts +await prisma.$transaction(async (tx) => { + const completion = await tx.completion.create({ data: ... }) + + await createOutboxService(prisma).createEvent(tx, { + aggregateId: completion.id, + aggregateType: 'Completion', + eventType: 'ModuleCompleted', + eventVersion: 1, + payload: { completionId: completion.id, userId }, + source: 'api.module.complete', + }) + + return completion +}) +``` + +If the transaction rolls back the event disappears with it, so an event can +never describe a write that did not happen. + +### 3. Write the handler + +```ts +export class RewardOnModuleCompleted implements OutboxEventHandler { + readonly name = 'rewards.on-module-completed' + readonly eventType = 'ModuleCompleted' + readonly eventVersion = 1 + readonly maxAttempts = 5 + + async handle(ctx: OutboxEventHandlerContext) { + const payload = ctx.payload as ModuleCompletedPayload + await rewardService.grant(payload.userId, payload.completionId) + + return { idempotencyKey: `${ctx.eventId}:${this.name}` } + } +} +``` + +`name` must be unique across the whole registry — it becomes `JobAttempt.jobType`. + +**Handlers must be idempotent.** A handler can run more than once for the same +event: after a crash mid-lease, or after an operator replays a dead-lettered +event. Make the side effect an upsert, or guard it on a key derived from +`ctx.eventId`. + +Throwing from `handle()` schedules a retry with exponential backoff. Returning +normally completes the attempt. + +### 4. Register it + +Add the handler in `src/jobs/handler-registrations.ts`, and add its event type +to `EMITTED_EVENT_TYPES` if the application emits it: + +```ts +registry.register(new RewardOnModuleCompleted()) +``` + +`registerOutboxHandlers()` throws at startup on a duplicate handler name, on a +handler whose event type has no schema, and on an emitted event type with no +handler — so a missing subscription fails loudly instead of leaving rows PENDING +forever. + +### What the relay guarantees + +- One `JobAttempt` per (event, handler). Several handlers may subscribe to the + same event type and each is tracked separately. +- An event becomes `PUBLISHED` only once **every** handler for its type has + completed. One failing handler holds the event back without blocking others. +- A handler that exhausts `maxAttempts` dead-letters its own job and the event, + leaving every other event type unaffected. +- An event with no registered handler is dead-lettered immediately and logged at + error level, rather than sitting `PENDING` unnoticed. + +### Operating dead letters + +```bash +pnpm outbox:replay list # dead-lettered events and last error +pnpm outbox:replay replay ... # reset to PENDING for another pass +``` + +Replay resets the dead-lettered `JobAttempt` rows and returns the event to +`PENDING`; the relay picks it up on its next tick. Completed handlers are not +re-run, and idempotent handlers make a repeated run harmless. + +### Where it runs + +The relay is a queue on the scheduled job runner +(`src/workers/scheduler.worker.ts`), registered as `outbox-relay`. There is no +per-domain worker process: adding a domain event means adding a handler. diff --git a/learnault-api@0.1.0 b/learnault-api@0.1.0 new file mode 100644 index 00000000..e69de29b diff --git a/package.json b/package.json index c30a7c2b..436419cc 100644 --- a/package.json +++ b/package.json @@ -38,13 +38,17 @@ "seed:reset": "tsx prisma/seed.ts --reset", "db:seed": "npm run seed", "db:studio": "prisma studio", - "worker:dev": "tsx src/workers/wallet-provisioning.worker.ts", + "scheduler": "node dist/workers/scheduler.worker.js", + "scheduler:dev": "tsx src/workers/scheduler.worker.ts", + "outbox:replay": "tsx src/workers/outbox-replay.ts", + "relay:verify": "bash scripts/relay-verification.sh", "stack:validate": "docker compose config --quiet", "stack:up": "docker compose up -d --build", "stack:down": "docker compose down", "stack:reset": "docker compose down -v", - "stack:logs": "docker compose logs -f api worker", - "stack:smoke": "bash scripts/stack-smoke-test.sh" + "stack:logs": "docker compose logs -f api scheduler", + "stack:smoke": "bash scripts/stack-smoke-test.sh", + "scheduler:verify": "bash scripts/scheduler-verification.sh" }, "dependencies": { "@prisma/adapter-pg": "^7.4.2", diff --git a/prisma/migrations/20260829120000_scheduled_job_runner/migration.sql b/prisma/migrations/20260829120000_scheduled_job_runner/migration.sql new file mode 100644 index 00000000..aa07da9e --- /dev/null +++ b/prisma/migrations/20260829120000_scheduled_job_runner/migration.sql @@ -0,0 +1,15 @@ +CREATE TABLE "queue_leases" ( + "id" TEXT NOT NULL, + "queueName" TEXT NOT NULL, + "leaseToken" TEXT, + "leasedUntil" TIMESTAMP(3), + "owner" TEXT, + "lastTickAt" TIMESTAMP(3), + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "queue_leases_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "queue_leases_queueName_key" ON "queue_leases"("queueName"); +CREATE INDEX "queue_leases_leasedUntil_idx" ON "queue_leases"("leasedUntil"); diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 2477b20b..7d3490b8 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -841,3 +841,17 @@ model RolledBackRecord { @@index([createdAt]) // For periodic cleanup @@map("rolled_back_records") } + +model QueueLease { + id String @id @default(uuid()) + queueName String @unique + leaseToken String? + leasedUntil DateTime? + owner String? + lastTickAt DateTime? + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + @@index([leasedUntil]) + @@map("queue_leases") +} diff --git a/scripts/relay-verification.sh b/scripts/relay-verification.sh new file mode 100644 index 00000000..48e5a667 --- /dev/null +++ b/scripts/relay-verification.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash +set -euo pipefail + +command -v docker >/dev/null 2>&1 || export PATH="$PATH:/c/Program Files/Docker/Docker/resources/bin" + +PGUSER_="${POSTGRES_USER:-learnault}" +PGDB_="${POSTGRES_DB:-learnault_dev}" +PGPORT_="${POSTGRES_PORT:-5433}" + +export DATABASE_URL="${DATABASE_URL:-postgresql://${PGUSER_}:learnault@localhost:${PGPORT_}/${PGDB_}?schema=public}" +export NODE_ENV="${NODE_ENV:-production}" +export LOG_LEVEL="${LOG_LEVEL:-info}" + +echo "==> Starting PostgreSQL" +docker compose up -d db >/dev/null 2>&1 +until docker compose exec -T db pg_isready -U "$PGUSER_" -d "$PGDB_" >/dev/null 2>&1; do sleep 1; done + +HAS_OUTBOX="$(docker compose exec -T db psql -U "$PGUSER_" -d "$PGDB_" -qAt \ + -c "SELECT to_regclass('public.outbox_events');" | tr -d '\r')" + +if [ -z "$HAS_OUTBOX" ]; then + echo "==> Syncing schema" + npx prisma db push --accept-data-loss >/dev/null 2>&1 +fi + +exec npx tsx scripts/relay-verification.ts diff --git a/scripts/relay-verification.ts b/scripts/relay-verification.ts new file mode 100644 index 00000000..1b6632f2 --- /dev/null +++ b/scripts/relay-verification.ts @@ -0,0 +1,230 @@ +import 'dotenv/config' +import { randomUUID } from 'crypto' +import prisma from '../src/config/database' +import { registerOutboxHandlers } from '../src/jobs/handler-registrations' +import { createOutboxRelay } from '../src/workers/outbox-relay' + +const TAG = 'relay-evidence' + +const relay = createOutboxRelay({ prisma, handlers: registerOutboxHandlers({ prisma }) }) + +function banner(title: string): void { + console.log('') + console.log('='.repeat(64)) + console.log(` ${title}`) + console.log('='.repeat(64)) +} + +async function cleanup(): Promise { + const users = await prisma.user.findMany({ + where: { username: { startsWith: TAG } }, + select: { id: true }, + }) + const ids = users.map(u => u.id) + + await prisma.outboxEvent.deleteMany({ + where: { source: { in: [TAG, 'relay.user-created'] } }, + }) + if (ids.length > 0) { + await prisma.walletProvisioningJob.deleteMany({ + where: { wallet: { userId: { in: ids } } }, + }) + await prisma.wallet.deleteMany({ where: { userId: { in: ids } } }) + await prisma.user.deleteMany({ where: { id: { in: ids } } }) + } +} + +async function makeUser(suffix: string, withConsent: boolean): Promise { + const user = await prisma.user.create({ + data: { + email: `${TAG}-${suffix}@example.com`, + username: `${TAG}-${suffix}`, + password: 'x', + role: 'LEARNER', + isVerified: true, + status: 'ACTIVE', + }, + }) + + if (withConsent) { + await prisma.consentRecord.create({ + data: { userId: user.id, purpose: 'custodial_wallet', policyVersion: '1', status: 'granted', source: 'api' }, + }) + } + + return user.id +} + +async function emit(eventType: string, aggregateId: string, payload: unknown): Promise { + const event = await prisma.outboxEvent.create({ + data: { + id: randomUUID(), + aggregateId, + aggregateType: 'User', + eventType, + eventVersion: 1, + payload: JSON.stringify(payload), + status: 'PENDING', + source: TAG, + }, + }) + + return event.id +} + +async function fastForwardBackoff(): Promise { + await prisma.jobAttempt.updateMany({ + where: { status: 'PENDING' }, + data: { availableAt: new Date() }, + }) +} + +async function drain(ticks: number): Promise { + for (let i = 0; i < ticks; i++) { + await fastForwardBackoff() + const summary = await relay.runOnce() + if ( + summary.materialized === 0 && + summary.dispatched === 0 && + summary.failed === 0 && + summary.unhandled === 0 + ) { + break + } + } +} + +async function showEvents(): Promise { + const events = await prisma.outboxEvent.findMany({ + where: { source: { in: [TAG, 'relay.user-created'] } }, + orderBy: { createdAt: 'asc' }, + select: { + id: true, + eventType: true, + status: true, + jobAttempts: { select: { jobType: true, status: true, attempt: true } }, + }, + }) + + for (const event of events) { + const jobs = event.jobAttempts + .map(j => `${j.jobType}=${j.status}(attempt ${j.attempt})`) + .join(' ') + console.log(` ${event.eventType.padEnd(28)} ${event.status.padEnd(12)} ${jobs}`) + } +} + +async function main(): Promise { + await cleanup() + + banner('A. Dispatch by event type') + const okUser = await makeUser('ok', true) + await emit('UserCreated', okUser, { + userId: okUser, + email: `${TAG}-ok@example.com`, + role: 'LEARNER', + }) + await drain(6) + console.log('') + await showEvents() + console.log('') + console.log(' One UserCreated event fanned out to its handler, which emitted') + console.log(' WalletProvisioningRequested; the relay dispatched that to a different handler.') + + banner('B. Event type with no registered handler') + await emit('ModuleCompleted', okUser, { + completionId: randomUUID(), + userId: okUser, + moduleId: randomUUID(), + score: 90, + }) + const summary = await relay.runOnce() + console.log('') + console.log(` relay summary: ${JSON.stringify(summary)}`) + const unhandled = await prisma.outboxEvent.findFirst({ + where: { source: TAG, eventType: 'ModuleCompleted' }, + select: { status: true }, + }) + console.log(` ModuleCompleted -> ${unhandled?.status} (dead-lettered, not left PENDING)`) + + banner('C. Failing handler dead-letters, then replays cleanly') + const blockedUser = await makeUser('blocked', false) + const blockedEvent = await emit('UserCreated', blockedUser, { + userId: blockedUser, + email: `${TAG}-blocked@example.com`, + role: 'LEARNER', + }) + console.log('') + console.log(' User has no custodial-wallet consent, so the handler keeps failing.') + await drain(12) + + const dead = await prisma.outboxEvent.findUnique({ + where: { id: blockedEvent }, + select: { status: true, jobAttempts: { select: { status: true, attempt: true, lastError: true } } }, + }) + const failedJob = dead?.jobAttempts[0] + console.log(` event -> ${dead?.status}`) + console.log(` job -> ${failedJob?.status} after ${failedJob?.attempt} attempts`) + console.log(` error -> ${failedJob?.lastError?.split('\n')[0]}`) + + const before = await prisma.wallet.count({ where: { userId: blockedUser } }) + + console.log('') + console.log(' Operator grants the missing consent, then replays:') + await prisma.consentRecord.create({ + data: { userId: blockedUser, purpose: 'custodial_wallet', policyVersion: '1', status: 'granted', source: 'api' }, + }) + await relay.replayDeadLetter(blockedEvent) + await drain(6) + + const replayed = await prisma.outboxEvent.findUnique({ + where: { id: blockedEvent }, + select: { status: true, jobAttempts: { select: { jobType: true, status: true } } }, + }) + const after = await prisma.wallet.count({ where: { userId: blockedUser } }) + + console.log(` event -> ${replayed?.status}`) + console.log(` jobs -> ${replayed?.jobAttempts.map(j => `${j.jobType}=${j.status}`).join(' ')}`) + console.log(` wallets for user: ${before} before replay, ${after} after`) + + banner('D. Re-delivering an already-published event is idempotent') + await prisma.jobAttempt.updateMany({ + where: { outboxEventId: blockedEvent }, + data: { status: 'PENDING', attempt: 0, availableAt: new Date(), leaseToken: null, leasedUntil: null }, + }) + await prisma.outboxEvent.update({ where: { id: blockedEvent }, data: { status: 'PENDING' } }) + console.log('') + console.log(' Forcing the same event through the relay a second time:') + await drain(6) + + const redelivered = await prisma.outboxEvent.findUnique({ + where: { id: blockedEvent }, + select: { status: true }, + }) + const afterRedelivery = await prisma.wallet.count({ where: { userId: blockedUser } }) + const walletJobs = await prisma.walletProvisioningJob.count({ + where: { wallet: { userId: blockedUser } }, + }) + + console.log(` event -> ${redelivered?.status}`) + console.log( + ` wallets for user: ${after} before re-delivery, ${afterRedelivery} after; provisioning jobs: ${walletJobs}` + ) + console.log(' The handler ran again and produced no second wallet.') + + banner('Final state') + await showEvents() + + await cleanup() + console.log('') + console.log('Evidence complete.') +} + +main() + .then(() => prisma.$disconnect()) + .then(() => process.exit(0)) + .catch(async (error) => { + console.error('FAILED:', error) + await prisma.$disconnect().catch(() => undefined) + process.exit(1) + }) diff --git a/scripts/scheduler-verification.sh b/scripts/scheduler-verification.sh new file mode 100755 index 00000000..df436b8b --- /dev/null +++ b/scripts/scheduler-verification.sh @@ -0,0 +1,101 @@ +#!/usr/bin/env bash +set -uo pipefail + +command -v docker >/dev/null 2>&1 || export PATH="$PATH:/c/Program Files/Docker/Docker/resources/bin" + +PGUSER_="${POSTGRES_USER:-learnault}" +PGDB_="${POSTGRES_DB:-learnault_dev}" +PGPORT_="${POSTGRES_PORT:-5433}" +INTERVAL="${SCHEDULER_INTERVAL_MS:-4000}" +BATCH="${BATCH_SIZE:-40}" +LOGDIR="$(mktemp -d)" + +export DATABASE_URL="postgresql://${PGUSER_}:learnault@localhost:${PGPORT_}/${PGDB_}?schema=public" +export NODE_ENV=production +export LOG_LEVEL=debug +export SCHEDULER_INTERVAL_MS="$INTERVAL" +export SCHEDULER_LEASE_MS=10000 + +q() { docker compose exec -T db psql -U "$PGUSER_" -d "$PGDB_" -qAt -c "$1" | tr -d '\r'; } + +echo "==> Starting PostgreSQL" +docker compose up -d db >/dev/null 2>&1 +until docker compose exec -T db pg_isready -U "$PGUSER_" -d "$PGDB_" >/dev/null 2>&1; do sleep 1; done + +if [ -z "$(q "SELECT to_regclass('public.queue_leases');")" ]; then + echo "==> Syncing schema" + npx prisma db push --accept-data-loss >/dev/null 2>&1 +fi + +q "DELETE FROM email_deliveries WHERE type='SCHED_EVIDENCE';" >/dev/null +q "DELETE FROM users WHERE username='sched-evidence';" >/dev/null +USR="$(q "INSERT INTO users (id,email,username,password,role,\"isVerified\",status,\"createdAt\",\"updatedAt\") + VALUES (gen_random_uuid(),'sched-evidence@example.com','sched-evidence','x','LEARNER',true,'ACTIVE',now(),now()) + RETURNING id;" | head -1)" + +echo "" +echo "============================================================" +echo " SCENARIO 1 — idle instance: failed delivery retried on time" +echo "============================================================" + +q "INSERT INTO email_deliveries (id,\"userId\",\"to\",subject,body,type,status,error,\"attemptCount\",\"maxAttempts\",\"nextAttemptAt\",\"createdAt\",\"updatedAt\") + VALUES (gen_random_uuid(),'$USR','idle@example.com','retry me','

x

','SCHED_EVIDENCE','pending','previous attempt failed',1,5,now()-interval '1 minute',now(),now());" >/dev/null + +printf 'API on :5000 : ' +if curl -s -m 2 http://localhost:5000/health/live >/dev/null 2>&1; then echo "RUNNING (stop it for a clean result)"; else echo "not running — no HTTP traffic is possible"; fi +printf 'before : %s\n' "$(q "SELECT 'status='||status||' attemptCount='||\"attemptCount\"||' error='||error FROM email_deliveries WHERE type='SCHED_EVIDENCE';")" + +./node_modules/.bin/tsx src/workers/scheduler.worker.ts > "$LOGDIR/a.log" 2>&1 & +PID_A=$! +echo "scheduler : started (interval ${INTERVAL}ms, no API process)" +sleep 10 + +printf 'after : %s\n' "$(q "SELECT 'status='||status||' attemptCount='||\"attemptCount\" FROM email_deliveries WHERE type='SCHED_EVIDENCE';")" +echo "email tick :" +grep '"queue":"email"' "$LOGDIR/a.log" | head -1 + +kill -TERM "$PID_A" 2>/dev/null; wait "$PID_A" 2>/dev/null +echo "after SIGTERM : $(q "SELECT count(*)||'/6 leases released' FROM queue_leases WHERE \"leaseToken\" IS NULL;")" + +S1="$(q "SELECT status FROM email_deliveries WHERE type='SCHED_EVIDENCE' LIMIT 1;")" +[ "$S1" = "sent" ] && echo "RESULT : PASS — drained on schedule with zero inbound HTTP" \ + || echo "RESULT : FAIL — status=$S1" + +echo "" +echo "============================================================" +echo " SCENARIO 2 — two replicas: no duplicate processing" +echo "============================================================" + +q "DELETE FROM email_deliveries WHERE type='SCHED_EVIDENCE';" >/dev/null +q "INSERT INTO email_deliveries (id,\"userId\",\"to\",subject,body,type,status,\"attemptCount\",\"maxAttempts\",\"nextAttemptAt\",\"createdAt\",\"updatedAt\") + SELECT gen_random_uuid(),'$USR','r'||g||'@example.com','batch '||g,'

x

','SCHED_EVIDENCE','pending',0,5,now()-interval '1 minute',now(),now() + FROM generate_series(1,$BATCH) g;" >/dev/null + +echo "queued : $(q "SELECT count(*) FROM email_deliveries WHERE type='SCHED_EVIDENCE';") due rows" + +SCHEDULER_OWNER_ID=replica-A ./node_modules/.bin/tsx src/workers/scheduler.worker.ts > "$LOGDIR/1.log" 2>&1 & +P1=$! +SCHEDULER_OWNER_ID=replica-B ./node_modules/.bin/tsx src/workers/scheduler.worker.ts > "$LOGDIR/2.log" 2>&1 & +P2=$! +echo "replicas : replica-A and replica-B running concurrently" +sleep 14 +kill -TERM "$P1" "$P2" 2>/dev/null; wait "$P1" "$P2" 2>/dev/null + +echo "attemptCounts : $(q "SELECT string_agg('attemptCount='||\"attemptCount\"||' -> '||c||' rows',', ') FROM (SELECT \"attemptCount\",count(*) c FROM email_deliveries WHERE type='SCHED_EVIDENCE' GROUP BY 1 ORDER BY 1) t;")" +A_SKIP=$(grep -c 'tick skipped' "$LOGDIR/1.log" 2>/dev/null || true) +B_SKIP=$(grep -c 'tick skipped' "$LOGDIR/2.log" 2>/dev/null || true) +echo "lease races : replica-A skipped ${A_SKIP:-0}, replica-B skipped ${B_SKIP:-0}" + +DONE_=$(q "SELECT count(*) FROM email_deliveries WHERE type='SCHED_EVIDENCE' AND status='sent';") +DUPE_=$(q "SELECT count(*) FROM email_deliveries WHERE type='SCHED_EVIDENCE' AND \"attemptCount\">1;") + +q "DELETE FROM email_deliveries WHERE type='SCHED_EVIDENCE';" >/dev/null +q "DELETE FROM users WHERE username='sched-evidence';" >/dev/null +rm -rf "$LOGDIR" + +if [ "$DUPE_" = "0" ] && [ "$DONE_" = "$BATCH" ]; then + echo "RESULT : PASS — $DONE_/$BATCH processed exactly once, 0 duplicates" + exit 0 +fi +echo "RESULT : FAIL — processed=$DONE_ duplicates=$DUPE_" +exit 1 diff --git a/src/audit/classification.ts b/src/audit/classification.ts index 628c37d3..18e70d3a 100644 --- a/src/audit/classification.ts +++ b/src/audit/classification.ts @@ -476,6 +476,18 @@ const RULES: readonly LifecycleRule[] = [ audited: false, notes: 'Tombstone marking an event as unprocessable. Written once, then only read.', }, + { + model: 'QueueLease', + table: 'queue_leases', + recordClass: RecordClass.MUTABLE, + category: DataCategory.OPERATIONAL, + retentionDays: Retention.INDEFINITE, + retentionAnchor: null, + onErasure: ErasureAction.RETAIN, + audited: false, + notes: + 'One row per recurring queue drain, reused by every scheduler tick. Holds a queue name, lease token, and holder id — no user data, so nothing to erase and nothing to age out.', + }, ] const BY_MODEL: ReadonlyMap = new Map( diff --git a/src/config/env.ts b/src/config/env.ts index b789b9ac..19605caa 100644 --- a/src/config/env.ts +++ b/src/config/env.ts @@ -35,7 +35,7 @@ export const env = { // Account lifecycle configurations DELETION_COOLING_OFF_DAYS: parseInt(process.env.DELETION_COOLING_OFF_DAYS || '30', 10), EXPORT_TTL_DAYS: parseInt(process.env.EXPORT_TTL_DAYS || '7', 10), - LIFECYCLE_SWEEP_INTERVAL_MS: parseInt(process.env.LIFECYCLE_SWEEP_INTERVAL_MS || '0', 10), // 0 = disabled + LIFECYCLE_SWEEP_INTERVAL_MS: parseInt(process.env.LIFECYCLE_SWEEP_INTERVAL_MS || '0', 10), // Data lifecycle / audit configurations — see docs/DATA_LIFECYCLE.md // HMAC key for the source-IP hash on audit events. Unset in production means diff --git a/src/config/scheduler.ts b/src/config/scheduler.ts new file mode 100644 index 00000000..fac95d97 --- /dev/null +++ b/src/config/scheduler.ts @@ -0,0 +1,72 @@ +import { config } from 'dotenv' +import os from 'os' + +config() + +const DEFAULT_INTERVAL_MS = 15_000 +const DEFAULT_LEASE_MS = 60_000 +const DEFAULT_SHUTDOWN_TIMEOUT_MS = 30_000 + +function toInt(value: string | undefined, fallback: number): number { + const parsed = parseInt(value ?? '', 10) + + return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback +} + +function toBool(value: string | undefined, fallback = false): boolean { + if (value === undefined || value === '') return fallback + + return value === 'true' || value === '1' +} + +function toList(value: string | undefined): string[] { + return (value ?? '') + .split(',') + .map(entry => entry.trim()) + .filter(entry => entry.length > 0) +} + +function envKeyFor(queueName: string): string { + return `SCHEDULER_${queueName.replace(/[^a-zA-Z0-9]+/g, '_').toUpperCase()}_INTERVAL_MS` +} + +export const schedulerConfig = { + intervalMs: toInt(process.env.SCHEDULER_INTERVAL_MS, DEFAULT_INTERVAL_MS), + leaseMs: toInt(process.env.SCHEDULER_LEASE_MS, DEFAULT_LEASE_MS), + shutdownTimeoutMs: toInt( + process.env.SCHEDULER_SHUTDOWN_TIMEOUT_MS, + DEFAULT_SHUTDOWN_TIMEOUT_MS + ), + inProcess: toBool(process.env.SCHEDULER_IN_PROCESS, false), + only: toList(process.env.SCHEDULER_QUEUES), + disabled: toList(process.env.SCHEDULER_DISABLED_QUEUES), + ownerId: process.env.SCHEDULER_OWNER_ID || `${os.hostname()}:${process.pid}`, + + isEnabled(queueName: string): boolean { + if (this.disabled.includes(queueName)) return false + + return this.only.length === 0 || this.only.includes(queueName) + }, + + intervalFor(queueName: string): number { + const override = process.env[envKeyFor(queueName)] + if (override !== undefined) { + return toInt(override, this.intervalMs) + } + + if (queueName === 'account-lifecycle') { + const legacy = parseInt(process.env.LIFECYCLE_SWEEP_INTERVAL_MS ?? '', 10) + if (Number.isFinite(legacy) && legacy > 0) { + return legacy + } + } + + return this.intervalMs + }, + + leaseFor(queueName: string): number { + return Math.max(this.leaseMs, this.intervalFor(queueName) * 2) + }, +} + +export type SchedulerConfig = typeof schedulerConfig diff --git a/src/controllers/account.controller.ts b/src/controllers/account.controller.ts index 05b97584..4c9d4234 100644 --- a/src/controllers/account.controller.ts +++ b/src/controllers/account.controller.ts @@ -82,8 +82,6 @@ export class AccountController { try { const userId = req.user!.id - this.sweepInBackground() - const result = await dataExportService.requestExport(userId) if (result.kind === 'duplicate') { @@ -466,8 +464,6 @@ export class AccountController { */ async getDeletionStatus(req: Request, res: Response): Promise { try { - this.sweepInBackground() - const request = await accountLifecycleService.getLatestDeletionRequest(req.user!.id) if (!request) { @@ -619,12 +615,6 @@ export class AccountController { } } - private sweepInBackground(): void { - accountLifecycleService.sweep().catch(err => - logger.error('Lifecycle sweep error:', err) - ) - } - private generateToken(userId: string, role: string): string { return issueAccessToken({ id: userId, role }) } diff --git a/src/controllers/auth.controller.ts b/src/controllers/auth.controller.ts index f4e3ad7f..97ae1237 100644 --- a/src/controllers/auth.controller.ts +++ b/src/controllers/auth.controller.ts @@ -4,6 +4,7 @@ import prisma from '../config/database' import { loginSchema, registerSchema, verifyEmailSchema, resendVerificationSchema, forgotPasswordSchema, resetPasswordSchema, otpRequestSchema, otpVerifySchema, refreshTokenSchema } from '../schemas/auth.schema' import { UserRole } from '../types/user.types' import { emailService } from '../services/email.service' +import { createOutboxService } from '../lib/transactions/outbox.service' import { otpService, normalizePhone, OtpPurpose } from '../services/otp.service' import { refreshTokenService } from '../services/refresh-token.service' import { comparePassword, hashPassword, needsRehash } from '../utils/password' @@ -154,25 +155,41 @@ export class AuthController { const hashedPassword = await hashPassword(password) - const user = await prisma.user.create({ - data: { - email, - username, - password: hashedPassword, - role: (role as any) || UserRole.LEARNER, - } - }) - - // Issue verification token const { rawToken, tokenHash } = generateVerificationToken() const expiresAt = new Date(Date.now() + VERIFICATION_TOKEN_EXPIRY_MS) - await prisma.verificationToken.create({ - data: { - userId: user.id, - tokenHash, - expiresAt, - } + const user = await prisma.$transaction(async (tx) => { + const created = await tx.user.create({ + data: { + email, + username, + password: hashedPassword, + role: (role as any) || UserRole.LEARNER, + } + }) + + await tx.verificationToken.create({ + data: { + userId: created.id, + tokenHash, + expiresAt, + } + }) + + await createOutboxService(prisma).createEvent(tx, { + aggregateId: created.id, + aggregateType: 'User', + eventType: 'UserCreated', + eventVersion: 1, + payload: { + userId: created.id, + email: created.email, + role: created.role, + }, + source: 'api.auth.register', + }) + + return created }) // Queue verification email via outbox diff --git a/src/jobs/handler-registrations.ts b/src/jobs/handler-registrations.ts new file mode 100644 index 00000000..6bc439ed --- /dev/null +++ b/src/jobs/handler-registrations.ts @@ -0,0 +1,54 @@ +import type { PrismaClient } from '@prisma/client' +import defaultPrisma from '../config/database' +import { + getOutboxHandlerRegistry, + OutboxHandlerRegistry, +} from '../lib/transactions/handler-registry' +import { registerBuiltInSchemas } from '../lib/transactions/event-schema' +import { InMemoryEnvelopeKms } from '../services/kms/in-memory-envelope-kms' +import { SdkStellarKeypairGenerator } from '../services/stellar-keypair.adapter' +import { PrismaWalletProvisioningRepository } from '../services/wallet-provisioning.repository' +import { UserCreatedHandler } from './user-created.handler' +import { WalletProvisioningOutboxHandler } from './wallet-provisioning.handler' +import { WalletProvisioningRequestedHandler } from './wallet-provisioning-requested.handler' + +export const EMITTED_EVENT_TYPES = [ + { eventType: 'UserCreated', eventVersion: 1 }, + { eventType: 'WalletProvisioningRequested', eventVersion: 1 }, +] + +let schemasRegistered = false + +export interface RegisterHandlersOptions { + prisma?: PrismaClient + registry?: OutboxHandlerRegistry +} + +export function registerOutboxHandlers( + options: RegisterHandlersOptions = {} +): OutboxHandlerRegistry { + const prisma = options.prisma ?? defaultPrisma + const registry = options.registry ?? getOutboxHandlerRegistry() + + if (!schemasRegistered) { + registerBuiltInSchemas() + schemasRegistered = true + } + + const repository = new PrismaWalletProvisioningRepository(prisma) + + registry.register(new UserCreatedHandler(prisma, repository)) + registry.register( + new WalletProvisioningRequestedHandler( + new WalletProvisioningOutboxHandler( + repository, + new InMemoryEnvelopeKms(), + new SdkStellarKeypairGenerator() + ) + ) + ) + + registry.assertHandlersFor(EMITTED_EVENT_TYPES) + + return registry +} diff --git a/src/jobs/user-created.handler.ts b/src/jobs/user-created.handler.ts new file mode 100644 index 00000000..725432c0 --- /dev/null +++ b/src/jobs/user-created.handler.ts @@ -0,0 +1,76 @@ +import type { Prisma, PrismaClient } from '@prisma/client' +import type { + OutboxEventHandler, + OutboxEventHandlerContext, + OutboxEventHandlerResult, +} from '../lib/transactions/types' +import { createOutboxService, OutboxService } from '../lib/transactions/outbox.service' +import type { WalletProvisioningRepository } from '../services/wallet-provisioning.repository' + +export interface UserCreatedPayload { + userId: string + email: string + role: 'ADMIN' | 'LEARNER' | 'INSTRUCTOR' +} + +export interface UserCreatedHandlerOptions { + network?: string + outboxService?: OutboxService +} + +export class UserCreatedHandler implements OutboxEventHandler { + readonly name = 'wallet.reserve-on-user-created' + readonly eventType = 'UserCreated' + readonly eventVersion = 1 + readonly maxAttempts = 5 + + private readonly network: string + private readonly outbox: OutboxService + + constructor( + private readonly prisma: PrismaClient, + private readonly repository: WalletProvisioningRepository, + options: UserCreatedHandlerOptions = {}, + ) { + this.network = options.network ?? 'TESTNET' + this.outbox = options.outboxService ?? createOutboxService(prisma) + } + + async handle(context: OutboxEventHandlerContext): Promise { + const payload = context.payload as UserCreatedPayload + + const wallet = await this.repository.reserveEligibleWallet(payload.userId, this.network) + + const alreadyRequested = await this.prisma.outboxEvent.findFirst({ + where: { + aggregateId: wallet.id, + eventType: 'WalletProvisioningRequested', + causedBy: context.eventId, + }, + select: { id: true }, + }) + + if (!alreadyRequested) { + await this.prisma.$transaction(async (tx: Prisma.TransactionClient) => { + await this.outbox.createEvent(tx, { + aggregateId: wallet.id, + aggregateType: 'Wallet', + eventType: 'WalletProvisioningRequested', + eventVersion: 1, + payload: { + walletId: wallet.id, + userId: payload.userId, + network: this.network, + }, + source: 'relay.user-created', + causedBy: context.eventId, + }) + }) + } + + return { + idempotencyKey: `${context.eventId}:${this.name}`, + result: { walletId: wallet.id, requested: !alreadyRequested }, + } + } +} diff --git a/src/jobs/wallet-provisioning-requested.handler.ts b/src/jobs/wallet-provisioning-requested.handler.ts new file mode 100644 index 00000000..b1d2dbf5 --- /dev/null +++ b/src/jobs/wallet-provisioning-requested.handler.ts @@ -0,0 +1,37 @@ +import type { + OutboxEventHandler, + OutboxEventHandlerContext, + OutboxEventHandlerResult, +} from '../lib/transactions/types' +import type { WalletProvisioningOutboxHandler } from './wallet-provisioning.handler' + +export interface WalletProvisioningRequestedPayload { + walletId: string + userId: string + network: string +} + +export class WalletProvisioningRequestedHandler implements OutboxEventHandler { + readonly name = 'wallet.provision' + readonly eventType = 'WalletProvisioningRequested' + readonly eventVersion = 1 + readonly maxAttempts = 5 + + constructor(private readonly handler: WalletProvisioningOutboxHandler) {} + + async handle(context: OutboxEventHandlerContext): Promise { + const payload = context.payload as WalletProvisioningRequestedPayload + const result = await this.handler.handleWallet(payload.walletId) + + if (result.kind === 'retry-scheduled') { + throw new Error( + `Wallet ${payload.walletId} provisioning failed with ${result.failureCode}` + ) + } + + return { + idempotencyKey: `${context.eventId}:${this.name}`, + result, + } + } +} diff --git a/src/jobs/wallet-provisioning.handler.ts b/src/jobs/wallet-provisioning.handler.ts index 94a48204..3708b50d 100644 --- a/src/jobs/wallet-provisioning.handler.ts +++ b/src/jobs/wallet-provisioning.handler.ts @@ -48,6 +48,23 @@ export class WalletProvisioningOutboxHandler { async handleNext(): Promise { const claimed = await this.repository.claimNext(this.now(), this.leaseMs) + + return this.handleClaimed(claimed) + } + + async handleWallet(walletId: string): Promise { + const claimed = await this.repository.claimByWalletId( + walletId, + this.now(), + this.leaseMs, + ) + + return this.handleClaimed(claimed) + } + + private async handleClaimed( + claimed: ClaimedWalletProvisioningJob | null, + ): Promise { if (!claimed) return { kind: 'idle' } const leaseToken = claimed.job.leaseToken diff --git a/src/lib/transactions/event-schema.ts b/src/lib/transactions/event-schema.ts index 5d154768..79bdef8d 100644 --- a/src/lib/transactions/event-schema.ts +++ b/src/lib/transactions/event-schema.ts @@ -148,6 +148,19 @@ export function registerBuiltInSchemas(): void { }) // Wallet domain events + registry.register({ + version: 1, + eventType: 'WalletProvisioningRequested', + validate: async (payload) => { + const schema = z.object({ + walletId: z.string().uuid(), + userId: z.string().uuid(), + network: z.string(), + }) + await schema.parseAsync(payload) + }, + }) + registry.register({ version: 1, eventType: 'WalletProvisioned', diff --git a/src/lib/transactions/handler-registry.ts b/src/lib/transactions/handler-registry.ts new file mode 100644 index 00000000..7cb55d0f --- /dev/null +++ b/src/lib/transactions/handler-registry.ts @@ -0,0 +1,112 @@ +import { EventSchemaRegistry, getEventSchemaRegistry } from './event-schema.js' +import type { OutboxEventHandler } from './types.js' + +export class DuplicateHandlerError extends Error { + constructor(name: string) { + super(`Outbox handler "${name}" is already registered`) + this.name = 'DuplicateHandlerError' + } +} + +export class UnknownEventTypeError extends Error { + constructor(handlerName: string, eventType: string, eventVersion: number) { + super( + `Outbox handler "${handlerName}" targets ${eventType} v${eventVersion}, ` + + 'which has no registered event schema' + ) + this.name = 'UnknownEventTypeError' + } +} + +export class UnhandledEventTypeError extends Error { + constructor(missing: string[]) { + super( + `No outbox handler registered for emitted event type(s): ${missing.join(', ')}. ` + + 'Register a handler or stop emitting the event.' + ) + this.name = 'UnhandledEventTypeError' + } +} + +function keyOf(eventType: string, eventVersion: number): string { + return `${eventType}:v${eventVersion}` +} + +export class OutboxHandlerRegistry { + private readonly byEvent = new Map() + private readonly byName = new Map() + + constructor(private readonly schemas: EventSchemaRegistry = getEventSchemaRegistry()) {} + + register(handler: OutboxEventHandler): void { + if (this.byName.has(handler.name)) { + throw new DuplicateHandlerError(handler.name) + } + + if (!this.schemas.has(handler.eventType, handler.eventVersion)) { + throw new UnknownEventTypeError(handler.name, handler.eventType, handler.eventVersion) + } + + const key = keyOf(handler.eventType, handler.eventVersion) + const existing = this.byEvent.get(key) ?? [] + existing.push(handler) + + this.byEvent.set(key, existing) + this.byName.set(handler.name, handler) + } + + handlersFor(eventType: string, eventVersion: number): OutboxEventHandler[] { + return this.byEvent.get(keyOf(eventType, eventVersion)) ?? [] + } + + handlerByName(name: string): OutboxEventHandler | undefined { + return this.byName.get(name) + } + + registeredNames(): string[] { + return [...this.byName.keys()].sort() + } + + describe(): Array<{ eventType: string; eventVersion: number; handlers: string[] }> { + return [...this.byEvent.entries()] + .map(([key, handlers]) => { + const [eventType, version] = key.split(':v') + + return { + eventType, + eventVersion: Number(version), + handlers: handlers.map(h => h.name).sort(), + } + }) + .sort((a, b) => a.eventType.localeCompare(b.eventType)) + } + + assertHandlersFor(emitted: Array<{ eventType: string; eventVersion: number }>): void { + const missing = emitted + .filter(e => this.handlersFor(e.eventType, e.eventVersion).length === 0) + .map(e => keyOf(e.eventType, e.eventVersion)) + + if (missing.length > 0) { + throw new UnhandledEventTypeError(missing) + } + } + + clear(): void { + this.byEvent.clear() + this.byName.clear() + } +} + +let instance: OutboxHandlerRegistry | null = null + +export function getOutboxHandlerRegistry(): OutboxHandlerRegistry { + if (!instance) { + instance = new OutboxHandlerRegistry() + } + + return instance +} + +export function resetOutboxHandlerRegistry(): void { + instance = null +} diff --git a/src/lib/transactions/index.ts b/src/lib/transactions/index.ts index c9a8d462..c539bfcd 100644 --- a/src/lib/transactions/index.ts +++ b/src/lib/transactions/index.ts @@ -19,3 +19,4 @@ export * from './types.js' export * from './outbox.service.js' export * from './job-lease.service.js' export * from './event-schema.js' +export * from './handler-registry.js' diff --git a/src/lib/transactions/job-lease.service.ts b/src/lib/transactions/job-lease.service.ts index 2d0976f3..3f3fae6a 100644 --- a/src/lib/transactions/job-lease.service.ts +++ b/src/lib/transactions/job-lease.service.ts @@ -10,9 +10,18 @@ */ import { PrismaClient } from '@prisma/client' -import { LeaseJobOptions, LeaseJobResult, JobResult, JobAttempt } from './types.js' +import { + LeaseJobOptions, + LeaseJobResult, + JobResult, + JobAttempt, + AcquireQueueLeaseOptions, + QueueLeaseResult, +} from './types.js' import { randomUUID } from 'crypto' +const DEFAULT_QUEUE_LEASE_MS = 60000 + export class JobLeaseService { constructor(private prisma: PrismaClient) {} @@ -307,6 +316,76 @@ export class JobLeaseService { return result.count } + async acquireQueueLease( + options: AcquireQueueLeaseOptions + ): Promise { + const leaseMs = options.leaseMs ?? DEFAULT_QUEUE_LEASE_MS + const leaseToken = randomUUID() + const owner = options.owner ?? null + + const rows = await this.prisma.$queryRaw>` + INSERT INTO "queue_leases" + ("id", "queueName", "leaseToken", "leasedUntil", "owner", "createdAt", "updatedAt") + VALUES ( + ${randomUUID()}, + ${options.queueName}, + ${leaseToken}, + now() + (${String(leaseMs)}::text || ' milliseconds')::interval, + ${owner}, + now(), + now() + ) + ON CONFLICT ("queueName") DO UPDATE + SET "leaseToken" = EXCLUDED."leaseToken", + "leasedUntil" = EXCLUDED."leasedUntil", + "owner" = EXCLUDED."owner", + "updatedAt" = now() + WHERE "queue_leases"."leasedUntil" IS NULL + OR "queue_leases"."leasedUntil" <= now() + RETURNING "leasedUntil" + ` + + if (rows.length === 0) { + return null + } + + return { + queueName: options.queueName, + leaseToken, + leasedUntil: rows[0].leasedUntil, + } + } + + async renewQueueLease( + queueName: string, + leaseToken: string, + leaseMs: number = DEFAULT_QUEUE_LEASE_MS + ): Promise { + const updated = await this.prisma.$executeRaw` + UPDATE "queue_leases" + SET "leasedUntil" = now() + (${String(leaseMs)}::text || ' milliseconds')::interval, + "updatedAt" = now() + WHERE "queueName" = ${queueName} + AND "leaseToken" = ${leaseToken} + ` + + return updated > 0 + } + + async releaseQueueLease(queueName: string, leaseToken: string): Promise { + const updated = await this.prisma.$executeRaw` + UPDATE "queue_leases" + SET "leaseToken" = NULL, + "leasedUntil" = NULL, + "lastTickAt" = now(), + "updatedAt" = now() + WHERE "queueName" = ${queueName} + AND "leaseToken" = ${leaseToken} + ` + + return updated > 0 + } + /** * Get dead-letter jobs for manual inspection and recovery * diff --git a/src/lib/transactions/types.ts b/src/lib/transactions/types.ts index 25ec87e9..ded1a7a0 100644 --- a/src/lib/transactions/types.ts +++ b/src/lib/transactions/types.ts @@ -160,3 +160,40 @@ export interface LeaseJobResult { attempt: number payload: unknown } + +export interface OutboxEventHandlerContext { + eventId: string + eventType: string + eventVersion: number + aggregateId: string + aggregateType: string + payload: unknown + attempt: number +} + +export interface OutboxEventHandlerResult { + idempotencyKey?: string + result?: unknown +} + +export interface OutboxEventHandler { + name: string + eventType: string + eventVersion: number + maxAttempts?: number + backoffBaseMs?: number + backoffMultiplier?: number + handle(context: OutboxEventHandlerContext): Promise +} + +export interface AcquireQueueLeaseOptions { + queueName: string + leaseMs?: number + owner?: string +} + +export interface QueueLeaseResult { + queueName: string + leaseToken: string + leasedUntil: Date +} diff --git a/src/server.ts b/src/server.ts index 087c2cc9..0c582154 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1,9 +1,9 @@ import { Server } from 'http' import app from './app' -import { env } from './config/env' -import { accountLifecycleService } from './services/account-lifecycle.service' +import { schedulerConfig } from './config/scheduler' import logger from './utils/logger' import prisma from './config/database' +import { createScheduledJobRunner, ScheduledJobRunner } from './workers/scheduled-job-runner' const PORT = process.env.PORT || 5000 const SHUTDOWN_TIMEOUT_MS = parseInt(process.env.SHUTDOWN_TIMEOUT_MS || '30000', 10) @@ -13,18 +13,14 @@ const server: Server = app.listen(PORT, () => { }) let isShuttingDown = false -let lifecycleSweepInterval: NodeJS.Timeout | null = null - -// Periodic lifecycle sweep (export generation, deletion finalization, artifact -// purge). Disabled when LIFECYCLE_SWEEP_INTERVAL_MS is 0 — the sweep still runs -// lazily from account endpoints, and a dedicated worker can call sweep() directly. -if (env.LIFECYCLE_SWEEP_INTERVAL_MS > 0) { - lifecycleSweepInterval = setInterval(() => { - accountLifecycleService.sweep().catch(err => - logger.error('Scheduled lifecycle sweep error:', err) - ) - }, env.LIFECYCLE_SWEEP_INTERVAL_MS) - lifecycleSweepInterval.unref() +let scheduler: ScheduledJobRunner | null = null + +if (schedulerConfig.inProcess) { + scheduler = createScheduledJobRunner({ prisma }) + scheduler.start() + logger.info( + `In-process scheduler enabled for queues: ${scheduler.registeredQueues.join(', ') || 'none'}` + ) } /** @@ -63,10 +59,10 @@ async function gracefulShutdown(signal: string): Promise { }) // 2. Stop background jobs - if (lifecycleSweepInterval) { - logger.info('Stopping lifecycle sweep interval...') - clearInterval(lifecycleSweepInterval) - lifecycleSweepInterval = null + if (scheduler) { + logger.info('Stopping in-process scheduler...') + await scheduler.stop() + scheduler = null } // 3. Close database connections diff --git a/src/services/account-lifecycle.service.ts b/src/services/account-lifecycle.service.ts index 720a616f..7f12498a 100644 --- a/src/services/account-lifecycle.service.ts +++ b/src/services/account-lifecycle.service.ts @@ -284,11 +284,6 @@ export class AccountLifecycleService { } } - /** - * Runs all due lifecycle work. Invoked lazily from account endpoints, - * optionally on an interval from server.ts, and callable from a future - * dedicated worker (docker/entrypoint-worker.sh). - */ async sweep(): Promise { const results = await Promise.allSettled([ this.processDue(), diff --git a/src/services/data-export.service.ts b/src/services/data-export.service.ts index e0b5ad5e..eaabe5bd 100644 --- a/src/services/data-export.service.ts +++ b/src/services/data-export.service.ts @@ -59,10 +59,6 @@ export class DataExportService { await auditService.record({ userId, action: AuditAction.EXPORT_REQUESTED, metadata: { requestId: request.id } }) - this.processQueue().catch(err => - logger.error('[DataExportService] Queue processing error:', err) - ) - return { kind: 'created', request } } diff --git a/src/services/email.service.ts b/src/services/email.service.ts index c8f6f5e6..5973e0fa 100644 --- a/src/services/email.service.ts +++ b/src/services/email.service.ts @@ -39,10 +39,6 @@ export class EmailService { }, }) - this.processQueue().catch(err => - logger.error('[EmailService] Queue processing error:', err) - ) - return delivery as unknown as EmailDeliveryRecord } diff --git a/src/services/notification.service.ts b/src/services/notification.service.ts index a706bceb..cdc6cfc0 100644 --- a/src/services/notification.service.ts +++ b/src/services/notification.service.ts @@ -1,5 +1,5 @@ import prisma from '../config/database' -import * as admin from 'firebase-admin' +import admin from 'firebase-admin' // Local type definition to avoid @prisma/client import at test time interface NotificationLog { @@ -94,11 +94,6 @@ export class NotificationService { data: { userId, type, title, body, status: 'pending', nextAttemptAt: new Date() } }) - // Process asynchronously – same pattern as webhook service - this.processQueue().catch(err => - console.error('[NotificationService] Queue processing error:', err) - ) - return log as unknown as NotificationLog } diff --git a/src/services/stellar-funding.service.ts b/src/services/stellar-funding.service.ts index 78e166fd..95f46ae1 100644 --- a/src/services/stellar-funding.service.ts +++ b/src/services/stellar-funding.service.ts @@ -1,7 +1,6 @@ import prisma from '../config/database' import { stellarConfig } from '../config/stellar' import { StellarService, StellarServiceError } from './stellar.service' -import logger from '../utils/logger' interface StellarFundingRecord { id: string @@ -33,11 +32,7 @@ export class StellarFundingService { }) if (existing) { - this.processQueue().catch((err) => - logger.error('[StellarFundingService] Queue processing error:', err) - ) - -return existing as unknown as StellarFundingRecord + return existing as unknown as StellarFundingRecord } const funding = await prisma.stellarFunding.create({ @@ -49,10 +44,6 @@ return existing as unknown as StellarFundingRecord }, }) - this.processQueue().catch((err) => - logger.error('[StellarFundingService] Queue processing error:', err) - ) - return funding as unknown as StellarFundingRecord } diff --git a/src/services/wallet-provisioning.repository.ts b/src/services/wallet-provisioning.repository.ts index f99e1738..02155ba2 100644 --- a/src/services/wallet-provisioning.repository.ts +++ b/src/services/wallet-provisioning.repository.ts @@ -14,6 +14,11 @@ export interface WalletProvisioningRepository { reserveEligibleWallet(userId: string, network: string): Promise getByUserId(userId: string): Promise claimNext(now: Date, leaseMs: number): Promise + claimByWalletId( + walletId: string, + now: Date, + leaseMs: number, + ): Promise complete( walletId: string, leaseToken: string, @@ -98,54 +103,76 @@ export class PrismaWalletProvisioningRepository implements WalletProvisioningRep }) for (const candidate of candidates) { - const leaseToken = randomUUID() - const leasedUntil = new Date(now.getTime() + leaseMs) + const claimed = await this.claimCandidate(candidate.id, now, leaseMs) + if (claimed) return claimed + } - const claimed = await this.prisma.$transaction(async (tx) => { - const result = await tx.walletProvisioningJob.updateMany({ - where: { id: candidate.id, ...this.claimableWhere(now) }, - data: { - status: 'PROCESSING', - leaseToken, - leasedUntil, - attempts: { increment: 1 }, - }, - }) + return null + } + + async claimByWalletId( + walletId: string, + now: Date, + leaseMs: number, + ): Promise { + const candidate = await this.prisma.walletProvisioningJob.findFirst({ + where: { walletId, ...this.claimableWhere(now) }, + select: { id: true }, + }) - if (result.count !== 1) return null + if (!candidate) return null - const job = await tx.walletProvisioningJob.findUniqueOrThrow({ - where: { id: candidate.id }, - include: { wallet: true }, - }) + return this.claimCandidate(candidate.id, now, leaseMs) + } - await tx.wallet.update({ - where: { id: job.walletId }, - data: { - status: 'PROVISIONING', - attemptCount: { increment: 1 }, - statusChangedAt: now, - }, - }) + private async claimCandidate( + jobId: string, + now: Date, + leaseMs: number, + ): Promise { + const leaseToken = randomUUID() + const leasedUntil = new Date(now.getTime() + leaseMs) - await tx.auditLog.create({ - data: { - userId: job.wallet.userId, - action: 'WALLET_PROVISIONING_ATTEMPTED', - metadata: JSON.stringify({ walletId: job.walletId, attempt: job.attempts }), - }, - }) + return this.prisma.$transaction(async (tx) => { + const result = await tx.walletProvisioningJob.updateMany({ + where: { id: jobId, ...this.claimableWhere(now) }, + data: { + status: 'PROCESSING', + leaseToken, + leasedUntil, + attempts: { increment: 1 }, + }, + }) + + if (result.count !== 1) return null - return { - job: { ...job, leaseToken }, - wallet: { ...job.wallet, status: 'PROVISIONING' }, - } as unknown as ClaimedWalletProvisioningJob + const job = await tx.walletProvisioningJob.findUniqueOrThrow({ + where: { id: jobId }, + include: { wallet: true }, }) - if (claimed) return claimed - } + await tx.wallet.update({ + where: { id: job.walletId }, + data: { + status: 'PROVISIONING', + attemptCount: { increment: 1 }, + statusChangedAt: now, + }, + }) - return null + await tx.auditLog.create({ + data: { + userId: job.wallet.userId, + action: 'WALLET_PROVISIONING_ATTEMPTED', + metadata: JSON.stringify({ walletId: job.walletId, attempt: job.attempts }), + }, + }) + + return { + job: { ...job, leaseToken }, + wallet: { ...job.wallet, status: 'PROVISIONING' }, + } as unknown as ClaimedWalletProvisioningJob + }) } async complete( diff --git a/src/services/webhook.service.ts b/src/services/webhook.service.ts index fdfadf8f..f0c69b5b 100644 --- a/src/services/webhook.service.ts +++ b/src/services/webhook.service.ts @@ -56,9 +56,6 @@ export class WebhookService { }) }) ) - - // Process asynchronously - this.processQueue().catch(err => console.error('[Webhook] Queue processing error:', err)) } /** diff --git a/src/workers/outbox-relay.ts b/src/workers/outbox-relay.ts new file mode 100644 index 00000000..4be70387 --- /dev/null +++ b/src/workers/outbox-relay.ts @@ -0,0 +1,247 @@ +import type { PrismaClient } from '@prisma/client' +import defaultPrisma from '../config/database' +import { + EventSchemaRegistry, + getEventSchemaRegistry, +} from '../lib/transactions/event-schema' +import { + getOutboxHandlerRegistry, + OutboxHandlerRegistry, +} from '../lib/transactions/handler-registry' +import { createJobLeaseService, JobLeaseService } from '../lib/transactions/job-lease.service' +import logger from '../utils/logger' + +export interface OutboxRelayOptions { + prisma?: PrismaClient + handlers?: OutboxHandlerRegistry + schemas?: EventSchemaRegistry + leaseService?: JobLeaseService + batchSize?: number + leaseMs?: number + log?: Pick +} + +export interface RelayTickSummary { + materialized: number + dispatched: number + failed: number + unhandled: number +} + +interface EventRow { + id: string + eventType: string + eventVersion: number + aggregateId: string + aggregateType: string + payload: string +} + +export class OutboxRelay { + private readonly prisma: PrismaClient + private readonly handlers: OutboxHandlerRegistry + private readonly schemas: EventSchemaRegistry + private readonly leases: JobLeaseService + private readonly batchSize: number + private readonly leaseMs: number + private readonly log: NonNullable + + constructor(options: OutboxRelayOptions = {}) { + this.prisma = options.prisma ?? defaultPrisma + this.handlers = options.handlers ?? getOutboxHandlerRegistry() + this.schemas = options.schemas ?? getEventSchemaRegistry() + this.leases = options.leaseService ?? createJobLeaseService(this.prisma) + this.batchSize = options.batchSize ?? 50 + this.leaseMs = options.leaseMs ?? 30_000 + this.log = options.log ?? logger + } + + async runOnce(): Promise { + const { materialized, unhandled } = await this.materializePending() + const { dispatched, failed } = await this.dispatchLeasable() + + return { materialized, dispatched, failed, unhandled } + } + + async materializePending(): Promise<{ materialized: number; unhandled: number }> { + const pending = (await this.prisma.outboxEvent.findMany({ + where: { status: 'PENDING', jobAttempts: { none: {} } }, + orderBy: { createdAt: 'asc' }, + take: this.batchSize, + select: { + id: true, + eventType: true, + eventVersion: true, + aggregateId: true, + aggregateType: true, + payload: true, + }, + })) as EventRow[] + + let materialized = 0 + let unhandled = 0 + + for (const event of pending) { + const handlers = this.handlers.handlersFor(event.eventType, event.eventVersion) + + if (handlers.length === 0) { + unhandled += 1 + this.log.error( + `[relay] no handler registered for ${event.eventType} v${event.eventVersion}; dead-lettering event ${event.id}` + ) + await this.prisma.outboxEvent.update({ + where: { id: event.id }, + data: { status: 'DEAD_LETTER' }, + }) + continue + } + + await this.prisma.$transaction( + handlers.map(handler => + this.prisma.jobAttempt.create({ + data: { + outboxEventId: event.id, + jobType: handler.name, + jobName: `${event.eventType} -> ${handler.name}`, + status: 'PENDING', + attempt: 0, + maxAttempts: handler.maxAttempts ?? 3, + backoffBaseMs: handler.backoffBaseMs ?? 1000, + backoffMultiplier: handler.backoffMultiplier ?? 2.0, + availableAt: new Date(), + }, + }) + ) + ) + + materialized += 1 + this.log.info( + `[relay] materialized ${handlers.length} job(s) for ${event.eventType} v${event.eventVersion} (${event.id})` + ) + } + + return { materialized, unhandled } + } + + async dispatchLeasable(): Promise<{ dispatched: number; failed: number }> { + let dispatched = 0 + let failed = 0 + + for (const name of this.handlers.registeredNames()) { + let drained = 0 + + while (drained < this.batchSize) { + const outcome = await this.dispatchOne(name) + if (outcome === 'idle') break + + drained += 1 + if (outcome === 'ok') dispatched += 1 + else failed += 1 + } + } + + return { dispatched, failed } + } + + private async dispatchOne(jobType: string): Promise<'ok' | 'failed' | 'idle'> { + const handler = this.handlers.handlerByName(jobType) + if (!handler) return 'idle' + + const lease = await this.leases.leaseJob({ jobType, maxLeaseMs: this.leaseMs }) + if (!lease) return 'idle' + + const job = await this.prisma.jobAttempt.findUnique({ + where: { id: lease.jobId }, + include: { outboxEvent: true }, + }) + + if (!job?.outboxEvent) { + await this.leases.failJob(lease.jobId, lease.leaseToken, 'Outbox event missing for job') + + return 'failed' + } + + const event = job.outboxEvent as unknown as EventRow + + try { + await this.schemas.validate(event.eventType, event.eventVersion, lease.payload) + + const result = await handler.handle({ + eventId: event.id, + eventType: event.eventType, + eventVersion: event.eventVersion, + aggregateId: event.aggregateId, + aggregateType: event.aggregateType, + payload: lease.payload, + attempt: lease.attempt, + }) + + await this.leases.completeJob(lease.jobId, lease.leaseToken, { + success: true, + idempotencyKey: result?.idempotencyKey ?? `${event.id}:${handler.name}`, + result: result?.result, + }) + + this.log.info( + `[relay] dispatched ${event.eventType} v${event.eventVersion} -> ${handler.name} (${event.id})` + ) + + return 'ok' + } catch (error) { + const message = error instanceof Error ? error.message : String(error) + this.log.error( + `[relay] handler ${handler.name} failed for ${event.eventType} (${event.id}) attempt ${lease.attempt + 1}: ${message}` + ) + await this.leases.failJob(lease.jobId, lease.leaseToken, error as Error) + + return 'failed' + } + } + + async replayDeadLetter(eventId: string): Promise { + const jobs = await this.prisma.jobAttempt.findMany({ + where: { outboxEventId: eventId, status: 'DEAD_LETTER' }, + select: { id: true }, + }) + + for (const job of jobs) { + await this.leases.resetJobForRetry(job.id) + } + + await this.prisma.outboxEvent.update({ + where: { id: eventId }, + data: { status: 'PENDING' }, + }) + + this.log.info(`[relay] replayed event ${eventId} (${jobs.length} job(s) reset)`) + + return jobs.length + } + + async deadLetterEvents(limit = 100): Promise> { + const events = await this.prisma.outboxEvent.findMany({ + where: { status: 'DEAD_LETTER' }, + orderBy: { createdAt: 'asc' }, + take: limit, + select: { + id: true, + eventType: true, + jobAttempts: { + where: { status: 'DEAD_LETTER' }, + select: { lastError: true }, + take: 1, + }, + }, + }) + + return events.map(event => ({ + id: event.id, + eventType: event.eventType, + lastError: event.jobAttempts[0]?.lastError ?? null, + })) + } +} + +export function createOutboxRelay(options: OutboxRelayOptions = {}): OutboxRelay { + return new OutboxRelay(options) +} diff --git a/src/workers/outbox-replay.ts b/src/workers/outbox-replay.ts new file mode 100644 index 00000000..280091aa --- /dev/null +++ b/src/workers/outbox-replay.ts @@ -0,0 +1,48 @@ +import 'dotenv/config' +import prisma from '../config/database' +import { registerOutboxHandlers } from '../jobs/handler-registrations' +import logger from '../utils/logger' +import { createOutboxRelay } from './outbox-relay' + +async function main(): Promise { + const [command, ...args] = process.argv.slice(2) + const relay = createOutboxRelay({ prisma, handlers: registerOutboxHandlers({ prisma }) }) + + if (command === 'list') { + const events = await relay.deadLetterEvents() + + if (events.length === 0) { + logger.info('[replay] no dead-lettered events') + } else { + for (const event of events) { + logger.info(`[replay] ${event.id} ${event.eventType} ${event.lastError ?? ''}`) + } + } + + return + } + + if (command === 'replay') { + if (args.length === 0) { + throw new Error('usage: outbox:replay replay [...]') + } + + for (const eventId of args) { + const reset = await relay.replayDeadLetter(eventId) + logger.info(`[replay] ${eventId}: ${reset} job(s) reset to PENDING`) + } + + return + } + + throw new Error('usage: outbox:replay [eventId...]') +} + +main() + .then(() => prisma.$disconnect()) + .then(() => process.exit(0)) + .catch(async (error) => { + logger.error('[replay] failed:', error) + await prisma.$disconnect().catch(() => undefined) + process.exit(1) + }) diff --git a/src/workers/queue-metrics.ts b/src/workers/queue-metrics.ts new file mode 100644 index 00000000..1f775f86 --- /dev/null +++ b/src/workers/queue-metrics.ts @@ -0,0 +1,113 @@ +import logger from '../utils/logger' + +export type TickOutcome = 'ran' | 'skipped' | 'failed' + +export interface QueueDepthSnapshot { + depth: number + due: number + oldestDueAt: Date | null +} + +export interface QueueTickSample { + queue: string + outcome: TickOutcome + durationMs: number + depth: number + due: number + lagMs: number + error?: string +} + +export interface QueueMetricsSnapshot { + queue: string + attempts: number + failures: number + skipped: number + depth: number + due: number + lagMs: number + lastDurationMs: number + lastRunAt: string | null + lastError: string | null +} + +const emptyMetrics = (queue: string): QueueMetricsSnapshot => ({ + queue, + attempts: 0, + failures: 0, + skipped: 0, + depth: 0, + due: 0, + lagMs: 0, + lastDurationMs: 0, + lastRunAt: null, + lastError: null, +}) + +export class QueueMetricsRegistry { + private readonly queues = new Map() + + record(sample: QueueTickSample): QueueMetricsSnapshot { + const current = this.queues.get(sample.queue) ?? emptyMetrics(sample.queue) + + const next: QueueMetricsSnapshot = { + ...current, + depth: sample.depth, + due: sample.due, + lagMs: sample.lagMs, + lastDurationMs: sample.durationMs, + } + + if (sample.outcome === 'skipped') { + next.skipped = current.skipped + 1 + } else { + next.attempts = current.attempts + 1 + next.lastRunAt = new Date().toISOString() + } + + if (sample.outcome === 'failed') { + next.failures = current.failures + 1 + next.lastError = sample.error ?? 'unknown error' + } else if (sample.outcome === 'ran') { + next.lastError = null + } + + this.queues.set(sample.queue, next) + this.emit(sample, next) + + return next + } + + snapshot(): QueueMetricsSnapshot[] { + return [...this.queues.values()].sort((a, b) => a.queue.localeCompare(b.queue)) + } + + reset(): void { + this.queues.clear() + } + + private emit(sample: QueueTickSample, totals: QueueMetricsSnapshot): void { + const meta = { + queue: sample.queue, + outcome: sample.outcome, + depth: totals.depth, + due: totals.due, + lagMs: totals.lagMs, + durationMs: totals.lastDurationMs, + attempts: totals.attempts, + failures: totals.failures, + skipped: totals.skipped, + ...(sample.error ? { error: sample.error } : {}), + } + + if (sample.outcome === 'failed') { + logger.error('[scheduler] queue tick failed', meta) + } else if (sample.outcome === 'skipped') { + logger.debug('[scheduler] queue tick skipped (lease held elsewhere)', meta) + } else { + logger.info('[scheduler] queue tick', meta) + } + } +} + +export const queueMetrics = new QueueMetricsRegistry() diff --git a/src/workers/queue-registry.ts b/src/workers/queue-registry.ts new file mode 100644 index 00000000..f20852e7 --- /dev/null +++ b/src/workers/queue-registry.ts @@ -0,0 +1,145 @@ +import type { PrismaClient } from '@prisma/client' +import defaultPrisma from '../config/database' +import { accountLifecycleService } from '../services/account-lifecycle.service' +import { dataExportService } from '../services/data-export.service' +import { emailService } from '../services/email.service' +import { NotificationService } from '../services/notification.service' +import { stellarFundingService } from '../services/stellar-funding.service' +import { WebhookService } from '../services/webhook.service' +import { DeletionStatus, ExportStatus } from '../types/account.types' +import { registerOutboxHandlers } from '../jobs/handler-registrations' +import { createOutboxRelay, type OutboxRelay } from './outbox-relay' +import type { QueueDepthSnapshot } from './queue-metrics' + +export interface ScheduledQueue { + name: string + drain(): Promise + inspect(): Promise +} + +interface DelegateDepthOptions { + pending: Record + dueField?: string +} + +interface CountableDelegate { + count(args: { where: Record }): Promise + findFirst(args: { + where: Record + orderBy: Record + select: Record + }): Promise | null> +} + +function delegateDepth( + delegate: CountableDelegate, + options: DelegateDepthOptions +): () => Promise { + const dueField = options.dueField ?? 'nextAttemptAt' + + return async () => { + const now = new Date() + const dueWhere = { + ...options.pending, + OR: [{ [dueField]: null }, { [dueField]: { lte: now } }], + } + + const [depth, due, oldest] = await Promise.all([ + delegate.count({ where: options.pending }), + delegate.count({ where: dueWhere }), + delegate.findFirst({ + where: dueWhere, + orderBy: { [dueField]: 'asc' }, + select: { [dueField]: true }, + }), + ]) + + return { depth, due, oldestDueAt: (oldest?.[dueField] as Date | null) ?? null } + } +} + +export interface QueueRegistryDeps { + prisma?: PrismaClient + notificationService?: Pick + webhookService?: Pick + outboxRelay?: OutboxRelay +} + +export function createDefaultQueues(deps: QueueRegistryDeps = {}): ScheduledQueue[] { + const prisma = deps.prisma ?? defaultPrisma + const notificationService = deps.notificationService ?? new NotificationService() + const webhookService = deps.webhookService ?? new WebhookService() + const outboxRelay = deps.outboxRelay ?? createOutboxRelay({ + prisma, + handlers: registerOutboxHandlers({ prisma }), + }) + + const db = prisma as unknown as Record + + return [ + { + name: 'email', + drain: () => emailService.processQueue(), + inspect: delegateDepth(db.emailDelivery, { pending: { status: 'pending' } }), + }, + { + name: 'notification', + drain: () => notificationService.processQueue(), + inspect: delegateDepth(db.notificationLog, { pending: { status: 'pending' } }), + }, + { + name: 'webhook', + drain: () => webhookService.processQueue(), + inspect: delegateDepth(db.webhookDelivery, { pending: { status: 'pending' } }), + }, + { + name: 'stellar-funding', + drain: () => stellarFundingService.processQueue(), + inspect: delegateDepth(db.stellarFunding, { + pending: { status: { in: ['pending', 'submitted'] } }, + }), + }, + { + name: 'data-export', + drain: async () => { + await dataExportService.processQueue() + await dataExportService.purgeExpired() + }, + inspect: delegateDepth(db.dataExportRequest, { + pending: { status: ExportStatus.PENDING }, + }), + }, + { + name: 'outbox-relay', + drain: async () => { + await outboxRelay.runOnce() + }, + inspect: async () => { + const [depth, oldest] = await Promise.all([ + (prisma as unknown as { outboxEvent: CountableDelegate }).outboxEvent.count({ + where: { status: 'PENDING' }, + }), + (prisma as unknown as { outboxEvent: CountableDelegate }).outboxEvent.findFirst({ + where: { status: 'PENDING' }, + orderBy: { createdAt: 'asc' }, + select: { createdAt: true }, + }), + ]) + + return { + depth, + due: depth, + oldestDueAt: (oldest?.createdAt as Date | null) ?? null, + } + }, + }, + { + name: 'account-lifecycle', + drain: () => accountLifecycleService.processDue(), + inspect: delegateDepth(db.accountDeletionRequest, { + pending: { status: DeletionStatus.PENDING }, + dueField: 'scheduledFor', + }), + }, + ] +} diff --git a/src/workers/scheduled-job-runner.ts b/src/workers/scheduled-job-runner.ts new file mode 100644 index 00000000..9d1946b9 --- /dev/null +++ b/src/workers/scheduled-job-runner.ts @@ -0,0 +1,242 @@ +import type { PrismaClient } from '@prisma/client' +import defaultPrisma from '../config/database' +import { schedulerConfig, type SchedulerConfig } from '../config/scheduler' +import { createJobLeaseService, JobLeaseService } from '../lib/transactions/job-lease.service' +import logger from '../utils/logger' +import { queueMetrics, QueueMetricsRegistry, type QueueDepthSnapshot, type TickOutcome } from './queue-metrics' +import { createDefaultQueues, type ScheduledQueue } from './queue-registry' + +export type QueueLeaseApi = Pick< + JobLeaseService, + 'acquireQueueLease' | 'renewQueueLease' | 'releaseQueueLease' +> + +export interface ScheduledJobRunnerOptions { + queues: ScheduledQueue[] + leaseService: QueueLeaseApi + config?: SchedulerConfig + metrics?: QueueMetricsRegistry + log?: Pick +} + +const EMPTY_DEPTH: QueueDepthSnapshot = { depth: 0, due: 0, oldestDueAt: null } + +export class ScheduledJobRunner { + private readonly queues: ScheduledQueue[] + private readonly leaseService: QueueLeaseApi + private readonly config: SchedulerConfig + private readonly metrics: QueueMetricsRegistry + private readonly log: NonNullable + + private readonly timers = new Map() + private readonly inFlight = new Map>() + private started = false + private stopping = false + + constructor(options: ScheduledJobRunnerOptions) { + this.config = options.config ?? schedulerConfig + this.queues = options.queues.filter(queue => this.config.isEnabled(queue.name)) + this.leaseService = options.leaseService + this.metrics = options.metrics ?? queueMetrics + this.log = options.log ?? logger + } + + get registeredQueues(): string[] { + return this.queues.map(queue => queue.name) + } + + start(): void { + if (this.started) return + this.started = true + this.stopping = false + + if (this.queues.length === 0) { + this.log.warn('[scheduler] no queues enabled; runner idle') + + return + } + + for (const queue of this.queues) { + this.log.info( + `[scheduler] registered queue "${queue.name}" (every ${this.config.intervalFor(queue.name)}ms)` + ) + this.schedule(queue, 0) + } + } + + async runOnce(): Promise> { + const outcomes: Record = {} + + for (const queue of this.queues) { + outcomes[queue.name] = await this.runQueue(queue) + } + + return outcomes + } + + async stop(): Promise { + if (!this.started || this.stopping) return + this.stopping = true + + for (const timer of this.timers.values()) { + clearTimeout(timer) + } + this.timers.clear() + + const pending = [...this.inFlight.values()] + if (pending.length > 0) { + this.log.info(`[scheduler] draining ${pending.length} in-flight tick(s)`) + const drained = await this.withDeadline( + Promise.allSettled(pending), + this.config.shutdownTimeoutMs + ) + + if (!drained) { + this.log.warn( + `[scheduler] shutdown deadline (${this.config.shutdownTimeoutMs}ms) reached; ` + + 'remaining leases expire on their own' + ) + } + } + + this.started = false + this.log.info('[scheduler] stopped') + } + + private schedule(queue: ScheduledQueue, delayMs: number): void { + const timer = setTimeout(() => { + void this.tick(queue) + }, delayMs) + + this.timers.set(queue.name, timer) + } + + private async tick(queue: ScheduledQueue): Promise { + this.timers.delete(queue.name) + if (this.stopping) return + + const pending = this.runQueue(queue) + this.inFlight.set(queue.name, pending) + + try { + await pending + } finally { + this.inFlight.delete(queue.name) + } + + if (!this.stopping) { + this.schedule(queue, this.config.intervalFor(queue.name)) + } + } + + private async runQueue(queue: ScheduledQueue): Promise { + const leaseMs = this.config.leaseFor(queue.name) + const before = await this.inspect(queue) + const lagMs = before.oldestDueAt + ? Math.max(0, Date.now() - before.oldestDueAt.getTime()) + : 0 + + let lease + try { + lease = await this.leaseService.acquireQueueLease({ + queueName: queue.name, + leaseMs, + owner: this.config.ownerId, + }) + } catch (error) { + return this.record(queue, 'failed', 0, before, lagMs, error) + } + + if (!lease) { + return this.record(queue, 'skipped', 0, before, lagMs) + } + + const heartbeat = setInterval(() => { + void this.leaseService + .renewQueueLease(queue.name, lease.leaseToken, leaseMs) + .catch(error => + this.log.warn(`[scheduler] failed to renew lease for "${queue.name}"`, error) + ) + }, Math.max(1_000, Math.floor(leaseMs / 2))) + + const startedAt = Date.now() + + try { + await queue.drain() + + return this.record(queue, 'ran', Date.now() - startedAt, before, lagMs) + } catch (error) { + return this.record(queue, 'failed', Date.now() - startedAt, before, lagMs, error) + } finally { + clearInterval(heartbeat) + await this.leaseService + .releaseQueueLease(queue.name, lease.leaseToken) + .catch(error => + this.log.warn(`[scheduler] failed to release lease for "${queue.name}"`, error) + ) + } + } + + private async inspect(queue: ScheduledQueue): Promise { + try { + return await queue.inspect() + } catch (error) { + this.log.warn(`[scheduler] depth probe failed for "${queue.name}"`, error) + + return EMPTY_DEPTH + } + } + + private record( + queue: ScheduledQueue, + outcome: TickOutcome, + durationMs: number, + depth: QueueDepthSnapshot, + lagMs: number, + error?: unknown + ): TickOutcome { + this.metrics.record({ + queue: queue.name, + outcome, + durationMs, + depth: depth.depth, + due: depth.due, + lagMs, + error: error === undefined ? undefined : toMessage(error), + }) + + return outcome + } + + private async withDeadline(work: Promise, timeoutMs: number): Promise { + let timer: NodeJS.Timeout | undefined + + const deadline = new Promise(resolve => { + timer = setTimeout(() => resolve(false), timeoutMs) + }) + + try { + return await Promise.race([work.then(() => true), deadline]) + } finally { + if (timer) clearTimeout(timer) + } + } +} + +function toMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error) +} + +export function createScheduledJobRunner( + overrides: Partial & { prisma?: PrismaClient } = {} +): ScheduledJobRunner { + const prisma = overrides.prisma ?? defaultPrisma + + return new ScheduledJobRunner({ + queues: overrides.queues ?? createDefaultQueues({ prisma }), + leaseService: overrides.leaseService ?? createJobLeaseService(prisma), + config: overrides.config, + metrics: overrides.metrics, + log: overrides.log, + }) +} diff --git a/src/workers/scheduler.worker.ts b/src/workers/scheduler.worker.ts new file mode 100644 index 00000000..47660c5d --- /dev/null +++ b/src/workers/scheduler.worker.ts @@ -0,0 +1,62 @@ +import 'dotenv/config' +import prisma from '../config/database' +import { schedulerConfig } from '../config/scheduler' +import logger from '../utils/logger' +import { createScheduledJobRunner } from './scheduled-job-runner' + +const runner = createScheduledJobRunner({ prisma }) + +let isShuttingDown = false + +async function gracefulShutdown(signal: string): Promise { + if (isShuttingDown) { + logger.warn('[scheduler] shutdown already in progress, ignoring additional signal') + + return + } + + isShuttingDown = true + logger.info(`[scheduler] received ${signal}, starting graceful shutdown...`) + + const forceExit = setTimeout(() => { + logger.error('[scheduler] shutdown deadline exceeded, forcing exit') + process.exit(1) + }, schedulerConfig.shutdownTimeoutMs + 5_000) + + try { + await runner.stop() + await prisma.$disconnect() + clearTimeout(forceExit) + logger.info('[scheduler] graceful shutdown completed') + process.exit(0) + } catch (error) { + logger.error('[scheduler] error during graceful shutdown:', error) + clearTimeout(forceExit) + process.exit(1) + } +} + +process.on('SIGTERM', () => void gracefulShutdown('SIGTERM')) +process.on('SIGINT', () => void gracefulShutdown('SIGINT')) + +process.on('uncaughtException', (error: Error) => { + logger.error('[scheduler] uncaught exception:', error) + void gracefulShutdown('uncaughtException') +}) + +process.on('unhandledRejection', (reason: unknown) => { + logger.error('[scheduler] unhandled rejection:', reason) + void gracefulShutdown('unhandledRejection') +}) + +logger.info( + `[scheduler] starting runner ${schedulerConfig.ownerId} ` + + `(base interval ${schedulerConfig.intervalMs}ms, lease ${schedulerConfig.leaseMs}ms)` +) + +runner.start() + +if (runner.registeredQueues.length === 0) { + logger.error('[scheduler] no queues registered; check SCHEDULER_QUEUES / SCHEDULER_DISABLED_QUEUES') + process.exit(1) +} diff --git a/src/workers/wallet-provisioning.worker.ts b/src/workers/wallet-provisioning.worker.ts deleted file mode 100644 index 77fd5804..00000000 --- a/src/workers/wallet-provisioning.worker.ts +++ /dev/null @@ -1,58 +0,0 @@ -import 'dotenv/config' -import { WalletProvisioningOutboxHandler } from '../jobs/wallet-provisioning.handler' -import { InMemoryEnvelopeKms } from '../services/kms/in-memory-envelope-kms' -import { SdkStellarKeypairGenerator } from '../services/stellar-keypair.adapter' -import { PrismaWalletProvisioningRepository } from '../services/wallet-provisioning.repository' -import prisma from '../config/database' - -const POLL_INTERVAL_MS = parseInt(process.env.WORKER_POLL_INTERVAL_MS ?? '5000', 10) - -// Development worker: drains the idempotent wallet-provisioning outbox using -// the in-memory KMS adapter. Production should swap in a real KMS adapter -// (e.g. AWS KMS) — the handler only depends on the KmsSecretStore interface. -const kms = new InMemoryEnvelopeKms() -const repository = new PrismaWalletProvisioningRepository(prisma) -const keypairGenerator = new SdkStellarKeypairGenerator() -const handler = new WalletProvisioningOutboxHandler(repository, kms, keypairGenerator) - -let isShuttingDown = false - -async function drainOnce(): Promise { - const result = await handler.handleNext() - if (result.kind !== 'idle') { - console.log(`[worker] ${result.kind}:`, JSON.stringify(result)) - } -} - -async function runLoop(): Promise { - console.log(`[worker] Wallet provisioning worker started (poll every ${POLL_INTERVAL_MS}ms)`) - while (!isShuttingDown) { - try { - await drainOnce() - } catch (error) { - console.error('[worker] Unhandled error while draining outbox:', error) - } - await new Promise((resolve) => setTimeout(resolve, POLL_INTERVAL_MS)) - } -} - -function gracefulShutdown(signal: string): void { - if (isShuttingDown) return - isShuttingDown = true - console.log(`[worker] Received ${signal}, shutting down...`) - prisma - .$disconnect() - .then(() => process.exit(0)) - .catch((error) => { - console.error('[worker] Error disconnecting Prisma:', error) - process.exit(1) - }) -} - -process.on('SIGTERM', () => gracefulShutdown('SIGTERM')) -process.on('SIGINT', () => gracefulShutdown('SIGINT')) - -runLoop().catch((error) => { - console.error('[worker] Fatal error:', error) - process.exit(1) -}) diff --git a/tests/auth.controller.test.ts b/tests/auth.controller.test.ts index 492fb813..bd494352 100644 --- a/tests/auth.controller.test.ts +++ b/tests/auth.controller.test.ts @@ -10,8 +10,8 @@ import { refreshTokenService } from '../src/services/refresh-token.service' const mockTokenHash = 'abc123def456hash' const mockRawToken = 'aaabbbcccddd00112233445566778899aabbccddeeff00112233445566778899' -vi.mock('../src/config/database', () => ({ - default: { +const { mockDb } = vi.hoisted(() => { + const mockDb: any = { user: { findFirst: vi.fn(), findUnique: vi.fn(), @@ -36,9 +36,19 @@ vi.mock('../src/config/database', () => ({ auditLog: { create: vi.fn(), }, - $transaction: vi.fn((args: any[]) => Promise.all(args)), - }, -})) + outboxEvent: { + create: vi.fn(), + }, + } + + mockDb.$transaction = vi.fn((arg: any) => + typeof arg === 'function' ? arg(mockDb) : Promise.all(arg) + ) + + return { mockDb } +}) + +vi.mock('../src/config/database', () => ({ default: mockDb })) vi.mock('bcryptjs', () => ({ default: { @@ -281,8 +291,8 @@ describe('AuthController', () => { status: 'USED', }) ;(prisma.user.update as any).mockResolvedValue({ id: '1', isVerified: true }) - ;(prisma.$transaction as any).mockImplementation( - async (args: any[]) => await Promise.all(args) + ;(prisma.$transaction as any).mockImplementation(async (arg: any) => + typeof arg === 'function' ? arg(prisma) : await Promise.all(arg) ) await authController.verifyEmail( diff --git a/tests/email.service.test.ts b/tests/email.service.test.ts index 712a2148..e4a27b04 100644 --- a/tests/email.service.test.ts +++ b/tests/email.service.test.ts @@ -29,7 +29,7 @@ describe('EmailService', () => { }) describe('queueEmail', () => { - it('should create an email delivery record and trigger queue processing', async () => { + it('should create a pending email delivery record without draining the queue', async () => { const mockDelivery = { id: 'del1', userId: 'user1', @@ -63,6 +63,7 @@ describe('EmailService', () => { }) ) expect(result).toEqual(mockDelivery) + expect(prisma.emailDelivery.findMany).not.toHaveBeenCalled() }) }) diff --git a/tests/lib/handler-registry.test.ts b/tests/lib/handler-registry.test.ts new file mode 100644 index 00000000..b9405263 --- /dev/null +++ b/tests/lib/handler-registry.test.ts @@ -0,0 +1,98 @@ +import { describe, it, expect, beforeEach } from 'vitest' +import { EventSchemaRegistry } from '../../src/lib/transactions/event-schema' +import { + DuplicateHandlerError, + OutboxHandlerRegistry, + UnhandledEventTypeError, + UnknownEventTypeError, +} from '../../src/lib/transactions/handler-registry' +import type { OutboxEventHandler } from '../../src/lib/transactions/types' + +function handler( + name: string, + eventType = 'UserCreated', + eventVersion = 1 +): OutboxEventHandler { + return { name, eventType, eventVersion, handle: async () => undefined } +} + +describe('OutboxHandlerRegistry', () => { + let schemas: EventSchemaRegistry + let registry: OutboxHandlerRegistry + + beforeEach(() => { + schemas = new EventSchemaRegistry() + schemas.register({ eventType: 'UserCreated', version: 1, validate: () => undefined }) + schemas.register({ eventType: 'WalletProvisioned', version: 1, validate: () => undefined }) + registry = new OutboxHandlerRegistry(schemas) + }) + + it('resolves handlers by event type and version', () => { + const a = handler('a') + registry.register(a) + + expect(registry.handlersFor('UserCreated', 1)).toEqual([a]) + expect(registry.handlersFor('UserCreated', 2)).toEqual([]) + expect(registry.handlerByName('a')).toBe(a) + }) + + it('supports several handlers for one event type', () => { + registry.register(handler('a')) + registry.register(handler('b')) + + expect(registry.handlersFor('UserCreated', 1).map(h => h.name)).toEqual(['a', 'b']) + expect(registry.registeredNames()).toEqual(['a', 'b']) + }) + + it('rejects a duplicate handler name', () => { + registry.register(handler('a')) + + expect(() => registry.register(handler('a', 'WalletProvisioned'))).toThrow( + DuplicateHandlerError + ) + }) + + it('rejects a handler for an event type with no registered schema', () => { + expect(() => registry.register(handler('a', 'NeverDeclared'))).toThrow( + UnknownEventTypeError + ) + }) + + it('rejects a handler for a version the schema registry does not know', () => { + expect(() => registry.register(handler('a', 'UserCreated', 7))).toThrow( + UnknownEventTypeError + ) + }) + + it('fails loudly when an emitted event type has no handler', () => { + registry.register(handler('a')) + + expect(() => + registry.assertHandlersFor([ + { eventType: 'UserCreated', eventVersion: 1 }, + { eventType: 'WalletProvisioned', eventVersion: 1 }, + ]) + ).toThrow(UnhandledEventTypeError) + + expect(() => + registry.assertHandlersFor([{ eventType: 'UserCreated', eventVersion: 1 }]) + ).not.toThrow() + }) + + it('names the missing event types in the startup error', () => { + expect(() => + registry.assertHandlersFor([{ eventType: 'WalletProvisioned', eventVersion: 1 }]) + ).toThrow(/WalletProvisioned:v1/) + }) + + it('describes what is registered for operator inspection', () => { + registry.register(handler('b')) + registry.register(handler('a')) + registry.register(handler('c', 'WalletProvisioned')) + + expect(registry.describe()).toEqual([ + { eventType: 'UserCreated', eventVersion: 1, handlers: ['a', 'b'] }, + { eventType: 'WalletProvisioned', eventVersion: 1, handlers: ['c'] }, + ]) + }) +}) diff --git a/tests/lib/job-lease-queue.test.ts b/tests/lib/job-lease-queue.test.ts new file mode 100644 index 00000000..60f90207 --- /dev/null +++ b/tests/lib/job-lease-queue.test.ts @@ -0,0 +1,106 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' +import { JobLeaseService } from '../../src/lib/transactions/job-lease.service' + +function sqlOf(call: unknown[]): string { + const [strings] = call as [TemplateStringsArray] + + return strings.join('?') +} + +describe('JobLeaseService queue leases', () => { + let queryRaw: ReturnType + let executeRaw: ReturnType + let service: JobLeaseService + + beforeEach(() => { + queryRaw = vi.fn() + executeRaw = vi.fn() + service = new JobLeaseService({ $queryRaw: queryRaw, $executeRaw: executeRaw } as any) + }) + + describe('acquireQueueLease', () => { + it('returns a token when the conditional upsert claims the queue', async () => { + const leasedUntil = new Date(Date.now() + 60_000) + queryRaw.mockResolvedValue([{ leasedUntil }]) + + const lease = await service.acquireQueueLease({ + queueName: 'email', + leaseMs: 60_000, + owner: 'scheduler@host:1', + }) + + expect(lease).not.toBeNull() + expect(lease!.queueName).toBe('email') + expect(lease!.leasedUntil).toBe(leasedUntil) + expect(lease!.leaseToken).toMatch(/^[0-9a-f-]{36}$/) + }) + + it('returns null when another holder still owns the lease', async () => { + queryRaw.mockResolvedValue([]) + + const lease = await service.acquireQueueLease({ queueName: 'email' }) + + expect(lease).toBeNull() + }) + + it('gates the upsert on the database clock, not the caller clock', async () => { + queryRaw.mockResolvedValue([{ leasedUntil: new Date() }]) + + await service.acquireQueueLease({ queueName: 'webhook', leaseMs: 30_000 }) + + const sql = sqlOf(queryRaw.mock.calls[0]) + expect(sql).toContain('ON CONFLICT ("queueName") DO UPDATE') + expect(sql).toContain('"queue_leases"."leasedUntil" IS NULL') + expect(sql).toContain('"queue_leases"."leasedUntil" <= now()') + + const values = queryRaw.mock.calls[0].slice(1) + expect(values).toContain('webhook') + expect(values).toContain('30000') + }) + + it('issues a distinct token per acquisition', async () => { + queryRaw.mockResolvedValue([{ leasedUntil: new Date() }]) + + const first = await service.acquireQueueLease({ queueName: 'email' }) + const second = await service.acquireQueueLease({ queueName: 'email' }) + + expect(first!.leaseToken).not.toBe(second!.leaseToken) + }) + }) + + describe('renewQueueLease', () => { + it('reports success only when the row still carries this token', async () => { + executeRaw.mockResolvedValueOnce(1) + await expect(service.renewQueueLease('email', 'token-a', 5_000)).resolves.toBe(true) + + executeRaw.mockResolvedValueOnce(0) + await expect(service.renewQueueLease('email', 'stale-token')).resolves.toBe(false) + + const values = executeRaw.mock.calls[0].slice(1) + expect(values).toContain('email') + expect(values).toContain('token-a') + }) + }) + + describe('releaseQueueLease', () => { + it('clears the lease scoped to the holding token', async () => { + executeRaw.mockResolvedValue(1) + + await expect(service.releaseQueueLease('data-export', 'token-a')).resolves.toBe(true) + + const sql = sqlOf(executeRaw.mock.calls[0]) + expect(sql).toContain('"leaseToken" = NULL') + expect(sql).toContain('"lastTickAt" = now()') + expect(sql).toContain('AND "leaseToken" =') + + const values = executeRaw.mock.calls[0].slice(1) + expect(values).toEqual(['data-export', 'token-a']) + }) + + it('does not release a lease a successor now holds', async () => { + executeRaw.mockResolvedValue(0) + + await expect(service.releaseQueueLease('data-export', 'stale')).resolves.toBe(false) + }) + }) +}) diff --git a/tests/notification.service.test.ts b/tests/notification.service.test.ts index fb03b14f..a11f2da3 100644 --- a/tests/notification.service.test.ts +++ b/tests/notification.service.test.ts @@ -2,20 +2,25 @@ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest' import { NotificationService } from '../src/services/notification.service' // Use vi.hoisted to mock dependencies before they are imported by the service -const { mockSendEachForMulticast } = vi.hoisted(() => ({ - mockSendEachForMulticast: vi.fn().mockResolvedValue({ failureCount: 0, responses: [] }) -})) +const { mockSendEachForMulticast, mockAdmin } = vi.hoisted(() => { + const mockSendEachForMulticast = vi.fn().mockResolvedValue({ failureCount: 0, responses: [] }) + + return { + mockSendEachForMulticast, + mockAdmin: { + apps: [{ name: 'mock-app' }], + initializeApp: vi.fn(), + credential: { + cert: vi.fn().mockReturnValue({}) + }, + messaging: vi.fn().mockReturnValue({ + sendEachForMulticast: mockSendEachForMulticast + }) + } + } +}) -vi.mock('firebase-admin', () => ({ - apps: [{ name: 'mock-app' }], - initializeApp: vi.fn(), - credential: { - cert: vi.fn().mockReturnValue({}) - }, - messaging: vi.fn().mockReturnValue({ - sendEachForMulticast: mockSendEachForMulticast - }) -})) +vi.mock('firebase-admin', () => ({ ...mockAdmin, default: mockAdmin })) // Use vi.hoisted to ensure these are available for vi.mock const { mockPrisma } = vi.hoisted(() => ({ diff --git a/tests/services/webhook.service.spec.ts b/tests/services/webhook.service.spec.ts index 25f98fd3..fd1fc78b 100644 --- a/tests/services/webhook.service.spec.ts +++ b/tests/services/webhook.service.spec.ts @@ -72,11 +72,12 @@ describe('WebhookService', () => { { id: 'ep1', url: 'https://ep1.com', secret: 's1', events: 'module.completed', isActive: true }, ]) mockPrismaInstance.webhookDelivery.create.mockResolvedValue({ id: 'd1' }) - mockPrismaInstance.webhookDelivery.findMany.mockResolvedValue([]) // for processQueue + mockPrismaInstance.webhookDelivery.findMany.mockResolvedValue([]) await service.queueEvent('module.completed', { foo: 'bar' }) expect(mockPrismaInstance.webhookDelivery.create).toHaveBeenCalledOnce() + expect(mockPrismaInstance.webhookDelivery.findMany).not.toHaveBeenCalled() const createCall = mockPrismaInstance.webhookDelivery.create.mock.calls[0][0] expect(createCall.data.eventType).toBe('module.completed') expect(JSON.parse(createCall.data.payload).data).toEqual({ foo: 'bar' }) diff --git a/tests/workers/outbox-relay.test.ts b/tests/workers/outbox-relay.test.ts new file mode 100644 index 00000000..bb6bac27 --- /dev/null +++ b/tests/workers/outbox-relay.test.ts @@ -0,0 +1,364 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' +import { EventSchemaRegistry } from '../../src/lib/transactions/event-schema' +import { OutboxHandlerRegistry } from '../../src/lib/transactions/handler-registry' +import { OutboxRelay } from '../../src/workers/outbox-relay' +import type { OutboxEventHandler } from '../../src/lib/transactions/types' + +vi.mock('../../src/utils/logger', () => ({ + default: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})) + +const silentLog = { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() } + +interface FakeEvent { + id: string + eventType: string + eventVersion: number + aggregateId: string + aggregateType: string + payload: string + status: string +} + +interface FakeJob { + id: string + outboxEventId: string + jobType: string + status: string + attempt: number + maxAttempts: number + lastError: string | null + availableAt: number +} + +class FakeDb { + events: FakeEvent[] = [] + jobs: FakeJob[] = [] + private seq = 0 + + addEvent(partial: Partial & { eventType: string; payload: unknown }): FakeEvent { + const event: FakeEvent = { + id: partial.id ?? `evt-${++this.seq}`, + eventType: partial.eventType, + eventVersion: partial.eventVersion ?? 1, + aggregateId: partial.aggregateId ?? 'agg-1', + aggregateType: partial.aggregateType ?? 'User', + payload: JSON.stringify(partial.payload), + status: partial.status ?? 'PENDING', + } + this.events.push(event) + + return event + } + + get outboxEvent() { + return { + findMany: async (args: any) => { + let rows = this.events.filter(e => e.status === args.where.status) + if (args.where.jobAttempts?.none) { + rows = rows.filter(e => !this.jobs.some(j => j.outboxEventId === e.id)) + } + + return rows.slice(0, args.take ?? rows.length) + }, + update: async (args: any) => { + const event = this.events.find(e => e.id === args.where.id)! + Object.assign(event, args.data) + + return event + }, + count: async () => this.events.filter(e => e.status === 'PENDING').length, + findFirst: async () => this.events.find(e => e.status === 'PENDING') ?? null, + } + } + + get jobAttempt() { + return { + create: async (args: any) => { + const job: FakeJob = { + id: `job-${++this.seq}`, + outboxEventId: args.data.outboxEventId, + jobType: args.data.jobType, + status: args.data.status, + attempt: args.data.attempt, + maxAttempts: args.data.maxAttempts, + lastError: null, + availableAt: Date.now(), + } + this.jobs.push(job) + + return job + }, + findUnique: async (args: any) => { + const job = this.jobs.find(j => j.id === args.where.id) + if (!job) return null + + return { + ...job, + outboxEvent: this.events.find(e => e.id === job.outboxEventId) ?? null, + } + }, + findMany: async (args: any) => + this.jobs.filter( + j => + j.outboxEventId === args.where.outboxEventId && j.status === args.where.status + ), + } + } + + async $transaction(arg: any) { + return Array.isArray(arg) ? Promise.all(arg) : arg(this) + } +} + +class FakeLeaseService { + constructor(private db: FakeDb) {} + + async leaseJob({ jobType }: { jobType: string }) { + const job = this.db.jobs.find( + j => j.jobType === jobType && j.status === 'PENDING' && j.availableAt <= Date.now() + ) + if (!job) return null + job.status = 'LEASED' + const event = this.db.events.find(e => e.id === job.outboxEventId)! + + return { + jobId: job.id, + leaseToken: `tok-${job.id}`, + leasedUntil: new Date(Date.now() + 30_000), + attempt: job.attempt, + payload: JSON.parse(event.payload), + } + } + + async completeJob(jobId: string) { + const job = this.db.jobs.find(j => j.id === jobId)! + job.status = 'COMPLETED' + const siblings = this.db.jobs.filter(j => j.outboxEventId === job.outboxEventId) + if (siblings.every(j => j.status === 'COMPLETED')) { + this.db.events.find(e => e.id === job.outboxEventId)!.status = 'PUBLISHED' + } + } + + async failJob(jobId: string, _token: string, error: Error | string) { + const job = this.db.jobs.find(j => j.id === jobId)! + job.attempt += 1 + job.lastError = error instanceof Error ? error.message : String(error) + + if (job.attempt >= job.maxAttempts) { + job.status = 'DEAD_LETTER' + this.db.events.find(e => e.id === job.outboxEventId)!.status = 'DEAD_LETTER' + } else { + job.status = 'PENDING' + job.availableAt = Date.now() + 60_000 + } + } + + async resetJobForRetry(jobId: string) { + const job = this.db.jobs.find(j => j.id === jobId)! + job.status = 'PENDING' + job.attempt = 0 + job.lastError = null + job.availableAt = Date.now() + } + + releaseBackoff() { + for (const job of this.db.jobs) job.availableAt = Date.now() + } +} + +function buildRelay(db: FakeDb, handlers: OutboxEventHandler[]) { + const schemas = new EventSchemaRegistry() + schemas.register({ eventType: 'UserCreated', version: 1, validate: () => undefined }) + schemas.register({ eventType: 'OrderPlaced', version: 1, validate: () => undefined }) + + const registry = new OutboxHandlerRegistry(schemas) + for (const handler of handlers) registry.register(handler) + + const leases = new FakeLeaseService(db) + const relay = new OutboxRelay({ + prisma: db as any, + handlers: registry, + schemas, + leaseService: leases as any, + log: silentLog, + }) + + return { relay, leases } +} + +describe('OutboxRelay', () => { + let db: FakeDb + + beforeEach(() => { + vi.clearAllMocks() + db = new FakeDb() + }) + + it('dispatches an event to its registered handler and publishes it', async () => { + const seen: unknown[] = [] + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + + const { relay } = buildRelay(db, [ + { + name: 'h1', + eventType: 'UserCreated', + eventVersion: 1, + handle: async ctx => { + seen.push(ctx.payload) + }, + }, + ]) + + const summary = await relay.runOnce() + + expect(summary).toMatchObject({ materialized: 1, dispatched: 1, failed: 0 }) + expect(seen).toEqual([{ userId: 'u1' }]) + expect(db.events[0].status).toBe('PUBLISHED') + }) + + it('routes each event only to handlers for its own type', async () => { + const userSeen: string[] = [] + const orderSeen: string[] = [] + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + db.addEvent({ eventType: 'OrderPlaced', payload: { orderId: 'o1' } }) + + const { relay } = buildRelay(db, [ + { + name: 'user-handler', + eventType: 'UserCreated', + eventVersion: 1, + handle: async ctx => { + userSeen.push(ctx.eventId) + }, + }, + { + name: 'order-handler', + eventType: 'OrderPlaced', + eventVersion: 1, + handle: async ctx => { + orderSeen.push(ctx.eventId) + }, + }, + ]) + + await relay.runOnce() + + expect(userSeen).toEqual(['evt-1']) + expect(orderSeen).toEqual(['evt-2']) + }) + + it('publishes only after every handler for the type succeeds', async () => { + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + let failNext = true + + const { relay, leases } = buildRelay(db, [ + { name: 'ok', eventType: 'UserCreated', eventVersion: 1, handle: async () => undefined }, + { + name: 'flaky', + eventType: 'UserCreated', + eventVersion: 1, + handle: async () => { + if (failNext) { + failNext = false + throw new Error('transient') + } + }, + }, + ]) + + await relay.runOnce() + expect(db.events[0].status).toBe('PENDING') + + leases.releaseBackoff() + await relay.runOnce() + expect(db.events[0].status).toBe('PUBLISHED') + }) + + it('dead-letters an event whose type has no registered handler', async () => { + db.addEvent({ eventType: 'OrderPlaced', payload: { orderId: 'o1' } }) + + const { relay } = buildRelay(db, [ + { name: 'user-handler', eventType: 'UserCreated', eventVersion: 1, handle: async () => undefined }, + ]) + + const summary = await relay.runOnce() + + expect(summary.unhandled).toBe(1) + expect(db.events[0].status).toBe('DEAD_LETTER') + expect(silentLog.error).toHaveBeenCalledWith(expect.stringContaining('no handler registered')) + }) + + it('dead-letters after max attempts without blocking other event types', async () => { + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + db.addEvent({ eventType: 'OrderPlaced', payload: { orderId: 'o1' } }) + + const { relay, leases } = buildRelay(db, [ + { + name: 'always-fails', + eventType: 'UserCreated', + eventVersion: 1, + maxAttempts: 2, + handle: async () => { + throw new Error('provider down') + }, + }, + { name: 'healthy', eventType: 'OrderPlaced', eventVersion: 1, handle: async () => undefined }, + ]) + + await relay.runOnce() + leases.releaseBackoff() + await relay.runOnce() + + const [userEvent, orderEvent] = db.events + expect(userEvent.status).toBe('DEAD_LETTER') + expect(orderEvent.status).toBe('PUBLISHED') + expect(db.jobs.find(j => j.jobType === 'always-fails')?.lastError).toContain('provider down') + }) + + it('replays a dead-lettered event without duplicating side effects', async () => { + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + const effects: string[] = [] + let healthy = false + + const { relay } = buildRelay(db, [ + { + name: 'recovers', + eventType: 'UserCreated', + eventVersion: 1, + maxAttempts: 1, + handle: async ctx => { + if (!healthy) throw new Error('downstream down') + effects.push(ctx.eventId) + }, + }, + ]) + + await relay.runOnce() + expect(db.events[0].status).toBe('DEAD_LETTER') + expect(effects).toEqual([]) + + healthy = true + await relay.replayDeadLetter('evt-1') + await relay.runOnce() + + expect(db.events[0].status).toBe('PUBLISHED') + expect(effects).toEqual(['evt-1']) + + await relay.runOnce() + expect(effects).toEqual(['evt-1']) + }) + + it('does not re-materialise jobs for an event it has already fanned out', async () => { + db.addEvent({ eventType: 'UserCreated', payload: { userId: 'u1' } }) + + const { relay } = buildRelay(db, [ + { name: 'h1', eventType: 'UserCreated', eventVersion: 1, handle: async () => undefined }, + ]) + + await relay.runOnce() + await relay.runOnce() + + expect(db.jobs).toHaveLength(1) + }) +}) diff --git a/tests/workers/scheduled-job-runner.test.ts b/tests/workers/scheduled-job-runner.test.ts new file mode 100644 index 00000000..2a2e6d0a --- /dev/null +++ b/tests/workers/scheduled-job-runner.test.ts @@ -0,0 +1,374 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' +import type { SchedulerConfig } from '../../src/config/scheduler' +import { QueueMetricsRegistry } from '../../src/workers/queue-metrics' +import { + ScheduledJobRunner, + type QueueLeaseApi, +} from '../../src/workers/scheduled-job-runner' +import type { ScheduledQueue } from '../../src/workers/queue-registry' + +vi.mock('../../src/utils/logger', () => ({ + default: { + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + debug: vi.fn(), + }, +})) + +const silentLog = { + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + debug: vi.fn(), +} + +function testConfig(overrides: Partial = {}): SchedulerConfig { + return { + intervalMs: 10, + leaseMs: 2_000, + shutdownTimeoutMs: 1_000, + inProcess: false, + only: [], + disabled: [], + ownerId: 'test-runner', + isEnabled(queueName: string) { + return ( + !this.disabled.includes(queueName) && + (this.only.length === 0 || this.only.includes(queueName)) + ) + }, + intervalFor() { + return this.intervalMs + }, + leaseFor() { + return this.leaseMs + }, + ...overrides, + } as SchedulerConfig +} + +class FakeLeaseStore implements QueueLeaseApi { + private readonly rows = new Map() + private seq = 0 + + acquireQueueLease = vi.fn( + async (options: { queueName: string; leaseMs?: number; owner?: string }) => { + const now = Date.now() + const leaseMs = options.leaseMs ?? 60_000 + const current = this.rows.get(options.queueName) + + if (current && current.until > now) { + return null + } + + const leaseToken = `${options.owner ?? 'anon'}-${++this.seq}` + this.rows.set(options.queueName, { token: leaseToken, until: now + leaseMs }) + + return { + queueName: options.queueName, + leaseToken, + leasedUntil: new Date(now + leaseMs), + } + } + ) + + renewQueueLease = vi.fn(async (queueName: string, leaseToken: string, leaseMs = 60_000) => { + const current = this.rows.get(queueName) + if (!current || current.token !== leaseToken) return false + current.until = Date.now() + leaseMs + + return true + }) + + releaseQueueLease = vi.fn(async (queueName: string, leaseToken: string) => { + const current = this.rows.get(queueName) + if (!current || current.token !== leaseToken) return false + this.rows.delete(queueName) + + return true + }) + + hold(queueName: string, forMs: number): void { + this.rows.set(queueName, { token: 'foreign-holder', until: Date.now() + forMs }) + } +} + +async function waitFor(predicate: () => boolean, timeoutMs = 2_000): Promise { + const deadline = Date.now() + timeoutMs + while (!predicate()) { + if (Date.now() > deadline) { + throw new Error('waitFor timed out') + } + await new Promise(resolve => setTimeout(resolve, 5)) + } +} + +interface FakeRow { + id: string + nextAttemptAt: Date +} + +function fakeQueue( + name: string, + rows: FakeRow[], + processed: string[], + drainImpl?: () => Promise +): ScheduledQueue { + return { + name, + drain: + drainImpl ?? + (async () => { + const now = Date.now() + const due = rows.filter(row => row.nextAttemptAt.getTime() <= now) + for (const row of due) { + rows.splice(rows.indexOf(row), 1) + processed.push(row.id) + } + }), + inspect: async () => { + const now = Date.now() + const due = rows.filter(row => row.nextAttemptAt.getTime() <= now) + const oldest = [...due].sort( + (a, b) => a.nextAttemptAt.getTime() - b.nextAttemptAt.getTime() + )[0] + + return { + depth: rows.length, + due: due.length, + oldestDueAt: oldest?.nextAttemptAt ?? null, + } + }, + } +} + +describe('ScheduledJobRunner', () => { + let leases: FakeLeaseStore + let metrics: QueueMetricsRegistry + + beforeEach(() => { + vi.clearAllMocks() + leases = new FakeLeaseStore() + metrics = new QueueMetricsRegistry() + }) + + it('drains due rows on a timer with no inbound traffic', async () => { + const processed: string[] = [] + const rows: FakeRow[] = [{ id: 'row-1', nextAttemptAt: new Date(Date.now() - 1) }] + + const runner = new ScheduledJobRunner({ + queues: [fakeQueue('email', rows, processed)], + leaseService: leases, + config: testConfig(), + metrics, + log: silentLog, + }) + + runner.start() + await waitFor(() => processed.length === 1) + await runner.stop() + + expect(processed).toEqual(['row-1']) + }) + + it('picks up a row whose nextAttemptAt falls due within one interval', async () => { + const processed: string[] = [] + const rows: FakeRow[] = [{ id: 'retry-1', nextAttemptAt: new Date(Date.now() + 40) }] + + const runner = new ScheduledJobRunner({ + queues: [fakeQueue('email', rows, processed)], + leaseService: leases, + config: testConfig({ intervalMs: 10 }), + metrics, + log: silentLog, + }) + + runner.start() + + await new Promise(resolve => setTimeout(resolve, 15)) + expect(processed).toEqual([]) + + await waitFor(() => processed.length === 1) + await runner.stop() + + expect(processed).toEqual(['retry-1']) + }) + + it('skips a tick when another holder owns the queue lease', async () => { + const processed: string[] = [] + const rows: FakeRow[] = [{ id: 'row-1', nextAttemptAt: new Date(Date.now() - 1) }] + leases.hold('email', 60_000) + + const runner = new ScheduledJobRunner({ + queues: [fakeQueue('email', rows, processed)], + leaseService: leases, + config: testConfig(), + metrics, + log: silentLog, + }) + + const outcomes = await runner.runOnce() + + expect(outcomes.email).toBe('skipped') + expect(processed).toEqual([]) + expect(metrics.snapshot()[0].skipped).toBe(1) + expect(metrics.snapshot()[0].attempts).toBe(0) + }) + + it('never lets two replicas process the same row twice', async () => { + const processed: string[] = [] + const rows: FakeRow[] = Array.from({ length: 25 }, (_, index) => ({ + id: `row-${index}`, + nextAttemptAt: new Date(Date.now() - 1), + })) + + const slowDrain = async () => { + const now = Date.now() + const due = rows.filter(row => row.nextAttemptAt.getTime() <= now) + for (const row of due) { + const index = rows.indexOf(row) + if (index === -1) continue + rows.splice(index, 1) + await new Promise(resolve => setTimeout(resolve, 1)) + processed.push(row.id) + } + } + + const makeRunner = (owner: string) => + new ScheduledJobRunner({ + queues: [fakeQueue('email', rows, processed, slowDrain)], + leaseService: leases, + config: testConfig({ ownerId: owner, intervalMs: 5 }), + metrics: new QueueMetricsRegistry(), + log: silentLog, + }) + + const replicaA = makeRunner('replica-a') + const replicaB = makeRunner('replica-b') + + replicaA.start() + replicaB.start() + + await waitFor(() => processed.length === 25, 5_000) + await Promise.all([replicaA.stop(), replicaB.stop()]) + + expect(new Set(processed).size).toBe(25) + expect(processed).toHaveLength(25) + }) + + it('releases the lease and records a failure when a drain throws', async () => { + const failing: ScheduledQueue = { + name: 'webhook', + drain: async () => { + throw new Error('provider down') + }, + inspect: async () => ({ depth: 3, due: 3, oldestDueAt: new Date(Date.now() - 5_000) }), + } + + const runner = new ScheduledJobRunner({ + queues: [failing], + leaseService: leases, + config: testConfig(), + metrics, + log: silentLog, + }) + + const outcomes = await runner.runOnce() + + expect(outcomes.webhook).toBe('failed') + expect(leases.releaseQueueLease).toHaveBeenCalledOnce() + + const [snapshot] = metrics.snapshot() + expect(snapshot.failures).toBe(1) + expect(snapshot.lastError).toBe('provider down') + + expect(await runner.runOnce()).toEqual({ webhook: 'failed' }) + expect(leases.acquireQueueLease).toHaveBeenCalledTimes(2) + }) + + it('emits depth, due, attempt, failure and lag metrics per queue', async () => { + const oldestDueAt = new Date(Date.now() - 30_000) + const queue: ScheduledQueue = { + name: 'notification', + drain: async () => undefined, + inspect: async () => ({ depth: 7, due: 4, oldestDueAt }), + } + + const runner = new ScheduledJobRunner({ + queues: [queue], + leaseService: leases, + config: testConfig(), + metrics, + log: silentLog, + }) + + await runner.runOnce() + + const [snapshot] = metrics.snapshot() + expect(snapshot.queue).toBe('notification') + expect(snapshot.depth).toBe(7) + expect(snapshot.due).toBe(4) + expect(snapshot.attempts).toBe(1) + expect(snapshot.failures).toBe(0) + expect(snapshot.lagMs).toBeGreaterThanOrEqual(30_000) + expect(snapshot.lastRunAt).not.toBeNull() + }) + + it('drains the in-flight tick and releases its lease on shutdown', async () => { + let started = false + + const queue: ScheduledQueue = { + name: 'data-export', + drain: async () => { + started = true + await new Promise(resolve => setTimeout(resolve, 60)) + }, + inspect: async () => ({ depth: 1, due: 1, oldestDueAt: null }), + } + + const runner = new ScheduledJobRunner({ + queues: [queue], + leaseService: leases, + config: testConfig(), + metrics, + log: silentLog, + }) + + runner.start() + await waitFor(() => started) + + await runner.stop() + + expect(leases.releaseQueueLease).toHaveBeenCalledOnce() + expect(await leases.acquireQueueLease({ queueName: 'data-export' })).not.toBeNull() + + const ticksAtStop = metrics.snapshot()[0].attempts + await new Promise(resolve => setTimeout(resolve, 40)) + expect(metrics.snapshot()[0].attempts).toBe(ticksAtStop) + }) + + it('honours the disabled and allow lists when registering queues', () => { + const queues = [ + fakeQueue('email', [], []), + fakeQueue('webhook', [], []), + fakeQueue('notification', [], []), + ] + + const disabled = new ScheduledJobRunner({ + queues, + leaseService: leases, + config: testConfig({ disabled: ['webhook'] }), + log: silentLog, + }) + expect(disabled.registeredQueues).toEqual(['email', 'notification']) + + const only = new ScheduledJobRunner({ + queues, + leaseService: leases, + config: testConfig({ only: ['notification'] }), + log: silentLog, + }) + expect(only.registeredQueues).toEqual(['notification']) + }) +}) diff --git a/tsx b/tsx new file mode 100644 index 00000000..e69de29b