taler-typescript-core

Wallet core logic and WebUIs for various components
Log | Files | Refs | Submodules | README | LICENSE

commit c88e885096bb5bcb67e76b1bd99f07f71adf166e
parent c95600fe0487e07caae0181eb73fd0fb94939970
Author: Florian Dold <dold@taler.net>
Date:   Fri,  4 Sep 2026 01:02:05 +0200

wallet-core: persist recoup and refresh recovery inputs

Save request commitments before deleting source records, materialize
successful outputs independently, and verify coin history before
correcting conflicts.

Diffstat:
Mpackages/taler-wallet-core/src/recoup.test.ts | 116+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mpackages/taler-wallet-core/src/recoup.ts | 368++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mpackages/taler-wallet-core/src/refresh.test.ts | 84+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/refresh.ts | 286+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Apackages/taler-wallet-core/src/shepherd.test.ts | 48++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/shepherd.ts | 42++++++++++++++++++++++++++++++++++++------
Mpackages/taler-wallet-core/src/transactions.ts | 15++++++++-------
7 files changed, 891 insertions(+), 68 deletions(-)

diff --git a/packages/taler-wallet-core/src/recoup.test.ts b/packages/taler-wallet-core/src/recoup.test.ts @@ -15,8 +15,49 @@ */ import assert from "node:assert"; import { test } from "node:test"; -import { WalletRecoupGroup } from "./db/records.js"; -import { scheduleRecoupRefresh } from "./recoup.js"; +import { DenomKeyType, encodeCrock } from "@gnu-taler/taler-util"; +import { nativeCrypto } from "./crypto/cryptoImplementation.js"; +import { + PlanchetStatus, + WalletPlanchet, + WalletRecoupGroup, +} from "./db/records.js"; +import { + deterministicRecoupWithdrawalGroupId, + reconstructWithdrawalBatchHPlanchets, + scheduleRecoupRefresh, +} from "./recoup.js"; +import { hashWithdrawalPlanchets } from "./withdraw.js"; + +function crock(byte: number, length = 32): string { + return encodeCrock(new Uint8Array(length).fill(byte)); +} + +function rsaPlanchet(index: number): WalletPlanchet { + return { + coinIdx: index, + coinEv: { + cipher: DenomKeyType.Rsa, + rsa_blinded_planchet: crock(index + 10), + }, + denomPubHash: crock(index + 20, 64), + planchetStatus: PlanchetStatus.WithdrawalDone, + } as WalletPlanchet; +} + +function csPlanchet(index: number): WalletPlanchet { + return { + coinIdx: index, + coinEv: { + cipher: DenomKeyType.ClauseSchnorr, + cs_nonce: crock(index + 1), + cs_blinded_c0: crock(index + 2), + cs_blinded_c1: crock(index + 3), + }, + denomPubHash: crock(index + 30, 64), + planchetStatus: PlanchetStatus.WithdrawalDone, + } as WalletPlanchet; +} test("recoup-refresh schedules and aggregates value on the old coin", () => { const group = { scheduleRefreshCoins: [] } as unknown as WalletRecoupGroup; @@ -30,3 +71,74 @@ test("recoup-refresh schedules and aggregates value on the old coin", () => { { coinPub: "another-old-coin", amount: "TESTKUDOS:1" }, ]); }); + +test("recoup reconstructs RSA and Clause-Schnorr withdrawal commitments", () => { + const planchets = [ + rsaPlanchet(0), + rsaPlanchet(1), + csPlanchet(2), + rsaPlanchet(3), + ]; + assert.strictEqual( + reconstructWithdrawalBatchHPlanchets(planchets, 1, 4), + hashWithdrawalPlanchets( + planchets.slice(0, 2).map((x) => x.coinEv), + planchets.slice(0, 2).map((x) => x.denomPubHash), + ), + ); + assert.strictEqual( + reconstructWithdrawalBatchHPlanchets(planchets, 2, 4), + hashWithdrawalPlanchets([planchets[2].coinEv], [planchets[2].denomPubHash]), + ); +}); + +test("Clause-Schnorr recoup requests retain all blinding inputs", async () => { + const exchangeWithdrawValues = { + cipher: DenomKeyType.ClauseSchnorr as const, + r_pub_0: crock(4), + r_pub_1: crock(5), + }; + const common = { + coinPub: crock(6), + coinPriv: crock(7), + blindingKey: crock(8), + denomPub: { + cipher: DenomKeyType.ClauseSchnorr as const, + age_mask: 0, + cs_public_key: crock(9), + }, + denomPubHash: crock(10, 64), + denomSig: { + cipher: DenomKeyType.ClauseSchnorr, + cs_signature: "sig", + } as never, + exchangeWithdrawValues, + ageCommitmentHash: crock(11), + nonce: crock(12), + }; + const withdrawRequest = await nativeCrypto.createRecoupRequest({ + ...common, + hPlanchets: crock(13, 64), + }); + assert.deepStrictEqual(withdrawRequest.ewv, exchangeWithdrawValues); + assert.strictEqual(withdrawRequest.h_planchets, crock(13, 64)); + assert.strictEqual(withdrawRequest.h_age_commitment, crock(11)); + assert.strictEqual(withdrawRequest.nonce, crock(12)); + + const refreshRequest = await nativeCrypto.createRecoupRefreshRequest(common); + assert.deepStrictEqual(refreshRequest.ewv, exchangeWithdrawValues); + assert.strictEqual(refreshRequest.h_age_commitment, crock(11)); + assert.strictEqual(refreshRequest.nonce, crock(12)); +}); + +test("recoup withdrawal identities are deterministic per reserve", () => { + const first = deterministicRecoupWithdrawalGroupId("group", "reserve-a"); + assert.strictEqual( + deterministicRecoupWithdrawalGroupId("group", "reserve-a"), + first, + ); + assert.notStrictEqual( + deterministicRecoupWithdrawalGroupId("group", "reserve-b"), + first, + ); +}); diff --git a/packages/taler-wallet-core/src/recoup.ts b/packages/taler-wallet-core/src/recoup.ts @@ -25,20 +25,29 @@ * Imports. */ import { + AgeRestriction, + AmountJson, AmountLike, + AmountString, Amounts, CoinStatus, + DenomKeyType, Logger, RefreshReason, TalerPreciseTimestamp, Transaction, + TransactionAction, TransactionIdStr, + TransactionMajorState, + TransactionState, TransactionType, WalletNotification, checkDbInvariant, encodeCrock, getRandomBytes, + hash, j2s, + stringToBytes, } from "@gnu-taler/taler-util"; import { PendingTaskType, @@ -51,6 +60,7 @@ import { RecoupOperationStatus, CoinSourceType, WithdrawalGroupStatus, + timestampPreciseFromDb, timestampPreciseToDb, WalletCoin, WalletRefreshCoinSource, @@ -58,8 +68,13 @@ import { WalletRecoupGroup, WithdrawalRecordType, } from "./db/records.js"; -import { requireExchangeCoinUseConfirmedOrThrow } from "./exchanges.js"; -import { createRefreshGroup } from "./refresh.js"; +import { + getExchangeDetailsInTx, + getLegacyScopesForTransaction, + getScopeForAllExchanges, + requireExchangeCoinUseConfirmedOrThrow, +} from "./exchanges.js"; +import { createRefreshGroup, recoverRefreshCoinNonce } from "./refresh.js"; import { constructTransactionIdentifier, makeTransactionActionUnsupportedError, @@ -69,7 +84,12 @@ import { getDenomInfo, walletExchangeClient, } from "./wallet.js"; -import { internalCreateWithdrawalGroup } from "./withdraw.js"; +import { + internalCreateWithdrawalGroup, + reconstructWithdrawalBatchHPlanchets, +} from "./withdraw.js"; + +export { reconstructWithdrawalBatchHPlanchets } from "./withdraw.js"; import { WalletDbTransaction } from "./db/transaction.js"; const logger = new Logger("operations/recoup.ts"); @@ -110,6 +130,9 @@ async function recoupRewardCoin( if (!recoupGroup) { return; } + if (recoupGroup.operationStatus !== RecoupOperationStatus.Pending) { + return; + } if (recoupGroup.recoupFinishedPerCoin[coinIdx]) { return; } @@ -117,6 +140,25 @@ async function recoupRewardCoin( }); } +async function markRecoupCoinPermanentlyFailed( + wex: WalletExecutionContext, + recoupGroupId: string, + coinIdx: number, + coinPub: string, +): Promise<void> { + await wex.runWalletDbTx(async (tx) => { + const group = await tx.getRecoupGroup(recoupGroupId); + if (group?.operationStatus !== RecoupOperationStatus.Pending) { + return; + } + group.failedCoinPubs ??= []; + if (!group.failedCoinPubs.includes(coinPub)) { + group.failedCoinPubs.push(coinPub); + } + await putGroupAsFinished(wex, tx, group, coinIdx); + }); +} + async function recoupRefreshCoin( wex: WalletExecutionContext, recoupGroupId: string, @@ -124,6 +166,25 @@ async function recoupRefreshCoin( coin: WalletCoin, cs: WalletRefreshCoinSource, ): Promise<void> { + let nonce = cs.nonce; + if ( + nonce == null && + coin.exchangeWithdrawValues.cipher === DenomKeyType.ClauseSchnorr + ) { + nonce = await recoverRefreshCoinNonce(wex, coin); + checkDbInvariant( + nonce != null, + `missing refresh nonce for recoup coin ${coin.coinPub}`, + ); + const recoveredNonce = nonce; + await wex.runWalletDbTx(async (tx) => { + const current = await tx.getCoin(coin.coinPub); + if (current?.coinSource.type === CoinSourceType.Refresh) { + current.coinSource.nonce = recoveredNonce; + await tx.upsertCoin(current); + } + }); + } const d = await wex.runWalletDbTx(async (tx) => { const denomInfo = await getDenomInfo(wex, tx, coin); if (!denomInfo) { @@ -145,6 +206,11 @@ async function recoupRefreshCoin( denomPub: d.denomInfo.denomPub, denomPubHash: coin.denomPubHash, denomSig: coin.denomSig, + exchangeWithdrawValues: coin.exchangeWithdrawValues, + ageCommitmentHash: coin.ageCommitmentProof + ? AgeRestriction.hashCommitment(coin.ageCommitmentProof.commitment) + : undefined, + nonce, }); logger.trace(`making recoup request for ${coin.coinPub}`); @@ -154,6 +220,18 @@ async function recoupRefreshCoin( coin.coinPub, recoupRequest, ); + if (recoupResp.case !== "ok") { + logger.warn( + `exchange permanently rejected recoup-refresh for ${coin.coinPub}: ${j2s(recoupResp.detail)}`, + ); + await markRecoupCoinPermanentlyFailed( + wex, + recoupGroupId, + coinIdx, + coin.coinPub, + ); + return; + } const recoupConfirmation = recoupResp.body; if (recoupConfirmation.old_coin_pub != cs.oldCoinPub) { @@ -165,6 +243,9 @@ async function recoupRefreshCoin( if (!recoupGroup) { return; } + if (recoupGroup.operationStatus !== RecoupOperationStatus.Pending) { + return; + } if (recoupGroup.recoupFinishedPerCoin[coinIdx]) { return; } @@ -228,11 +309,45 @@ export async function recoupWithdrawCoin( cs: WalletWithdrawCoinSource, ): Promise<void> { const reservePub = cs.reservePub; - const denomInfo = await wex.runWalletDbTx(async (tx) => { + const recoupInputs = await wex.runWalletDbTx(async (tx) => { const denomInfo = await getDenomInfo(wex, tx, coin); - return denomInfo; + if (!denomInfo) { + return; + } + const planchet = await tx.getPlanchetByGroupAndIndex( + cs.withdrawalGroupId, + cs.coinIndex, + ); + let hPlanchets = cs.hPlanchets; + if (!hPlanchets) { + const withdrawalGroup = await tx.getWithdrawalGroup(cs.withdrawalGroupId); + if (withdrawalGroup?.denomsSel) { + const numCoins = withdrawalGroup.denomsSel.selectedDenoms.reduce( + (n, x) => n + x.count, + 0, + ); + hPlanchets = reconstructWithdrawalBatchHPlanchets( + await tx.getPlanchetsByGroup(cs.withdrawalGroupId), + cs.coinIndex, + numCoins, + ); + } + } + checkDbInvariant( + !!hPlanchets, + `missing withdrawal commitment for recoup coin ${coin.coinPub}`, + ); + return { + denomInfo, + hPlanchets, + nonce: + cs.nonce ?? + (planchet?.coinEv.cipher === DenomKeyType.ClauseSchnorr + ? planchet.coinEv.cs_nonce + : undefined), + }; }); - if (!denomInfo) { + if (!recoupInputs) { // FIXME: We should at least emit some pending operation / warning for this? return; } @@ -243,9 +358,16 @@ export async function recoupWithdrawCoin( blindingKey: coin.blindingKey, coinPriv: coin.coinPriv, coinPub: coin.coinPub, - denomPub: denomInfo.denomPub, + denomPub: recoupInputs.denomInfo.denomPub, denomPubHash: coin.denomPubHash, denomSig: coin.denomSig, + exchangeWithdrawValues: coin.exchangeWithdrawValues, + hPlanchets: recoupInputs.hPlanchets, + withdrawCommitmentHash: cs.withdrawCommitmentHash, + ageCommitmentHash: coin.ageCommitmentProof + ? AgeRestriction.hashCommitment(coin.ageCommitmentProof.commitment) + : undefined, + nonce: recoupInputs.nonce, }); logger.trace(`requesting recoup for coin ${coin.coinPub}`); const exchangeClient = walletExchangeClient(coin.exchangeBaseUrl, wex); @@ -254,6 +376,18 @@ export async function recoupWithdrawCoin( coin.coinPub, recoupRequest, ); + if (recoupResp.case !== "ok") { + logger.warn( + `exchange permanently rejected recoup for ${coin.coinPub}: ${j2s(recoupResp.detail)}`, + ); + await markRecoupCoinPermanentlyFailed( + wex, + recoupGroupId, + coinIdx, + coin.coinPub, + ); + return; + } const recoupConfirmation = recoupResp.body; logger.trace(`got recoup confirmation ${j2s(recoupConfirmation)}`); @@ -268,6 +402,9 @@ export async function recoupWithdrawCoin( if (!recoupGroup) { return; } + if (recoupGroup.operationStatus !== RecoupOperationStatus.Pending) { + return; + } if (recoupGroup.recoupFinishedPerCoin[coinIdx]) { return; } @@ -276,6 +413,10 @@ export async function recoupWithdrawCoin( return; } updatedCoin.status = CoinStatus.Dormant; + recoupGroup.successfulCoinPubs ??= []; + if (!recoupGroup.successfulCoinPubs.includes(updatedCoin.coinPub)) { + recoupGroup.successfulCoinPubs.push(updatedCoin.coinPub); + } await tx.upsertCoin(updatedCoin); await putGroupAsFinished(wex, tx, recoupGroup, coinIdx); }); @@ -301,6 +442,10 @@ export async function processRecoupGroup( if (!recoupGroup) { return TaskRunResult.finished(); } + + if (recoupGroup.operationStatus !== RecoupOperationStatus.Pending) { + return TaskRunResult.finished(); + } if (recoupGroup.timestampFinished) { logger.trace("recoup group finished"); return TaskRunResult.finished(); @@ -331,47 +476,46 @@ export async function processRecoupGroup( logger.info("all recoups of recoup group are finished"); - const reserveSet = new Set<string>(); + const reserveAmounts = new Map<string, AmountJson>(); const reservePrivMap: Record<string, string> = {}; await wex.runWalletDbTx(async (tx) => { - const coins = await tx.getCoinsByPubs(recoupGroup.coinPubs); + const coins = await tx.getCoinsByPubs(recoupGroup.successfulCoinPubs ?? []); const coinsByPub = new Map(coins.map((coin) => [coin.coinPub, coin])); const reservePubs: string[] = []; - for (const coinPub of recoupGroup.coinPubs) { + for (const coinPub of recoupGroup.successfulCoinPubs ?? []) { const coin = coinsByPub.get(coinPub); if (!coin) { throw Error(`Coin ${coinPub} not found, can't request recoup`); } if (coin.coinSource.type === CoinSourceType.Withdraw) { reservePubs.push(coin.coinSource.reservePub); + const denom = await getDenomInfo(wex, tx, coin); + checkDbInvariant(!!denom, `denomination for ${coinPub} is missing`); + const previous = reserveAmounts.get(coin.coinSource.reservePub); + reserveAmounts.set( + coin.coinSource.reservePub, + previous + ? Amounts.add(previous, denom.value).amount + : Amounts.jsonifyAmount(denom.value), + ); } } const uniqueReservePubs = [...new Set(reservePubs)]; const reserves = await tx.getReservesByPubs(uniqueReservePubs); for (const reserve of reserves) { - reserveSet.add(reserve.reservePub); reservePrivMap[reserve.reservePub] = reserve.reservePriv; } }); - for (const reservePub of reserveSet) { - logger.info(`querying reserve status for recoup of ${reservePub}`); - - const exchangeClient = walletExchangeClient( - recoupGroup.exchangeBaseUrl, - wex, - ); - const reserveResp = await exchangeClient.getReserveStatus(reservePub); - if (reserveResp.case !== "ok") { - throw Error( - `could not query reserve status for recoup (${reserveResp.case})`, - ); - } - const result = reserveResp.body; + for (const [reservePub, recoveredAmount] of reserveAmounts) { await internalCreateWithdrawalGroup(wex, { - amount: Amounts.parseOrThrow(result.balance), + amount: recoveredAmount, exchangeBaseUrl: recoupGroup.exchangeBaseUrl, reserveStatus: WithdrawalGroupStatus.PendingQueryingStatus, + forcedWithdrawalGroupId: deterministicRecoupWithdrawalGroupId( + recoupGroupId, + reservePub, + ), reserveKeyPair: { pub: reservePub, priv: reservePrivMap[reservePub], @@ -389,8 +533,13 @@ export async function processRecoupGroup( if (!rg2) { return; } + if (rg2.operationStatus !== RecoupOperationStatus.Pending) { + return; + } rg2.timestampFinished = timestampPreciseToDb(TalerPreciseTimestamp.now()); - rg2.operationStatus = RecoupOperationStatus.Finished; + rg2.operationStatus = rg2.failedCoinPubs?.length + ? RecoupOperationStatus.Failed + : RecoupOperationStatus.Finished; if (rg2.scheduleRefreshCoins.length > 0) { await createRefreshGroup( wex, @@ -410,6 +559,17 @@ export async function processRecoupGroup( return TaskRunResult.finished(); } +export function deterministicRecoupWithdrawalGroupId( + recoupGroupId: string, + reservePub: string, +): string { + return encodeCrock( + hash( + stringToBytes(`recoup-withdraw\0${recoupGroupId}\0${reservePub}`), + ).slice(0, 32), + ); +} + export class RecoupTransactionContext implements TransactionContext { public transactionId: TransactionIdStr; public taskId: TaskIdStr; @@ -428,30 +588,55 @@ export class RecoupTransactionContext implements TransactionContext { }); } - /** - * A recoup group is not a materialized transaction yet: lookupFullTransaction - * below, getContextForTransaction, rematerializeTransactions and the - * shepherd's retry notifications all leave recoup out. Writing a meta record - * anyway would put an entry in the transaction list that nothing can render, - * and getTransactions would then fail for the *whole* wallet. - * - * So only ever remove the record here. That also repairs wallets that - * already stored one before this was fixed. - */ async updateTransactionMeta(tx: WalletDbTransaction): Promise<void> { - await tx.deleteTransactionMeta(this.transactionId); + const rec = await tx.getRecoupGroup(this.recoupGroupId); + if (!rec) { + await tx.deleteTransactionMeta(this.transactionId); + return; + } + const summary = await getRecoupSummary(this.wex, tx, rec); + await tx.upsertTransactionMeta({ + transactionId: this.transactionId, + status: rec.operationStatus, + timestamp: rec.timestampStarted, + currency: summary.currency, + exchanges: [rec.exchangeBaseUrl], + legacyScopes: await getLegacyScopesForTransaction(tx, { + currency: summary.currency, + exchanges: [rec.exchangeBaseUrl], + coinPubs: rec.coinPubs, + }), + }); } userAbortTransaction(): Promise<void> { throw makeTransactionActionUnsupportedError(this.transactionId, "abort"); } - userSuspendTransaction(): Promise<void> { - throw makeTransactionActionUnsupportedError(this.transactionId, "suspend"); + async userSuspendTransaction(): Promise<void> { + await this.wex.runWalletDbTx(async (tx) => { + const rec = await tx.getRecoupGroup(this.recoupGroupId); + if (rec?.operationStatus !== RecoupOperationStatus.Pending) { + return; + } + rec.operationStatus = RecoupOperationStatus.Suspended; + await tx.upsertRecoupGroup(rec); + await this.updateTransactionMeta(tx); + }); + this.wex.taskScheduler.stopShepherdTask(this.taskId); } - userResumeTransaction(): Promise<void> { - throw makeTransactionActionUnsupportedError(this.transactionId, "resume"); + async userResumeTransaction(): Promise<void> { + await this.wex.runWalletDbTx(async (tx) => { + const rec = await tx.getRecoupGroup(this.recoupGroupId); + if (rec?.operationStatus !== RecoupOperationStatus.Suspended) { + return; + } + rec.operationStatus = RecoupOperationStatus.Pending; + await tx.upsertRecoupGroup(rec); + await this.updateTransactionMeta(tx); + }); + await this.wex.taskScheduler.resetTask(this.taskId); } userFailTransaction(): Promise<void> { @@ -480,14 +665,85 @@ export class RecoupTransactionContext implements TransactionContext { return { notifs }; } - lookupFullTransaction( + async lookupFullTransaction( tx: WalletDbTransaction, ): Promise<Transaction | undefined> { - throw makeTransactionActionUnsupportedError( - this.transactionId, - "lookup", - "recoup transactions are not materialized and cannot be looked up", - ); + const rec = await tx.getRecoupGroup(this.recoupGroupId); + if (!rec) { + return undefined; + } + const summary = await getRecoupSummary(this.wex, tx, rec); + const ort = await tx.getOperationRetry(this.taskId); + return { + type: TransactionType.Recoup, + transactionId: this.transactionId, + timestamp: timestampPreciseFromDb(rec.timestampStarted), + txState: computeRecoupTransactionState(rec), + stId: rec.operationStatus, + txActions: computeRecoupTransactionActions(rec), + scopes: await getScopeForAllExchanges(tx, [rec.exchangeBaseUrl]), + amountRaw: summary.amount, + // Recoup moves value from a revoked coin to its source reserve/coin; + // the resulting withdrawal/refresh reports the spendable balance gain. + amountEffective: Amounts.stringify( + Amounts.zeroOfCurrency(summary.currency), + ), + ...(ort?.lastError ? { error: ort.lastError } : {}), + }; + } +} + +async function getRecoupSummary( + wex: WalletExecutionContext, + tx: WalletDbTransaction, + rec: WalletRecoupGroup, +): Promise<{ currency: string; amount: AmountString }> { + const exchangeDetails = await getExchangeDetailsInTx(tx, rec.exchangeBaseUrl); + checkDbInvariant( + !!exchangeDetails, + `missing exchange details for recoup ${rec.recoupGroupId}`, + ); + const amounts: AmountString[] = []; + for (const coin of await tx.getCoinsByPubs(rec.coinPubs)) { + const denom = await getDenomInfo(wex, tx, coin); + if (denom) { + amounts.push(denom.value); + } + } + return { + currency: exchangeDetails.currency, + amount: Amounts.stringify( + Amounts.sumOrZero(exchangeDetails.currency, amounts).amount, + ), + }; +} + +export function computeRecoupTransactionState( + rec: WalletRecoupGroup, +): TransactionState { + switch (rec.operationStatus) { + case RecoupOperationStatus.Pending: + return { major: TransactionMajorState.Pending, working: true }; + case RecoupOperationStatus.Suspended: + return { major: TransactionMajorState.Suspended }; + case RecoupOperationStatus.Finished: + return { major: TransactionMajorState.Done }; + case RecoupOperationStatus.Failed: + return { major: TransactionMajorState.Failed }; + } +} + +export function computeRecoupTransactionActions( + rec: WalletRecoupGroup, +): TransactionAction[] { + switch (rec.operationStatus) { + case RecoupOperationStatus.Pending: + return [TransactionAction.Retry, TransactionAction.Suspend]; + case RecoupOperationStatus.Suspended: + return [TransactionAction.Resume]; + case RecoupOperationStatus.Finished: + case RecoupOperationStatus.Failed: + return [TransactionAction.Delete]; } } @@ -507,6 +763,8 @@ export async function createRecoupGroup( timestampFinished: undefined, timestampStarted: timestampPreciseToDb(TalerPreciseTimestamp.now()), recoupFinishedPerCoin: coinPubs.map(() => false), + failedCoinPubs: [], + successfulCoinPubs: [], scheduleRefreshCoins: [], operationStatus: RecoupOperationStatus.Pending, }; @@ -520,6 +778,17 @@ export async function createRecoupGroup( await putGroupAsFinished(wex, tx, recoupGroup, coinIdx); continue; } + // Only coins with locally spendable value can contribute anything to a + // recoup. In particular, a fully spent (Dormant) coin can legitimately + // be rejected by the exchange and must not block recovery of its fresh + // siblings forever. + if ( + coin.status !== CoinStatus.Fresh && + coin.status !== CoinStatus.FreshSuspended + ) { + await putGroupAsFinished(wex, tx, recoupGroup, coinIdx); + continue; + } await tx.upsertCoin(coin); } @@ -547,6 +816,9 @@ async function processRecoupForCoin( if (recoupGroup.timestampFinished) { return; } + if (recoupGroup.operationStatus !== RecoupOperationStatus.Pending) { + return; + } if (recoupGroup.recoupFinishedPerCoin[coinIdx]) { return; } diff --git a/packages/taler-wallet-core/src/refresh.test.ts b/packages/taler-wallet-core/src/refresh.test.ts @@ -29,6 +29,7 @@ import { getTotalRefreshCostInternal, RefreshTransactionContext, requireValidNorevealIndex, + validateAndRecomputeCoinHistoryBalance, } from "./refresh.js"; import { RefreshCoinStatus, @@ -134,6 +135,9 @@ test("deleting a live refresh removes its pending output count", async () => { async getRefreshGroup() { return storedRefreshGroup; }, + async getCoinsBySourceTransaction() { + return []; + }, async getRefreshSessionsByGroup() { return [session]; }, @@ -178,3 +182,83 @@ test("an impossible refresh costs the full remaining amount", () => { amountLeft, ); }); + +test("coin history balance is recomputed from typed operations", () => { + const history = { + h_denom_pub: "denom", + balance: "TESTKUDOS:8", + history: [ + { + type: "DEPOSIT", + history_offset: 0, + amount: "TESTKUDOS:3", + h_denom_pub: "denom", + }, + { + type: "REFUND", + history_offset: 1, + amount: "TESTKUDOS:1", + }, + ], + } as never; + assert.strictEqual( + validateAndRecomputeCoinHistoryBalance("denom", "TESTKUDOS:10", history), + "TESTKUDOS:8", + ); +}); + +test("coin history rejects unbound or non-conserving aggregates", () => { + const response = { + h_denom_pub: "denom", + balance: "TESTKUDOS:0", + history: [ + { + type: "MELT", + history_offset: 0, + amount: "TESTKUDOS:3", + old_denom_pub_h: "denom", + }, + ], + } as never; + assert.throws( + () => + validateAndRecomputeCoinHistoryBalance("denom", "TESTKUDOS:10", response), + /does not match recomputed/, + ); + assert.throws( + () => + validateAndRecomputeCoinHistoryBalance( + "another-denom", + "TESTKUDOS:10", + response, + ), + /denomination does not match/, + ); +}); + +test("coin history accounts for the waived fee of refunded purse deposits", () => { + const response = { + h_denom_pub: "denom", + balance: "TESTKUDOS:10.08", + history: [ + { + type: "PURSE-DEPOSIT", + history_offset: 0, + amount: "TESTKUDOS:3", + deposit_fee: "TESTKUDOS:0.1", + refunded: true, + h_denom_pub: "denom", + }, + { + type: "PURSE-REFUND", + history_offset: 1, + amount: "TESTKUDOS:2.98", + refund_fee: "TESTKUDOS:0.02", + }, + ], + } as never; + assert.strictEqual( + validateAndRecomputeCoinHistoryBalance("denom", "TESTKUDOS:10", response), + "TESTKUDOS:10.08", + ); +}); diff --git a/packages/taler-wallet-core/src/refresh.ts b/packages/taler-wallet-core/src/refresh.ts @@ -27,12 +27,14 @@ import { AgeRestriction, AmountJson, AmountLike, + AmountString, Amounts, assertUnreachable, BlindedDenominationSignature, checkDbInvariant, checkLogicInvariant, CoinRefreshRequest, + CoinHistoryResponse, CoinStatus, DenominationInfo, DenomKeyType, @@ -89,7 +91,10 @@ import { TaskRunResultType, TransactionContext, } from "./common.js"; -import { requireValidDirectExchangeRefundConfirmation } from "./exchange-signatures.js"; +import { + requireValidDirectExchangeRefundConfirmation, + requireValidExchangeCoinHistory, +} from "./exchange-signatures.js"; import { RefreshNewDenomInfo } from "./crypto/cryptoTypes.js"; import { CryptoApiStoppedError } from "./crypto/workers/crypto-dispatcher.js"; import { @@ -254,6 +259,33 @@ export class RefreshTransactionContext implements TransactionContext { } async userDeleteTransaction(): Promise<void> { + const sourceCoins = await this.wex.runWalletDbTx((tx) => + tx.getCoinsBySourceTransaction(this.transactionId), + ); + for (const coin of sourceCoins) { + if ( + coin.coinSource.type !== CoinSourceType.Refresh || + coin.exchangeWithdrawValues.cipher !== DenomKeyType.ClauseSchnorr || + coin.coinSource.nonce != null + ) { + continue; + } + const nonce = await recoverRefreshCoinNonce(this.wex, coin); + if (!nonce) { + throw makeTransactionActionUnsupportedError( + this.transactionId, + "delete", + "refresh source material is still needed for recoup", + ); + } + await this.wex.runWalletDbTx(async (tx) => { + const current = await tx.getCoin(coin.coinPub); + if (current?.coinSource.type === CoinSourceType.Refresh) { + current.coinSource.nonce = nonce; + await tx.upsertCoin(current); + } + }); + } await this.wex.runWalletDbTx(async (tx) => { return this.deleteTransactionInTx(tx); }); @@ -267,6 +299,21 @@ export class RefreshTransactionContext implements TransactionContext { ); return; } + for (const coin of await tx.getCoinsBySourceTransaction( + this.transactionId, + )) { + if ( + coin.coinSource.type === CoinSourceType.Refresh && + coin.exchangeWithdrawValues.cipher === DenomKeyType.ClauseSchnorr && + coin.coinSource.nonce == null + ) { + throw makeTransactionActionUnsupportedError( + this.transactionId, + "delete", + "refresh source material is still needed for recoup", + ); + } + } const sessions = await tx.getRefreshSessionsByGroup(rg.refreshGroupId); for (const s of sessions) { const coinStatus = rg.statusPerCoin[s.coinIndex]; @@ -794,6 +841,80 @@ async function deriveRefreshSession( } /** + * Recover the Clause-Schnorr nonce for a coin created by an older wallet that + * did not copy it into the coin source record. The refresh seed and selected + * cut-and-choose batch make the planchets deterministic. + */ +export async function recoverRefreshCoinNonce( + wex: WalletExecutionContext, + coin: WalletCoin, +): Promise<string | undefined> { + if (coin.coinSource.type !== CoinSourceType.Refresh) { + return undefined; + } + if (coin.coinSource.nonce != null) { + return coin.coinSource.nonce; + } + if (coin.exchangeWithdrawValues.cipher !== DenomKeyType.ClauseSchnorr) { + return undefined; + } + const src = coin.coinSource; + const inputs = await wex.runWalletDbTx(async (tx) => { + const group = await tx.getRefreshGroup(src.refreshGroupId); + if (!group) { + return undefined; + } + const coinIndex = group.oldCoinPubs.indexOf(src.oldCoinPub); + if (coinIndex < 0) { + return undefined; + } + const session = await tx.getRefreshSession(src.refreshGroupId, coinIndex); + const oldCoin = await tx.getCoin(src.oldCoinPub); + if (!session || !oldCoin || session.norevealIndex == null) { + return undefined; + } + const oldDenom = await getDenomInfo(wex, tx, oldCoin); + if (!oldDenom) { + return undefined; + } + const newCoinDenoms: RefreshNewDenomInfo[] = []; + for (const nd of session.newDenoms) { + const denom = await getDenomInfo(wex, tx, { + exchangeMasterPub: oldCoin.exchangeMasterPub, + denomPubHash: nd.denomPubHash, + }); + if (!denom) { + return undefined; + } + newCoinDenoms.push({ + count: nd.count, + denomPub: denom.denomPub, + denomPubHash: denom.denomPubHash, + feeWithdraw: denom.feeWithdraw, + value: Amounts.stringify(denom.value), + }); + } + return { session, oldCoin, oldDenom, newCoinDenoms }; + }); + if (!inputs) { + return undefined; + } + const derived = await deriveRefreshSession( + wex, + inputs.session, + inputs.oldCoin, + inputs.oldDenom, + inputs.newCoinDenoms, + ); + const planchet = derived.planchets[inputs.session.norevealIndex!]?.find( + (x) => x.coinPub === coin.coinPub, + ); + return planchet?.coinEv.cipher === DenomKeyType.ClauseSchnorr + ? planchet.coinEv.cs_nonce + : undefined; +} + +/** * Run the melt step of a refresh session. * * If the melt step succeeds or fails permanently, @@ -1229,8 +1350,22 @@ async function handleRefreshMeltConflict( const historyJson = historyResp.body; logger.info(`coin history: ${j2s(historyJson)}`); - - // FIXME: If response seems wrong, report to auditor (in the future!); + const denomination = await ctx.wex.runWalletDbTx(async (tx) => { + const denom = await getDenomInfo(ctx.wex, tx, oldCoin); + checkDbInvariant(!!denom, "missing denomination for history recovery"); + return denom; + }); + await requireValidExchangeCoinHistory(ctx.wex, { + exchangeBaseUrl: oldCoin.exchangeBaseUrl, + coinPub: oldCoin.coinPub, + denomination, + response: historyJson, + }); + const correctedBalance = validateAndRecomputeCoinHistoryBalance( + oldCoin.denomPubHash, + denomination.value, + historyJson, + ); await ctx.wex.runWalletDbTx(async (tx) => { const [rg, h] = await ctx.getRecordHandle(tx); @@ -1243,7 +1378,7 @@ async function handleRefreshMeltConflict( if (rg.statusPerCoin[coinIndex] !== RefreshCoinStatus.Pending) { return; } - if (Amounts.isZero(historyJson.balance)) { + if (Amounts.isZero(correctedBalance)) { const refreshSession = await tx.getRefreshSession( ctx.refreshGroupId, coinIndex, @@ -1273,7 +1408,7 @@ async function handleRefreshMeltConflict( } } else { // Try again with new denoms! - rg.inputPerCoin[coinIndex] = historyJson.balance; + rg.inputPerCoin[coinIndex] = correctedBalance; const refreshSession = await tx.getRefreshSession( ctx.refreshGroupId, coinIndex, @@ -1283,6 +1418,19 @@ async function handleRefreshMeltConflict( } await destroyRefreshSession(tx, rg, refreshSession); await tx.deleteRefreshSession(ctx.refreshGroupId, coinIndex); + const correctedOutput = await calculateRefreshOutput( + ctx.wex, + tx, + rg.currency, + rg.oldCoinPubs.map((coinPub, i) => ({ + coinPub, + amount: rg.inputPerCoin[i], + })), + ); + rg.expectedOutputPerCoin = correctedOutput.outputPerCoin.map((x) => + Amounts.stringify(x), + ); + rg.infoPerExchange = correctedOutput.perExchangeInfo; await initRefreshSession(ctx.wex, tx, rg, coinIndex); // The new session was computed from the corrected input amount, // so that amount must be stored as well. @@ -1291,6 +1439,130 @@ async function handleRefreshMeltConflict( }); } +/** + * Recompute a coin balance from the typed history instead of trusting the + * unauthenticated top-level aggregate. + */ +export function validateAndRecomputeCoinHistoryBalance( + denomPubHash: string, + denomValue: AmountString, + response: CoinHistoryResponse, +): AmountString { + if (response.h_denom_pub !== denomPubHash) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "coin history denomination does not match the requested coin", + ); + } + let balance = Amounts.parseOrThrow(denomValue); + const offsets = new Set<number>(); + for (const item of response.history) { + if ( + !Number.isInteger(item.history_offset) || + item.history_offset < 0 || + offsets.has(item.history_offset) + ) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "coin history contains an invalid or duplicate offset", + ); + } + offsets.add(item.history_offset); + let amount: AmountString; + let addsValue = false; + switch (item.type) { + case "DEPOSIT": + case "RECOUP-WITHDRAW": + case "RECOUP-REFRESH": + case "PURSE-DEPOSIT": + amount = item.amount; + break; + case "MELT": + amount = item.amount; + break; + case "RESERVE-OPEN-DEPOSIT": + amount = item.coin_contribution; + break; + case "REFUND": + case "RECOUP-REFRESH-RECEIVER": + case "PURSE-REFUND": + amount = item.amount; + addsValue = true; + break; + default: + assertUnreachable(item); + } + if (item.type === "PURSE-DEPOSIT" && item.refunded) { + if ( + !Amounts.isSameCurrency( + Amounts.currencyOf(item.deposit_fee), + Amounts.currencyOf(denomValue), + ) + ) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "refunded purse deposit has a fee in another currency", + ); + } + const waivedFeeBalance = Amounts.add(balance, item.deposit_fee); + if (waivedFeeBalance.saturated) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "refunded purse deposit fee overflows the coin balance", + ); + } + balance = waivedFeeBalance.amount; + } + if ( + "h_denom_pub" in item && + item.h_denom_pub !== undefined && + item.h_denom_pub !== denomPubHash + ) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + `coin history ${item.type} entry has a different denomination`, + ); + } + if ( + !Amounts.isSameCurrency( + Amounts.currencyOf(amount), + Amounts.currencyOf(denomValue), + ) + ) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + `coin history ${item.type} entry has a different currency`, + ); + } + const result = addsValue + ? Amounts.add(balance, amount) + : Amounts.sub(balance, amount); + if (result.saturated) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + `coin history ${item.type} entry violates value conservation`, + ); + } + balance = result.amount; + } + const recomputed = Amounts.stringify(balance); + if (Amounts.cmp(recomputed, response.balance) !== 0) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + `coin history balance ${response.balance} does not match recomputed ${recomputed}`, + ); + } + return recomputed; +} + async function handleRefreshMeltNotFound( ctx: RefreshTransactionContext, coinIndex: number, @@ -1518,6 +1790,10 @@ async function refreshReveal( type: CoinSourceType.Refresh, refreshGroupId, oldCoinPub: refreshGroup.oldCoinPubs[coinIndex], + nonce: + pc.coinEv.cipher === DenomKeyType.ClauseSchnorr + ? pc.coinEv.cs_nonce + : undefined, }, sourceTransactionId: ctx.transactionId, coinEvHash: pc.coinEvHash, diff --git a/packages/taler-wallet-core/src/shepherd.test.ts b/packages/taler-wallet-core/src/shepherd.test.ts @@ -0,0 +1,48 @@ +/* + This file is part of GNU Taler + (C) 2026 Taler Systems S.A. + + GNU Taler is free software; you can redistribute it and/or modify it under the + terms of the GNU General Public License as published by the Free Software + Foundation; either version 3, or (at your option) any later version. +*/ + +import assert from "node:assert/strict"; +import test from "node:test"; +import { PendingTaskType, constructTaskIdentifier } from "./common.js"; +import { WalletDbTransaction } from "./db/transaction.js"; +import { getActiveTaskIds } from "./shepherd.js"; +import { InternalWalletState } from "./wallet.js"; + +test("expired-transaction cleanup is always scheduled", async () => { + const empty = async () => []; + const tx = { + getActiveWithdrawalGroups: empty, + getActiveDepositGroups: empty, + getActiveRefreshGroups: empty, + getActivePurchases: empty, + getActivePeerPushDebits: empty, + getActivePeerPushCredits: empty, + getActivePeerPullDebits: empty, + getActivePeerPullCredits: empty, + getActiveRecoupGroups: empty, + getExchanges: empty, + } as unknown as WalletDbTransaction; + const ws = { + async runStandaloneWalletDbTx<T>( + f: (transaction: WalletDbTransaction) => Promise<T>, + ): Promise<T> { + return f(tx); + }, + } as unknown as InternalWalletState; + + const result = await getActiveTaskIds(ws); + + assert.ok( + result.taskIds.includes( + constructTaskIdentifier({ + tag: PendingTaskType.CleanupExpiredTransactions, + }), + ), + ); +}); diff --git a/packages/taler-wallet-core/src/shepherd.ts b/packages/taler-wallet-core/src/shepherd.ts @@ -17,6 +17,7 @@ /** * Imports. */ +import { computeRecoupTransactionState } from "./recoup.js"; import { AbsoluteTime, AsyncCondition, @@ -902,12 +903,13 @@ async function taskToRetryNotification( case PendingTaskType.PeerPushDebit: case PendingTaskType.Purchase: return makeTransactionRetryNotification(ws, tx, pendingTaskId, e); - case PendingTaskType.Recoup: case PendingTaskType.ValidateDenoms: case PendingTaskType.CleanupExpiredTransactions: case PendingTaskType.ExchangeWalletKyc: case PendingTaskType.ExchangeAutoRefresh: return undefined; + case PendingTaskType.Recoup: + return makeTransactionRetryNotification(ws, tx, pendingTaskId, e); default: assertUnreachable(parsedTaskId); } @@ -1014,8 +1016,16 @@ async function getTransactionState( stId: rec.operationStatus, }; } - case TransactionType.Recoup: - throw Error("not yet supported"); + case TransactionType.Recoup: { + const rec = await tx.getRecoupGroup(parsedTxId.recoupGroupId); + if (!rec) { + return undefined; + } + return { + txState: computeRecoupTransactionState(rec), + stId: rec.operationStatus, + }; + } case TransactionType.DenomLoss: { const rec = await tx.getDenomLossEvent(parsedTxId.denomLossEventId); if (!rec) { @@ -1152,6 +1162,11 @@ export function convertTaskToTransactionId( tag: TransactionType.Payment, proposalId: parsedTaskId.proposalId, }); + case PendingTaskType.Recoup: + return constructTransactionIdentifier({ + tag: TransactionType.Recoup, + recoupGroupId: parsedTaskId.recoupGroupId, + }); default: return undefined; } @@ -1326,6 +1341,14 @@ export async function getActiveTaskIds( constructTaskIdentifier({ tag: PendingTaskType.ValidateDenoms }), ); + // Transaction expiry is maintenance rather than a liveness task, but it + // must always be present so old proposals are eventually removed. + res.taskIds.push( + constructTaskIdentifier({ + tag: PendingTaskType.CleanupExpiredTransactions, + }), + ); + // FIXME: Recoup! }); @@ -1335,7 +1358,8 @@ export async function getActiveTaskIds( export async function processCleanupExpiredTransactions( wex: WalletExecutionContext, ): Promise<TaskRunResult> { - await wex.runWalletDbTx(async (tx) => { + const proposalIds = await wex.runWalletDbTx(async (tx) => { + const readyForDeletion: string[] = []; const expired = await tx.getPurchasesByStatus(PurchaseStatus.Expired); for (const exp of expired) { if (!exp.timestampExpired) { @@ -1349,10 +1373,16 @@ export async function processCleanupExpiredTransactions( ), ) ) { - const ctx = new PayMerchantTransactionContext(wex, exp.proposalId); - await ctx.userDeleteTransaction(); + readyForDeletion.push(exp.proposalId); } } + return readyForDeletion; }); + // Deletion opens its own transaction. Run it after the scan transaction + // has closed, since the native SQLite queue is intentionally non-reentrant. + for (const proposalId of proposalIds) { + const ctx = new PayMerchantTransactionContext(wex, proposalId); + await ctx.userDeleteTransaction(); + } return TaskRunResult.runAgainAfter({ hours: 24 }); } diff --git a/packages/taler-wallet-core/src/transactions.ts b/packages/taler-wallet-core/src/transactions.ts @@ -72,6 +72,7 @@ import { PeerPullDebitTransactionContext } from "./pay-peer-pull-debit.js"; import { PeerPushCreditTransactionContext } from "./pay-peer-push-credit.js"; import { PeerPushDebitTransactionContext } from "./pay-peer-push-debit.js"; import { RefreshTransactionContext } from "./refresh.js"; +import { RecoupTransactionContext } from "./recoup.js"; import type { WalletExecutionContext } from "./wallet.js"; import { WithdrawTransactionContext } from "./withdraw.js"; import { WalletDbTransaction } from "./db/transaction.js"; @@ -188,7 +189,7 @@ export function makeTransactionNotFoundError( */ export function makeTransactionActionUnsupportedError( transactionId: string, - action: "abort" | "suspend" | "resume" | "fail" | "lookup", + action: "abort" | "suspend" | "resume" | "fail" | "lookup" | "delete", hint?: string, ): TalerError { return TalerError.fromDetail( @@ -765,6 +766,11 @@ export async function rematerializeTransactions( await ctx.updateTransactionMeta(tx); } + for (const x of await tx.listAllRecoupGroups()) { + const ctx = new RecoupTransactionContext(wex, x.recoupGroupId); + await ctx.updateTransactionMeta(tx); + } + for (const x of await tx.listAllWithdrawalGroups()) { const ctx = new WithdrawTransactionContext(wex, x.withdrawalGroupId); await ctx.updateTransactionMeta(tx); @@ -1057,12 +1063,7 @@ async function getContextForTransaction( case TransactionType.Refund: return new RefundTransactionContext(wex, tx.refundGroupId); case TransactionType.Recoup: - //return new RecoupTransactionContext(ws, tx.recoupGroupId); - throw TalerError.fromDetail( - TalerErrorCode.WALLET_TRANSACTION_ACTION_UNSUPPORTED, - { transactionId, action: "lookup" }, - "recoup transactions cannot be inspected yet", - ); + return new RecoupTransactionContext(wex, tx.recoupGroupId); case TransactionType.DenomLoss: return new DenomLossTransactionContext(wex, tx.denomLossEventId); default: