taler-typescript-core

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

commit 7060c6b2c79c7104a1c6b9b8fe72b400ba52cb7f
parent e86e3c1100657a0e375dedac77755046ee96d435
Author: Florian Dold <dold@taler.net>
Date:   Fri,  4 Sep 2026 01:01:41 +0200

wallet-core: persist merge reserves for peer credits

Store the merge reserve for each push or pull credit, verify purse
creation and merge responses, and paginate reserve history during
recovery.

Diffstat:
Mpackages/taler-wallet-core/src/db/indexeddb/fixups.ts | 8+++++++-
Mpackages/taler-wallet-core/src/db/indexeddb/schema.ts | 16++++++++++++++++
Mpackages/taler-wallet-core/src/db/indexeddb/transaction.ts | 6++++++
Mpackages/taler-wallet-core/src/db/records.ts | 41+++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/db/sqlite/schema.ts | 38+++++++++++++++++++++++++++++++++++++-
Mpackages/taler-wallet-core/src/db/sqlite/transaction.ts | 65++++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Mpackages/taler-wallet-core/src/db/testing/conformance-cases.ts | 7+++++++
Apackages/taler-wallet-core/src/pay-peer-common.test.ts | 176+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/pay-peer-common.ts | 234+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/pay-peer-pull-credit.test.ts | 32++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/pay-peer-pull-credit.ts | 322++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mpackages/taler-wallet-core/src/pay-peer-push-credit.test.ts | 82+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/pay-peer-push-credit.ts | 337++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
13 files changed, 1228 insertions(+), 136 deletions(-)

