Compare commits

...
3 Commits
Author SHA1 Message Date
nobswebdev b3e683cd71 Merge pull request 'Slave/btc rbf edge case fix' (#5) from slave/btc-rbf-edge-case-fix into master
Reviewed-on: #5
2026-09-08 17:02:47 +00:00
nobswebdev 4458edb0d1 Prune stale unconfirmed invoice payments during wallet polling.
Refactor processInvoice into upsert and prune helpers so absent 0-conf payments are removed when the wallet no longer reports them.
2026-09-08 18:30:16 +02:00
nobswebdev dcd499664a Filter superseded Bitcoin transfers in Electrum wallet client.
Detect RBF conflicts via shared input outpoints so invoice polling no longer ingests both the original and replacement mempool transactions.
2026-09-08 18:30:06 +02:00
9 changed files with 354 additions and 48 deletions
@@ -34,13 +34,17 @@ describe('ElectrumWalletRpcClient', () => {
tx_hash: 'abc123', tx_hash: 'abc123',
height: 800_000 height: 800_000
}, },
'50000', {
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
},
'bc1qtest',
800_002 800_002
) )
).toEqual({ ).toEqual({
txHash: 'abc123', txHash: 'abc123',
amountAtomic: '50000', amountAtomic: '50000',
confirmations: 3 confirmations: 3,
inputOutpoints: []
}); });
}); });
@@ -51,13 +55,17 @@ describe('ElectrumWalletRpcClient', () => {
tx_hash: 'abc123', tx_hash: 'abc123',
height: 0 height: 0
}, },
'50000', {
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
},
'bc1qtest',
800_002 800_002
) )
).toEqual({ ).toEqual({
txHash: 'abc123', txHash: 'abc123',
amountAtomic: '50000', amountAtomic: '50000',
confirmations: 0 confirmations: 0,
inputOutpoints: []
}); });
}); });
@@ -68,11 +76,61 @@ describe('ElectrumWalletRpcClient', () => {
tx_hash: 'abc123', tx_hash: 'abc123',
height: 800_000 height: 800_000
}, },
'0', {
outputs: [{ address: 'bc1qother', value_sats: 10_000 }]
},
'bc1qtest',
800_002 800_002
) )
).toBeNull(); ).toBeNull();
}); });
it('maps input outpoints from transaction inputs', () => {
expect(
clientTest.mapIncomingTransfer(
{
tx_hash: 'abc123',
height: 800_000
},
{
inputs: [
{ prevout_hash: 'abc123', prevout_n: 0 },
{ prevout_hash: 'def456', prevout_n: 2 }
],
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
},
'bc1qtest',
800_002
)
).toEqual({
txHash: 'abc123',
amountAtomic: '50000',
confirmations: 3,
inputOutpoints: ['abc123:0', 'def456:2']
});
});
it('skips coinbase-like inputs without prevout data', () => {
expect(
clientTest.mapIncomingTransfer(
{
tx_hash: 'abc123',
height: 800_000
},
{
inputs: [{}, { prevout_hash: '', prevout_n: 0 }],
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
},
'bc1qtest',
800_002
)
).toEqual({
txHash: 'abc123',
amountAtomic: '50000',
confirmations: 3,
inputOutpoints: []
});
});
}); });
describe('sumOutputValueAtomic', () => { describe('sumOutputValueAtomic', () => {
@@ -125,6 +183,7 @@ describe('ElectrumWalletRpcClient', () => {
jsonrpc: '2.0', jsonrpc: '2.0',
id: 'nullcart', id: 'nullcart',
result: { result: {
inputs: [{ prevout_hash: 'input123', prevout_n: 0 }],
outputs: [ outputs: [
{ address: 'bc1qother', value_sats: 10_000 }, { address: 'bc1qother', value_sats: 10_000 },
{ address: 'bc1qtest', value_sats: 50_000 } { address: 'bc1qtest', value_sats: 50_000 }
@@ -137,7 +196,8 @@ describe('ElectrumWalletRpcClient', () => {
{ {
txHash: 'abc123', txHash: 'abc123',
amountAtomic: '50000', amountAtomic: '50000',
confirmations: 3 confirmations: 3,
inputOutpoints: ['input123:0']
} }
]); ]);
@@ -164,6 +224,105 @@ describe('ElectrumWalletRpcClient', () => {
expect.any(Object) expect.any(Object)
); );
}); });
it('drops superseded unconfirmed transfers that share inputs with a confirmed replacement', async () => {
mockedAxios.post
.mockResolvedValueOnce({
data: {
jsonrpc: '2.0',
id: 'nullcart',
result: [
{ tx_hash: 'original', height: 0 },
{ tx_hash: 'replacement', height: 800_000 }
]
}
})
.mockResolvedValueOnce({
data: {
jsonrpc: '2.0',
id: 'nullcart',
result: '01000000'
}
})
.mockResolvedValueOnce({
data: {
jsonrpc: '2.0',
id: 'nullcart',
result: {
inputs: [{ prevout_hash: 'shared-input', prevout_n: 0 }],
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
}
}
})
.mockResolvedValueOnce({
data: {
jsonrpc: '2.0',
id: 'nullcart',
result: '02000000'
}
})
.mockResolvedValueOnce({
data: {
jsonrpc: '2.0',
id: 'nullcart',
result: {
inputs: [{ prevout_hash: 'shared-input', prevout_n: 0 }],
outputs: [{ address: 'bc1qtest', value_sats: 50_000 }]
}
}
});
await expect(client.getIncomingTransfers('bc1qtest', 800_002)).resolves.toEqual([
{
txHash: 'replacement',
amountAtomic: '50000',
confirmations: 3,
inputOutpoints: ['shared-input:0']
}
]);
});
});
describe('filterSupersededBitcoinTransfers', () => {
it('keeps unrelated transfers unchanged', () => {
const transfers = [
{ txHash: 'a', amountAtomic: '1', confirmations: 0, inputOutpoints: ['in1:0'] },
{ txHash: 'b', amountAtomic: '1', confirmations: 3, inputOutpoints: ['in2:1'] }
];
expect(clientTest.filterSupersededBitcoinTransfers(transfers)).toEqual(transfers);
});
it('drops an unconfirmed transfer superseded by a confirmed replacement', () => {
const transfers = [
{ txHash: 'original', amountAtomic: '1', confirmations: 0, inputOutpoints: ['in1:0'] },
{ txHash: 'replacement', amountAtomic: '1', confirmations: 2, inputOutpoints: ['in1:0'] }
];
expect(clientTest.filterSupersededBitcoinTransfers(transfers)).toEqual([
{ txHash: 'replacement', amountAtomic: '1', confirmations: 2, inputOutpoints: ['in1:0'] }
]);
});
it('keeps the later unconfirmed transfer when both conflict before confirmation', () => {
const transfers = [
{ txHash: 'original', amountAtomic: '1', confirmations: 0, inputOutpoints: ['in1:0'] },
{ txHash: 'replacement', amountAtomic: '1', confirmations: 0, inputOutpoints: ['in1:0'] }
];
expect(clientTest.filterSupersededBitcoinTransfers(transfers)).toEqual([
{ txHash: 'replacement', amountAtomic: '1', confirmations: 0, inputOutpoints: ['in1:0'] }
]);
});
it('keeps transfers without input outpoints', () => {
const transfers = [
{ txHash: 'coinbase', amountAtomic: '1', confirmations: 0, inputOutpoints: [] },
{ txHash: 'payment', amountAtomic: '1', confirmations: 1, inputOutpoints: ['in1:0'] }
];
expect(clientTest.filterSupersededBitcoinTransfers(transfers)).toEqual(transfers);
});
}); });
describe('createAddress', () => { describe('createAddress', () => {
@@ -200,21 +200,26 @@ export class ElectrumWalletRpcClient {
const serializedTransaction = await this.call<string>('gettransaction', { txid: txHash }); const serializedTransaction = await this.call<string>('gettransaction', { txid: txHash });
const transaction = await this.deserializeTransaction(serializedTransaction); const transaction = await this.deserializeTransaction(serializedTransaction);
const amountAtomic = this.sumOutputValueAtomic(transaction, address);
return this.mapIncomingTransfer(entry, amountAtomic, blockHeight); return this.mapIncomingTransfer(entry, transaction, address, blockHeight);
}) })
); );
return transfers.filter((transfer): transfer is ElectrumWalletIncomingTransfer => transfer !== null); const resolvedTransfers = transfers.filter(
(transfer): transfer is ElectrumWalletIncomingTransfer => transfer !== null
);
return this.filterSupersededBitcoinTransfers(resolvedTransfers);
} }
private mapIncomingTransfer( private mapIncomingTransfer(
entry: ElectrumWalletAddressHistoryEntry, entry: ElectrumWalletAddressHistoryEntry,
amountAtomic: string, transaction: ElectrumWalletDeserializedTransaction,
address: string,
blockHeight: number | null blockHeight: number | null
): ElectrumWalletIncomingTransfer | null { ): ElectrumWalletIncomingTransfer | null {
const txHash = entry.tx_hash; const txHash = entry.tx_hash;
const amountAtomic = this.sumOutputValueAtomic(transaction, address);
if (!txHash || amountAtomic === '0') { if (!txHash || amountAtomic === '0') {
return null; return null;
@@ -226,7 +231,60 @@ export class ElectrumWalletRpcClient {
return { return {
txHash, txHash,
amountAtomic, amountAtomic,
confirmations confirmations,
inputOutpoints: this.extractInputOutpoints(transaction)
}; };
} }
private filterSupersededBitcoinTransfers(
transfers: ElectrumWalletIncomingTransfer[]
): ElectrumWalletIncomingTransfer[] {
return transfers.filter((transfer, index) => {
const isSuperseded = transfers.some((other, otherIndex) => {
if (otherIndex === index) {
return false;
}
if (!this.sharesInputOutpoint(transfer.inputOutpoints, other.inputOutpoints)) {
return false;
}
if (other.confirmations > transfer.confirmations) {
return true;
}
return other.confirmations === transfer.confirmations && otherIndex > index;
});
return !isSuperseded;
});
}
private extractInputOutpoints(transaction: ElectrumWalletDeserializedTransaction): string[] {
if (!Array.isArray(transaction.inputs)) {
return [];
}
return transaction.inputs.flatMap(input => {
if (typeof input.prevout_hash !== 'string' || input.prevout_hash.length === 0) {
return [];
}
if (typeof input.prevout_n !== 'number') {
return [];
}
return [`${input.prevout_hash}:${input.prevout_n}`];
});
}
private sharesInputOutpoint(left: readonly string[], right: readonly string[]): boolean {
if (left.length === 0 || right.length === 0) {
return false;
}
const rightOutpoints = new Set(right);
return left.some(outpoint => rightOutpoints.has(outpoint));
}
} }
@@ -0,0 +1,4 @@
export type ElectrumWalletDeserializedInput = {
prevout_hash?: string;
prevout_n?: number;
};
@@ -0,0 +1,4 @@
export type ElectrumWalletDeserializedOutput = {
address?: string;
value_sats: number;
};
@@ -1,8 +1,7 @@
export type ElectrumWalletDeserializedOutput = { import type { ElectrumWalletDeserializedInput } from './ElectrumWalletDeserializedInput';
address?: string; import type { ElectrumWalletDeserializedOutput } from './ElectrumWalletDeserializedOutput';
value_sats: number;
};
export type ElectrumWalletDeserializedTransaction = { export type ElectrumWalletDeserializedTransaction = {
inputs?: ElectrumWalletDeserializedInput[];
outputs: ElectrumWalletDeserializedOutput[]; outputs: ElectrumWalletDeserializedOutput[];
}; };
@@ -2,4 +2,5 @@ export type ElectrumWalletIncomingTransfer = {
txHash: string; txHash: string;
amountAtomic: string; amountAtomic: string;
confirmations: number; confirmations: number;
inputOutpoints: string[];
}; };
@@ -5,8 +5,12 @@ import type { ElectrumWalletIncomingTransfer } from './ElectrumWalletIncomingTra
export type ElectrumWalletRpcClientTest = { export type ElectrumWalletRpcClientTest = {
mapIncomingTransfer: ( mapIncomingTransfer: (
entry: ElectrumWalletAddressHistoryEntry, entry: ElectrumWalletAddressHistoryEntry,
amountAtomic: string, transaction: ElectrumWalletDeserializedTransaction,
address: string,
blockHeight: number | null blockHeight: number | null
) => ElectrumWalletIncomingTransfer | null; ) => ElectrumWalletIncomingTransfer | null;
filterSupersededBitcoinTransfers: (
transfers: ElectrumWalletIncomingTransfer[]
) => ElectrumWalletIncomingTransfer[];
sumOutputValueAtomic: (transaction: ElectrumWalletDeserializedTransaction, address: string) => string; sumOutputValueAtomic: (transaction: ElectrumWalletDeserializedTransaction, address: string) => string;
}; };
@@ -87,6 +87,7 @@ describe('InvoicePaymentService', () => {
}; };
let paymentRepo: { let paymentRepo: {
update: jest.Mock; update: jest.Mock;
delete: jest.Mock;
createQueryBuilder: jest.Mock; createQueryBuilder: jest.Mock;
}; };
let insertQueryBuilder: { let insertQueryBuilder: {
@@ -133,6 +134,7 @@ describe('InvoicePaymentService', () => {
paymentRepo = { paymentRepo = {
update: jest.fn().mockResolvedValue(undefined), update: jest.fn().mockResolvedValue(undefined),
delete: jest.fn().mockResolvedValue(undefined),
createQueryBuilder: jest.fn().mockReturnValue(insertQueryBuilder) createQueryBuilder: jest.fn().mockReturnValue(insertQueryBuilder)
}; };
@@ -384,13 +386,35 @@ describe('InvoicePaymentService', () => {
expect(paymentRepo.update).not.toHaveBeenCalled(); expect(paymentRepo.update).not.toHaveBeenCalled();
}); });
it('does nothing when there are no transfers to process', async () => { it('does not mutate payments when there are no transfers and no existing payments', async () => {
transactionalInvoiceQueryBuilder.getOne.mockResolvedValue(buildXmrInvoice()); transactionalInvoiceQueryBuilder.getOne.mockResolvedValue(buildXmrInvoice());
await service.processInvoice('invoice-1', []); await service.processInvoice('invoice-1', []);
expect(paymentRepo.createQueryBuilder).not.toHaveBeenCalled(); expect(paymentRepo.createQueryBuilder).not.toHaveBeenCalled();
expect(paymentRepo.update).not.toHaveBeenCalled(); expect(paymentRepo.update).not.toHaveBeenCalled();
expect(paymentRepo.delete).not.toHaveBeenCalled();
});
it('removes unconfirmed payments that are no longer reported when transfers are empty', async () => {
transactionalInvoiceQueryBuilder.getOne.mockResolvedValue(
buildBtcInvoice({
payments: [
{
id: 'payment-ghost',
txHash: 'ghost',
amountAtomic: '50000',
confirmations: 0
} as InvoicePayment
]
})
);
await service.processInvoice('invoice-btc-1', []);
expect(paymentRepo.delete).toHaveBeenCalledWith('payment-ghost');
expect(paymentRepo.createQueryBuilder).not.toHaveBeenCalled();
expect(paymentRepo.update).not.toHaveBeenCalled();
}); });
it('skips transfers below the configured minimum', async () => { it('skips transfers below the configured minimum', async () => {
@@ -528,5 +552,34 @@ describe('InvoicePaymentService', () => {
confirmations: 1 confirmations: 1
}); });
}); });
it('removes unconfirmed payments that are no longer reported', async () => {
transactionalInvoiceQueryBuilder.getOne.mockResolvedValue(
buildBtcInvoice({
payments: [
{
id: 'payment-original',
txHash: 'original',
amountAtomic: '50000',
confirmations: 0
} as InvoicePayment,
{
id: 'payment-replacement',
txHash: 'replacement',
amountAtomic: '50000',
confirmations: 3
} as InvoicePayment
]
})
);
await service.processInvoice('invoice-btc-1', [
buildBtcTransfer({ txHash: 'replacement', amountAtomic: '50000', confirmations: 3 })
]);
expect(paymentRepo.delete).toHaveBeenCalledWith('payment-original');
expect(paymentRepo.update).not.toHaveBeenCalled();
expect(paymentRepo.createQueryBuilder).not.toHaveBeenCalled();
});
}); });
}); });
@@ -135,8 +135,6 @@ export class InvoicePaymentService {
} }
private async processInvoice(invoiceId: string, transfers: InvoiceIncomingTransfer[]): Promise<void> { private async processInvoice(invoiceId: string, transfers: InvoiceIncomingTransfer[]): Promise<void> {
const { minByMethod } = this.configService.get('invoice') as Config['invoice'];
await this.dataSource.transaction(async manager => { await this.dataSource.transaction(async manager => {
const invoiceRepo = manager.getRepository(Invoice); const invoiceRepo = manager.getRepository(Invoice);
const paymentRepo = manager.getRepository(InvoicePayment); const paymentRepo = manager.getRepository(InvoicePayment);
@@ -152,6 +150,19 @@ export class InvoicePaymentService {
return; return;
} }
await this.upsertIncomingPayments(paymentRepo, invoice, transfers);
await this.pruneAbsentUnconfirmedPayments(paymentRepo, invoice, transfers);
});
}
private async upsertIncomingPayments(
paymentRepo: Repository<InvoicePayment>,
invoice: Invoice,
transfers: InvoiceIncomingTransfer[]
): Promise<void> {
const { minByMethod } = this.configService.get('invoice') as Config['invoice'];
const minIncomingAtomic = minByMethod[invoice.paymentMethod]; const minIncomingAtomic = minByMethod[invoice.paymentMethod];
const knownByTxHash = new Map((invoice.payments ?? []).map(payment => [payment.txHash, payment])); const knownByTxHash = new Map((invoice.payments ?? []).map(payment => [payment.txHash, payment]));
@@ -174,7 +185,7 @@ export class InvoicePaymentService {
.createQueryBuilder() .createQueryBuilder()
.insert() .insert()
.values({ .values({
invoice: { id: invoiceId }, invoice: { id: invoice.id },
txHash: transfer.txHash, txHash: transfer.txHash,
amountAtomic: transfer.amountAtomic, amountAtomic: transfer.amountAtomic,
confirmations: transfer.confirmations confirmations: transfer.confirmations
@@ -182,6 +193,19 @@ export class InvoicePaymentService {
.orIgnore() .orIgnore()
.execute(); .execute();
} }
}); }
private async pruneAbsentUnconfirmedPayments(
paymentRepo: Repository<InvoicePayment>,
invoice: Invoice,
transfers: InvoiceIncomingTransfer[]
): Promise<void> {
const activeTxHashes = new Set(transfers.map(transfer => transfer.txHash));
for (const payment of invoice.payments ?? []) {
if (payment.confirmations === 0 && !activeTxHashes.has(payment.txHash)) {
await paymentRepo.delete(payment.id);
}
}
} }
} }