diff --git a/__tests__/interruptedOperations.test.ts b/__tests__/interruptedOperations.test.ts new file mode 100644 index 00000000..1ae3915d --- /dev/null +++ b/__tests__/interruptedOperations.test.ts @@ -0,0 +1,214 @@ +/** + * The interrupted-operation resolver's decision table. + * + * Startup holds a reservation whose process died mid-swap or mid-melt; the resolver + * asks the mint what it did and settles by that. Real ProofsStore/MintsStore on the + * real database (reservations, rollback and commit run for real); the mint calls, + * TransferOperationApi.refresh and the queue are stubbed. + * + * @jest-environment node + */ +jest.mock('../src/services/logService', () => ({ + log: {debug: jest.fn(), error: jest.fn(), info: jest.fn(), trace: jest.fn(), warn: jest.fn()}, +})) +jest.mock('../src/services/nostrService', () => ({NostrClient: {getFirstTagValue: jest.fn()}})) +jest.mock('../src/services/syncQueueService', () => ({SyncQueue: {addPrioritizedTask: jest.fn()}})) +jest.mock('../src/services/wallet/operations/transferOperationApi', () => ({ + TransferOperationApi: {refresh: jest.fn()}, +})) + +const mockWalletStore = { + getProofsStatesFromMint: jest.fn(), + getCachedSeed: jest.fn(async () => new Uint8Array(64)), + restore: jest.fn(), +} +const mockTransactionsStore = {findById: jest.fn()} +// One root for the whole file: the resolver captures its stores at import time. +let mockRoot: any + +jest.mock('../src/models', () => ({ + get rootStoreInstance() { + return mockRoot + }, +})) + +import {applySnapshot, types} from 'mobx-state-tree' +import {MintsStoreModel} from '../src/models/MintsStore' +import {ProofsStoreModel} from '../src/models/ProofsStore' +import {TransactionModel, TransactionStatus} from '../src/models/Transaction' +import {Database} from '../src/services/db' +import {NetworkError} from '../src/utils/AppError' +import {TransferOperationApi} from '../src/services/wallet/operations/transferOperationApi' + +const TestRoot = types + .model('RootStore', { + mintsStore: types.optional(MintsStoreModel, {}), + proofsStore: types.optional(ProofsStoreModel, {}), + // A real Transaction node: commitReservation mirrors its tx update via setProp. + tx: types.maybe(TransactionModel), + }) + .volatile(() => ({transactionsStore: mockTransactionsStore, walletStore: mockWalletStore})) + +const MINT_URL = 'https://mint.test' +const KEYSET = 'keyset1' +const TX_ID = 30 + +const proof = (secret: string, amount: number) => ({ + id: KEYSET, amount, secret, C: 'C' + secret, unit: 'sat', tId: 1, mintUrl: MINT_URL, state: 'UNSPENT', +}) + +const txSnapshot = { + id: TX_ID, type: 'TRANSFER', amount: 90, unit: 'sat', mint: MINT_URL, + status: TransactionStatus.DRAFT, data: JSON.stringify([{status: 'DRAFT'}]), +} +const lastAudit = (tx: any) => JSON.parse(tx.data).at(-1) + +/** The mint reports the locked inputs in `bucket`; anything else (restored outputs) UNSPENT. */ +const mintSays = (bucket: 'SPENT' | 'PENDING' | 'UNSPENT') => + mockWalletStore.getProofsStatesFromMint.mockImplementation(async (_u: string, _unit: string, ps: any[]) => { + const inputs = ps.filter(p => p.secret.startsWith('in')) + const others = ps.filter(p => !p.secret.startsWith('in')) + return { + SPENT: bucket === 'SPENT' ? inputs : [], + PENDING: bucket === 'PENDING' ? inputs : [], + UNSPENT: bucket === 'UNSPENT' ? [...inputs, ...others] : others, + } + }) + +/** The state a process leaves when it dies inside the operation, after a restart. */ +function interrupted(operationType: string, counters?: {start: number; count: number; next: number}) { + Database.getInstance().executeBatch([ + ['DELETE FROM proofs'], ['DELETE FROM reservations'], ['DELETE FROM mint_counters'], + ['DELETE FROM transactions'], + [ + 'INSERT INTO transactions (id, type, amount, fee, unit, mint, status, data, createdAt) VALUES (?, ?, ?, 0, ?, ?, ?, ?, ?)', + [TX_ID, 'TRANSFER', 90, 'sat', MINT_URL, txSnapshot.status, txSnapshot.data, new Date().toISOString()], + ], + ]) + const proofs = [proof('in1', 64), proof('in2', 32), proof('free', 4)] + mockRoot.proofsStore.interruptedReservations.clear() + applySnapshot(mockRoot, { + mintsStore: { + mints: [{ + id: 'mint1111', mintUrl: MINT_URL, units: ['sat'], + keysets: [{id: KEYSET, unit: 'sat', active: true}], + proofsCounters: [{keyset: KEYSET, unit: 'sat', counter: 40}], + }], + } as any, + proofsStore: {proofs: Object.fromEntries(proofs.map(p => [p.secret, p])) as any}, + tx: txSnapshot, + } as any) + const {proofsStore} = mockRoot + Database.addOrUpdateProofs([...proofsStore.proofs.values()], 'UNSPENT') + + const reservation = proofsStore.reserve([proofsStore.getBySecret('in1'), proofsStore.getBySecret('in2')], { + transactionId: TX_ID, mintUrl: MINT_URL, unit: 'sat', operationType, rollbackTo: 'UNSPENT', + }) + if (counters) Database.setReservationCounters(reservation.id, {keysetId: KEYSET, ...counters}) + Database.addInFlightRequest(TX_ID, {amount: 90, proofs: []}) + + proofsStore.recoverOrphanReservations() // the restart + const tx = mockRoot.tx + mockTransactionsStore.findById.mockReturnValue(tx) + jest.spyOn(tx, 'update') + return {proofsStore, tx, reservation} +} + +mockRoot = TestRoot.create({}) +// Imported after the root exists, since it destructures the stores on load. +const {InterruptedOperationService} = require('../src/services/wallet/operations/interruptedOperations') +const run = () => InterruptedOperationService.resolveInterruptedOperationsTask() +const state = (s: string) => mockRoot.proofsStore.getBySecret(s)?.state +const openRows = () => Database.getOpenReservations().length + +beforeEach(() => { + jest.clearAllMocks() +}) + +test('inputs UNSPENT at the mint: nothing happened — roll back, tx REVERTED', async () => { + const {proofsStore, tx} = interrupted('transfer-swap', {start: 40, count: 3, next: 43}) + mintSays('UNSPENT') + + await run() + + expect([state('in1'), state('in2')]).toEqual(['UNSPENT', 'UNSPENT']) + expect(openRows()).toBe(0) + expect(tx.status).toBe(TransactionStatus.REVERTED) + expect(lastAudit(tx)).toMatchObject({status: 'REVERTED', interrupted: true}) + expect(Database.getInFlightRequest(TX_ID)).toBeUndefined() + expect(proofsStore.interruptedReservations.size).toBe(0) +}) + +test('swap executed: outputs restored from the recorded range, tx REVERTED', async () => { + const {proofsStore, tx} = interrupted('transfer-swap', {start: 40, count: 3, next: 43}) + mintSays('SPENT') + // The range can also hold outputs the wallet already has (counter reuse) — skipped. + mockWalletStore.restore.mockResolvedValue({ + proofs: [ + {id: KEYSET, amount: 64, secret: 'out1', C: 'Cout1'}, + {id: KEYSET, amount: 32, secret: 'out2', C: 'Cout2'}, + {id: KEYSET, amount: 4, secret: 'free', C: 'Cfree'}, + ], + }) + + await run() + + expect(mockWalletStore.restore).toHaveBeenCalledWith(MINT_URL, expect.any(Uint8Array), { + indexFrom: 40, indexTo: 43, keysetId: KEYSET, unit: 'sat', + }) + expect([state('in1'), state('in2')]).toEqual(['SPENT', 'SPENT']) + expect([state('out1'), state('out2')]).toEqual(['UNSPENT', 'UNSPENT']) + expect(proofsStore.getBySecret('out1').tId).toBe(TX_ID) + expect(state('free')).toBe('UNSPENT') + expect(proofsStore.getMintBalance(MINT_URL).balances.sat).toBe(64 + 32 + 4) + // Never derive from that range again. + expect(mockRoot.mintsStore.findByUrl(MINT_URL).getProofsCounter(KEYSET).counter).toBe(43) + expect(tx.status).toBe(TransactionStatus.REVERTED) + expect(lastAudit(tx)).toMatchObject({interrupted: true, restoredAmount: 96}) + expect(openRows()).toBe(0) + expect(Database.getInFlightRequest(TX_ID)).toBeUndefined() +}) + +test('swap executed but no range recorded (pre-v36 row): inputs SPENT, tx ERROR', async () => { + const {tx} = interrupted('send-online-swap') + mintSays('SPENT') + + await run() + + expect(mockWalletStore.restore).not.toHaveBeenCalled() + expect([state('in1'), state('in2')]).toEqual(['SPENT', 'SPENT']) + expect(tx.status).toBe(TransactionStatus.ERROR) + expect(lastAudit(tx).message).toMatch(/seed recovery/) + expect(openRows()).toBe(0) +}) + +test.each(['SPENT', 'PENDING'] as const)( + 'melt with inputs %s: handed to refresh as an ordinary PENDING transfer', + async bucket => { + const {proofsStore, tx} = interrupted('transfer-melt') + mintSays(bucket) + + await run() + + expect(openRows()).toBe(0) + expect([state('in1'), state('in2')]).toEqual(['PENDING', 'PENDING']) + expect(tx.status).toBe(TransactionStatus.PENDING) + expect(TransferOperationApi.refresh).toHaveBeenCalledWith(TX_ID) + // No longer held: the regular pending sweep may now see these proofs. + expect(proofsStore.isHeldByInterruptedOperation('in1')).toBe(false) + }, +) + +test('mint unreachable: stays held for the next sweep, mint marked OFFLINE', async () => { + const {proofsStore, tx} = interrupted('transfer-swap', {start: 40, count: 3, next: 43}) + mockWalletStore.getProofsStatesFromMint.mockRejectedValue(new NetworkError('Network request failed')) + + const result = await run() + + expect(result.errors).toHaveLength(1) + expect([state('in1'), state('in2')]).toEqual(['PENDING', 'PENDING']) + expect(openRows()).toBe(1) + expect(proofsStore.interruptedReservations.size).toBe(1) + expect(tx.update).not.toHaveBeenCalled() + expect(mockRoot.mintsStore.findByUrl(MINT_URL).status).toBe('OFFLINE') +}) diff --git a/__tests__/interruptedReservations.test.ts b/__tests__/interruptedReservations.test.ts index 100e57f2..bad11777 100644 --- a/__tests__/interruptedReservations.test.ts +++ b/__tests__/interruptedReservations.test.ts @@ -90,7 +90,10 @@ describe('recoverOrphanReservations', () => { expect(dbState(secret)).toBe('PENDING') } expect(Database.getOpenReservations().map(r => r.id)).toEqual([swap.id]) - expect([...proofsStore.interruptedReservationIds]).toEqual([swap.id]) + expect([...proofsStore.interruptedReservations.keys()]).toEqual([swap.id]) + expect(proofsStore.isHeldByInterruptedOperation('swapIn1')).toBe(true) + expect(proofsStore.isHeldByInterruptedOperation('free')).toBe(false) + expect(proofsStore.isTransactionInterrupted(20)).toBe(true) // Held proofs are not spendable. expect(proofsStore.getMintBalance(MINT_URL)!.balances.sat).toBe(4 + 8) @@ -103,7 +106,7 @@ describe('recoverOrphanReservations', () => { expect(proofsStore.recoverOrphanReservations()).toEqual({recoveredCount: 0, heldCount: 1}) proofsStore.releaseInterruptedReservation(swap.id) - expect(proofsStore.interruptedReservationIds.size).toBe(0) + expect(proofsStore.interruptedReservations.size).toBe(0) }) test.each(['transfer-swap', 'transfer-melt', 'transfer-melt-after-swap', 'send-online-swap'])( @@ -117,7 +120,65 @@ describe('recoverOrphanReservations', () => { proofsStore.recoverOrphanReservations() - expect(proofsStore.interruptedReservationIds.has(swap.id)).toBe(true) + expect(proofsStore.interruptedReservations.has(swap.id)).toBe(true) }, ) }) + +describe('revertAbandonedDrafts', () => { + const insertTx = (id: number, type: string, status: string, quote: string | null = null) => + Database.getInstance().execute( + `INSERT INTO transactions (id, type, amount, fee, unit, mint, status, quote, data, createdAt) + VALUES (?, ?, 1, 0, 'sat', ?, ?, ?, ?, ?)`, + [id, type, MINT_URL, status, quote, JSON.stringify([{status: 'DRAFT'}]), new Date().toISOString()], + ) + const txRow = (id: number) => + Database.getInstance().execute('SELECT status, data FROM transactions WHERE id = ?', [id]).rows?.item(0) + + function setup() { + const {proofsStore, swap} = crashedMidOperations() + Database.getInstance().execute('DELETE FROM transactions') + + insertTx(40, 'TRANSFER', 'DRAFT', 'quote-40') // died after its preemptive swap committed + insertTx(41, 'TRANSFER', 'DRAFT') // Nostr invoice waiting for the user: no quote yet + insertTx(20, 'TRANSFER', 'DRAFT', 'quote-20') // owns the held transfer-swap reservation + insertTx(43, 'TOPUP', 'PREPARED', 'quote-43') // invoice issued, waiting for payment + insertTx(44, 'TRANSFER', 'EXECUTING', 'quote-44') // a mint call may have happened + insertTx(45, 'SEND', 'PREPARED') + + // The swap's outputs, committed PENDING under tx 40 before the melt reserved them. + const out = proofsStore.getBySecret('free')! + out.setProp('state', 'PENDING') + out.setProp('tId', 40) + Database.addOrUpdateProofs([out], 'PENDING') + + proofsStore.recoverOrphanReservations() + return {proofsStore, swap} + } + + test('reverts abandoned transfers and sends, releasing their PENDING proofs', () => { + const {proofsStore} = setup() + + expect(proofsStore.revertAbandonedDrafts()).toEqual({revertedCount: 2}) + + for (const id of [40, 45]) { + expect(txRow(id).status).toBe('REVERTED') + expect(JSON.parse(txRow(id).data).at(-1)).toMatchObject({status: 'REVERTED', interrupted: true}) + } + expect(JSON.parse(txRow(40).data).at(-1).releasedAmount).toBe(4) + expect(proofsStore.getBySecret('free')!.state).toBe('UNSPENT') + expect(dbState('free')).toBe('UNSPENT') + }) + + test('leaves alone what may still be live or waiting', () => { + const {proofsStore} = setup() + + proofsStore.revertAbandonedDrafts() + + expect(txRow(41).status).toBe('DRAFT') // Nostr invoice + expect(txRow(20).status).toBe('DRAFT') // held for the resolver + expect(txRow(43).status).toBe('PREPARED') // topup + expect(txRow(44).status).toBe('EXECUTING') + expect(proofsStore.getBySecret('swapIn1')!.state).toBe('PENDING') + }) +}) diff --git a/src/models/ProofsStore.ts b/src/models/ProofsStore.ts index 537031ee..1770cfd3 100644 --- a/src/models/ProofsStore.ts +++ b/src/models/ProofsStore.ts @@ -7,6 +7,7 @@ import { } from 'mobx-state-tree' import { withSetPropAction } from './helpers/withSetPropAction' import { ProofModel, Proof, ProofRecord, ProofState } from './Proof' + import { TransactionData, TransactionStatus } from './Transaction' import { log } from '../services/logService' import { getRootStore } from './helpers/getRootStore' import AppError, { Err } from '../utils/AppError' @@ -33,11 +34,24 @@ import { // the interrupted-operation resolver. In memory only: the rows themselves are in // SQLite, and the next launch rebuilds this list from them. .volatile(() => ({ - interruptedReservationIds: new Set(), + interruptedReservations: new Map}>(), })) // ───────────────────── VIEWS ───────────────────── .views(self => ({ + /** Locked by an interrupted operation awaiting the resolver — no one else may settle it. */ + isHeldByInterruptedOperation(secret: string): boolean { + for (const held of self.interruptedReservations.values()) { + if (held.secrets.has(secret)) return true + } + return false + }, + isTransactionInterrupted(transactionId: number): boolean { + for (const held of self.interruptedReservations.values()) { + if (held.transactionId === transactionId) return true + } + return false + }, getBySecret(secret: string): Proof | undefined { return self.proofs.get(secret) }, @@ -668,7 +682,10 @@ import { let recoveredCount = 0 for (const orphan of orphans) { if (INTERRUPTIBLE_OPERATION_TYPES.has(orphan.operationType)) { - self.interruptedReservationIds.add(orphan.id) + self.interruptedReservations.set(orphan.id, { + transactionId: orphan.transactionId, + secrets: new Set(orphan.lockedProofs.map(p => p.secret)), + }) log.warn('[recoverOrphanReservations] Holding interrupted operation for mint check', { id: orphan.id, transactionId: orphan.transactionId, @@ -698,12 +715,58 @@ import { } } - return { recoveredCount, heldCount: self.interruptedReservationIds.size } + return { recoveredCount, heldCount: self.interruptedReservations.size } + }, + + /** + * Close out outgoing operations a previous process abandoned before reaching + * the mint (see Database.getAbandonedDraftTransactions): tx → REVERTED with an + * `interrupted` audit entry, and any proofs still PENDING under it released. + * Those can only be a preemptive swap's outputs, committed PENDING just before + * the process died and before the melt reserved them — fresh, unspent ecash + * that nothing else would ever release. + * + * Startup only, after recoverOrphanReservations and before any operation can + * start. Writes the database directly: transactions are loaded afterwards. + */ + revertAbandonedDrafts(): { revertedCount: number } { + let revertedCount = 0 + for (const draft of Database.getAbandonedDraftTransactions()) { + try { + const pending = self + .getByTransactionId(draft.id) + .filter(p => p.state === 'PENDING' && !self.isHeldByInterruptedOperation(p.secret)) + if (pending.length > 0) { + Database.addOrUpdateProofs(pending, 'UNSPENT') + for (const p of pending) p.state = 'UNSPENT' + } + + let data: TransactionData[] = [] + try { + data = JSON.parse(draft.data) + } catch {} + data.push({ + status: TransactionStatus.REVERTED, + interrupted: true, + message: 'Interrupted before reaching the mint. Nothing was paid.', + ...(pending.length > 0 && {releasedAmount: pending.reduce((sum, p) => sum + p.amount, 0)}), + createdAt: new Date(), + }) + Database.updateTransaction(draft.id, { + status: TransactionStatus.REVERTED, + data: JSON.stringify(data), + }) + revertedCount++ + } catch (e: any) { + log.error('[revertAbandonedDrafts] Could not revert', {transactionId: draft.id, error: e.message}) + } + } + return { revertedCount } }, /** The resolver settled this interrupted reservation; stop tracking it. */ releaseInterruptedReservation(reservationId: string): void { - self.interruptedReservationIds.delete(reservationId) + self.interruptedReservations.delete(reservationId) }, })) diff --git a/src/models/helpers/setupRootStore.ts b/src/models/helpers/setupRootStore.ts index 75fdd5d7..97122f5a 100644 --- a/src/models/helpers/setupRootStore.ts +++ b/src/models/helpers/setupRootStore.ts @@ -142,6 +142,12 @@ export async function setupRootStore(rootStore: RootStore, opts: SetupRootStoreO if (recoveredCount > 0 || heldCount > 0) { log.warn('[setupRootStore] Orphan proof reservations', {rolledBack: recoveredCount, heldForMintCheck: heldCount}) } + + // Before transactions load, so they come up with their final status. + const { revertedCount } = proofsStore.revertAbandonedDrafts() + if (revertedCount > 0) { + log.warn('[setupRootStore] Reverted abandoned draft transactions', {revertedCount}) + } } const orphansRecovered = performance.now() diff --git a/src/screens/WalletScreen.tsx b/src/screens/WalletScreen.tsx index 96f0cdb3..a2296920 100644 --- a/src/screens/WalletScreen.tsx +++ b/src/screens/WalletScreen.tsx @@ -343,6 +343,9 @@ export const WalletScreen = observer(function WalletScreen({ route }: Props) { if (nowInSec - lastMintCheckRef.current > MINT_CHECK_INTERVAL) { lastMintCheckRef.current = nowInSec + // Settle operations a killed process left mid-way. The sweeps below + // skip what these hold, so ordering between them does not matter. + WalletTask.resolveInterruptedQueue() WalletTask.handleInFlightQueue() WalletTask.handlePendingQueue() await WalletTask.syncStateWithAllMintsQueueAwaitable({proofState: 'PENDING'}) diff --git a/src/services/db/index.ts b/src/services/db/index.ts index 4f64ef61..59f71380 100644 --- a/src/services/db/index.ts +++ b/src/services/db/index.ts @@ -10,6 +10,7 @@ import {getDatabaseVersion} from './migrations' import { getTransactionsCount, getTransactionById, + getAbandonedDraftTransactions, getLastTransactionBy, getRecentTransactionsByUnitAsync, getTransactionsAsync, @@ -117,6 +118,7 @@ export const Database = { cleanAll, getTransactionsCount, getTransactionById, + getAbandonedDraftTransactions, getLastTransactionBy, getRecentTransactionsByUnitAsync, getTransactionsAsync, diff --git a/src/services/db/transactionsRepo.ts b/src/services/db/transactionsRepo.ts index 7cbf8db9..c3d53640 100644 --- a/src/services/db/transactionsRepo.ts +++ b/src/services/db/transactionsRepo.ts @@ -4,6 +4,7 @@ import { Transaction, TransactionDirection, TransactionStatus, + TransactionType, } from '../../models/Transaction' import AppError, {Err} from '../../utils/AppError' import {log} from '../logService' @@ -517,6 +518,43 @@ export const getPendingAmount = function () { } +/** + * Outgoing operations a previous process abandoned before reaching the mint: still + * DRAFT or PREPARED, with no open reservation left (startup already rolled back or + * is holding every reservation). Run at startup only, before any operation can + * start — a live prepare() is DRAFT without a reservation too. + * + * TRANSFERs qualify only once prepare() stamped a quote: a DRAFT with no quote is an + * invoice received over Nostr, legitimately waiting for the user to pay it. EXECUTING + * is deliberately excluded — a mint call may have happened, so it cannot be declared + * unpaid without asking the mint. + */ +export const getAbandonedDraftTransactions = function (): Array<{id: number; data: string}> { + try { + const {rows} = getInstance().execute( + `SELECT id, data FROM transactions + WHERE status IN (?, ?) + AND (type = ? OR (type IN (?, ?) AND quote IS NOT NULL)) + AND id NOT IN (SELECT transactionId FROM reservations)`, + [ + TransactionStatus.DRAFT, + TransactionStatus.PREPARED, + TransactionType.SEND, + TransactionType.TRANSFER, + TransactionType.TRANSFER_ONCHAIN, + ], + ) + const result: Array<{id: number; data: string}> = [] + for (let i = 0; i < (rows?.length ?? 0); i++) { + const row = rows!.item(i) + result.push({id: row.id, data: row.data}) + } + return result + } catch (e: any) { + throw dbError('Could not read abandoned draft transactions', e) + } +} + export const getTransactionById = function (id: number) { try { const query = ` diff --git a/src/services/wallet/operations/inFlightOperations.ts b/src/services/wallet/operations/inFlightOperations.ts index 327456ee..afbbd9f4 100644 --- a/src/services/wallet/operations/inFlightOperations.ts +++ b/src/services/wallet/operations/inFlightOperations.ts @@ -56,6 +56,12 @@ const handleInFlightByMintTask = async (mint: Mint): Promise = continue } + // Held by an interrupted operation: the resolver settles it and removes this + // record. Replaying here would race it over the same inputs. + if (proofsStore.isTransactionInterrupted(tx.id)) { + continue + } + // Already resolved elsewhere (e.g. sync confirmed the proofs SPENT and finalized // the tx): nothing to recover. Drop the lingering in-flight request so it isn't // retried on every sweep. diff --git a/src/services/wallet/operations/interruptedOperations.ts b/src/services/wallet/operations/interruptedOperations.ts new file mode 100644 index 00000000..01343e77 --- /dev/null +++ b/src/services/wallet/operations/interruptedOperations.ts @@ -0,0 +1,240 @@ +/** + * Settles operations a previous process left open mid-way (INTERRUPTIBLE_OPERATION_TYPES). + * + * Startup holds these reservations instead of rolling them back, because the mint may + * already have consumed the locked proofs. Here the mint is asked, and each is settled + * by what it actually did: + * + * inputs UNSPENT → nothing reached the mint: roll back, tx REVERTED + * melt, inputs PENDING/SPENT → recreate the state execute() leaves for an async melt + * (tx PENDING, inputs PENDING, row committed) and let + * TransferOperationApi.refresh settle it on quote state + * swap, inputs SPENT → restore the outputs (NUT-09) from the counter range + * recorded before the request; inputs SPENT, outputs + * UNSPENT, tx REVERTED. Without a range (row opened + * before v36): tx ERROR, the user must run seed recovery + * + * A mint that cannot be reached leaves the reservation held for the next sweep. + */ +import {isAlive} from 'mobx-state-tree' +import {rootStoreInstance} from '../../../models' +import {MintStatus} from '../../../models/Mint' +import {Proof} from '../../../models/Proof' +import {Transaction, TransactionData, TransactionStatus} from '../../../models/Transaction' +import {NetworkError} from '../../../utils/AppError' +import {log} from '../../logService' +import {Database, ReservationRow, ReservationTransactionUpdate} from '../../sqlite' +import {SyncQueue} from '../../syncQueueService' +import {CashuUtils} from '../../cashu/cashuUtils' +import {ProofReservation} from '../proofReservation' +import {WalletTaskResult} from '../types' +import {TransferOperationApi} from './transferOperationApi' + +const {mintsStore, proofsStore, transactionsStore, walletStore} = rootStoreInstance + +export const RESOLVE_INTERRUPTED_TASK = 'resolveInterruptedOperationsTask' + +const MELT_OPERATION_TYPES = new Set(['transfer-melt', 'transfer-melt-after-swap']) + +type Outcome = 'nothing-reached-mint' | 'melt-handed-over' | 'swap-restored' | 'swap-lost' | 'unresolved' + +const resolveInterruptedOperationsTask = async function (): Promise { + const held = proofsStore.interruptedReservations + const rows = Database.getOpenReservations().filter(r => held.has(r.id)) + + // A held row that is gone was settled some other way; nothing left to do. + for (const id of [...held.keys()]) { + if (!rows.some(r => r.id === id)) proofsStore.releaseInterruptedReservation(id) + } + + const errors: string[] = [] + for (const row of rows) { + try { + const outcome = await _resolve(row) + if (outcome !== 'unresolved') proofsStore.releaseInterruptedReservation(row.id) + log.warn('[resolveInterruptedOperationsTask] Interrupted operation', { + reservationId: row.id, + transactionId: row.transactionId, + operationType: row.operationType, + outcome, + }) + } catch (e: any) { + // Held until the next sweep — most often the mint is unreachable. + errors.push(`tId=${row.transactionId}: ${e.message}`) + log.error('[resolveInterruptedOperationsTask] Could not resolve, will retry', { + transactionId: row.transactionId, + operationType: row.operationType, + error: e.message, + }) + } + } + + return { + taskFunction: RESOLVE_INTERRUPTED_TASK, + message: `Resolved ${rows.length - errors.length} of ${rows.length} interrupted operations`, + errors, + } +} + +async function _resolve(row: ReservationRow): Promise { + const mint = (row.mintId ? mintsStore.findById(row.mintId) : undefined) ?? mintsStore.findByUrl(row.mintUrl) + if (!mint) { + log.error('[resolveInterruptedOperationsTask] Mint no longer in wallet', {transactionId: row.transactionId}) + return 'unresolved' + } + + const reservation: ProofReservation = { + id: row.id, + transactionId: row.transactionId, + mintId: row.mintId, + mintUrl: row.mintUrl, + unit: row.unit as ProofReservation['unit'], + operationType: row.operationType, + lockedProofs: row.lockedProofs, + } + const tx = transactionsStore.findById(row.transactionId) + const locked = row.lockedProofs + .map(snap => proofsStore.getBySecret(snap.secret)) + .filter((p): p is Proof => !!p && isAlive(p)) + + const states = await _checkStates(mint, row.unit, locked) + const spent = _nodes(locked, states.SPENT) + const pendingAtMint = _nodes(locked, states.PENDING) + const unspent = _nodes(locked, states.UNSPENT) + + // ── The mint never took the inputs ────────────────────────────────── + if (spent.length === 0 && pendingAtMint.length === 0) { + proofsStore.rollbackReservation(reservation) + // Before the tx write: a stale in-flight record must never outlive this, + // or the in-flight sweep would replay a request the mint never saw. + Database.removeInFlightRequest(row.transactionId) + if (MELT_OPERATION_TYPES.has(row.operationType)) Database.removeMeltRecovery(row.transactionId) + if (tx) { + const {status, data} = _audit( + tx, + TransactionStatus.REVERTED, + 'Interrupted before the mint processed it. Nothing was paid.', + ) + tx.update({status, data}) + } + return 'nothing-reached-mint' + } + + // ── Melt: hand over to the async-melt machinery ───────────────────── + if (MELT_OPERATION_TYPES.has(row.operationType)) { + proofsStore.commitReservation(reservation, { + ...(tx && { + transactionUpdate: _audit( + tx, + TransactionStatus.PENDING, + 'Interrupted while paying. Checking the payment with the mint.', + ), + }), + }) + try { + await TransferOperationApi.refresh(row.transactionId) + } catch (e: any) { + // Now an ordinary PENDING transfer: the pending sweep keeps refreshing it. + log.warn('[resolveInterruptedOperationsTask] Melt refresh failed, left PENDING', {error: e.message}) + } + return 'melt-handed-over' + } + + // ── Swap ──────────────────────────────────────────────────────────── + if (pendingAtMint.length > 0) return 'unresolved' + + const counters = row.counters + if (!counters) { + proofsStore.commitReservation(reservation, { + toSpent: spent, + toUnspent: unspent, + ...(tx && { + transactionUpdate: _audit( + tx, + TransactionStatus.ERROR, + 'Interrupted after the mint executed the swap. Its ecash could not be restored automatically: run seed recovery for this mint.', + ), + }), + }) + Database.removeInFlightRequest(row.transactionId) + return 'swap-lost' + } + + const seed: Uint8Array = await walletStore.getCachedSeed() + const {proofs: restored} = await walletStore.restore(mint.mintUrl, seed, { + indexFrom: counters.start, + indexTo: counters.start + counters.count, + keysetId: counters.keysetId, + unit: row.unit as ProofReservation['unit'], + }) + + // Skip anything already in the wallet: after a counter reuse the range can hold + // outputs of another, completed operation. + const fresh = restored.filter((p: any) => !proofsStore.getBySecret(p.secret)) + const restoredUnspent = fresh.length > 0 ? (await _checkStates(mint, row.unit, fresh as any)).UNSPENT : [] + const restoredAmount = CashuUtils.getProofsAmount(restoredUnspent) + + // Never derive from this range again. + const counter = mint.getProofsCounterByKeysetId(counters.keysetId) + if (counter.counter < counters.next) counter.setProofsCounter(counters.next) + + proofsStore.commitReservation(reservation, { + toSpent: spent, + toUnspent: unspent, + newProofs: restoredUnspent.length > 0 ? [{proofs: restoredUnspent, state: 'UNSPENT', tId: row.transactionId}] : [], + ...(tx && { + transactionUpdate: _audit( + tx, + TransactionStatus.REVERTED, + 'Interrupted after the mint executed the swap. Its ecash was restored. Nothing was paid.', + {restoredAmount}, + ), + }), + }) + Database.removeInFlightRequest(row.transactionId) + return 'swap-restored' +} + +async function _checkStates(mint: {mintUrl: string; setStatus: (s: MintStatus) => void}, unit: string, proofs: Proof[]) { + try { + const states = await walletStore.getProofsStatesFromMint(mint.mintUrl, unit as ProofReservation['unit'], proofs) + mint.setStatus(MintStatus.ONLINE) + return states + } catch (e: any) { + if (e instanceof NetworkError) mint.setStatus(MintStatus.OFFLINE) + throw e + } +} + +/** The locked MST nodes whose secrets the mint reported in `bucket`. */ +function _nodes(locked: Proof[], bucket: Array<{secret: string}>): Proof[] { + const secrets = new Set(bucket.map(p => p.secret)) + return locked.filter(p => secrets.has(p.secret)) +} + +/** Final status plus an audit-trail entry marked `interrupted`, as one tx update. */ +function _audit( + tx: Transaction, + status: TransactionStatus, + message: string, + extra: Record = {}, +): ReservationTransactionUpdate & {status: TransactionStatus; data: string} { + let data: TransactionData[] = [] + try { + data = JSON.parse(tx.data) + } catch {} + data.push({status, interrupted: true, message, ...extra, createdAt: new Date()}) + return {id: tx.id, status, data: JSON.stringify(data)} +} + +const resolveInterruptedQueue = function (): void { + if (proofsStore.interruptedReservations.size === 0) return + // Order against the other sweeps does not matter: sync and the in-flight sweep + // skip anything a held reservation owns. + SyncQueue.addPrioritizedTask(`${RESOLVE_INTERRUPTED_TASK}-${Date.now()}`, resolveInterruptedOperationsTask) +} + +export const InterruptedOperationService = { + resolveInterruptedQueue, + resolveInterruptedOperationsTask, +} diff --git a/src/services/wallet/operations/syncOperations.ts b/src/services/wallet/operations/syncOperations.ts index 63feadc2..cceba3b3 100644 --- a/src/services/wallet/operations/syncOperations.ts +++ b/src/services/wallet/operations/syncOperations.ts @@ -62,7 +62,12 @@ const syncStateWithMintTask = async function ( const revertedTxIds: number[] = [] try { - const aliveProofs = proofsToSync.filter(p => isAlive(p)) + // Proofs held by an interrupted operation are the resolver's alone: the + // outcomes below (SPENT → finalize as success, UNSPENT → ignore) are wrong + // for a swap or melt whose process died mid-flight. + const aliveProofs = proofsToSync.filter( + p => isAlive(p) && !proofsStore.isHeldByInterruptedOperation(p.secret), + ) if (aliveProofs.length === 0) { const message = `No ${proofState === 'PENDING' ? 'pending ' : ''}proofs to sync with mint` diff --git a/src/services/walletService.ts b/src/services/walletService.ts index 482c0c5d..1190af9a 100644 --- a/src/services/walletService.ts +++ b/src/services/walletService.ts @@ -18,6 +18,7 @@ import {SendOperationService} from './wallet/operations/sendOperations' import {ReceiveOperationService} from './wallet/operations/receiveOperations' import {SyncOperationService} from './wallet/operations/syncOperations' import {InFlightOperationService} from './wallet/operations/inFlightOperations' +import {InterruptedOperationService} from './wallet/operations/interruptedOperations' import {PendingOperationService} from './wallet/operations/pendingOperations' import {NostrOperationService} from './wallet/operations/nostrOperations' import {RevertOperationService} from './wallet/operations/revertOperations' @@ -68,6 +69,7 @@ type WalletTaskService = { proofState: ProofState }, ) => Promise + resolveInterruptedQueue: () => void handleInFlightQueue: () => Promise handlePendingQueue: () => Promise handleClaimQueue: () => Promise @@ -184,6 +186,7 @@ export const WalletTask: WalletTaskService = { swapAllQueue: SyncOperationService.swapAllQueue, swapByDenominationQueue: SyncOperationService.swapByDenominationQueue, // In-flight recovery + resolveInterruptedQueue: InterruptedOperationService.resolveInterruptedQueue, handleInFlightQueue: InFlightOperationService.handleInFlightQueue, // Pending orchestration handlePendingQueue: PendingOperationService.handlePendingQueue,