diff --git a/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md b/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md new file mode 100644 index 0000000..92b1ea5 --- /dev/null +++ b/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md @@ -0,0 +1,57 @@ +# Transfer lifecycle concurrency + +Transfer creation already uses actor-scoped idempotency. Terminal lifecycle +mutations (claim / cancel) add optimistic resource versions, a per-transfer +lease, and an idempotent settlement worker. + +## Client contract + +Every single-transfer response carries an ETag containing the integer transfer +version, for example: + + ETag: "1" + +To claim or cancel a transfer, clients send both: + + If-Match: "1" + Idempotency-Key: + +The server reserves the actor/key pair before provider work. A repeated request +with the same key replays the first terminal result. A request based on an old +version returns HTTP 409 with expected version, actual version, and current +status. A missing If-Match returns HTTP 428. Invalid state transitions also +return 409. + +Creation starts at version 1. Every lifecycle, archive, or unarchive mutation +increments the version. + +## Commit order + +Terminal mutations follow this order: + +1. reserve the actor-scoped operation key; +2. acquire an exclusive per-transfer lifecycle lease; +3. compare expected version and allowed transition; +4. prepare the provider-side artifact via the settlement worker (stable + operation key); +5. compare-and-set status + version (single commit path); +6. complete the replay receipt and append the audit event; +7. release the lease. + +A provider failure occurs before the local terminal commit, so the transfer +remains pending and both the lease and reservation are released for a safe +retry. The lease is what closes the double-settlement window when two different +operation keys race: only one caller may prepare provider work for a given +transfer at a time. + +## Storage boundary + +The current demo store is process-local. Version check plus mutation is a +compare-and-set within this process. When the store moves to a database, +preserve the contract atomically with an update constrained by transfer id and +version, and move lifecycle idempotency, leases, and settlement receipts to the +same durable/shared boundary so multiple workers share reservation and replay +state. + +Regression coverage lives in `test/transferLifecycleConcurrency.test.js` and +`test/transferLifecycleHttp.test.js`. diff --git a/src/controllers/transferController.js b/src/controllers/transferController.js index 4fb2e72..cf66282 100644 --- a/src/controllers/transferController.js +++ b/src/controllers/transferController.js @@ -20,11 +20,11 @@ const MAX_IDEMPOTENCY_KEY_LENGTH = 255; * @returns {string} * @throws {ApiError} 400 when the header is missing or unusable. */ -function requireIdempotencyKey(req) { +function requireIdempotencyKey(req, purpose = 'create a transfer') { const raw = req.get('Idempotency-Key'); if (typeof raw !== 'string' || raw.trim() === '') { throw ApiError.badRequest( - 'Idempotency-Key header is required to create a transfer' + `Idempotency-Key header is required to ${purpose}` ); } const key = raw.trim(); @@ -36,6 +36,41 @@ function requireIdempotencyKey(req) { return key; } +/** + * Parse the strong ETag version used as the optimistic precondition for + * terminal transfer mutations. + * @param {import('express').Request} req + * @returns {number} + */ +function requireTransferVersion(req) { + const raw = req.get('If-Match'); + if (typeof raw !== 'string' || raw.trim() === '') { + throw new ApiError( + 428, + 'If-Match header is required for transfer lifecycle mutations' + ); + } + + const match = /^"([1-9][0-9]*)"$/.exec(raw.trim()); + if (!match) { + throw ApiError.badRequest( + 'If-Match must contain the quoted transfer version, for example "1"' + ); + } + return Number(match[1]); +} + +/** + * Return a transfer with its current version as a strong ETag. + * @param {import('express').Response} res + * @param {object} transfer + * @param {number} [status] + */ +function sendTransfer(res, transfer, status = 200) { + res.set('ETag', `"${transfer.version}"`); + res.status(status).json(transfer); +} + /** * Transfer controllers. */ @@ -70,7 +105,7 @@ function createTransfer(req, res) { // did. Replaying the stored result means replaying all of it; downgrading the // status would make a successful retry look different from the response it is // standing in for. - res.status(201).json(transfer); + sendTransfer(res, transfer, 201); } /** @@ -127,7 +162,7 @@ function getStats(req, res) { */ function getTransfer(req, res) { const transfer = transferService.getTransferOrThrow(req.params.id); - res.json(transfer); + sendTransfer(res, transfer); } /** @@ -135,8 +170,14 @@ function getTransfer(req, res) { * Mark a transfer as claimed by the recipient. */ function claimTransfer(req, res) { - const transfer = transferService.claimTransfer(req.params.id, req.id); - res.json(transfer); + const expectedVersion = requireTransferVersion(req); + const key = requireIdempotencyKey(req, 'claim a transfer'); + const transfer = transferService.claimTransfer(req.params.id, req.id, { + actor: req.token, + key, + expectedVersion, + }); + sendTransfer(res, transfer); } /** @@ -144,8 +185,14 @@ function claimTransfer(req, res) { * Cancel a pending transfer. */ function cancelTransfer(req, res) { - const transfer = transferService.cancelTransfer(req.params.id, req.id); - res.json(transfer); + const expectedVersion = requireTransferVersion(req); + const key = requireIdempotencyKey(req, 'cancel a transfer'); + const transfer = transferService.cancelTransfer(req.params.id, req.id, { + actor: req.token, + key, + expectedVersion, + }); + sendTransfer(res, transfer); } /** @@ -154,7 +201,7 @@ function cancelTransfer(req, res) { */ function archiveTransfer(req, res) { const transfer = transferService.archiveTransfer(req.params.id); - res.json(transfer); + sendTransfer(res, transfer); } /** @@ -163,7 +210,7 @@ function archiveTransfer(req, res) { */ function unarchiveTransfer(req, res) { const transfer = transferService.unarchiveTransfer(req.params.id); - res.json(transfer); + sendTransfer(res, transfer); } module.exports = { diff --git a/src/services/settlementWorker.js b/src/services/settlementWorker.js new file mode 100644 index 0000000..771b23d --- /dev/null +++ b/src/services/settlementWorker.js @@ -0,0 +1,56 @@ +'use strict'; + +const { store } = require('../store'); +const stellarService = require('./stellarService'); + +/** + * Idempotent settlement worker for terminal claim operations. + * + * Provider work is keyed by a stable operation id. Retries — including after a + * worker-module reload that shares the same process store — return the first + * settlement receipt instead of creating a second claimable balance. + * + * Receipts live in the shared store so their lifetime matches transfers: a + * restart that clears transfers also clears receipts, keeping the two + * consistent. When the store becomes durable, keep receipts on the same + * boundary as the transfer row that commits the claim. + */ + +/** + * Settle a claim against the payment provider exactly once for `operationId`. + * @param {string} operationId - stable id spanning retries of one claim attempt + * @returns {{ operationId: string, claimableBalanceId: string, settledAt: string }} + */ +function settleClaim(operationId) { + if (typeof operationId !== 'string' || operationId.trim() === '') { + throw new Error('settlementWorker.settleClaim requires a non-empty operationId'); + } + + const existing = store.settlementReceipts.get(operationId); + if (existing) { + return existing; + } + + const claimableBalanceId = stellarService.createClaimableBalanceId(operationId); + const receipt = { + operationId, + claimableBalanceId, + settledAt: new Date().toISOString(), + }; + store.settlementReceipts.set(operationId, receipt); + return receipt; +} + +/** + * Look up a prior settlement receipt without contacting the provider. + * @param {string} operationId + * @returns {object|null} + */ +function getReceipt(operationId) { + return store.settlementReceipts.get(operationId) || null; +} + +module.exports = { + settleClaim, + getReceipt, +}; diff --git a/src/services/stellarService.js b/src/services/stellarService.js index e77cea9..e875891 100644 --- a/src/services/stellarService.js +++ b/src/services/stellarService.js @@ -4,6 +4,12 @@ const config = require('../config'); const { prefixedId } = require('../utils/ids'); const logger = require('../utils/logger'); +// Mock provider receipts live outside application transfer state. A real +// payment provider offers the same property through idempotency keys: retrying +// an ambiguous request returns the first settlement artifact. +const paymentReceipts = new Map(); +const claimableBalanceReceipts = new Map(); + /** * Mock Stellar integration. * Real RemitFlow would submit path payments to the Stellar network. @@ -16,26 +22,51 @@ const logger = require('../utils/logger'); * @param {object} params * @param {number} params.amount * @param {string} params.currency + * @param {string} [params.idempotencyKey] * @returns {{ txHash: string, network: string, ledger: number }} */ -function submitPayment({ amount, currency }) { +function submitPayment({ amount, currency, idempotencyKey }) { + if (idempotencyKey && paymentReceipts.has(idempotencyKey)) { + return paymentReceipts.get(idempotencyKey); + } + logger.debug(`Submitting mock Stellar payment of ${amount} ${currency}`); - return { + const result = { txHash: prefixedId('stellar').replace('stellar_', ''), network: config.stellar.network, ledger: Math.floor(Date.now() / 1000), }; + if (idempotencyKey) { + paymentReceipts.set(idempotencyKey, result); + } + return result; } /** * Generate a mock claimable-balance id used when a recipient claims funds. + * @param {string} [operationId] - stable provider operation key for retries * @returns {string} */ -function createClaimableBalanceId() { - return prefixedId('cb'); +function createClaimableBalanceId(operationId) { + if (operationId && claimableBalanceReceipts.has(operationId)) { + return claimableBalanceReceipts.get(operationId); + } + + const id = prefixedId('cb'); + if (operationId) { + claimableBalanceReceipts.set(operationId, id); + } + return id; +} + +/** Test helper: drop in-module provider receipts without resetting the store. */ +function resetProviderReceipts() { + paymentReceipts.clear(); + claimableBalanceReceipts.clear(); } module.exports = { submitPayment, createClaimableBalanceId, + resetProviderReceipts, }; diff --git a/src/services/transferService.js b/src/services/transferService.js index 1517f62..da558d8 100644 --- a/src/services/transferService.js +++ b/src/services/transferService.js @@ -6,6 +6,7 @@ const ApiError = require('../utils/ApiError'); const { TRANSFER_STATUS, TRANSFER_TRANSITIONS } = require('../config/constants'); const quoteService = require('./quoteService'); const stellarService = require('./stellarService'); +const settlementWorker = require('./settlementWorker'); const idempotencyService = require('./idempotencyService'); const auditService = require('./auditService'); const config = require('../config'); @@ -247,6 +248,9 @@ function createTransferUnchecked(data, requestId, idempotency) { const settlement = stellarService.submitPayment({ amount: quote.sendAmount, currency: quote.from, + idempotencyKey: idempotency + ? `create:${idempotency.actor}:${idempotency.key}` + : undefined, }); const transfer = { @@ -260,6 +264,7 @@ function createTransferUnchecked(data, requestId, idempotency) { rate: quote.rate, receiveAmount: quote.receiveAmount, status: TRANSFER_STATUS.PENDING, + version: 1, stellar: settlement, createdAt: new Date().toISOString(), updatedAt: null, @@ -296,62 +301,314 @@ function createTransferUnchecked(data, requestId, idempotency) { } /** - * Move a transfer to a new status if the transition is allowed. + * Snapshot a transfer before storing it as an idempotent replay result. * @param {object} transfer - * @param {string} nextStatus * @returns {object} */ -function transition(transfer, nextStatus) { +function snapshotTransfer(transfer) { + return JSON.parse(JSON.stringify(transfer)); +} + +/** + * Ensure every transfer carries a positive integer resource version. + * Older in-memory records created before versioning default to 1. + * @param {object} transfer + * @returns {number} + */ +function currentVersion(transfer) { + if (!Number.isInteger(transfer.version) || transfer.version < 1) { + transfer.version = 1; + } + return transfer.version; +} + +/** + * Validate a status transition without mutating state. + * @param {object} transfer + * @param {string} nextStatus + */ +function assertTransitionAllowed(transfer, nextStatus) { const allowed = TRANSFER_TRANSITIONS[transfer.status] || []; if (!allowed.includes(nextStatus)) { throw ApiError.conflict( - `Cannot change transfer from ${transfer.status} to ${nextStatus}` + `Cannot change transfer from ${transfer.status} to ${nextStatus}`, + { + transferId: transfer.id, + currentStatus: transfer.status, + requestedStatus: nextStatus, + version: currentVersion(transfer), + } ); } +} + +/** + * Compare the caller's observed version with current state. + * @param {object} transfer + * @param {number} expectedVersion + */ +function assertExpectedVersion(transfer, expectedVersion) { + if (!Number.isInteger(expectedVersion) || expectedVersion < 1) { + throw ApiError.badRequest('expectedVersion must be a positive integer'); + } + const actual = currentVersion(transfer); + if (actual !== expectedVersion) { + throw ApiError.conflict('Transfer version conflict', { + transferId: transfer.id, + expectedVersion, + actualVersion: actual, + currentStatus: transfer.status, + }); + } +} + +/** + * Compare-and-set a lifecycle transition on a transfer. + * + * Checks expected version and allowed transition, then commits status, + * version, and optional field mutations in one synchronous step. This is the + * only write path for terminal claim/cancel mutations. + * + * @param {object} transfer + * @param {object} options + * @param {number} options.expectedVersion + * @param {string} options.nextStatus + * @param {(transfer: object) => void} [options.mutate] + * @returns {object} the mutated transfer + */ +function compareAndSetTransition(transfer, { expectedVersion, nextStatus, mutate }) { + assertExpectedVersion(transfer, expectedVersion); + assertTransitionAllowed(transfer, nextStatus); + + const beforeVersion = currentVersion(transfer); transfer.status = nextStatus; + transfer.version = beforeVersion + 1; transfer.updatedAt = nextTimestamp(transfer.updatedAt); + if (typeof mutate === 'function') { + mutate(transfer); + } return transfer; } /** - * Mark a transfer as claimed by the recipient. + * Acquire an exclusive lease for one transfer while a lifecycle mutation runs. + * Two different operation keys cannot both prepare provider work for the same + * transfer — that is the original double-settlement window. + * @param {string} transferId + * @param {{ action: string, actor: string, key: string, token: string }} lease + */ +function acquireLifecycleLease(transferId, lease) { + const existing = store.lifecycleLeases.get(transferId); + if (existing) { + throw ApiError.conflict( + 'Transfer lifecycle operation already in progress', + { + transferId, + heldAction: existing.action, + requestedAction: lease.action, + } + ); + } + store.lifecycleLeases.set(transferId, lease); +} + +/** + * Release a lifecycle lease when the holder still owns it. + * @param {string} transferId + * @param {string} token + */ +function releaseLifecycleLease(transferId, token) { + const existing = store.lifecycleLeases.get(transferId); + if (existing && existing.token === token) { + store.lifecycleLeases.delete(transferId); + } +} + +/** + * Execute one terminal lifecycle mutation with optimistic concurrency, a + * transfer-scoped lease, and an actor-scoped idempotency reservation. + * + * Order of operations: + * 1. reserve the actor/key pair (replay if already completed); + * 2. acquire the per-transfer lease; + * 3. check expected version + allowed transition; + * 4. prepare provider work with a stable operation id; + * 5. compare-and-set status/version (and apply prepared fields); + * 6. complete the idempotency receipt and audit. + * + * Provider failure happens before the local terminal commit, so status/version + * stay unchanged and both the lease and reservation are released for retry. + * * @param {string} id - * @param {string} [requestId] - optional correlation id for audit logging + * @param {object} spec + * @param {string} spec.action + * @param {string} spec.nextStatus + * @param {(operationId: string) => unknown} [spec.prepare] + * @param {(transfer: object, prepared: unknown) => void} [spec.applyPrepared] + * @param {string} spec.auditAction + * @param {(transfer: object) => object} spec.auditPayload + * @param {string} [requestId] + * @param {{actor: string, key: string, expectedVersion: number}} [lifecycle] * @returns {object} */ -function claimTransfer(id, requestId) { - const transfer = getTransferOrThrow(id); - transition(transfer, TRANSFER_STATUS.CLAIMED); - transfer.claimableBalanceId = stellarService.createClaimableBalanceId(); +function executeLifecycleMutation(id, spec, requestId, lifecycle) { + const initial = getTransferOrThrow(id); + const context = lifecycle || { + actor: 'internal', + key: `${spec.action}:${id}:${currentVersion(initial)}`, + expectedVersion: currentVersion(initial), + }; - auditService.addEntry({ - action: 'transfer.claimed', - resourceId: transfer.id, - payload: { claimableBalanceId: transfer.claimableBalanceId }, - requestId, + if (!context.actor || !context.key) { + throw ApiError.badRequest( + 'Lifecycle mutations require an actor and idempotency key' + ); + } + + const fingerprint = idempotencyService.fingerprint({ + transferId: id, + action: spec.action, + expectedVersion: context.expectedVersion, }); + const scopedKey = `${id}:${spec.action}:${context.key}`; + const reservation = idempotencyService.begin( + store.lifecycleIdempotency, + context.actor, + scopedKey, + fingerprint + ); + + if (reservation.status === 'replay') { + return reservation.result; + } - return transfer; + const leaseToken = `${context.actor}:${scopedKey}:${Date.now()}`; + let leased = false; + let committed = false; + + try { + acquireLifecycleLease(id, { + action: spec.action, + actor: context.actor, + key: context.key, + token: leaseToken, + }); + leased = true; + + const transfer = getTransferOrThrow(id); + assertExpectedVersion(transfer, context.expectedVersion); + assertTransitionAllowed(transfer, spec.nextStatus); + + const beforeStatus = transfer.status; + const beforeVersion = currentVersion(transfer); + const operationId = `${id}:${spec.action}:${context.actor}:${context.key}`; + const prepared = spec.prepare ? spec.prepare(operationId) : undefined; + + // Re-check immediately before CAS so a re-entrant adapter that somehow + // mutated state during prepare cannot commit over newer state. The lease + // already blocks competing lifecycle keys; this is defense in depth. + if (transfer.status !== beforeStatus || currentVersion(transfer) !== beforeVersion) { + throw ApiError.conflict( + 'Transfer changed while lifecycle operation was in progress', + { + transferId: id, + expectedVersion: beforeVersion, + actualVersion: currentVersion(transfer), + currentStatus: transfer.status, + } + ); + } + + compareAndSetTransition(transfer, { + expectedVersion: beforeVersion, + nextStatus: spec.nextStatus, + mutate: spec.applyPrepared + ? (target) => spec.applyPrepared(target, prepared) + : undefined, + }); + + const result = snapshotTransfer(transfer); + idempotencyService.complete( + store.lifecycleIdempotency, + context.actor, + scopedKey, + result + ); + committed = true; + + auditService.addEntry({ + action: spec.auditAction, + resourceId: transfer.id, + payload: spec.auditPayload(transfer), + requestId, + }); + + return result; + } catch (err) { + if (!committed) { + idempotencyService.release( + store.lifecycleIdempotency, + context.actor, + scopedKey + ); + } + throw err; + } finally { + if (leased) { + releaseLifecycleLease(id, leaseToken); + } + } } /** - * Cancel a pending transfer. + * Mark a transfer as claimed by the recipient (settlement). * @param {string} id - * @param {string} [requestId] - optional correlation id for audit logging + * @param {string} [requestId] + * @param {{actor: string, key: string, expectedVersion: number}} [lifecycle] * @returns {object} */ -function cancelTransfer(id, requestId) { - const transfer = getTransferOrThrow(id); - transition(transfer, TRANSFER_STATUS.CANCELLED); - - auditService.addEntry({ - action: 'transfer.cancelled', - resourceId: transfer.id, - payload: {}, +function claimTransfer(id, requestId, lifecycle) { + return executeLifecycleMutation( + id, + { + action: 'claim', + nextStatus: TRANSFER_STATUS.CLAIMED, + prepare: (operationId) => settlementWorker.settleClaim(operationId), + applyPrepared: (transfer, receipt) => { + transfer.claimableBalanceId = receipt.claimableBalanceId; + transfer.settlementOperationId = receipt.operationId; + }, + auditAction: 'transfer.claimed', + auditPayload: (transfer) => ({ + claimableBalanceId: transfer.claimableBalanceId, + settlementOperationId: transfer.settlementOperationId, + version: transfer.version, + }), + }, requestId, - }); + lifecycle + ); +} - return transfer; +/** + * Cancel a pending transfer. + * @param {string} id + * @param {string} [requestId] + * @param {{actor: string, key: string, expectedVersion: number}} [lifecycle] + * @returns {object} + */ +function cancelTransfer(id, requestId, lifecycle) { + return executeLifecycleMutation( + id, + { + action: 'cancel', + nextStatus: TRANSFER_STATUS.CANCELLED, + auditAction: 'transfer.cancelled', + auditPayload: (transfer) => ({ version: transfer.version }), + }, + requestId, + lifecycle + ); } /** @@ -367,6 +624,7 @@ function archiveTransfer(id) { const timestamp = nextTimestamp(transfer.updatedAt); transfer.archivedAt = timestamp; transfer.updatedAt = timestamp; + transfer.version = currentVersion(transfer) + 1; } return transfer; } @@ -383,6 +641,7 @@ function unarchiveTransfer(id) { } transfer.archivedAt = null; transfer.updatedAt = nextTimestamp(transfer.updatedAt); + transfer.version = currentVersion(transfer) + 1; return transfer; } @@ -398,4 +657,7 @@ module.exports = { cancelTransfer, archiveTransfer, unarchiveTransfer, + compareAndSetTransition, + assertTransitionAllowed, + currentVersion, }; diff --git a/src/store/index.js b/src/store/index.js index 0024fac..f944e39 100644 --- a/src/store/index.js +++ b/src/store/index.js @@ -19,10 +19,21 @@ const store = { * valid for the life of the process. */ transferIndex: new OrderedIndex({ sortKeyOf: (transfer) => transfer.createdAt }), - // Keyed by " ". Lives here rather than in a module + // Keyed by "\0". Lives here rather than in a module // local so it shares the transfers' lifetime: a replay can never outlive the // transfer it would replay. idempotency: new Map(), + // Actor-scoped receipts for terminal claim/cancel mutations. Isolated from + // create idempotency so a retry of a claim replays the terminal result + // without colliding with the key used to create the transfer. + lifecycleIdempotency: new Map(), + // Per-transfer leases held while a lifecycle mutation is between reservation + // and commit. Prevents two different operation keys from both calling the + // provider for the same transfer (the double-settlement window). + lifecycleLeases: new Map(), + // Provider settlement receipts keyed by stable operation id. Shared with the + // settlement worker so a module reload still returns the first artifact. + settlementReceipts: new Map(), }; /** Remove all records from the store. Primarily used in tests/seeding. */ @@ -31,6 +42,9 @@ function reset() { store.transfers.clear(); store.transferIndex.reset(); store.idempotency.clear(); + store.lifecycleIdempotency.clear(); + store.lifecycleLeases.clear(); + store.settlementReceipts.clear(); auditService.reset(); } diff --git a/test/requireScope.test.js b/test/requireScope.test.js index 33b723d..ad3cbdc 100644 --- a/test/requireScope.test.js +++ b/test/requireScope.test.js @@ -305,10 +305,14 @@ test('full transfer lifecycle: create → claim with correct scopes', async () = assert.equal(readRes.status, 200); assert.equal(readRes.body.id, id); - // Claim with write token + // Claim with write token (If-Match + Idempotency-Key required for lifecycle mutations) const claimRes = await fetchJson(`/api/transfers/${id}/claim`, { method: 'POST', - headers: authHeader('test-token-admin'), + headers: { + ...authHeader('test-token-admin'), + 'If-Match': `"${createRes.body.version}"`, + 'Idempotency-Key': 'idem-requireScope-claim', + }, }); assert.equal(claimRes.status, 200); assert.equal(claimRes.body.status, 'claimed'); @@ -325,7 +329,11 @@ test('full transfer lifecycle: create → cancel with correct scopes', async () const cancelRes = await fetchJson(`/api/transfers/${id}/cancel`, { method: 'POST', - headers: authHeader('test-token-transfers'), + headers: { + ...authHeader('test-token-transfers'), + 'If-Match': `"${createRes.body.version}"`, + 'Idempotency-Key': 'idem-requireScope-cancel', + }, }); assert.equal(cancelRes.status, 200); assert.equal(cancelRes.body.status, 'cancelled'); diff --git a/test/transferLifecycleConcurrency.test.js b/test/transferLifecycleConcurrency.test.js new file mode 100644 index 0000000..de68455 --- /dev/null +++ b/test/transferLifecycleConcurrency.test.js @@ -0,0 +1,346 @@ +'use strict'; + +const { test, beforeEach, afterEach } = require('node:test'); +const assert = require('node:assert/strict'); + +const { store, reset } = require('../src/store'); +const transferService = require('../src/services/transferService'); +const stellarService = require('../src/services/stellarService'); +const settlementWorker = require('../src/services/settlementWorker'); +const ApiError = require('../src/utils/ApiError'); +const { TRANSFER_STATUS, TRANSFER_TRANSITIONS } = require('../src/config/constants'); + +const PAYLOAD = { + senderName: 'Alice', + recipientName: 'Bob', + amount: 100, + from: 'USD', + to: 'EUR', +}; + +function lifecycle(key, version, actor = 'test-token-admin') { + return { actor, key, expectedVersion: version }; +} + +let settleCalls; +const realSettleClaim = settlementWorker.settleClaim; +const realCreateClaimableBalanceId = stellarService.createClaimableBalanceId; + +beforeEach(() => { + reset(); + settleCalls = 0; + settlementWorker.settleClaim = (operationId) => { + settleCalls += 1; + return realSettleClaim(operationId); + }; +}); + +afterEach(() => { + settlementWorker.settleClaim = realSettleClaim; + stellarService.createClaimableBalanceId = realCreateClaimableBalanceId; +}); + +// ============================================================================ +// State machine +// ============================================================================ + +test('only transitions listed in TRANSFER_TRANSITIONS may commit', () => { + const transfer = transferService.createTransfer(PAYLOAD); + assert.equal(transfer.version, 1); + assert.deepEqual( + TRANSFER_TRANSITIONS[TRANSFER_STATUS.PENDING], + [TRANSFER_STATUS.CLAIMED, TRANSFER_STATUS.CANCELLED] + ); + assert.deepEqual(TRANSFER_TRANSITIONS[TRANSFER_STATUS.CLAIMED], []); + assert.deepEqual(TRANSFER_TRANSITIONS[TRANSFER_STATUS.CANCELLED], []); + + assert.throws( + () => + transferService.compareAndSetTransition(transfer, { + expectedVersion: 1, + nextStatus: 'pending', + }), + (err) => err instanceof ApiError && err.statusCode === 409 + ); +}); + +test('compareAndSetTransition rejects a stale expected version without mutating', () => { + const transfer = transferService.createTransfer(PAYLOAD); + transferService.claimTransfer(transfer.id, 'req', lifecycle('claim-v', 1)); + + assert.throws( + () => + transferService.compareAndSetTransition(store.transfers.get(transfer.id), { + expectedVersion: 1, + nextStatus: TRANSFER_STATUS.CANCELLED, + }), + (err) => + err instanceof ApiError + && err.statusCode === 409 + && err.details.actualVersion === 2 + ); + + assert.equal(store.transfers.get(transfer.id).status, 'claimed'); + assert.equal(store.transfers.get(transfer.id).version, 2); +}); + +// ============================================================================ +// One terminal outcome / race +// ============================================================================ + +test('one terminal outcome wins when claim and cancel race from the same version', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const claimed = transferService.claimTransfer( + transfer.id, + 'req-claim', + lifecycle('race-claim', transfer.version) + ); + + assert.equal(claimed.status, 'claimed'); + assert.equal(claimed.version, 2); + + assert.throws( + () => + transferService.cancelTransfer( + transfer.id, + 'req-cancel', + lifecycle('race-cancel', 1) + ), + (err) => + err instanceof ApiError + && err.statusCode === 409 + && err.details.actualVersion === 2 + ); + + assert.equal(store.transfers.get(transfer.id).status, 'claimed'); + assert.equal(store.transfers.get(transfer.id).version, 2); +}); + +test('reentrant cancel during claim prepare loses; only one terminal outcome commits', () => { + // Reproduce the original failure window: a competing mutation arrives while + // the first operation has reserved work but has not yet committed. + const transfer = transferService.createTransfer(PAYLOAD); + let nestedError = null; + + settlementWorker.settleClaim = (operationId) => { + settleCalls += 1; + try { + transferService.cancelTransfer( + transfer.id, + 'req-nested-cancel', + lifecycle('nested-cancel', 1) + ); + nestedError = false; + } catch (err) { + nestedError = err; + } + return realSettleClaim(operationId); + }; + + const claimed = transferService.claimTransfer( + transfer.id, + 'req-claim', + lifecycle('outer-claim', 1) + ); + + assert.ok(nestedError instanceof ApiError, 'nested cancel must be refused'); + assert.equal(nestedError.statusCode, 409); + assert.match(nestedError.message, /already in progress|version conflict/i); + assert.equal(claimed.status, 'claimed'); + assert.equal(store.transfers.get(transfer.id).status, 'claimed'); + assert.equal(store.transfers.get(transfer.id).version, 2); + assert.equal(settleCalls, 1); +}); + +test('two concurrent claim keys cannot double-settle the same transfer', () => { + const transfer = transferService.createTransfer(PAYLOAD); + let nestedError = null; + + settlementWorker.settleClaim = (operationId) => { + settleCalls += 1; + if (settleCalls === 1) { + try { + transferService.claimTransfer( + transfer.id, + 'req-other', + lifecycle('other-claim-key', 1) + ); + nestedError = false; + } catch (err) { + nestedError = err; + } + } + return realSettleClaim(operationId); + }; + + const first = transferService.claimTransfer( + transfer.id, + 'req-1', + lifecycle('first-claim-key', 1) + ); + + assert.ok(nestedError instanceof ApiError); + assert.equal(nestedError.statusCode, 409); + assert.match(nestedError.message, /already in progress/i); + assert.equal(first.status, 'claimed'); + assert.equal(settleCalls, 1, 'provider settlement must run exactly once'); + assert.equal(store.settlementReceipts.size, 1); + assert.equal(store.transfers.get(transfer.id).claimableBalanceId, first.claimableBalanceId); +}); + +// ============================================================================ +// Duplicate callback / idempotency +// ============================================================================ + +test('duplicate claim callback replays the first provider artifact', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('provider-callback-42', transfer.version); + + const first = transferService.claimTransfer(transfer.id, 'req-1', ctx); + const duplicate = transferService.claimTransfer(transfer.id, 'req-2', ctx); + + assert.deepEqual(duplicate, first); + assert.equal(duplicate.claimableBalanceId, first.claimableBalanceId); + assert.equal(store.transfers.get(transfer.id).version, 2); + assert.equal(settleCalls, 1); +}); + +test('a new claim key after terminal success cannot settle twice', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const first = transferService.claimTransfer( + transfer.id, + 'req-1', + lifecycle('claim-once', 1) + ); + + assert.throws( + () => + transferService.claimTransfer( + transfer.id, + 'req-2', + lifecycle('claim-again', 2) + ), + (err) => + err instanceof ApiError + && err.statusCode === 409 + && /Cannot change transfer from claimed/.test(err.message) + ); + + assert.equal(settleCalls, 1); + assert.equal(store.transfers.get(transfer.id).claimableBalanceId, first.claimableBalanceId); +}); + +// ============================================================================ +// Worker restart +// ============================================================================ + +test('settlement worker reload keeps shared receipts and lifecycle replay', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('worker-restart', transfer.version); + const first = transferService.claimTransfer(transfer.id, 'req-1', ctx); + const operationId = first.settlementOperationId; + + delete require.cache[require.resolve('../src/services/settlementWorker')]; + delete require.cache[require.resolve('../src/services/transferService')]; + const restartedWorker = require('../src/services/settlementWorker'); + const restartedService = require('../src/services/transferService'); + + const receipt = restartedWorker.settleClaim(operationId); + assert.equal(receipt.claimableBalanceId, first.claimableBalanceId); + + const retry = restartedService.claimTransfer(transfer.id, 'req-2', ctx); + assert.deepEqual(retry, first); + assert.equal(store.transfers.get(transfer.id).version, 2); +}); + +// ============================================================================ +// Rollback / provider failure +// ============================================================================ + +test('provider failure rolls back before terminal commit and releases retry reservation', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('provider-failure', transfer.version); + + settlementWorker.settleClaim = () => { + settleCalls += 1; + throw new Error('provider unavailable'); + }; + + assert.throws( + () => transferService.claimTransfer(transfer.id, 'req-1', ctx), + /provider unavailable/ + ); + + const afterFailure = store.transfers.get(transfer.id); + assert.equal(afterFailure.status, 'pending'); + assert.equal(afterFailure.version, 1); + assert.equal(store.lifecycleIdempotency.size, 0); + assert.equal(store.lifecycleLeases.size, 0); + assert.equal(store.settlementReceipts.size, 0); + + settlementWorker.settleClaim = (operationId) => { + settleCalls += 1; + return realSettleClaim(operationId); + }; + const recovered = transferService.claimTransfer(transfer.id, 'req-2', ctx); + assert.equal(recovered.status, 'claimed'); + assert.equal(recovered.version, 2); +}); + +test('completed cancellation is idempotent and cannot be overwritten', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('cancel-once', transfer.version); + const first = transferService.cancelTransfer(transfer.id, 'req-1', ctx); + const retry = transferService.cancelTransfer(transfer.id, 'req-2', ctx); + + assert.deepEqual(retry, first); + assert.equal(first.status, 'cancelled'); + + assert.throws( + () => + transferService.claimTransfer( + transfer.id, + 'req-3', + lifecycle('claim-after-cancel', 1) + ), + (err) => err instanceof ApiError && err.statusCode === 409 + ); +}); + +// ============================================================================ +// Archive version bump vs stale lifecycle +// ============================================================================ + +test('archive advances version so a stale claim cannot silently commit', () => { + const transfer = transferService.createTransfer(PAYLOAD); + transferService.archiveTransfer(transfer.id); + assert.equal(store.transfers.get(transfer.id).version, 2); + + assert.throws( + () => + transferService.claimTransfer( + transfer.id, + 'req-stale', + lifecycle('stale-after-archive', 1) + ), + (err) => + err instanceof ApiError + && err.statusCode === 409 + && err.details.actualVersion === 2 + ); + + assert.equal(store.transfers.get(transfer.id).status, 'pending'); + assert.equal(settleCalls, 0); +}); + +// ============================================================================ +// Backwards compatibility for internal callers +// ============================================================================ + +test('claimTransfer without lifecycle context still works for internal callers', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const claimed = transferService.claimTransfer(transfer.id, 'internal-req'); + assert.equal(claimed.status, 'claimed'); + assert.equal(claimed.version, 2); + assert.equal(transfer.status, 'claimed'); +}); diff --git a/test/transferLifecycleHttp.test.js b/test/transferLifecycleHttp.test.js new file mode 100644 index 0000000..1579b0f --- /dev/null +++ b/test/transferLifecycleHttp.test.js @@ -0,0 +1,162 @@ +'use strict'; + +const { test, before, after, beforeEach } = require('node:test'); +const assert = require('node:assert/strict'); + +process.env.NODE_ENV = 'test'; + +const createApp = require('../src/app'); +const { store, reset } = require('../src/store'); + +let server; +let baseUrl; + +before(() => { + const app = createApp(); + return new Promise((resolve) => { + server = app.listen(0, () => { + baseUrl = `http://127.0.0.1:${server.address().port}`; + resolve(); + }); + }); +}); + +after(() => { + if (server) { + server.close(); + } +}); + +beforeEach(() => { + reset(); +}); + +const BODY = { + senderName: 'Alice', + recipientName: 'Bob', + amount: 100, + from: 'USD', + to: 'EUR', +}; + +async function createTransfer(key = 'create-1') { + const res = await fetch(`${baseUrl}/api/transfers`, { + method: 'POST', + headers: { + Authorization: 'Bearer test-token-admin', + 'Content-Type': 'application/json', + 'Idempotency-Key': key, + }, + body: JSON.stringify(BODY), + }); + const etag = res.headers.get('etag'); + return { status: res.status, body: await res.json(), etag }; +} + +async function lifecyclePost(path, { ifMatch, idempotencyKey }) { + const headers = { + Authorization: 'Bearer test-token-admin', + 'Content-Type': 'application/json', + }; + if (ifMatch !== null && ifMatch !== undefined) { + headers['If-Match'] = ifMatch; + } + if (idempotencyKey !== null && idempotencyKey !== undefined) { + headers['Idempotency-Key'] = idempotencyKey; + } + const res = await fetch(`${baseUrl}${path}`, { + method: 'POST', + headers, + }); + const etag = res.headers.get('etag'); + return { status: res.status, body: await res.json(), etag }; +} + +test('GET transfer exposes version as a strong ETag', async () => { + const created = await createTransfer('etag-create'); + assert.equal(created.status, 201); + assert.equal(created.etag, '"1"'); + assert.equal(created.body.version, 1); + + const res = await fetch(`${baseUrl}/api/transfers/${created.body.id}`, { + headers: { Authorization: 'Bearer test-token-admin' }, + }); + assert.equal(res.status, 200); + assert.equal(res.headers.get('etag'), '"1"'); +}); + +test('claim without If-Match returns 428', async () => { + const created = await createTransfer('no-if-match'); + const { status, body } = await lifecyclePost( + `/api/transfers/${created.body.id}/claim`, + { ifMatch: null, idempotencyKey: 'claim-1' } + ); + assert.equal(status, 428); + assert.match(body.error.message, /If-Match/i); + assert.equal(store.transfers.get(created.body.id).status, 'pending'); +}); + +test('claim without Idempotency-Key returns 400', async () => { + const created = await createTransfer('no-idem'); + const { status, body } = await lifecyclePost( + `/api/transfers/${created.body.id}/claim`, + { ifMatch: '"1"', idempotencyKey: null } + ); + assert.equal(status, 400); + assert.match(body.error.message, /Idempotency-Key/i); +}); + +test('claim with stale If-Match returns 409 conflict details', async () => { + const created = await createTransfer('stale-claim'); + const first = await lifecyclePost(`/api/transfers/${created.body.id}/claim`, { + ifMatch: '"1"', + idempotencyKey: 'claim-ok', + }); + assert.equal(first.status, 200); + assert.equal(first.etag, '"2"'); + + const stale = await lifecyclePost(`/api/transfers/${created.body.id}/cancel`, { + ifMatch: '"1"', + idempotencyKey: 'cancel-stale', + }); + assert.equal(stale.status, 409); + assert.equal(stale.body.error.details.actualVersion, 2); + assert.equal(stale.body.error.details.currentStatus, 'claimed'); + assert.equal(store.transfers.get(created.body.id).status, 'claimed'); +}); + +test('retried claim with the same Idempotency-Key replays the first result', async () => { + const created = await createTransfer('replay-claim'); + const first = await lifecyclePost(`/api/transfers/${created.body.id}/claim`, { + ifMatch: '"1"', + idempotencyKey: 'stable-claim', + }); + const second = await lifecyclePost(`/api/transfers/${created.body.id}/claim`, { + ifMatch: '"1"', + idempotencyKey: 'stable-claim', + }); + + assert.equal(first.status, 200); + assert.equal(second.status, 200); + assert.equal(second.body.id, first.body.id); + assert.equal(second.body.claimableBalanceId, first.body.claimableBalanceId); + assert.equal(second.body.version, 2); + assert.equal(store.transfers.size, 1); +}); + +test('cancel then claim from the cancelled version is rejected', async () => { + const created = await createTransfer('cancel-then-claim'); + const cancelled = await lifecyclePost( + `/api/transfers/${created.body.id}/cancel`, + { ifMatch: '"1"', idempotencyKey: 'cancel-ok' } + ); + assert.equal(cancelled.status, 200); + assert.equal(cancelled.body.status, 'cancelled'); + + const claim = await lifecyclePost(`/api/transfers/${created.body.id}/claim`, { + ifMatch: '"2"', + idempotencyKey: 'claim-after', + }); + assert.equal(claim.status, 409); + assert.match(claim.body.error.message, /Cannot change transfer from cancelled/); +});