Skip to content
Open
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
57 changes: 57 additions & 0 deletions docs/TRANSFER_LIFECYCLE_CONCURRENCY.md
Original file line number Diff line number Diff line change
@@ -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: <stable operation 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`.
67 changes: 57 additions & 10 deletions src/controllers/transferController.js
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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.
*/
Expand Down Expand Up @@ -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);
}

/**
Expand Down Expand Up @@ -127,25 +162,37 @@ function getStats(req, res) {
*/
function getTransfer(req, res) {
const transfer = transferService.getTransferOrThrow(req.params.id);
res.json(transfer);
sendTransfer(res, transfer);
}

/**
* POST /api/transfers/:id/claim
* 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);
}

/**
* POST /api/transfers/:id/cancel
* 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);
}

/**
Expand All @@ -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);
}

/**
Expand All @@ -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 = {
Expand Down
56 changes: 56 additions & 0 deletions src/services/settlementWorker.js
Original file line number Diff line number Diff line change
@@ -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,
};
39 changes: 35 additions & 4 deletions src/services/stellarService.js
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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,
};
Loading