From e414263325e007d92fb47b85abfec50f5976d565 Mon Sep 17 00:00:00 2001 From: danieloche635-bit Date: Mon, 28 Sep 2026 03:26:05 +0000 Subject: [PATCH] feat: job queue retries, cron recurring billing, and split payments (#917, #918, #952) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements three payment-platform capabilities, each as a self-contained, tested domain module following the existing BaseService/Result + repository conventions, with an HTTP surface and docs. #952 — background job queue (backend/src/services/job-queue) Retries with configurable exponential back-off (+ optional full jitter), a dead-letter queue with replay, metrics, and lifecycle events. The caller drives the queue, so retry timing is deterministic and testable. #918 — cron-based recurring billing (backend/src/services/recurring-billing) Schedules validated against cron-parser with IANA timezone checks, run-due invoicing with start/end/maxRuns bounds, pause/resume/cancel/reschedule, and a sweep wired into the existing scheduled-tasks registry. #917 — split payments (backend/src/services/split-payments) Exact largest-remainder allocation on integer basis points, so platform fee + recipient shares reconcile to the payment down to the minor unit. Plans, execution, preview, per-recipient minimums and analytics. Tests: 66 new unit tests (job queue 16, recurring billing 25, split payments 25), all passing. New/changed files add no tsc errors over the repo baseline and eslint is clean. 🤖 Generated with Codebuff Co-Authored-By: Codebuff --- backend/docs/JOB_QUEUE.md | 102 +++++ backend/docs/RECURRING_BILLING.md | 109 +++++ backend/docs/SPLIT_PAYMENTS.md | 108 +++++ backend/src/config/scheduled-tasks.ts | 20 + backend/src/index.ts | 6 + backend/src/routes/job-queue.ts | 92 ++++ backend/src/routes/recurring-billing.ts | 214 +++++++++ backend/src/routes/split-payments.ts | 188 ++++++++ backend/src/services/job-queue/backoff.ts | 57 +++ backend/src/services/job-queue/index.ts | 20 + .../src/services/job-queue/jobQueue.test.ts | 244 ++++++++++ backend/src/services/job-queue/jobQueue.ts | 378 ++++++++++++++++ backend/src/services/job-queue/types.ts | 109 +++++ .../src/services/recurring-billing/index.ts | 28 ++ .../recurring-billing.test.ts | 312 +++++++++++++ .../recurringBillingService.ts | 417 ++++++++++++++++++ .../services/recurring-billing/schedule.ts | 204 +++++++++ .../src/services/recurring-billing/store.ts | 97 ++++ .../src/services/recurring-billing/types.ts | 137 ++++++ .../src/services/split-payments/allocation.ts | 241 ++++++++++ backend/src/services/split-payments/index.ts | 27 ++ .../split-payments/split-payments.test.ts | 352 +++++++++++++++ .../split-payments/splitPaymentService.ts | 276 ++++++++++++ backend/src/services/split-payments/store.ts | 78 ++++ backend/src/services/split-payments/types.ts | 145 ++++++ 25 files changed, 3961 insertions(+) create mode 100644 backend/docs/JOB_QUEUE.md create mode 100644 backend/docs/RECURRING_BILLING.md create mode 100644 backend/docs/SPLIT_PAYMENTS.md create mode 100644 backend/src/routes/job-queue.ts create mode 100644 backend/src/routes/recurring-billing.ts create mode 100644 backend/src/routes/split-payments.ts create mode 100644 backend/src/services/job-queue/backoff.ts create mode 100644 backend/src/services/job-queue/index.ts create mode 100644 backend/src/services/job-queue/jobQueue.test.ts create mode 100644 backend/src/services/job-queue/jobQueue.ts create mode 100644 backend/src/services/job-queue/types.ts create mode 100644 backend/src/services/recurring-billing/index.ts create mode 100644 backend/src/services/recurring-billing/recurring-billing.test.ts create mode 100644 backend/src/services/recurring-billing/recurringBillingService.ts create mode 100644 backend/src/services/recurring-billing/schedule.ts create mode 100644 backend/src/services/recurring-billing/store.ts create mode 100644 backend/src/services/recurring-billing/types.ts create mode 100644 backend/src/services/split-payments/allocation.ts create mode 100644 backend/src/services/split-payments/index.ts create mode 100644 backend/src/services/split-payments/split-payments.test.ts create mode 100644 backend/src/services/split-payments/splitPaymentService.ts create mode 100644 backend/src/services/split-payments/store.ts create mode 100644 backend/src/services/split-payments/types.ts diff --git a/backend/docs/JOB_QUEUE.md b/backend/docs/JOB_QUEUE.md new file mode 100644 index 00000000..e97b3587 --- /dev/null +++ b/backend/docs/JOB_QUEUE.md @@ -0,0 +1,102 @@ +# Background Job Queue (Issue #952) + +A reliable background job queue with automatic retries, exponential back-off, +configurable limits and dead-letter handling. + +## Concepts + +| Concept | Description | +|---------|-------------| +| **Job** | A named unit of work plus a payload, tracked by a `JobRecord`. | +| **Handler** | The async function that performs a job's work. Rejecting signals failure. | +| **Attempt** | One execution of a handler. `attempts` is the number started so far. | +| **Retry policy** | `maxAttempts`, `initialDelayMs`, `maxDelayMs`, `multiplier`, `jitter`. | +| **DLQ** | Dead-letter queue — terminal failures kept for inspection and replay. | + +Job states: `pending → processing → completed`, or on failure `pending` +(retry scheduled) and finally `dead` when attempts are exhausted. + +## Reliability guarantees + +- **At-least-once** — a job is marked `completed` only after its handler resolves. +- **Bounded retries** — a failure is rescheduled with exponential back-off until + `maxAttempts` (counts the first attempt) is exhausted. +- **No silent loss** — every terminal failure is written to the DLQ with the + last error message, surfaced through `metrics()` and lifecycle events. +- **Deterministic scheduling** — the queue is driven by the caller (`drain()`), + so retry timing is testable and no hidden timers run in unit tests. + +## Usage + +```ts +import { jobQueue } from './services/job-queue/index.js'; + +// 1. Register a handler once at startup. +jobQueue.registerHandler('send-receipt', async (payload: { paymentId: string }) => { + await receipts.send(payload.paymentId); +}); + +// 2. Enqueue work. +jobQueue.enqueue('send-receipt', { paymentId: 'pay_123' }); + +// 3. Drive the queue — from a worker loop, a request, or a scheduler. +const summary = await jobQueue.drain(); +// { processed: 1, completed: 0, retried: 1, deadLettered: 0 } + +// Or start an interval-based polling loop (unref'd, safe in tests). +jobQueue.start(1_000); +``` + +### Custom retry policy + +```ts +import { JobQueue } from './services/job-queue/index.js'; + +const queue = new JobQueue({ + retry: { maxAttempts: 5, initialDelayMs: 500, maxDelayMs: 30_000, multiplier: 2, jitter: true }, +}); +``` + +Back-off delay for a failed attempt is +`min(initialDelayMs * multiplier ** (attempt - 1), maxDelayMs)`, optionally +multiplied by a full-jitter random factor in `[0, 1)`. + +### Dead-letter handling + +```ts +const dlq = jobQueue.getDeadLetters(); // [{ job, failedAt, reason }] +jobQueue.requeueDeadLetter(dlq[0].job.id); // reset attempts → pending +``` + +## Events + +Pass an `onEvent` sink to observe lifecycle transitions: +`job.enqueued`, `job.started`, `job.completed`, `job.failed`, `job.retrying`, +`job.dead_lettered`. + +## REST admin surface + +Mounted at `/api/v1/job-queue`. + +| Method | Path | Purpose | +|--------|------|---------| +| `POST` | `/enqueue` | Enqueue a registered job (`{ name, payload }`). | +| `GET` | `/metrics` | Queue depth counters. | +| `GET` | `/jobs?state=&name=` | List jobs, filterable by state/name. | +| `GET` | `/dead-letters` | Inspect the DLQ. | +| `POST` | `/dead-letters/:id/requeue` | Requeue a dead-lettered job. | + +```bash +curl -X POST http://localhost:3000/api/v1/job-queue/enqueue \ + -H 'Content-Type: application/json' \ + -d '{ "name": "send-receipt", "payload": { "paymentId": "pay_123" } }' + +curl http://localhost:3000/api/v1/job-queue/metrics +``` + +## Persistence + +`JobQueue` holds jobs in memory; the retry/DLQ semantics are transport +agnostic. For production durability, back the same contract with Redis +(this repo already depends on BullMQ) or the `JobRecord` Prisma model, keeping +`drain()` as the worker entry point. diff --git a/backend/docs/RECURRING_BILLING.md b/backend/docs/RECURRING_BILLING.md new file mode 100644 index 00000000..5627989c --- /dev/null +++ b/backend/docs/RECURRING_BILLING.md @@ -0,0 +1,109 @@ +# Recurring Payment Schedules (Issue #918) + +Cron-based recurring billing. A schedule bills a customer a fixed amount on a +cron cadence; each due run materialises a pending invoice. + +## Concepts + +| Concept | Description | +|---------|-------------| +| **Schedule** | `customerId` + `amount` + `currency` + a `cronExpression` (or `preset`) evaluated in a `timezone`. | +| **Preset** | `hourly`, `daily`, `weekly`, `monthly`, `yearly` — expanded to a canonical cron expression. | +| **Next run** | First cron occurrence after the effective start, stored as `nextRunAt`. | +| **Invoice** | A `pending` charge created when a schedule's `nextRunAt` comes due. | +| **Bounds** | Optional `startAt`, `endAt` and `maxRuns` (invoices cap). | + +Schedule states: `active → paused ⇄ active`, plus `cancelled` and `completed` +(no further runs possible). + +## Cadence + +`PRESET_CRONS`: + +| Preset | Cron | +|--------|------| +| `hourly` | `0 * * * *` | +| `daily` | `0 0 * * *` | +| `weekly` | `0 0 * * 0` | +| `monthly` | `0 0 1 * *` | +| `yearly` | `0 0 1 1 *` | + +Provide a raw `cronExpression` for a custom cadence (e.g. `*/15 * * * *`). +`cronExpression` and `preset` are mutually exclusive; the cron wins if both are +present at the service layer. Cron expressions are validated on create/update, +and timezones are validated as IANA zones (`Intl.DateTimeFormat`). + +## Behaviour + +- **First run** — the first cron occurrence on/after `max(startAt, now)`; a past + `startAt` never backfills. +- **Advance** — after each run, `nextRunAt` is recomputed from the run's due + instant. `runDue` produces at most one invoice per schedule per call, so a + schedule that is behind catches up on subsequent sweeps. +- **Completion** — the schedule flips to `completed` when `maxRuns` is reached or + the next run would fall after `endAt`. +- **Failure safety** — if invoice persistence throws, the schedule is left + untouched, a `recurring_invoice.failed` event fires, and the run is retried on + the next sweep. + +## REST API + +Mounted at `/api/v1/recurring-payments`. + +| Method | Path | Purpose | +|--------|------|---------| +| `POST` | `/schedules` | Create a schedule. | +| `GET` | `/schedules?tenantId=&status=&customerId=&merchantId=` | List schedules. | +| `GET` | `/schedules/:id?tenantId=` | Fetch a schedule. | +| `GET` | `/schedules/:id/upcoming?tenantId=&count=` | Preview the next runs. | +| `GET` | `/schedules/:id/invoices?tenantId=` | List generated invoices. | +| `POST` | `/schedules/:id/pause?tenantId=` | Pause billing. | +| `POST` | `/schedules/:id/resume?tenantId=` | Resume billing (recomputes next run). | +| `POST` | `/schedules/:id/resume` → `/cancel?tenantId=` | Cancel billing. | +| `POST` | `/schedules/:id/reschedule` | Change cadence/timezone. | +| `POST` | `/run-due` | Bill every due schedule. | + +### Create a schedule + +```bash +curl -X POST http://localhost:3000/api/v1/recurring-payments/schedules \ + -H 'Content-Type: application/json' \ + -d '{ + "tenantId": "tenant-1", + "customerId": "cust-9", + "amount": 49.99, + "currency": "USD", + "preset": "monthly", + "startAt": "2026-10-01T00:00:00.000Z", + "maxRuns": 12 + }' +``` + +### Run due billing + +```bash +curl -X POST http://localhost:3000/api/v1/recurring-payments/run-due \ + -H 'Content-Type: application/json' \ + -d '{ "at": "2026-11-01T00:00:00.000Z" }' +``` + +Wire `/run-due` to a scheduler (the repo already has a cron/BullMQ scheduler in +`src/config/scheduled-tasks.ts` and `src/services/bullmq-scheduler.ts`), e.g. +every 15 minutes. + +## Events + +`RecurringBillingService` publishes through the injectable +`RecurringBillingPublisher`: + +- `recurring_schedule.created` +- `recurring_schedule.paused` / `resumed` / `rescheduled` +- `recurring_schedule.cancelled` / `completed` +- `recurring_invoice.generated` / `failed` + +## Persistence + +The service depends on `RecurringScheduleRepository` / +`RecurringInvoiceRepository`; `InMemoryRecurringBillingStore` backs both for +tests and local development. Implement the same interfaces with Prisma +(`recurring_schedules`, `recurring_invoices`) for production. diff --git a/backend/docs/SPLIT_PAYMENTS.md b/backend/docs/SPLIT_PAYMENTS.md new file mode 100644 index 00000000..b11fad4e --- /dev/null +++ b/backend/docs/SPLIT_PAYMENTS.md @@ -0,0 +1,108 @@ +# Split Payments (Issue #917) + +Split a single payment across a platform fee and one or more recipients, with +an allocation that reconciles to the payment exactly. + +## Concepts + +| Concept | Description | +|---------|-------------| +| **Split plan** | A reusable recipe: a set of recipients and a platform fee for a merchant. | +| **Recipient** | `recipientId` + `walletAddress` + a `percentage` share (optionally a `minimumAmount`). | +| **Basis points** | Percentages are normalised to integer basis points (`shareBps`, 1 bp = 0.01%) before any money maths. | +| **Execution** | Applying a plan to a concrete `paymentId` + `totalAmount`, producing per-recipient distributions. | + +Plan states: `active` and `archived` (archived plans cannot be executed). + +## Rounding guarantee + +Percentages must sum to **exactly 100** (the platform fee counts toward that +total). Allocation runs on integer minor units with the largest-remainder +(Hamilton) method, so: + +``` +platformFeeMinor + Σ share.amountMinor === totalMinor +unallocatedMinor === 0 +``` + +No money is created or lost to rounding, even for amounts like `33.33 / 33.33 / +33.34` or one-cent splits across many recipients. + +``` +$100, 2.5% fee, recipients 65% / 32.5% +→ fee 2.50 + 65.00 + 32.50 = 100.00 +``` + +## Validation + +`validateSplitPlan` rejects requests before anything is persisted: + +- missing `tenantId`, empty recipient list, or more than `maxRecipients` (25) +- duplicate or missing `recipientId`, missing `walletAddress` +- a `percentage` outside `(0, 100]`, or a negative `minimumAmount` +- percentages + platform fee that do not sum to exactly 100 +- an unsupported currency (`USD`, `EUR`, `GBP`, `XLM`, `USDC`) + +## Minimum amounts + +A recipient whose allocated amount is below `minimumAmount` is marked +`skipped: true` with a reason, so the caller can hold or reroute that share. + +## REST API + +Mounted at `/api/v1/split-payments`. + +| Method | Path | Purpose | +|--------|------|---------| +| `POST` | `/plans` | Create a split plan. | +| `GET` | `/plans?tenantId=&status=&merchantId=` | List plans. | +| `GET` | `/plans/:id?tenantId=` | Fetch a plan. | +| `POST` | `/plans/:id/archive?tenantId=` | Archive a plan. | +| `POST` | `/plans/:id/execute` | Execute a payment against the plan. | +| `GET` | `/plans/:id/executions?tenantId=` | List executions. | +| `GET` | `/plans/:id/summary?tenantId=` | Execution summary (processed, fees, skips). | +| `GET` | `/plans/:id/preview?tenantId=&totalAmount=` | Preview a split without executing. | + +### Create a plan + +```bash +curl -X POST http://localhost:3000/api/v1/split-payments/plans \ + -H 'Content-Type: application/json' \ + -d '{ + "tenantId": "tenant-1", + "currency": "USD", + "platformFeePercentage": 2.5, + "recipients": [ + { "recipientId": "creator", "walletAddress": "GA...", "percentage": 65 }, + { "recipientId": "affiliate", "walletAddress": "GB...", "percentage": 32.5 } + ] + }' +``` + +### Execute a payment + +```bash +curl -X POST http://localhost:3000/api/v1/split-payments/plans//execute \ + -H 'Content-Type: application/json' \ + -d '{ "tenantId": "tenant-1", "paymentId": "pay_123", "totalAmount": 199.99 }' +``` + +Each execution returns the exact per-recipient `distributions` and the +`platformFeeAmount`; their sum equals `totalAmount`. + +## Events + +`SplitPaymentService` publishes through the injectable `SplitEventPublisher`: + +- `split_plan.created` +- `split_plan.archived` +- `split.executed` + +## Persistence + +The service depends on `SplitPlanRepository` / `SplitExecutionRepository`; +`InMemorySplitStore` backs both for tests and local development. Implement the +same interfaces with Prisma (`split_plans`, `split_executions`) for production. + +> Note: the legacy `/api/v1/splits` endpoints expose a percentage-of-total model +> and remain unchanged. This module is the exact-allocation engine for #917. diff --git a/backend/src/config/scheduled-tasks.ts b/backend/src/config/scheduled-tasks.ts index e2273299..c3f9a05a 100644 --- a/backend/src/config/scheduled-tasks.ts +++ b/backend/src/config/scheduled-tasks.ts @@ -19,6 +19,7 @@ import { markOverdueRequests } from '../services/gdpr.js'; import { sandboxCleanupJobs } from '../jobs/sandbox-cleanup.js'; import { SubscriptionService } from '../services/subscription.service.js'; import { SubscriptionProcessor } from '../jobs/subscription-processor.js'; +import { recurringBillingService } from '../services/recurring-billing/index.js'; import { ethers } from 'ethers'; // --------------------------------------------------------------------------- @@ -123,6 +124,25 @@ const RAW_TASKS: Omit & { defaultSchedule: string await processor.processPendingRenewals(); }, }, + { + id: 'recurring-payments-run-due', + name: 'Run due recurring payments', + description: 'Generates invoices for cron-based recurring payment schedules whose next run is due (Issue #918).', + defaultSchedule: '*/15 * * * *', + timeoutMs: 2 * 60 * 1000, + handler: async () => { + const result = await recurringBillingService.runDue(); + if (!result.ok) { + console.error(`[jobs] recurring billing sweep failed: ${result.error.message}`); + return; + } + if (result.value.processed > 0) { + console.log( + `[jobs] recurring billing: processed ${result.value.processed} schedule(s), generated ${result.value.invoices.length} invoice(s)`, + ); + } + }, + }, { id: 'gdpr-deadline-check', name: 'GDPR 30-day deadline enforcement', diff --git a/backend/src/index.ts b/backend/src/index.ts index 87e32d2c..049a8a37 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -109,6 +109,9 @@ import { workspacesRouter } from './routes/workspaces.js'; import { teamsRouter } from './routes/teams.js'; import { merchantAuditRouter } from './routes/merchant-audit.js'; import { installmentsRouter } from './routes/installments.js'; +import { jobQueueRouter } from './routes/job-queue.js'; +import { recurringBillingRouter } from './routes/recurring-billing.js'; +import { splitPaymentsRouter } from './routes/split-payments.js'; // Validate environment variables at startup validateEnv(); @@ -267,6 +270,9 @@ apiV1Router.use('/workspaces', workspacesRouter); apiV1Router.use('/teams', teamsRouter); apiV1Router.use('/merchant-audit', merchantAuditRouter); apiV1Router.use('/installments', installmentsRouter); +apiV1Router.use('/job-queue', jobQueueRouter); +apiV1Router.use('/recurring-payments', recurringBillingRouter); +apiV1Router.use('/split-payments', splitPaymentsRouter); apiV1Router.get('/compression/metrics', (_req, res) => { res.json(getCompressionMetrics()); }); diff --git a/backend/src/routes/job-queue.ts b/backend/src/routes/job-queue.ts new file mode 100644 index 00000000..0a074a61 --- /dev/null +++ b/backend/src/routes/job-queue.ts @@ -0,0 +1,92 @@ +/** + * job-queue.ts — Issue #952: Background job queue with retries + * + * Observability and administration surface for the shared `jobQueue`. + * Handlers are registered by application wiring; this router only enqueues + * already-registered jobs and exposes queue state. + * + * POST /job-queue/enqueue — enqueue a registered job + * GET /job-queue/metrics — queue depth counters + * GET /job-queue/jobs?state=&name= — list jobs (filterable) + * GET /job-queue/dead-letters — inspect the DLQ + * POST /job-queue/dead-letters/:id/requeue — retry a dead-lettered job + */ +import { Router, type Request, type Response } from 'express'; +import { z } from 'zod'; + +import { asyncHandler } from '../middleware/errorHandler.js'; +import { jobQueue } from '../services/job-queue/index.js'; +import type { JobState } from '../services/job-queue/index.js'; + +export const jobQueueRouter = Router(); +const firstParam = (value: string | string[] | undefined): string => + Array.isArray(value) ? value[0] ?? '' : value ?? ''; + +const jobStateSchema = z.enum(['pending', 'processing', 'completed', 'failed', 'dead']); + +const enqueueSchema = z.object({ + name: z.string().min(1), + payload: z.unknown().optional(), +}); + +// ── POST /job-queue/enqueue ────────────────────────────────────────────────── + +jobQueueRouter.post( + '/enqueue', + asyncHandler(async (req: Request, res: Response) => { + const parsed = enqueueSchema.safeParse(req.body); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + if (!jobQueue.hasHandler(parsed.data.name)) { + return res.status(404).json({ success: false, error: `No handler registered for job "${parsed.data.name}"` }); + } + + const job = jobQueue.enqueue(parsed.data.name, parsed.data.payload ?? {}); + return res.status(202).json({ success: true, data: job }); + }), +); + +// ── GET /job-queue/metrics ─────────────────────────────────────────────────── + +jobQueueRouter.get( + '/metrics', + asyncHandler(async (_req: Request, res: Response) => { + return res.json({ success: true, data: jobQueue.metrics() }); + }), +); + +// ── GET /job-queue/jobs ────────────────────────────────────────────────────── + +jobQueueRouter.get( + '/jobs', + asyncHandler(async (req: Request, res: Response) => { + const stateParse = jobStateSchema.safeParse(req.query.state); + const state: JobState | undefined = stateParse.success ? stateParse.data : undefined; + const name = (req.query.name as string) || undefined; + + return res.json({ success: true, data: jobQueue.listJobs({ state, name }) }); + }), +); + +// ── GET /job-queue/dead-letters ────────────────────────────────────────────── + +jobQueueRouter.get( + '/dead-letters', + asyncHandler(async (_req: Request, res: Response) => { + return res.json({ success: true, data: jobQueue.getDeadLetters() }); + }), +); + +// ── POST /job-queue/dead-letters/:id/requeue ───────────────────────────────── + +jobQueueRouter.post( + '/dead-letters/:id/requeue', + asyncHandler(async (req: Request, res: Response) => { + const requeued = jobQueue.requeueDeadLetter(firstParam(req.params.id)); + if (!requeued) { + return res.status(404).json({ success: false, error: 'Dead letter not found' }); + } + return res.json({ success: true, data: requeued }); + }), +); diff --git a/backend/src/routes/recurring-billing.ts b/backend/src/routes/recurring-billing.ts new file mode 100644 index 00000000..7189e3b6 --- /dev/null +++ b/backend/src/routes/recurring-billing.ts @@ -0,0 +1,214 @@ +/** + * recurring-billing.ts — Issue #918: Recurring payment schedules with + * cron-based billing + * + * REST surface for recurring schedules. + * + * POST /recurring-payments/schedules — create a schedule + * GET /recurring-payments/schedules?tenantId=&status= — list schedules + * GET /recurring-payments/schedules/:id?tenantId= — fetch a schedule + * GET /recurring-payments/schedules/:id/upcoming?tenantId=&count= — next runs + * GET /recurring-payments/schedules/:id/invoices?tenantId= — generated invoices + * POST /recurring-payments/schedules/:id/pause?tenantId= — pause billing + * POST /recurring-payments/schedules/:id/resume?tenantId= — resume billing + * POST /recurring-payments/schedules/:id/cancel?tenantId= — cancel billing + * POST /recurring-payments/schedules/:id/reschedule — change cadence + * POST /recurring-payments/run-due — bill all due schedules + */ +import { Router, type Request, type Response } from 'express'; +import { z } from 'zod'; + +import { asyncHandler } from '../middleware/errorHandler.js'; +import { recurringBillingService } from '../services/recurring-billing/index.js'; +import type { ServiceError } from '../lib/result.js'; + +export const recurringBillingRouter = Router(); +const firstParam = (value: string | string[] | undefined): string => + Array.isArray(value) ? value[0] ?? '' : value ?? ''; + +const presetSchema = z.enum(['hourly', 'daily', 'weekly', 'monthly', 'yearly']); +const statusSchema = z.enum(['active', 'paused', 'cancelled', 'completed']); + +const createSchema = z + .object({ + tenantId: z.string().min(1), + customerId: z.string().min(1), + amount: z.number().positive(), + currency: z.string().length(3).optional(), + cronExpression: z.string().min(1).optional(), + preset: presetSchema.optional(), + timezone: z.string().min(1).optional(), + startAt: z.string().datetime().optional(), + endAt: z.string().datetime().optional(), + maxRuns: z.number().int().positive().optional(), + merchantId: z.string().min(1).optional(), + name: z.string().min(1).optional(), + metadata: z.record(z.unknown()).optional(), + }) + .refine((value) => !(value.cronExpression && value.preset), { + message: 'Provide either cronExpression or preset, not both', + path: ['cronExpression'], + }); + +const rescheduleSchema = z.object({ + cronExpression: z.string().min(1).optional(), + preset: presetSchema.optional(), + timezone: z.string().min(1).optional(), +}); + +const cancelSchema = z.object({ reason: z.string().max(280).optional() }); +const runDueSchema = z.object({ at: z.string().datetime().optional() }); + +function requireTenant(req: Request, res: Response): string | null { + const tenantId = (req.query.tenantId ?? req.body?.tenantId) as string | undefined; + if (!tenantId) { + res.status(400).json({ success: false, error: 'tenantId is required' }); + return null; + } + return tenantId; +} + +function sendResult(res: Response, result: { ok: true; value: T } | { ok: false; error: ServiceError }) { + if (result.ok) { + return res.json({ success: true, data: result.value }); + } + return res.status(result.error.statusCode ?? 400).json({ + success: false, + error: result.error.message, + code: result.error.code, + }); +} + +// ── POST /recurring-payments/schedules ─────────────────────────────────────── + +recurringBillingRouter.post( + '/schedules', + asyncHandler(async (req: Request, res: Response) => { + const parsed = createSchema.safeParse(req.body); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + const result = await recurringBillingService.createSchedule(parsed.data); + if (!result.ok) { + return sendResult(res, result); + } + return res.status(201).json({ success: true, data: result.value }); + }), +); + +// ── GET /recurring-payments/schedules ──────────────────────────────────────── + +recurringBillingRouter.get( + '/schedules', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + + const status = statusSchema.safeParse(req.query.status).success + ? (req.query.status as z.infer) + : undefined; + const result = await recurringBillingService.listSchedules(tenantId, { + status, + customerId: (req.query.customerId as string) || undefined, + merchantId: (req.query.merchantId as string) || undefined, + }); + return sendResult(res, result); + }), +); + +// ── GET /recurring-payments/schedules/:id/upcoming ─────────────────────────── + +recurringBillingRouter.get( + '/schedules/:id/upcoming', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const count = req.query.count ? Number(req.query.count) : 5; + const result = await recurringBillingService.previewUpcoming(tenantId, firstParam(req.params.id), count); + return sendResult(res, result); + }), +); + +// ── GET /recurring-payments/schedules/:id/invoices ─────────────────────────── + +recurringBillingRouter.get( + '/schedules/:id/invoices', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const result = await recurringBillingService.listInvoices(tenantId, firstParam(req.params.id)); + return sendResult(res, result); + }), +); + +// ── GET /recurring-payments/schedules/:id ──────────────────────────────────── + +recurringBillingRouter.get( + '/schedules/:id', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const result = await recurringBillingService.getSchedule(tenantId, firstParam(req.params.id)); + return sendResult(res, result); + }), +); + +// ── Lifecycle transitions ──────────────────────────────────────────────────── + +recurringBillingRouter.post( + '/schedules/:id/pause', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await recurringBillingService.pauseSchedule(tenantId, firstParam(req.params.id))); + }), +); + +recurringBillingRouter.post( + '/schedules/:id/resume', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await recurringBillingService.resumeSchedule(tenantId, firstParam(req.params.id))); + }), +); + +recurringBillingRouter.post( + '/schedules/:id/cancel', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const parsed = cancelSchema.safeParse(req.body ?? {}); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + return sendResult(res, await recurringBillingService.cancelSchedule(tenantId, firstParam(req.params.id), parsed.data.reason)); + }), +); + +recurringBillingRouter.post( + '/schedules/:id/reschedule', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const parsed = rescheduleSchema.safeParse(req.body ?? {}); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + return sendResult(res, await recurringBillingService.reschedule(tenantId, firstParam(req.params.id), parsed.data)); + }), +); + +// ── POST /recurring-payments/run-due ───────────────────────────────────────── + +recurringBillingRouter.post( + '/run-due', + asyncHandler(async (req: Request, res: Response) => { + const parsed = runDueSchema.safeParse(req.body ?? {}); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + const result = await recurringBillingService.runDue(parsed.data.at ?? new Date()); + return sendResult(res, result); + }), +); diff --git a/backend/src/routes/split-payments.ts b/backend/src/routes/split-payments.ts new file mode 100644 index 00000000..c1f0c86e --- /dev/null +++ b/backend/src/routes/split-payments.ts @@ -0,0 +1,188 @@ +/** + * split-payments.ts — Issue #917: Split payments between multiple recipients + * + * REST surface for split plans. + * + * POST /split-payments/plans — create a plan + * GET /split-payments/plans?tenantId=&status=&merchantId= — list plans + * GET /split-payments/plans/:id?tenantId= — fetch a plan + * POST /split-payments/plans/:id/archive?tenantId= — archive a plan + * POST /split-payments/plans/:id/execute — execute a payment + * GET /split-payments/plans/:id/executions?tenantId= — list executions + * GET /split-payments/plans/:id/summary?tenantId= — execution summary + * GET /split-payments/plans/:id/preview?tenantId=&totalAmount= — preview split + */ +import { Router, type Request, type Response } from 'express'; +import { z } from 'zod'; + +import { asyncHandler } from '../middleware/errorHandler.js'; +import { splitPaymentService } from '../services/split-payments/index.js'; +import type { ServiceError } from '../lib/result.js'; + +export const splitPaymentsRouter = Router(); +const firstParam = (value: string | string[] | undefined): string => + Array.isArray(value) ? value[0] ?? '' : value ?? ''; + +const statusSchema = z.enum(['active', 'archived']); + +const recipientSchema = z.object({ + recipientId: z.string().min(1), + walletAddress: z.string().min(1), + percentage: z.number().positive().max(100), + minimumAmount: z.number().nonnegative().optional(), + label: z.string().min(1).optional(), +}); + +const createPlanSchema = z.object({ + tenantId: z.string().min(1), + recipients: z.array(recipientSchema).min(1), + platformFeePercentage: z.number().min(0).max(100).optional(), + currency: z.string().length(3).optional(), + merchantId: z.string().min(1).optional(), + name: z.string().min(1).optional(), + metadata: z.record(z.unknown()).optional(), +}); + +const executeSchema = z.object({ + tenantId: z.string().min(1).optional(), + paymentId: z.string().min(1), + totalAmount: z.number().positive(), + currency: z.string().length(3).optional(), +}); + +function requireTenant(req: Request, res: Response): string | null { + const tenantId = (req.query.tenantId ?? req.body?.tenantId) as string | undefined; + if (!tenantId) { + res.status(400).json({ success: false, error: 'tenantId is required' }); + return null; + } + return tenantId; +} + +function sendResult(res: Response, result: { ok: true; value: T } | { ok: false; error: ServiceError }) { + if (result.ok) { + return res.json({ success: true, data: result.value }); + } + return res.status(result.error.statusCode ?? 400).json({ + success: false, + error: result.error.message, + code: result.error.code, + }); +} + +// ── POST /split-payments/plans ─────────────────────────────────────────────── + +splitPaymentsRouter.post( + '/plans', + asyncHandler(async (req: Request, res: Response) => { + const parsed = createPlanSchema.safeParse(req.body); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + const result = await splitPaymentService.createPlan(parsed.data); + if (!result.ok) { + return sendResult(res, result); + } + return res.status(201).json({ success: true, data: result.value }); + }), +); + +// ── GET /split-payments/plans ──────────────────────────────────────────────── + +splitPaymentsRouter.get( + '/plans', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + + const status = statusSchema.safeParse(req.query.status).success + ? (req.query.status as z.infer) + : undefined; + const result = await splitPaymentService.listPlans(tenantId, { + status, + merchantId: (req.query.merchantId as string) || undefined, + }); + return sendResult(res, result); + }), +); + +// ── GET /split-payments/plans/:id/summary ──────────────────────────────────── + +splitPaymentsRouter.get( + '/plans/:id/summary', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await splitPaymentService.getExecutionSummary(tenantId, firstParam(req.params.id))); + }), +); + +// ── GET /split-payments/plans/:id/executions ───────────────────────────────── + +splitPaymentsRouter.get( + '/plans/:id/executions', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await splitPaymentService.listExecutions(tenantId, firstParam(req.params.id))); + }), +); + +// ── GET /split-payments/plans/:id/preview ──────────────────────────────────── + +splitPaymentsRouter.get( + '/plans/:id/preview', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + const totalAmount = Number(req.query.totalAmount); + return sendResult(res, await splitPaymentService.previewAllocation(tenantId, firstParam(req.params.id), totalAmount)); + }), +); + +// ── GET /split-payments/plans/:id ──────────────────────────────────────────── + +splitPaymentsRouter.get( + '/plans/:id', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await splitPaymentService.getPlan(tenantId, firstParam(req.params.id))); + }), +); + +// ── POST /split-payments/plans/:id/archive ─────────────────────────────────── + +splitPaymentsRouter.post( + '/plans/:id/archive', + asyncHandler(async (req: Request, res: Response) => { + const tenantId = requireTenant(req, res); + if (!tenantId) return; + return sendResult(res, await splitPaymentService.archivePlan(tenantId, firstParam(req.params.id))); + }), +); + +// ── POST /split-payments/plans/:id/execute ─────────────────────────────────── + +splitPaymentsRouter.post( + '/plans/:id/execute', + asyncHandler(async (req: Request, res: Response) => { + const parsed = executeSchema.safeParse(req.body); + if (!parsed.success) { + return res.status(400).json({ success: false, error: 'Invalid payload', details: parsed.error.format() }); + } + const tenantId = parsed.data.tenantId ?? ((req.query.tenantId as string) || undefined); + if (!tenantId) { + return res.status(400).json({ success: false, error: 'tenantId is required' }); + } + const result = await splitPaymentService.executeSplit(tenantId, firstParam(req.params.id), { + paymentId: parsed.data.paymentId, + totalAmount: parsed.data.totalAmount, + currency: parsed.data.currency, + }); + if (!result.ok) { + return sendResult(res, result); + } + return res.status(201).json({ success: true, data: result.value }); + }), +); diff --git a/backend/src/services/job-queue/backoff.ts b/backend/src/services/job-queue/backoff.ts new file mode 100644 index 00000000..ea900e87 --- /dev/null +++ b/backend/src/services/job-queue/backoff.ts @@ -0,0 +1,57 @@ +/** + * backoff.ts — Issue #952: Background job queue with retries + * + * Pure helpers for computing retry delays. Keeping the maths side-effect free + * makes the schedule deterministic and directly unit-testable: the caller can + * inject a `random` source to remove jitter from assertions. + */ +import type { RetryPolicy } from './types.js'; + +export const DEFAULT_RETRY_POLICY: RetryPolicy = { + maxAttempts: 3, + initialDelayMs: 1_000, + maxDelayMs: 60_000, + multiplier: 2, + jitter: false, +}; + +/** + * Delay before retrying after `failedAttempt` consecutive failures. + * + * The un-jittered delay is `initialDelayMs * multiplier ** (failedAttempt - 1)`, + * capped at `maxDelayMs`. When the policy enables jitter the delay is + * multiplied by a value in `[0, 1)` (full jitter) to spread retries out and + * avoid a thundering herd. + * + * @param failedAttempt 1-based number of the attempt that just failed. + * @param policy Retry policy; defaults to `DEFAULT_RETRY_POLICY`. + * @param random Injectable RNG in `[0, 1)`, defaults to `Math.random`. + */ +export function computeBackoffDelayMs( + failedAttempt: number, + policy: RetryPolicy = DEFAULT_RETRY_POLICY, + random: () => number = Math.random, +): number { + if (failedAttempt < 1) { + throw new Error('failedAttempt must be >= 1'); + } + if (policy.maxAttempts < 1) { + throw new Error('maxAttempts must be >= 1'); + } + if (policy.initialDelayMs < 0 || policy.maxDelayMs < 0) { + throw new Error('backoff delays must be non-negative'); + } + if (policy.multiplier < 1) { + throw new Error('multiplier must be >= 1'); + } + + const exponential = policy.initialDelayMs * Math.pow(policy.multiplier, failedAttempt - 1); + const capped = Math.min(exponential, policy.maxDelayMs); + const delay = policy.jitter ? capped * random() : capped; + return Math.max(0, Math.round(delay)); +} + +/** True when another attempt is still permitted for the given attempt count. */ +export function canRetry(attempts: number, policy: RetryPolicy): boolean { + return attempts < policy.maxAttempts; +} diff --git a/backend/src/services/job-queue/index.ts b/backend/src/services/job-queue/index.ts new file mode 100644 index 00000000..746fa075 --- /dev/null +++ b/backend/src/services/job-queue/index.ts @@ -0,0 +1,20 @@ +/** + * index.ts — Issue #952: Background job queue with retries + * + * Public surface of the background job queue domain. Import `jobQueue` (the + * shared instance) in application wiring and register handlers with + * `jobQueue.registerHandler('name', fn)` before enqueuing that job. + */ +export * from './types.js'; +export { DEFAULT_RETRY_POLICY, computeBackoffDelayMs, canRetry } from './backoff.js'; +export { + JobQueue, + type JobQueueOptions, + type DrainOptions, + type JobFilter, +} from './jobQueue.js'; + +import { JobQueue } from './jobQueue.js'; + +/** Shared queue instance used by the HTTP admin surface and workers. */ +export const jobQueue = new JobQueue({ name: 'default' }); diff --git a/backend/src/services/job-queue/jobQueue.test.ts b/backend/src/services/job-queue/jobQueue.test.ts new file mode 100644 index 00000000..2058fc2c --- /dev/null +++ b/backend/src/services/job-queue/jobQueue.test.ts @@ -0,0 +1,244 @@ +/** + * jobQueue.test.ts — Issue #952: Background job queue with retries + * + * Covers the success path, graceful failure, exponential back-off with + * configurable limits, dead-letter handling, requeueing, metrics and events. + */ +import { describe, expect, it, vi } from 'vitest'; + +import { JobQueue, type JobQueueEvent, type RetryPolicy } from './index.js'; +import { computeBackoffDelayMs, canRetry, DEFAULT_RETRY_POLICY } from './backoff.js'; + +const T0 = new Date('2026-09-28T00:00:00.000Z'); + +const silentLogger = { info: () => undefined, warn: () => undefined, error: () => undefined }; + +function makeQueue(overrides: Partial = {}, clock = { current: T0 }) { + const events: JobQueueEvent[] = []; + let counter = 0; + const queue = new JobQueue({ + retry: { maxAttempts: 3, initialDelayMs: 1_000, maxDelayMs: 10_000, multiplier: 2, jitter: false, ...overrides }, + now: () => clock.current, + idFactory: () => `job-${++counter}`, + onEvent: (event) => events.push(event), + logger: silentLogger, + }); + return { queue, events, clock }; +} + +describe('computeBackoffDelayMs', () => { + it('grows exponentially from the initial delay', () => { + const policy: RetryPolicy = { ...DEFAULT_RETRY_POLICY, initialDelayMs: 1_000, multiplier: 2, maxDelayMs: 100_000 }; + expect(computeBackoffDelayMs(1, policy)).toBe(1_000); + expect(computeBackoffDelayMs(2, policy)).toBe(2_000); + expect(computeBackoffDelayMs(3, policy)).toBe(4_000); + }); + + it('caps the delay at maxDelayMs', () => { + const policy: RetryPolicy = { ...DEFAULT_RETRY_POLICY, initialDelayMs: 1_000, multiplier: 10, maxDelayMs: 5_000 }; + expect(computeBackoffDelayMs(4, policy)).toBe(5_000); + }); + + it('applies full jitter using the injected RNG', () => { + const policy: RetryPolicy = { ...DEFAULT_RETRY_POLICY, initialDelayMs: 1_000, multiplier: 2, maxDelayMs: 100_000, jitter: true }; + expect(computeBackoffDelayMs(2, policy, () => 0.5)).toBe(1_000); + expect(computeBackoffDelayMs(2, policy, () => 0)).toBe(0); + }); + + it('rejects invalid inputs', () => { + expect(() => computeBackoffDelayMs(0)).toThrow(); + expect(() => computeBackoffDelayMs(1, { ...DEFAULT_RETRY_POLICY, maxAttempts: 0 })).toThrow(); + expect(() => computeBackoffDelayMs(1, { ...DEFAULT_RETRY_POLICY, multiplier: 0 })).toThrow(); + }); + + it('reports retry eligibility', () => { + const policy: RetryPolicy = { ...DEFAULT_RETRY_POLICY, maxAttempts: 2 }; + expect(canRetry(0, policy)).toBe(true); + expect(canRetry(1, policy)).toBe(true); + expect(canRetry(2, policy)).toBe(false); + }); +}); + +describe('JobQueue registration & enqueue', () => { + it('runs a registered handler and completes the job', async () => { + const { queue, events } = makeQueue(); + const handler = vi.fn().mockResolvedValue(undefined); + queue.registerHandler('send-email', handler); + + const job = queue.enqueue('send-email', { to: 'a@b.c' }); + expect(job.state).toBe('pending'); + + const summary = await queue.drain(); + expect(summary).toEqual({ processed: 1, completed: 1, retried: 0, deadLettered: 0 }); + expect(handler).toHaveBeenCalledWith({ to: 'a@b.c' }, { jobId: 'job-1', name: 'send-email', attempt: 1 }); + + const stored = queue.getJob('job-1'); + expect(stored?.state).toBe('completed'); + expect(stored?.attempts).toBe(1); + + const types = events.map((event) => event.type); + expect(types).toEqual(['job.enqueued', 'job.started', 'job.completed']); + }); + + it('rejects enqueuing an unregistered job name', () => { + const { queue } = makeQueue(); + expect(() => queue.enqueue('nope', {})).toThrow(/Unknown job/); + }); + + it('drains nothing when the queue is empty', async () => { + const { queue } = makeQueue(); + const summary = await queue.drain(); + expect(summary.processed).toBe(0); + }); +}); + +describe('JobQueue retries with exponential back-off', () => { + it('retries a failing job and succeeds on the final attempt', async () => { + const { queue, events, clock } = makeQueue(); + const handler = vi + .fn() + .mockRejectedValueOnce(new Error('flaky 1')) + .mockRejectedValueOnce(new Error('flaky 2')) + .mockResolvedValue(undefined); + queue.registerHandler('flaky', handler); + queue.enqueue('flaky', {}); + + const first = await queue.drain(); + expect(first).toEqual({ processed: 1, completed: 0, retried: 1, deadLettered: 0 }); + expect(queue.getJob('job-1')?.state).toBe('pending'); + expect(queue.getJob('job-1')?.nextAttemptAt).toBe('2026-09-28T00:00:01.000Z'); + + // Before the back-off elapses the job is not eligible. + clock.current = new Date('2026-09-28T00:00:00.500Z'); + expect((await queue.drain()).processed).toBe(0); + + // First retry is due at +1s. + clock.current = new Date('2026-09-28T00:00:01.000Z'); + expect((await queue.drain()).retried).toBe(1); + expect(queue.getJob('job-1')?.nextAttemptAt).toBe('2026-09-28T00:00:03.000Z'); + + // Second retry is due at +3s and succeeds. + clock.current = new Date('2026-09-28T00:00:03.000Z'); + const third = await queue.drain(); + expect(third.completed).toBe(1); + expect(queue.getJob('job-1')?.state).toBe('completed'); + expect(queue.getJob('job-1')?.attempts).toBe(3); + expect(handler).toHaveBeenCalledTimes(3); + + expect(events.map((event) => event.type)).toContain('job.retrying'); + }); + + it('honours a custom back-off limit', async () => { + const { queue, clock } = makeQueue({ maxAttempts: 5, initialDelayMs: 100, maxDelayMs: 150, multiplier: 3 }); + queue.registerHandler('fail', () => { + throw new Error('always'); + }); + queue.enqueue('fail', {}); + + await queue.drain(); + expect(queue.getJob('job-1')?.nextAttemptAt).toBe('2026-09-28T00:00:00.100Z'); + + clock.current = new Date('2026-09-28T00:00:00.100Z'); + await queue.drain(); + // 100 * 3 = 300, capped at 150. + expect(queue.getJob('job-1')?.nextAttemptAt).toBe('2026-09-28T00:00:00.250Z'); + }); +}); + +describe('JobQueue dead-letter handling', () => { + it('moves a job to the DLQ once attempts are exhausted', async () => { + const { queue, events } = makeQueue({ maxAttempts: 2, initialDelayMs: 0 }); + queue.registerHandler('always-fail', () => Promise.reject(new Error('boom'))); + queue.enqueue('always-fail', { id: 7 }); + + await queue.drain(); + expect(queue.getJob('job-1')?.state).toBe('pending'); + + const second = await queue.drain(); + expect(second.deadLettered).toBe(1); + + const stored = queue.getJob('job-1'); + expect(stored?.state).toBe('dead'); + expect(stored?.attempts).toBe(2); + expect(stored?.lastError).toBe('boom'); + + const dlq = queue.getDeadLetters(); + expect(dlq).toHaveLength(1); + expect(dlq[0].reason).toBe('boom'); + expect(dlq[0].job.id).toBe('job-1'); + expect(events.map((event) => event.type)).toContain('job.dead_lettered'); + }); + + it('dead-letters a job whose handler is missing at run time', async () => { + const { queue } = makeQueue(); + queue.registerHandler('temp', () => undefined); + queue.enqueue('temp', {}); + // Simulate the handler being unregistered between enqueue and drain. + (queue as unknown as { handlers: Map }).handlers.delete('temp'); + + const summary = await queue.drain(); + expect(summary.deadLettered).toBe(1); + expect(queue.getDeadLetters()[0].reason).toMatch(/No handler/); + }); + + it('requeues a dead-lettered job and processes it again', async () => { + const { queue } = makeQueue({ maxAttempts: 1 }); + let shouldFail = true; + queue.registerHandler('switchy', () => { + if (shouldFail) throw new Error('first time'); + }); + queue.enqueue('switchy', {}); + + await queue.drain(); + expect(queue.getJob('job-1')?.state).toBe('dead'); + + shouldFail = false; + const requeued = queue.requeueDeadLetter('job-1'); + expect(requeued?.state).toBe('pending'); + expect(requeued?.attempts).toBe(0); + expect(queue.getDeadLetters()).toHaveLength(0); + + const summary = await queue.drain(); + expect(summary.completed).toBe(1); + expect(queue.getJob('job-1')?.state).toBe('completed'); + }); + + it('returns null when requeueing an unknown dead letter', () => { + const { queue } = makeQueue(); + expect(queue.requeueDeadLetter('missing')).toBeNull(); + }); +}); + +describe('JobQueue metrics & listing', () => { + it('tracks counts across states', async () => { + const { queue } = makeQueue({ maxAttempts: 1 }); + queue.registerHandler('ok', () => undefined); + queue.registerHandler('bad', () => Promise.reject(new Error('nope'))); + + queue.enqueue('ok', {}); + queue.enqueue('bad', {}); + await queue.drain(); + + const metrics = queue.metrics(); + expect(metrics.completed).toBe(1); + expect(metrics.dead).toBe(1); + expect(metrics.pending).toBe(0); + expect(metrics.failed).toBe(1); + expect(metrics.totalEnqueued).toBe(2); + expect(metrics.totalCompleted).toBe(1); + expect(metrics.totalDeadLettered).toBe(1); + }); + + it('filters jobs by state and name', async () => { + const { queue } = makeQueue(); + queue.registerHandler('a', () => undefined); + queue.registerHandler('b', () => undefined); + queue.enqueue('a', {}); + queue.enqueue('b', {}); + await queue.drain(); + + expect(queue.listJobs({ state: 'completed' })).toHaveLength(2); + expect(queue.listJobs({ name: 'a' })).toHaveLength(1); + expect(queue.listJobs({ name: 'a', state: 'completed' })[0].name).toBe('a'); + }); +}); diff --git a/backend/src/services/job-queue/jobQueue.ts b/backend/src/services/job-queue/jobQueue.ts new file mode 100644 index 00000000..fe46d4aa --- /dev/null +++ b/backend/src/services/job-queue/jobQueue.ts @@ -0,0 +1,378 @@ +/** + * jobQueue.ts — Issue #952: Background job queue with retries + * + * A small, dependency-free background job queue. It is driven by the caller + * (`drain()` is invoked by a timer, a request handler, or a test) which keeps + * the retry/back-off behaviour fully deterministic and lets the same code run + * on top of an in-memory store today and a durable transport later. + * + * Guarantees: + * - At-least-once execution: a job is only marked `completed` after its + * handler resolves. + * - Bounded retries: failures are rescheduled with exponential back-off until + * `maxAttempts` is exhausted, after which the job moves to the DLQ. + * - No silent loss: every terminal failure is recorded in the DLQ with the + * last error message and is observable through `metrics()` / events. + */ +import { randomUUID } from 'node:crypto'; + +import { DEFAULT_RETRY_POLICY, canRetry, computeBackoffDelayMs } from './backoff.js'; +import type { + DeadLetterEntry, + JobContext, + JobHandler, + JobQueueEvent, + JobQueueLogger, + JobQueueMetrics, + JobQueueSummary, + JobRecord, + JobState, + RetryPolicy, +} from './types.js'; + +export interface JobQueueOptions { + /** Logical name, used in log lines. */ + name?: string; + /** Overrides for the default retry policy. */ + retry?: Partial; + logger?: JobQueueLogger; + /** Injectable clock, defaults to `() => new Date()`. */ + now?: () => Date; + /** Injectable id generator, defaults to `randomUUID`. */ + idFactory?: () => string; + /** Injectable RNG for jitter, defaults to `Math.random`. */ + random?: () => number; + /** Sink for lifecycle events (metrics, audit, webhooks…). */ + onEvent?: (event: JobQueueEvent) => void; +} + +export interface DrainOptions { + /** Reference time; defaults to the queue clock. */ + now?: Date; + /** Upper bound on jobs considered in a single pass. Defaults to Infinity. */ + maxJobs?: number; +} + +export interface JobFilter { + state?: JobState; + name?: string; +} + +const consoleLogger: JobQueueLogger = { + info: (message, meta) => console.info(`[job-queue] ${message}`, meta ?? ''), + warn: (message, meta) => console.warn(`[job-queue] ${message}`, meta ?? ''), + error: (message, meta) => console.error(`[job-queue] ${message}`, meta ?? ''), +}; + +export class JobQueue { + private readonly handlers = new Map>(); + private readonly jobs = new Map(); + private readonly deadLetters: DeadLetterEntry[] = []; + private readonly retry: RetryPolicy; + private readonly logger: JobQueueLogger; + private readonly now: () => Date; + private readonly idFactory: () => string; + private readonly random: () => number; + private readonly onEvent?: (event: JobQueueEvent) => void; + + private totalEnqueued = 0; + private totalCompleted = 0; + private totalFailedAttempts = 0; + private totalDeadLettered = 0; + private timer?: ReturnType; + + constructor(options: JobQueueOptions = {}) { + this.retry = { ...DEFAULT_RETRY_POLICY, ...options.retry }; + this.logger = options.logger ?? consoleLogger; + this.now = options.now ?? (() => new Date()); + this.idFactory = options.idFactory ?? (() => randomUUID()); + this.random = options.random ?? Math.random; + this.onEvent = options.onEvent; + } + + /** Register (or replace) the handler for a job name. */ + registerHandler(name: string, handler: JobHandler): void { + if (!name) { + throw new Error('Job name is required'); + } + this.handlers.set(name, handler as unknown as JobHandler); + } + + hasHandler(name: string): boolean { + return this.handlers.has(name); + } + + /** Enqueue a job for asynchronous processing. */ + enqueue(name: string, payload: TPayload): JobRecord { + if (!this.handlers.has(name)) { + throw new Error(`Unknown job: ${name}`); + } + + const timestamp = this.now().toISOString(); + const record: JobRecord = { + id: this.idFactory(), + name, + payload, + state: 'pending', + attempts: 0, + maxAttempts: this.retry.maxAttempts, + createdAt: timestamp, + updatedAt: timestamp, + startedAt: null, + finishedAt: null, + lastError: null, + nextAttemptAt: null, + }; + + this.jobs.set(record.id, record as JobRecord); + this.totalEnqueued += 1; + this.emit({ + type: 'job.enqueued', + jobId: record.id, + name, + attempt: 0, + occurredAt: timestamp, + }); + return record; + } + + getJob(id: string): JobRecord | undefined { + return this.jobs.get(id); + } + + listJobs(filter: JobFilter = {}): JobRecord[] { + return Array.from(this.jobs.values()) + .filter((job) => (filter.state ? job.state === filter.state : true)) + .filter((job) => (filter.name ? job.name === filter.name : true)) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt) || a.id.localeCompare(b.id)); + } + + /** Snapshot of the dead-letter queue. */ + getDeadLetters(): DeadLetterEntry[] { + return this.deadLetters.map((entry) => ({ ...entry, job: { ...entry.job } })); + } + + /** + * Move a dead-lettered job back to `pending` so it can be processed again. + * Resets the attempt counter. + */ + requeueDeadLetter(jobId: string): JobRecord | null { + const index = this.deadLetters.findIndex((entry) => entry.job.id === jobId); + if (index === -1) { + return null; + } + const [entry] = this.deadLetters.splice(index, 1); + const job = this.jobs.get(jobId); + if (!job || !entry) { + return null; + } + job.state = 'pending'; + job.attempts = 0; + job.nextAttemptAt = null; + job.lastError = null; + job.finishedAt = null; + job.updatedAt = this.now().toISOString(); + this.emit({ + type: 'job.enqueued', + jobId: job.id, + name: job.name, + attempt: 0, + occurredAt: job.updatedAt, + }); + return job; + } + + metrics(): JobQueueMetrics { + let pending = 0; + let processing = 0; + let completed = 0; + let dead = 0; + for (const job of this.jobs.values()) { + if (job.state === 'pending') pending += 1; + else if (job.state === 'processing') processing += 1; + else if (job.state === 'completed') completed += 1; + else if (job.state === 'dead') dead += 1; + } + return { + pending, + processing, + completed, + dead, + failed: this.totalFailedAttempts, + totalEnqueued: this.totalEnqueued, + totalCompleted: this.totalCompleted, + totalDeadLettered: this.totalDeadLettered, + }; + } + + /** + * Process every job whose next attempt is due, up to `maxJobs`. + * + * A job that fails and is rescheduled within the same pass is not retried + * again until the next `drain()` call, which prevents tight retry loops. + */ + async drain(options: DrainOptions = {}): Promise { + const now = options.now ?? this.now(); + const maxJobs = options.maxJobs ?? Number.POSITIVE_INFINITY; + const seen = new Set(); + + let processed = 0; + let completed = 0; + let retried = 0; + let deadLettered = 0; + + while (processed < maxJobs) { + const job = this.nextEligible(now, seen); + if (!job) { + break; + } + seen.add(job.id); + + await this.runAttempt(job, now); + processed += 1; + if (job.state === 'completed') completed += 1; + else if (job.state === 'dead') deadLettered += 1; + else retried += 1; + } + + return { processed, completed, retried, deadLettered }; + } + + /** Run a single due job. Returns false when nothing was eligible. */ + async processNext(now: Date = this.now()): Promise { + const job = this.nextEligible(now, new Set()); + if (!job) { + return false; + } + await this.runAttempt(job, now); + return true; + } + + /** Start a polling loop that drains the queue every `intervalMs`. */ + start(intervalMs = 1_000): void { + if (this.timer) { + return; + } + this.timer = setInterval(() => { + void this.drain().catch((error) => { + this.logger.error('drain failed', { error: (error as Error).message }); + }); + }, intervalMs); + if (typeof this.timer.unref === 'function') { + this.timer.unref(); + } + } + + stop(): void { + if (this.timer) { + clearInterval(this.timer); + this.timer = undefined; + } + } + + // ── Internals ────────────────────────────────────────────────────────────── + + private nextEligible(now: Date, exclude: Set): JobRecord | undefined { + for (const job of this.jobs.values()) { + if (job.state !== 'pending' || exclude.has(job.id)) { + continue; + } + if (!job.nextAttemptAt || new Date(job.nextAttemptAt).getTime() <= now.getTime()) { + return job; + } + } + return undefined; + } + + private async runAttempt(job: JobRecord, now: Date): Promise { + const handler = this.handlers.get(job.name); + if (!handler) { + this.deadLetter(job, now, `No handler registered for "${job.name}"`); + return; + } + + const attempt = job.attempts + 1; + job.state = 'processing'; + job.attempts = attempt; + job.startedAt = this.now().toISOString(); + job.updatedAt = job.startedAt; + this.emit({ type: 'job.started', jobId: job.id, name: job.name, attempt, occurredAt: job.updatedAt }); + + const context: JobContext = { jobId: job.id, name: job.name, attempt }; + + try { + await handler(job.payload, context); + job.state = 'completed'; + job.finishedAt = this.now().toISOString(); + job.updatedAt = job.finishedAt; + job.lastError = null; + job.nextAttemptAt = null; + this.totalCompleted += 1; + this.emit({ + type: 'job.completed', + jobId: job.id, + name: job.name, + attempt, + occurredAt: job.updatedAt, + }); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + this.totalFailedAttempts += 1; + job.lastError = message; + job.updatedAt = this.now().toISOString(); + this.emit({ + type: 'job.failed', + jobId: job.id, + name: job.name, + attempt, + occurredAt: job.updatedAt, + error: message, + }); + this.logger.warn('job attempt failed', { jobId: job.id, name: job.name, attempt, error: message }); + + if (canRetry(job.attempts, this.retry)) { + const delay = computeBackoffDelayMs(job.attempts, this.retry, this.random); + job.state = 'pending'; + job.nextAttemptAt = new Date(now.getTime() + delay).toISOString(); + this.emit({ + type: 'job.retrying', + jobId: job.id, + name: job.name, + attempt, + occurredAt: job.updatedAt, + error: message, + }); + } else { + this.deadLetter(job, now, message); + } + } + } + + private deadLetter(job: JobRecord, now: Date, reason: string): void { + job.state = 'dead'; + job.lastError = reason; + job.finishedAt = this.now().toISOString(); + job.updatedAt = job.finishedAt; + job.nextAttemptAt = null; + this.totalDeadLettered += 1; + this.deadLetters.push({ job: { ...job }, failedAt: now.toISOString(), reason }); + this.emit({ + type: 'job.dead_lettered', + jobId: job.id, + name: job.name, + attempt: job.attempts, + occurredAt: job.updatedAt, + error: reason, + }); + this.logger.error('job moved to dead-letter queue', { + jobId: job.id, + name: job.name, + attempts: job.attempts, + reason, + }); + } + + private emit(event: JobQueueEvent): void { + this.onEvent?.(event); + } +} diff --git a/backend/src/services/job-queue/types.ts b/backend/src/services/job-queue/types.ts new file mode 100644 index 00000000..ca10f196 --- /dev/null +++ b/backend/src/services/job-queue/types.ts @@ -0,0 +1,109 @@ +/** + * types.ts — Issue #952: Background job queue with retries + * + * Domain types for a durable background job queue. A job is enqueued with a + * payload, executed by a registered handler, and — on failure — retried with + * exponential back-off until it either succeeds or is moved to the + * dead-letter queue (DLQ). + */ + +/** Lifecycle of a queued job. */ +export type JobState = 'pending' | 'processing' | 'completed' | 'failed' | 'dead'; + +/** + * Retry configuration. `maxAttempts` counts the initial attempt, so a value of + * 1 means "never retry". + */ +export interface RetryPolicy { + /** Total attempts, including the first. Must be >= 1. */ + maxAttempts: number; + /** Delay before the first retry, in milliseconds. */ + initialDelayMs: number; + /** Upper bound applied after exponential growth, in milliseconds. */ + maxDelayMs: number; + /** Exponential growth factor applied per consecutive failure. */ + multiplier: number; + /** When true, apply full jitter to the computed delay. */ + jitter: boolean; +} + +/** Context handed to a handler on each attempt. */ +export interface JobContext { + jobId: string; + name: string; + /** 1-based attempt number. */ + attempt: number; +} + +/** A handler performs the job's work and rejects to signal a failure. */ +export type JobHandler = ( + payload: TPayload, + context: JobContext, +) => Promise | void; + +export interface JobRecord { + id: string; + name: string; + payload: TPayload; + state: JobState; + /** Number of attempts started so far. */ + attempts: number; + maxAttempts: number; + createdAt: string; + updatedAt: string; + startedAt?: string | null; + finishedAt?: string | null; + /** Error message from the most recent failed attempt. */ + lastError?: string | null; + /** When the next attempt is eligible to run (ISO-8601), or null. */ + nextAttemptAt?: string | null; +} + +export interface DeadLetterEntry { + job: JobRecord; + failedAt: string; + reason: string; +} + +export interface JobQueueMetrics { + pending: number; + processing: number; + completed: number; + failed: number; + dead: number; + totalEnqueued: number; + totalCompleted: number; + totalDeadLettered: number; +} + +export type JobQueueEventType = + | 'job.enqueued' + | 'job.started' + | 'job.completed' + | 'job.failed' + | 'job.retrying' + | 'job.dead_lettered'; + +export interface JobQueueEvent { + type: JobQueueEventType; + jobId: string; + name: string; + attempt: number; + occurredAt: string; + error?: string; +} + +export interface JobQueueSummary { + /** Jobs with `nextAttemptAt <= now` that were considered. */ + processed: number; + completed: number; + retried: number; + deadLettered: number; +} + +/** Minimal logger contract so the queue stays decoupled from a logging lib. */ +export interface JobQueueLogger { + info(message: string, meta?: Record): void; + warn(message: string, meta?: Record): void; + error(message: string, meta?: Record): void; +} diff --git a/backend/src/services/recurring-billing/index.ts b/backend/src/services/recurring-billing/index.ts new file mode 100644 index 00000000..43710677 --- /dev/null +++ b/backend/src/services/recurring-billing/index.ts @@ -0,0 +1,28 @@ +/** + * index.ts — Issue #918: Recurring payment schedules with cron-based billing + * + * Public surface of the recurring-billing domain. + */ +export * from './types.js'; +export { + DEFAULT_RECURRING_BILLING_CONFIG, + PRESET_CRONS, + isValidCronExpression, + isValidTimezone, + nextRunAfter, + upcomingRuns, + resolveCronExpression, + roundCurrency, + validateRecurringSchedule, +} from './schedule.js'; +export { + InMemoryRecurringBillingStore, + type RecurringInvoiceRepository, + type RecurringScheduleRepository, +} from './store.js'; +export { + RecurringBillingService, + recurringBillingService, + type RecurringBillingServiceOptions, + type RescheduleInput, +} from './recurringBillingService.js'; diff --git a/backend/src/services/recurring-billing/recurring-billing.test.ts b/backend/src/services/recurring-billing/recurring-billing.test.ts new file mode 100644 index 00000000..8ffd2794 --- /dev/null +++ b/backend/src/services/recurring-billing/recurring-billing.test.ts @@ -0,0 +1,312 @@ +/** + * recurring-billing.test.ts — Issue #918: Recurring payment schedules with + * cron-based billing + * + * Covers cron resolution/validation, schedule creation, pause/resume/cancel, + * rescheduling, preview, and due-run invoicing (success, cap and failure). + */ +import { describe, expect, it, beforeEach } from 'vitest'; + +import { + DEFAULT_RECURRING_BILLING_CONFIG, + InMemoryRecurringBillingStore, + RecurringBillingService, + isValidCronExpression, + isValidTimezone, + nextRunAfter, + upcomingRuns, + validateRecurringSchedule, +} from './index.js'; +import type { + RecurringBillingEvent, + RecurringBillingPublisher, + RecurringInvoice, + RecurringInvoiceRepository, +} from './index.js'; + +class CapturingPublisher implements RecurringBillingPublisher { + events: RecurringBillingEvent[] = []; + publish(event: RecurringBillingEvent): void { + this.events.push(event); + } +} + +const NOW = new Date('2026-09-28T12:00:00.000Z'); + +function makeService( + options: { publisher?: RecurringBillingPublisher; clock?: { current: Date }; invoiceRepository?: RecurringInvoiceRepository } = {}, +) { + const clock = options.clock ?? { current: NOW }; + let counter = 0; + const store = new InMemoryRecurringBillingStore(); + const service = new RecurringBillingService({ + repository: store, + invoiceRepository: options.invoiceRepository ?? store, + publisher: options.publisher ?? { publish: () => undefined }, + now: () => clock.current, + idFactory: (() => { + const prefixes = ['sched', 'inv']; + return () => { + counter += 1; + return `${prefixes[counter % prefixes.length]}-${counter}`; + }; + })(), + }); + return { service, clock }; +} + +async function createDaily(service: RecurringBillingService, overrides: Record = {}) { + const result = await service.createSchedule({ + tenantId: 'tenant-1', + customerId: 'cust-1', + amount: 49.99, + currency: 'USD', + preset: 'daily', + ...overrides, + }); + if (!result.ok) throw new Error(`createSchedule failed: ${result.error.message}`); + return result.value; +} + +describe('cron helpers', () => { + it('validates cron expressions and timezones', () => { + expect(isValidCronExpression('0 0 * * *')).toBe(true); + expect(isValidCronExpression('not a cron')).toBe(false); + expect(isValidTimezone('UTC')).toBe(true); + expect(isValidTimezone('America/New_York')).toBe(true); + expect(isValidTimezone('Mars/Phobos')).toBe(false); + }); + + it('computes the next occurrence strictly after a reference time', () => { + expect(nextRunAfter('0 0 * * *', 'UTC', NOW)?.toISOString()).toBe('2026-09-29T00:00:00.000Z'); + expect(nextRunAfter('0 0 1 * *', 'UTC', NOW)?.toISOString()).toBe('2026-10-01T00:00:00.000Z'); + expect(nextRunAfter('invalid', 'UTC', NOW)).toBeNull(); + }); + + it('lists multiple upcoming runs', () => { + const runs = upcomingRuns('0 0 * * *', 'UTC', 3, NOW).map((date) => date.toISOString()); + expect(runs).toEqual([ + '2026-09-29T00:00:00.000Z', + '2026-09-30T00:00:00.000Z', + '2026-10-01T00:00:00.000Z', + ]); + expect(upcomingRuns('0 0 * * *', 'UTC', 0, NOW)).toEqual([]); + }); +}); + +describe('validateRecurringSchedule', () => { + const base = { tenantId: 'tenant-1', customerId: 'cust-1', amount: 10 }; + + it('applies defaults for currency, timezone and cadence', () => { + const result = validateRecurringSchedule(base); + expect(result.ok).toBe(true); + if (!result.ok) return; + expect(result.value.currency).toBe('USD'); + expect(result.value.timezone).toBe('UTC'); + expect(result.value.cronExpression).toBe('0 0 1 * *'); // monthly preset + }); + + it('expands presets and accepts raw cron', () => { + const weekly = validateRecurringSchedule({ ...base, preset: 'weekly' }); + expect(weekly.ok && weekly.value.cronExpression).toBe('0 0 * * 0'); + const raw = validateRecurringSchedule({ ...base, cronExpression: '*/15 * * * *' }); + expect(raw.ok && raw.value.cronExpression).toBe('*/15 * * * *'); + }); + + it('rejects missing tenant/customer', () => { + expect(validateRecurringSchedule({ ...base, tenantId: '' }).ok).toBe(false); + expect(validateRecurringSchedule({ ...base, customerId: '' }).ok).toBe(false); + }); + + it('rejects non-positive and out-of-range amounts', () => { + expect(validateRecurringSchedule({ ...base, amount: 0 }).ok).toBe(false); + expect(validateRecurringSchedule({ ...base, amount: -1 }).ok).toBe(false); + expect(validateRecurringSchedule({ ...base, amount: DEFAULT_RECURRING_BILLING_CONFIG.maxAmount + 1 }).ok).toBe(false); + }); + + it('rejects unsupported currencies', () => { + expect(validateRecurringSchedule({ ...base, currency: 'JPY' }).ok).toBe(false); + }); + + it('rejects invalid cron expressions', () => { + expect(validateRecurringSchedule({ ...base, cronExpression: 'nope' }).ok).toBe(false); + }); + + it('rejects an end date before the start date', () => { + expect( + validateRecurringSchedule({ + ...base, + startAt: '2026-10-01T00:00:00Z', + endAt: '2026-09-01T00:00:00Z', + }).ok, + ).toBe(false); + }); + + it('rejects a non-positive maxRuns', () => { + expect(validateRecurringSchedule({ ...base, maxRuns: 0 }).ok).toBe(false); + expect(validateRecurringSchedule({ ...base, maxRuns: 2.5 }).ok).toBe(false); + }); +}); + +describe('RecurringBillingService lifecycle', () => { + let publisher: CapturingPublisher; + let service: RecurringBillingService; + + beforeEach(() => { + publisher = new CapturingPublisher(); + service = makeService({ publisher }).service; + }); + + it('creates an active schedule with the first run in the future', async () => { + const schedule = await createDaily(service); + expect(schedule.status).toBe('active'); + expect(schedule.cronExpression).toBe('0 0 * * *'); + expect(schedule.nextRunAt).toBe('2026-09-29T00:00:00.000Z'); + expect(schedule.runCount).toBe(0); + expect(publisher.events.map((event) => event.type)).toContain('recurring_schedule.created'); + }); + + it('honours an explicit start date as the first run boundary', async () => { + const schedule = await createDaily(service, { startAt: '2026-10-05T00:00:00.000Z' }); + expect(schedule.nextRunAt).toBe('2026-10-05T00:00:00.000Z'); + }); + + it('scopes reads to the owning tenant', async () => { + const schedule = await createDaily(service); + const other = await service.getSchedule('tenant-2', schedule.id); + expect(other.ok).toBe(false); + }); + + it('filters schedules by status', async () => { + const schedule = await createDaily(service); + await service.pauseSchedule('tenant-1', schedule.id); + + const active = await service.listSchedules('tenant-1', { status: 'active' }); + expect(active.ok && active.value).toHaveLength(0); + const paused = await service.listSchedules('tenant-1', { status: 'paused' }); + expect(paused.ok && paused.value).toHaveLength(1); + }); + + it('pauses and resumes a schedule', async () => { + const schedule = await createDaily(service); + const paused = await service.pauseSchedule('tenant-1', schedule.id); + expect(paused.ok && paused.value.status).toBe('paused'); + + // Resume recomputes the next run from "now". + const resumed = await service.resumeSchedule('tenant-1', schedule.id); + expect(resumed.ok && resumed.value.status).toBe('active'); + expect(resumed.ok && resumed.value.nextRunAt).toBe('2026-09-29T00:00:00.000Z'); + expect(publisher.events.map((event) => event.type)).toContain('recurring_schedule.resumed'); + }); + + it('rejects pausing a non-active schedule', async () => { + const schedule = await createDaily(service); + await service.pauseSchedule('tenant-1', schedule.id); + const again = await service.pauseSchedule('tenant-1', schedule.id); + expect(again.ok).toBe(false); + if (again.ok) return; + expect(again.error.code).toBe('CONFLICT'); + }); + + it('cancels a schedule and clears its next run', async () => { + const schedule = await createDaily(service); + const cancelled = await service.cancelSchedule('tenant-1', schedule.id, 'customer asked'); + expect(cancelled.ok && cancelled.value.status).toBe('cancelled'); + expect(cancelled.ok && cancelled.value.nextRunAt).toBeNull(); + expect(publisher.events.map((event) => event.type)).toContain('recurring_schedule.cancelled'); + }); + + it('reschedules to a new cron cadence', async () => { + const schedule = await createDaily(service); + const rescheduled = await service.reschedule('tenant-1', schedule.id, { preset: 'monthly' }); + expect(rescheduled.ok && rescheduled.value.cronExpression).toBe('0 0 1 * *'); + expect(rescheduled.ok && rescheduled.value.nextRunAt).toBe('2026-10-01T00:00:00.000Z'); + }); + + it('previews upcoming runs', async () => { + const schedule = await createDaily(service); + const preview = await service.previewUpcoming('tenant-1', schedule.id, 2); + expect(preview.ok && preview.value.runs).toEqual([ + '2026-09-29T00:00:00.000Z', + '2026-09-30T00:00:00.000Z', + ]); + }); +}); + +describe('RecurringBillingService.runDue', () => { + it('generates an invoice and advances the schedule', async () => { + const publisher = new CapturingPublisher(); + const { service, clock } = makeService({ publisher }); + const schedule = await createDaily(service); + + clock.current = new Date('2026-09-29T00:00:00.000Z'); + const result = await service.runDue(clock.current); + + expect(result.ok && result.value.processed).toBe(1); + if (!result.ok) return; + expect(result.value.invoices).toHaveLength(1); + expect(result.value.invoices[0].dueAt).toBe('2026-09-29T00:00:00.000Z'); + expect(result.value.invoices[0].status).toBe('pending'); + + const refreshed = await service.getSchedule('tenant-1', schedule.id); + expect(refreshed.ok && refreshed.value.runCount).toBe(1); + expect(refreshed.ok && refreshed.value.nextRunAt).toBe('2026-09-30T00:00:00.000Z'); + expect(publisher.events.map((event) => event.type)).toContain('recurring_invoice.generated'); + + const invoices = await service.listInvoices('tenant-1', schedule.id); + expect(invoices.ok && invoices.value).toHaveLength(1); + }); + + it('completes a schedule once maxRuns is reached', async () => { + const { service, clock } = makeService(); + const schedule = await createDaily(service, { maxRuns: 2 }); + + clock.current = new Date('2026-09-29T00:00:00.000Z'); + await service.runDue(clock.current); + clock.current = new Date('2026-09-30T00:00:00.000Z'); + const second = await service.runDue(clock.current); + + expect(second.ok && second.value.completedScheduleIds).toEqual([schedule.id]); + const refreshed = await service.getSchedule('tenant-1', schedule.id); + expect(refreshed.ok && refreshed.value.status).toBe('completed'); + expect(refreshed.ok && refreshed.value.nextRunAt).toBeNull(); + }); + + it('completes a schedule when the next run passes endAt', async () => { + const { service, clock } = makeService(); + const schedule = await createDaily(service, { endAt: '2026-09-29T00:00:00.000Z' }); + + clock.current = new Date('2026-09-29T00:00:00.000Z'); + const result = await service.runDue(clock.current); + + expect(result.ok && result.value.completedScheduleIds).toEqual([schedule.id]); + const refreshed = await service.getSchedule('tenant-1', schedule.id); + expect(refreshed.ok && refreshed.value.status).toBe('completed'); + }); + + it('does nothing when no schedule is due', async () => { + const { service } = makeService(); + await createDaily(service); + const result = await service.runDue(NOW); + expect(result.ok && result.value.processed).toBe(0); + }); + + it('emits a failure event when invoice persistence fails', async () => { + const publisher = new CapturingPublisher(); + const failingInvoiceRepo: RecurringInvoiceRepository = { + createInvoice: () => Promise.reject(new Error('db down')), + listBySchedule: async () => [] as RecurringInvoice[], + }; + const { service, clock } = makeService({ publisher, invoiceRepository: failingInvoiceRepo }); + const schedule = await createDaily(service); + + clock.current = new Date('2026-09-29T00:00:00.000Z'); + const result = await service.runDue(clock.current); + + expect(result.ok && result.value.invoices).toHaveLength(0); + expect(publisher.events.map((event) => event.type)).toContain('recurring_invoice.failed'); + // The schedule is left untouched so the run can be retried. + const refreshed = await service.getSchedule('tenant-1', schedule.id); + expect(refreshed.ok && refreshed.value.nextRunAt).toBe('2026-09-29T00:00:00.000Z'); + }); +}); diff --git a/backend/src/services/recurring-billing/recurringBillingService.ts b/backend/src/services/recurring-billing/recurringBillingService.ts new file mode 100644 index 00000000..28e4d674 --- /dev/null +++ b/backend/src/services/recurring-billing/recurringBillingService.ts @@ -0,0 +1,417 @@ +/** + * recurringBillingService.ts — Issue #918: Recurring payment schedules with + * cron-based billing + * + * Coordinates schedule lifecycle and due-run invoicing. All cron/date maths + * lives in `schedule.ts`; this class only sequences state transitions and + * persistence so it stays easy to test with a stubbed store. + */ +import { randomUUID } from 'node:crypto'; + +import { BaseService } from '../BaseService.js'; +import type { Result } from '../../lib/result.js'; +import { + DEFAULT_RECURRING_BILLING_CONFIG, + PRESET_CRONS, + nextRunAfter, + upcomingRuns, + validateRecurringSchedule, +} from './schedule.js'; +import { + InMemoryRecurringBillingStore, + type RecurringInvoiceRepository, + type RecurringScheduleRepository, +} from './store.js'; +import type { + BillingPreset, + CreateRecurringScheduleInput, + RecurringBillingConfig, + RecurringBillingEvent, + RecurringBillingPublisher, + RecurringInvoice, + RecurringSchedule, + RecurringScheduleFilter, + RunDueResult, +} from './types.js'; + +export interface RecurringBillingServiceOptions { + repository?: RecurringScheduleRepository; + invoiceRepository?: RecurringInvoiceRepository; + config?: RecurringBillingConfig; + publisher?: RecurringBillingPublisher; + now?: () => Date; + idFactory?: () => string; +} + +export interface RescheduleInput { + cronExpression?: string; + preset?: BillingPreset; + timezone?: string; +} + +/** Default publisher that logs; swap for an event-bus adapter in production. */ +class LoggingRecurringBillingPublisher implements RecurringBillingPublisher { + publish(event: RecurringBillingEvent): void { + console.info('[recurring-billing] event', event.type, event.scheduleId); + } +} + +export class RecurringBillingService extends BaseService { + private readonly repository: RecurringScheduleRepository; + private readonly invoiceRepository: RecurringInvoiceRepository; + private readonly config: RecurringBillingConfig; + private readonly publisher: RecurringBillingPublisher; + private readonly now: () => Date; + private readonly idFactory: () => string; + + constructor(options: RecurringBillingServiceOptions = {}) { + super(); + const store = new InMemoryRecurringBillingStore(); + this.repository = options.repository ?? store; + this.invoiceRepository = options.invoiceRepository ?? store; + this.config = options.config ?? DEFAULT_RECURRING_BILLING_CONFIG; + this.publisher = options.publisher ?? new LoggingRecurringBillingPublisher(); + this.now = options.now ?? (() => new Date()); + this.idFactory = options.idFactory ?? (() => randomUUID()); + } + + getConfig(): RecurringBillingConfig { + return { ...this.config, supportedCurrencies: [...this.config.supportedCurrencies] }; + } + + /** Create a schedule and compute its first billing instant. */ + async createSchedule(input: CreateRecurringScheduleInput): Promise> { + const validated = validateRecurringSchedule(input, this.config); + if (!validated.ok) { + return validated; + } + + const request = validated.value; + const now = this.now(); + const timestamp = now.toISOString(); + + // The schedule starts at max(startAt, now): a past start date must not + // backfill, a future one delays the first invoice. + const effectiveFrom = request.startAt.getTime() > now.getTime() ? request.startAt : now; + // Inclusive of `effectiveFrom` itself when it lands exactly on a tick. + const first = nextRunAfter( + request.cronExpression, + request.timezone, + new Date(effectiveFrom.getTime() - 1), + ); + + const withinEnd = first != null && (request.endAt == null || first.getTime() <= request.endAt.getTime()); + const nextRunAt = withinEnd && first ? first.toISOString() : null; + + const schedule: RecurringSchedule = { + id: this.idFactory(), + tenantId: request.tenantId, + customerId: request.customerId, + merchantId: request.merchantId ?? null, + name: request.name ?? null, + cronExpression: request.cronExpression, + timezone: request.timezone, + amount: request.amount, + currency: request.currency, + status: nextRunAt ? 'active' : 'completed', + startAt: request.startAt.toISOString(), + endAt: request.endAt ? request.endAt.toISOString() : null, + maxRuns: request.maxRuns, + runCount: 0, + lastRunAt: null, + nextRunAt, + metadata: request.metadata, + createdAt: timestamp, + updatedAt: timestamp, + }; + + const saved = await this.repository.create(schedule); + await this.publisher.publish({ + type: 'recurring_schedule.created', + scheduleId: saved.id, + tenantId: saved.tenantId, + occurredAt: timestamp, + data: { cronExpression: saved.cronExpression, nextRunAt: saved.nextRunAt }, + }); + return this.ok(saved); + } + + async getSchedule(tenantId: string, scheduleId: string): Promise> { + const schedule = await this.repository.findById(scheduleId); + if (!schedule || schedule.tenantId !== tenantId) { + return this.notFoundFailure('Recurring schedule', scheduleId); + } + return this.ok(schedule); + } + + async listSchedules( + tenantId: string, + filter?: RecurringScheduleFilter, + ): Promise> { + if (!tenantId) { + return this.validationFailure('tenantId is required'); + } + return this.ok(await this.repository.list(tenantId, filter)); + } + + /** Preview the next `count` billing instants without mutating the schedule. */ + async previewUpcoming( + tenantId: string, + scheduleId: string, + count = 5, + ): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + const bounded = Math.min(Math.max(1, count), this.config.maxUpcomingPreview); + const schedule = found.value; + const from = schedule.nextRunAt ? new Date(new Date(schedule.nextRunAt).getTime() - 1) : this.now(); + const runs = upcomingRuns(schedule.cronExpression, schedule.timezone, bounded, from).map((date) => + date.toISOString(), + ); + return this.ok({ runs }); + } + + async pauseSchedule(tenantId: string, scheduleId: string): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + const schedule = found.value; + if (schedule.status !== 'active') { + return this.conflictFailure(`Cannot pause a ${schedule.status} schedule`); + } + schedule.status = 'paused'; + schedule.updatedAt = this.now().toISOString(); + const saved = await this.repository.update(schedule); + await this.publisher.publish({ + type: 'recurring_schedule.paused', + scheduleId: saved.id, + tenantId: saved.tenantId, + occurredAt: saved.updatedAt, + }); + return this.ok(saved); + } + + async resumeSchedule(tenantId: string, scheduleId: string): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + const schedule = found.value; + if (schedule.status !== 'paused') { + return this.conflictFailure(`Cannot resume a ${schedule.status} schedule`); + } + + const now = this.now(); + const next = nextRunAfter(schedule.cronExpression, schedule.timezone, now); + const withinEnd = + next != null && (schedule.endAt == null || next.getTime() <= new Date(schedule.endAt).getTime()); + const runsExhausted = schedule.maxRuns != null && schedule.runCount >= schedule.maxRuns; + + schedule.status = withinEnd && !runsExhausted ? 'active' : 'completed'; + schedule.nextRunAt = schedule.status === 'active' && next ? next.toISOString() : null; + schedule.updatedAt = now.toISOString(); + + const saved = await this.repository.update(schedule); + await this.publisher.publish({ + type: 'recurring_schedule.resumed', + scheduleId: saved.id, + tenantId: saved.tenantId, + occurredAt: saved.updatedAt, + data: { nextRunAt: saved.nextRunAt }, + }); + return this.ok(saved); + } + + async cancelSchedule( + tenantId: string, + scheduleId: string, + reason?: string, + ): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + const schedule = found.value; + if (schedule.status === 'cancelled') { + return this.conflictFailure('Schedule is already cancelled'); + } + if (schedule.status === 'completed') { + return this.conflictFailure('Cannot cancel a completed schedule'); + } + + schedule.status = 'cancelled'; + schedule.nextRunAt = null; + schedule.updatedAt = this.now().toISOString(); + + const saved = await this.repository.update(schedule); + await this.publisher.publish({ + type: 'recurring_schedule.cancelled', + scheduleId: saved.id, + tenantId: saved.tenantId, + occurredAt: saved.updatedAt, + data: { reason: reason ?? null }, + }); + return this.ok(saved); + } + + /** Change the cadence of an existing schedule and recompute its next run. */ + async reschedule( + tenantId: string, + scheduleId: string, + input: RescheduleInput, + ): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + const schedule = found.value; + if (schedule.status === 'cancelled') { + return this.conflictFailure('Cannot reschedule a cancelled schedule'); + } + + const effectiveCron = input.cronExpression ?? (input.preset ? PRESET_CRONS[input.preset] : schedule.cronExpression); + + const validated = validateRecurringSchedule( + { + tenantId: schedule.tenantId, + customerId: schedule.customerId, + amount: schedule.amount, + currency: schedule.currency, + cronExpression: effectiveCron, + timezone: input.timezone ?? schedule.timezone, + startAt: schedule.startAt, + endAt: schedule.endAt ?? undefined, + maxRuns: schedule.maxRuns ?? undefined, + merchantId: schedule.merchantId ?? undefined, + name: schedule.name ?? undefined, + metadata: schedule.metadata, + }, + this.config, + ); + if (!validated.ok) { + return validated; + } + + const request = validated.value; + const now = this.now(); + const next = nextRunAfter(request.cronExpression, request.timezone, now); + const withinEnd = + next != null && (request.endAt == null || next.getTime() <= request.endAt.getTime()); + const runsExhausted = schedule.maxRuns != null && schedule.runCount >= schedule.maxRuns; + + schedule.cronExpression = request.cronExpression; + schedule.timezone = request.timezone; + schedule.status = withinEnd && !runsExhausted ? 'active' : 'completed'; + schedule.nextRunAt = schedule.status === 'active' && next ? next.toISOString() : null; + schedule.updatedAt = now.toISOString(); + + const saved = await this.repository.update(schedule); + await this.publisher.publish({ + type: 'recurring_schedule.rescheduled', + scheduleId: saved.id, + tenantId: saved.tenantId, + occurredAt: saved.updatedAt, + data: { cronExpression: saved.cronExpression, nextRunAt: saved.nextRunAt }, + }); + return this.ok(saved); + } + + /** + * Generate invoices for every active schedule whose next run is due. + * One invoice is produced per due schedule per call; a schedule that is + * behind catch up on subsequent sweeps. + */ + async runDue(at: Date | string = new Date()): Promise> { + const now = at instanceof Date ? at : new Date(at); + if (Number.isNaN(now.getTime())) { + return this.validationFailure('runDue timestamp must be a valid date'); + } + + const due = await this.repository.findDue(undefined, now); + const invoices: RecurringInvoice[] = []; + const completedScheduleIds: string[] = []; + + for (const schedule of due) { + if (!schedule.nextRunAt) { + continue; + } + const dueAt = new Date(schedule.nextRunAt); + + const invoice: RecurringInvoice = { + id: this.idFactory(), + scheduleId: schedule.id, + tenantId: schedule.tenantId, + amount: schedule.amount, + currency: schedule.currency, + status: 'pending', + periodStart: schedule.lastRunAt ?? schedule.createdAt, + dueAt: dueAt.toISOString(), + createdAt: now.toISOString(), + }; + + try { + await this.invoiceRepository.createInvoice(invoice); + } catch (error) { + await this.publisher.publish({ + type: 'recurring_invoice.failed', + scheduleId: schedule.id, + tenantId: schedule.tenantId, + occurredAt: now.toISOString(), + data: { dueAt: invoice.dueAt, error: error instanceof Error ? error.message : String(error) }, + }); + continue; + } + + invoices.push(invoice); + schedule.runCount += 1; + schedule.lastRunAt = dueAt.toISOString(); + + const next = nextRunAfter(schedule.cronExpression, schedule.timezone, dueAt); + const runsExhausted = schedule.maxRuns != null && schedule.runCount >= schedule.maxRuns; + const pastEnd = + next == null || (schedule.endAt != null && next.getTime() > new Date(schedule.endAt).getTime()); + + if (runsExhausted || pastEnd) { + schedule.status = 'completed'; + schedule.nextRunAt = null; + completedScheduleIds.push(schedule.id); + } else { + schedule.nextRunAt = next.toISOString(); + } + schedule.updatedAt = now.toISOString(); + await this.repository.update(schedule); + + await this.publisher.publish({ + type: 'recurring_invoice.generated', + scheduleId: schedule.id, + tenantId: schedule.tenantId, + occurredAt: now.toISOString(), + data: { invoiceId: invoice.id, dueAt: invoice.dueAt, runCount: schedule.runCount }, + }); + if (schedule.status === 'completed') { + await this.publisher.publish({ + type: 'recurring_schedule.completed', + scheduleId: schedule.id, + tenantId: schedule.tenantId, + occurredAt: schedule.updatedAt, + data: { runCount: schedule.runCount }, + }); + } + } + + return this.ok({ processed: due.length, invoices, completedScheduleIds }); + } + + async listInvoices(tenantId: string, scheduleId: string): Promise> { + const found = await this.getSchedule(tenantId, scheduleId); + if (!found.ok) { + return found; + } + return this.ok(await this.invoiceRepository.listBySchedule(scheduleId)); + } +} + +export const recurringBillingService = new RecurringBillingService(); diff --git a/backend/src/services/recurring-billing/schedule.ts b/backend/src/services/recurring-billing/schedule.ts new file mode 100644 index 00000000..cee5a98e --- /dev/null +++ b/backend/src/services/recurring-billing/schedule.ts @@ -0,0 +1,204 @@ +/** + * schedule.ts — Issue #918: Recurring payment schedules with cron-based billing + * + * Pure, side-effect-free helpers that validate a recurring-schedule request and + * resolve its cron cadence. All date maths lives here so it is directly + * unit-testable and excluded from clock/IO concerns. + */ +import cronParser from 'cron-parser'; + +import { err, ok, type Result } from '../../lib/result.js'; +import type { + BillingPreset, + CreateRecurringScheduleInput, + NormalizedRecurringSchedule, + RecurringBillingConfig, +} from './types.js'; + +export const DEFAULT_RECURRING_BILLING_CONFIG: RecurringBillingConfig = { + supportedCurrencies: ['USD', 'EUR', 'GBP', 'XLM', 'USDC'], + minAmount: 0.01, + maxAmount: 1_000_000, + defaultCurrency: 'USD', + defaultTimezone: 'UTC', + defaultPreset: 'monthly', + maxUpcomingPreview: 50, +}; + +/** Canonical cron expressions for the supported presets. */ +export const PRESET_CRONS: Record = { + hourly: '0 * * * *', + daily: '0 0 * * *', + weekly: '0 0 * * 0', + monthly: '0 0 1 * *', + yearly: '0 0 1 1 *', +}; + +/** Round to 2 decimals using a half-up rule that tolerates float error. */ +export function roundCurrency(value: number): number { + return Math.round((value + Number.EPSILON) * 100) / 100; +} + +/** True when `expression` is a cron string cron-parser can evaluate. */ +export function isValidCronExpression(expression: string, timezone = 'UTC'): boolean { + if (!expression || typeof expression !== 'string') { + return false; + } + try { + cronParser.parseExpression(expression, { tz: timezone }); + return true; + } catch { + return false; + } +} + +/** + * True when `timezone` is a valid IANA zone. Validated through `Intl` because + * cron-parser silently accepts unknown zones instead of rejecting them. + */ +export function isValidTimezone(timezone: string): boolean { + if (!timezone || typeof timezone !== 'string') { + return false; + } + try { + new Intl.DateTimeFormat('en-US', { timeZone: timezone }); + return true; + } catch { + return false; + } +} + +/** + * First cron occurrence strictly after `after`. Returns null when the + * expression/timezone cannot be evaluated. + */ +export function nextRunAfter(cronExpression: string, timezone: string, after: Date): Date | null { + try { + const interval = cronParser.parseExpression(cronExpression, { tz: timezone, currentDate: after }); + return interval.next().toDate(); + } catch { + return null; + } +} + +/** The next `count` occurrences strictly after `after`. */ +export function upcomingRuns( + cronExpression: string, + timezone: string, + count: number, + after: Date, +): Date[] { + const runs: Date[] = []; + if (count <= 0) { + return runs; + } + try { + const interval = cronParser.parseExpression(cronExpression, { tz: timezone, currentDate: after }); + for (let i = 0; i < count; i += 1) { + runs.push(interval.next().toDate()); + } + } catch { + return []; + } + return runs; +} + +/** Resolve the effective cron expression for a request (raw cron wins). */ +export function resolveCronExpression(input: Pick, config = DEFAULT_RECURRING_BILLING_CONFIG): string { + if (input.cronExpression) { + return input.cronExpression.trim(); + } + return PRESET_CRONS[input.preset ?? config.defaultPreset]; +} + +/** + * Validate and normalise a create request. Returns a `Result` so callers can + * surface a precise 400 instead of throwing deep inside the service. + */ +export function validateRecurringSchedule( + input: CreateRecurringScheduleInput, + config: RecurringBillingConfig = DEFAULT_RECURRING_BILLING_CONFIG, +): Result { + if (!input || typeof input !== 'object') { + return err({ code: 'VALIDATION_ERROR', message: 'Schedule request body is required', statusCode: 400 }); + } + + if (!input.tenantId || typeof input.tenantId !== 'string') { + return err({ code: 'VALIDATION_ERROR', message: 'tenantId is required', statusCode: 400 }); + } + if (!input.customerId || typeof input.customerId !== 'string') { + return err({ code: 'VALIDATION_ERROR', message: 'customerId is required', statusCode: 400 }); + } + + const amount = roundCurrency(Number(input.amount)); + if (!Number.isFinite(amount) || amount <= 0) { + return err({ code: 'VALIDATION_ERROR', message: 'amount must be a positive number', statusCode: 400 }); + } + if (amount < config.minAmount) { + return err({ code: 'VALIDATION_ERROR', message: `amount must be at least ${config.minAmount}`, statusCode: 400 }); + } + if (amount > config.maxAmount) { + return err({ code: 'VALIDATION_ERROR', message: `amount must not exceed ${config.maxAmount}`, statusCode: 400 }); + } + + const currency = (input.currency ?? config.defaultCurrency).toUpperCase(); + if (!config.supportedCurrencies.includes(currency)) { + return err({ + code: 'VALIDATION_ERROR', + message: `currency must be one of: ${config.supportedCurrencies.join(', ')}`, + statusCode: 400, + }); + } + + const cronExpression = resolveCronExpression(input, config); + const timezone = input.timezone ?? config.defaultTimezone; + if (!isValidTimezone(timezone)) { + return err({ code: 'VALIDATION_ERROR', message: `timezone "${timezone}" is not a valid IANA zone`, statusCode: 400 }); + } + if (!isValidCronExpression(cronExpression, timezone)) { + return err({ + code: 'VALIDATION_ERROR', + message: `cronExpression "${cronExpression}" is invalid${input.cronExpression ? '' : ` for preset "${input.preset ?? config.defaultPreset}"`}`, + statusCode: 400, + }); + } + + const startAt = input.startAt ? new Date(input.startAt) : new Date(); + if (Number.isNaN(startAt.getTime())) { + return err({ code: 'VALIDATION_ERROR', message: 'startAt must be a valid date', statusCode: 400 }); + } + + let endAt: Date | null = null; + if (input.endAt != null) { + endAt = new Date(input.endAt); + if (Number.isNaN(endAt.getTime())) { + return err({ code: 'VALIDATION_ERROR', message: 'endAt must be a valid date', statusCode: 400 }); + } + if (endAt.getTime() < startAt.getTime()) { + return err({ code: 'VALIDATION_ERROR', message: 'endAt must be on or after startAt', statusCode: 400 }); + } + } + + let maxRuns: number | null = null; + if (input.maxRuns != null) { + if (!Number.isInteger(input.maxRuns) || input.maxRuns < 1) { + return err({ code: 'VALIDATION_ERROR', message: 'maxRuns must be a positive integer', statusCode: 400 }); + } + maxRuns = input.maxRuns; + } + + return ok({ + tenantId: input.tenantId, + customerId: input.customerId, + amount, + currency, + cronExpression, + timezone, + startAt, + endAt, + maxRuns, + merchantId: input.merchantId, + name: input.name, + metadata: input.metadata, + }); +} diff --git a/backend/src/services/recurring-billing/store.ts b/backend/src/services/recurring-billing/store.ts new file mode 100644 index 00000000..debc4a6f --- /dev/null +++ b/backend/src/services/recurring-billing/store.ts @@ -0,0 +1,97 @@ +/** + * store.ts — Issue #918: Recurring payment schedules with cron-based billing + * + * Persistence boundary for recurring schedules and their generated invoices. + * The service depends on these interfaces, so production can back them with + * Prisma while tests use the deterministic in-memory implementation below. + */ +import type { RecurringInvoice, RecurringSchedule, RecurringScheduleFilter } from './types.js'; + +export interface RecurringScheduleRepository { + create(schedule: RecurringSchedule): Promise; + findById(id: string): Promise; + update(schedule: RecurringSchedule): Promise; + list(tenantId: string, filter?: RecurringScheduleFilter): Promise; + /** + * Active schedules whose `nextRunAt` is at or before `before`. When + * `tenantId` is omitted all tenants are scanned (used by the billing sweep). + */ + findDue(tenantId: string | undefined, before: Date): Promise; +} + +export interface RecurringInvoiceRepository { + /** Named distinctly from the schedule `create` so one class can back both. */ + createInvoice(invoice: RecurringInvoice): Promise; + listBySchedule(scheduleId: string): Promise; +} + +function cloneSchedule(schedule: RecurringSchedule): RecurringSchedule { + return { + ...schedule, + metadata: schedule.metadata ? { ...schedule.metadata } : schedule.metadata, + }; +} + +export class InMemoryRecurringBillingStore + implements RecurringScheduleRepository, RecurringInvoiceRepository +{ + private readonly schedules = new Map(); + private readonly invoices = new Map(); + + async create(schedule: RecurringSchedule): Promise { + const stored = cloneSchedule(schedule); + this.schedules.set(stored.id, stored); + if (!this.invoices.has(stored.id)) { + this.invoices.set(stored.id, []); + } + return cloneSchedule(stored); + } + + async findById(id: string): Promise { + const found = this.schedules.get(id); + return found ? cloneSchedule(found) : null; + } + + async update(schedule: RecurringSchedule): Promise { + const stored = cloneSchedule(schedule); + this.schedules.set(stored.id, stored); + return cloneSchedule(stored); + } + + async list(tenantId: string, filter?: RecurringScheduleFilter): Promise { + return Array.from(this.schedules.values()) + .filter((schedule) => schedule.tenantId === tenantId) + .filter((schedule) => (filter?.status ? schedule.status === filter.status : true)) + .filter((schedule) => (filter?.customerId ? schedule.customerId === filter.customerId : true)) + .filter((schedule) => (filter?.merchantId ? schedule.merchantId === filter.merchantId : true)) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt) || a.id.localeCompare(b.id)) + .map(cloneSchedule); + } + + async findDue(tenantId: string | undefined, before: Date): Promise { + const cutoff = before.getTime(); + return Array.from(this.schedules.values()) + .filter((schedule) => (tenantId ? schedule.tenantId === tenantId : true)) + .filter((schedule) => schedule.status === 'active') + .filter((schedule) => schedule.nextRunAt != null && new Date(schedule.nextRunAt).getTime() <= cutoff) + .sort((a, b) => (a.nextRunAt ?? '').localeCompare(b.nextRunAt ?? '') || a.id.localeCompare(b.id)) + .map(cloneSchedule); + } + + async createInvoice(invoice: RecurringInvoice): Promise { + const existing = this.invoices.get(invoice.scheduleId) ?? []; + const stored = { ...invoice }; + existing.push(stored); + this.invoices.set(invoice.scheduleId, existing); + return { ...stored }; + } + + async listBySchedule(scheduleId: string): Promise { + return (this.invoices.get(scheduleId) ?? []).map((invoice) => ({ ...invoice })); + } + + /** Test/introspection helper: every invoice across all schedules. */ + allInvoices(): RecurringInvoice[] { + return Array.from(this.invoices.values()).flatMap((list) => list.map((invoice) => ({ ...invoice }))); + } +} diff --git a/backend/src/services/recurring-billing/types.ts b/backend/src/services/recurring-billing/types.ts new file mode 100644 index 00000000..5c1875bc --- /dev/null +++ b/backend/src/services/recurring-billing/types.ts @@ -0,0 +1,137 @@ +/** + * types.ts — Issue #918: Recurring payment schedules with cron-based billing + * + * Domain types for billing a customer on a cron cadence. A `RecurringSchedule` + * describes *who* is billed *how much* and *when* (a cron expression plus an + * IANA timezone); each due run materialises a `RecurringInvoice`. + */ + +/** Convenience presets that expand to a canonical cron expression. */ +export type BillingPreset = 'hourly' | 'daily' | 'weekly' | 'monthly' | 'yearly'; + +/** Lifecycle of a recurring schedule. */ +export type RecurringScheduleStatus = 'active' | 'paused' | 'cancelled' | 'completed'; + +/** Lifecycle of a generated invoice. */ +export type RecurringInvoiceStatus = 'pending' | 'paid' | 'failed' | 'void'; + +export interface RecurringSchedule { + id: string; + tenantId: string; + customerId: string; + merchantId?: string | null; + name?: string | null; + /** Cron expression that defines the billing cadence. */ + cronExpression: string; + /** IANA timezone the cron expression is evaluated in. */ + timezone: string; + amount: number; + currency: string; + status: RecurringScheduleStatus; + /** When the schedule becomes eligible to bill (ISO-8601). */ + startAt: string; + /** Optional end boundary; no invoice is generated after this instant. */ + endAt?: string | null; + /** Optional cap on the number of invoices generated. */ + maxRuns?: number | null; + runCount: number; + lastRunAt?: string | null; + /** Next invoice instant, or null once no further runs are possible. */ + nextRunAt: string | null; + metadata?: Record; + createdAt: string; + updatedAt: string; +} + +export interface CreateRecurringScheduleInput { + tenantId: string; + customerId: string; + amount: number; + currency?: string; + /** Raw cron expression. Mutually exclusive with `preset`. */ + cronExpression?: string; + /** Preset cadence. Ignored when `cronExpression` is provided. */ + preset?: BillingPreset; + timezone?: string; + startAt?: string | Date; + endAt?: string | Date; + maxRuns?: number; + merchantId?: string; + name?: string; + metadata?: Record; +} + +export interface NormalizedRecurringSchedule { + tenantId: string; + customerId: string; + amount: number; + currency: string; + cronExpression: string; + timezone: string; + startAt: Date; + endAt: Date | null; + maxRuns: number | null; + merchantId?: string; + name?: string; + metadata?: Record; +} + +export interface RecurringInvoice { + id: string; + scheduleId: string; + tenantId: string; + amount: number; + currency: string; + status: RecurringInvoiceStatus; + /** Start of the billed period (previous run, or schedule creation). */ + periodStart: string; + /** When the invoice falls due (the schedule's `nextRunAt`). */ + dueAt: string; + createdAt: string; +} + +export interface RecurringScheduleFilter { + status?: RecurringScheduleStatus; + customerId?: string; + merchantId?: string; +} + +export interface RunDueResult { + processed: number; + invoices: RecurringInvoice[]; + completedScheduleIds: string[]; +} + +export interface RecurringBillingConfig { + supportedCurrencies: string[]; + minAmount: number; + maxAmount: number; + defaultCurrency: string; + defaultTimezone: string; + defaultPreset: BillingPreset; + /** Upper bound for `previewUpcoming`. */ + maxUpcomingPreview: number; +} + +export type RecurringBillingEventType = + | 'recurring_schedule.created' + | 'recurring_schedule.paused' + | 'recurring_schedule.resumed' + | 'recurring_schedule.cancelled' + | 'recurring_schedule.completed' + | 'recurring_schedule.rescheduled' + | 'recurring_invoice.generated' + | 'recurring_invoice.failed'; + +export interface RecurringBillingEvent { + type: RecurringBillingEventType; + scheduleId: string; + tenantId: string; + occurredAt: string; + data?: Record; +} + +/** Pluggable publisher so the service stays transport-agnostic. */ +export interface RecurringBillingPublisher { + publish(event: RecurringBillingEvent): Promise | void; +} diff --git a/backend/src/services/split-payments/allocation.ts b/backend/src/services/split-payments/allocation.ts new file mode 100644 index 00000000..b31629cd --- /dev/null +++ b/backend/src/services/split-payments/allocation.ts @@ -0,0 +1,241 @@ +/** + * allocation.ts — Issue #917: Split payments between multiple recipients + * + * Pure, side-effect-free helpers that validate a split request and allocate a + * payment down to the last minor unit. Shares are converted to integer basis + * points and distributed with the largest-remainder (Hamilton) method, so: + * + * platformFeeMinor + Σ share.amountMinor === totalMinor + * + * and no money is created or lost to rounding. + */ +import { err, ok, type Result } from '../../lib/result.js'; +import type { + CreateSplitPlanInput, + NormalizedSplitPlan, + SplitAllocation, + SplitAllocationShare, + SplitConfig, + SplitRecipient, + SplitRecipientInput, +} from './types.js'; + +export const DEFAULT_SPLIT_CONFIG: SplitConfig = { + supportedCurrencies: ['USD', 'EUR', 'GBP', 'XLM', 'USDC'], + defaultCurrency: 'USD', + maxRecipients: 25, +}; + +const BPS_TOTAL = 10_000; + +/** Round to 2 decimals using a half-up rule that tolerates float error. */ +export function roundCurrency(value: number): number { + return Math.round((value + Number.EPSILON) * 100) / 100; +} + +/** Convert a percentage to integer basis points (1 bp = 0.01%). */ +export function toBasisPoints(percentage: number): number { + return Math.round(percentage * 100); +} + +/** Convert basis points back to a percentage. */ +export function fromBasisPoints(bps: number): number { + return bps / 100; +} + +/** + * Largest-remainder allocation of `totalMinor` across `weights`. + * The returned integers sum to exactly `totalMinor`. + */ +export function allocateMinorUnits(totalMinor: number, weights: number[]): number[] { + const totalWeight = weights.reduce((sum, weight) => sum + weight, 0); + if (totalMinor <= 0 || totalWeight <= 0) { + return weights.map(() => 0); + } + + const exact = weights.map((weight) => (totalMinor * weight) / totalWeight); + const minors = exact.map((value) => Math.floor(value)); + let remaining = totalMinor - minors.reduce((sum, value) => sum + value, 0); + + const byRemainder = exact + .map((value, index) => ({ index, remainder: value - Math.floor(value) })) + .sort((a, b) => b.remainder - a.remainder || a.index - b.index); + + for (const entry of byRemainder) { + if (remaining <= 0) break; + minors[entry.index] += 1; + remaining -= 1; + } + + return minors; +} + +/** + * Validate recipient shares (plus the optional platform fee). The declared + * percentages must sum to exactly 100 so the whole payment is allocated. + */ +export function validateSplitPlan( + input: CreateSplitPlanInput, + config: SplitConfig = DEFAULT_SPLIT_CONFIG, +): Result { + if (!input || typeof input !== 'object') { + return err({ code: 'VALIDATION_ERROR', message: 'Split request body is required', statusCode: 400 }); + } + if (!input.tenantId || typeof input.tenantId !== 'string') { + return err({ code: 'VALIDATION_ERROR', message: 'tenantId is required', statusCode: 400 }); + } + if (!Array.isArray(input.recipients) || input.recipients.length === 0) { + return err({ code: 'VALIDATION_ERROR', message: 'At least one split recipient is required', statusCode: 400 }); + } + if (input.recipients.length > config.maxRecipients) { + return err({ + code: 'VALIDATION_ERROR', + message: `A split may have at most ${config.maxRecipients} recipients`, + statusCode: 400, + }); + } + + const currency = (input.currency ?? config.defaultCurrency).toUpperCase(); + if (!config.supportedCurrencies.includes(currency)) { + return err({ + code: 'VALIDATION_ERROR', + message: `currency must be one of: ${config.supportedCurrencies.join(', ')}`, + statusCode: 400, + }); + } + + const platformFeePercentage = roundCurrency(input.platformFeePercentage ?? 0); + if (!Number.isFinite(platformFeePercentage) || platformFeePercentage < 0 || platformFeePercentage > 100) { + return err({ + code: 'VALIDATION_ERROR', + message: 'platformFeePercentage must be between 0 and 100', + statusCode: 400, + }); + } + + const recipients: SplitRecipient[] = []; + const seen = new Set(); + + for (const raw of input.recipients as SplitRecipientInput[]) { + if (!raw.recipientId || typeof raw.recipientId !== 'string') { + return err({ code: 'VALIDATION_ERROR', message: 'Each recipient requires a recipientId', statusCode: 400 }); + } + if (seen.has(raw.recipientId)) { + return err({ + code: 'VALIDATION_ERROR', + message: `Duplicate recipientId "${raw.recipientId}"`, + statusCode: 400, + }); + } + seen.add(raw.recipientId); + + if (!raw.walletAddress || typeof raw.walletAddress !== 'string') { + return err({ + code: 'VALIDATION_ERROR', + message: `Recipient "${raw.recipientId}" requires a walletAddress`, + statusCode: 400, + }); + } + + const percentage = roundCurrency(raw.percentage); + if (!Number.isFinite(percentage) || percentage <= 0 || percentage > 100) { + return err({ + code: 'VALIDATION_ERROR', + message: `Recipient "${raw.recipientId}" percentage must be greater than 0 and at most 100`, + statusCode: 400, + }); + } + + const minimumAmount = roundCurrency(raw.minimumAmount ?? 0); + if (!Number.isFinite(minimumAmount) || minimumAmount < 0) { + return err({ + code: 'VALIDATION_ERROR', + message: `Recipient "${raw.recipientId}" minimumAmount must be a non-negative number`, + statusCode: 400, + }); + } + + recipients.push({ + recipientId: raw.recipientId, + walletAddress: raw.walletAddress, + percentage, + shareBps: toBasisPoints(percentage), + minimumAmount, + label: raw.label ?? null, + }); + } + + const platformFeeBps = toBasisPoints(platformFeePercentage); + const recipientBps = recipients.reduce((sum, recipient) => sum + recipient.shareBps, 0); + const totalBps = recipientBps + platformFeeBps; + + if (totalBps !== BPS_TOTAL) { + return err({ + code: 'VALIDATION_ERROR', + message: `Recipient percentages plus platform fee must sum to exactly 100 (got ${fromBasisPoints(totalBps)})`, + statusCode: 400, + }); + } + + return ok({ + tenantId: input.tenantId, + currency, + platformFeePercentage, + platformFeeBps, + recipients, + merchantId: input.merchantId, + name: input.name, + metadata: input.metadata, + }); +} + +export interface AllocateSplitParams { + totalAmount: number; + platformFeePercentage: number; + platformFeeBps: number; + recipients: SplitRecipient[]; +} + +/** + * Allocate a payment across the platform fee and recipients. + * Assumes the recipients were normalised by `validateSplitPlan`. + */ +export function allocateSplit(params: AllocateSplitParams): SplitAllocation { + const totalMinor = Math.round(params.totalAmount * 100); + if (!Number.isFinite(totalMinor) || totalMinor <= 0) { + throw new Error('totalAmount must be a positive number'); + } + + const weights = [params.platformFeeBps, ...params.recipients.map((recipient) => recipient.shareBps)]; + const minors = allocateMinorUnits(totalMinor, weights); + const platformFeeMinor = minors[0] ?? 0; + + const shares: SplitAllocationShare[] = params.recipients.map((recipient, index) => { + const amountMinor = minors[index + 1] ?? 0; + const amount = roundCurrency(amountMinor / 100); + const skipped = amount < recipient.minimumAmount; + return { + recipientId: recipient.recipientId, + walletAddress: recipient.walletAddress, + percentage: recipient.percentage, + amount, + amountMinor, + skipped, + reason: skipped ? 'Below minimum amount' : undefined, + }; + }); + + const allocatedMinor = platformFeeMinor + shares.reduce((sum, share) => sum + share.amountMinor, 0); + + return { + totalAmount: roundCurrency(params.totalAmount), + totalMinor, + platformFeePercentage: params.platformFeePercentage, + platformFeeBps: params.platformFeeBps, + platformFeeAmount: roundCurrency(platformFeeMinor / 100), + platformFeeMinor, + shares, + allocatedMinor, + unallocatedMinor: totalMinor - allocatedMinor, + }; +} diff --git a/backend/src/services/split-payments/index.ts b/backend/src/services/split-payments/index.ts new file mode 100644 index 00000000..fe67e241 --- /dev/null +++ b/backend/src/services/split-payments/index.ts @@ -0,0 +1,27 @@ +/** + * index.ts — Issue #917: Split payments between multiple recipients + * + * Public surface of the split-payments domain. + */ +export * from './types.js'; +export { + DEFAULT_SPLIT_CONFIG, + allocateMinorUnits, + allocateSplit, + fromBasisPoints, + roundCurrency, + toBasisPoints, + validateSplitPlan, + type AllocateSplitParams, +} from './allocation.js'; +export { + InMemorySplitStore, + type SplitExecutionRepository, + type SplitPlanRepository, +} from './store.js'; +export { + SplitPaymentService, + splitPaymentService, + type ExecuteSplitInput, + type SplitPaymentServiceOptions, +} from './splitPaymentService.js'; diff --git a/backend/src/services/split-payments/split-payments.test.ts b/backend/src/services/split-payments/split-payments.test.ts new file mode 100644 index 00000000..570f00c5 --- /dev/null +++ b/backend/src/services/split-payments/split-payments.test.ts @@ -0,0 +1,352 @@ +/** + * split-payments.test.ts — Issue #917: Split payments between multiple + * recipients + * + * Covers basis-point conversion, exact minor-unit allocation, request + * validation, plan lifecycle and execution reconciliation. + */ +import { describe, expect, it, beforeEach } from 'vitest'; + +import { + InMemorySplitStore, + SplitPaymentService, + allocateMinorUnits, + allocateSplit, + fromBasisPoints, + toBasisPoints, + validateSplitPlan, +} from './index.js'; +import type { SplitEvent, SplitEventPublisher, SplitRecipientInput } from './index.js'; + +class CapturingPublisher implements SplitEventPublisher { + events: SplitEvent[] = []; + publish(event: SplitEvent): void { + this.events.push(event); + } +} + +const NOW = new Date('2026-09-28T00:00:00.000Z'); + +function makeService(publisher?: SplitEventPublisher) { + const store = new InMemorySplitStore(); + let counter = 0; + const service = new SplitPaymentService({ + repository: store, + executionRepository: store, + publisher: publisher ?? { publish: () => undefined }, + now: () => NOW, + idFactory: () => `id-${++counter}`, + }); + return service; +} + +function recipient(overrides: Partial = {}): SplitRecipientInput { + return { recipientId: 'r1', walletAddress: 'GABC', percentage: 50, ...overrides }; +} + +function sumMinor(result: ReturnType): number { + return result.platformFeeMinor + result.shares.reduce((sum, share) => sum + share.amountMinor, 0); +} + +describe('basis-point conversion', () => { + it('round-trips percentages through basis points', () => { + expect(toBasisPoints(47.5)).toBe(4750); + expect(toBasisPoints(33.33)).toBe(3333); + expect(fromBasisPoints(4750)).toBe(47.5); + }); +}); + +describe('allocateMinorUnits', () => { + it('allocates the whole total with no remainder lost', () => { + expect(allocateMinorUnits(1000, [3333, 3333, 3334])).toEqual([333, 333, 334]); + expect(allocateMinorUnits(100, [1, 1, 1])).toEqual([34, 33, 33]); + expect(allocateMinorUnits(7, [1, 1, 1, 1, 1, 1, 1])).toEqual([1, 1, 1, 1, 1, 1, 1]); + }); + + it('returns zeros for non-positive inputs', () => { + expect(allocateMinorUnits(0, [1, 1])).toEqual([0, 0]); + expect(allocateMinorUnits(100, [0, 0])).toEqual([0, 0]); + }); +}); + +describe('validateSplitPlan', () => { + const base = { tenantId: 'tenant-1', recipients: [recipient({ percentage: 60 }), recipient({ recipientId: 'r2', percentage: 40 })] }; + + it('normalises a valid plan', () => { + const result = validateSplitPlan({ ...base, platformFeePercentage: 0 }); + expect(result.ok).toBe(true); + if (!result.ok) return; + expect(result.value.currency).toBe('USD'); + expect(result.value.platformFeeBps).toBe(0); + expect(result.value.recipients.map((r) => r.shareBps)).toEqual([6000, 4000]); + }); + + it('rejects an empty recipient list', () => { + expect(validateSplitPlan({ tenantId: 't', recipients: [] }).ok).toBe(false); + }); + + it('rejects duplicate or missing recipient ids', () => { + expect( + validateSplitPlan({ + tenantId: 't', + recipients: [recipient({ percentage: 50 }), recipient({ percentage: 50 })], + }).ok, + ).toBe(false); + expect( + validateSplitPlan({ tenantId: 't', recipients: [recipient({ recipientId: '', percentage: 100 })] }).ok, + ).toBe(false); + }); + + it('rejects a missing wallet address', () => { + expect( + validateSplitPlan({ tenantId: 't', recipients: [recipient({ walletAddress: '', percentage: 100 })] }).ok, + ).toBe(false); + }); + + it('rejects percentages that do not sum to exactly 100', () => { + const short = validateSplitPlan({ tenantId: 't', recipients: [recipient({ percentage: 90 })] }); + expect(short.ok).toBe(false); + if (short.ok) return; + expect(short.error.code).toBe('VALIDATION_ERROR'); + + expect( + validateSplitPlan({ + tenantId: 't', + recipients: [recipient({ recipientId: 'a', percentage: 60 }), recipient({ recipientId: 'b', percentage: 60 })], + }).ok, + ).toBe(false); + }); + + it('rejects out-of-range percentages and fees', () => { + expect(validateSplitPlan({ tenantId: 't', recipients: [recipient({ percentage: 0 })] }).ok).toBe(false); + expect( + validateSplitPlan({ ...base, platformFeePercentage: 101 }).ok, + ).toBe(false); + expect( + validateSplitPlan({ ...base, platformFeePercentage: -1 }).ok, + ).toBe(false); + }); + + it('rejects unsupported currencies and too many recipients', () => { + expect(validateSplitPlan({ ...base, currency: 'JPY' }).ok).toBe(false); + const many = Array.from({ length: 26 }, (_, index) => recipient({ recipientId: `r${index}`, percentage: 100 / 26 })); + expect(validateSplitPlan({ tenantId: 't', recipients: many }).ok).toBe(false); + }); +}); + +describe('allocateSplit', () => { + it('splits fee and recipients exactly', () => { + const validated = validateSplitPlan({ + tenantId: 't', + platformFeePercentage: 5, + recipients: [ + recipient({ recipientId: 'a', percentage: 47.5, walletAddress: 'GA' }), + recipient({ recipientId: 'b', percentage: 47.5, walletAddress: 'GB' }), + ], + }); + if (!validated.ok) throw new Error('validation failed'); + + const allocation = allocateSplit({ + totalAmount: 100, + platformFeePercentage: validated.value.platformFeePercentage, + platformFeeBps: validated.value.platformFeeBps, + recipients: validated.value.recipients, + }); + + expect(allocation.platformFeeAmount).toBe(5); + expect(allocation.shares.map((share) => share.amount)).toEqual([47.5, 47.5]); + expect(allocation.allocatedMinor).toBe(10_000); + expect(allocation.unallocatedMinor).toBe(0); + }); + + it('absorbs rounding remainders so nothing is lost', () => { + const validated = validateSplitPlan({ + tenantId: 't', + recipients: [ + recipient({ recipientId: 'a', percentage: 33.33 }), + recipient({ recipientId: 'b', percentage: 33.33 }), + recipient({ recipientId: 'c', percentage: 33.34 }), + ], + }); + if (!validated.ok) throw new Error('validation failed'); + + const allocation = allocateSplit({ + totalAmount: 10, + platformFeePercentage: 0, + platformFeeBps: 0, + recipients: validated.value.recipients, + }); + + expect(allocation.shares.map((share) => share.amount)).toEqual([3.33, 3.33, 3.34]); + expect(sumMinor(allocation)).toBe(1000); + expect(allocation.unallocatedMinor).toBe(0); + }); + + it('flags recipients below their minimum amount', () => { + const validated = validateSplitPlan({ + tenantId: 't', + recipients: [ + recipient({ recipientId: 'small', percentage: 50, minimumAmount: 10 }), + recipient({ recipientId: 'ok', percentage: 50 }), + ], + }); + if (!validated.ok) throw new Error('validation failed'); + + const allocation = allocateSplit({ + totalAmount: 1, + platformFeePercentage: 0, + platformFeeBps: 0, + recipients: validated.value.recipients, + }); + + expect(allocation.shares[0].skipped).toBe(true); + expect(allocation.shares[0].reason).toBe('Below minimum amount'); + expect(allocation.shares[1].skipped).toBe(false); + }); + + it('throws for a non-positive total', () => { + expect(() => + allocateSplit({ totalAmount: 0, platformFeePercentage: 0, platformFeeBps: 0, recipients: [] }), + ).toThrow(/positive/); + }); +}); + +describe('SplitPaymentService', () => { + let publisher: CapturingPublisher; + let service: SplitPaymentService; + + beforeEach(() => { + publisher = new CapturingPublisher(); + service = makeService(publisher); + }); + + async function createPlan(overrides: Record = {}) { + const result = await service.createPlan({ + tenantId: 'tenant-1', + recipients: [ + recipient({ recipientId: 'a', percentage: 65, walletAddress: 'GA' }), + recipient({ recipientId: 'b', percentage: 32.5, walletAddress: 'GB' }), + ], + platformFeePercentage: 2.5, + ...overrides, + }); + if (!result.ok) throw new Error(`createPlan failed: ${result.error.message}`); + return result.value; + } + + it('creates an active plan and emits an event', async () => { + const plan = await createPlan(); + expect(plan.status).toBe('active'); + expect(plan.platformFeeBps).toBe(250); + expect(plan.recipients).toHaveLength(2); + expect(publisher.events.map((event) => event.type)).toContain('split_plan.created'); + }); + + it('does not persist a plan when validation fails', async () => { + const result = await service.createPlan({ + tenantId: 'tenant-1', + recipients: [recipient({ percentage: 90 })], + }); + expect(result.ok).toBe(false); + const listed = await service.listPlans('tenant-1'); + expect(listed.ok && listed.value).toHaveLength(0); + }); + + it('scopes reads to the owning tenant', async () => { + const plan = await createPlan(); + const other = await service.getPlan('tenant-2', plan.id); + expect(other.ok).toBe(false); + if (other.ok) return; + expect(other.error.code).toBe('NOT_FOUND'); + }); + + it('lists plans filtered by status', async () => { + const plan = await createPlan(); + await service.archivePlan('tenant-1', plan.id); + + const active = await service.listPlans('tenant-1', { status: 'active' }); + expect(active.ok && active.value).toHaveLength(0); + const archived = await service.listPlans('tenant-1', { status: 'archived' }); + expect(archived.ok && archived.value).toHaveLength(1); + }); + + it('rejects archiving a plan twice', async () => { + const plan = await createPlan(); + await service.archivePlan('tenant-1', plan.id); + const again = await service.archivePlan('tenant-1', plan.id); + expect(again.ok).toBe(false); + if (again.ok) return; + expect(again.error.code).toBe('CONFLICT'); + }); + + it('executes a payment and reconciles the full amount', async () => { + const plan = await createPlan(); + const executed = await service.executeSplit('tenant-1', plan.id, { + paymentId: 'pay_1', + totalAmount: 199.99, + }); + + expect(executed.ok).toBe(true); + if (!executed.ok) return; + const execution = executed.value; + expect(execution.totalAmount).toBe(199.99); + expect(execution.platformFeeAmount).toBe(5); + expect(execution.distributions.map((share) => share.amount)).toEqual([129.99, 65]); + expect(execution.allocatedMinor).toBe(19_999); + expect(publisher.events.map((event) => event.type)).toContain('split.executed'); + }); + + it('rejects executing an archived plan', async () => { + const plan = await createPlan(); + await service.archivePlan('tenant-1', plan.id); + const result = await service.executeSplit('tenant-1', plan.id, { paymentId: 'p', totalAmount: 10 }); + expect(result.ok).toBe(false); + if (result.ok) return; + expect(result.error.code).toBe('CONFLICT'); + }); + + it('rejects a currency that does not match the plan', async () => { + const plan = await createPlan({ currency: 'USD' }); + const result = await service.executeSplit('tenant-1', plan.id, { + paymentId: 'p', + totalAmount: 10, + currency: 'EUR', + }); + expect(result.ok).toBe(false); + if (result.ok) return; + expect(result.error.code).toBe('VALIDATION_ERROR'); + }); + + it('rejects executing an unknown plan', async () => { + const result = await service.executeSplit('tenant-1', 'missing', { paymentId: 'p', totalAmount: 10 }); + expect(result.ok).toBe(false); + if (result.ok) return; + expect(result.error.code).toBe('NOT_FOUND'); + }); + + it('previews an allocation without persisting an execution', async () => { + const plan = await createPlan(); + const preview = await service.previewAllocation('tenant-1', plan.id, 100); + + expect(preview.ok).toBe(true); + if (!preview.ok) return; + expect(preview.value.platformFeeAmount).toBe(2.5); + expect(preview.value.shares.map((share) => share.amount)).toEqual([65, 32.5]); + + const executions = await service.listExecutions('tenant-1', plan.id); + expect(executions.ok && executions.value).toHaveLength(0); + }); + + it('summarises executions', async () => { + const plan = await createPlan(); + await service.executeSplit('tenant-1', plan.id, { paymentId: 'p1', totalAmount: 100 }); + await service.executeSplit('tenant-1', plan.id, { paymentId: 'p2', totalAmount: 200 }); + + const summary = await service.getExecutionSummary('tenant-1', plan.id); + expect(summary.ok).toBe(true); + if (!summary.ok) return; + expect(summary.value.executionCount).toBe(2); + expect(summary.value.totalProcessed).toBe(300); + expect(summary.value.totalPlatformFees).toBe(7.5); + }); +}); diff --git a/backend/src/services/split-payments/splitPaymentService.ts b/backend/src/services/split-payments/splitPaymentService.ts new file mode 100644 index 00000000..4c6b4f1c --- /dev/null +++ b/backend/src/services/split-payments/splitPaymentService.ts @@ -0,0 +1,276 @@ +/** + * splitPaymentService.ts — Issue #917: Split payments between multiple + * recipients + * + * Coordinates split-plan lifecycle and payment execution. All validation and + * money maths lives in `allocation.ts`; this class only sequences state + * transitions and persistence so it stays easy to test with a stubbed store. + */ +import { randomUUID } from 'node:crypto'; + +import { BaseService } from '../BaseService.js'; +import type { Result } from '../../lib/result.js'; +import { + DEFAULT_SPLIT_CONFIG, + allocateSplit, + roundCurrency, + validateSplitPlan, +} from './allocation.js'; +import { + InMemorySplitStore, + type SplitExecutionRepository, + type SplitPlanRepository, +} from './store.js'; +import type { + CreateSplitPlanInput, + SplitAllocation, + SplitConfig, + SplitEvent, + SplitEventPublisher, + SplitExecution, + SplitExecutionSummary, + SplitPlan, + SplitPlanFilter, +} from './types.js'; + +export interface SplitPaymentServiceOptions { + repository?: SplitPlanRepository; + executionRepository?: SplitExecutionRepository; + config?: SplitConfig; + publisher?: SplitEventPublisher; + now?: () => Date; + idFactory?: () => string; +} + +export interface ExecuteSplitInput { + paymentId: string; + totalAmount: number; + currency?: string; +} + +/** Default publisher that logs; swap for an event-bus adapter in production. */ +class LoggingSplitEventPublisher implements SplitEventPublisher { + publish(event: SplitEvent): void { + console.info('[split-payments] event', event.type, event.planId); + } +} + +export class SplitPaymentService extends BaseService { + private readonly repository: SplitPlanRepository; + private readonly executionRepository: SplitExecutionRepository; + private readonly config: SplitConfig; + private readonly publisher: SplitEventPublisher; + private readonly now: () => Date; + private readonly idFactory: () => string; + + constructor(options: SplitPaymentServiceOptions = {}) { + super(); + const store = new InMemorySplitStore(); + this.repository = options.repository ?? store; + this.executionRepository = options.executionRepository ?? store; + this.config = options.config ?? DEFAULT_SPLIT_CONFIG; + this.publisher = options.publisher ?? new LoggingSplitEventPublisher(); + this.now = options.now ?? (() => new Date()); + this.idFactory = options.idFactory ?? (() => randomUUID()); + } + + getConfig(): SplitConfig { + return { ...this.config, supportedCurrencies: [...this.config.supportedCurrencies] }; + } + + /** Create a split plan. Percentages must allocate exactly 100. */ + async createPlan(input: CreateSplitPlanInput): Promise> { + const validated = validateSplitPlan(input, this.config); + if (!validated.ok) { + return validated; + } + const request = validated.value; + const timestamp = this.now().toISOString(); + + const plan: SplitPlan = { + id: this.idFactory(), + tenantId: request.tenantId, + merchantId: request.merchantId ?? null, + name: request.name ?? null, + currency: request.currency, + platformFeePercentage: request.platformFeePercentage, + platformFeeBps: request.platformFeeBps, + recipients: request.recipients, + status: 'active', + metadata: request.metadata, + createdAt: timestamp, + updatedAt: timestamp, + }; + + const saved = await this.repository.create(plan); + await this.publisher.publish({ + type: 'split_plan.created', + planId: saved.id, + tenantId: saved.tenantId, + occurredAt: timestamp, + data: { recipientCount: saved.recipients.length }, + }); + return this.ok(saved); + } + + async getPlan(tenantId: string, planId: string): Promise> { + const plan = await this.repository.findById(planId); + if (!plan || plan.tenantId !== tenantId) { + return this.notFoundFailure('Split plan', planId); + } + return this.ok(plan); + } + + async listPlans(tenantId: string, filter?: SplitPlanFilter): Promise> { + if (!tenantId) { + return this.validationFailure('tenantId is required'); + } + return this.ok(await this.repository.list(tenantId, filter)); + } + + async archivePlan(tenantId: string, planId: string): Promise> { + const found = await this.getPlan(tenantId, planId); + if (!found.ok) { + return found; + } + const plan = found.value; + if (plan.status === 'archived') { + return this.conflictFailure('Split plan is already archived'); + } + plan.status = 'archived'; + plan.updatedAt = this.now().toISOString(); + + const saved = await this.repository.update(plan); + await this.publisher.publish({ + type: 'split_plan.archived', + planId: saved.id, + tenantId: saved.tenantId, + occurredAt: saved.updatedAt, + }); + return this.ok(saved); + } + + /** Preview how a given amount would be distributed — no side effects. */ + async previewAllocation( + tenantId: string, + planId: string, + totalAmount: number, + ): Promise> { + const found = await this.getPlan(tenantId, planId); + if (!found.ok) { + return found; + } + if (!Number.isFinite(totalAmount) || totalAmount <= 0) { + return this.validationFailure('totalAmount must be a positive number'); + } + const plan = found.value; + return this.ok( + allocateSplit({ + totalAmount, + platformFeePercentage: plan.platformFeePercentage, + platformFeeBps: plan.platformFeeBps, + recipients: plan.recipients, + }), + ); + } + + /** + * Execute a payment against a split plan. The allocation reconciles to the + * payment exactly: `platformFeeAmount + Σ distributions === totalAmount`. + */ + async executeSplit( + tenantId: string, + planId: string, + input: ExecuteSplitInput, + ): Promise> { + const found = await this.getPlan(tenantId, planId); + if (!found.ok) { + return found; + } + const plan = found.value; + + if (plan.status === 'archived') { + return this.conflictFailure('Cannot execute an archived split plan'); + } + if (!input.paymentId || typeof input.paymentId !== 'string') { + return this.validationFailure('paymentId is required'); + } + if (!Number.isFinite(input.totalAmount) || input.totalAmount <= 0) { + return this.validationFailure('totalAmount must be a positive number'); + } + if (input.currency && input.currency.toUpperCase() !== plan.currency) { + return this.validationFailure( + `currency "${input.currency.toUpperCase()}" does not match the plan currency "${plan.currency}"`, + ); + } + + const allocation = allocateSplit({ + totalAmount: input.totalAmount, + platformFeePercentage: plan.platformFeePercentage, + platformFeeBps: plan.platformFeeBps, + recipients: plan.recipients, + }); + + const execution: SplitExecution = { + id: this.idFactory(), + planId: plan.id, + tenantId: plan.tenantId, + paymentId: input.paymentId, + totalAmount: allocation.totalAmount, + currency: plan.currency, + platformFeeAmount: allocation.platformFeeAmount, + distributions: allocation.shares, + allocatedMinor: allocation.allocatedMinor, + executedAt: this.now().toISOString(), + }; + + const saved = await this.executionRepository.createExecution(execution); + await this.publisher.publish({ + type: 'split.executed', + planId: saved.planId, + tenantId: saved.tenantId, + occurredAt: saved.executedAt, + data: { + executionId: saved.id, + paymentId: saved.paymentId, + totalAmount: saved.totalAmount, + allocatedMinor: saved.allocatedMinor, + }, + }); + return this.ok(saved); + } + + async listExecutions(tenantId: string, planId: string): Promise> { + const found = await this.getPlan(tenantId, planId); + if (!found.ok) { + return found; + } + return this.ok(await this.executionRepository.listByPlan(planId)); + } + + async getExecutionSummary(tenantId: string, planId: string): Promise> { + const found = await this.getPlan(tenantId, planId); + if (!found.ok) { + return found; + } + const executions = await this.executionRepository.listByPlan(planId); + const totalProcessed = roundCurrency(executions.reduce((sum, execution) => sum + execution.totalAmount, 0)); + const totalPlatformFees = roundCurrency( + executions.reduce((sum, execution) => sum + execution.platformFeeAmount, 0), + ); + const skippedDistributions = executions.reduce( + (sum, execution) => sum + execution.distributions.filter((share) => share.skipped).length, + 0, + ); + + return this.ok({ + planId, + executionCount: executions.length, + totalProcessed, + totalPlatformFees, + skippedDistributions, + }); + } +} + +export const splitPaymentService = new SplitPaymentService(); diff --git a/backend/src/services/split-payments/store.ts b/backend/src/services/split-payments/store.ts new file mode 100644 index 00000000..93e98606 --- /dev/null +++ b/backend/src/services/split-payments/store.ts @@ -0,0 +1,78 @@ +/** + * store.ts — Issue #917: Split payments between multiple recipients + * + * Persistence boundary for split plans and their executions. The service + * depends on these interfaces, so production can back them with Prisma while + * tests use the deterministic in-memory implementation below. + */ +import type { SplitExecution, SplitPlan, SplitPlanFilter } from './types.js'; + +export interface SplitPlanRepository { + create(plan: SplitPlan): Promise; + findById(id: string): Promise; + update(plan: SplitPlan): Promise; + list(tenantId: string, filter?: SplitPlanFilter): Promise; +} + +export interface SplitExecutionRepository { + createExecution(execution: SplitExecution): Promise; + listByPlan(planId: string): Promise; +} + +function clonePlan(plan: SplitPlan): SplitPlan { + return { + ...plan, + recipients: plan.recipients.map((recipient) => ({ ...recipient })), + metadata: plan.metadata ? { ...plan.metadata } : plan.metadata, + }; +} + +function cloneExecution(execution: SplitExecution): SplitExecution { + return { + ...execution, + distributions: execution.distributions.map((share) => ({ ...share })), + }; +} + +export class InMemorySplitStore implements SplitPlanRepository, SplitExecutionRepository { + private readonly plans = new Map(); + private readonly executions = new Map(); + + async create(plan: SplitPlan): Promise { + const stored = clonePlan(plan); + this.plans.set(stored.id, stored); + return clonePlan(stored); + } + + async findById(id: string): Promise { + const found = this.plans.get(id); + return found ? clonePlan(found) : null; + } + + async update(plan: SplitPlan): Promise { + const stored = clonePlan(plan); + this.plans.set(stored.id, stored); + return clonePlan(stored); + } + + async list(tenantId: string, filter?: SplitPlanFilter): Promise { + return Array.from(this.plans.values()) + .filter((plan) => plan.tenantId === tenantId) + .filter((plan) => (filter?.status ? plan.status === filter.status : true)) + .filter((plan) => (filter?.merchantId ? plan.merchantId === filter.merchantId : true)) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt) || a.id.localeCompare(b.id)) + .map(clonePlan); + } + + async createExecution(execution: SplitExecution): Promise { + const existing = this.executions.get(execution.planId) ?? []; + const stored = cloneExecution(execution); + existing.push(stored); + this.executions.set(execution.planId, existing); + return cloneExecution(stored); + } + + async listByPlan(planId: string): Promise { + return (this.executions.get(planId) ?? []).map(cloneExecution); + } +} diff --git a/backend/src/services/split-payments/types.ts b/backend/src/services/split-payments/types.ts new file mode 100644 index 00000000..a456692e --- /dev/null +++ b/backend/src/services/split-payments/types.ts @@ -0,0 +1,145 @@ +/** + * types.ts — Issue #917: Split payments between multiple recipients + * + * Domain types for splitting a single payment across a platform fee and one or + * more recipients. Shares are declared as percentages but normalised to + * integer basis points (`shareBps`) so allocation is exact down to the minor + * unit and never drifts through floating point. + */ + +export type SplitPlanStatus = 'active' | 'archived'; + +/** Recipient as supplied by the caller. */ +export interface SplitRecipientInput { + recipientId: string; + walletAddress: string; + /** Share of the payment as a percentage in (0, 100]. */ + percentage: number; + /** Allocations below this amount (major units) are marked `skipped`. */ + minimumAmount?: number; + label?: string; +} + +/** Recipient after normalisation into basis points. */ +export interface SplitRecipient { + recipientId: string; + walletAddress: string; + percentage: number; + /** Integer share in basis points (1 bp = 0.01%). */ + shareBps: number; + minimumAmount: number; + label?: string | null; +} + +export interface SplitPlan { + id: string; + tenantId: string; + merchantId?: string | null; + name?: string | null; + currency: string; + platformFeePercentage: number; + platformFeeBps: number; + recipients: SplitRecipient[]; + status: SplitPlanStatus; + metadata?: Record; + createdAt: string; + updatedAt: string; +} + +export interface CreateSplitPlanInput { + tenantId: string; + recipients: SplitRecipientInput[]; + platformFeePercentage?: number; + currency?: string; + merchantId?: string; + name?: string; + metadata?: Record; +} + +export interface SplitPlanFilter { + status?: SplitPlanStatus; + merchantId?: string; +} + +/** One recipient's (or the platform fee's) slice of an executed payment. */ +export interface SplitAllocationShare { + recipientId: string; + walletAddress: string; + percentage: number; + /** Exact amount, in major currency units. */ + amount: number; + /** Exact amount, in integer minor units (e.g. cents). */ + amountMinor: number; + skipped: boolean; + reason?: string; +} + +export interface SplitAllocation { + totalAmount: number; + totalMinor: number; + platformFeePercentage: number; + platformFeeBps: number; + platformFeeAmount: number; + platformFeeMinor: number; + shares: SplitAllocationShare[]; + /** `platformFeeMinor` plus every share's `amountMinor`. */ + allocatedMinor: number; + /** Always 0 for a validated split; exposed for observability. */ + unallocatedMinor: number; +} + +export interface SplitExecution { + id: string; + planId: string; + tenantId: string; + paymentId: string; + totalAmount: number; + currency: string; + platformFeeAmount: number; + distributions: SplitAllocationShare[]; + allocatedMinor: number; + executedAt: string; +} + +export interface SplitExecutionSummary { + planId: string; + executionCount: number; + totalProcessed: number; + totalPlatformFees: number; + skippedDistributions: number; +} + +export interface SplitConfig { + supportedCurrencies: string[]; + defaultCurrency: string; + maxRecipients: number; +} + +export interface NormalizedSplitPlan { + tenantId: string; + currency: string; + platformFeePercentage: number; + platformFeeBps: number; + recipients: SplitRecipient[]; + merchantId?: string; + name?: string; + metadata?: Record; +} + +export type SplitEventType = + | 'split_plan.created' + | 'split_plan.archived' + | 'split.executed'; + +export interface SplitEvent { + type: SplitEventType; + planId: string; + tenantId: string; + occurredAt: string; + data?: Record; +} + +/** Pluggable publisher so the service stays transport-agnostic. */ +export interface SplitEventPublisher { + publish(event: SplitEvent): Promise | void; +}