diff --git a/CHANGELOG.md b/CHANGELOG.md index 62c73a98cf2..ee85f2183c2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,11 +1,3 @@ -## Unreleased - -### Fixed - -- **platform-wallet:** Allow reconciliation of already-loaded asset locks with atomic tracked-lock persistence without requiring unrelated wallet restore support; retain nonterminal recovery state and typed consumption errors. - -- **platform-wallet-storage:** Restore confirmed Core spend and finality state on SQLite load so old funding transactions cannot make already-spent outputs selectable again. - ## [4.2.0-beta.7](https://github.com/dashpay/platform/compare/v4.2.0-beta.6...v4.2.0-beta.7) (2026-09-29) diff --git a/Cargo.lock b/Cargo.lock index a593de4f7f5..bdb7d2c4477 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5350,7 +5350,6 @@ dependencies = [ "chacha20poly1305", "chrono", "clap", - "dash-async", "dash-sdk", "dashcore", "dbus-secret-service-keyring-store", diff --git a/packages/rs-platform-wallet-storage/Cargo.toml b/packages/rs-platform-wallet-storage/Cargo.toml index fe1ed8c06f6..17dcc10136e 100644 --- a/packages/rs-platform-wallet-storage/Cargo.toml +++ b/packages/rs-platform-wallet-storage/Cargo.toml @@ -41,7 +41,6 @@ platform-wallet = { path = "../rs-platform-wallet", features = [ "eddsa", ], optional = true } serde = { version = "1", features = ["derive"], optional = true } -dash-async = { path = "../rs-dash-async", optional = true } key-wallet = { workspace = true, optional = true } dashcore = { workspace = true, optional = true } dpp = { path = "../rs-dpp", optional = true } @@ -225,7 +224,6 @@ sqlite = [ "dep:platform-wallet", "dep:serde", "dep:key-wallet", - "dep:dash-async", "dep:dashcore", "dep:dpp", "dep:dash-sdk", diff --git a/packages/rs-platform-wallet-storage/migrations/V019__core_transaction_accounting.rs b/packages/rs-platform-wallet-storage/migrations/V019__core_transaction_accounting.rs index b92003e8545..d563fd867f0 100644 --- a/packages/rs-platform-wallet-storage/migrations/V019__core_transaction_accounting.rs +++ b/packages/rs-platform-wallet-storage/migrations/V019__core_transaction_accounting.rs @@ -1,4 +1,5 @@ -//! Index raw inputs so late output ownership can repair the spending history. +//! Index raw inputs so late output ownership can repair the spending history, +//! and keep each record's pre-repair blob so every repair stays reversible. pub fn migration() -> String { "CREATE TABLE core_transaction_inputs ( @@ -9,6 +10,13 @@ pub fn migration() -> String { FOREIGN KEY (wallet_id, txid) REFERENCES core_transactions(wallet_id, txid) ON DELETE CASCADE ); CREATE INDEX idx_core_transaction_inputs_outpoint - ON core_transaction_inputs(wallet_id, outpoint);" + ON core_transaction_inputs(wallet_id, outpoint); + CREATE TABLE core_transaction_record_originals ( + wallet_id BLOB NOT NULL, + txid BLOB NOT NULL, + record_blob BLOB NOT NULL, + PRIMARY KEY (wallet_id, txid), + FOREIGN KEY (wallet_id) REFERENCES wallets(wallet_id) ON DELETE CASCADE + );" .to_owned() } diff --git a/packages/rs-platform-wallet-storage/src/sqlite/error.rs b/packages/rs-platform-wallet-storage/src/sqlite/error.rs index 2a7a30414b9..9f843edf7bb 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/error.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/error.rs @@ -34,10 +34,6 @@ pub enum AutoBackupOperation { /// Errors produced by the wallet-storage SQLite backend. #[derive(Debug, thiserror::Error)] pub enum WalletStorageError { - /// Confirmed Core history could not be replayed into the restored wallet. - #[error("could not restore confirmed Core history: {0}")] - CoreHistoryReplay(#[source] dash_async::AsyncError), - /// File-system I/O error reaching the database or backup files. #[error("io error")] Io(#[from] std::io::Error), @@ -710,8 +706,7 @@ impl WalletStorageError { // `ToSqlConversionFailure`, `InvalidColumnIndex`) — is a // logic bug, not a contention failure. Self::Sqlite(_) => false, - Self::CoreHistoryReplay(_) - | Self::Io(_) + Self::Io(_) | Self::Migration(_) | Self::IntegrityCheckFailed { .. } | Self::IntegrityCheckRunFailed { .. } @@ -876,7 +871,6 @@ impl WalletStorageError { | Self::UnownedIdentityHasRegistrationIndex { .. } | Self::EmptyUtxoScript { .. } | Self::EmptyPoolAddressScript { .. } - | Self::CoreHistoryReplay(_) | Self::DatabasePathIsSymlink { .. } => PersistenceErrorKind::Fatal, } } @@ -897,7 +891,6 @@ impl WalletStorageError { }, Self::Sqlite(_) => "sqlite_other", Self::FlushRetryable { .. } => "flush_retryable", - Self::CoreHistoryReplay(_) => "core_history_replay", Self::Io(_) => "io", Self::Migration(_) => "migration", Self::IntegrityCheckFailed { .. } => "integrity_check_failed", diff --git a/packages/rs-platform-wallet-storage/src/sqlite/migrations.rs b/packages/rs-platform-wallet-storage/src/sqlite/migrations.rs index e152a68a97a..ddd18caad25 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/migrations.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/migrations.rs @@ -10,6 +10,7 @@ use crate::sqlite::error::WalletStorageError; use refinery_core::error::WrapMigrationError; mod legacy_v008; +mod legacy_v019; // Generates a `migrations` module with `runner()`; path is relative to // the crate root. @@ -78,7 +79,7 @@ impl refinery_core::traits::sync::Transaction for MigrationTransaction<'_> { } else if query == self.pool_sql { legacy_v008::convert_pools(&self.tx)?; } else if query == self.history_sql { - super::schema::core_history::migrate(&self.tx)?; + legacy_v019::repair_history(&self.tx)?; } count += 1; } diff --git a/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs b/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs new file mode 100644 index 00000000000..e1bf7fcd5dc --- /dev/null +++ b/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs @@ -0,0 +1,569 @@ +//! Frozen V019 history repair. It backfills the raw-input index and corrects +//! stored accounting for databases migrating from V018 or earlier. +//! +//! Frozen on purpose: the live per-round repair in `schema::core_history` may +//! evolve, but V019 must keep doing exactly what it did when it shipped. Only +//! the low-level blob codec is shared; `TransactionRecord`'s encoding is owned +//! upstream (key-wallet) and cannot be frozen here. Do not edit. + +use std::collections::BTreeMap; + +use dashcore::hashes::Hash; +use dashcore::{Address, OutPoint, ScriptBuf, Txid}; +use key_wallet::managed_account::transaction_record::{ + InputDetail, OutputDetail, OutputRole, TransactionDirection, TransactionRecord, +}; +use key_wallet::transaction_checking::{TransactionContext, TransactionType}; +use platform_wallet::wallet::platform_wallet::WalletId; +use rusqlite::{params, OptionalExtension, Transaction}; + +use crate::sqlite::error::WalletStorageError; +use crate::sqlite::schema::{blob, id32, wallets}; +use crate::sqlite::util::safe_cast::i64_to_u64; + +/// Read a stored record; undecodable bytes are an error the caller classifies. +fn read_record( + tx: &Transaction<'_>, + wallet_id: &WalletId, + txid: &Txid, +) -> Result, WalletStorageError> { + let stored: Option)>> = tx + .query_row( + "SELECT length(record_blob), record_blob FROM core_transactions \ + WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + |row| { + Ok(match row.get::<_, Option>(0)? { + Some(len) => Some((len, row.get(1)?)), + None => None, + }) + }, + ) + .optional()?; + let Some(Some((len, payload))) = stored else { + return Ok(None); + }; + blob::check_size(len)?; + let record: TransactionRecord = blob::decode(&payload)?; + if record.txid != *txid { + return Err(WalletStorageError::blob_decode( + "transaction record names another transaction", + )); + } + Ok(Some(record)) +} + +/// Whether `error` reports stored bytes that cannot be decoded, not a database failure. +fn is_unreadable(error: &WalletStorageError) -> bool { + matches!( + error, + WalletStorageError::BincodeDecode { .. } + | WalletStorageError::BlobDecode { .. } + | WalletStorageError::BlobTooLarge { .. } + | WalletStorageError::HashDecode { .. } + | WalletStorageError::IntegerOverflow { .. } + ) +} + +/// Drop an undecodable record that a Core resync re-delivers, and force that resync. +/// +/// Only a block-confirmed record (typed `height` set) is re-delivered by a +/// filter rescan; an unconfirmed one may never be seen again, so it keeps the +/// migration failing rather than losing it. Lowering `synced_height` to just +/// below the wallet's birth height makes the next SPV start rescan the wallet +/// from its birth, which re-records the transaction and re-applies its spends. +/// The spent marks and outputs it already produced stay as they are: +/// conservative until the rescan confirms them. +fn drop_for_resync( + tx: &Transaction<'_>, + wallet_id: &WalletId, + txid: &Txid, + error: WalletStorageError, +) -> Result<(), WalletStorageError> { + let height: Option = tx.query_row( + "SELECT height FROM core_transactions WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + |row| row.get(0), + )?; + if height.is_none() { + return Err(error); + } + let birth_height: i64 = tx.query_row( + "SELECT birth_height FROM wallets WHERE wallet_id = ?1", + params![wallet_id.as_slice()], + |row| row.get(0), + )?; + let rescan_from = (birth_height - 1).max(0); + tx.execute( + "DELETE FROM core_transaction_inputs WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + )?; + tx.execute( + "DELETE FROM core_transactions WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + )?; + tx.execute( + "UPDATE core_sync_state SET synced_height = MIN(COALESCE(synced_height, ?2), ?2) \ + WHERE wallet_id = ?1", + params![wallet_id.as_slice(), rescan_from], + )?; + tracing::warn!( + wallet_id = %hex::encode(wallet_id), + %txid, + %error, + rescan_from, + "dropped an undecodable confirmed transaction record; Core history rescans from birth" + ); + Ok(()) +} + +/// Index raw inputs independently of when their ownership becomes known. +fn index_record( + tx: &Transaction<'_>, + wallet_id: &WalletId, + record: &TransactionRecord, +) -> Result<(), WalletStorageError> { + let mut stmt = tx.prepare_cached( + "INSERT OR IGNORE INTO core_transaction_inputs (wallet_id, txid, outpoint) VALUES (?1, ?2, ?3)", + )?; + for input in &record.transaction.input { + stmt.execute(params![ + wallet_id.as_slice(), + record.txid.as_byte_array().as_slice(), + blob::encode_outpoint(&input.previous_output)?, + ])?; + } + Ok(()) +} + +fn network( + tx: &Transaction<'_>, + wallet_id: &WalletId, +) -> Result { + let label: String = tx.query_row( + "SELECT network FROM wallets WHERE wallet_id = ?1", + params![wallet_id.as_slice()], + |r| r.get(0), + )?; + wallets::parse_network(&label) + .ok_or_else(|| WalletStorageError::blob_decode("wallets.network is unknown")) +} + +fn owned_output( + tx: &Transaction<'_>, + wallet_id: &WalletId, + outpoint: &OutPoint, + network: dashcore::Network, +) -> Result, WalletStorageError> { + let mut stmt = tx.prepare_cached("SELECT value, length(script), script FROM core_utxos WHERE wallet_id = ?1 AND outpoint = ?2 AND is_sweep_placeholder = 0")?; + let mut rows = stmt.query(params![ + wallet_id.as_slice(), + blob::encode_outpoint(outpoint)? + ])?; + let Some(row) = rows.next()? else { + return Ok(None); + }; + let value = i64_to_u64("core_utxos.value", row.get(0)?)?; + blob::check_size(row.get(1)?)?; + let script: Vec = row.get(2)?; + if contact_only_script(tx, wallet_id, &script)? { + return Ok(None); + } + let address = Address::from_script(&ScriptBuf::from_bytes(script), network)?; + Ok(Some((value, address))) +} + +/// Whether `script` is tracked only by a contact's watch-only (DashPay external) chain. +fn contact_only_script( + conn: &Transaction<'_>, + wallet_id: &WalletId, + script: &[u8], +) -> Result { + Ok(conn.query_row( + "SELECT EXISTS(SELECT 1 FROM core_address_pool WHERE wallet_id = ?1 AND script = ?2) AND NOT EXISTS(SELECT 1 FROM core_address_pool WHERE wallet_id = ?1 AND script = ?2 AND account_type != 'dashpay_external')", + params![wallet_id.as_slice(), script], |r| r.get(0))?) +} + +fn repair_record( + tx: &Transaction<'_>, + wallet_id: &WalletId, + mut record: TransactionRecord, + network: dashcore::Network, +) -> Result<(), WalletStorageError> { + let record_txid = record.txid; + let txid = &record_txid; + let original = blob::encode(&record)?; + let mut inputs = BTreeMap::new(); + for detail in record.input_details.drain(..) { + if !contact_only_script(tx, wallet_id, detail.address.script_pubkey().as_bytes())? { + inputs.insert(detail.index, detail); + } + } + for (index, input) in record.transaction.input.iter().enumerate() { + if let Some((value, address)) = + owned_output(tx, wallet_id, &input.previous_output, network)? + { + inputs.insert( + index as u32, + InputDetail { + index: index as u32, + value, + address, + }, + ); + // Stale mempool rows cannot overrule a later sweep's release. + if !matches!(record.context, TransactionContext::Mempool) { + // Record the spender so the mark stays attributable and + // reversible; an existing claim by another spender stands. + tx.execute( + "UPDATE core_utxos SET spent = 1, \ + spent_in_txid = CASE WHEN spent = 1 AND spent_in_txid IS NOT NULL \ + THEN spent_in_txid ELSE ?3 END \ + WHERE wallet_id = ?1 AND outpoint = ?2", + params![ + wallet_id.as_slice(), + blob::encode_outpoint(&input.previous_output)?, + txid.as_byte_array().as_slice() + ], + )?; + } + } + } + let mut outputs = BTreeMap::new(); + for mut detail in record.output_details.drain(..) { + if let Some(address) = &detail.address { + if contact_only_script(tx, wallet_id, address.script_pubkey().as_bytes())? { + detail.role = OutputRole::Sent; + } + } + outputs.insert(detail.index, detail); + } + for (index, output) in record.transaction.output.iter().enumerate() { + let index = index as u32; + if let Some((_, address)) = owned_output( + tx, + wallet_id, + &OutPoint { + txid: *txid, + vout: index, + }, + network, + )? { + let role = outputs.get(&index).map_or(OutputRole::Received, |d| { + if d.role == OutputRole::Change { + OutputRole::Change + } else { + OutputRole::Received + } + }); + outputs.insert( + index, + OutputDetail { + index, + role, + address: Some(address), + value: output.value, + }, + ); + } + } + // Empty metadata is not accounting evidence (e.g. confirmation-only placeholders). + if inputs.is_empty() && outputs.is_empty() { + return Ok(()); + } + let received: i128 = outputs + .values() + .filter(|d| matches!(d.role, OutputRole::Received | OutputRole::Change)) + .map(|d| i128::from(d.value)) + .sum(); + let spent: i128 = inputs.values().map(|d| i128::from(d.value)).sum(); + record.net_amount = i64::try_from(received - spent).map_err(|_| { + WalletStorageError::blob_decode("wallet transaction net amount exceeds i64") + })?; + let has_ours = outputs + .values() + .any(|d| matches!(d.role, OutputRole::Received | OutputRole::Change)); + let has_external = record + .transaction + .output + .iter() + .enumerate() + .any(|(i, output)| { + !output.script_pubkey.is_op_return() + && !outputs.get(&(i as u32)).is_some_and(|d| { + matches!( + d.role, + OutputRole::Received | OutputRole::Change | OutputRole::Unspendable + ) + }) + }); + record.direction = if record.transaction_type == TransactionType::CoinJoin { + TransactionDirection::CoinJoin + } else if inputs.is_empty() { + TransactionDirection::Incoming + } else if !has_external && (has_ours || record.transaction_type == TransactionType::AssetLock) { + TransactionDirection::Internal + } else { + TransactionDirection::Outgoing + }; + record.input_details = inputs.into_values().collect(); + record.output_details = outputs.into_values().collect(); + let repaired = blob::encode(&record)?; + if repaired != original { + // Append-only: the first pre-repair blob is kept verbatim and never + // replaced, so a wrong repair can always be undone. + tx.execute( + "INSERT OR IGNORE INTO core_transaction_record_originals (wallet_id, txid, record_blob) \ + SELECT wallet_id, txid, record_blob FROM core_transactions \ + WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + )?; + tx.execute( + "UPDATE core_transactions SET record_blob = ?1 WHERE wallet_id = ?2 AND txid = ?3", + params![ + repaired, + wallet_id.as_slice(), + txid.as_byte_array().as_slice() + ], + )?; + } + Ok(()) +} + +pub(super) fn repair_history(tx: &Transaction<'_>) -> Result<(), WalletStorageError> { + // Keys are collected first so drops never race the open cursor. + let mut keys = Vec::new(); + { + let mut stmt = tx.prepare_cached("SELECT length(wallet_id), wallet_id, length(txid), txid FROM core_transactions WHERE record_blob IS NOT NULL")?; + let mut rows = stmt.query([])?; + while let Some(row) = rows.next()? { + blob::check_fixed_width(row.get(0)?, 32, "core_transactions.wallet_id")?; + let wallet_id: Vec = row.get(1)?; + let wallet_id = id32("core_transactions.wallet_id", &wallet_id)?; + blob::check_fixed_width(row.get(2)?, 32, "core_transactions.txid")?; + let txid: Vec = row.get(3)?; + keys.push((wallet_id, Txid::from_slice(&txid)?)); + } + } + for (wallet_id, txid) in keys { + match read_record(tx, &wallet_id, &txid) { + Ok(Some(record)) => { + index_record(tx, &wallet_id, &record)?; + repair_record(tx, &wallet_id, record, network(tx, &wallet_id)?)?; + } + Ok(None) => {} + Err(error) if is_unreadable(&error) => drop_for_resync(tx, &wallet_id, &txid, error)?, + Err(error) => return Err(error), + } + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use dashcore::address::Payload; + use dashcore::{BlockHash, PubkeyHash, Transaction as CoreTransaction, TxIn, TxOut}; + use key_wallet::account::{AccountType, StandardAccountType}; + use key_wallet::transaction_checking::BlockInfo; + use key_wallet::Utxo; + use platform_wallet::changeset::CoreChangeSet; + use rusqlite::Connection; + + use super::*; + use crate::sqlite::schema::core_state; + + fn address(marker: u8) -> Address { + Address::new( + dashcore::Network::Testnet, + Payload::PubkeyHash(PubkeyHash::from_byte_array([marker; 20])), + ) + } + + fn utxo(outpoint: OutPoint, value: u64, address: Address) -> Utxo { + Utxo { + outpoint, + txout: TxOut { + value, + script_pubkey: address.script_pubkey(), + }, + address, + height: 100, + is_coinbase: false, + is_confirmed: true, + is_instantlocked: false, + is_locked: false, + is_trusted: false, + } + } + + /// Pins V019's observable result on a V018-shaped database: the repaired + /// record, the preserved original, the input index and the spent marks. + #[test] + fn should_pin_v019_repair_of_a_v018_database() { + let mut conn = Connection::open_in_memory().unwrap(); + crate::sqlite::migrations::run(&mut conn).unwrap(); + let wallet_id = [0xC1u8; 32]; + conn.execute( + "INSERT INTO wallets (wallet_id, network, birth_height) VALUES (?1, 'testnet', 0)", + params![&wallet_id[..]], + ) + .unwrap(); + let (own, change, contact, external) = (address(1), address(2), address(3), address(4)); + conn.execute( + "INSERT INTO core_address_pool (wallet_id, account_type, account_index, pool_type, address_index, script) \ + VALUES (?1, 'dashpay_external', 0, 0, 0, ?2)", + params![&wallet_id[..], contact.script_pubkey().as_bytes()], + ) + .unwrap(); + let funding = OutPoint::new(Txid::from_byte_array([0x71; 32]), 0); + let body = CoreTransaction { + version: 1, + lock_time: 0, + input: vec![TxIn { + previous_output: funding, + ..Default::default() + }], + output: vec![ + TxOut { + value: 30_000, + script_pubkey: change.script_pubkey(), + }, + TxOut { + value: 50_000, + script_pubkey: contact.script_pubkey(), + }, + TxOut { + value: 15_000, + script_pubkey: external.script_pubkey(), + }, + ], + special_transaction_payload: None, + }; + let txid = body.txid(); + // What an old build stored: the input was not known to be ours and + // the contact's output was credited as received. + let original = TransactionRecord::new( + body, + AccountType::Standard { + index: 0, + standard_account_type: StandardAccountType::BIP44Account, + }, + TransactionContext::InBlock(BlockInfo::new(101, BlockHash::all_zeros(), 7)), + TransactionType::Standard, + TransactionDirection::Incoming, + Vec::new(), + vec![ + OutputDetail { + index: 0, + role: OutputRole::Change, + address: Some(change.clone()), + value: 30_000, + }, + OutputDetail { + index: 1, + role: OutputRole::Received, + address: Some(contact), + value: 50_000, + }, + ], + 80_000, + ); + { + let tx = conn.transaction().unwrap(); + core_state::apply( + &tx, + &wallet_id, + &CoreChangeSet { + new_utxos: vec![ + utxo(funding, 100_000, own.clone()), + utxo(OutPoint::new(txid, 0), 30_000, change.clone()), + ], + ..Default::default() + }, + ) + .unwrap(); + tx.execute( + "INSERT OR REPLACE INTO core_transactions (wallet_id, txid, height, finalized, record_blob) \ + VALUES (?1, ?2, 101, 1, ?3)", + params![ + &wallet_id[..], + txid.as_byte_array().as_slice(), + blob::encode(&original).unwrap() + ], + ) + .unwrap(); + tx.execute_batch( + "DROP TABLE core_transaction_inputs; DROP TABLE core_transaction_record_originals; \ + DELETE FROM refinery_schema_history WHERE version >= 19;", + ) + .unwrap(); + tx.commit().unwrap(); + } + + crate::sqlite::migrations::run(&mut conn).unwrap(); + + let mut expected = original.clone(); + expected.input_details = vec![InputDetail { + index: 0, + value: 100_000, + address: own, + }]; + expected.output_details = vec![ + OutputDetail { + index: 0, + role: OutputRole::Change, + address: Some(change), + value: 30_000, + }, + OutputDetail { + index: 1, + role: OutputRole::Sent, + address: Some(address(3)), + value: 50_000, + }, + ]; + expected.net_amount = -70_000; + expected.direction = TransactionDirection::Outgoing; + let read_blob = |sql: &str| -> Vec { + conn.query_row( + sql, + params![&wallet_id[..], txid.as_byte_array().as_slice()], + |r| r.get(0), + ) + .unwrap() + }; + assert_eq!( + read_blob( + "SELECT record_blob FROM core_transactions WHERE wallet_id = ?1 AND txid = ?2" + ), + blob::encode(&expected).unwrap() + ); + assert_eq!( + read_blob( + "SELECT record_blob FROM core_transaction_record_originals WHERE wallet_id = ?1 AND txid = ?2" + ), + blob::encode(&original).unwrap() + ); + let (spent, spender): (bool, Option>) = conn + .query_row( + "SELECT spent, spent_in_txid FROM core_utxos WHERE wallet_id = ?1 AND outpoint = ?2", + params![&wallet_id[..], blob::encode_outpoint(&funding).unwrap()], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert!(spent); + assert_eq!(spender.as_deref(), Some(txid.as_byte_array().as_slice())); + let indexed: i64 = conn + .query_row( + "SELECT count(*) FROM core_transaction_inputs WHERE wallet_id = ?1 AND txid = ?2 AND outpoint = ?3", + params![ + &wallet_id[..], + txid.as_byte_array().as_slice(), + blob::encode_outpoint(&funding).unwrap() + ], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(indexed, 1); + } +} diff --git a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs index 28623c535e2..4654f63ff1f 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs @@ -1308,11 +1308,10 @@ impl PlatformWalletPersistence for SqlitePersister { // transaction. The current schema also has lossless token balances, // invitations, account pools, tracked asset locks, and // deferred-contact-crypto queue rows. - // Do NOT attest WALLET_RESTORE (and therefore not provider restore): - // token balances and the DashPay overlay have no load readers, so a - // full restore remains lossy. Shielded viewing keys are native when - // compiled in; notes, nullifiers, and sync state use ShieldedStore. + // WALLET_RESTORE covers the Core snapshot; auxiliary token, profile, + // invitation, and deferred-contact readback remains separate. let capabilities = PersistenceCapabilities::ATOMIC_CHANGESETS + .union(PersistenceCapabilities::WALLET_RESTORE) .union(PersistenceCapabilities::INVITATIONS) .union(PersistenceCapabilities::ASSET_LOCK_FUNDING_INDICES) .union(PersistenceCapabilities::UNSIGNED_TOKEN_STORAGE) @@ -1515,7 +1514,11 @@ impl PlatformWalletPersistence for SqlitePersister { /// # } /// ``` fn load(&self) -> Result { - let conn = self.conn().map_err(PersistenceError::from)?; + let mut connection = self.conn().map_err(PersistenceError::from)?; + let conn = connection + .transaction() + .map_err(WalletStorageError::from) + .map_err(PersistenceError::from)?; let ctx = LoadCtx::new(self.config.load_policy); // Cleared up front so a failed load never leaves the previous load's // snapshot behind, masquerading as this one's verdict. @@ -1605,6 +1608,9 @@ impl PlatformWalletPersistence for SqlitePersister { unimplemented_rows = degradation.unimplemented_rows, "load() summary" ); + conn.commit() + .map_err(WalletStorageError::from) + .map_err(PersistenceError::from)?; self.replace_load_degradation(degradation); Ok(state) } @@ -1802,12 +1808,7 @@ fn load_one_wallet( key_wallet::account::account_collection::AccountCollection::new(), ) } else { - build_wallet(network, wallet_id, &account_manifest).map_err(|e| { - PersistenceError::backend(format!( - "watch-only wallet rebuild failed for {}: {e}", - hex::encode(wallet_id) - )) - })? + build_wallet(network, wallet_id, &account_manifest).map_err(PersistenceError::from)? }; // TODO(insert-wallet-id-recompute): confirm whether key_wallet's // insert_wallet recomputes wallet_id — see PR's existing Deferred @@ -1830,28 +1831,23 @@ fn load_one_wallet( &used_core_addresses, ctx, ) - .map_err(|e| { - PersistenceError::backend(format!( - "core-state rehydration failed for {}: {e}", - hex::encode(wallet_id) - )) - })?; + .map_err(PersistenceError::from)?; if account_manifest .provider .iter() .any(|entry| entry.account_type == key_wallet::account::AccountType::ProviderPlatformKeys) { restore_provider_platform_node_pool(&mut wallet_info, conn, &wallet_id, network, ctx) - .map_err(|e| { - PersistenceError::backend(format!( - "platform-node pool rehydration failed for {}: {e}", - hex::encode(wallet_id) - )) - })?; - } - let (wallet, wallet_info) = - super::rehydrate::restore_confirmed_transactions(wallet_info, wallet, core_state.records) .map_err(PersistenceError::from)?; + } + schema::core_pool::restore_pools(conn, &wallet_id, &mut wallet_info, ctx) + .map_err(PersistenceError::from)?; + let mut wallet = wallet; + super::rehydrate::restore_recorded_transactions( + &mut wallet_info, + &mut wallet, + core_state.records, + ); Ok(platform_wallet::changeset::ClientWalletStartState { wallet, wallet_info, @@ -2418,6 +2414,9 @@ mod tests { "tracked_masternodes", // load_tracked_masternodes ]; const INFRASTRUCTURE: &[&str] = &["refinery_schema_history"]; + // Append-only archive of pre-repair history blobs. Never loaded: it + // exists so a wrong history repair can be undone by hand. + const RETAINED_FOR_RECOVERY: &[&str] = &["core_transaction_record_originals"]; // `load()` rehydrates these only with the `shielded` feature on, so // the classification follows the build rather than claiming one. #[cfg(feature = "shielded")] @@ -2459,6 +2458,7 @@ mod tests { && !READ_BY_A_DEDICATED_API.contains(&table.as_str()) && !LOAD_UNIMPLEMENTED_TABLES.contains(&table.as_str()) && !INFRASTRUCTURE.contains(&table.as_str()) + && !RETAINED_FOR_RECOVERY.contains(&table.as_str()) && !FEATURE_GATED.contains(&table.as_str()) && !NOT_REHYDRATED_WITHOUT_FEATURE.contains(&table.as_str()) }) diff --git a/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs b/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs index db27b6106ce..b2f578c74ef 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs @@ -264,8 +264,9 @@ pub(crate) fn restore_provider_platform_node_pool( /// Coinbase-maturity nuance re-warms on sync. `is_instantlocked` is NOT /// among them: it is rebuilt from `core_instant_locks` above, for every /// UTXO a replayed lock covers. -/// - **Transaction records**: the SQLite loader replays confirmed history -/// through the wallet checker after this projection. +/// - **Transaction records**: the SQLite loader replays recorded history +/// (confirmed, then unconfirmed) through the wallet checker after this +/// projection. /// /// # Errors /// @@ -420,12 +421,15 @@ pub fn apply_persisted_core_state( Ok(()) } -/// Restore spend and finality guards without re-crediting outputs excluded by persistence. -pub(crate) fn restore_confirmed_transactions( - mut wallet_info: ManagedWalletInfo, - mut wallet: Wallet, +/// Restore spend reservations and finality guards without re-crediting outputs excluded by persistence. +/// +/// Unconfirmed (mempool / InstantSend) spends are replayed too: without them a +/// redelivered funding transaction would re-credit an output they reserve. +pub(crate) fn restore_recorded_transactions( + wallet_info: &mut ManagedWalletInfo, + wallet: &mut Wallet, records: Vec, -) -> Result<(Wallet, ManagedWalletInfo), WalletStorageError> { +) { // Where the load projection parked each unspent outpoint; its keys are // the outputs persistence still considers unspent. let placed: HashMap = wallet_info @@ -437,72 +441,119 @@ pub(crate) fn restore_confirmed_transactions( account.utxos.keys().map(move |outpoint| (*outpoint, owner)) }) .collect(); - let mut confirmed: Vec<_> = records + if records.is_empty() { + return; + } + // TODO(bound-load-history-replay): every stored record is replayed on each + // load; bounding it to records above the last chain lock needs care so + // finality and spend guards for older records are not lost. + let replay = replay_order(records); + + // Kept only to undo a replay the checker suspended part-way through. + let (info_before, wallet_before) = (wallet_info.clone(), wallet.clone()); + let completed = poll_ready(async { + for record in replay { + wallet_info + .check_core_transaction(&record.transaction, record.context, wallet, true, false) + .await; + } + }) + .is_some(); + if !completed { + // Degrade to the pre-replay projection: persisted spends stay + // excluded, only redelivery guards are missing until the next sync. + tracing::error!( + wallet_id = %hex::encode(wallet_info.wallet_id), + "transaction checker suspended during load replay; restored spend guards skipped" + ); + *wallet_info = info_before; + *wallet = wallet_before; + return; + } + + let spent: HashSet<_> = wallet_info + .observed_spent_outpoints() + .keys() + .copied() + .collect(); + // Replay credits an output to the account whose pool derives it. When + // that differs from the load-time fallback, the fallback copy is a + // duplicate: drop it so each outpoint lives in exactly one account. + let misplaced: HashSet<(OutPoint, AccountType)> = wallet_info + .accounts + .all_funding_accounts() .into_iter() - .filter(|record| record.block_info().is_some()) + .flat_map(|account| { + let owner = funds_account_type(account); + let placed = &placed; + account.utxos.keys().filter_map(move |outpoint| { + placed + .get(outpoint) + .filter(|parked| **parked != owner) + .map(|parked| (*outpoint, *parked)) + }) + }) .collect(); - if confirmed.is_empty() { - return Ok((wallet, wallet_info)); + for account in wallet_info.accounts.all_funding_accounts_mut() { + let owner = funds_account_type(account); + account.utxos.retain(|outpoint, _| { + placed.contains_key(outpoint) + && !spent.contains(outpoint) + && !misplaced.contains(&(*outpoint, owner)) + }); + } + // Finalize replayed records before a sync checkpoint can prune their spend guards. + if let Some(chain_lock) = wallet_info.metadata.last_applied_chain_lock.clone() { + wallet_info.apply_chain_lock(chain_lock); + } + wallet_info.update_balance(); +} + +/// Poll `future` once, returning its output only if it completed without suspending. +/// +/// The wallet checker is `async` only by trait shape: it never awaits, so it +/// completes on the first poll and load needs no async runtime. A test pins +/// that; an upstream change that adds a real await fails it. +fn poll_ready(future: F) -> Option { + let mut future = std::pin::pin!(future); + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + match future.as_mut().poll(&mut cx) { + std::task::Poll::Ready(output) => Some(output), + std::task::Poll::Pending => None, } - confirmed.sort_by_key(|record| { +} + +/// Order records as the chain would deliver them: confirmed by block position, +/// then unconfirmed ones with every in-set parent ahead of its children. +fn replay_order(records: Vec) -> Vec { + let (mut ordered, mut pending): (Vec<_>, Vec<_>) = records + .into_iter() + .partition(|record| record.block_info().is_some()); + ordered.sort_by_key(|record| { record .block_info() .map(|block| (block.height(), block.position())) }); - - // The checker only mutates in-memory state; no network requests or persistence. - dash_async::block_on(async move { - for record in confirmed { - wallet_info - .check_core_transaction( - &record.transaction, - record.context, - &mut wallet, - true, - false, - ) - .await; - } - - let spent: HashSet<_> = wallet_info - .observed_spent_outpoints() - .keys() - .copied() - .collect(); - // Replay credits an output to the account whose pool derives it. When - // that differs from the load-time fallback, the fallback copy is a - // duplicate: drop it so each outpoint lives in exactly one account. - let misplaced: HashSet<(OutPoint, AccountType)> = wallet_info - .accounts - .all_funding_accounts() - .into_iter() - .flat_map(|account| { - let owner = funds_account_type(account); - let placed = &placed; - account.utxos.keys().filter_map(move |outpoint| { - placed - .get(outpoint) - .filter(|parked| **parked != owner) - .map(|parked| (*outpoint, *parked)) - }) + // Deterministic start; parents are then pulled forward in rounds. + pending.sort_by_key(|record| record.txid); + while !pending.is_empty() { + let waiting: HashSet<_> = pending.iter().map(|record| record.txid).collect(); + let (ready, blocked): (Vec<_>, Vec<_>) = pending.into_iter().partition(|record| { + record.transaction.input.iter().all(|input| { + input.previous_output.txid == record.txid + || !waiting.contains(&input.previous_output.txid) }) - .collect(); - for account in wallet_info.accounts.all_funding_accounts_mut() { - let owner = funds_account_type(account); - account.utxos.retain(|outpoint, _| { - placed.contains_key(outpoint) - && !spent.contains(outpoint) - && !misplaced.contains(&(*outpoint, owner)) - }); - } - // Finalize replayed records before a sync checkpoint can prune their spend guards. - if let Some(chain_lock) = wallet_info.metadata.last_applied_chain_lock.clone() { - wallet_info.apply_chain_lock(chain_lock); + }); + if ready.is_empty() { + // Unreachable for real transactions (txids cannot form a cycle); + // keep the rest rather than drop a reservation. + ordered.extend(blocked); + break; } - wallet_info.update_balance(); - (wallet, wallet_info) - }) - .map_err(WalletStorageError::CoreHistoryReplay) + ordered.extend(ready); + pending = blocked; + } + ordered } /// Account identity of a funds account, stable across replay mutations. @@ -3256,4 +3307,62 @@ mod tests { "the restored UTXO must carry instant-locked status, not wait for the next sync" ); } + + /// Load replays history without an async runtime by polling the checker + /// once. If upstream ever makes it suspend, this fails in CI instead of + /// load silently skipping the restored spend guards in production. + #[test] + fn should_complete_transaction_checker_on_first_poll() { + use dashcore::hashes::Hash; + use dashcore::{OutPoint, Transaction, TxIn, TxOut, Txid}; + use key_wallet::transaction_checking::{ + BlockInfo, TransactionContext, WalletTransactionChecker, + }; + use key_wallet::wallet::initialization::WalletAccountCreationOptions; + + let mut wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + let mut info = ManagedWalletInfo::from_wallet(&wallet, 0); + let xpub = wallet.accounts.standard_bip44_accounts[&0].account_xpub; + let address = info + .accounts + .standard_bip44_accounts + .get_mut(&0) + .unwrap() + .next_receive_address(Some(&xpub), true) + .unwrap(); + let funding = Transaction { + version: 1, + lock_time: 0, + input: vec![TxIn { + previous_output: OutPoint::new(Txid::from_byte_array([9; 32]), 0), + ..Default::default() + }], + output: vec![TxOut { + value: 1_000, + script_pubkey: address.script_pubkey(), + }], + special_transaction_payload: None, + }; + let spend = Transaction { + version: 1, + lock_time: 0, + input: vec![TxIn { + previous_output: OutPoint::new(funding.txid(), 0), + ..Default::default() + }], + output: Vec::new(), + special_transaction_payload: None, + }; + let block = + TransactionContext::InBlock(BlockInfo::new(1, dashcore::BlockHash::all_zeros(), 1)); + for (tx, context) in [(&funding, block), (&spend, TransactionContext::Mempool)] { + let result = + poll_ready(info.check_core_transaction(tx, context, &mut wallet, true, false)); + assert!( + result.is_some_and(|r| r.is_relevant), + "the checker must complete on its first poll" + ); + } + } } diff --git a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs index dd7081c4668..73567111f25 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs @@ -24,9 +24,7 @@ pub(super) fn preserve_known_details( incoming: &TransactionRecord, ) -> Result { let mut merged = incoming.clone(); - let Some(previous) = - core_state::get_tx_record(tx, wallet_id, &incoming.txid, &LoadCtx::strict())? - else { + let Some(previous) = prior_record(tx, wallet_id, &incoming.txid)? else { return Ok(merged); }; if previous.transaction != incoming.transaction { @@ -61,6 +59,15 @@ pub(super) fn preserve_known_details( Ok(merged) } +/// Read a stored record strictly: corrupt history is an error, never skipped. +fn prior_record( + conn: &Connection, + wallet_id: &WalletId, + txid: &Txid, +) -> Result, WalletStorageError> { + core_state::get_tx_record(conn, wallet_id, txid, &LoadCtx::strict()) +} + /// Index raw inputs independently of when their ownership becomes known. pub(super) fn index_record( tx: &Transaction<'_>, @@ -165,8 +172,7 @@ fn repair_record( txid: &Txid, network: dashcore::Network, ) -> Result<(), WalletStorageError> { - let Some(mut record) = core_state::get_tx_record(tx, wallet_id, txid, &LoadCtx::strict())? - else { + let Some(mut record) = prior_record(tx, wallet_id, txid)? else { return Ok(()); }; let original = blob::encode(&record)?; @@ -190,11 +196,21 @@ fn repair_record( ); // Stale mempool rows cannot overrule a later sweep's release. if !matches!(record.context, TransactionContext::Mempool) { + // Record the spender so the mark stays attributable and + // reversible; an existing claim by another spender stands. + // TODO(release-repair-spends-after-reorg): release rows whose + // `spent_in_txid` spender is reorged out and never re-mined; + // needs verification of how upstream downgrades a stored + // record's context on reorg. tx.execute( - "UPDATE core_utxos SET spent = 1 WHERE wallet_id = ?1 AND outpoint = ?2", + "UPDATE core_utxos SET spent = 1, \ + spent_in_txid = CASE WHEN spent = 1 AND spent_in_txid IS NOT NULL \ + THEN spent_in_txid ELSE ?3 END \ + WHERE wallet_id = ?1 AND outpoint = ?2", params![ wallet_id.as_slice(), - blob::encode_outpoint(&input.previous_output)? + blob::encode_outpoint(&input.previous_output)?, + txid.as_byte_array().as_slice() ], )?; } @@ -268,19 +284,24 @@ fn repair_record( ) }) }); - record.direction = if record.transaction_type == TransactionType::CoinJoin { - TransactionDirection::CoinJoin - } else if inputs.is_empty() { - TransactionDirection::Incoming - } else if !has_external && (has_ours || record.transaction_type == TransactionType::AssetLock) { - TransactionDirection::Internal - } else { - TransactionDirection::Outgoing - }; + record.direction = repaired_direction( + record.transaction_type, + !inputs.is_empty(), + has_ours, + has_external, + ); record.input_details = inputs.into_values().collect(); record.output_details = outputs.into_values().collect(); let repaired = blob::encode(&record)?; if repaired != original { + // Append-only: the first pre-repair blob is kept verbatim and never + // replaced, so a wrong repair can always be undone. + tx.execute( + "INSERT OR IGNORE INTO core_transaction_record_originals (wallet_id, txid, record_blob) \ + SELECT wallet_id, txid, record_blob FROM core_transactions \ + WHERE wallet_id = ?1 AND txid = ?2", + params![wallet_id.as_slice(), txid.as_byte_array().as_slice()], + )?; tx.execute( "UPDATE core_transactions SET record_blob = ?1 WHERE wallet_id = ?2 AND txid = ?3", params![ @@ -293,22 +314,55 @@ fn repair_record( Ok(()) } -/// Backfill the input index and correct existing history in the migration transaction. -pub(crate) fn migrate(tx: &Transaction<'_>) -> Result<(), WalletStorageError> { - let mut stmt = tx.prepare_cached("SELECT length(wallet_id), wallet_id, length(txid), txid FROM core_transactions WHERE record_blob IS NOT NULL")?; - let mut rows = stmt.query([])?; - while let Some(row) = rows.next()? { - blob::check_fixed_width(row.get(0)?, 32, "core_transactions.wallet_id")?; - let wallet_id: Vec = row.get(1)?; - let wallet_id = super::id32("core_transactions.wallet_id", &wallet_id)?; - blob::check_fixed_width(row.get(2)?, 32, "core_transactions.txid")?; - let txid: Vec = row.get(3)?; - let txid = Txid::from_slice(&txid)?; - if let Some(record) = core_state::get_tx_record(tx, &wallet_id, &txid, &LoadCtx::strict())? - { - index_record(tx, &wallet_id, &record)?; - repair_record(tx, &wallet_id, &txid, network(tx, &wallet_id)?)?; +/// Direction of a repaired record. The Swift SDK's +/// `PersistentTransaction.reconciledAccounting` applies the same rule; keep +/// both in step (each side tests the same case table). +/// +/// `has_external` counts every output that is neither ours nor an OP_RETURN +/// burn, so an asset lock is internal only when nothing leaves the wallet. +fn repaired_direction( + transaction_type: TransactionType, + spends_ours: bool, + has_ours: bool, + has_external: bool, +) -> TransactionDirection { + if transaction_type == TransactionType::CoinJoin { + TransactionDirection::CoinJoin + } else if !spends_ours { + TransactionDirection::Incoming + } else if !has_external && (has_ours || transaction_type == TransactionType::AssetLock) { + TransactionDirection::Internal + } else { + TransactionDirection::Outgoing + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Shared with the Swift SDK's `TransactionAccountingTests` direction table. + #[test] + fn should_classify_repaired_direction_like_the_swift_sdk() { + use TransactionDirection::{CoinJoin, Incoming, Internal, Outgoing}; + use TransactionType::{AssetLock, Standard}; + // (type, spends ours, has owned output, has external output, expected) + let cases = [ + (Standard, true, true, false, Internal), + (Standard, true, true, true, Outgoing), + (Standard, true, false, false, Outgoing), + (AssetLock, true, false, false, Internal), + (AssetLock, true, true, false, Internal), + (AssetLock, true, false, true, Outgoing), + (Standard, false, true, false, Incoming), + (TransactionType::CoinJoin, true, true, false, CoinJoin), + ]; + for (kind, spends_ours, has_ours, has_external, expected) in cases { + assert_eq!( + repaired_direction(kind, spends_ours, has_ours, has_external), + expected, + "{kind:?} spends_ours={spends_ours} has_ours={has_ours} has_external={has_external}" + ); } } - Ok(()) } diff --git a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_pool.rs b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_pool.rs index f5b3d4056d3..2c725f23359 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_pool.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_pool.rs @@ -16,12 +16,13 @@ use rusqlite::{params, Connection, Transaction}; use platform_wallet::changeset::AccountAddressPoolEntry; use platform_wallet::wallet::platform_wallet::WalletId; -use key_wallet::managed_account::address_pool::{AddressPoolType, PublicKeyType}; +use key_wallet::managed_account::address_pool::{AddressPoolType, AddressState, PublicKeyType}; use crate::sqlite::error::WalletStorageError; -use crate::sqlite::load_ctx::{LoadCtx, LoadSite}; +use crate::sqlite::load_ctx::{LoadCtx, LoadSite, SiteCoords}; use crate::sqlite::schema::accounts; use crate::sqlite::schema::blob; +use crate::sqlite::util::safe_cast::{i64_to_u32, i64_to_u64}; /// Stored `pool_type` discriminant. Kept in the primary key so an External /// and an Internal pool never collide at the same `address_index`. @@ -483,6 +484,188 @@ pub fn load_used_addresses_with_ctx( Ok(out) } +/// Restore every persisted Core pool, including handed-out and pre-derived keys. +pub(crate) fn restore_pools( + conn: &Connection, + wallet_id: &WalletId, + wallet_info: &mut key_wallet::wallet::ManagedWalletInfo, + ctx: &LoadCtx, +) -> Result<(), WalletStorageError> { + for mut account in wallet_info.accounts.all_accounts_mut() { + let account_type = account.managed_account_type().to_account_type(); + let label = accounts::account_type_db_label(&account_type); + let (user, friend) = accounts::account_dashpay_ids(&account_type); + for pool in account.managed_account_type_mut().address_pools_mut() { + let mut stmt = conn.prepare( + "SELECT address_index, length(script), script, length(public_key), public_key, \ + key_type, used, reserved_at FROM core_address_pool \ + WHERE wallet_id = ?1 AND account_type = ?2 AND account_index = ?3 \ + AND key_class = ?4 AND user_identity_id = ?5 AND friend_identity_id = ?6 \ + AND pool_type = ?7 ORDER BY address_index", + )?; + let mut rows = stmt.query(params![ + wallet_id.as_slice(), + label, + i64::from(accounts::account_index(&account_type)), + i64::from(accounts::account_key_class(&account_type)), + user.as_slice(), + friend.as_slice(), + pool_type_to_i64(pool.pool_type) + ])?; + while let Some(row) = rows.next()? { + let index = i64_to_u32("core_address_pool.address_index", row.get(0)?)?; + blob::check_size(row.get(1)?)?; + let script = dashcore::ScriptBuf::from_bytes(row.get(2)?); + let address = match dashcore::Address::from_script(&script, pool.network) { + Ok(address) => address, + Err(source) => { + // Earlier readers already account for used and typed platform rows. + if row.get::<_, bool>(6)? + || (account_type + == key_wallet::account::AccountType::ProviderPlatformKeys + && (row.get::<_, Option>(3)?.is_some() + || row.get::<_, Option>(5)?.is_some())) + { + continue; + } + ctx.tolerate_at( + LoadSite::UndecodableAddressScript, + SiteCoords { + wallet_id: Some(*wallet_id), + account_type: &account_type, + affected: 1, + detail: Some(&pool.pool_type), + }, + WalletStorageError::from(source), + )?; + continue; + } + }; + let key_len: Option = row.get(3)?; + let key_type: Option = row.get(5)?; + let public_key = match (key_len, key_type) { + (None, None) => None, + (Some(len), Some(kind)) => { + let expected_kind = match account_type { + key_wallet::account::AccountType::ProviderOperatorKeys => KEY_TYPE_BLS, + key_wallet::account::AccountType::ProviderPlatformKeys => { + KEY_TYPE_EDDSA + } + _ => KEY_TYPE_ECDSA, + }; + if kind != expected_kind { + return Err(WalletStorageError::blob_decode( + "core pool key type does not match its account", + )); + } + let expected = match kind { + KEY_TYPE_ECDSA => ECDSA_PUBLIC_KEY_LEN, + KEY_TYPE_EDDSA => EDDSA_PUBLIC_KEY_LEN, + KEY_TYPE_BLS => BLS_PUBLIC_KEY_LEN, + _ => { + return Err(WalletStorageError::blob_decode( + "core_address_pool.key_type is outside 0..=2", + )) + } + }; + blob::check_fixed_width( + len, + expected, + "core_address_pool.public_key has the wrong length for key_type", + )?; + let bytes = row.get(4)?; + Some(match kind { + KEY_TYPE_ECDSA => PublicKeyType::ECDSA(bytes), + KEY_TYPE_EDDSA => PublicKeyType::EdDSA(bytes), + _ => PublicKeyType::BLS(bytes), + }) + } + _ => { + return Err(WalletStorageError::blob_decode( + "core_address_pool.public_key and key_type nullability differ", + )) + } + }; + let used: bool = row.get(6)?; + let reserved: Option = row.get(7)?; + let reserved = reserved + .map(|at| i64_to_u64("core_address_pool.reserved_at", at)) + .transpose()?; + let child = if pool.pool_type == AddressPoolType::AbsentHardened { + key_wallet::bip32::ChildNumber::from_hardened_idx(index) + } else { + key_wallet::bip32::ChildNumber::from_normal_idx(index) + } + .map_err(|source| WalletStorageError::AccountRecordInvalid { + e: key_wallet::error::Error::Bip32(source), + })?; + let mut path = pool.base_path.clone(); + path.push(child); + let previous = pool.addresses.get(&index); + if previous.is_some_and(|info| { + info.address != address + || public_key + .as_ref() + .zip(info.public_key.as_ref()) + .is_some_and(|(a, b)| a != b) + }) { + return Err(WalletStorageError::blob_decode( + "core pool row conflicts with its derived address or key", + )); + } + if pool + .address_index + .get(&address) + .is_some_and(|other| *other != index) + { + return Err(WalletStorageError::blob_decode( + "core pool address has multiple derivation indices", + )); + } + // Chain-derived usage takes precedence over an older pool snapshot. + let used = used || previous.is_some_and(|info| info.is_used()); + let state = if used { + AddressState::Used + } else if let Some(at) = reserved { + AddressState::Reserved { at } + } else { + AddressState::Available + }; + let info = key_wallet::AddressInfo { + address, + script_pubkey: script, + public_key: public_key + .or_else(|| previous.and_then(|info| info.public_key.clone())), + index, + path, + state, + tx_count: previous.map_or(0, |info| info.tx_count), + total_received: previous.map_or(0, |info| info.total_received), + total_sent: previous.map_or(0, |info| info.total_sent), + balance: previous.map_or(0, |info| info.balance), + label: previous.and_then(|info| info.label.clone()), + metadata: previous.map_or_else(Default::default, |info| info.metadata.clone()), + }; + if let Some(previous) = pool.addresses.remove(&index) { + pool.address_index.remove(&previous.address); + pool.script_pubkey_index.remove(&previous.script_pubkey); + } + pool.address_index.insert(info.address.clone(), index); + pool.script_pubkey_index + .insert(info.script_pubkey.clone(), index); + pool.addresses.insert(index, info); + pool.highest_generated = + Some(pool.highest_generated.map_or(index, |old| old.max(index))); + if used { + pool.used_indices.insert(index); + } + pool.highest_used = pool.used_indices.iter().max().copied(); + } + } + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; diff --git a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_state.rs b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_state.rs index a6b1c32c6a2..b7af25af2c2 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/schema/core_state.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_state.rs @@ -1510,6 +1510,36 @@ mod tests { ) .unwrap(); assert!(spent, "late funding must not resurrect the spent coin"); + let spender: Option> = tx + .query_row( + "SELECT spent_in_txid FROM core_utxos WHERE wallet_id = ?1 AND outpoint = ?2", + params![ + &wallet_id[..], + blob::encode_outpoint(&funding.outpoint).unwrap() + ], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + spender.as_deref(), + Some(AsRef::<[u8]>::as_ref(&spending.txid)), + "a repair mark names its spender so it can be reverted" + ); + let original: Vec = tx + .query_row( + "SELECT record_blob FROM core_transaction_record_originals \ + WHERE wallet_id = ?1 AND txid = ?2", + params![&wallet_id[..], AsRef::<[u8]>::as_ref(&spending.txid)], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + blob::decode::(&original) + .unwrap() + .net_amount, + 90_000, + "the pre-repair record is kept verbatim" + ); apply( &tx, &wallet_id, @@ -1567,7 +1597,7 @@ mod tests { ], ) .unwrap(); - tx.execute_batch("DROP TABLE IF EXISTS core_transaction_inputs; DELETE FROM refinery_schema_history WHERE version >= 19;").unwrap(); + tx.execute_batch("DROP TABLE IF EXISTS core_transaction_inputs; DROP TABLE IF EXISTS core_transaction_record_originals; DELETE FROM refinery_schema_history WHERE version >= 19;").unwrap(); tx.commit().unwrap(); crate::sqlite::migrations::run(&mut conn).unwrap(); let repaired = get_tx_record(&conn, &wallet_id, &spending.txid, &LoadCtx::strict()) @@ -1597,21 +1627,52 @@ mod tests { .net_amount, -150_000 ); + let original: Vec = conn + .query_row( + "SELECT record_blob FROM core_transaction_record_originals \ + WHERE wallet_id = ?1 AND txid = ?2", + params![&wallet_id[..], AsRef::<[u8]>::as_ref(&spending.txid)], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + original, + blob::encode(&spending).unwrap(), + "V019 keeps the record it rewrote" + ); } - #[test] - fn should_roll_back_history_migration_on_corrupt_record() { + /// A V018 database whose wallet 0xAD has one corrupt record at `height`. + fn v018_with_corrupt_record(height: Option) -> (Connection, [u8; 32]) { let mut conn = Connection::open_in_memory().unwrap(); crate::sqlite::migrations::run(&mut conn).unwrap(); - conn.execute_batch("DROP TABLE core_transaction_inputs; DELETE FROM refinery_schema_history WHERE version >= 19;").unwrap(); + conn.execute_batch("DROP TABLE core_transaction_inputs; DROP TABLE core_transaction_record_originals; DELETE FROM refinery_schema_history WHERE version >= 19;").unwrap(); let wallet_id = [0xADu8; 32]; conn.execute( - "INSERT INTO wallets (wallet_id, network, birth_height) VALUES (?1, 'testnet', 0)", + "INSERT INTO wallets (wallet_id, network, birth_height) VALUES (?1, 'testnet', 100)", + params![&wallet_id[..]], + ) + .unwrap(); + conn.execute( + "INSERT INTO core_sync_state (wallet_id, last_processed_height, synced_height) VALUES (?1, 500, 500)", params![&wallet_id[..]], ) .unwrap(); - conn.execute("INSERT INTO core_transactions (wallet_id, txid, finalized, record_blob) VALUES (?1, ?2, 0, ?3)", params![&wallet_id[..], &[0u8;32][..], &[0xffu8][..]]).unwrap(); - assert!(crate::sqlite::migrations::run(&mut conn).is_err()); + conn.execute( + "INSERT INTO core_transactions (wallet_id, txid, height, finalized, record_blob) VALUES (?1, ?2, ?3, 0, ?4)", + params![&wallet_id[..], &[0u8; 32][..], height, &[0xffu8][..]], + ) + .unwrap(); + (conn, wallet_id) + } + + #[test] + fn should_fail_history_migration_on_corrupt_unconfirmed_record() { + let (mut conn, _) = v018_with_corrupt_record(None); + assert!( + crate::sqlite::migrations::run(&mut conn).is_err(), + "a resync cannot restore an unconfirmed record, so it must not be dropped" + ); let tables: i64 = conn .query_row( "SELECT count(*) FROM sqlite_master WHERE name = 'core_transaction_inputs'", @@ -1631,7 +1692,85 @@ mod tests { } #[test] - fn should_keep_uncredited_outputs_spent() { + fn should_drop_corrupt_confirmed_record_and_rescan_its_wallet() { + let (mut conn, wallet_id) = v018_with_corrupt_record(Some(150)); + let kept = transaction_record( + Txid::from_byte_array([0x11; 32]), + TransactionContext::InBlock(BlockInfo::new(160, BlockHash::all_zeros(), 1)), + ); + conn.execute( + "INSERT INTO core_transactions (wallet_id, txid, height, finalized, record_blob) VALUES (?1, ?2, 160, 1, ?3)", + params![&wallet_id[..], AsRef::<[u8]>::as_ref(&kept.txid), blob::encode(&kept).unwrap()], + ) + .unwrap(); + crate::sqlite::migrations::run(&mut conn) + .expect("a re-deliverable corrupt record must not block opening"); + let rows: Vec> = conn + .prepare_cached("SELECT txid FROM core_transactions WHERE wallet_id = ?1") + .unwrap() + .query_map(params![&wallet_id[..]], |r| r.get(0)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!(rows, vec![AsRef::<[u8]>::as_ref(&kept.txid).to_vec()]); + let (last_processed, synced): (i64, i64) = conn + .query_row( + "SELECT last_processed_height, synced_height FROM core_sync_state WHERE wallet_id = ?1", + params![&wallet_id[..]], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!( + synced, 99, + "the filter checkpoint rewinds to just below birth" + ); + assert_eq!( + last_processed, 500, + "the processed watermark stays monotonic" + ); + let (cs, _) = load_state( + &conn, + &wallet_id, + dashcore::Network::Testnet, + &LoadCtx::strict(), + ) + .unwrap(); + assert_eq!(cs.synced_height, Some(99), "load hands the rewind to SPV"); + } + + #[test] + fn should_reject_corrupt_prior_record_when_storing_the_transaction() { + let mut conn = Connection::open_in_memory().unwrap(); + crate::sqlite::migrations::run(&mut conn).unwrap(); + let wallet_id = [0xA9u8; 32]; + conn.execute( + "INSERT INTO wallets (wallet_id, network, birth_height) VALUES (?1, 'testnet', 0)", + params![&wallet_id[..]], + ) + .unwrap(); + let record = transaction_record(Txid::all_zeros(), TransactionContext::Mempool); + conn.execute( + "INSERT INTO core_transactions (wallet_id, txid, finalized, record_blob) VALUES (?1, ?2, 0, ?3)", + params![&wallet_id[..], AsRef::<[u8]>::as_ref(&record.txid), &[0xffu8][..]], + ) + .unwrap(); + let tx = conn.transaction().unwrap(); + assert!( + apply( + &tx, + &wallet_id, + &CoreChangeSet { + records: vec![record], + ..Default::default() + }, + ) + .is_err(), + "normal operation treats corrupt stored history as an error" + ); + } + + #[test] + fn should_mark_observed_spent_and_doomed_outputs_spent() { use platform_wallet::changeset::changeset::UtxoCreditVerdict; let mut conn = Connection::open_in_memory().unwrap(); crate::sqlite::migrations::run(&mut conn).unwrap(); @@ -1669,6 +1808,60 @@ mod tests { } } + #[test] + fn should_keep_prior_spent_flag_for_uncredited_outputs() { + use platform_wallet::changeset::changeset::UtxoCreditVerdict; + let mut conn = Connection::open_in_memory().unwrap(); + crate::sqlite::migrations::run(&mut conn).unwrap(); + let wallet_id = [0xA8u8; 32]; + conn.execute( + "INSERT INTO wallets (wallet_id, network, birth_height) VALUES (?1, 'testnet', 0)", + params![&wallet_id[..]], + ) + .unwrap(); + let tx = conn.transaction().unwrap(); + let stored_spent = |outpoint: &OutPoint| -> Option { + tx.query_row( + "SELECT spent FROM core_utxos WHERE wallet_id = ?1 AND outpoint = ?2", + params![&wallet_id[..], blob::encode_outpoint(outpoint).unwrap()], + |row| row.get(0), + ) + .optional() + .unwrap() + }; + let uncredited = |utxo: Utxo| CoreChangeSet { + utxo_credit_verdicts: [(utxo.outpoint, UtxoCreditVerdict::Uncredited)].into(), + new_utxos: vec![utxo], + ..Default::default() + }; + + let fresh = sample_utxo(Txid::from_byte_array([3; 32]), 100, true); + apply(&tx, &wallet_id, &uncredited(fresh.clone())).unwrap(); + assert_eq!( + stored_spent(&fresh.outpoint), + None, + "never materialize an uncredited output" + ); + + let known = sample_utxo(Txid::from_byte_array([4; 32]), 100, true); + apply( + &tx, + &wallet_id, + &CoreChangeSet { + new_utxos: vec![known.clone()], + utxo_credit_verdicts: [(known.outpoint, UtxoCreditVerdict::Doomed)].into(), + ..Default::default() + }, + ) + .unwrap(); + apply(&tx, &wallet_id, &uncredited(known.clone())).unwrap(); + assert_eq!( + stored_spent(&known.outpoint), + Some(true), + "an uncredited replay keeps the spend" + ); + } + #[test] fn should_exclude_historical_contact_outputs_from_accounting() { use key_wallet::managed_account::transaction_record::{OutputDetail, OutputRole}; @@ -1772,7 +1965,7 @@ mod tests { let wallet_id = [0xB2u8; 32]; let (contact, own) = stage_contact_only_utxo(&conn, &wallet_id); conn.execute_batch( - "DROP TABLE core_transaction_inputs; DELETE FROM refinery_schema_history WHERE version >= 19;", + "DROP TABLE core_transaction_inputs; DROP TABLE core_transaction_record_originals; DELETE FROM refinery_schema_history WHERE version >= 19;", ) .unwrap(); crate::sqlite::migrations::run(&mut conn).unwrap(); diff --git a/packages/rs-platform-wallet-storage/tests/sqlite_error_classification.rs b/packages/rs-platform-wallet-storage/tests/sqlite_error_classification.rs index 935b204a2bc..36d46cb6073 100644 --- a/packages/rs-platform-wallet-storage/tests/sqlite_error_classification.rs +++ b/packages/rs-platform-wallet-storage/tests/sqlite_error_classification.rs @@ -351,7 +351,6 @@ fn samples() -> Vec { highest_used: Some(u32::MAX - 5), gap_limit: 20, }, - WalletStorageError::CoreHistoryReplay(dash_async::AsyncError::Generic("test".into())), WalletStorageError::DatabasePathIsSymlink { path: PathBuf::from("/tmp/wallet.db"), }, @@ -494,7 +493,6 @@ fn tc_p2_005_is_transient_table() { WalletStorageError::EmptyPoolAddressScript { .. } => { (false, "empty_pool_address_script") } - WalletStorageError::CoreHistoryReplay(_) => (false, "core_history_replay"), WalletStorageError::DatabasePathIsSymlink { .. } => (false, "database_path_is_symlink"), } } diff --git a/packages/rs-platform-wallet-storage/tests/sqlite_persist_roundtrip.rs b/packages/rs-platform-wallet-storage/tests/sqlite_persist_roundtrip.rs index a1419d5add9..b8ac27f8fc0 100644 --- a/packages/rs-platform-wallet-storage/tests/sqlite_persist_roundtrip.rs +++ b/packages/rs-platform-wallet-storage/tests/sqlite_persist_roundtrip.rs @@ -489,7 +489,7 @@ async fn tc010b_recovered_from_chain_lock_roundtrip() { ), "SQLite must support reconciliation of the asset-lock rows it restores" ); - assert!(!persister + assert!(persister .persistence_capabilities() .contains(platform_wallet::changeset::PersistenceCapabilities::WALLET_RESTORE)); let w = wid(0xFB); diff --git a/packages/rs-platform-wallet-storage/tests/sqlite_schema_pinning.rs b/packages/rs-platform-wallet-storage/tests/sqlite_schema_pinning.rs index 511f1b612d0..d27c6cce39b 100644 --- a/packages/rs-platform-wallet-storage/tests/sqlite_schema_pinning.rs +++ b/packages/rs-platform-wallet-storage/tests/sqlite_schema_pinning.rs @@ -21,7 +21,7 @@ const EXPECTED_ID_FINGERPRINT: &str = /// Bump it only when ADDING a migration file; a body change on an already /// applied migration is a defect, not a golden to refresh. const EXPECTED_SQL_FINGERPRINT: &str = - "50cbaacf66de2622f65f60699115a7a154345c250659b546a1d9a0658b575a97"; + "6023660fb488d9d3a98bb069bcb9950bdf125c3e8a8264f2a3dd3bd36a0844b8"; /// The migrations merged `v4.2-dev` already ships. Refinery keys /// `refinery_schema_history` by version and validates an applied migration's diff --git a/packages/rs-platform-wallet-storage/tests/sqlite_spent_rehydration.rs b/packages/rs-platform-wallet-storage/tests/sqlite_spent_rehydration.rs index 46076a9bb99..3beff7ba942 100644 --- a/packages/rs-platform-wallet-storage/tests/sqlite_spent_rehydration.rs +++ b/packages/rs-platform-wallet-storage/tests/sqlite_spent_rehydration.rs @@ -45,6 +45,13 @@ struct Fixture { impl Fixture { async fn new(spend_context: TransactionContext) -> Self { + Self::with_funding(block(100), spend_context).await + } + + async fn with_funding( + funding_context: TransactionContext, + spend_context: TransactionContext, + ) -> Self { let mut wallet = Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); let mut info = ManagedWalletInfo::from_wallet(&wallet, 0); @@ -91,7 +98,7 @@ impl Fixture { special_transaction_payload: None, }; let funding_result = info - .check_core_transaction(&funding, block(100), &mut wallet, true, true) + .check_core_transaction(&funding, funding_context, &mut wallet, true, true) .await; let coins: Vec<_> = info.accounts.standard_bip44_accounts[&0] .utxos @@ -170,6 +177,19 @@ impl Fixture { assert_eq!(selection.selected[0].outpoint, self.available); } + fn assert_spent_stored(&self, expected: bool) { + let spent: bool = self + .persister + .lock_conn_for_test() + .query_row( + "SELECT spent FROM core_utxos WHERE substr(outpoint, 2, 32) = ?1 AND value = 100000", + [self.spent.txid.as_byte_array().as_slice()], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(spent, expected, "stored spent flag of the reserved input"); + } + async fn redeliver(&self, wallet: &mut Wallet, info: &mut ManagedWalletInfo) { let result = info .check_core_transaction(&self.funding, block(100), wallet, true, true) @@ -235,9 +255,23 @@ async fn should_reconcile_stale_unspent_projection_against_confirmed_history() { #[tokio::test] async fn should_not_release_inputs_reserved_by_unconfirmed_spend() { let fixture = Fixture::new(TransactionContext::Mempool).await; - let (_, info) = fixture.load(); + let (mut wallet, mut info) = fixture.load(); fixture.assert_spent_excluded(&info); assert!(!info.observed_spent_outpoints().contains_key(&fixture.spent)); + fixture.redeliver(&mut wallet, &mut info).await; + fixture.assert_spent_stored(true); +} + +#[tokio::test] +async fn should_keep_unconfirmed_funding_reserved_when_it_confirms_after_reload() { + let fixture = + Fixture::with_funding(TransactionContext::Mempool, TransactionContext::Mempool).await; + let (mut wallet, mut info) = fixture.load(); + assert!(!info.accounts.standard_bip44_accounts[&0] + .utxos + .contains_key(&fixture.spent)); + fixture.redeliver(&mut wallet, &mut info).await; + fixture.assert_spent_stored(true); } #[tokio::test] diff --git a/packages/rs-platform-wallet-storage/tests/sqlite_wallet_restore.rs b/packages/rs-platform-wallet-storage/tests/sqlite_wallet_restore.rs new file mode 100644 index 00000000000..48d472de68e --- /dev/null +++ b/packages/rs-platform-wallet-storage/tests/sqlite_wallet_restore.rs @@ -0,0 +1,613 @@ +//! Core snapshot restoration through SQLite reopen and real manager hydration. + +mod common; + +use std::sync::Arc; + +use dashcore::{hashes::Hash, BlockHash, Network, OutPoint, Transaction, TxIn, TxOut, Txid}; +use key_wallet::{ + account::AccountType, + managed_account::address_pool::{AddressState, KeySource}, + transaction_checking::{BlockInfo, TransactionContext, WalletTransactionChecker}, + wallet::{initialization::WalletAccountCreationOptions, ManagedWalletInfo, Wallet}, +}; +use platform_wallet::{ + changeset::{ + AccountAddressPoolEntry, AccountRegistrationEntry, CoreChangeSet, PersistenceCapabilities, + PlatformWalletChangeSet, PlatformWalletPersistence, ProviderKeyAccountEntry, + ProviderKeyExtendedPubKey, WalletMetadataEntry, + }, + events::EventHandler, + PlatformEventHandler, PlatformWalletManager, +}; +use platform_wallet_storage::{SqlitePersister, SqlitePersisterConfig}; + +struct NoopHandler; +impl EventHandler for NoopHandler {} +impl PlatformEventHandler for NoopHandler {} + +fn registration(wallet: &Wallet, info: &ManagedWalletInfo) -> PlatformWalletChangeSet { + let mut provider = Vec::new(); + if let Some(account) = wallet + .accounts + .bls_account_of_type(AccountType::ProviderOperatorKeys) + { + provider.push(ProviderKeyAccountEntry { + account_type: AccountType::ProviderOperatorKeys, + extended_public_key: ProviderKeyExtendedPubKey::Bls(account.bls_public_key.clone()), + }); + } + if let Some(account) = wallet + .accounts + .eddsa_account_of_type(AccountType::ProviderPlatformKeys) + { + provider.push(ProviderKeyAccountEntry { + account_type: AccountType::ProviderPlatformKeys, + extended_public_key: ProviderKeyExtendedPubKey::EdDSA( + account.ed25519_public_key.clone(), + ), + }); + } + let mut pools = Vec::new(); + for account in info.accounts.all_accounts() { + let managed = account.managed_account_type(); + for pool in managed.address_pools() { + if !pool.addresses.is_empty() { + pools.push(AccountAddressPoolEntry { + account_type: managed.to_account_type(), + pool_type: pool.pool_type, + addresses: pool.addresses.values().cloned().collect(), + }); + } + } + } + PlatformWalletChangeSet { + wallet_metadata: Some(WalletMetadataEntry { + network: Network::Testnet, + wallet_group_id: [0; 32], + birth_height: 1, + }), + account_registrations: wallet + .accounts + .all_accounts() + .into_iter() + .map(|account| AccountRegistrationEntry { + account_type: account.account_type, + account_xpub: account.account_xpub, + }) + .collect(), + provider_key_account_registrations: provider, + account_address_pools: pools, + ..Default::default() + } +} + +async fn manager(persister: SqlitePersister) -> PlatformWalletManager { + let manager = PlatformWalletManager::new( + Arc::new( + dash_sdk::SdkBuilder::new_mock() + .with_network(Network::Testnet) + .build() + .unwrap(), + ), + Arc::new(persister), + Arc::new(NoopHandler), + ); + manager + .load_from_persistor() + .await + .expect("manager hydration"); + manager +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn should_restore_funding_indices_and_reserved_addresses_before_handing_out_keys() { + let mut wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + wallet + .add_account( + AccountType::IdentityTopUp { + registration_index: 7, + }, + None, + ) + .unwrap(); + let mut info = ManagedWalletInfo::from_wallet(&wallet, 1); + let platform_keys = + platform_wallet::wallet::provider_key_at_index::derive_platform_node_public_keys( + &wallet, + Network::Testnet, + 3, + ) + .unwrap(); + platform_wallet::wallet::provider_key_at_index::populate_platform_node_pool( + &mut info, + &platform_keys, + Network::Testnet, + ) + .unwrap(); + let mut expected = Vec::new(); + for mut account in info.accounts.all_accounts_mut() { + let account_type = account.managed_account_type().to_account_type(); + let source = wallet + .accounts + .account_of_type(account_type) + .map(|keys| KeySource::Public(keys.account_xpub)); + for pool in account.managed_account_type_mut().address_pools_mut() { + let reserved_index = if let Some(source) = &source { + pool.generate_addresses(40, source, true).unwrap(); + for index in 0..40 { + assert!(pool.mark_index_used(index)); + } + let reserved = pool.next_unused_and_reserve(source, 12345).unwrap(); + pool.address_index[&reserved] + } else { + if pool.addresses.is_empty() { + continue; + } + assert!(pool.mark_index_used(0)); + pool.addresses.get_mut(&1).unwrap().state = AddressState::Reserved { at: 12345 }; + 1 + }; + expected.push(( + account_type, + pool.pool_type, + pool.addresses.clone(), + reserved_index, + )); + } + } + for required in [ + AccountType::IdentityRegistration, + AccountType::IdentityTopUp { + registration_index: 7, + }, + AccountType::IdentityTopUpNotBoundToIdentity, + AccountType::IdentityInvitation, + AccountType::AssetLockAddressTopUp, + AccountType::AssetLockShieldedAddressTopUp, + AccountType::ProviderPlatformKeys, + AccountType::ProviderOperatorKeys, + ] { + assert!( + expected + .iter() + .any(|(account, _, _, _)| *account == required), + "missing required role {required:?}" + ); + } + let (persister, _tmp, path) = common::fresh_persister(); + persister + .store(wallet.wallet_id, registration(&wallet, &info)) + .unwrap(); + drop(persister); + let manager = manager(SqlitePersister::open(SqlitePersisterConfig::new(path)).unwrap()).await; + let wm = manager.wallet_manager_arc(); + let mut wm = wm.write().await; + let restored = wm.get_wallet_info_mut(&wallet.wallet_id).unwrap(); + for (account_type, pool_type, addresses, reserved_index) in expected { + let source = wallet + .accounts + .account_of_type(account_type) + .map(|keys| KeySource::Public(keys.account_xpub)); + let mut accounts = restored.core_wallet.accounts.all_accounts_mut(); + let account = accounts + .iter_mut() + .find(|account| account.managed_account_type().to_account_type() == account_type) + .unwrap(); + let mut pools = account.managed_account_type_mut().address_pools_mut(); + let pool = pools + .iter_mut() + .find(|pool| pool.pool_type == pool_type) + .unwrap(); + for (index, expected) in addresses { + let actual = pool + .addresses + .get(&index) + .expect("persisted beyond-gap address"); + assert_eq!( + actual.state, expected.state, + "{account_type:?} {pool_type:?} #{index}" + ); + assert_eq!(actual.address, expected.address); + assert_eq!(actual.path, expected.path); + assert_eq!(actual.public_key, expected.public_key); + } + assert_eq!( + pool.addresses[&reserved_index].state, + AddressState::Reserved { at: 12345 } + ); + if let Some(source) = source { + let next = pool.next_unused(&source, true).unwrap(); + assert!( + pool.address_index[&next] > reserved_index, + "must not reuse funded or handed-out keys" + ); + } + } + drop(wm); + assert!(manager.shutdown().await.all_clean()); +} + +fn transaction(input: OutPoint, outputs: Vec) -> Transaction { + Transaction { + version: 1, + lock_time: 0, + input: vec![TxIn { + previous_output: input, + ..Default::default() + }], + output: outputs, + special_transaction_payload: None, + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn should_restore_pending_sends_in_dependency_order_without_resending_incoming() { + let mut wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + let mut info = ManagedWalletInfo::from_wallet(&wallet, 1); + let xpub = wallet.accounts.standard_bip44_accounts[&0].account_xpub; + let address = info + .accounts + .standard_bip44_accounts + .get_mut(&0) + .unwrap() + .next_receive_address(Some(&xpub), true) + .unwrap(); + let output = |value| TxOut { + value, + script_pubkey: address.script_pubkey(), + }; + let funding = transaction( + OutPoint::new(Txid::from_byte_array([1; 32]), 0), + vec![output(100_000)], + ); + let parent = transaction(OutPoint::new(funding.txid(), 0), vec![output(90_000)]); + let child = transaction(OutPoint::new(parent.txid(), 0), vec![output(80_000)]); + let incoming = transaction( + OutPoint::new(Txid::from_byte_array([2; 32]), 0), + vec![output(5_000)], + ); + let (persister, _tmp, path) = common::fresh_persister(); + persister + .store(wallet.wallet_id, registration(&wallet, &info)) + .unwrap(); + for (tx, context) in [ + ( + &funding, + TransactionContext::InBlock(BlockInfo::new( + 100, + BlockHash::from_byte_array([3; 32]), + 1, + )), + ), + (&parent, TransactionContext::Mempool), + (&child, TransactionContext::Mempool), + (&incoming, TransactionContext::Mempool), + ] { + let result = info + .check_core_transaction(tx, context, &mut wallet, true, true) + .await; + assert!(result.is_relevant); + persister + .store( + wallet.wallet_id, + PlatformWalletChangeSet { + core: Some(CoreChangeSet { + records: result.new_records, + new_utxos: info.accounts.standard_bip44_accounts[&0] + .utxos + .values() + .cloned() + .collect(), + synced_height: Some(100), + last_processed_height: Some(100), + ..Default::default() + }), + ..Default::default() + }, + ) + .unwrap(); + } + drop(persister); + let persister = SqlitePersister::open(SqlitePersisterConfig::new(path)).unwrap(); + let start = persister.load().unwrap(); + let pending = &start.wallets[&wallet.wallet_id].unconfirmed_outgoing_txs; + assert!( + pending.is_empty(), + "restoration must not schedule transactions for rebroadcast" + ); + let manager = manager(persister).await; + let wm = manager.wallet_manager_arc(); + let wm = wm.read().await; + let restored = &wm.get_wallet_info(&wallet.wallet_id).unwrap().core_wallet; + let utxos = &restored.accounts.standard_bip44_accounts[&0].utxos; + assert!(!utxos.contains_key(&OutPoint::new(funding.txid(), 0))); + assert!(!utxos.contains_key(&OutPoint::new(parent.txid(), 0))); + assert!(utxos.contains_key(&OutPoint::new(child.txid(), 0))); + assert!(utxos.contains_key(&OutPoint::new(incoming.txid(), 0))); + assert_eq!(utxos.values().map(|utxo| utxo.value()).sum::(), 85_000); + use key_wallet::managed_account::managed_account_trait::ManagedAccountTrait; + let before = info.accounts.standard_bip44_accounts[&0] + .managed_account_type() + .address_pools()[0] + .addresses + .values() + .find(|entry| entry.address == address) + .unwrap(); + let after = restored.accounts.standard_bip44_accounts[&0] + .managed_account_type() + .address_pools()[0] + .addresses + .values() + .find(|entry| entry.address == address) + .unwrap(); + assert_eq!( + ( + after.tx_count, + after.total_received, + after.total_sent, + after.balance + ), + ( + before.tx_count, + before.total_received, + before.total_sent, + before.balance + ) + ); + drop(wm); + assert!(manager.shutdown().await.all_clean()); +} + +#[test] +fn should_advertise_full_core_wallet_restore() { + let (persister, _tmp, _path) = common::fresh_persister(); + assert!(persister + .persistence_capabilities() + .contains(PersistenceCapabilities::WALLET_RESTORE)); +} + +#[test] +fn should_read_one_snapshot_despite_a_concurrent_commit() { + use rusqlite::hooks::{AuthAction, Authorization}; + use std::sync::atomic::{AtomicBool, Ordering}; + + let (persister, _tmp, path) = common::fresh_persister(); + let wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + let info = ManagedWalletInfo::from_wallet(&wallet, 1); + let mut cs = registration(&wallet, &info); + cs.core = Some(CoreChangeSet { + synced_height: Some(100), + ..Default::default() + }); + persister.store(wallet.wallet_id, cs).unwrap(); + let writer = rusqlite::Connection::open(path).unwrap(); + let changed = Arc::new(AtomicBool::new(false)); + let changed_callback = Arc::clone(&changed); + persister + .lock_conn_for_test() + .authorizer(Some(move |ctx: rusqlite::hooks::AuthContext<'_>| { + if matches!( + ctx.action, + AuthAction::Read { + table_name: "core_utxos", + .. + } + ) && !changed_callback.swap(true, Ordering::SeqCst) + { + writer + .execute("UPDATE core_sync_state SET synced_height = 200", []) + .unwrap(); + } + Authorization::Allow + })) + .unwrap(); + let state = persister.load().unwrap(); + assert!(changed.load(Ordering::SeqCst)); + assert_eq!( + state.wallets[&wallet.wallet_id] + .wallet_info + .metadata + .synced_height, + 100, + "one load must retain its original read snapshot" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn should_reject_malformed_pool_keys_with_typed_source_before_manager_hydration() { + use platform_wallet::changeset::PersistenceError; + use platform_wallet_storage::WalletStorageError; + + let (persister, _tmp, path) = common::fresh_persister(); + let wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + let info = ManagedWalletInfo::from_wallet(&wallet, 1); + persister + .store(wallet.wallet_id, registration(&wallet, &info)) + .unwrap(); + persister.lock_conn_for_test().execute( + "UPDATE core_address_pool SET public_key = X'00' WHERE account_type = 'identity_invitation'", + [], + ).unwrap(); + drop(persister); + let persister = SqlitePersister::open(SqlitePersisterConfig::new(path)).unwrap(); + let error = persister + .load() + .expect_err("invalid key must fail the full snapshot"); + let PersistenceError::Backend { source, .. } = error else { + panic!("typed backend error") + }; + assert!(source.downcast_ref::().is_some()); + let manager = PlatformWalletManager::new( + Arc::new( + dash_sdk::SdkBuilder::new_mock() + .with_network(Network::Testnet) + .build() + .unwrap(), + ), + Arc::new(persister), + Arc::new(NoopHandler), + ); + assert!(manager.load_from_persistor().await.is_err()); + assert!(manager + .wallet_manager_arc() + .read() + .await + .get_wallet_info(&wallet.wallet_id) + .is_none()); + assert!(manager.shutdown().await.all_clean()); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn should_restore_provider_key_transaction_history_and_positions() { + use dashcore::blockdata::transaction::special_transaction::{ + provider_registration::{ProviderMasternodeType, ProviderRegistrationPayload}, + provider_update_registrar::ProviderUpdateRegistrarPayload, + TransactionPayload, + }; + use key_wallet::managed_account::address_pool::PublicKeyType; + use key_wallet::managed_account::managed_account_trait::ManagedAccountTrait; + + let mut wallet = + Wallet::new_random(Network::Testnet, WalletAccountCreationOptions::Default).unwrap(); + let mut info = ManagedWalletInfo::from_wallet(&wallet, 1); + for role in [ + AccountType::ProviderOwnerKeys, + AccountType::ProviderVotingKeys, + ] { + let source = KeySource::Public(wallet.accounts.account_of_type(role).unwrap().account_xpub); + let mut accounts = info.accounts.all_accounts_mut(); + let account = accounts + .iter_mut() + .find(|account| account.managed_account_type().to_account_type() == role) + .unwrap(); + for pool in account.managed_account_type_mut().address_pools_mut() { + pool.generate_addresses(45, &source, true).unwrap(); + } + } + let hash_at = |role| { + let accounts = info.accounts.all_accounts(); + let account = accounts + .iter() + .find(|account| account.managed_account_type().to_account_type() == role) + .unwrap(); + let address = &account.managed_account_type().address_pools()[0].addresses[&40].address; + let dashcore::address::Payload::PubkeyHash(hash) = address.payload() else { + panic!("P2PKH") + }; + *hash + }; + let owner = hash_at(AccountType::ProviderOwnerKeys); + let voting = hash_at(AccountType::ProviderVotingKeys); + let operator_pool = info + .accounts + .provider_operator_keys + .as_ref() + .unwrap() + .managed_account_type() + .address_pools()[0]; + let Some(PublicKeyType::BLS(operator)) = &operator_pool.addresses[&0].public_key else { + panic!("BLS key") + }; + let operator = dashcore::bls_sig_utils::BLSPublicKey::from( + <[u8; 48]>::try_from(operator.as_slice()).unwrap(), + ); + let special = |payload| Transaction { + version: 3, + lock_time: 0, + input: vec![], + output: vec![], + special_transaction_payload: Some(payload), + }; + let registration_tx = special(TransactionPayload::ProviderRegistrationPayloadType( + ProviderRegistrationPayload { + version: 2, + masternode_type: ProviderMasternodeType::Regular, + masternode_mode: 0, + collateral_outpoint: OutPoint::null(), + service_address: "127.0.0.1:19999".parse().unwrap(), + owner_key_hash: owner, + operator_public_key: operator, + voting_key_hash: voting, + operator_reward: 0, + script_payout: dashcore::ScriptBuf::new(), + inputs_hash: [0; 32].into(), + signature: vec![], + platform_node_id: None, + platform_p2p_port: None, + platform_http_port: None, + }, + )); + let pro_tx_hash = registration_tx.txid(); + let registrar = special(TransactionPayload::ProviderUpdateRegistrarPayloadType( + ProviderUpdateRegistrarPayload { + version: 2, + pro_tx_hash, + provider_mode: 0, + operator_public_key: operator, + voting_key_hash: voting, + script_payout: dashcore::ScriptBuf::new(), + inputs_hash: [0; 32].into(), + payload_sig: vec![], + }, + )); + let (persister, _tmp, path) = common::fresh_persister(); + persister + .store(wallet.wallet_id, registration(&wallet, &info)) + .unwrap(); + let mut expected = Vec::new(); + for (position, tx) in [registration_tx, registrar].into_iter().enumerate() { + let context = TransactionContext::InBlock(BlockInfo::new( + 100, + BlockHash::from_byte_array([7; 32]), + position as u32, + )); + let result = info + .check_core_transaction(&tx, context, &mut wallet, true, true) + .await; + assert!( + result.is_relevant, + "provider payload #{position} must match owned keys" + ); + assert!(!result.new_records.is_empty()); + expected.extend(result.new_records.clone()); + persister + .store( + wallet.wallet_id, + PlatformWalletChangeSet { + core: Some(CoreChangeSet { + records: result.new_records, + ..Default::default() + }), + ..Default::default() + }, + ) + .unwrap(); + } + drop(persister); + let manager = manager(SqlitePersister::open(SqlitePersisterConfig::new(path)).unwrap()).await; + let wm = manager.wallet_manager_arc(); + let wm = wm.read().await; + let restored = &wm.get_wallet_info(&wallet.wallet_id).unwrap().core_wallet; + for expected in expected { + let accounts = restored.accounts.all_accounts(); + let account = accounts + .iter() + .find(|account| { + account.managed_account_type().to_account_type() == expected.account_type + }) + .unwrap(); + let actual = account + .transactions() + .get(&expected.txid) + .expect("provider history"); + assert_eq!(actual.transaction, expected.transaction); + assert_eq!(actual.context, expected.context); + } + drop(wm); + assert!(manager.shutdown().await.all_clean()); +} diff --git a/packages/rs-platform-wallet/src/changeset/changeset.rs b/packages/rs-platform-wallet/src/changeset/changeset.rs index 0c9a8988d75..f678c407003 100644 --- a/packages/rs-platform-wallet/src/changeset/changeset.rs +++ b/packages/rs-platform-wallet/src/changeset/changeset.rs @@ -626,6 +626,10 @@ fn coalesce_newest_wins( } /// Keep the last correction per account before folding account contributions. +// TODO(unify-record-coalescing): `coalesce_newest_wins` (the `Merge` path) +// can regress a confirmed context to an older mempool slice; this helper +// keeps the highest context rank. Unify them once `Merge` semantics are +// reviewed. pub(crate) fn coalesce_account_records(records: &mut Vec) { let mut positions = BTreeMap::new(); for mut record in std::mem::take(records) { diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift index 0fe3efe8f2e..f79cbbebdb2 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift @@ -241,8 +241,9 @@ public final class PersistentTransaction { let wallets = Set((inputs + outputs).filter(PlatformWalletPersistenceHandler.isWalletOwnedTxo) .compactMap { PlatformWalletPersistenceHandler.resolvedWalletId(of: $0) }) let hasUnownedTxos = (inputs + outputs).contains { !PlatformWalletPersistenceHandler.isWalletOwnedTxo($0) } - if wallets.count == 1, wallets.contains(walletId), !hasUnownedTxos { return netAmount } + // Unresolved inputs make any amount provisional, the stored one included. guard pendingInputs.isEmpty else { return nil } + if wallets.count == 1, wallets.contains(walletId), !hasUnownedTxos { return netAmount } let walletInputs = owned(inputs) let walletOutputs = owned(outputs) guard !walletInputs.isEmpty || !walletOutputs.isEmpty else { return nil } @@ -256,12 +257,13 @@ public final class PersistentTransaction { public func direction(for walletId: Data) -> UInt32 { let wallets = Set((inputs + outputs).filter(PlatformWalletPersistenceHandler.isWalletOwnedTxo) .compactMap { PlatformWalletPersistenceHandler.resolvedWalletId(of: $0) }) - guard wallets.count > 1, direction != 3, transactionTypeKind != 1, !isAssetLock else { return direction } + guard wallets.count > 1, direction != CoreDirectionCode.coinJoin, + typedKind != .coinJoin, !isAssetLock else { return direction } let spendsOurs = inputs.contains { PlatformWalletPersistenceHandler.isWalletOwnedTxo($0) && PlatformWalletPersistenceHandler.resolvedWalletId(of: $0) == walletId } - return spendsOurs ? 1 : 0 + return spendsOurs ? CoreDirectionCode.outgoing : CoreDirectionCode.incoming } /// Format the wallet's Core value movement in DASH. @@ -287,11 +289,15 @@ public final class PersistentTransaction { guard let received = total(ownedOutputAmounts), let spent = total(inputs.map(\.amount)) else { return nil } + // Same rule as the Rust repair (`core_history::repaired_direction`): + // internal only when nothing leaves the wallet and something stays in + // it, or an asset lock burns into Platform. Both sides test one table. let direction: UInt32 - if previousDirection == 3 { direction = 3 } - else if inputs.isEmpty { direction = 0 } - else if isAssetLock || allOutputsOwned { direction = 2 } - else { direction = 1 } + if previousDirection == CoreDirectionCode.coinJoin { direction = CoreDirectionCode.coinJoin } + else if inputs.isEmpty { direction = CoreDirectionCode.incoming } + else if allOutputsOwned && (!ownedOutputAmounts.isEmpty || isAssetLock) { + direction = CoreDirectionCode.internalTransfer + } else { direction = CoreDirectionCode.outgoing } return (received - spent, direction) } @@ -423,6 +429,15 @@ public final class PersistentTransaction { /// "pre-feature / not-populated" sentinel and is NOT a case in this /// enum — `TransactionTypeKind(rawValue: 0xFF)` returns `nil`, which /// the accessors treat as "unknown" so no branch fires falsely. +/// Wire values of `PersistentTransaction.direction`, matching `directionName` +/// and the FFI's `TransactionDirection` discriminants. +enum CoreDirectionCode { + static let incoming: UInt32 = 0 + static let outgoing: UInt32 = 1 + static let internalTransfer: UInt32 = 2 + static let coinJoin: UInt32 = 3 +} + public enum TransactionTypeKind: UInt8 { case standard = 0 case coinJoin = 1 diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift index 354e2cf32e8..2893e42c4a2 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift @@ -2547,8 +2547,10 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { let blockHashBytes = hashData(tx.block_hash) record.blockHash = blockHashBytes.allSatisfy { $0 == 0 } ? nil : blockHashBytes // A context-only recovery record has zero accounting; a funded asset lock burns Core value. - let preserveLockAccounting = tx.transaction_type_kind == 6 && tx.net_amount == 0 && !tx.has_fee - && record.netAmount != 0 && record.inputs.contains(where: Self.isWalletOwnedTxo) + // A stored debit is itself the proof we funded it: its inputs may not be linked yet. + let preserveLockAccounting = tx.transaction_type_kind == TransactionTypeKind.assetLock.rawValue + && tx.net_amount == 0 && !tx.has_fee + && record.netAmount < 0 if !preserveLockAccounting { record.direction = tx.direction } if let typeName = tx.transaction_type { record.transactionType = String(cString: typeName) @@ -3427,8 +3429,19 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { if roundAdvancedFinalityBoundary { collectFinalizedSweptTombstones(walletId: walletId) } + // Display-only accounting: a failure here must not fail the + // round, and is reported under its own event, not save_failed. + // Recomputed values that did land are consistent on their own. do { try reconcileTransactionAccounting(Array(accountingDirty.values)) + } catch { + SDKLogger.event( + "persistence_transaction_accounting_failed", category: .persistence, severity: .error, + fields: ["phase": .publicText("round"), "wallet_reference": .reference(walletId)], + error: error + ) + } + do { try backgroundContext.save() committedRoundGeneration &+= 1 SDKLogger.event( @@ -6810,7 +6823,8 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { /// Repair only fully resolved spends; missing prevouts are not evidence of external ownership. func reconcileTransactionAccounting( - _ transactions: [PersistentTransaction], txos: [Data: PersistentTxo]? = nil + _ transactions: [PersistentTransaction], txos: [Data: PersistentTxo]? = nil, + addresses: [String: PersistentCoreAddress]? = nil ) throws { for transaction in transactions where !transaction.isDeleted { guard let transactionNetwork = network @@ -6837,6 +6851,10 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { var amounts: [UInt64] = [] var allOutputsOwned = true let ownedVouts = Set(transaction.outputs.filter { !$0.isDeleted && Self.isWalletOwnedTxo($0) }.map(\.vout)) + // An address-matched output with no linked TXO carries no wallet of + // its own, so it counts only for the spending wallets: another local + // wallet's credit must not hide inside the sender's scalar. + let spendingWallets = Set(inputs.compactMap { Self.resolvedWalletId(of: $0) }) for (index, output) in decoded.outputs.enumerated() { // OP_RETURN burns (including asset locks) are not spendable Core outputs. if output.scriptPubkey.first == 0x6a { continue } @@ -6845,8 +6863,10 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { let descriptor = FetchDescriptor(predicate: #Predicate { $0.address == address }) let owner: PersistentCoreAddress? if let cached = roundIndex?.coreAddressesByAddress[address] { owner = cached } + else if let addresses { owner = addresses[address] } else { owner = try modelFetcher.fetch(descriptor, in: backgroundContext).first } - if let account = owner?.account, account.accountType != Self.dashpayExternalAccountTypeTag { + if let account = owner?.account, account.accountType != Self.dashpayExternalAccountTypeTag, + spendingWallets.contains(account.wallet.walletId) { belongs = true } } @@ -6855,7 +6875,8 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { } if let accounting = PersistentTransaction.reconciledAccounting( inputs: inputs, ownedOutputAmounts: amounts, allOutputsOwned: allOutputsOwned, - previousDirection: transaction.transactionTypeKind == 1 ? 3 : transaction.direction, + previousDirection: transaction.typedKind == .coinJoin + ? CoreDirectionCode.coinJoin : transaction.direction, isAssetLock: transaction.isAssetLock ) { transaction.netAmount = accounting.netAmount @@ -6912,6 +6933,10 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { // Mid-round the context holds another round's staged writes: saving // would commit half of it and rolling back would silently drop it. // That round reconciles its own dirty rows; the next load repairs the rest. + // TODO(persist-accounting-backfill-marker): this pass re-reads all + // history on every launch; a persisted completion marker needs a + // SwiftData shape change (a new live schema version) or a side + // channel, which is a product decision. if !inChangeset { do { let walletIds = Set(wallets.map(\.walletId)) @@ -6923,15 +6948,22 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { } } let txos = try modelFetcher.fetch(FetchDescriptor(), in: backgroundContext) - try reconcileTransactionAccounting(transactions, txos: Dictionary(uniqueKeysWithValues: txos.map { ($0.outpoint, $0) })) + // One read instead of one per unlinked output. + let addresses = try modelFetcher.fetch(FetchDescriptor(), in: backgroundContext) + try reconcileTransactionAccounting( + transactions, + txos: Dictionary(uniqueKeysWithValues: txos.map { ($0.outpoint, $0) }), + addresses: Dictionary(addresses.map { ($0.address, $0) }, uniquingKeysWith: { first, _ in first }) + ) try backgroundContext.save() } catch { + // Display-only accounting must never block restoring wallets: + // drop this pass's edits and restore on the stored values. backgroundContext.rollback() SDKLogger.event( - "persistence_wallet_load_failed", category: .persistence, severity: .error, - fields: ["phase": .publicText("transaction_accounting")], error: error + "persistence_transaction_accounting_failed", category: .persistence, severity: .error, + fields: ["phase": .publicText("load")], error: error ) - return (nil, 0, true) } } let restorable = wallets.filter { wallet in diff --git a/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionDetailView.swift b/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionDetailView.swift index 37bd50aedad..c854fd5ce4d 100644 --- a/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionDetailView.swift +++ b/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionDetailView.swift @@ -4,7 +4,9 @@ import SwiftDashSDK struct TransactionDetailView: View { let transaction: PersistentTransaction var walletId: Data? = nil - private var netAmount: Int64 { walletId.flatMap { transaction.netAmount(for: $0) } ?? transaction.netAmount } + /// `nil` while this wallet's amount is unresolved — the same state the + /// amount label shows as "Amount unavailable", so fee and amount agree. + private var netAmount: Int64? { walletId.map { transaction.netAmount(for: $0) } ?? transaction.netAmount } private var direction: UInt32 { walletId.map { transaction.direction(for: $0) } ?? transaction.direction } /// Asset-lock payload funding amount, excluding the Core transaction fee. var assetLockAmountDuffs: Int64? = nil @@ -201,7 +203,7 @@ struct TransactionDetailView: View { ) } - if let fee = formattedFee, netAmount < 0 { + if let fee = formattedFee, let amount = netAmount, amount < 0 { TransactionDetailRow( label: "Network Fee", value: fee diff --git a/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionListView.swift b/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionListView.swift index 4e0033ff220..1f34b9467e8 100644 --- a/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionListView.swift +++ b/packages/swift-sdk/SwiftExampleApp/SwiftExampleApp/Core/Views/TransactionListView.swift @@ -222,7 +222,9 @@ struct TransactionListView: View { struct TransactionRowView: View { let transaction: PersistentTransaction var walletId: Data? = nil - private var netAmount: Int64 { walletId.flatMap { transaction.netAmount(for: $0) } ?? transaction.netAmount } + /// `nil` while this wallet's amount is unresolved — the same state the + /// amount label shows as "Amount unavailable", so fee and amount agree. + private var netAmount: Int64? { walletId.map { transaction.netAmount(for: $0) } ?? transaction.netAmount } private var direction: UInt32 { walletId.map { transaction.direction(for: $0) } ?? transaction.direction } /// Asset-lock payload funding amount, excluding the Core transaction fee. var assetLockAmountDuffs: Int64? = nil @@ -400,7 +402,7 @@ struct TransactionRowView: View { .font(.headline) .foregroundColor(typeColor) - if let fee = transaction.fee, netAmount < 0 { + if let fee = transaction.fee, let amount = netAmount, amount < 0 { Text("Fee: \(formatFee(fee))") .font(.caption2) .foregroundColor(.secondary) diff --git a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index ccde14865f4..26b690a985c 100644 --- a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift +++ b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift @@ -31,7 +31,7 @@ final class TransactionAccountingTests: XCTestCase { inputs: [spent], ownedOutputAmounts: [99], allOutputsOwned: true, previousDirection: 0, isAssetLock: false )?.direction, 2) let lock = PersistentTransaction.reconciledAccounting( - inputs: [spent], ownedOutputAmounts: [], allOutputsOwned: false, previousDirection: 2, isAssetLock: true + inputs: [spent], ownedOutputAmounts: [], allOutputsOwned: true, previousDirection: 2, isAssetLock: true ) XCTAssertEqual(lock?.netAmount, -100) XCTAssertEqual(lock?.direction, 2) @@ -40,6 +40,30 @@ final class TransactionAccountingTests: XCTestCase { )?.direction, 3) } + /// Same case table as the Rust repair's + /// `should_classify_repaired_direction_like_the_swift_sdk`. + func testShouldClassifyDirectionLikeTheRustRepair() { + let spent = input(100) + // (spends ours, owned output amounts, all outputs owned, asset lock, previous direction, expected) + let cases: [(Bool, [UInt64], Bool, Bool, UInt32, UInt32)] = [ + (true, [99], true, false, 0, 2), + (true, [40], false, false, 0, 1), + (true, [], true, false, 0, 1), + (true, [], true, true, 0, 2), + (true, [40], true, true, 0, 2), + (true, [], false, true, 0, 1), + (false, [40], true, false, 0, 0), + (true, [99], true, false, 3, 3), + ] + for (index, (spendsOurs, owned, allOwned, isLock, previous, expected)) in cases.enumerated() { + let result = PersistentTransaction.reconciledAccounting( + inputs: spendsOurs ? [spent] : [], ownedOutputAmounts: owned, + allOutputsOwned: allOwned, previousDirection: previous, isAssetLock: isLock + ) + XCTAssertEqual(result?.direction, expected, "case \(index)") + } + } + func testShouldRejectOverflowInsteadOfWrappingHistory() { XCTAssertNil(PersistentTransaction.reconciledAccounting( inputs: [input(UInt64.max)], ownedOutputAmounts: [], allOutputsOwned: false, previousDirection: 0, isAssetLock: false @@ -54,7 +78,8 @@ final class TransactionAccountingTests: XCTestCase { XCTAssertEqual(tx.netAmount(for: first.walletId), -100) XCTAssertEqual(tx.netAmount(for: other.walletId), -200) } - private func serializedSpend(inputs: [Data], outputValue: UInt64 = 40) -> Data { + /// `burn` makes the single output an OP_RETURN, as an asset lock's is. + private func serializedSpend(inputs: [Data], outputValue: UInt64 = 40, burn: Bool = false) -> Data { var bytes = Data([2, 0, 0, 0, UInt8(inputs.count)]) for txid in inputs { bytes.append(txid) @@ -62,7 +87,8 @@ final class TransactionAccountingTests: XCTestCase { } bytes.append(1) withUnsafeBytes(of: outputValue.littleEndian) { bytes.append(contentsOf: $0) } - bytes.append(contentsOf: [0, 0, 0, 0, 0]) + bytes.append(contentsOf: burn ? [1, 0x6a] : [0]) + bytes.append(contentsOf: [0, 0, 0, 0]) return bytes } @@ -139,6 +165,18 @@ final class TransactionAccountingTests: XCTestCase { XCTAssertEqual(repaired?.netAmount, -100) } + func testShouldRestoreWalletsWhenLoadTimeAccountingFails() throws { + let container = try DashModelContainer.createInMemory() + container.mainContext.insert(PersistentWallet(walletId: Data(repeating: 1, count: 32), network: .testnet)) + try container.mainContext.save() + let injector = FetchFaultInjector(faulting: PersistentCoreAddress.self) + let handler = PlatformWalletPersistenceHandler( + modelContainer: container, network: .testnet, modelFetcher: injector + ) + XCTAssertFalse(handler.loadWalletList().errored, "display accounting must not block restore") + XCTAssertTrue(injector.observedReads.contains("PersistentCoreAddress")) + } + func testShouldPreserveAccountingWhenSomePrevoutsAreMissing() throws { let container = try DashModelContainer.createInMemory() let context = container.mainContext @@ -242,7 +280,7 @@ final class TransactionAccountingTests: XCTestCase { let walletId = Data(repeating: 1, count: 32) context.insert(PersistentWallet(walletId: walletId, network: .testnet)) let spenderId = Data(repeating: 3, count: 32) - let bytes = serializedSpend(inputs: [walletId, Data(repeating: 2, count: 32)], outputValue: 0) + let bytes = serializedSpend(inputs: [walletId, Data(repeating: 2, count: 32)], outputValue: 0, burn: true) let spender = PersistentTransaction(txid: spenderId, transactionData: bytes, direction: 2, netAmount: -200) spender.transactionTypeKind = 6 let coin = input(100) @@ -258,6 +296,26 @@ final class TransactionAccountingTests: XCTestCase { XCTAssertEqual(row.context, 2) } + func testShouldPreserveAssetLockDebitWhenNoInputIsLinkedYet() throws { + let container = try DashModelContainer.createInMemory() + let context = container.mainContext + let walletId = Data(repeating: 1, count: 32) + context.insert(PersistentWallet(walletId: walletId, network: .testnet)) + let spenderId = Data(repeating: 3, count: 32) + let bytes = serializedSpend(inputs: [Data(repeating: 2, count: 32)], outputValue: 0, burn: true) + let lock = PersistentTransaction(txid: spenderId, transactionData: bytes, direction: 2, netAmount: -200) + lock.transactionTypeKind = 6 + lock.fee = 7 + context.insert(lock) + try context.save() + let handler = PlatformWalletPersistenceHandler(modelContainer: container, network: .testnet) + persist(handler, walletId: walletId, txid: spenderId, bytes: bytes, kind: 6) + let row = try XCTUnwrap(ModelContext(container).fetch(FetchDescriptor()).first { $0.txid == spenderId }) + XCTAssertEqual(row.netAmount, -200, "a synthetic zero update must not erase the debit") + XCTAssertEqual(row.fee, 7) + XCTAssertEqual(row.direction, 2) + } + func testShouldRepairNoChangeAssetLockToFullCoreDebit() throws { let container = try DashModelContainer.createInMemory() let walletId = Data(repeating: 1, count: 32) @@ -267,13 +325,65 @@ final class TransactionAccountingTests: XCTestCase { let spenderId = Data(repeating: 3, count: 32) persist(handler, walletId: walletId, txid: walletId, outputs: [(walletId, 100)]) persist(handler, walletId: walletId, txid: spenderId, - bytes: serializedSpend(inputs: [walletId], outputValue: 0), kind: 6, inputTxids: [walletId]) + bytes: serializedSpend(inputs: [walletId], outputValue: 0, burn: true), kind: 6, inputTxids: [walletId]) let row = try XCTUnwrap(ModelContext(container).fetch(FetchDescriptor()).first { $0.txid == spenderId }) XCTAssertEqual(row.netAmount, -100) XCTAssertEqual(row.direction, 2) XCTAssertTrue(row.isAssetLock) } + func testShouldNotReportStoredAmountWhileInputsArePending() { + let walletId = Data(repeating: 1, count: 32) + let tx = PersistentTransaction(txid: Data(repeating: 3, count: 32), transactionData: Data(), netAmount: 40) + let change = PersistentTxo(transaction: tx, vout: 0, amount: 40, address: "", height: 1) + change.walletId = walletId + tx.outputs = [change] + XCTAssertEqual(tx.netAmount(for: walletId), 40) + tx.pendingInputs = [PersistentPendingInput( + outpoint: Data(repeating: 9, count: 36), inputIndex: 0, + spendingTxid: tx.txid, spendingTransaction: tx, walletId: walletId + )] + XCTAssertNil(tx.netAmount(for: walletId), "a missing input makes the stored amount provisional") + } + + func testShouldNotCountAnotherLocalWalletsUnlinkedOutputForTheSender() throws { + let container = try DashModelContainer.createInMemory() + let context = container.mainContext + let senderId = Data(repeating: 1, count: 32) + let receiver = PersistentWallet(walletId: Data(repeating: 2, count: 32), network: .testnet) + let receiverAccount = PersistentAccount( + wallet: receiver, accountType: 0, accountIndex: 0, accountTypeName: "Standard BIP44 Account" + ) + // P2PKH to pubkey hash 0x05 x 20 on testnet: B's address with no TXO row yet. + let receiverAddress = PersistentCoreAddress( + address: "yLmzEvw3frCPS4cyRmFFeKbt64fUPzMwFh", poolTypeTag: 0, addressIndex: 0, derivationPath: "" + ) + receiverAddress.account = receiverAccount + context.insert(PersistentWallet(walletId: senderId, network: .testnet)) + context.insert(receiver) + context.insert(receiverAccount) + context.insert(receiverAddress) + var bytes = Data([2, 0, 0, 0, 1]) + bytes.append(senderId) + bytes.append(contentsOf: [0, 0, 0, 0, 0, 255, 255, 255, 255, 1]) + withUnsafeBytes(of: UInt64(40).littleEndian) { bytes.append(contentsOf: $0) } + bytes.append(contentsOf: [25, 0x76, 0xa9, 0x14] + [UInt8](repeating: 5, count: 20) + [0x88, 0xac]) + bytes.append(contentsOf: [0, 0, 0, 0]) + let spender = PersistentTransaction( + txid: Data(repeating: 3, count: 32), transactionData: bytes, direction: 1, netAmount: -100 + ) + let coin = input(100) + coin.spendingTransaction = spender + context.insert(coin) + context.insert(spender) + try context.save() + let handler = PlatformWalletPersistenceHandler(modelContainer: container, network: .testnet) + XCTAssertFalse(handler.loadWalletList().errored) + let row = try XCTUnwrap(ModelContext(container).fetch(FetchDescriptor()).first { $0.txid == spender.txid }) + XCTAssertEqual(row.netAmount, -100, "the receiving wallet's credit is not the sender's") + XCTAssertEqual(row.netAmount(for: senderId), -100) + } + func testShouldExcludePersistedContactOutputsFromOwnedAccounting() throws { let container = try DashModelContainer.createInMemory() let context = container.mainContext