diff --git a/dash-spv-ffi/src/callbacks.rs b/dash-spv-ffi/src/callbacks.rs index daa30940c..8f3a5df5c 100644 --- a/dash-spv-ffi/src/callbacks.rs +++ b/dash-spv-ffi/src/callbacks.rs @@ -726,7 +726,10 @@ impl Drop for FFIDerivedAddress { /// /// Fires when a wallet-relevant transaction is first seen off-chain — either /// in the mempool, or directly via an InstantSend lock (in that case the -/// record's `context` is `InstantSend(..)`). +/// record's `context` is `InstantSend(..)`). Also fires when an off-chain +/// funding arrival corrects an earlier transaction's accounting. Consumers must +/// upsert by wallet, account and txid; a correction retains the spender's context. +/// A plain InstantSend lock on a known mempool transaction emits only the lock callback. /// /// All pointer parameters are borrowed and only valid for the duration of the /// callback. `balance` is the wallet's balance *after* the transaction was diff --git a/key-wallet-manager/src/event_tests.rs b/key-wallet-manager/src/event_tests.rs index afccc35ae..f8f6c0324 100644 --- a/key-wallet-manager/src/event_tests.rs +++ b/key-wallet-manager/src/event_tests.rs @@ -2119,3 +2119,157 @@ async fn dropped_persistence_consumer_does_not_wedge_emission() { "broadcast delivery must be unaffected by a lost persistence consumer" ); } + +#[tokio::test] +async fn late_funding_mempool_publishes_complete_correction() { + use key_wallet::managed_account::transaction_record::{OutputRole, TransactionDirection}; + + for locked_funding in [false, true] { + let (mut manager, wallet_id, addr) = setup_manager_with_wallet(); + let mut funding = create_tx_paying_to(&addr, 0xa1); + funding.output.push(funding.output[0].clone()); + let earlier_funding = create_tx_paying_to(&addr, 0xa2); + let mut spender = create_tx_paying_to(&addr, 0xa3); + spender.input = [ + OutPoint::new(funding.txid(), 1), + OutPoint::new(funding.txid(), 0), + OutPoint::new(earlier_funding.txid(), 0), + ] + .into_iter() + .map(|previous_output| TxIn { + previous_output, + ..spender.input[0].clone() + }) + .collect(); + spender.output[0].value = 3 * TX_AMOUNT - 2000; + spender.output.push(TxOut { + value: 1000, + script_pubkey: Address::dummy(Network::Testnet, 999).script_pubkey(), + }); + manager.process_mempool_transaction(&spender, None).await; + manager.process_mempool_transaction(&earlier_funding, None).await; + let account = manager.wallet_infos[&wallet_id].first_bip44_managed_account().unwrap(); + assert_eq!(account.transactions()[&spender.txid()].input_details.len(), 1); + + let mut rx = manager.subscribe_events(); + let lock = locked_funding.then(|| dummy_instant_lock(funding.txid())); + manager.process_mempool_transaction(&funding, lock.clone()).await; + let events = drain_events(&mut rx); + let corrections: Vec<_> = events + .iter() + .filter_map(|event| match event { + WalletEvent::TransactionDetected { + wallet_id: id, + record, + .. + } if *id == wallet_id && record.txid == spender.txid() => Some(record), + _ => None, + }) + .collect(); + assert_eq!(corrections.len(), 1, "publish one complete correction"); + let corrected = corrections[0]; + assert_eq!(corrected.net_amount, -2000); + assert_eq!(corrected.direction, TransactionDirection::Outgoing); + assert_eq!(corrected.context, TransactionContext::Mempool); + assert_eq!( + corrected.input_details.iter().map(|d| (d.index, d.value)).collect::>(), + vec![(0, TX_AMOUNT), (1, TX_AMOUNT), (2, TX_AMOUNT)] + ); + assert_eq!(corrected.output_details.len(), 2); + assert_eq!(corrected.output_details[1].role, OutputRole::Sent); + assert!( + !events.iter().any(|event| matches!(event, + WalletEvent::TransactionInstantLocked { txid, .. } if *txid == spender.txid() + )), + "the funding lock belongs only to the funding transaction" + ); + manager.process_mempool_transaction(&funding, lock).await; + assert_no_events(&mut rx); + + // Funding confirmation must not attribute the same inputs again. + let block = make_block(vec![funding, earlier_funding], 0xa4, 100); + manager + .process_block_for_wallets(&block, block.block_hash(), 1, &BTreeSet::from([wallet_id])) + .await; + let account = manager.wallet_infos[&wallet_id].first_bip44_managed_account().unwrap(); + let stored = &account.transactions()[&spender.txid()]; + assert_eq!(stored.net_amount, -2000); + assert_eq!(stored.input_details.len(), 3); + assert_eq!(account.utxos.len(), 1); + for input in &spender.input { + assert!(!account.utxos.contains_key(&input.previous_output)); + } + assert!(drain_events(&mut rx).iter().all(|event| !matches!(event, + WalletEvent::BlockProcessed { updated, .. } + if updated.iter().any(|r| r.txid == spender.txid()) + ))); + + manager + .process_mempool_transaction(&spender, Some(dummy_instant_lock(spender.txid()))) + .await; + let events = drain_events(&mut rx); + assert_eq!(events.len(), 1, "a lock alone must not redetect the spender"); + assert!(matches!(&events[0], WalletEvent::TransactionInstantLocked { txid, .. } + if *txid == spender.txid())); + } +} + +#[tokio::test] +async fn late_funding_block_publishes_spender_correction() { + use key_wallet::managed_account::transaction_record::{OutputRole, TransactionDirection}; + + for change_value in [0, TX_AMOUNT - 2000] { + let (mut manager, wallet_id, addr) = setup_manager_with_wallet(); + let change = manager + .wallet_infos + .get_mut(&wallet_id) + .unwrap() + .first_bip44_managed_account_mut() + .unwrap() + .next_change_address(None, true) + .unwrap(); + let funding = create_tx_paying_to(&addr, 0xd5); + let mut spender = create_tx_paying_to(&change, 0xd6); + spender.input[0].previous_output = OutPoint::new(funding.txid(), 0); + spender.output[0].value = change_value; + let spend_block = make_block(vec![spender.clone()], 0xd7, 200); + let wallets = BTreeSet::from([wallet_id]); + manager + .process_block_for_wallets(&spend_block, spend_block.block_hash(), 2, &wallets) + .await; + let mut rx = manager.subscribe_events(); + let fund_block = make_block(vec![funding.clone()], 0xd8, 100); + let result = manager + .process_block_for_wallets(&fund_block, fund_block.block_hash(), 1, &wallets) + .await; + assert!(result.reapply_heights.is_empty()); + let events = drain_events(&mut rx); + let updated = events + .iter() + .find_map(|event| match event { + WalletEvent::BlockProcessed { + updated, + .. + } => Some(updated), + _ => None, + }) + .expect("block correction event"); + assert_eq!(updated.len(), 1); + let record = &updated[0]; + assert_eq!(record.txid, spender.txid()); + assert_eq!(record.net_amount, change_value as i64 - TX_AMOUNT as i64); + assert_eq!(record.input_details.len(), 1); + assert_eq!(record.input_details[0].value, TX_AMOUNT); + assert_eq!(record.input_details[0].index, 0); + assert_eq!(record.direction, TransactionDirection::Internal); + assert_eq!(record.output_details[0].role, OutputRole::Change); + assert_eq!(record.context.block_info().unwrap().height(), 2); + assert_eq!(record.context.block_info().unwrap().block_hash(), spend_block.block_hash()); + let account = manager.wallet_infos[&wallet_id].first_bip44_managed_account().unwrap(); + assert_eq!(account.transactions()[&spender.txid()].net_amount, record.net_amount); + assert!(!account.utxos.contains_key(&OutPoint::new(funding.txid(), 0))); + + manager.process_block_for_wallets(&fund_block, fund_block.block_hash(), 1, &wallets).await; + assert_no_events(&mut rx); + } +} diff --git a/key-wallet-manager/src/events.rs b/key-wallet-manager/src/events.rs index e044f0796..4613abb4a 100644 --- a/key-wallet-manager/src/events.rs +++ b/key-wallet-manager/src/events.rs @@ -180,9 +180,9 @@ fn format_account_balances(map: &BTreeMap) -> St /// consumers can persist the record(s) and balance atomically. #[derive(Debug, Clone)] pub enum WalletEvent { - /// First time the wallet sees an off-chain wallet-relevant transaction - /// (mempool, or directly via an InstantSend lock — in that case - /// `record.context` is `InstantSend(..)`). + /// An off-chain transaction was detected or its accounting details were corrected. + /// Corrections retain the transaction's own confirmation context; consumers upsert by account and txid. + /// Lock-only updates use `TransactionInstantLocked`. TransactionDetected { /// ID of the affected wallet. wallet_id: WalletId, diff --git a/key-wallet-manager/src/process_block.rs b/key-wallet-manager/src/process_block.rs index 5d9b61092..742b92eab 100644 --- a/key-wallet-manager/src/process_block.rs +++ b/key-wallet-manager/src/process_block.rs @@ -272,26 +272,33 @@ impl WalletInterface for WalletM per_wallet_released ); - if let Some(lock) = instant_lock { - for (wallet_id, records) in per_wallet_updated_records { - if records.is_empty() { - continue; + for (wallet_id, records) in per_wallet_updated_records { + let Some(info) = self.wallet_infos.get(&wallet_id) else { + continue; + }; + let balance = info.balance(); + let account_balances = + per_wallet_account_diff.get(&wallet_id).cloned().unwrap_or_default(); + for record in records { + let txid = record.txid; + // The arriving tx only changes lock status; other txids carry late-input corrections. + if txid != tx.txid() { + self.emit_event(WalletEvent::TransactionDetected { + wallet_id, + record: Box::new(record), + balance, + account_balances: account_balances.clone(), + addresses_derived: Vec::new(), + }); } - let Some(info) = self.wallet_infos.get(&wallet_id) else { - continue; - }; - let balance = info.balance(); - let account_balances = - per_wallet_account_diff.get(&wallet_id).cloned().unwrap_or_default(); - for record in records { - let event = WalletEvent::TransactionInstantLocked { + if let Some(lock) = instant_lock.as_ref().filter(|lock| lock.txid == txid) { + self.emit_event(WalletEvent::TransactionInstantLocked { wallet_id, - txid: record.txid, + txid, instant_lock: lock.clone(), balance, account_balances: account_balances.clone(), - }; - self.emit_event(event); + }); } } } @@ -695,30 +702,114 @@ mod tests { #[tokio::test] async fn test_funding_after_its_spend_asks_to_reapply_the_spend_block() { - let (mut manager, wallet_id, addr) = setup_manager_with_wallet(); - let funding = create_tx_paying_to(&addr, 0xaa); - let spend = spend_first_output_of(&funding); - let wallets = BTreeSet::from([wallet_id]); - - let mut spend_block = make_block(vec![spend]); - spend_block.header.nonce = 1; - let funding_block = make_block(vec![funding]); - - manager - .process_block_for_wallets(&spend_block, spend_block.block_hash(), 200, &wallets) - .await; - let result = manager - .process_block_for_wallets(&funding_block, funding_block.block_hash(), 100, &wallets) - .await; - assert_eq!(result.reapply_heights, BTreeMap::from([(wallet_id, BTreeSet::from([200]))])); - - manager - .process_block_for_wallets(&spend_block, spend_block.block_hash(), 200, &wallets) - .await; - let again = manager - .process_block_for_wallets(&funding_block, funding_block.block_hash(), 100, &wallets) - .await; - assert!(again.reapply_heights.is_empty()); + use key_wallet::managed_account::transaction_record::TransactionDirection; + + for (sibling, finalized) in [(false, false), (true, false), (true, true)] { + let (mut manager, wallet_id, addr) = setup_manager_with_wallet(); + let funding = create_tx_paying_to(&addr, 0xaa); + let mut spend = spend_first_output_of(&funding); + spend.output[0].value = TX_AMOUNT - 2000; + if sibling { + spend.output[0].script_pubkey = coinjoin_account(&manager, &wallet_id) + .managed_account_type() + .address_pools()[0] + .address_at_index(0) + .unwrap() + .script_pubkey(); + } + let spent_outpoint = OutPoint::new(funding.txid(), 0); + let spender_txid = spend.txid(); + let wallets = BTreeSet::from([wallet_id]); + let mut spend_block = make_block(vec![spend]); + spend_block.header.nonce = 1; + let funding_block = make_block(vec![funding]); + + manager + .process_block_for_wallets(&spend_block, spend_block.block_hash(), 200, &wallets) + .await; + assert_eq!( + coinjoin_account(&manager, &wallet_id).has_transaction(&spender_txid), + sibling + ); + if finalized { + manager.apply_chain_lock(ChainLock::dummy(200)); + } + let result = manager + .process_block_for_wallets( + &funding_block, + funding_block.block_hash(), + 100, + &wallets, + ) + .await; + assert_eq!( + result.reapply_heights, + BTreeMap::from([(wallet_id, BTreeSet::from([200]))]) + ); + let account = manager.wallet_infos[&wallet_id].first_bip44_managed_account().unwrap(); + assert!(!account.transactions().contains_key(&spender_txid)); + assert!(!account.utxos.contains_key(&spent_outpoint)); + + let mut rx = manager.subscribe_events(); + for (wallet, heights) in result.reapply_heights { + for height in heights { + let replay = manager + .process_block_for_wallets( + &spend_block, + spend_block.block_hash(), + height, + &BTreeSet::from([wallet]), + ) + .await; + assert!(replay.reapply_heights.is_empty()); + } + } + let events = drain_events(&mut rx); + let recorded = events + .iter() + .find_map(|event| match event { + WalletEvent::BlockProcessed { + inserted, + .. + } => inserted.iter().find(|r| { + r.txid == spender_txid + && matches!(r.account_type, AccountType::Standard { .. }) + }), + _ => None, + }) + .expect("replay publishes the missing funding-account record"); + assert_eq!(recorded.net_amount, -(TX_AMOUNT as i64)); + assert_eq!(recorded.direction, TransactionDirection::Outgoing); + assert_eq!(recorded.input_details.len(), 1); + assert_eq!(recorded.input_details[0].value, TX_AMOUNT); + assert_eq!(recorded.context.block_info().unwrap().height(), 200); + assert_eq!(recorded.context.is_chain_locked(), finalized); + let account = manager.wallet_infos[&wallet_id].first_bip44_managed_account().unwrap(); + assert!(!account.utxos.contains_key(&spent_outpoint)); + assert_eq!( + account.transactions().contains_key(&spender_txid), + !finalized || cfg!(feature = "keep-finalized-transactions") + ); + if sibling { + assert_eq!(manager.wallet_infos[&wallet_id].balance().total(), TX_AMOUNT - 2000); + if let Some(record) = + coinjoin_account(&manager, &wallet_id).transactions().get(&spender_txid) + { + assert_eq!(record.net_amount, (TX_AMOUNT - 2000) as i64); + assert!(record.input_details.is_empty()); + } + } + let again = manager + .process_block_for_wallets( + &funding_block, + funding_block.block_hash(), + 100, + &wallets, + ) + .await; + assert!(again.reapply_heights.is_empty()); + assert_no_events(&mut rx); + } } #[tokio::test] diff --git a/key-wallet/src/managed_account/managed_core_funds_account.rs b/key-wallet/src/managed_account/managed_core_funds_account.rs index e075f977c..93bfa8e9e 100644 --- a/key-wallet/src/managed_account/managed_core_funds_account.rs +++ b/key-wallet/src/managed_account/managed_core_funds_account.rs @@ -20,9 +20,7 @@ use crate::managed_account::managed_account_trait::ManagedAccountTrait; use crate::managed_account::managed_account_type::ManagedAccountType; use crate::managed_account::managed_core_keys_account::ManagedCoreKeysAccount; use crate::managed_account::reservation::{ReservationSet, ReservationToken}; -use crate::managed_account::transaction_record::{ - InputDetail, OutputDetail, OutputRole, TransactionDirection, -}; +use crate::managed_account::transaction_record::{InputDetail, OutputDetail, OutputRole}; use crate::transaction_checking::transaction_router::TransactionType; use crate::transaction_checking::{AccountMatch, TransactionContext}; use crate::utxo::Utxo; @@ -190,7 +188,7 @@ impl ManagedCoreFundsAccount { } /// Check if an outpoint was spent by a previously recorded transaction. - fn is_outpoint_spent(&self, outpoint: &OutPoint) -> bool { + pub(crate) fn is_outpoint_spent(&self, outpoint: &OutPoint) -> bool { self.spent_outpoints.contains(outpoint) } @@ -889,20 +887,8 @@ impl ManagedCoreFundsAccount { }); } - // Determine direction - let has_sent = output_details.iter().any(|d| d.role == OutputRole::Sent); - let has_our_outputs = output_details - .iter() - .any(|d| d.role == OutputRole::Received || d.role == OutputRole::Change); - let direction = if transaction_type == TransactionType::CoinJoin { - TransactionDirection::CoinJoin - } else if !has_sent && has_inputs && has_our_outputs { - TransactionDirection::Internal - } else if has_inputs { - TransactionDirection::Outgoing - } else { - TransactionDirection::Incoming - }; + let direction = + TransactionRecord::direction_for(transaction_type, has_inputs, &output_details); let tx_record = TransactionRecord::new( tx.clone(), @@ -1376,6 +1362,7 @@ mod conflict_sweep_walk_tests { use super::*; use crate::account::AccountType; use crate::account::StandardAccountType; + use crate::managed_account::transaction_record::TransactionDirection; use crate::transaction_checking::BlockInfo; use dashcore::ephemerealdata::instant_lock::InstantLock; use dashcore::hashes::Hash; diff --git a/key-wallet/src/managed_account/transaction_record.rs b/key-wallet/src/managed_account/transaction_record.rs index b51aee6f1..ff200f6b6 100644 --- a/key-wallet/src/managed_account/transaction_record.rs +++ b/key-wallet/src/managed_account/transaction_record.rs @@ -131,6 +131,44 @@ impl TransactionRecord { } } + /// Derive account-local flow from attributed inputs and owned outputs. + pub(crate) fn recompute_net_and_direction(&mut self) { + let owned: i64 = self + .output_details + .iter() + .filter(|o| matches!(o.role, OutputRole::Received | OutputRole::Change)) + .map(|o| o.value as i64) + .sum(); + let spent: i64 = self.input_details.iter().map(|i| i.value as i64).sum(); + self.net_amount = owned - spent; + self.direction = Self::direction_for( + self.transaction_type, + !self.input_details.is_empty(), + &self.output_details, + ); + } + + /// Classify account-local flow consistently for initial records and late-input corrections. + pub(crate) fn direction_for( + transaction_type: TransactionType, + has_inputs: bool, + output_details: &[OutputDetail], + ) -> TransactionDirection { + let has_sent = output_details.iter().any(|d| d.role == OutputRole::Sent); + let has_our_outputs = output_details + .iter() + .any(|d| matches!(d.role, OutputRole::Received | OutputRole::Change)); + if transaction_type == TransactionType::CoinJoin { + TransactionDirection::CoinJoin + } else if !has_sent && has_inputs && has_our_outputs { + TransactionDirection::Internal + } else if has_inputs { + TransactionDirection::Outgoing + } else { + TransactionDirection::Incoming + } + } + /// Calculate the number of confirmations based on current chain height pub fn confirmations(&self, current_height: u32) -> u32 { match self.context.block_info() { diff --git a/key-wallet/src/transaction_checking/wallet_checker.rs b/key-wallet/src/transaction_checking/wallet_checker.rs index 81f1adf1b..e5ab647db 100644 --- a/key-wallet/src/transaction_checking/wallet_checker.rs +++ b/key-wallet/src/transaction_checking/wallet_checker.rs @@ -6,12 +6,14 @@ pub(crate) use super::account_checker::TransactionCheckResult; use super::transaction_context::TransactionContext; use super::transaction_router::{AccountTypeToCheck, TransactionRouter}; +use crate::managed_account::managed_account_trait::ManagedAccountTrait; +use crate::managed_account::transaction_record::{InputDetail, OutputDetail, OutputRole}; use crate::wallet::managed_wallet_info::wallet_info_interface::WalletInfoInterface; use crate::wallet::managed_wallet_info::ManagedWalletInfo; use crate::{KeySource, Wallet}; use async_trait::async_trait; use dashcore::blockdata::transaction::Transaction; -use dashcore::{Amount, SignedAmount}; +use dashcore::{Address, Amount, OutPoint, SignedAmount}; /// Extension trait for ManagedWalletInfo to add transaction checking capabilities #[async_trait] @@ -29,6 +31,11 @@ pub trait WalletTransactionChecker { /// Callers that batch multiple transactions (e.g. block processing) can pass `false` /// and refresh once at the end via `update_last_processed_height`. /// + /// Late funding corrects accounting only while the spender's full record is retained. + /// By default, ChainLocks prune funds records to txids; later funding cannot correct + /// those records or previously emitted copies. Enable `keep-finalized-transactions` + /// before processing for corrections after finalization, at the cost of retaining history. + /// /// The context parameter indicates where the transaction comes from (mempool, block, etc.) /// async fn check_core_transaction( @@ -42,6 +49,63 @@ pub trait WalletTransactionChecker { } impl ManagedWalletInfo { + /// Correct late inputs in existing records; missing account slices use block replay. + fn attribute_late_inputs(&mut self, tx: &Transaction, result: &mut TransactionCheckResult) { + let txid = tx.txid(); + for mut account in self.accounts.all_accounts_mut() { + let Some(funds) = account.as_funds_mut() else { + continue; + }; + for (vout, output) in tx.output.iter().enumerate() { + let outpoint = OutPoint::new(txid, vout as u32); + if !funds.is_outpoint_spent(&outpoint) { + continue; + } + let Ok(address) = Address::from_script(&output.script_pubkey, self.network) else { + continue; + }; + if !funds.contains_address(&address) { + continue; + } + for record in funds.transactions_mut().values_mut() { + let Some(index) = + record.transaction.input.iter().position(|i| i.previous_output == outpoint) + else { + continue; + }; + if record.input_details.iter().any(|d| d.index == index as u32) { + continue; + } + record.input_details.push(InputDetail { + index: index as u32, + value: output.value, + address: address.clone(), + }); + record.input_details.sort_by_key(|d| d.index); + // Records without known inputs omit foreign outputs. + for (index, output) in record.transaction.output.iter().enumerate() { + if record.output_details.iter().all(|d| d.index != index as u32) { + record.output_details.push(OutputDetail { + index: index as u32, + role: OutputRole::Sent, + address: Address::from_script(&output.script_pubkey, self.network) + .ok(), + value: output.value, + }); + } + } + record.output_details.sort_by_key(|d| d.index); + record.recompute_net_and_direction(); + result + .updated_records + .retain(|r| r.txid != record.txid || r.account_type != record.account_type); + result.updated_records.push(record.clone()); + result.state_modified = true; + } + } + } + } + /// Promote records whose first sighting already consumed the live UTXO /// evidence used by transaction relevance checks. /// @@ -329,6 +393,8 @@ impl WalletTransactionChecker for ManagedWalletInfo { } } + self.attribute_late_inputs(tx, &mut result); + if is_new { // Populate dedup sets when a tx arrives with an initial IS status if context.is_instant_send() {