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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
102 changes: 102 additions & 0 deletions backend/docs/JOB_QUEUE.md
Original file line number Diff line number Diff line change
@@ -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.
109 changes: 109 additions & 0 deletions backend/docs/RECURRING_BILLING.md
Original file line number Diff line number Diff line change
@@ -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.
108 changes: 108 additions & 0 deletions backend/docs/SPLIT_PAYMENTS.md
Original file line number Diff line number Diff line change
@@ -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/<planId>/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.
20 changes: 20 additions & 0 deletions backend/src/config/scheduled-tasks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -123,6 +124,25 @@ const RAW_TASKS: Omit<ScheduledTaskMeta, 'schedule'> & { 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',
Expand Down
6 changes: 6 additions & 0 deletions backend/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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());
});
Expand Down
Loading
Loading