From 2122c473006d0a49e3bfd1c649e21592f6d91933 Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:16 -0400 Subject: [PATCH 1/6] fix: make transfer lifecycle concurrency-safe --- src/services/transferService.js | 221 +++++++++++++++++++++++++++----- 1 file changed, 190 insertions(+), 31 deletions(-) diff --git a/src/services/transferService.js b/src/services/transferService.js index 1517f62..3740a5c 100644 --- a/src/services/transferService.js +++ b/src/services/transferService.js @@ -247,6 +247,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 +263,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 +300,215 @@ 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)); +} + +/** + * 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: transfer.version, + } ); } - transfer.status = nextStatus; - transfer.updatedAt = nextTimestamp(transfer.updatedAt); - return transfer; } /** - * Mark a transfer as claimed by the recipient. + * 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'); + } + if (transfer.version !== expectedVersion) { + throw ApiError.conflict('Transfer version conflict', { + transferId: transfer.id, + expectedVersion, + actualVersion: transfer.version, + currentStatus: transfer.status, + }); + } +} + +/** + * Execute one terminal lifecycle mutation with optimistic concurrency and an + * actor-scoped idempotency reservation. + * + * Provider preparation happens before the local terminal commit. Provider + * calls receive a stable operation id so retrying an ambiguous request reuses + * the first remote artifact. If preparation fails, status/version stay + * unchanged and the reservation is released for a safe 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}:${initial.version}`, + expectedVersion: initial.version, + }; - 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; + let committed = false; + try { + const transfer = getTransferOrThrow(id); + assertExpectedVersion(transfer, context.expectedVersion); + assertTransitionAllowed(transfer, spec.nextStatus); + + const beforeStatus = transfer.status; + const beforeVersion = transfer.version; + const operationId = + `${id}:${spec.action}:${context.actor}:${context.key}`; + const prepared = spec.prepare ? spec.prepare(operationId) : undefined; + + // Provider adapters are synchronous today. Re-check immediately before + // commit so a future re-entrant adapter cannot commit over newer state. + if (transfer.status !== beforeStatus || transfer.version !== beforeVersion) { + throw ApiError.conflict( + 'Transfer changed while lifecycle operation was in progress', + { + transferId: id, + expectedVersion: beforeVersion, + actualVersion: transfer.version, + currentStatus: transfer.status, + } + ); + } + + transfer.status = spec.nextStatus; + transfer.version = beforeVersion + 1; + transfer.updatedAt = nextTimestamp(transfer.updatedAt); + if (spec.applyPrepared) { + spec.applyPrepared(transfer, prepared); + } + + 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; + } } /** - * Cancel a pending transfer. + * Mark a transfer as claimed by the recipient. * @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) => + stellarService.createClaimableBalanceId(operationId), + applyPrepared: (transfer, claimableBalanceId) => { + transfer.claimableBalanceId = claimableBalanceId; + }, + auditAction: 'transfer.claimed', + auditPayload: (transfer) => ({ + claimableBalanceId: transfer.claimableBalanceId, + 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 +524,7 @@ function archiveTransfer(id) { const timestamp = nextTimestamp(transfer.updatedAt); transfer.archivedAt = timestamp; transfer.updatedAt = timestamp; + transfer.version = (transfer.version || 1) + 1; } return transfer; } @@ -383,6 +541,7 @@ function unarchiveTransfer(id) { } transfer.archivedAt = null; transfer.updatedAt = nextTimestamp(transfer.updatedAt); + transfer.version = (transfer.version || 1) + 1; return transfer; } From 13346a76ba788a51fd92e6afce4a3324b4b07045 Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:19 -0400 Subject: [PATCH 2/6] fix: make transfer lifecycle concurrency-safe --- src/controllers/transferController.js | 63 +++++++++++++++++++++++---- 1 file changed, 55 insertions(+), 8 deletions(-) diff --git a/src/controllers/transferController.js b/src/controllers/transferController.js index 4fb2e72..d1423a7 100644 --- a/src/controllers/transferController.js +++ b/src/controllers/transferController.js @@ -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); + 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); + 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 = { From 57e9e37f4d0c536fcf2a1211b8add3deb925ed22 Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:22 -0400 Subject: [PATCH 3/6] fix: make transfer lifecycle concurrency-safe --- src/store/index.js | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/src/store/index.js b/src/store/index.js index 0024fac..b723801 100644 --- a/src/store/index.js +++ b/src/store/index.js @@ -23,6 +23,10 @@ const store = { // local so it shares the transfers' lifetime: a replay can never outlive the // transfer it would replay. idempotency: new Map(), + // Terminal lifecycle operation receipts are isolated from create idempotency. + // A claim/cancel retry replays the exact terminal result instead of running + // provider work twice. + lifecycleIdempotency: new Map(), }; /** Remove all records from the store. Primarily used in tests/seeding. */ @@ -31,6 +35,7 @@ function reset() { store.transfers.clear(); store.transferIndex.reset(); store.idempotency.clear(); + store.lifecycleIdempotency.clear(); auditService.reset(); } From c378bca86a9e2a1b4f24e7fafbc9886c14759087 Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:25 -0400 Subject: [PATCH 4/6] fix: make transfer lifecycle concurrency-safe --- src/services/stellarService.js | 30 ++++++++++++++++++++++++++---- 1 file changed, 26 insertions(+), 4 deletions(-) diff --git a/src/services/stellarService.js b/src/services/stellarService.js index e77cea9..1a62da4 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. @@ -18,21 +24,37 @@ const logger = require('../utils/logger'); * @param {string} params.currency * @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. * @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; } module.exports = { From f57d5371e2ce59a8ad9e0fb9165e19570c1c6f4d Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:27 -0400 Subject: [PATCH 5/6] fix: make transfer lifecycle concurrency-safe --- test/transferLifecycleConcurrency.test.js | 124 ++++++++++++++++++++++ 1 file changed, 124 insertions(+) create mode 100644 test/transferLifecycleConcurrency.test.js diff --git a/test/transferLifecycleConcurrency.test.js b/test/transferLifecycleConcurrency.test.js new file mode 100644 index 0000000..5a35a74 --- /dev/null +++ b/test/transferLifecycleConcurrency.test.js @@ -0,0 +1,124 @@ +'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 ApiError = require('../src/utils/ApiError'); + +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 }; +} + +beforeEach(() => { + reset(); +}); + +const realCreateClaimableBalanceId = stellarService.createClaimableBalanceId; +afterEach(() => { + stellarService.createClaimableBalanceId = realCreateClaimableBalanceId; +}); + +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('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); +}); + +test('service-worker reload keeps the shared operation receipt replayable', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('worker-restart', transfer.version); + const first = transferService.claimTransfer(transfer.id, 'req-1', ctx); + + delete require.cache[require.resolve('../src/services/transferService')]; + const restartedService = require('../src/services/transferService'); + const retry = restartedService.claimTransfer(transfer.id, 'req-2', ctx); + + assert.deepEqual(retry, first); + assert.equal(store.transfers.get(transfer.id).version, 2); +}); + +test('provider failure rolls back before terminal commit and releases retry reservation', () => { + const transfer = transferService.createTransfer(PAYLOAD); + const ctx = lifecycle('provider-failure', transfer.version); + + stellarService.createClaimableBalanceId = () => { + 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); + + stellarService.createClaimableBalanceId = realCreateClaimableBalanceId; + 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 + ); +}); From 373c6fb23d130df6f754ba7ad850c105e9ab0d94 Mon Sep 17 00:00:00 2001 From: tokenjunkielabs Date: Sun, 20 Sep 2026 19:05:29 -0400 Subject: [PATCH 6/6] fix: make transfer lifecycle concurrency-safe --- docs/TRANSFER_LIFECYCLE_CONCURRENCY.md | 54 ++++++++++++++++++++++++++ 1 file changed, 54 insertions(+) create mode 100644 docs/TRANSFER_LIFECYCLE_CONCURRENCY.md diff --git a/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md b/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md new file mode 100644 index 0000000..6b8f886 --- /dev/null +++ b/docs/TRANSFER_LIFECYCLE_CONCURRENCY.md @@ -0,0 +1,54 @@ +# Transfer lifecycle concurrency + +Transfer creation already uses actor-scoped idempotency. Terminal lifecycle +mutations now add optimistic resource versions. + +## 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. compare expected version and allowed transition; +3. prepare the provider-side artifact with a stable provider operation key; +4. re-check observed status and version; +5. commit status, version, provider result, and replay receipt; +6. append the audit event. + +A provider failure occurs before the local terminal commit, so the transfer +remains pending and the reservation is released for a safe retry. The mock +Stellar adapter models provider-side idempotency by remembering operation +receipts independently from application transfer state. + +## Storage boundary + +The current demo store is process-local. Its mutation is synchronous, so the +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 records +to the same durable/shared boundary so multiple workers share reservation and +replay state. + +The regression source in test/transferLifecycleConcurrency.test.js covers the +state-machine race, duplicate provider callback, service-worker reload, +provider rollback/retry, and one-terminal-outcome behavior.