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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -316,13 +316,15 @@ export class TransactionRequestService {
});
}

async getUsedSettlementTxIds(userId: number): Promise<string[]> {
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<void> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ describe('RealUnitJobService', () => {
transactionRequestService = createMock<TransactionRequestService>();

jest.spyOn(realunitService, 'getRealuAsset').mockResolvedValue(realuAsset);
jest.spyOn(transactionRequestService, 'getUsedSettlementTxIds').mockResolvedValue([]);
jest.spyOn(transactionRequestService, 'getUsedSettlements').mockResolvedValue([]);

const module: TestingModule = await Test.createTestingModule({
imports: [TestSharedModule],
Expand Down Expand Up @@ -115,16 +115,94 @@ 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();

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 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);
Expand Down
80 changes: 61 additions & 19 deletions src/subdomains/supporting/realunit/realunit-job.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,11 @@ export class RealUnitJobService {
if (!openQuotes.length) return;

const historyCache = new Map<string, HistoryEventDto[]>();
// per user: settlement txs already consumed by earlier runs (persisted) or earlier in this run
const usedTxIdsByUser = new Map<number, Set<string>>();
// 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<number, Map<string, number>>();

for (const quote of openQuotes) {
try {
Expand All @@ -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(
Expand All @@ -83,4 +77,52 @@ export class RealUnitJobService {
async reconcilePendingTransfers(): Promise<void> {
await this.realunitService.reconcilePendingTransfers();
}

// --- HELPER METHODS --- //

private async getConsumedSettlements(userId: number): Promise<Map<string, number>> {
const settlements = await this.transactionRequestService.getUsedSettlements(userId);

const consumed = new Map<string, number>();
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<string, number>,
address: string,
expectedShares: number,
minTimestamp: Date,
): HistoryEventDto | undefined {
const incomingTransfers = Util.sort(
history.filter((e) => e.transfer && Util.equalsIgnoreCase(e.transfer.to, address)),
'timestamp',
);

const seen = new Map<string, number>();

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}`;
}
}