diff --git a/packages/taler-wallet-core/src/db/indexeddb/fixups.ts b/packages/taler-wallet-core/src/db/indexeddb/fixups.ts @@ -921,6 +921,12 @@ async function remapAndDeleteReserve( await tx.peerPullCredit.put(p); } }); + await tx.peerPushCredit.iter().forEachAsync(async (p) => { + if (p.mergeReserveRowId === droppedRowId) { + p.mergeReserveRowId = retainedRowId; + await tx.peerPushCredit.put(p); + } + }); await tx.reserves.delete(droppedRowId); } @@ -931,7 +937,7 @@ async function remapAndDeleteReserve( * instead of upserting, leaving several rows with identical key material * under different row ids. Databases like that exist in the wild. The * lowest row id is kept and the two references to reserve rows (the - * exchange's current merge reserve and peer-pull-credit merge reserves) are + * exchange's current merge reserve and peer-credit merge reserves) are * remapped onto it. * * Rows are removed only when every field except rowId matches the kept row. diff --git a/packages/taler-wallet-core/src/db/indexeddb/schema.ts b/packages/taler-wallet-core/src/db/indexeddb/schema.ts @@ -26,6 +26,7 @@ import { HashCodeString, MailboxConfiguration, MailboxMessageRecord, + PurseCreateSuccessResponse, ScopeInfo, TalerErrorDetail, } from "@gnu-taler/taler-util"; @@ -304,6 +305,8 @@ export interface PeerPullCreditRecord { mergeReserveRowId: number; + purseCreateProof?: PurseCreateSuccessResponse; + /** * Status of the peer pull payment initiation. */ @@ -364,6 +367,8 @@ export interface PeerPushCreditRecord { */ withdrawalGroupId: string | undefined; + mergeReserveRowId?: number; + /** * Currency of the peer push payment credit transaction. * @@ -392,6 +397,12 @@ export interface PeerPullPaymentCoinSelection { * Non undefined after status === "Accepted" */ totalCost: AmountString | undefined; + + /** Number of leading entries confirmed by a signed exchange response. */ + depositedCoinCount?: number; + + /** Latest purse balance covered by a verified deposit confirmation. */ + confirmedPurseBalance?: AmountString; } /** @@ -427,6 +438,11 @@ export interface PeerPullPaymentIncomingRecord { abortRefreshGroupId?: string; + cleanupFinalStatus?: + | PeerPullDebitRecordStatus.Aborted + | PeerPullDebitRecordStatus.Expired + | PeerPullDebitRecordStatus.Failed; + abortReason?: TalerErrorDetail; failReason?: TalerErrorDetail; diff --git a/packages/taler-wallet-core/src/db/indexeddb/transaction.ts b/packages/taler-wallet-core/src/db/indexeddb/transaction.ts @@ -1644,6 +1644,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { contractEncNonce: r.contractEncNonce, mergeTimestamp: r.mergeTimestamp, mergeReserveRowId: r.mergeReserveRowId, + purseCreateProof: r.purseCreateProof, status: r.status, kycPaytoHash: r.kycPaytoHash, kycAccessToken: r.kycAccessToken, @@ -1674,6 +1675,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { contractEncNonce: rec.contractEncNonce, mergeTimestamp: rec.mergeTimestamp, mergeReserveRowId: rec.mergeReserveRowId, + purseCreateProof: rec.purseCreateProof, status: rec.status, kycPaytoHash: rec.kycPaytoHash, kycAccessToken: rec.kycAccessToken, @@ -1775,6 +1777,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { abortReason: r.abortReason, failReason: r.failReason, withdrawalGroupId: r.withdrawalGroupId, + mergeReserveRowId: r.mergeReserveRowId, currency: r.currency, kycPaytoHash: r.kycPaytoHash, kycAccessToken: r.kycAccessToken, @@ -1801,6 +1804,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { abortReason: rec.abortReason, failReason: rec.failReason, withdrawalGroupId: rec.withdrawalGroupId, + mergeReserveRowId: rec.mergeReserveRowId, currency: rec.currency, kycPaytoHash: rec.kycPaytoHash, kycAccessToken: rec.kycAccessToken, @@ -1851,6 +1855,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { status: r.status, totalCostEstimated: r.totalCostEstimated, abortRefreshGroupId: r.abortRefreshGroupId, + cleanupFinalStatus: r.cleanupFinalStatus, abortReason: r.abortReason, failReason: r.failReason, coinSel: r.coinSel, @@ -1870,6 +1875,7 @@ export class IdbWalletTransaction implements WalletDbTransaction { status: rec.status, totalCostEstimated: rec.totalCostEstimated, abortRefreshGroupId: rec.abortRefreshGroupId, + cleanupFinalStatus: rec.cleanupFinalStatus, abortReason: rec.abortReason, failReason: rec.failReason, coinSel: rec.coinSel, diff --git a/packages/taler-wallet-core/src/db/records.ts b/packages/taler-wallet-core/src/db/records.ts @@ -59,6 +59,7 @@ import { hash, stringToBytes, canonicalJson, + PurseCreateSuccessResponse, } from "@gnu-taler/taler-util"; import { DbPreciseTimestamp, @@ -315,6 +316,12 @@ export interface WalletRecoupGroup { */ recoupFinishedPerCoin: boolean[]; + /** Coins whose exchange rejected recoup permanently. */ + failedCoinPubs?: string[]; + + /** Direct-withdraw coins whose value was credited to their reserve. */ + successfulCoinPubs?: string[]; + /** * Public keys of coins that should be scheduled for refreshing * after all individual recoups are done. @@ -777,6 +784,8 @@ export interface WalletDonationSummary { currency: string; amountReceiptsAvailable: AmountString; amountReceiptsSubmitted: AmountString; + /** Total in the latest statement requested from this Donau, if any. */ + amountStatement?: AmountString; } /** @@ -1472,12 +1481,24 @@ export interface WalletWithdrawCoinSource { * Reserve public key for the reserve we got this coin from. */ reservePub: string; + + /** Hash over the batch planchets, needed for protocol-v26+ recoup. */ + hPlanchets?: string; + + /** Legacy protocol-v25 withdrawal commitment, when available. */ + withdrawCommitmentHash?: string; + + /** Clause-Schnorr nonce used for the blinded coin envelope. */ + nonce?: string; } export interface WalletRefreshCoinSource { type: CoinSourceType.Refresh; refreshGroupId: string; oldCoinPub: string; + + /** Clause-Schnorr nonce used for the blinded coin envelope. */ + nonce?: string; } export interface WalletRewardCoinSource { @@ -2277,6 +2298,9 @@ export enum PeerPushCreditStatus { DialogProposed = 0x0101_0000, + /** Reconcile a possibly committed merge after the user aborts. */ + AbortingMerge = 0x0103_0000, + /** Waiting for the rejected purse to expire and refund its payer. */ FinalizingKycHardLimit = 0x0200_0000, @@ -2290,9 +2314,11 @@ export enum PeerPushCreditStatus { export enum PeerPullDebitRecordStatus { PendingDeposit = 0x0100_0001, AbortingRefresh = 0x0103_0001, + AbortingReconcile = 0x0103_0002, SuspendedDeposit = 0x0110_0001, SuspendedAbortingRefresh = 0x0113_0001, + SuspendedAbortingReconcile = 0x0113_0002, DialogProposed = 0x0101_0001, @@ -2511,6 +2537,9 @@ export interface WalletPeerPullCredit { mergeReserveRowId: number; + /** Exchange proof binding this invoice's purse metadata. */ + purseCreateProof?: PurseCreateSuccessResponse; + /** * Status of the peer pull payment initiation. */ @@ -2574,6 +2603,12 @@ export interface WalletPeerPushCredit { withdrawalGroupId: string | undefined; /** + * Reserve selected for the purse merge. Optional only for legacy records; + * current records bind it before the merge request can be sent. + */ + mergeReserveRowId?: number; + + /** * Currency of the peer push payment credit transaction. * * Mandatory in current schema version, optional for compatibility @@ -2629,6 +2664,12 @@ export interface WalletPeerPullDebit { abortRefreshGroupId?: string; + /** Terminal state to enter after safe coin recovery. */ + cleanupFinalStatus?: + | PeerPullDebitRecordStatus.Aborted + | PeerPullDebitRecordStatus.Expired + | PeerPullDebitRecordStatus.Failed; + abortReason?: TalerErrorDetail; failReason?: TalerErrorDetail; diff --git a/packages/taler-wallet-core/src/db/sqlite/schema.ts b/packages/taler-wallet-core/src/db/sqlite/schema.ts @@ -85,7 +85,7 @@ * * Bump this when adding a migration to {@link schemaMigrations}. */ -export const SQLITE_SCHEMA_VERSION = 15; +export const SQLITE_SCHEMA_VERSION = 20; /** * Tables of the IndexedDB emulation, children before parents. @@ -1427,6 +1427,42 @@ export const schemaMigrations: SchemaMigration[] = [ "ALTER TABLE denom_loss_events ADD COLUMN exchange_master_pub BLOB", ], }, + { + version: 16, + name: "donau-statement-total", + statements: [ + "ALTER TABLE donation_summaries ADD COLUMN amount_statement TEXT", + ], + }, + { + version: 17, + name: "push-credit-merge-reserve", + statements: [ + "ALTER TABLE peer_push_credit ADD COLUMN merge_reserve_row_id INTEGER REFERENCES reserves(row_id)", + ], + }, + { + version: 18, + name: "pull-credit-purse-proof", + statements: [ + "ALTER TABLE peer_pull_credit ADD COLUMN purse_create_proof TEXT", + ], + }, + { + version: 19, + name: "pull-debit-cleanup-final-status", + statements: [ + "ALTER TABLE peer_pull_debit ADD COLUMN cleanup_final_status INTEGER", + ], + }, + { + version: 20, + name: "recoup-per-coin-failures", + statements: [ + "ALTER TABLE recoup_groups ADD COLUMN failed_coin_pubs TEXT", + "ALTER TABLE recoup_groups ADD COLUMN successful_coin_pubs TEXT", + ], + }, ]; /** Native tables that contain wallet records (not schema bookkeeping). */ diff --git a/packages/taler-wallet-core/src/db/sqlite/transaction.ts b/packages/taler-wallet-core/src/db/sqlite/transaction.ts @@ -3106,6 +3106,10 @@ export class SqliteWalletTransaction implements WalletDbTransaction { contractTermsHash: dbToCrock(row.contract_terms_hash), status: num(row.status), withdrawalGroupId: optStr(row.withdrawal_group_id), + mergeReserveRowId: + row.merge_reserve_row_id == null + ? undefined + : num(row.merge_reserve_row_id), currency: optStr(row.currency), ...(row.abort_reason != null ? { abortReason: dbToJson(row.abort_reason) } @@ -3153,12 +3157,12 @@ export class SqliteWalletTransaction implements WalletDbTransaction { peer_push_credit_id, exchange_base_url, purse_pub, merge_priv, contract_priv, timestamp, estimated_amount_effective, contract_terms_hash, status, abort_reason, fail_reason, - withdrawal_group_id, currency, kyc_payto_hash, kyc_access_token, + withdrawal_group_id, merge_reserve_row_id, currency, kyc_payto_hash, kyc_access_token, kyc_last_check_status, kyc_last_check_code, kyc_last_rule_gen, kyc_last_aml_review, kyc_last_deny ) VALUES ( $id, $url, $ppub, $mpriv, $cpriv, $ts, $eae, $cth, $status, - $abort, $fail, $wgid, $cur, $kph, $kat, $klcs, $klcc, $klrg, + $abort, $fail, $wgid, $mrri, $cur, $kph, $kat, $klcs, $klcc, $klrg, $klar, $kld ) ON CONFLICT(peer_push_credit_id) DO UPDATE SET @@ -3173,6 +3177,7 @@ export class SqliteWalletTransaction implements WalletDbTransaction { abort_reason = excluded.abort_reason, fail_reason = excluded.fail_reason, withdrawal_group_id = excluded.withdrawal_group_id, + merge_reserve_row_id = excluded.merge_reserve_row_id, currency = excluded.currency, kyc_payto_hash = excluded.kyc_payto_hash, kyc_access_token = excluded.kyc_access_token, @@ -3194,6 +3199,7 @@ export class SqliteWalletTransaction implements WalletDbTransaction { abort: rec.abortReason === undefined ? null : jsonToDb(rec.abortReason), fail: rec.failReason === undefined ? null : jsonToDb(rec.failReason), wgid: rec.withdrawalGroupId ?? null, + mrri: rec.mergeReserveRowId ?? null, cur: rec.currency ?? null, kph: optCrockToDb(rec.kycPaytoHash), kat: rec.kycAccessToken ?? null, @@ -3258,6 +3264,9 @@ export class SqliteWalletTransaction implements WalletDbTransaction { ...(row.abort_refresh_group_id != null ? { abortRefreshGroupId: str(row.abort_refresh_group_id) } : undefined), + ...(row.cleanup_final_status != null + ? { cleanupFinalStatus: num(row.cleanup_final_status) } + : undefined), ...(row.abort_reason != null ? { abortReason: dbToJson(row.abort_reason) } : undefined), @@ -3286,10 +3295,10 @@ export class SqliteWalletTransaction implements WalletDbTransaction { peer_pull_debit_id, purse_pub, exchange_base_url, amount, contract_terms_hash, timestamp_created, contract_priv, status, total_cost_estimated, abort_refresh_group_id, abort_reason, - fail_reason, coin_sel + fail_reason, coin_sel, cleanup_final_status ) VALUES ( $id, $ppub, $url, $amt, $cth, $created, $cpriv, $status, $tce, - $argi, $abort, $fail, $csel + $argi, $abort, $fail, $csel, $cfs ) ON CONFLICT(peer_pull_debit_id) DO UPDATE SET purse_pub = excluded.purse_pub, @@ -3303,7 +3312,8 @@ export class SqliteWalletTransaction implements WalletDbTransaction { abort_refresh_group_id = excluded.abort_refresh_group_id, abort_reason = excluded.abort_reason, fail_reason = excluded.fail_reason, - coin_sel = excluded.coin_sel`, + coin_sel = excluded.coin_sel, + cleanup_final_status = excluded.cleanup_final_status`, { id: rec.peerPullDebitId, ppub: crockToDb(rec.pursePub), @@ -3318,6 +3328,7 @@ export class SqliteWalletTransaction implements WalletDbTransaction { abort: rec.abortReason === undefined ? null : jsonToDb(rec.abortReason), fail: rec.failReason === undefined ? null : jsonToDb(rec.failReason), csel: rec.coinSel === undefined ? null : jsonToDb(rec.coinSel), + cfs: rec.cleanupFinalStatus ?? null, }, ); } @@ -3363,6 +3374,9 @@ export class SqliteWalletTransaction implements WalletDbTransaction { contractEncNonce: dbToCrock(row.contract_enc_nonce), mergeTimestamp: dbTimestamp(row.merge_timestamp), mergeReserveRowId: num(row.merge_reserve_row_id), + ...(row.purse_create_proof != null + ? { purseCreateProof: dbToJson(row.purse_create_proof) } + : undefined), status: num(row.status), withdrawalGroupId: optStr(row.withdrawal_group_id), ...(row.kyc_payto_hash != null @@ -3414,11 +3428,11 @@ export class SqliteWalletTransaction implements WalletDbTransaction { merge_reserve_row_id, status, kyc_payto_hash, kyc_access_token, kyc_last_check_status, kyc_last_check_code, kyc_last_rule_gen, kyc_last_aml_review, kyc_last_deny, abort_reason, fail_reason, - withdrawal_group_id + withdrawal_group_id, purse_create_proof ) VALUES ( $pub, $url, $amt, $eae, $ppriv, $cth, $mpub, $mpriv, $cpub, $cpriv, $nonce, $mts, $mrri, $status, $kph, $kat, $klcs, $klcc, - $klrg, $klar, $kld, $abort, $fail, $wgid + $klrg, $klar, $kld, $abort, $fail, $wgid, $pcp ) ON CONFLICT(purse_pub) DO UPDATE SET exchange_base_url = excluded.exchange_base_url, @@ -3443,7 +3457,8 @@ export class SqliteWalletTransaction implements WalletDbTransaction { kyc_last_deny = excluded.kyc_last_deny, abort_reason = excluded.abort_reason, fail_reason = excluded.fail_reason, - withdrawal_group_id = excluded.withdrawal_group_id`, + withdrawal_group_id = excluded.withdrawal_group_id, + purse_create_proof = excluded.purse_create_proof`, { pub: crockToDb(rec.pursePub), url: rec.exchangeBaseUrl, @@ -3469,6 +3484,10 @@ export class SqliteWalletTransaction implements WalletDbTransaction { abort: rec.abortReason === undefined ? null : jsonToDb(rec.abortReason), fail: rec.failReason === undefined ? null : jsonToDb(rec.failReason), wgid: rec.withdrawalGroupId ?? null, + pcp: + rec.purseCreateProof === undefined + ? null + : jsonToDb(rec.purseCreateProof), }, ); } @@ -4307,6 +4326,9 @@ export class SqliteWalletTransaction implements WalletDbTransaction { currency: str(row.currency), amountReceiptsAvailable: dbAmount(row.amount_receipts_available), amountReceiptsSubmitted: dbAmount(row.amount_receipts_submitted), + ...(row.amount_statement != null + ? { amountStatement: dbAmount(row.amount_statement) } + : undefined), ...(row.legal_domain != null ? { legalDomain: str(row.legal_domain) } : undefined), @@ -4335,12 +4357,13 @@ export class SqliteWalletTransaction implements WalletDbTransaction { await this.run( `INSERT INTO donation_summaries ( donau_base_url, year, currency, legal_domain, - amount_receipts_available, amount_receipts_submitted - ) VALUES ($url, $year, $cur, $ld, $avail, $sub) + amount_receipts_available, amount_receipts_submitted, amount_statement + ) VALUES ($url, $year, $cur, $ld, $avail, $sub, $stmt) ON CONFLICT(donau_base_url, year, currency) DO UPDATE SET legal_domain = excluded.legal_domain, amount_receipts_available = excluded.amount_receipts_available, - amount_receipts_submitted = excluded.amount_receipts_submitted`, + amount_receipts_submitted = excluded.amount_receipts_submitted, + amount_statement = excluded.amount_statement`, { url: rec.donauBaseUrl, year: rec.year, @@ -4348,6 +4371,7 @@ export class SqliteWalletTransaction implements WalletDbTransaction { ld: rec.legalDomain ?? null, avail: rec.amountReceiptsAvailable, sub: rec.amountReceiptsSubmitted, + stmt: rec.amountStatement ?? null, }, ); } @@ -5279,6 +5303,14 @@ export class SqliteWalletTransaction implements WalletDbTransaction { : dbTimestamp(row.timestamp_finished), coinPubs: dbToJson(row.coin_pubs), recoupFinishedPerCoin: dbToJson(row.recoup_finished_per_coin), + failedCoinPubs: + row.failed_coin_pubs == null + ? undefined + : dbToJson(row.failed_coin_pubs), + successfulCoinPubs: + row.successful_coin_pubs == null + ? undefined + : dbToJson(row.successful_coin_pubs), scheduleRefreshCoins: dbToJson(row.schedule_refresh_coins), }; } @@ -5303,8 +5335,9 @@ export class SqliteWalletTransaction implements WalletDbTransaction { `INSERT INTO recoup_groups ( recoup_group_id, exchange_base_url, operation_status, timestamp_started, timestamp_finished, coin_pubs, - recoup_finished_per_coin, schedule_refresh_coins - ) VALUES ($id, $url, $status, $started, $finished, $pubs, $fin, $sched) + recoup_finished_per_coin, failed_coin_pubs, successful_coin_pubs, + schedule_refresh_coins + ) VALUES ($id, $url, $status, $started, $finished, $pubs, $fin, $failed, $success, $sched) ON CONFLICT(recoup_group_id) DO UPDATE SET exchange_base_url = excluded.exchange_base_url, operation_status = excluded.operation_status, @@ -5312,6 +5345,8 @@ export class SqliteWalletTransaction implements WalletDbTransaction { timestamp_finished = excluded.timestamp_finished, coin_pubs = excluded.coin_pubs, recoup_finished_per_coin = excluded.recoup_finished_per_coin, + failed_coin_pubs = excluded.failed_coin_pubs, + successful_coin_pubs = excluded.successful_coin_pubs, schedule_refresh_coins = excluded.schedule_refresh_coins`, { id: rec.recoupGroupId, @@ -5321,6 +5356,10 @@ export class SqliteWalletTransaction implements WalletDbTransaction { finished: rec.timestampFinished ?? null, pubs: jsonToDb(rec.coinPubs), fin: jsonToDb(rec.recoupFinishedPerCoin), + failed: rec.failedCoinPubs ? jsonToDb(rec.failedCoinPubs) : null, + success: rec.successfulCoinPubs + ? jsonToDb(rec.successfulCoinPubs) + : null, sched: jsonToDb(rec.scheduleRefreshCoins), }, ); diff --git a/packages/taler-wallet-core/src/db/testing/conformance-cases.ts b/packages/taler-wallet-core/src/db/testing/conformance-cases.ts @@ -684,6 +684,12 @@ function makePeerPullCredit(pursePub: string): WalletPeerPullCredit { contractEncNonce: ck("nonce"), mergeTimestamp: tsPrecise(1000), mergeReserveRowId: 1, + purseCreateProof: { + total_deposited: amt("TESTKUDOS:0"), + exchange_timestamp: { t_s: 1001 }, + exchange_pub: ck("purse-proof-pub"), + exchange_sig: ckOfSize("purse-proof-sig", 64), + }, status: PeerPullPaymentCreditStatus.PendingCreatePurse, withdrawalGroupId: undefined, }; @@ -3983,6 +3989,7 @@ export const conformanceCases: ConformanceCase[] = [ depositedCoinCount: 1, confirmedPurseBalance: amt("TESTKUDOS:1"), }; + rec.cleanupFinalStatus = PeerPullDebitRecordStatus.Aborted; await runner.runReadWriteTx((tx) => tx.upsertPeerPullDebit(rec)); const got = await runner.runReadWriteTx((tx) => tx.getPeerPullDebit("ppld-1"), diff --git a/packages/taler-wallet-core/src/pay-peer-common.test.ts b/packages/taler-wallet-core/src/pay-peer-common.test.ts @@ -0,0 +1,176 @@ +/* + 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 { + CancellationToken, + HttpStatusCode, + ReserveHistoryMergeEntry, + TalerProtocolTimestamp, + WalletAccountMergeFlags, + encodeCrock, +} from "@gnu-taler/taler-util"; +import type { + HttpRequestLibrary, + HttpRequestOptions, + HttpResponse, +} from "@gnu-taler/taler-util/http"; +import assert from "node:assert/strict"; +import test from "node:test"; +import { + ExpectedPurseMergeHistoryEntry, + getMergedPurseAmountFromReserveHistory, + validateExpectedPurseMergeHistoryEntry, +} from "./pay-peer-common.js"; +import { WalletExecutionContext } from "./wallet.js"; + +const mergeTimestamp = TalerProtocolTimestamp.fromSeconds(1_800_000_000); +const purseExpiration = TalerProtocolTimestamp.fromSeconds(1_800_086_400); + +const expected: ExpectedPurseMergeHistoryEntry = { + pursePub: "purse-pub", + contractTermsHash: "contract-hash", + mergePub: "merge-pub", + mergeTimestamp, + purseExpiration, + purseAmount: "TESTKUDOS:2", + purseFee: "TESTKUDOS:0", + flags: WalletAccountMergeFlags.MergeFullyPaidPurse, + minAge: 0, + reserveSig: "reserve-sig", +}; + +const validMergeEntry: ReserveHistoryMergeEntry = { + type: "MERGE", + history_offset: 2, + h_contract_terms: expected.contractTermsHash, + merge_pub: expected.mergePub, + min_age: expected.minAge, + flags: expected.flags, + purse_pub: expected.pursePub, + reserve_sig: expected.reserveSig, + merge_timestamp: expected.mergeTimestamp, + purse_expiration: expected.purseExpiration, + purse_fee: expected.purseFee, + amount: expected.purseAmount, + merged: true, +}; + +test("reserve history merge validation binds every signed field", () => { + validateExpectedPurseMergeHistoryEntry(validMergeEntry, expected); + for (const mutation of [ + { reserve_sig: "wrong-signature" }, + { merge_pub: "wrong-merge-key" }, + { flags: WalletAccountMergeFlags.CreateFromPurseQuota }, + { min_age: 1 }, + { merge_timestamp: TalerProtocolTimestamp.fromSeconds(1_800_000_001) }, + { purse_expiration: TalerProtocolTimestamp.fromSeconds(1_800_086_401) }, + { purse_fee: "TESTKUDOS:0.1" }, + { amount: "TESTKUDOS:3" }, + ] as unknown as Partial<ReserveHistoryMergeEntry>[]) { + assert.throws(() => + validateExpectedPurseMergeHistoryEntry( + { ...validMergeEntry, ...mutation }, + expected, + ), + ); + } +}); + +function response(url: string, status: number, body?: unknown): HttpResponse { + return { + requestUrl: url, + requestMethod: "GET", + status, + headers: { + get(name: string): string | null { + return name.toLowerCase() === "content-type" && body !== undefined + ? "application/json" + : null; + }, + } as HttpResponse["headers"], + json: async () => body, + text: async () => (body === undefined ? "" : JSON.stringify(body)), + bytes: async () => new Uint8Array(), + }; +} + +test("shared reserve recovery traverses all history pages", async () => { + const requests: Array<{ url: string; signature?: string }> = []; + const http: HttpRequestLibrary = { + async fetch(url: string, options?: HttpRequestOptions) { + requests.push({ + url, + signature: options?.headers?.["Taler-Reserve-History-Signature"], + }); + const start = new URL(url).searchParams.get("start"); + if (start === null) { + return response(url, HttpStatusCode.Ok, { + balance: "TESTKUDOS:7", + history: [{ type: "CREDIT", history_offset: 1 }], + }); + } + if (start === "1") { + return response(url, HttpStatusCode.Ok, { + balance: "TESTKUDOS:7", + history: [validMergeEntry], + }); + } + assert.strictEqual(start, "2"); + return response(url, HttpStatusCode.NoContent); + }, + }; + const wex = { + http, + cancellationToken: CancellationToken.CONTINUE, + ws: { longpollQueue: undefined }, + cryptoApi: { + async eddsaGetPublic() { + return { pub: expected.mergePub }; + }, + async signPurseMerge() { + return { mergeSig: "merge-sig", accountSig: expected.reserveSig }; + }, + async signReserveHistoryReq(req: { startOffset: number }) { + return { sig: `history-sig-${req.startOffset}` }; + }, + }, + } as unknown as WalletExecutionContext; + + const result = await getMergedPurseAmountFromReserveHistory(wex, { + exchangeBaseUrl: "https://exchange.example/", + reserve: { + reservePub: encodeCrock(new Uint8Array(32).fill(1)), + reservePriv: "reserve-priv", + }, + pursePub: expected.pursePub, + contractTermsHash: expected.contractTermsHash, + mergePriv: "merge-priv", + mergeTimestamp: expected.mergeTimestamp, + purseExpiration: expected.purseExpiration, + purseAmount: expected.purseAmount, + purseFee: expected.purseFee, + flags: expected.flags, + }); + + assert.deepStrictEqual(result, { + status: "merged", + amount: expected.purseAmount, + }); + assert.deepStrictEqual( + requests.map((x) => [ + new URL(x.url).searchParams.get("start"), + x.signature, + ]), + [ + [null, "history-sig-0"], + ["1", "history-sig-1"], + ["2", "history-sig-2"], + ], + ); +}); diff --git a/packages/taler-wallet-core/src/pay-peer-common.ts b/packages/taler-wallet-core/src/pay-peer-common.ts @@ -18,12 +18,19 @@ import { AbsoluteTime, AmountJson, + AmountString, Amounts, ExchangePurseStatus, + HttpStatusCode, + ReserveHistoryMergeEntry, SelectedProspectiveCoin, + TalerError, + TalerErrorCode, TalerProtocolTimestamp, + WalletAccountMergeFlags, checkDbInvariant, } from "@gnu-taler/taler-util"; +import { reservePaytoFromExchange } from "./common.js"; import { WalletReserve } from "./db/records.js"; import { SpendCoinDetails } from "./crypto/cryptoImplementation.js"; import { DbPeerPushPaymentCoinSelection } from "./db/indexeddb/schema.js"; @@ -36,6 +43,7 @@ import { denomRefKey, getDenomInfos, WalletExecutionContext, + walletExchangeClient, } from "./wallet.js"; import { updateWithdrawalDenomsForExchange } from "./withdraw.js"; import { WalletDbTransaction } from "./db/transaction.js"; @@ -172,6 +180,232 @@ export async function getMergeReserveInfo( return mergeReserveRecord; } +export type MergedPurseHistoryResult = + | { status: "merged"; amount: AmountString } + | { status: "not-merged" }; + +export interface ExpectedPurseMergeHistoryEntry { + pursePub: string; + contractTermsHash: string; + mergePub: string; + mergeTimestamp: TalerProtocolTimestamp; + purseExpiration: TalerProtocolTimestamp; + purseAmount: AmountString; + purseFee: AmountString; + flags: WalletAccountMergeFlags; + minAge: number; + reserveSig: string; +} + +function reserveHistoryProtocolViolation(message: string): TalerError { + return TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + message, + ); +} + +/** Validate the locally signed fields of this operation's MERGE entry. */ +export function validateExpectedPurseMergeHistoryEntry( + entry: ReserveHistoryMergeEntry, + expected: ExpectedPurseMergeHistoryEntry, +): void { + if (entry.purse_pub !== expected.pursePub) { + throw reserveHistoryProtocolViolation( + "reserve history merge has the wrong purse public key", + ); + } + if ( + entry.h_contract_terms !== expected.contractTermsHash || + entry.merge_pub !== expected.mergePub || + entry.merge_timestamp.t_s !== expected.mergeTimestamp.t_s || + entry.purse_expiration.t_s !== expected.purseExpiration.t_s || + entry.flags !== expected.flags || + entry.min_age !== expected.minAge || + entry.reserve_sig !== expected.reserveSig || + Amounts.cmp(entry.purse_fee, expected.purseFee) !== 0 || + Amounts.cmp(entry.amount, expected.purseAmount) !== 0 + ) { + throw reserveHistoryProtocolViolation( + "reserve history merge does not match the wallet's signed intent", + ); + } +} + +/** + * Reconcile an ambiguous/purged purse through the reserve's authenticated + * history. The shared reserve is the stable P2P AML/KYC identity, so each + * operation is matched by its purse and all fields covered by our reserve + * signature. Every page is fetched before absence is considered conclusive. + */ +export async function getMergedPurseAmountFromReserveHistory( + wex: WalletExecutionContext, + args: { + exchangeBaseUrl: string; + reserve: WalletReserve; + pursePub: string; + contractTermsHash: string; + mergePriv: string; + mergeTimestamp: TalerProtocolTimestamp; + purseExpiration: TalerProtocolTimestamp; + purseAmount: AmountString; + purseFee: AmountString; + flags: WalletAccountMergeFlags; + minAge?: number; + /** Locally signed legacy variants accepted for pre-migration records. */ + alternativePurseTerms?: Array<{ + purseFee: AmountString; + flags: WalletAccountMergeFlags; + }>; + }, +): Promise<MergedPurseHistoryResult> { + const minAge = args.minAge ?? 0; + const reservePayto = reservePaytoFromExchange( + args.exchangeBaseUrl, + args.reserve.reservePub, + ); + const mergeKey = await wex.cryptoApi.eddsaGetPublic({ priv: args.mergePriv }); + const purseTerms = [ + { purseFee: args.purseFee, flags: args.flags }, + ...(args.alternativePurseTerms ?? []), + ]; + const expectedEntries = await Promise.all( + purseTerms.map(async (terms): Promise<ExpectedPurseMergeHistoryEntry> => { + const expectedSigs = await wex.cryptoApi.signPurseMerge({ + contractTermsHash: args.contractTermsHash, + flags: terms.flags, + mergePriv: args.mergePriv, + mergeTimestamp: args.mergeTimestamp, + purseAmount: args.purseAmount, + purseExpiration: args.purseExpiration, + purseFee: terms.purseFee, + pursePub: args.pursePub, + reservePayto, + reservePriv: args.reserve.reservePriv, + }); + return { + pursePub: args.pursePub, + contractTermsHash: args.contractTermsHash, + mergePub: mergeKey.pub, + mergeTimestamp: args.mergeTimestamp, + purseExpiration: args.purseExpiration, + purseAmount: args.purseAmount, + purseFee: terms.purseFee, + flags: terms.flags, + minAge, + reserveSig: expectedSigs.accountSig, + }; + }), + ); + const exchangeClient = walletExchangeClient(args.exchangeBaseUrl, wex); + const seenOffsets = new Set<number>(); + let startOffset = 0; + let matchingEntry: ReserveHistoryMergeEntry | undefined; + for (;;) { + const sig = await wex.cryptoApi.signReserveHistoryReq({ + reservePriv: args.reserve.reservePriv, + startOffset, + }); + const resp = await exchangeClient.getReserveHistory( + args.reserve.reservePub, + sig.sig, + startOffset, + ); + switch (resp.case) { + case HttpStatusCode.NotFound: + if (startOffset !== 0) { + throw reserveHistoryProtocolViolation( + "reserve disappeared while its history was being paginated", + ); + } + return { status: "not-merged" }; + case HttpStatusCode.NoContent: + case HttpStatusCode.NotModified: + return matchingEntry?.merged + ? { status: "merged", amount: matchingEntry.amount } + : { status: "not-merged" }; + case HttpStatusCode.Forbidden: + throw Error("exchange rejected authenticated reserve history request"); + case "ok": + break; + default: + throw Error("unexpected reserve history response"); + } + + if ( + Amounts.currencyOf(resp.body.balance) !== + Amounts.currencyOf(args.purseAmount) + ) { + throw reserveHistoryProtocolViolation( + "reserve history balance uses the wrong currency", + ); + } + + let nextOffset = startOffset; + for (const entry of resp.body.history) { + if ( + !Number.isSafeInteger(entry.history_offset) || + entry.history_offset <= startOffset || + seenOffsets.has(entry.history_offset) + ) { + throw reserveHistoryProtocolViolation( + "reserve history contains a duplicate or non-increasing offset", + ); + } + seenOffsets.add(entry.history_offset); + nextOffset = Math.max(nextOffset, entry.history_offset); + if (entry.type !== "MERGE" || entry.purse_pub !== args.pursePub) { + continue; + } + if (matchingEntry) { + throw reserveHistoryProtocolViolation( + "reserve history contains duplicate merge entries for one purse", + ); + } + let valid = false; + for (const expected of expectedEntries) { + try { + validateExpectedPurseMergeHistoryEntry(entry, expected); + valid = true; + break; + } catch { + // Try another locally reproducible legacy signature variant. + } + } + if (!valid) { + throw reserveHistoryProtocolViolation( + "reserve history merge does not match any locally signed intent", + ); + } + matchingEntry = entry; + } + if (nextOffset === startOffset) { + throw reserveHistoryProtocolViolation( + "reserve history pagination made no progress", + ); + } + startOffset = nextOffset; + } +} + +/** Compute the post-merge reserve balance used by KYC BALANCE rules. */ +export async function getReserveBalanceAfterCredit( + wex: WalletExecutionContext, + exchangeBaseUrl: string, + reserve: WalletReserve, + credit: AmountString, +): Promise<AmountString> { + const status = await walletExchangeClient( + exchangeBaseUrl, + wex, + ).getReserveStatus(reserve.reservePub); + const current = + status.case === "ok" + ? Amounts.parseOrThrow(status.body.balance) + : Amounts.zeroOfAmount(credit); + return Amounts.stringify(Amounts.add(current, credit).amount); +} + /** Check if a purse is merged */ export function isPurseMerged(purse: ExchangePurseStatus): boolean { const mergeTimestamp = purse.merge_timestamp; diff --git a/packages/taler-wallet-core/src/pay-peer-pull-credit.test.ts b/packages/taler-wallet-core/src/pay-peer-pull-credit.test.ts @@ -15,6 +15,7 @@ */ import { + AmountString, TalerProtocolTimestamp, TransactionAction, TransactionMajorState, @@ -30,7 +31,9 @@ import { WalletDbTransaction } from "./db/transaction.js"; import { computePeerPullCreditTransactionActions, computePeerPullCreditTransactionState, + getPullCreditPurseCreationTerms, PeerPullCreditTransactionContext, + statusAfterPullCreditCreate, statusAfterPullCreditDeleteConflict, } from "./pay-peer-pull-credit.js"; import { WalletExecutionContext } from "./wallet.js"; @@ -131,6 +134,35 @@ test("delete conflict continues a pull credit whose purse was merged", () => { ); }); +test("pull-credit purse creation uses the reserve quota", () => { + assert.deepStrictEqual( + getPullCreditPurseCreationTerms("TESTKUDOS:12" as AmountString), + { + flags: 2, + purseFee: "TESTKUDOS:0", + }, + ); +}); + +test("late pull-credit create confirmations preserve cleanup", () => { + assert.strictEqual( + statusAfterPullCreditCreate(PeerPullPaymentCreditStatus.PendingCreatePurse), + PeerPullPaymentCreditStatus.PendingReady, + ); + assert.strictEqual( + statusAfterPullCreditCreate(PeerPullPaymentCreditStatus.Aborted), + PeerPullPaymentCreditStatus.AbortingDeletePurse, + ); + assert.strictEqual( + statusAfterPullCreditCreate(PeerPullPaymentCreditStatus.Expired), + PeerPullPaymentCreditStatus.ExpiredDeletePurse, + ); + assert.strictEqual( + statusAfterPullCreditCreate(PeerPullPaymentCreditStatus.FailedKycHardLimit), + PeerPullPaymentCreditStatus.FinalizingKycHardLimit, + ); +}); + test("pull-credit hard-KYC cleanup has explicit finalizing and failed states", () => { assert.deepStrictEqual( computePeerPullCreditTransactionState( diff --git a/packages/taler-wallet-core/src/pay-peer-pull-credit.ts b/packages/taler-wallet-core/src/pay-peer-pull-credit.ts @@ -16,6 +16,7 @@ import { AbsoluteTime, + AmountString, Amounts, CheckPeerPullCreditRequest, CheckPeerPullCreditResponse, @@ -101,7 +102,9 @@ import { runKycCheckAlgo, } from "./kyc.js"; import { + getMergedPurseAmountFromReserveHistory, getMergeReserveInfo, + getReserveBalanceAfterCredit, isPurseGoneByExpiration, isPurseDeposited, isPurseMerged, @@ -121,10 +124,43 @@ import { const logger = new Logger("pay-peer-pull-credit.ts"); +function makePeerPullPaymentUri(rec: WalletPeerPullCredit): string { + const proof = rec.purseCreateProof; + return TalerUris.stringify({ + type: TalerUriAction.PayPull, + exchangeBaseUrl: rec.exchangeBaseUrl as HostPortPath, + contractPriv: rec.contractPriv, + purseCreateProof: proof + ? { + totalDeposited: proof.total_deposited, + exchangeTimestamp: proof.exchange_timestamp, + exchangeSig: proof.exchange_sig, + exchangePub: proof.exchange_pub, + } + : undefined, + }); +} + const defaultPeerPullExpiration = Duration.toTalerProtocolDuration( Duration.fromSpec({ days: 1 }), ); +/** + * Pull-credit purses are created against the merge reserve's free purse + * quota. That reserve is deliberately not funded before the payer merges + * the purse, so promising to pay a purse fee from it is both misleading and + * generally impossible. + */ +export function getPullCreditPurseCreationTerms(amount: AmountString): { + flags: WalletAccountMergeFlags; + purseFee: AmountString; +} { + return { + flags: WalletAccountMergeFlags.CreateFromPurseQuota, + purseFee: Amounts.stringify(Amounts.zeroOfAmount(amount)), + }; +} + function requireFutureFinitePurseExpiration( purseExpiration: TalerProtocolTimestamp, ): void { @@ -278,11 +314,7 @@ export class PeerPullCreditTransactionContext implements TransactionContext { summary: peerContractTerms.summary, iconId: peerContractTerms.icon_id, }, - talerUri: TalerUris.stringify({ - type: TalerUriAction.PayPull, - exchangeBaseUrl: wsr.exchangeBaseUrl as HostPortPath, // FIXME: change record type - contractPriv: wsr.wgInfo.contractPriv, - }), + talerUri: makePeerPullPaymentUri(pullCredit), transactionId: this.transactionId, abortReason: pullCredit.abortReason, failReason: pullCredit.failReason, @@ -316,11 +348,7 @@ export class PeerPullCreditTransactionContext implements TransactionContext { summary: peerContractTerms.summary, iconId: peerContractTerms.icon_id, }, - talerUri: TalerUris.stringify({ - type: TalerUriAction.PayPull, - exchangeBaseUrl: pullCredit.exchangeBaseUrl as HostPortPath, // FIXME: change record type - contractPriv: pullCredit.contractPriv, - }), + talerUri: makePeerPullPaymentUri(pullCredit), transactionId: this.transactionId, kycUrl, kycAccessToken: pullCredit.kycAccessToken, @@ -601,6 +629,98 @@ async function beginPeerPullCreditExpirationCleanup( }); } +async function recoverPeerPullCreditFromReserve( + wex: WalletExecutionContext, + pullIni: WalletPeerPullCredit, + reason: string, + confirmedAmount?: AmountString, +): Promise<TaskRunResult | undefined> { + const reserve = await wex.runWalletDbTx((tx) => + tx.getReserve(pullIni.mergeReserveRowId), + ); + if (!reserve) { + throw Error("reserve for peer pull credit not found in wallet DB"); + } + let recoveredAmount = confirmedAmount; + if (recoveredAmount == null) { + const contractRecord = await wex.runWalletDbTx((tx) => + tx.getContractTerms(pullIni.contractTermsHash), + ); + checkDbInvariant( + !!contractRecord, + "contract terms for peer pull credit are missing", + ); + const contractTerms = contractRecord.contractTermsRaw as PeerContractTerms; + const purseCreationTerms = getPullCreditPurseCreationTerms(pullIni.amount); + const historyResult = await getMergedPurseAmountFromReserveHistory(wex, { + exchangeBaseUrl: pullIni.exchangeBaseUrl, + reserve, + pursePub: pullIni.pursePub, + contractTermsHash: pullIni.contractTermsHash, + mergePriv: pullIni.mergePriv, + mergeTimestamp: TalerPreciseTimestamp.round( + timestampPreciseFromDb(pullIni.mergeTimestamp), + ), + purseExpiration: contractTerms.purse_expiration, + purseAmount: pullIni.amount, + purseFee: purseCreationTerms.purseFee, + flags: purseCreationTerms.flags, + alternativePurseTerms: [ + { + purseFee: Amounts.stringify(Amounts.zeroOfAmount(pullIni.amount)), + flags: WalletAccountMergeFlags.CreateWithPurseFee, + }, + ], + }); + if (historyResult.status === "not-merged") { + return undefined; + } + recoveredAmount = historyResult.amount; + } + if (Amounts.cmp(recoveredAmount, pullIni.amount) < 0) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "exchange confirmed a peer merge below the contract amount", + ); + } + + await internalCreateWithdrawalGroup(wex, { + amount: Amounts.parseOrThrow(recoveredAmount), + wgInfo: { + withdrawalType: WithdrawalRecordType.PeerPullCredit, + contractPriv: pullIni.contractPriv, + }, + forcedWithdrawalGroupId: pullIni.withdrawalGroupId, + exchangeBaseUrl: pullIni.exchangeBaseUrl, + reserveStatus: WithdrawalGroupStatus.PendingQueryingStatus, + reserveKeyPair: { + priv: reserve.reservePriv, + pub: reserve.reservePub, + }, + }); + + const ctx = new PeerPullCreditTransactionContext(wex, pullIni.pursePub); + await wex.runWalletDbTx(async (tx) => { + const [rec, h] = await ctx.getRecordHandle(tx); + if (!rec) { + return; + } + switch (rec.status) { + case PeerPullPaymentCreditStatus.PendingReady: + case PeerPullPaymentCreditStatus.AbortingDeletePurse: + case PeerPullPaymentCreditStatus.ExpiredDeletePurse: + case PeerPullPaymentCreditStatus.FinalizingKycHardLimit: + rec.status = PeerPullPaymentCreditStatus.PendingWithdrawing; + break; + default: + return; + } + await h.update(rec, reason); + }); + return TaskRunResult.progress(); +} + async function processPendingReady( wex: WalletExecutionContext, pullIni: WalletPeerPullCredit, @@ -633,8 +753,13 @@ async function processPendingReady( }); return TaskRunResult.finished(); case HttpStatusCode.NotFound: - await ctx.failTransaction(pullIni.status, resp.detail); - return TaskRunResult.finished(); + return ( + (await recoverPeerPullCreditFromReserve( + wex, + pullIni, + "ready-purse-purged-recovery", + )) ?? TaskRunResult.backoff() + ); default: assertUnreachable(resp); } @@ -650,43 +775,14 @@ async function processPendingReady( return TaskRunResult.longpollReturnedPending(); } - const reserve = await wex.runWalletDbTx((tx) => - tx.getReserve(pullIni.mergeReserveRowId), + return ( + (await recoverPeerPullCreditFromReserve( + wex, + pullIni, + "ready-merged", + resp.body.balance, + )) ?? TaskRunResult.backoff() ); - - if (!reserve) { - throw Error("reserve for peer pull credit not found in wallet DB"); - } - - await internalCreateWithdrawalGroup(wex, { - amount: Amounts.parseOrThrow(pullIni.amount), - wgInfo: { - withdrawalType: WithdrawalRecordType.PeerPullCredit, - contractPriv: pullIni.contractPriv, - }, - forcedWithdrawalGroupId: pullIni.withdrawalGroupId, - exchangeBaseUrl: pullIni.exchangeBaseUrl, - reserveStatus: WithdrawalGroupStatus.PendingQueryingStatus, - reserveKeyPair: { - priv: reserve.reservePriv, - pub: reserve.reservePub, - }, - }); - await wex.runWalletDbTx(async (tx) => { - const [rec, h] = await ctx.getRecordHandle(tx); - if (!rec) { - return; - } - switch (rec.status) { - case PeerPullPaymentCreditStatus.PendingReady: - rec.status = PeerPullPaymentCreditStatus.PendingWithdrawing; - break; - default: - return; - } - await h.update(rec, "ready-merged"); - }); - return TaskRunResult.progress(); } async function processPendingMergeKycRequired( @@ -696,13 +792,19 @@ async function processPendingMergeKycRequired( const ctx = new PeerPullCreditTransactionContext(wex, pullIni.pursePub); const { exchangeBaseUrl, kycPaytoHash } = pullIni; - // FIXME: What if this changes? Should be part of the p2p record - const mergeReserveInfo = await getMergeReserveInfo(wex, { - exchangeBaseUrl, - }); + const mergeReserveInfo = await wex.runWalletDbTx((tx) => + tx.getReserve(pullIni.mergeReserveRowId), + ); + checkDbInvariant(!!mergeReserveInfo, "bound pull-credit reserve is missing"); const accountPub = mergeReserveInfo.reservePub; const accountPriv = mergeReserveInfo.reservePriv; + const balanceAfterOperation = await getReserveBalanceAfterCredit( + wex, + exchangeBaseUrl, + mergeReserveInfo, + pullIni.amount, + ); let myKycState: GenericKycStatusReq | undefined; @@ -712,6 +814,7 @@ async function processPendingMergeKycRequired( accountPub, // FIXME: Is this the correct amount? amount: pullIni.estimatedAmountEffective, + balanceAfterOperation, exchangeBaseUrl, operation: "MERGE", paytoHash: kycPaytoHash, @@ -788,8 +891,18 @@ async function processPeerPullCreditAbortingDeletePurse( let completionReason: TalerErrorDetail | undefined; switch (resp.case) { case "ok": - case HttpStatusCode.NotFound: break; + case HttpStatusCode.NotFound: { + const recovery = await recoverPeerPullCreditFromReserve( + wex, + peerPullIni, + "delete-purse-purged-recovery", + ); + if (recovery) { + return recovery; + } + break; + } case HttpStatusCode.Forbidden: if (!hardLimitRecovery) { await ctx.failTransaction(peerPullIni.status, resp.detail); @@ -809,13 +922,30 @@ async function processPeerPullCreditAbortingDeletePurse( // If the payer won the race with our abort, continue through the // normal reserve-withdrawal path instead of retrying DELETE. if (isPurseMerged(statusResp.body)) { - completionStatus = PeerPullPaymentCreditStatus.PendingReady; + return ( + (await recoverPeerPullCreditFromReserve( + wex, + peerPullIni, + "delete-conflict-merged-recovery", + statusResp.body.balance, + )) ?? TaskRunResult.backoff() + ); } break; case HttpStatusCode.Gone: - case HttpStatusCode.NotFound: // No merged purse remains, so there is nothing left to receive. break; + case HttpStatusCode.NotFound: { + const recovery = await recoverPeerPullCreditFromReserve( + wex, + peerPullIni, + "delete-conflict-purse-purged-recovery", + ); + if (recovery) { + return recovery; + } + break; + } default: assertUnreachable(statusResp); } @@ -834,23 +964,16 @@ async function processPeerPullCreditAbortingDeletePurse( if (completionReason) { rec.failReason = completionReason; } - if (completionStatus === PeerPullPaymentCreditStatus.PendingReady) { - delete rec.abortReason; - } await h.update( rec, - completionStatus === PeerPullPaymentCreditStatus.PendingReady - ? "abort-purse-merged" - : hardLimitRecovery - ? "kyc-hard-limit-delete-purse" - : expirationRecovery - ? "expire-delete-purse" - : "aborting-delete-purse", + hardLimitRecovery + ? "kyc-hard-limit-delete-purse" + : expirationRecovery + ? "expire-delete-purse" + : "aborting-delete-purse", ); }); - return completionStatus === PeerPullPaymentCreditStatus.PendingReady - ? TaskRunResult.progress() - : TaskRunResult.finished(); + return TaskRunResult.finished(); } export function statusAfterPullCreditDeleteConflict( @@ -979,7 +1102,7 @@ async function processPeerPullCreditCreatePurse( return TaskRunResult.progress(); } - const purseFee = Amounts.stringify(Amounts.zeroOfAmount(pullIni.amount)); + const purseCreationTerms = getPullCreditPurseCreationTerms(pullIni.amount); const mergeReserve = await wex.runWalletDbTx(async (tx) => tx.getReserve(pullIni.mergeReserveRowId), @@ -1006,12 +1129,12 @@ async function processPeerPullCreditCreatePurse( const sigRes = await wex.cryptoApi.signReservePurseCreate({ contractTermsHash: pullIni.contractTermsHash, - flags: WalletAccountMergeFlags.CreateWithPurseFee, + flags: purseCreationTerms.flags, mergePriv: pullIni.mergePriv, mergeTimestamp: TalerPreciseTimestamp.round(mergeTimestamp), purseAmount: pullIni.amount, purseExpiration: purseExpiration, - purseFee: purseFee, + purseFee: purseCreationTerms.purseFee, pursePriv: pullIni.pursePriv, pursePub: pullIni.pursePub, reservePayto, @@ -1025,7 +1148,6 @@ async function processPeerPullCreditCreatePurse( merge_pub: pullIni.mergePub, min_age: 0, purse_expiration: purseExpiration, - purse_fee: purseFee, purse_pub: pullIni.pursePub, purse_sig: sigRes.purseSig, purse_value: pullIni.amount, @@ -1074,7 +1196,8 @@ async function processPeerPullCreditCreatePurse( await ctx.failTransaction(pullIni.status, { code: resp.body.code }); return TaskRunResult.finished(); case HttpStatusCode.PaymentRequired: - throw Error(`unexpected reserve merge response ${resp.case}`); + await ctx.failTransaction(pullIni.status, resp.detail); + return TaskRunResult.finished(); case TalerErrorCode.EXCHANGE_RESERVES_PURSE_EXPIRATION_BEFORE_NOW: await beginPeerPullCreditExpirationCleanup( wex, @@ -1091,10 +1214,34 @@ async function processPeerPullCreditCreatePurse( if (!rec) { return; } - rec.status = PeerPullPaymentCreditStatus.PendingReady; - await h.update(rec, "create-purse"); + const nextStatus = statusAfterPullCreditCreate(rec.status); + rec.purseCreateProof = resp.body; + rec.status = nextStatus; + await h.update( + rec, + nextStatus === PeerPullPaymentCreditStatus.PendingReady + ? "create-purse" + : "late-create-purse-cleanup", + ); }); - return TaskRunResult.backoff(); + return TaskRunResult.progress(); +} + +export function statusAfterPullCreditCreate( + status: PeerPullPaymentCreditStatus, +): PeerPullPaymentCreditStatus { + switch (status) { + case PeerPullPaymentCreditStatus.PendingCreatePurse: + return PeerPullPaymentCreditStatus.PendingReady; + case PeerPullPaymentCreditStatus.Aborted: + return PeerPullPaymentCreditStatus.AbortingDeletePurse; + case PeerPullPaymentCreditStatus.Expired: + return PeerPullPaymentCreditStatus.ExpiredDeletePurse; + case PeerPullPaymentCreditStatus.FailedKycHardLimit: + return PeerPullPaymentCreditStatus.FinalizingKycHardLimit; + default: + return status; + } } export async function processPeerPullCredit( @@ -1278,7 +1425,7 @@ async function handlePeerPullCreditKycRequired( await ctx.wex.runWalletDbTx(async (tx) => { const [rec, h] = await ctx.getRecordHandle(tx); - if (!rec) { + if (rec?.status !== PeerPullPaymentCreditStatus.PendingCreatePurse) { return; } logger.info(`setting peer-pull-credit kyc payto hash to ${kycPaytoHash}`); @@ -1312,9 +1459,6 @@ export async function internalCheckPeerPullCredit( wex: WalletExecutionContext, req: CheckPeerPullCreditRequest, ): Promise<CheckPeerPullCreditResponse> { - // FIXME: We don't support exchanges with purse fees yet. - // Select an exchange where we have money in the specified currency - const instructedAmount = Amounts.parseOrThrow(req.amount); const currency = instructedAmount.currency; @@ -1443,7 +1587,7 @@ async function internalInitiatePeerPullPayment( } const mergeReserveInfo = await getMergeReserveInfo(wex, { - exchangeBaseUrl: exchangeBaseUrl, + exchangeBaseUrl, }); const pursePair = await wex.cryptoApi.createEddsaKeypair({}); @@ -1513,14 +1657,20 @@ async function internalInitiatePeerPullPayment( ); }); + const initialRecord = await wex.runWalletDbTx((tx) => + tx.getPeerPullCredit(pursePair.pub), + ); + checkDbInvariant(!!initialRecord, "new peer pull credit is missing"); + await processPeerPullCreditCreatePurse(wex, initialRecord); wex.taskScheduler.startShepherdTask(ctx.taskId); + const currentRecord = await wex.runWalletDbTx((tx) => + tx.getPeerPullCredit(pursePair.pub), + ); return { - talerUri: TalerUris.stringify({ - type: TalerUriAction.PayPull, - exchangeBaseUrl: exchangeBaseUrl as HostPortPath, // FIXME: change record type, - contractPriv: contractKeyPair.priv, - }), + talerUri: currentRecord?.purseCreateProof + ? makePeerPullPaymentUri(currentRecord) + : undefined, transactionId: ctx.transactionId, }; } diff --git a/packages/taler-wallet-core/src/pay-peer-push-credit.test.ts b/packages/taler-wallet-core/src/pay-peer-push-credit.test.ts @@ -22,10 +22,14 @@ import { import assert from "node:assert"; import { test } from "node:test"; import { PeerPushCreditStatus, WalletPeerPushCredit } from "./db/records.js"; +import { WalletDbTransaction } from "./db/transaction.js"; import { computePeerPushCreditTransactionActions, computePeerPushCreditTransactionState, + peerPushCreditStatusAfterNoMerge, + PeerPushCreditTransactionContext, } from "./pay-peer-push-credit.js"; +import { WalletExecutionContext } from "./wallet.js"; function makeRecord(status: PeerPushCreditStatus): WalletPeerPushCredit { return { @@ -70,3 +74,81 @@ test("push-credit hard-KYC cleanup has explicit finalizing and failed states", ( [TransactionAction.Retry], ); }); + +test("aborting a possibly in-flight push-credit merge keeps reconciling", async () => { + let record = makeRecord(PeerPushCreditStatus.PendingMerge); + const tx = { + async getExchange() { + return undefined; + }, + async listExchangeDetailsByBaseUrl() { + return []; + }, + async getPeerPushCredit(): Promise<WalletPeerPushCredit> { + return record; + }, + async upsertPeerPushCredit(updated: WalletPeerPushCredit): Promise<void> { + record = updated; + }, + async deletePeerPushCredit(): Promise<void> {}, + async upsertTransactionMeta(): Promise<void> {}, + async deleteTransactionMeta(): Promise<void> {}, + notify(): void {}, + } as unknown as WalletDbTransaction; + let reset = false; + let stopped = false; + const wex = { + async runWalletDbTx<T>( + f: (transaction: WalletDbTransaction) => Promise<T>, + ): Promise<T> { + return f(tx); + }, + taskScheduler: { + async resetTask(): Promise<void> { + reset = true; + }, + stopShepherdTask(): void { + stopped = true; + }, + }, + } as unknown as WalletExecutionContext; + + await new PeerPushCreditTransactionContext( + wex, + record.peerPushCreditId, + ).userAbortTransaction(); + + assert.strictEqual(record.status, PeerPushCreditStatus.AbortingMerge); + assert.strictEqual(reset, true); + assert.strictEqual(stopped, false); + assert.deepStrictEqual(computePeerPushCreditTransactionState(record), { + major: TransactionMajorState.Finalizing, + minor: TransactionMinorState.Merge, + working: true, + }); + assert.deepStrictEqual(computePeerPushCreditTransactionActions(record), [ + TransactionAction.Retry, + ]); +}); + +test("complete shared-reserve history terminates conclusive no-merge states", () => { + assert.strictEqual( + peerPushCreditStatusAfterNoMerge(PeerPushCreditStatus.AbortingMerge, false), + PeerPushCreditStatus.Aborted, + ); + assert.strictEqual( + peerPushCreditStatusAfterNoMerge( + PeerPushCreditStatus.FinalizingKycHardLimit, + false, + ), + PeerPushCreditStatus.FailedKycHardLimit, + ); + assert.strictEqual( + peerPushCreditStatusAfterNoMerge(PeerPushCreditStatus.PendingMerge, true), + PeerPushCreditStatus.Expired, + ); + assert.strictEqual( + peerPushCreditStatusAfterNoMerge(PeerPushCreditStatus.PendingMerge, false), + undefined, + ); +}); diff --git a/packages/taler-wallet-core/src/pay-peer-push-credit.ts b/packages/taler-wallet-core/src/pay-peer-push-credit.ts @@ -103,7 +103,9 @@ import { runKycCheckAlgo, } from "./kyc.js"; import { + getMergedPurseAmountFromReserveHistory, getMergeReserveInfo, + getReserveBalanceAfterCredit, isPurseGoneByExpiration, isPurseMerged, } from "./pay-peer-common.js"; @@ -124,7 +126,10 @@ import { waitWithdrawalFinal, } from "./withdraw.js"; import { WalletDbTransaction } from "./db/transaction.js"; -import { requireValidExchangePurseStatus } from "./exchange-signatures.js"; +import { + requireValidExchangePurseMergeConfirmation, + requireValidExchangePurseStatus, +} from "./exchange-signatures.js"; const logger = new Logger("pay-peer-push-credit.ts"); @@ -344,6 +349,7 @@ export class PeerPushCreditTransactionContext implements TransactionContext { case PeerPushCreditStatus.Aborted: case PeerPushCreditStatus.Expired: case PeerPushCreditStatus.FinalizingKycHardLimit: + case PeerPushCreditStatus.AbortingMerge: return; case PeerPushCreditStatus.PendingBalanceKycRequired: rec.status = PeerPushCreditStatus.SuspendedBalanceKycRequired; @@ -370,10 +376,10 @@ export class PeerPushCreditTransactionContext implements TransactionContext { } async userAbortTransaction(): Promise<void> { - await this.wex.runWalletDbTx(async (tx) => { + const shouldReconcile = await this.wex.runWalletDbTx(async (tx) => { const [rec, h] = await this.getRecordHandle(tx); if (!rec) { - return; + return false; } switch (rec.status) { case PeerPushCreditStatus.Failed: @@ -381,27 +387,37 @@ export class PeerPushCreditTransactionContext implements TransactionContext { case PeerPushCreditStatus.Aborted: case PeerPushCreditStatus.Done: case PeerPushCreditStatus.FinalizingKycHardLimit: - return; - case PeerPushCreditStatus.SuspendedMerge: + case PeerPushCreditStatus.PendingWithdrawing: + case PeerPushCreditStatus.SuspendedWithdrawing: + return false; + case PeerPushCreditStatus.AbortingMerge: + return true; case PeerPushCreditStatus.DialogProposed: + rec.status = PeerPushCreditStatus.Aborted; + break; + case PeerPushCreditStatus.SuspendedMerge: case PeerPushCreditStatus.SuspendedMergeKycRequired: - case PeerPushCreditStatus.SuspendedWithdrawing: case PeerPushCreditStatus.PendingBalanceKycRequired: case PeerPushCreditStatus.SuspendedBalanceKycRequired: - case PeerPushCreditStatus.PendingWithdrawing: case PeerPushCreditStatus.PendingMergeKycRequired: case PeerPushCreditStatus.PendingMerge: case PeerPushCreditStatus.PendingBalanceKycInit: case PeerPushCreditStatus.SuspendedBalanceKycInit: - case PeerPushCreditStatus.Expired: - rec.status = PeerPushCreditStatus.Aborted; + rec.status = PeerPushCreditStatus.AbortingMerge; break; + case PeerPushCreditStatus.Expired: + return false; default: assertUnreachable(rec.status); } await h.update(rec, "abort"); + return rec.status === PeerPushCreditStatus.AbortingMerge; }); - this.wex.taskScheduler.stopShepherdTask(this.taskId); + if (shouldReconcile) { + await this.wex.taskScheduler.resetTask(this.taskId); + } else { + this.wex.taskScheduler.stopShepherdTask(this.taskId); + } } async userResumeTransaction(): Promise<void> { @@ -423,6 +439,7 @@ export class PeerPushCreditTransactionContext implements TransactionContext { case PeerPushCreditStatus.FailedKycHardLimit: case PeerPushCreditStatus.Expired: case PeerPushCreditStatus.FinalizingKycHardLimit: + case PeerPushCreditStatus.AbortingMerge: return; case PeerPushCreditStatus.SuspendedMerge: rec.status = PeerPushCreditStatus.PendingMerge; @@ -493,6 +510,7 @@ export class PeerPushCreditTransactionContext implements TransactionContext { case PeerPushCreditStatus.SuspendedBalanceKycRequired: case PeerPushCreditStatus.PendingBalanceKycInit: case PeerPushCreditStatus.SuspendedBalanceKycInit: + case PeerPushCreditStatus.AbortingMerge: // The current state does not advertise Fail. return; default: @@ -762,6 +780,49 @@ async function internalPreparePeerPushCredit( }; } +/** + * Bind a push-credit operation to one merge reserve before any merge request + * is sent. Legacy records acquire the binding on their next retry. + */ +async function getBoundPeerPushCreditMergeReserve( + wex: WalletExecutionContext, + peerInc: WalletPeerPushCredit, +): Promise<WalletReserve> { + let reserveRowId = peerInc.mergeReserveRowId; + if (reserveRowId == null) { + const candidate = await getMergeReserveInfo(wex, { + exchangeBaseUrl: peerInc.exchangeBaseUrl, + }); + checkDbInvariant( + candidate.rowId != null, + "merge reserve for peer push credit has no row id", + ); + const ctx = new PeerPushCreditTransactionContext( + wex, + peerInc.peerPushCreditId, + ); + reserveRowId = await wex.runWalletDbTx(async (tx) => { + const [rec, h] = await ctx.getRecordHandle(tx); + checkDbInvariant( + !!rec, + "peer push credit disappeared while binding reserve", + ); + if (rec.mergeReserveRowId == null) { + rec.mergeReserveRowId = candidate.rowId; + await h.update(rec, "bind-merge-reserve"); + } + return rec.mergeReserveRowId; + }); + } + checkDbInvariant( + reserveRowId != null, + "peer push credit merge reserve binding is missing", + ); + const reserve = await wex.runWalletDbTx((tx) => tx.getReserve(reserveRowId)); + checkDbInvariant(!!reserve, "bound peer push credit reserve is missing"); + return reserve; +} + async function processPeerPushDebitMergeKyc( wex: WalletExecutionContext, peerInc: WalletPeerPushCredit, @@ -775,13 +836,19 @@ async function processPeerPushDebitMergeKyc( peerInc.peerPushCreditId, ); const { exchangeBaseUrl } = peerInc; - // FIXME: What if this changes? Should be part of the p2p record - const mergeReserveInfo = await getMergeReserveInfo(wex, { - exchangeBaseUrl, - }); + const mergeReserveInfo = await getBoundPeerPushCreditMergeReserve( + wex, + peerInc, + ); const accountPub = mergeReserveInfo.reservePub; const accountPriv = mergeReserveInfo.reservePriv; + const balanceAfterOperation = await getReserveBalanceAfterCredit( + wex, + exchangeBaseUrl, + mergeReserveInfo, + contractTerms.amount, + ); let myKycState: GenericKycStatusReq | undefined; @@ -791,6 +858,7 @@ async function processPeerPushDebitMergeKyc( accountPub, // FIXME: Is this the correct amount? amount: peerInc.estimatedAmountEffective, + balanceAfterOperation, exchangeBaseUrl, operation: "MERGE", paytoHash: peerInc.kycPaytoHash, @@ -847,6 +915,7 @@ async function processPeerPushDebitMergeKyc( async function processPeerPushCreditKycHardLimitRecovery( wex: WalletExecutionContext, peerInc: WalletPeerPushCredit, + contractTerms: PeerContractTerms, ): Promise<TaskRunResult> { const exchangeClient = walletExchangeClient(peerInc.exchangeBaseUrl, wex); const statusResp = await exchangeClient.getPurseStatusAtMerge( @@ -869,11 +938,19 @@ async function processPeerPushCreditKycHardLimitRecovery( } break; case HttpStatusCode.Gone: - case HttpStatusCode.NotFound: // Unmerged purse value is refunded to its depositing coins when the // purse expires or is deleted by its owner. nextStatus = PeerPushCreditStatus.FailedKycHardLimit; break; + case HttpStatusCode.NotFound: + // Purging removes the evidence needed to distinguish a committed merge + // from an unmerged purse. Recover the deterministic reserve first. + return continuePeerPushCreditAfterMerge( + wex, + peerInc, + undefined, + await getBoundPeerPushCreditMergeReserve(wex, peerInc), + ); default: assertUnreachable(statusResp); } @@ -915,6 +992,15 @@ async function transitionPeerPushCreditKycRequired( if (!peerInc) { return TaskRunResult.finished(); } + switch (peerInc.status) { + case PeerPushCreditStatus.PendingMerge: + case PeerPushCreditStatus.PendingMergeKycRequired: + break; + case PeerPushCreditStatus.AbortingMerge: + return TaskRunResult.progress(); + default: + return TaskRunResult.finished(); + } peerInc.kycPaytoHash = kycPending.h_payto; peerInc.status = PeerPushCreditStatus.PendingMergeKycRequired; peerInc.kycLastDeny = timestampPreciseToDb(TalerPreciseTimestamp.now()); @@ -926,9 +1012,46 @@ async function transitionPeerPushCreditKycRequired( async function continuePeerPushCreditAfterMerge( wex: WalletExecutionContext, peerInc: WalletPeerPushCredit, - amount: ReturnType<typeof Amounts.parseOrThrow>, + confirmedAmount: ReturnType<typeof Amounts.parseOrThrow> | undefined, mergeReserveInfo: WalletReserve, ): Promise<TaskRunResult> { + const contractRecord = await wex.runWalletDbTx((tx) => + tx.getContractTerms(peerInc.contractTermsHash), + ); + checkDbInvariant(!!contractRecord, "peer push contract terms are missing"); + const contractTerms = contractRecord.contractTermsRaw as PeerContractTerms; + const target = Amounts.parseOrThrow(contractTerms.amount); + let amount = confirmedAmount; + if (amount == null) { + const mergeTimestamp = AbsoluteTime.toProtocolTimestamp( + AbsoluteTime.fromPreciseTimestamp( + timestampPreciseFromDb(peerInc.timestamp), + ), + ); + const historyResult = await getMergedPurseAmountFromReserveHistory(wex, { + exchangeBaseUrl: peerInc.exchangeBaseUrl, + reserve: mergeReserveInfo, + pursePub: peerInc.pursePub, + contractTermsHash: peerInc.contractTermsHash, + mergePriv: peerInc.mergePriv, + mergeTimestamp, + purseExpiration: contractTerms.purse_expiration, + purseAmount: Amounts.stringify(target), + purseFee: Amounts.stringify(Amounts.zeroOfAmount(target)), + flags: WalletAccountMergeFlags.MergeFullyPaidPurse, + }); + if (historyResult.status === "not-merged") { + return finishPeerPushCreditWithoutMerge(wex, peerInc, contractTerms); + } + amount = Amounts.parseOrThrow(historyResult.amount); + } + if (Amounts.cmp(amount, target) < 0) { + throw TalerError.fromDetail( + TalerErrorCode.WALLET_TRANSACTION_PROTOCOL_VIOLATION, + {}, + "exchange confirmed a peer merge below the contract amount", + ); + } const ctx = new PeerPushCreditTransactionContext( wex, peerInc.peerPushCreditId, @@ -956,7 +1079,9 @@ async function continuePeerPushCreditAfterMerge( case PeerPushCreditStatus.PendingMerge: case PeerPushCreditStatus.PendingMergeKycRequired: case PeerPushCreditStatus.PendingBalanceKycRequired: - case PeerPushCreditStatus.PendingBalanceKycInit: { + case PeerPushCreditStatus.PendingBalanceKycInit: + case PeerPushCreditStatus.AbortingMerge: + case PeerPushCreditStatus.FinalizingKycHardLimit: { current.status = PeerPushCreditStatus.PendingWithdrawing; const wgCreateRes = await internalPerformCreateWithdrawalGroup( wex, @@ -975,6 +1100,56 @@ async function continuePeerPushCreditAfterMerge( return TaskRunResult.backoff(); } +async function finishPeerPushCreditWithoutMerge( + wex: WalletExecutionContext, + peerInc: WalletPeerPushCredit, + contractTerms: PeerContractTerms, +): Promise<TaskRunResult> { + const ctx = new PeerPushCreditTransactionContext( + wex, + peerInc.peerPushCreditId, + ); + let finished = false; + await wex.runWalletDbTx(async (tx) => { + const [rec, h] = await ctx.getRecordHandle(tx); + if (!rec) { + finished = true; + return; + } + const nextStatus = peerPushCreditStatusAfterNoMerge( + rec.status, + isPurseGoneByExpiration(contractTerms.purse_expiration), + ); + if (nextStatus == null) { + finished = !isActivePreMergePeerPushCreditStatus(rec.status); + return; + } + rec.status = nextStatus; + finished = true; + await h.update(rec, "reserve-history-no-merge"); + }); + return finished ? TaskRunResult.finished() : TaskRunResult.backoff(); +} + +export function peerPushCreditStatusAfterNoMerge( + status: PeerPushCreditStatus, + purseExpired: boolean, +): PeerPushCreditStatus | undefined { + switch (status) { + case PeerPushCreditStatus.AbortingMerge: + return PeerPushCreditStatus.Aborted; + case PeerPushCreditStatus.FinalizingKycHardLimit: + return PeerPushCreditStatus.FailedKycHardLimit; + case PeerPushCreditStatus.PendingMerge: + case PeerPushCreditStatus.PendingMergeKycRequired: + case PeerPushCreditStatus.PendingBalanceKycRequired: + case PeerPushCreditStatus.PendingBalanceKycInit: + return purseExpired ? PeerPushCreditStatus.Expired : undefined; + default: + return undefined; + } +} + function isActivePreMergePeerPushCreditStatus( status: PeerPushCreditStatus, ): boolean { @@ -983,6 +1158,7 @@ function isActivePreMergePeerPushCreditStatus( case PeerPushCreditStatus.PendingMergeKycRequired: case PeerPushCreditStatus.PendingBalanceKycRequired: case PeerPushCreditStatus.PendingBalanceKycInit: + case PeerPushCreditStatus.AbortingMerge: return true; default: return false; @@ -1012,10 +1188,8 @@ async function reconcileExpiredPeerPushCredit( return continuePeerPushCreditAfterMerge( wex, peerInc, - Amounts.parseOrThrow(contractTerms.amount), - await getMergeReserveInfo(wex, { - exchangeBaseUrl: peerInc.exchangeBaseUrl, - }), + Amounts.parseOrThrow(statusResp.body.balance), + await getBoundPeerPushCreditMergeReserve(wex, peerInc), ); case HttpStatusCode.Gone: { const ctx = new PeerPushCreditTransactionContext( @@ -1032,17 +1206,74 @@ async function reconcileExpiredPeerPushCredit( }); return TaskRunResult.finished(); } - case HttpStatusCode.NotFound: { - const ctx = new PeerPushCreditTransactionContext( + case HttpStatusCode.NotFound: + // The purse has already been purged, so reconcile through the reserve + // associated with our deterministic withdrawal instead of losing the + // incoming balance. + return continuePeerPushCreditAfterMerge( wex, - peerInc.peerPushCreditId, + peerInc, + undefined, + await getBoundPeerPushCreditMergeReserve(wex, peerInc), + ); + default: + assertUnreachable(statusResp); + } +} + +async function processAbortingPeerPushCredit( + wex: WalletExecutionContext, + peerInc: WalletPeerPushCredit, + contractTerms: PeerContractTerms, +): Promise<TaskRunResult> { + const statusResp = await walletExchangeClient( + peerInc.exchangeBaseUrl, + wex, + ).getPurseStatusAtMerge(peerInc.pursePub, true); + switch (statusResp.case) { + case "ok": + await requireValidExchangePurseStatus( + wex, + peerInc.exchangeBaseUrl, + statusResp.body, + ); + if (isPurseMerged(statusResp.body)) { + return continuePeerPushCreditAfterMerge( + wex, + peerInc, + Amounts.parseOrThrow(statusResp.body.balance), + await getBoundPeerPushCreditMergeReserve(wex, peerInc), + ); + } + break; + case HttpStatusCode.Gone: + break; + case HttpStatusCode.NotFound: + // A purged purse is ambiguous. Querying the merge reserve through the + // deterministic withdrawal is the only remaining recovery path. + return continuePeerPushCreditAfterMerge( + wex, + peerInc, + undefined, + await getBoundPeerPushCreditMergeReserve(wex, peerInc), ); - await ctx.failTransaction(peerInc.status, statusResp.detail); - return TaskRunResult.finished(); - } default: assertUnreachable(statusResp); } + + const ctx = new PeerPushCreditTransactionContext( + wex, + peerInc.peerPushCreditId, + ); + await wex.runWalletDbTx(async (tx) => { + const [rec, h] = await ctx.getRecordHandle(tx); + if (rec?.status !== PeerPushCreditStatus.AbortingMerge) { + return; + } + rec.status = PeerPushCreditStatus.Aborted; + await h.update(rec, "abort-merge-not-committed"); + }); + return TaskRunResult.finished(); } async function processPendingMerge( @@ -1090,10 +1321,10 @@ async function processPendingMerge( const amount = Amounts.parseOrThrow(contractTerms.amount); - // FIXME: What if this changes? Should be part of the p2p record - const mergeReserveInfo = await getMergeReserveInfo(wex, { - exchangeBaseUrl: peerInc.exchangeBaseUrl, - }); + const mergeReserveInfo = await getBoundPeerPushCreditMergeReserve( + wex, + peerInc, + ); const timestamp = timestampPreciseFromDb(peerInc.timestamp); @@ -1136,10 +1367,20 @@ async function processPendingMerge( peerInc.pursePub, mergeReq, ); + let confirmedAmount: ReturnType<typeof Amounts.parseOrThrow> | undefined; switch (mergeResp.case) { case "ok": logger.trace(`merge response: ${j2s(mergeResp.body)}`); + await requireValidExchangePurseMergeConfirmation(wex, { + exchangeBaseUrl: peerInc.exchangeBaseUrl, + pursePub: peerInc.pursePub, + reservePub: mergeReserveInfo.reservePub, + contractTermsHash: peerInc.contractTermsHash, + purseExpiration: contractTerms.purse_expiration, + response: mergeResp.body, + }); + confirmedAmount = Amounts.parseOrThrow(mergeResp.body.merge_amount); break; case HttpStatusCode.UnavailableForLegalReasons: { const kycLegiNeededResp = mergeResp.body; @@ -1189,6 +1430,9 @@ async function processPendingMerge( rec.status = PeerPushCreditStatus.Failed; break; } + case PeerPushCreditStatus.AbortingMerge: + rec.status = PeerPushCreditStatus.Aborted; + break; default: return; } @@ -1223,6 +1467,9 @@ async function processPendingMerge( } break; } + case PeerPushCreditStatus.AbortingMerge: + rec.status = PeerPushCreditStatus.Aborted; + break; default: return; } @@ -1230,8 +1477,14 @@ async function processPendingMerge( }); return TaskRunResult.finished(); } - case HttpStatusCode.Forbidden: case HttpStatusCode.NotFound: + return continuePeerPushCreditAfterMerge( + wex, + peerInc, + undefined, + mergeReserveInfo, + ); + case HttpStatusCode.Forbidden: await ctx.failTransaction(peerInc.status, mergeResp.detail); return TaskRunResult.finished(); default: @@ -1241,7 +1494,7 @@ async function processPendingMerge( return continuePeerPushCreditAfterMerge( wex, peerInc, - amount, + confirmedAmount, mergeReserveInfo, ); } @@ -1318,10 +1571,10 @@ async function processPendingWithdrawing( case WithdrawalGroupStatus.SuspendedRegisteringBank: case WithdrawalGroupStatus.SuspendedWaitConfirmBank: case WithdrawalGroupStatus.FinalizingKycHardLimit: - // The withdrawal's own task handles the KYC; the credit waits for it. case WithdrawalGroupStatus.PendingBalanceKyc: case WithdrawalGroupStatus.PendingBalanceKycInit: case WithdrawalGroupStatus.PendingKyc: + // The withdrawal's own task handles the KYC; the credit waits for it. return TaskRunResult.backoff(); case WithdrawalGroupStatus.DialogProposed: throw Error( @@ -1514,9 +1767,15 @@ export async function processPeerPushCredit( } return processPeerPushDebitMergeKyc(wex, peerInc, contractTerms); case PeerPushCreditStatus.FinalizingKycHardLimit: - return processPeerPushCreditKycHardLimitRecovery(wex, peerInc); + return processPeerPushCreditKycHardLimitRecovery( + wex, + peerInc, + contractTerms, + ); case PeerPushCreditStatus.PendingMerge: return processPendingMerge(wex, peerInc, contractTerms); + case PeerPushCreditStatus.AbortingMerge: + return processAbortingPeerPushCredit(wex, peerInc, contractTerms); case PeerPushCreditStatus.PendingWithdrawing: return processPendingWithdrawing(wex, peerInc); case PeerPushCreditStatus.PendingBalanceKycInit: @@ -1736,6 +1995,12 @@ export function computePeerPushCreditTransactionState( major: TransactionMajorState.Dialog, minor: TransactionMinorState.Proposed, }; + case PeerPushCreditStatus.AbortingMerge: + return { + major: TransactionMajorState.Finalizing, + minor: TransactionMinorState.Merge, + working: true, + }; case PeerPushCreditStatus.PendingMerge: return { major: TransactionMajorState.Pending, @@ -1852,6 +2117,8 @@ export function computePeerPushCreditTransactionActions( return [TransactionAction.Delete]; case PeerPushCreditStatus.FinalizingKycHardLimit: return [TransactionAction.Retry]; + case PeerPushCreditStatus.AbortingMerge: + return [TransactionAction.Retry]; case PeerPushCreditStatus.FailedKycHardLimit: return [TransactionAction.Delete]; case PeerPushCreditStatus.PendingMergeKycRequired: