From 51118d88e272dbb24dca56c08359d5d73c8729ee Mon Sep 17 00:00:00 2001 From: Blume1977 Date: Wed, 29 Jul 2026 16:08:07 +0200 Subject: [PATCH 1/2] fix(realunit): complete all quotes of a batch settlement tx The issuer may settle multiple purchases in a single on-chain tx with one transfer event each. Quote completion deduplicated consumed settlements per tx hash, so only the first quote of a batch was ever completed and the remaining quotes were stuck in WaitingForPayment although the shares had arrived. Consumption is now tracked per transfer event, identified by its (tx hash, share amount) pairing and counted, which also reconstructs the consumed event of already completed requests and thereby heals stuck quotes on the next cron run. --- .../entities/transaction-request.entity.ts | 3 +- .../services/transaction-request.service.ts | 8 +- .../__tests__/realunit-job.service.spec.ts | 67 +++++++++++++++- .../realunit/realunit-job.service.ts | 77 ++++++++++++++----- 4 files changed, 129 insertions(+), 26 deletions(-) diff --git a/src/subdomains/supporting/payment/entities/transaction-request.entity.ts b/src/subdomains/supporting/payment/entities/transaction-request.entity.ts index 4b1c563e75..bb6fc841a9 100644 --- a/src/subdomains/supporting/payment/entities/transaction-request.entity.ts +++ b/src/subdomains/supporting/payment/entities/transaction-request.entity.ts @@ -127,7 +127,8 @@ export class TransactionRequest extends IEntity { aktionariatResponse?: string; // tx hash of the on-chain transfer that settled this request (set by the settlement job); - // each settlement tx may complete at most one request per user + // a settlement tx may contain multiple transfer events (batch settlement), each of which + // may complete at most one request per user @Column({ length: 256, nullable: true }) settlementTxId?: string; diff --git a/src/subdomains/supporting/payment/services/transaction-request.service.ts b/src/subdomains/supporting/payment/services/transaction-request.service.ts index 9dc7560b60..85281e9866 100644 --- a/src/subdomains/supporting/payment/services/transaction-request.service.ts +++ b/src/subdomains/supporting/payment/services/transaction-request.service.ts @@ -316,13 +316,15 @@ export class TransactionRequestService { }); } - async getUsedSettlementTxIds(userId: number): Promise { + async getUsedSettlements(userId: number): Promise<{ settlementTxId: string; estimatedAmount: number }[]> { return this.transactionRequestRepo .find({ where: { user: { id: userId }, settlementTxId: Not(IsNull()) }, - select: { settlementTxId: true }, + select: { settlementTxId: true, estimatedAmount: true }, }) - .then((requests) => requests.map((r) => r.settlementTxId)); + .then((requests) => + requests.map((r) => ({ settlementTxId: r.settlementTxId, estimatedAmount: r.estimatedAmount })), + ); } async updateEstimatedAmount(id: number, estimatedAmount: number): Promise { diff --git a/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts b/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts index 4cb7b05246..9a4a46855b 100644 --- a/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts +++ b/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts @@ -39,7 +39,7 @@ describe('RealUnitJobService', () => { transactionRequestService = createMock(); jest.spyOn(realunitService, 'getRealuAsset').mockResolvedValue(realuAsset); - jest.spyOn(transactionRequestService, 'getUsedSettlementTxIds').mockResolvedValue([]); + jest.spyOn(transactionRequestService, 'getUsedSettlements').mockResolvedValue([]); const module: TestingModule = await Test.createTestingModule({ imports: [TestSharedModule], @@ -115,9 +115,11 @@ describe('RealUnitJobService', () => { expect(transactionRequestService.complete).toHaveBeenCalledWith(10, '0xSettlementTx'); }); - it('should not reuse a settlement tx that already completed a quote in an earlier run', async () => { + it('should not reuse a settlement transfer that already completed a quote in an earlier run', async () => { jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([quote] as any); - jest.spyOn(transactionRequestService, 'getUsedSettlementTxIds').mockResolvedValue(['0xSettlementTx']); + jest + .spyOn(transactionRequestService, 'getUsedSettlements') + .mockResolvedValue([{ settlementTxId: '0xSettlementTx', estimatedAmount: 72.123 }]); mockHistory([settlementEvent]); await service.completeSettledQuotes(); @@ -125,6 +127,65 @@ describe('RealUnitJobService', () => { expect(transactionRequestService.complete).not.toHaveBeenCalled(); }); + it('should complete multiple quotes settled in a single batch tx', async () => { + const smallQuote = { ...quote, id: 10, estimatedAmount: 219.71 }; + const largeQuote = { ...quote, id: 11, estimatedAmount: 22047 }; + jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([smallQuote, largeQuote] as any); + mockHistory([ + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '219' } }, + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '22047' } }, + ]); + + await service.completeSettledQuotes(); + + expect(transactionRequestService.complete).toHaveBeenCalledTimes(2); + expect(transactionRequestService.complete).toHaveBeenCalledWith(10, '0xBatchTx'); + expect(transactionRequestService.complete).toHaveBeenCalledWith(11, '0xBatchTx'); + }); + + it('should complete a quote from a batch tx whose other transfer already settled an earlier request', async () => { + const largeQuote = { ...quote, id: 11, estimatedAmount: 22047 }; + jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([largeQuote] as any); + jest + .spyOn(transactionRequestService, 'getUsedSettlements') + .mockResolvedValue([{ settlementTxId: '0xBatchTx', estimatedAmount: 219.71 }]); + mockHistory([ + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '219' } }, + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '22047' } }, + ]); + + await service.completeSettledQuotes(); + + expect(transactionRequestService.complete).toHaveBeenCalledWith(11, '0xBatchTx'); + }); + + it('should not reuse a same-amount transfer within a batch tx across runs', async () => { + jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([quote] as any); + jest + .spyOn(transactionRequestService, 'getUsedSettlements') + .mockResolvedValue([{ settlementTxId: '0xBatchTx', estimatedAmount: 72.9 }]); + mockHistory([{ ...settlementEvent, txHash: '0xBatchTx' }]); + + await service.completeSettledQuotes(); + + expect(transactionRequestService.complete).not.toHaveBeenCalled(); + }); + + it('should complete a second same-amount quote when the batch tx contains two matching transfers', async () => { + jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([quote] as any); + jest + .spyOn(transactionRequestService, 'getUsedSettlements') + .mockResolvedValue([{ settlementTxId: '0xBatchTx', estimatedAmount: 72.9 }]); + mockHistory([ + { ...settlementEvent, txHash: '0xBatchTx' }, + { ...settlementEvent, txHash: '0xBatchTx', timestamp: new Date('2026-06-30T09:04:00Z') }, + ]); + + await service.completeSettledQuotes(); + + expect(transactionRequestService.complete).toHaveBeenCalledWith(10, '0xBatchTx'); + }); + it('should match the oldest unused settlement transfer', async () => { const laterEvent = { ...settlementEvent, txHash: '0xLaterTx', timestamp: new Date('2026-07-01T12:00:00Z') }; jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([quote] as any); diff --git a/src/subdomains/supporting/realunit/realunit-job.service.ts b/src/subdomains/supporting/realunit/realunit-job.service.ts index ec7da7b76c..5b9f893669 100644 --- a/src/subdomains/supporting/realunit/realunit-job.service.ts +++ b/src/subdomains/supporting/realunit/realunit-job.service.ts @@ -27,8 +27,11 @@ export class RealUnitJobService { if (!openQuotes.length) return; const historyCache = new Map(); - // per user: settlement txs already consumed by earlier runs (persisted) or earlier in this run - const usedTxIdsByUser = new Map>(); + // per user: settlement transfers already consumed by completed requests (persisted) or earlier in this + // run. The issuer may settle multiple purchases in a single tx (one transfer event each), so consumption + // is tracked per transfer event, not per tx. The history carries no per-event id, so a consumed event is + // identified by its (tx hash, share amount) pairing and counted to also cover same-amount settlements. + const consumedByUser = new Map>(); for (const quote of openQuotes) { try { @@ -41,27 +44,18 @@ export class RealUnitJobService { historyCache.set(address, history); } - let usedTxIds = usedTxIdsByUser.get(quote.user.id); - if (!usedTxIds) { - usedTxIds = new Set(await this.transactionRequestService.getUsedSettlementTxIds(quote.user.id)); - usedTxIdsByUser.set(quote.user.id, usedTxIds); + let consumed = consumedByUser.get(quote.user.id); + if (!consumed) { + consumed = await this.getConsumedSettlements(quote.user.id); + consumedByUser.set(quote.user.id, consumed); } - // quotes are ordered oldest-first, so match the oldest unused settlement transfer - const settlement = history - .filter( - (e) => - e.transfer && - !usedTxIds.has(e.txHash) && - Util.equalsIgnoreCase(e.transfer.to, address) && - Number(e.transfer.value) === expectedShares && - e.timestamp >= quote.created, - ) - .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()) - .at(0); + // quotes are ordered oldest-first, so match the oldest unconsumed settlement transfer + const settlement = this.findUnconsumedSettlement(history, consumed, address, expectedShares, quote.created); if (!settlement) continue; - usedTxIds.add(settlement.txHash); + const key = this.settlementKey(settlement.txHash, expectedShares); + consumed.set(key, (consumed.get(key) ?? 0) + 1); await this.transactionRequestService.complete(quote.id, settlement.txHash); this.logger.info( @@ -76,6 +70,51 @@ export class RealUnitJobService { } } + private async getConsumedSettlements(userId: number): Promise> { + const settlements = await this.transactionRequestService.getUsedSettlements(userId); + + const consumed = new Map(); + for (const settlement of settlements) { + const key = this.settlementKey(settlement.settlementTxId, Math.floor(settlement.estimatedAmount)); + consumed.set(key, (consumed.get(key) ?? 0) + 1); + } + + return consumed; + } + + // Walks all incoming transfers oldest-first and treats the first n events of each (tx hash, share amount) + // pairing as consumed, where n is the number of settlements already recorded for that pairing — so a batch + // settlement tx can complete one request per contained transfer event, but never the same event twice. + private findUnconsumedSettlement( + history: HistoryEventDto[], + consumed: Map, + address: string, + expectedShares: number, + minTimestamp: Date, + ): HistoryEventDto | undefined { + const incomingTransfers = history + .filter((e) => e.transfer && Util.equalsIgnoreCase(e.transfer.to, address)) + .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); + + const seen = new Map(); + + for (const event of incomingTransfers) { + const shares = Number(event.transfer.value); + const key = this.settlementKey(event.txHash, shares); + const position = (seen.get(key) ?? 0) + 1; + seen.set(key, position); + + if (position <= (consumed.get(key) ?? 0)) continue; + if (shares === expectedShares && event.timestamp >= minTimestamp) return event; + } + + return undefined; + } + + private settlementKey(txHash: string, shares: number): string { + return `${txHash.toLowerCase()}|${shares}`; + } + // Resolves RealUnit W2W transfer requests stuck in PROCESSING after a crash/restart between the // atomic claim and the broadcast/callback in confirmTransfer — see // RealUnitService.reconcilePendingTransfers for the actual reconciliation logic. From 79eefa9fc8ab36bc3cf3d5789ce4c0651a27efce Mon Sep 17 00:00:00 2001 From: TaprootFreak <142087526+TaprootFreak@users.noreply.github.com> Date: Wed, 29 Jul 2026 17:24:08 +0200 Subject: [PATCH 2/2] test(realunit): cover the (tx hash, amount) settlement pairing MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reducing the settlement key to the plain tx hash — the very semantics this branch replaces — left all existing specs green, so the amount component of the pairing was never exercised. Add the case a tx-hash-only match gets wrong: the already consumed transfer is not the first event of the batch tx. Also use Util.sort instead of a hand-rolled comparator and move the private helpers behind the HELPER METHODS section marker, as in the sibling job services. --- .../__tests__/realunit-job.service.spec.ts | 17 +++++++++++++ .../realunit/realunit-job.service.ts | 25 +++++++++++-------- 2 files changed, 31 insertions(+), 11 deletions(-) diff --git a/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts b/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts index 9a4a46855b..b42d03b385 100644 --- a/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts +++ b/src/subdomains/supporting/realunit/__tests__/realunit-job.service.spec.ts @@ -186,6 +186,23 @@ describe('RealUnitJobService', () => { expect(transactionRequestService.complete).toHaveBeenCalledWith(10, '0xBatchTx'); }); + it('should complete a quote when the consumed transfer is not the first event of the batch tx', async () => { + const largeQuote = { ...quote, id: 11, estimatedAmount: 22047 }; + jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([largeQuote] as any); + jest + .spyOn(transactionRequestService, 'getUsedSettlements') + .mockResolvedValue([{ settlementTxId: '0xBatchTx', estimatedAmount: 219.71 }]); + // the consumed transfer is the second event here, so a tx-hash-only match would skip the wrong one + mockHistory([ + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '22047' } }, + { ...settlementEvent, txHash: '0xBatchTx', transfer: { ...settlementEvent.transfer, value: '219' } }, + ]); + + await service.completeSettledQuotes(); + + expect(transactionRequestService.complete).toHaveBeenCalledWith(11, '0xBatchTx'); + }); + it('should match the oldest unused settlement transfer', async () => { const laterEvent = { ...settlementEvent, txHash: '0xLaterTx', timestamp: new Date('2026-07-01T12:00:00Z') }; jest.spyOn(transactionRequestService, 'getOpenBuyQuotes').mockResolvedValue([quote] as any); diff --git a/src/subdomains/supporting/realunit/realunit-job.service.ts b/src/subdomains/supporting/realunit/realunit-job.service.ts index 5b9f893669..b75475ebaf 100644 --- a/src/subdomains/supporting/realunit/realunit-job.service.ts +++ b/src/subdomains/supporting/realunit/realunit-job.service.ts @@ -70,6 +70,16 @@ export class RealUnitJobService { } } + // Resolves RealUnit W2W transfer requests stuck in PROCESSING after a crash/restart between the + // atomic claim and the broadcast/callback in confirmTransfer — see + // RealUnitService.reconcilePendingTransfers for the actual reconciliation logic. + @DfxCron(CronExpression.EVERY_5_MINUTES, { process: Process.REALUNIT_TRANSFER_RECONCILIATION, timeout: 1800 }) + async reconcilePendingTransfers(): Promise { + await this.realunitService.reconcilePendingTransfers(); + } + + // --- HELPER METHODS --- // + private async getConsumedSettlements(userId: number): Promise> { const settlements = await this.transactionRequestService.getUsedSettlements(userId); @@ -92,9 +102,10 @@ export class RealUnitJobService { expectedShares: number, minTimestamp: Date, ): HistoryEventDto | undefined { - const incomingTransfers = history - .filter((e) => e.transfer && Util.equalsIgnoreCase(e.transfer.to, address)) - .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); + const incomingTransfers = Util.sort( + history.filter((e) => e.transfer && Util.equalsIgnoreCase(e.transfer.to, address)), + 'timestamp', + ); const seen = new Map(); @@ -114,12 +125,4 @@ export class RealUnitJobService { private settlementKey(txHash: string, shares: number): string { return `${txHash.toLowerCase()}|${shares}`; } - - // Resolves RealUnit W2W transfer requests stuck in PROCESSING after a crash/restart between the - // atomic claim and the broadcast/callback in confirmTransfer — see - // RealUnitService.reconcilePendingTransfers for the actual reconciliation logic. - @DfxCron(CronExpression.EVERY_5_MINUTES, { process: Process.REALUNIT_TRANSFER_RECONCILIATION, timeout: 1800 }) - async reconcilePendingTransfers(): Promise { - await this.realunitService.reconcilePendingTransfers(); - } }