diff --git a/.changeset/wait-for-transaction-leaks.md b/.changeset/wait-for-transaction-leaks.md new file mode 100644 index 0000000..0282593 --- /dev/null +++ b/.changeset/wait-for-transaction-leaks.md @@ -0,0 +1,9 @@ +--- +'@bigmi/core': patch +--- + +Keep the shared block watcher alive when a `waitForTransaction` times out. When a wait reached its `retryCount` timeout and the same block callback then found the transaction confirmed or replaced, or failed to fetch the block, its internal `done` ran a second time. That second pass ran the watcher's cleanup while another wait on the same client still listened, so that wait and every later wait on the client never settled. It could also settle a newer wait on the same txId with the stale error of the first one. `done` now runs at most once, and an observer's `unwatch` is a no-op once its own listener is gone. + +Release the observers of a settled `waitForTransaction`. A wait that joined another wait on the same txId, such as a resumed run, stayed in the module-level `listenersCache` with its callbacks after both settled, and a third wait on that txId joined it and never settled. Every finished wait also left an empty `listenersCache` key behind, one per txId. Each wait now removes its own listener when it settles. When the last listener leaves, `observe` drops the key and its `cleanupCache` entry and runs the cleanup, which removes the observer from the shared block watcher. A wait now joins another wait on the same txId only when both use the same `confirmations`, `pollingInterval`, `retryCount`, numeric `retryDelay` and `senderAddress` (a `retryDelay` function is not compared), so it no longer runs with the options of the wait it joined, such as resolving at 1 confirmation when it asked for 3. + +Make the `timeout` option of `waitForTransaction` end the wait. The timer only rejected the promise: it was never cleared and the block watcher kept polling, so a wait whose `getblockcount` never succeeded polled until the page or process ended. The timeout now rejects and removes only its own wait, so another wait on the same txId with a longer timeout or none keeps waiting, and the shared block watcher stops once no wait listens to it. Every way a wait settles clears its timer. diff --git a/.changeset/wait-for-transaction-viem-parity.md b/.changeset/wait-for-transaction-viem-parity.md new file mode 100644 index 0000000..7acd780 --- /dev/null +++ b/.changeset/wait-for-transaction-viem-parity.md @@ -0,0 +1,11 @@ +--- +'@bigmi/core': patch +--- + +Stop `waitForTransaction` from reporting the awaited transaction as its own replacement. When `getrawtransaction` still reported the transaction unconfirmed but `getblock` already listed it, the replacement search matched the transaction itself, because it spends the same inputs, and called `onReplaced` with it. A replacement that `getrawtransaction` reported with no confirmations also resolved the wait at once, and a replacement that needed more confirmations resolved it in a later block without `onReplaced`. The search now skips the awaited txid, a replacement without confirmations keeps the wait polling, and `onReplaced` fires once, when the replacement has enough confirmations. It reports the awaited transaction as `replacedTransaction` and finds the reason against it, also when the replacement replaced an earlier replacement that left the chain. + +Settle `withRetry` when `shouldRetry` or the `delay` function throws. The throw rejected an internal attempt that nothing handled, so the returned promise stayed pending and the error surfaced only as an unhandled rejection. `withRetry` now rejects with that error. A `retryDelay` function of `waitForTransaction` that threw hung the wait in the same way; the wait now rejects with that error. + +Wait for every confirmation of a mined transaction in `waitForTransaction`. The `retryCount` block budget also counted the blocks after the transaction was mined, so a wait for more confirmations than the budget had left rejected with `WaitForTransactionReceiptTimeoutError` while the transaction was confirming: with the default `retryCount` of 10, a wait for 6 confirmations rejected if the transaction was mined 6 or more blocks after the wait started. The budget now counts only the blocks in which the transaction, or the replacement it tracks, is not mined or the height of its block is unknown. An unmined transaction still rejects on the same block as before. + +Evict the least recently used key from the internal `LruMap` that holds the in-flight requests of request deduplication. A key that was set again kept its old place, a read moved a key to the newest place only when its value was not `undefined`, an empty-string key was never evicted, and the eviction read the keys through the subclass, which viem reports can give a stale iterator on iOS 18 JavaScriptCore. A key that is set again or read now becomes the newest, the oldest key is evicted also when it is an empty string, and the eviction reads the keys of the base `Map`. This matters only when more than 8192 deduplicated requests are in flight at once. diff --git a/packages/core/src/actions/waitForTransaction.spec.ts b/packages/core/src/actions/waitForTransaction.spec.ts new file mode 100644 index 0000000..0753572 --- /dev/null +++ b/packages/core/src/actions/waitForTransaction.spec.ts @@ -0,0 +1,1203 @@ +import { address, Block, Transaction } from 'bitcoinjs-lib' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { BlockNotFoundError } from '../errors/block.js' +import { WaitForTransactionReceiptTimeoutError } from '../errors/transaction.js' +import { createClient } from '../factories/createClient.js' +import { custom } from '../transports/custom.js' +import { cleanupCache, listenersCache } from '../utils/observe.js' +import { waitForTransaction } from './waitForTransaction.js' + +const POLLING_INTERVAL = 1_000 +const SENDER = 'bc1qsender' + +let txSeed = 0 + +/** + * A parseable non-coinbase transaction. It spends output 0 of the outpoint + * `spends`, so two transactions with the same `spends` replace each other. + */ +function makeTx( + spends = ++txSeed, + outputByte = 0 +): { txId: string; txHex: string } { + const prevHash = new Uint8Array(32) + new DataView(prevHash.buffer).setUint32(0, spends) + const tx = new Transaction() + tx.version = 2 + tx.addInput(prevHash, 0, 0xfffffffd) + tx.addOutput( + Uint8Array.from([0x00, 0x14, ...new Uint8Array(20).fill(outputByte)]), + 10_000n + ) + return { txId: tx.getId(), txHex: tx.toHex() } +} + +/** A block with a coinbase and `txs`. */ +function makeBlockHex(txs: Transaction[] = []): string { + const coinbase = new Transaction() + coinbase.version = 1 + coinbase.addInput( + new Uint8Array(32), + 0xffffffff, + 0xffffffff, + Uint8Array.from([1, 1]) + ) + coinbase.addOutput(Uint8Array.from([0x6a]), 0n) + const block = new Block() + block.version = 1 + block.prevHash = new Uint8Array(32) + block.merkleRoot = Block.calculateMerkleRoot([coinbase, ...txs]) + block.timestamp = 0 + block.bits = 0 + block.nonce = 0 + block.transactions = [coinbase, ...txs] + return block.toHex() +} + +type MockTx = { hex: string; confirmedAt?: number } + +/** + * A client on an in-memory chain. `state.height` is the tip, `state.txs` maps + * a txid to its raw hex and the height it confirms at, and `state.blockHex` is + * every block (only a coinbase by default, so no replacement is found). + */ +function createMockChain({ + fail, + growingConfirmations = false, + omitMempoolConfirmations = false, + blockStatsWithoutHeight = false, + caseInsensitiveTxIds = false, +}: { + fail?: (method: string, params: unknown[]) => boolean + /** + * Count the confirmations of a transaction from its block to the tip, and + * let `getblockstats` return the height of that block. When off, a + * confirmed transaction has 1 confirmation and `getblockstats` the tip. + */ + growingConfirmations?: boolean + /** + * Answer for a transaction that is not confirmed without `confirmations`, + * as Bitcoin Core does for a mempool transaction. When off, the answer has + * `confirmations: 0`. + */ + omitMempoolConfirmations?: boolean + /** Answer `getblockstats` without `height`. */ + blockStatsWithoutHeight?: boolean + /** Find a txid in either case in `getrawtransaction`, as Bitcoin Core does. */ + caseInsensitiveTxIds?: boolean +} = {}) { + const state = { + height: 100, + txs: new Map(), + blockHex: makeBlockHex(), + calls: {} as Record, + } + const request = async ({ + method, + params, + }: { + method: string + params: unknown[] + }) => { + state.calls[method] = (state.calls[method] ?? 0) + 1 + if (fail?.(method, params)) { + throw new Error(`${method} failed`) + } + switch (method) { + case 'getblockcount': + return state.height + case 'getrawtransaction': { + const txId = caseInsensitiveTxIds + ? (params[0] as string).toLowerCase() + : (params[0] as string) + const tx = state.txs.get(txId) + if (!tx) { + throw new Error('No such mempool or blockchain transaction') + } + const confirmed = + tx.confirmedAt !== undefined && state.height >= tx.confirmedAt + if (confirmed && growingConfirmations) { + return { + txid: txId, + hex: tx.hex, + // The hash carries the height, for `getblockstats`. + blockhash: `bh${tx.confirmedAt}`, + confirmations: state.height - tx.confirmedAt! + 1, + } + } + if (!confirmed && omitMempoolConfirmations) { + return { txid: txId, hex: tx.hex } + } + return confirmed + ? { txid: txId, hex: tx.hex, blockhash: 'bh', confirmations: 1 } + : { txid: txId, hex: tx.hex, confirmations: 0 } + } + case 'getblockstats': + if (blockStatsWithoutHeight) { + return {} + } + return growingConfirmations + ? { height: Number((params[0] as string).slice(2)) } + : { height: state.height } + case 'getblockhash': + return 'bh' + case 'getblock': + return state.blockHex + default: + throw new Error(`Unexpected method ${method}`) + } + } + const client = createClient({ + // No transport retries: they would add timers of their own. + transport: custom({ request }, { retryCount: 0 }), + pollingInterval: POLLING_INTERVAL, + }) + return { client, state } +} + +type Tracked = { + status: 'pending' | 'resolved' | 'rejected' + error?: unknown +} + +/** Records how a promise settles, so a test can check it without waiting. */ +function track(promise: Promise): Tracked { + const tracked: Tracked = { status: 'pending' } + promise.then( + () => { + tracked.status = 'resolved' + }, + (error) => { + tracked.status = 'rejected' + tracked.error = error + } + ) + return tracked +} + +/** Every observer key that `client` still holds in the shared caches. */ +function cacheKeysOf(client: { uid: string }): string[] { + const uid = JSON.stringify(client.uid) + return [...listenersCache.keys(), ...cleanupCache.keys()].filter((key) => + key.includes(uid) + ) +} + +/** The observer key of a wait with `SENDER` and the default options. */ +const waitId = (client: { uid: string }, txId: string) => + JSON.stringify([ + 'waitForTransaction', + client.uid, + txId, + { + confirmations: 1, + pollingInterval: POLLING_INTERVAL, + retryCount: 10, + retryDelay: 3_000, + senderAddress: SENDER, + }, + ]) + +const watchId = (client: { uid: string }) => + JSON.stringify(['watchBlockNumber', client.uid, true, true, POLLING_INTERVAL]) + +const advance = (ms: number) => vi.advanceTimersByTimeAsync(ms) + +describe('waitForTransaction', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.clearAllTimers() + vi.useRealTimers() + }) + + describe('timeout branch (count > retryCount)', () => { + // `retryCount: 1` makes the branch come on the 3rd block callback; the + // default (10) reaches the same branch on the 12th. + it('keeps the shared block watcher alive when the timed-out transaction confirms in that block', async () => { + const { client, state } = createMockChain() + const txA = makeTx() + const txB = makeTx() + const txC = makeTx() + state.txs.set(txA.txId, { hex: txA.txHex }) + state.txs.set(txB.txId, { hex: txB.txHex }) + const wait = (tx: { txId: string; txHex: string }) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 1, + }) + ) + + // A starts the shared block watcher: callback 1 at height 100. + const a = wait(txA) + await advance(0) + // Callback 2. + state.height = 101 + await advance(POLLING_INTERVAL) + // B joins the shared block watcher. + const b = wait(txB) + await advance(0) + expect(listenersCache.get(watchId(client))).toHaveLength(2) + + // A's callback 3 times out, and A is confirmed in that same block. + state.txs.get(txA.txId)!.confirmedAt = 102 + state.height = 102 + await advance(POLLING_INTERVAL) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + // B confirms in the next block. + state.txs.get(txB.txId)!.confirmedAt = 103 + state.height = 103 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('resolved') + + // A later wait on the same client gets its own block watcher. + state.txs.set(txC.txId, { hex: txC.txHex, confirmedAt: 0 }) + const c = wait(txC) + await advance(POLLING_INTERVAL) + expect(c.status).toBe('resolved') + + expect(listenersCache.get(watchId(client)) ?? []).toEqual([]) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not settle or stop a later wait on the same txId', async () => { + // getblockhash for height 102 (A's timeout block) always fails. + const { client, state } = createMockChain({ + fail: (method, params) => + method === 'getblockhash' && params[0] === 102, + }) + const tx = makeTx() + const txC = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = (t: { txId: string; txHex: string }) => + track( + waitForTransaction(client, { + ...t, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 1, + retryDelay: 300, + }) + ) + + // Callbacks 1 and 2 at heights 100 and 101. + const a = wait(tx) + await advance(0) + state.height = 101 + await advance(POLLING_INTERVAL) + // Callback 3 times out. + state.height = 102 + await advance(POLLING_INTERVAL) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + // The route is resumed: a wait on the same txId, inside the 300 ms in + // which A's timeout callback could still retry getBlock(102). + state.height = 103 + const b = wait(tx) + await advance(400) + expect(b.error).not.toBeInstanceOf(BlockNotFoundError) + expect(b.status).toBe('pending') + + // B's transaction confirms in the next block. + state.txs.get(tx.txId)!.confirmedAt = 104 + state.height = 104 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('resolved') + + // A later wait on the same client. + state.txs.set(txC.txId, { hex: txC.txHex, confirmedAt: 0 }) + state.height = 105 + const c = wait(txC) + await advance(POLLING_INTERVAL) + expect(c.status).toBe('resolved') + + expect(listenersCache.get(watchId(client)) ?? []).toEqual([]) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('rejects an unmined transaction on the 12th block with the default retryCount', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + }) + ) + + // Blocks 100 to 110: 11 callbacks, each looks the transaction up once. + await advance(0) + for (let height = 101; height <= 110; height++) { + state.height = height + await advance(POLLING_INTERVAL) + } + expect(wait.status).toBe('pending') + expect(state.calls.getrawtransaction).toBe(11) + + // Block 111: the 12th callback is past the budget. + state.height = 111 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('rejected') + expect(wait.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + expect(state.calls.getrawtransaction).toBe(11) + }) + + it('counts the callback of a missed block while the transaction is unmined', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 2, + }) + ) + + // Callback 1 at block 100. + await advance(0) + // One poll finds block 103, so 101 and 102 are missed blocks. Only the + // callback of 101 runs: the callbacks of 102 and 103 start while it + // still looks the transaction up, and return without counting. + state.height = 103 + await advance(POLLING_INTERVAL) + expect(state.calls.getrawtransaction).toBe(2) + // Callback 3 at block 104 is the last one in the budget. + state.height = 104 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('pending') + expect(state.calls.getrawtransaction).toBe(3) + + state.height = 105 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('rejected') + expect(wait.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + }) + + it('waits for every confirmation of a mined transaction past the budget', async () => { + const { client, state } = createMockChain({ growingConfirmations: true }) + const tx = makeTx() + // Mined at block 106; the wait starts at block 100. + state.txs.set(tx.txId, { hex: tx.txHex, confirmedAt: 106 }) + const promise = waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + confirmations: 6, + }) + const wait = track(promise) + + // Blocks 100 to 110: 11 callbacks, 5 of them after the transaction is + // mined. + await advance(0) + for (let height = 101; height <= 110; height++) { + state.height = height + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('pending') + } + + // Block 111 is the sixth confirmation. + state.height = 111 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: tx.txId }) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('waits for every confirmation of a mined replacement past the budget', async () => { + const { client, state } = createMockChain({ growingConfirmations: true }) + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + state.txs.set(original.txId, { hex: original.txHex }) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + ...original, + senderAddress: SENDER, + onReplaced, + confirmations: 6, + // The lookups of a missing transaction end within one poll. + retryDelay: 10, + }) + const wait = track(promise) + + // Blocks 100 to 110: 11 callbacks. At block 106 the original leaves + // the mempool, and its replacement is mined in that block. + await advance(0) + for (let height = 101; height <= 110; height++) { + if (height === 106) { + state.txs.delete(original.txId) + state.txs.set(replacement.txId, { + hex: replacement.txHex, + confirmedAt: 106, + }) + state.blockHex = makeBlockHex([ + Transaction.fromHex(replacement.txHex), + ]) + } + state.height = height + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('pending') + } + + // Block 111 is the sixth confirmation of the replacement. + state.height = 111 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: replacement.txId }) + expect(onReplaced).toHaveBeenCalledTimes(1) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('keeps the budget while the height of the block of a mined transaction is unknown', async () => { + const { client, state } = createMockChain({ + blockStatsWithoutHeight: true, + }) + const tx = makeTx() + // Mined at block 101, but getblockstats does not give its height. + state.txs.set(tx.txId, { hex: tx.txHex, confirmedAt: 101 }) + const wait = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + confirmations: 3, + }) + ) + + // Blocks 100 to 110: 11 callbacks, and all of them count. + await advance(0) + for (let height = 101; height <= 110; height++) { + state.height = height + await advance(POLLING_INTERVAL) + } + expect(wait.status).toBe('pending') + + // Block 111: the 12th callback is past the budget. + state.height = 111 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('rejected') + expect(wait.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + }) + }) + + describe('observer cleanup', () => { + it('removes every wait on one txId when the transaction confirms', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const id = waitId(client, tx.txId) + const wait = () => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + }) + ) + + // A drives the observer; B (a resumed run) joins it. + const a = wait() + await advance(0) + const b = wait() + await advance(0) + expect(listenersCache.get(id)).toHaveLength(2) + + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(a.status).toBe('resolved') + expect(b.status).toBe('resolved') + expect(listenersCache.has(id)).toBe(false) + + // A third wait on the same txId starts its own observer and settles. + const c = wait() + await advance(POLLING_INTERVAL) + expect(c.status).toBe('resolved') + + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('leaves no key behind for waits on distinct txIds', async () => { + const { client, state } = createMockChain() + const waits: Tracked[] = [] + for (let i = 0; i < 100; i++) { + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex, confirmedAt: 0 }) + waits.push( + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + }) + ) + ) + } + + await advance(POLLING_INTERVAL) + expect(waits.every((wait) => wait.status === 'resolved')).toBe(true) + + const waitKeys = cacheKeysOf(client).filter((key) => + key.includes('waitForTransaction') + ) + expect(waitKeys).toEqual([]) + expect(cacheKeysOf(client)).toEqual([]) + }) + }) + + describe('confirmations option', () => { + it('waits for its own confirmations next to a wait on the same txId that needs fewer', async () => { + const { client, state } = createMockChain({ growingConfirmations: true }) + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = (confirmations: number) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + confirmations, + }) + ) + + // A needs 1 confirmation; B starts after it and needs 3. + const a = wait(1) + await advance(0) + const b = wait(3) + await advance(0) + + // The first confirmation. + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(a.status).toBe('resolved') + expect(b.status).toBe('pending') + + // The second confirmation. + state.height = 102 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('pending') + + // The third confirmation. + state.height = 103 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('resolved') + expect(cacheKeysOf(client)).toEqual([]) + }) + }) + + describe('timeout option', () => { + it('stops polling when the timeout expires', async () => { + const { client, state } = createMockChain({ + fail: (method) => method === 'getblockcount', + }) + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + + const wait = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + // Off the polling grid, so the deadline is not on a poll tick. + timeout: 10_500, + }) + ) + + await advance(10_499) + expect(wait.status).toBe('pending') + expect(state.calls.getblockcount).toBeGreaterThan(0) + + await advance(1) + expect(wait.status).toBe('rejected') + expect(wait.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + const calls = state.calls.getblockcount + await advance(POLLING_INTERVAL * 5) + expect(state.calls.getblockcount).toBe(calls) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('clears the timer when the transaction confirms first', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex, confirmedAt: 0 }) + + const wait = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + timeout: 60_000, + }) + ) + await advance(POLLING_INTERVAL) + + expect(wait.status).toBe('resolved') + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('rejects only the joined wait when its own timeout expires', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const id = waitId(client, tx.txId) + const wait = (timeout?: number) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + timeout, + }) + ) + + // A drives the observer with no timeout; B joins it with one. + const a = wait() + await advance(0) + const b = wait(5_500) + await advance(5_500) + expect(b.status).toBe('rejected') + expect(b.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + expect(a.status).toBe('pending') + expect(listenersCache.get(id)).toHaveLength(1) + + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(a.status).toBe('resolved') + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('resolves a joined wait with a longer timeout after the driver times out', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = (timeout: number) => + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + timeout, + }) + + // A drives the observer; B joins it with a longer timeout. + const a = track(wait(1_500)) + await advance(0) + const promiseB = wait(20_000) + const b = track(promiseB) + await advance(1_500) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + expect(b.status).toBe('pending') + + // The transaction confirms in the block that the poll at 3 s finds. + await advance(1_000) + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('resolved') + await expect(promiseB).resolves.toMatchObject({ txid: tx.txId }) + + // Let the stopped poll's last sleep run out. + await advance(POLLING_INTERVAL * 3) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('resolves a joined wait without a timeout after the driver times out', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const id = waitId(client, tx.txId) + const wait = (timeout?: number) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + timeout, + }) + ) + + // A drives the observer with a timeout; B joins it with none. + const a = wait(1_500) + await advance(0) + const b = wait() + await advance(1_500) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + expect(b.status).toBe('pending') + expect(listenersCache.get(id)).toHaveLength(1) + + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(b.status).toBe('resolved') + + // Let the stopped poll's last sleep run out. + await advance(POLLING_INTERVAL * 3) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('clears every timer when the retryCount timeout settles a driver and a joiner', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = (timeout: number) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 1, + timeout, + }) + ) + + // A drives the observer; B joins it. Both set a long timeout. + const a = wait(3_600_000) + await advance(0) + const b = wait(7_200_000) + await advance(0) + // A's callback 3 (count 2 > retryCount 1) rejects every wait. + state.height = 101 + await advance(POLLING_INTERVAL) + state.height = 102 + await advance(POLLING_INTERVAL) + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + expect(b.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + // Let the stopped poll's last sleep run out. + await advance(POLLING_INTERVAL * 3) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not settle a newer wait on the same txId when an in-flight callback fails after the timeout', async () => { + // getblockhash for height 100 always fails, so A's first callback + // retries getBlock(100) until about 6 s. + const { client, state } = createMockChain({ + fail: (method, params) => + method === 'getblockhash' && params[0] === 100, + }) + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + + const a = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 3, + retryDelay: 2_000, + timeout: 1_500, + }) + ) + await advance(1_500) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + // A resumed run: C starts a new observer on the same txId. Its first + // callback is at height 101, so only A's callback touches block 100. + state.height = 101 + const c = track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + }) + ) + // A's getBlock(100) retries end with BlockNotFoundError. + await advance(8_000) + expect(c.error).toBeUndefined() + expect(c.status).toBe('pending') + + state.txs.get(tx.txId)!.confirmedAt = 102 + state.height = 102 + await advance(POLLING_INTERVAL) + expect(c.status).toBe('resolved') + + await advance(POLLING_INTERVAL * 3) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not settle a newer wait with the same options when an in-flight callback fails after the timeout', async () => { + // getblockhash for height 100 always fails, so A's first callback + // retries getBlock(100) until about 6 s. + const { client, state } = createMockChain({ + fail: (method, params) => + method === 'getblockhash' && params[0] === 100, + }) + const tx = makeTx() + state.txs.set(tx.txId, { hex: tx.txHex }) + const wait = (timeout?: number) => + track( + waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced: () => {}, + retryCount: 3, + retryDelay: 2_000, + timeout, + }) + ) + + const a = wait(1_500) + await advance(1_500) + expect(a.status).toBe('rejected') + expect(a.error).toBeInstanceOf(WaitForTransactionReceiptTimeoutError) + + // A resumed run: C has A's options, so it starts a new observer on + // A's key. Its first callback is at height 101, so only A's callback + // touches block 100. + state.height = 101 + const c = wait() + // A's getBlock(100) retries end with BlockNotFoundError. + await advance(8_000) + expect(c.error).toBeUndefined() + expect(c.status).toBe('pending') + + state.txs.get(tx.txId)!.confirmedAt = 102 + state.height = 102 + await advance(POLLING_INTERVAL) + expect(c.status).toBe('resolved') + + await advance(POLLING_INTERVAL * 3) + expect(vi.getTimerCount()).toBe(0) + expect(cacheKeysOf(client)).toEqual([]) + }) + }) + + describe('replacement', () => { + it('rejects with the error that onReplaced throws', async () => { + const { client, state } = createMockChain() + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + // The original left the mempool; its replacement is in the tip block. + state.txs.set(replacement.txId, { + hex: replacement.txHex, + confirmedAt: 0, + }) + state.blockHex = makeBlockHex([Transaction.fromHex(replacement.txHex)]) + const error = new Error('onReplaced failed') + + const wait = track( + waitForTransaction(client, { + ...original, + senderAddress: SENDER, + onReplaced: () => { + throw error + }, + retryCount: 0, + }) + ) + await advance(POLLING_INTERVAL) + + expect(wait.status).toBe('rejected') + expect(wait.error).toBe(error) + expect(listenersCache.get(watchId(client)) ?? []).toEqual([]) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not report the awaited transaction as its own replacement', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + // getrawtransaction still reports the transaction unconfirmed, but + // getblock already lists it. + state.txs.set(tx.txId, { hex: tx.txHex }) + state.blockHex = makeBlockHex([Transaction.fromHex(tx.txHex)]) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + ...tx, + senderAddress: SENDER, + onReplaced, + }) + const wait = track(promise) + + await advance(0) + expect(wait.status).toBe('pending') + expect(onReplaced).not.toHaveBeenCalled() + + // getrawtransaction reports it confirmed in the next block. + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: tx.txId }) + expect(onReplaced).not.toHaveBeenCalled() + expect(cacheKeysOf(client)).toEqual([]) + }) + + it.each([0, undefined])( + 'waits until a replacement that getrawtransaction reports unconfirmed is confirmed (confirmations %s)', + async (mempoolConfirmations) => { + const { client, state } = createMockChain({ + omitMempoolConfirmations: mempoolConfirmations === undefined, + }) + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + // The original left the mempool. Its replacement is in the tip block, + // but getrawtransaction reports it unconfirmed. + state.txs.set(replacement.txId, { hex: replacement.txHex }) + state.blockHex = makeBlockHex([Transaction.fromHex(replacement.txHex)]) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + ...original, + senderAddress: SENDER, + onReplaced, + // Enough block callbacks to find the replacement, check it again and + // resolve; the lookups of the original end within one poll. + retryCount: 3, + retryDelay: 100, + }) + const wait = track(promise) + + // Callback 1 finds the replacement. + await advance(POLLING_INTERVAL / 2) + expect(wait.status).toBe('pending') + expect(onReplaced).not.toHaveBeenCalled() + + // Callback 2: still unconfirmed. + state.height = 101 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('pending') + expect(onReplaced).not.toHaveBeenCalled() + + // Callback 3: confirmed. + state.txs.get(replacement.txId)!.confirmedAt = 102 + state.height = 102 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: replacement.txId }) + expect(onReplaced).toHaveBeenCalledTimes(1) + const [{ reason, replacedTransaction, transaction }] = + onReplaced.mock.calls[0]! + expect(reason).toBe('replaced') + expect(replacedTransaction.getId()).toBe(original.txId) + expect(transaction.txid).toBe(replacement.txId) + expect(cacheKeysOf(client)).toEqual([]) + } + ) + + it('reports a replacement that needs more confirmations once it has them', async () => { + const { client, state } = createMockChain({ growingConfirmations: true }) + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + // The original left the mempool. Its replacement is mined in the tip + // block, so it has 1 of the 2 confirmations. + state.txs.set(replacement.txId, { + hex: replacement.txHex, + confirmedAt: 100, + }) + state.blockHex = makeBlockHex([Transaction.fromHex(replacement.txHex)]) + const onReplaced = vi.fn() + const wait = track( + waitForTransaction(client, { + ...original, + senderAddress: SENDER, + onReplaced, + confirmations: 2, + // The lookups of the original end within one poll. + retryCount: 3, + retryDelay: 100, + }) + ) + + // Callback 1 finds the replacement with 1 confirmation. + await advance(POLLING_INTERVAL / 2) + expect(wait.status).toBe('pending') + expect(onReplaced).not.toHaveBeenCalled() + + // Callback 2: the second confirmation. + state.height = 101 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + expect(onReplaced).toHaveBeenCalledTimes(1) + expect(onReplaced.mock.calls[0]![0].transaction.txid).toBe( + replacement.txId + ) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('compares a replacement of a replacement with the awaited transaction', async () => { + const { client, state } = createMockChain() + const spends = ++txSeed + const original = makeTx(spends) + // A cancel pays the sender, and a fee bump of the cancel pays the + // sender less. + const cancel = makeTx(spends, 2) + const senderAddress = address.fromOutputScript( + Transaction.fromHex(cancel.txHex).outs[0]!.script + ) + const bump = Transaction.fromHex(cancel.txHex) + bump.outs[0]!.value = 9_000n + const bumpedCancel = { txId: bump.getId(), txHex: bump.toHex() } + // The original left the mempool. getblock lists the cancel in block + // 100, but getrawtransaction has it only in the mempool. + state.txs.set(cancel.txId, { hex: cancel.txHex }) + state.blockHex = makeBlockHex([Transaction.fromHex(cancel.txHex)]) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + ...original, + senderAddress, + onReplaced, + // The lookups of a missing transaction end within one poll. + retryCount: 3, + retryDelay: 100, + }) + const wait = track(promise) + + // Callback 1 finds the cancel with 0 confirmations. + await advance(POLLING_INTERVAL / 2) + expect(wait.status).toBe('pending') + + // Block 100 is reorged out, and the bumped cancel is mined in 101. + state.txs.delete(cancel.txId) + state.txs.set(bumpedCancel.txId, { + hex: bumpedCancel.txHex, + confirmedAt: 101, + }) + state.blockHex = makeBlockHex([bump]) + state.height = 101 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: bumpedCancel.txId }) + expect(onReplaced).toHaveBeenCalledTimes(1) + const [{ reason, replacedTransaction, transaction }] = + onReplaced.mock.calls[0]! + expect(reason).toBe('cancelled') + expect(replacedTransaction.getId()).toBe(original.txId) + expect(transaction.txid).toBe(bumpedCancel.txId) + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not report the awaited transaction as a replacement of its replacement', async () => { + const { client, state } = createMockChain() + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + // The original left the mempool. getblock lists its replacement in + // block 100, but getrawtransaction has it only in the mempool. + state.txs.set(replacement.txId, { hex: replacement.txHex }) + state.blockHex = makeBlockHex([Transaction.fromHex(replacement.txHex)]) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + ...original, + senderAddress: SENDER, + onReplaced, + // The lookups of a missing transaction end within one poll. + retryCount: 3, + retryDelay: 100, + }) + const wait = track(promise) + + // Callback 1 tracks the replacement with no confirmations. + await advance(POLLING_INTERVAL / 2) + expect(wait.status).toBe('pending') + + // Block 100 is reorged out, and the original is mined in block 101. + state.txs.delete(replacement.txId) + state.txs.set(original.txId, { hex: original.txHex, confirmedAt: 101 }) + state.blockHex = makeBlockHex([Transaction.fromHex(original.txHex)]) + state.height = 101 + await advance(POLLING_INTERVAL) + expect(onReplaced).not.toHaveBeenCalled() + + // The next callback looks the original up again and finds it mined. + state.height = 102 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: original.txId }) + expect(onReplaced).not.toHaveBeenCalled() + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('does not report the awaited transaction as a replacement of its replacement for a txId in upper case', async () => { + const { client, state } = createMockChain({ caseInsensitiveTxIds: true }) + const spends = ++txSeed + const original = makeTx(spends) + const replacement = makeTx(spends, 1) + // The original left the mempool. getblock lists its replacement in + // block 100, but getrawtransaction has it only in the mempool. + state.txs.set(replacement.txId, { hex: replacement.txHex }) + state.blockHex = makeBlockHex([Transaction.fromHex(replacement.txHex)]) + const onReplaced = vi.fn() + const promise = waitForTransaction(client, { + txId: original.txId.toUpperCase(), + txHex: original.txHex, + senderAddress: SENDER, + onReplaced, + // The lookups of a missing transaction end within one poll. + retryCount: 3, + retryDelay: 100, + }) + const wait = track(promise) + + // Callback 1 tracks the replacement with no confirmations. + await advance(POLLING_INTERVAL / 2) + expect(wait.status).toBe('pending') + + // Block 100 is reorged out, and the original is mined in block 101. + state.txs.delete(replacement.txId) + state.txs.set(original.txId, { hex: original.txHex, confirmedAt: 101 }) + state.blockHex = makeBlockHex([Transaction.fromHex(original.txHex)]) + state.height = 101 + await advance(POLLING_INTERVAL) + expect(onReplaced).not.toHaveBeenCalled() + + // The next callback looks the original up again and finds it mined. + state.height = 102 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: original.txId }) + expect(onReplaced).not.toHaveBeenCalled() + expect(cacheKeysOf(client)).toEqual([]) + }) + + it('resolves when txHex does not parse but the node knows the transaction', async () => { + const { client, state } = createMockChain() + const tx = makeTx() + // getrawtransaction answers with the hex of the transaction, so the + // wait needs txHex only to find a replacement. + state.txs.set(tx.txId, { hex: tx.txHex }) + const promise = waitForTransaction(client, { + txId: tx.txId, + txHex: 'zz', + senderAddress: SENDER, + onReplaced: () => {}, + }) + const wait = track(promise) + + // Callback 1: in the mempool, and the search finds no replacement. + await advance(0) + expect(wait.status).toBe('pending') + + state.txs.get(tx.txId)!.confirmedAt = 101 + state.height = 101 + await advance(POLLING_INTERVAL) + expect(wait.status).toBe('resolved') + await expect(promise).resolves.toMatchObject({ txid: tx.txId }) + expect(cacheKeysOf(client)).toEqual([]) + }) + }) +}) diff --git a/packages/core/src/actions/waitForTransaction.ts b/packages/core/src/actions/waitForTransaction.ts index 5a0d8e2..402132a 100644 --- a/packages/core/src/actions/waitForTransaction.ts +++ b/packages/core/src/actions/waitForTransaction.ts @@ -56,12 +56,17 @@ export type WaitForTransactionReceiptParameters = { */ pollingInterval?: number | undefined /** - * Number of times to retry if the transaction or block is not found. - * @default 6 (exponential backoff) + * Number of times to retry a failed lookup of the transaction or of a block. + * It is also the block budget: once `retryCount + 1` block callbacks have + * counted, the next one rejects with `WaitForTransactionReceiptTimeoutError`. + * A callback counts while the transaction is not mined or the height of its + * block is unknown. + * @default 10 */ retryCount?: number /** - * Time to wait (in ms) between retries. + * Time to wait (in ms) between the retries of a lookup. + * @default 3_000 */ retryDelay?: ((config: { count: number; error: Error }) => number) | number /** Optional timeout (in milliseconds) to wait before stopping polling. */ @@ -105,28 +110,82 @@ export async function waitForTransaction( timeout, }: WaitForTransactionReceiptParameters ): Promise { - const observerId = stringify(['waitForTransaction', client.uid, txId]) + const observerId = stringify([ + 'waitForTransaction', + client.uid, + txId, + // The first wait's closure decides how every wait on its observer + // confirms, polls, retries and labels a replacement, so only waits with + // the same options share one. A `retryDelay` function cannot be compared + // and is left out. Each wait owns its `timeout`, so it is left out too. + { + confirmations, + pollingInterval, + retryCount, + retryDelay: typeof retryDelay === 'number' ? retryDelay : undefined, + senderAddress, + }, + ]) + // `getId()` gives a txid in lower case; the node takes `txId` in either case. + const awaitedTxId = txId.toLowerCase() let count = 0 let transaction: UTXOTransaction | undefined + // The height of the block of `transaction`, once `getblockstats` gives it. + let minedHeight: number | undefined let replacedTransaction: Transaction | undefined + // The replacement that `transaction` tracks once one is found. It is + // reported when `transaction` has enough confirmations, which can be in a + // later block. + let replacement: Omit | undefined let retrying = false return new Promise((resolve, reject) => { - if (timeout) { - setTimeout( - () => - reject( - new WaitForTransactionReceiptTimeoutError({ hash: txId as never }) - ), - timeout - ) + let timer: ReturnType | undefined + // Every way a wait settles goes through here: it clears its own timer and + // removes only its own listener. When the last wait leaves, `observe` runs + // the cleanup returned below, which unwatches the shared block watcher. + const settle = (fn: () => void) => { + clearTimeout(timer) + _unobserve() + fn() } const _unobserve = observe( observerId, - { onReplaced, resolve, reject }, + { + onReplaced, + resolve: (transaction: WaitForTransactionReceiptReturnType) => + settle(() => resolve(transaction)), + reject: (error: unknown) => settle(() => reject(error)), + }, (emit) => { + // Settles the waits at most once. A callback that is still in flight + // must not emit a second time. + let finished = false + const done = (fn: () => void) => { + if (finished) { + return + } + finished = true + try { + fn() + } catch (error) { + // A throwing `onReplaced` still settles the wait. + emit.reject(error) + } + } + + // Resolves with the tracked transaction, and first reports the + // replacement it tracks, if any. + const resolveWith = (transaction: UTXOTransaction) => + done(() => { + if (replacement) { + emit.onReplaced?.({ ...replacement, transaction }) + } + emit.resolve(transaction) + }) + const _unwatch = getAction( client, watchBlockNumber, @@ -136,12 +195,6 @@ export async function waitForTransaction( emitOnBegin: true, pollingInterval, async onBlockNumber(blockNumber_) { - const done = (fn: () => void) => { - _unwatch() - fn() - _unobserve() - } - let blockNumber = blockNumber_ if (retrying) { @@ -155,6 +208,7 @@ export async function waitForTransaction( }) ) ) + return } try { @@ -169,6 +223,7 @@ export async function waitForTransaction( blockHash: transaction.blockhash, stats: ['height'], }) + minedHeight = blockStats.height || undefined if ( confirmations > 1 && (!blockStats.height || @@ -176,7 +231,7 @@ export async function waitForTransaction( ) { return } - done(() => emit.resolve(transaction!)) + resolveWith(transaction) return } @@ -206,6 +261,7 @@ export async function waitForTransaction( blockHash: transaction.blockhash, stats: ['height'], }) + minedHeight = blockStats.height || undefined if (blockStats.height) { blockNumber = blockStats.height } @@ -224,7 +280,7 @@ export async function waitForTransaction( return } - done(() => emit.resolve(transaction!)) + resolveWith(transaction) } catch (err) { // If the receipt is not found, the transaction will be pending. // We need to check if it has potentially been replaced. @@ -273,6 +329,8 @@ export async function waitForTransaction( } let replacementTransaction: Transaction | undefined + let originalTransactionInBlock = false + const replacedTransactionId = replacedTransaction.getId() for (const tx of block.transactions!) { if (tx.isCoinbase()) { @@ -288,15 +346,35 @@ export async function waitForTransaction( const vout = input.index const inputId = `${txid}:${vout}` if (replacedTransactionInputs.has(inputId)) { - replacementTransaction = tx + // The tracked and the awaited transaction spend the + // same inputs, and a provider can list one in a block + // before getrawtransaction reports it mined. Neither + // is a replacement: a later callback finds it mined. + // `awaitedTxId` is the awaited one; the tracked one + // differs from it once a replacement is tracked. + const id = tx.getId() + if (id === awaitedTxId) { + originalTransactionInBlock = true + } else if (id !== replacedTransactionId) { + replacementTransaction = tx + } break } } - if (replacementTransaction) { + if (replacementTransaction || originalTransactionInBlock) { break } } + // The awaited transaction is back after its tracked + // replacement left the chain: track the awaited one again, + // so the next callback looks it up by `txId`. + if (originalTransactionInBlock) { + transaction = undefined + replacement = undefined + return + } + // If we couldn't find a replacement transaction, continue polling. if (!replacementTransaction) { return @@ -311,14 +389,6 @@ export async function waitForTransaction( txId: replacementTransaction.getId(), }) - // Check if we have enough confirmations. If not, continue polling. - if ( - transaction.confirmations && - transaction.confirmations < confirmations - ) { - return - } - let reason: ReplacementReason = 'replaced' // Function to get output addresses @@ -337,9 +407,12 @@ export async function waitForTransaction( return addresses } - // Get the recipient addresses from the original transaction + // Get the recipient addresses from the original transaction. + // That is the awaited one, also when the tracked transaction + // is an earlier replacement that left the chain. + const originalTransaction = Transaction.fromHex(txHex) const originalOutputAddresses = - getOutputAddresses(replacedTransaction) + getOutputAddresses(originalTransaction) // Get the recipient addresses from the replacement transaction const replacementOutputAddresses = getOutputAddresses( @@ -362,14 +435,22 @@ export async function waitForTransaction( reason = 'cancelled' } - done(() => { - emit.onReplaced?.({ - reason, - replacedTransaction: replacedTransaction!, - transaction: transaction!, - }) - emit.resolve(transaction!) - }) + replacement = { + reason, + replacedTransaction: originalTransaction, + } + + // Check if we have enough confirmations. If not, continue + // polling. A replacement with no confirmations is not mined + // yet, so a later callback checks it again. + if ( + !transaction.confirmations || + transaction.confirmations < confirmations + ) { + return + } + + resolveWith(transaction) } catch (err_) { done(() => emit.reject(err_)) } @@ -377,11 +458,37 @@ export async function waitForTransaction( done(() => emit.reject(err)) } } finally { - count++ + // The budget counts only blocks in which the tracked transaction + // is not mined, or the height of its block is unknown. A mined + // one waits for its confirmations, which can take more blocks + // than `retryCount`. + if (!transaction?.blockhash || !minedHeight) { + count++ + } } }, }) + + // `observe` runs this when the last wait on this observer leaves. It + // also ends this observer: a callback that is still in flight must + // not emit to a newer wait that starts a new observer on the same id. + return () => { + finished = true + _unwatch() + } } ) + + if (timeout) { + timer = setTimeout( + () => + settle(() => + reject( + new WaitForTransactionReceiptTimeoutError({ hash: txId as never }) + ) + ), + timeout + ) + } }) } diff --git a/packages/core/src/utils/lru.spec.ts b/packages/core/src/utils/lru.spec.ts new file mode 100644 index 0000000..f93e1d7 --- /dev/null +++ b/packages/core/src/utils/lru.spec.ts @@ -0,0 +1,47 @@ +import { describe, expect, it } from 'vitest' +import { LruMap } from './lru.js' + +describe('LruMap', () => { + it('makes a key that is set again the newest', () => { + const cache = new LruMap(2) + cache.set('a', 1) + cache.set('b', 2) + cache.set('a', 3) + cache.set('c', 4) + + expect([...cache.keys()]).toEqual(['a', 'c']) + expect(cache.get('a')).toBe(3) + }) + + it.each([1, undefined])( + 'makes a key that is read the newest (value %s)', + (value) => { + const cache = new LruMap(2) + cache.set('a', value) + cache.set('b', 2) + cache.get('a') + cache.set('c', 3) + + expect([...cache.keys()]).toEqual(['a', 'c']) + } + ) + + it('evicts an empty-string key when it is the oldest', () => { + const cache = new LruMap(1) + cache.set('', 1) + cache.set('x', 2) + + expect([...cache.keys()]).toEqual(['x']) + }) + + it('keeps at most maxSize keys', () => { + const cache = new LruMap(100) + for (let i = 0; i < 1_000; i++) { + cache.set(`key${i}`, i) + } + + expect(cache.size).toBe(100) + expect(cache.has('key899')).toBe(false) + expect(cache.has('key900')).toBe(true) + }) +}) diff --git a/packages/core/src/utils/lru.ts b/packages/core/src/utils/lru.ts index 93ce747..57239e3 100644 --- a/packages/core/src/utils/lru.ts +++ b/packages/core/src/utils/lru.ts @@ -14,20 +14,26 @@ export class LruMap extends Map { override get(key: string): value | undefined { const value = super.get(key) - if (super.has(key) && value !== undefined) { - this.delete(key) - super.set(key, value) + if (super.has(key)) { + super.delete(key) + super.set(key, value as value) } return value } override set(key: string, value: value): this { + // Delete first, so a key that is set again becomes the newest. + if (super.has(key)) { + super.delete(key) + } super.set(key, value) if (this.maxSize && this.size > this.maxSize) { - const firstKey = this.keys().next().value - if (firstKey) { - this.delete(firstKey) + // viem reports a stale iterator of a Map subclass on iOS 18 + // JavaScriptCore; `super` avoids it. An empty string is a key too. + const firstKey = super.keys().next().value + if (firstKey !== undefined) { + super.delete(firstKey) } } return this diff --git a/packages/core/src/utils/observe.spec.ts b/packages/core/src/utils/observe.spec.ts new file mode 100644 index 0000000..f59575a --- /dev/null +++ b/packages/core/src/utils/observe.spec.ts @@ -0,0 +1,65 @@ +import { describe, expect, it, vi } from 'vitest' +import { cleanupCache, listenersCache, observe } from './observe.js' + +let observerCount = 0 + +/** Two listeners on one observer id; the first one starts `fn`. */ +function setup() { + const id = `observe.spec.${++observerCount}` + const cleanup = vi.fn() + let emit: { onData: (data: number) => void } | undefined + const fn = vi.fn((emit_: { onData: (data: number) => void }) => { + emit = emit_ + return cleanup + }) + const onDataA = vi.fn() + const onDataB = vi.fn() + const unwatchA = observe(id, { onData: onDataA }, fn) + const unwatchB = observe(id, { onData: onDataB }, fn) + return { id, cleanup, emit: emit!, fn, onDataA, onDataB, unwatchA, unwatchB } +} + +describe('observe', () => { + it('runs the cleanup and drops the key when the last listener leaves', () => { + const { id, cleanup, fn, unwatchA, unwatchB } = setup() + expect(fn).toHaveBeenCalledTimes(1) + + unwatchA() + expect(cleanup).not.toHaveBeenCalled() + expect(listenersCache.get(id)).toHaveLength(1) + + unwatchB() + expect(cleanup).toHaveBeenCalledTimes(1) + expect(listenersCache.has(id)).toBe(false) + expect(cleanupCache.has(id)).toBe(false) + }) + + it('ignores a repeated unwatch', () => { + const { id, cleanup, emit, onDataB, unwatchA, unwatchB } = setup() + + unwatchA() + unwatchA() + expect(cleanup).not.toHaveBeenCalled() + expect(listenersCache.get(id)).toHaveLength(1) + emit.onData(1) + expect(onDataB).toHaveBeenCalledWith(1) + + unwatchB() + unwatchB() + expect(cleanup).toHaveBeenCalledTimes(1) + }) + + it('starts a new observer on the same id after the last listener leaves', () => { + const { id, cleanup, fn, unwatchA, unwatchB } = setup() + unwatchA() + unwatchB() + + const unwatchC = observe(id, { onData: vi.fn() }, fn) + expect(fn).toHaveBeenCalledTimes(2) + expect(cleanupCache.get(id)).toBe(cleanup) + + unwatchC() + expect(cleanup).toHaveBeenCalledTimes(2) + expect(listenersCache.has(id)).toBe(false) + }) +}) diff --git a/packages/core/src/utils/observe.ts b/packages/core/src/utils/observe.ts index e22f6d5..8ff9e0f 100644 --- a/packages/core/src/utils/observe.ts +++ b/packages/core/src/utils/observe.ts @@ -29,20 +29,24 @@ export function observe( const getListeners = () => listenersCache.get(observerId) || [] - const unsubscribe = () => { - const listeners = getListeners() - listenersCache.set( - observerId, - listeners.filter((cb: any) => cb.id !== callbackId) - ) - } - const unwatch = () => { - const cleanup = cleanupCache.get(observerId) - if (getListeners().length === 1 && cleanup) { - cleanup() + const listeners = getListeners() + // A late or repeated call must not run the cleanup of the observers that + // are still listening. + if (!listeners.some((cb) => cb.id === callbackId)) { + return + } + const remaining = listeners.filter((cb) => cb.id !== callbackId) + if (remaining.length > 0) { + listenersCache.set(observerId, remaining) + return } - unsubscribe() + // The last listener is gone: drop the key and its cleanup, so neither + // map keeps one entry per finished observer. + const cleanup = cleanupCache.get(observerId) + listenersCache.delete(observerId) + cleanupCache.delete(observerId) + cleanup?.() } const listeners = getListeners() diff --git a/packages/core/src/utils/withRetry.spec.ts b/packages/core/src/utils/withRetry.spec.ts new file mode 100644 index 0000000..f43ca74 --- /dev/null +++ b/packages/core/src/utils/withRetry.spec.ts @@ -0,0 +1,150 @@ +import { describe, expect, it } from 'vitest' +import { withRetry } from './withRetry.js' + +type Outcome = { + status: 'pending' | 'resolved' | 'rejected' + value?: unknown + error?: unknown + /** Every unhandled rejection while the promise ran. */ + unhandled: unknown[] +} + +/** + * Records how `promise` settles. With no delay, `withRetry` runs on + * microtasks only, so a timer later it has settled or never will, and Node + * has reported every unhandled rejection. + */ +async function outcomeOf(promise: Promise): Promise { + const outcome: Outcome = { status: 'pending', unhandled: [] } + const onUnhandled = (reason: unknown) => { + outcome.unhandled.push(reason) + } + process.on('unhandledRejection', onUnhandled) + promise.then( + (value) => { + outcome.status = 'resolved' + outcome.value = value + }, + (error) => { + outcome.status = 'rejected' + outcome.error = error + } + ) + try { + await new Promise((resolve) => setTimeout(resolve, 10)) + return outcome + } finally { + process.off('unhandledRejection', onUnhandled) + } +} + +/** A function that fails `failures` times, then returns `'ok'`. */ +function failing(failures = Number.POSITIVE_INFINITY) { + let calls = 0 + const fn = async () => { + calls++ + if (calls <= failures) { + throw new Error(`fn failed ${calls}`) + } + return 'ok' + } + return { fn, calls: () => calls } +} + +describe('withRetry', () => { + it('rejects with the error that shouldRetry throws after a retry', async () => { + const error = new Error('shouldRetry failed') + const outcome = await outcomeOf( + withRetry(failing().fn, { + delay: 0, + shouldRetry: ({ count }) => { + if (count === 0) { + return true + } + throw error + }, + }) + ) + + expect(outcome.status).toBe('rejected') + expect(outcome.error).toBe(error) + expect(outcome.unhandled).toEqual([]) + }) + + it('rejects with the error that an async shouldRetry rejects with', async () => { + const error = new Error('shouldRetry failed') + const outcome = await outcomeOf( + withRetry(failing().fn, { + delay: 0, + shouldRetry: async ({ count }) => { + if (count === 0) { + return true + } + throw error + }, + }) + ) + + expect(outcome.status).toBe('rejected') + expect(outcome.error).toBe(error) + expect(outcome.unhandled).toEqual([]) + }) + + it('rejects with the error that the delay function throws after a retry', async () => { + const error = new Error('delay failed') + const outcome = await outcomeOf( + withRetry(failing().fn, { + delay: ({ count }) => { + if (count === 0) { + return 0 + } + throw error + }, + }) + ) + + expect(outcome.status).toBe('rejected') + expect(outcome.error).toBe(error) + expect(outcome.unhandled).toEqual([]) + }) + + it('retries and then resolves', async () => { + const { fn, calls } = failing(2) + const delays: number[] = [] + const outcome = await outcomeOf( + withRetry(fn, { + delay: ({ count }) => { + delays.push(count) + return 0 + }, + retryCount: 2, + }) + ) + + expect(outcome.status).toBe('resolved') + expect(outcome.value).toBe('ok') + expect(calls()).toBe(3) + expect(delays).toEqual([0, 1]) + }) + + it('rejects with the last error when the retries run out', async () => { + const { fn, calls } = failing() + const outcome = await outcomeOf(withRetry(fn, { delay: 0, retryCount: 2 })) + + expect(outcome.status).toBe('rejected') + expect((outcome.error as Error).message).toBe('fn failed 3') + expect(calls()).toBe(3) + expect(outcome.unhandled).toEqual([]) + }) + + it('does not retry when shouldRetry returns false', async () => { + const { fn, calls } = failing() + const outcome = await outcomeOf( + withRetry(fn, { delay: 0, shouldRetry: () => false }) + ) + + expect(outcome.status).toBe('rejected') + expect((outcome.error as Error).message).toBe('fn failed 1') + expect(calls()).toBe(1) + }) +}) diff --git a/packages/core/src/utils/withRetry.ts b/packages/core/src/utils/withRetry.ts index a4a0079..3e71ed0 100644 --- a/packages/core/src/utils/withRetry.ts +++ b/packages/core/src/utils/withRetry.ts @@ -32,14 +32,14 @@ export function withRetry( }: WithRetryParameters = {} ): Promise { return new Promise((resolve, reject) => { - const attemptRetry = async ({ count = 0 } = {}) => { - const retry = async ({ error }: { error: Error }) => { + const attemptRetry = async ({ count = 0 } = {}): Promise => { + const retry = async ({ error }: { error: Error }): Promise => { const delay = typeof delay_ === 'function' ? delay_({ count, error }) : delay_ if (delay) { await wait(delay) } - attemptRetry({ count: count + 1 }) + return attemptRetry({ count: count + 1 }) } try { @@ -55,6 +55,8 @@ export function withRetry( reject(err) } } - attemptRetry() + // A throw from `shouldRetry` or `delay` rejects an attempt; pass it on + // instead of leaving the promise pending. + attemptRetry().catch(reject) }) }