Hnadle correctly interrupted and abandoned transaction statuses and proof states.

This commit is contained in:
minibits-cash
2026-09-28 22:44:31 +02:00
parent 7edaaa2257
commit 6f6e723185
11 changed files with 649 additions and 8 deletions
+214
View File
@@ -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')
})
+64 -3
View File
@@ -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')
})
})
+67 -4
View File
@@ -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<string>(),
interruptedReservations: new Map<string, {transactionId: number; secrets: Set<string>}>(),
}))
// ───────────────────── 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)
},
}))
+6
View File
@@ -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()
+3
View File
@@ -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'})
+2
View File
@@ -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,
+38
View File
@@ -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 = `
@@ -56,6 +56,12 @@ const handleInFlightByMintTask = async (mint: Mint): Promise<WalletTaskResult> =
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.
@@ -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<WalletTaskResult> {
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<Outcome> {
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<string, unknown> = {},
): 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,
}
@@ -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`
+3
View File
@@ -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<SyncStateTaskResult | void>
resolveInterruptedQueue: () => void
handleInFlightQueue: () => Promise<void>
handlePendingQueue: () => Promise<void>
handleClaimQueue: () => Promise<void>
@@ -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,