From 63a7f8e09d17f7e6ffd51d6953d87aedbf840c08 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:31:07 +0000 Subject: [PATCH 01/13] fix(platform-wallet-storage): skip unreadable records in history repair V019 and the per-round history repair decoded stored records strictly, so one undecodable record_blob aborted the migration (blocking every open, Recovery included) and made each later write of that txid fail. History repair is best-effort accounting: unreadable stored data is now logged and skipped. Each record is repaired inside its own savepoint so a skipped one leaves no partial writes, the corrupt row is left untouched, and a later write of the same transaction replaces it. Co-Authored-By: Claude Opus 5.5 --- .../src/sqlite/schema/core_history.rs | 113 ++++++++++++++---- .../src/sqlite/schema/core_state.rs | 51 ++++++-- 2 files changed, 133 insertions(+), 31 deletions(-) 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..c74636b2117 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,15 +24,17 @@ 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 { - return Err(WalletStorageError::blob_decode( - "same transaction id has different raw transaction bodies", - )); + // A txid commits to its body, so the stored copy is corrupt; the + // incoming record replaces it rather than wedging every later write. + tracing::warn!( + txid = %incoming.txid, + "stored transaction body disagrees with its txid; replacing it" + ); + return Ok(merged); } let mut inputs: BTreeMap<_, _> = previous .input_details @@ -61,6 +63,57 @@ pub(super) fn preserve_known_details( Ok(merged) } +/// Read a stored record as repair evidence; an unreadable one is no evidence. +/// +/// History repair is best-effort accounting and must never block opening the +/// database or storing new state, so undecodable or drifted rows are skipped. +fn prior_record( + conn: &Connection, + wallet_id: &WalletId, + txid: &Txid, +) -> Result, WalletStorageError> { + match core_state::get_tx_record(conn, wallet_id, txid, &LoadCtx::recovery()) { + Err(error) if is_unreadable(&error) => { + tracing::warn!(%txid, %error, "skipping unreadable transaction record in history repair"); + Ok(None) + } + result => result, + } +} + +/// 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 { .. } + ) +} + +/// Run one record's repair atomically; unreadable stored data skips it instead of failing. +fn repair_best_effort( + tx: &Transaction<'_>, + txid: &Txid, + repair: impl FnOnce() -> Result<(), WalletStorageError>, +) -> Result<(), WalletStorageError> { + tx.execute_batch("SAVEPOINT core_history_repair")?; + let result = repair(); + if result.is_err() { + tx.execute_batch("ROLLBACK TO core_history_repair")?; + } + tx.execute_batch("RELEASE core_history_repair")?; + match result { + Err(error) if is_unreadable(&error) => { + tracing::warn!(%txid, %error, "skipping history repair of unreadable stored data"); + Ok(()) + } + result => result, + } +} + /// Index raw inputs independently of when their ownership becomes known. pub(super) fn index_record( tx: &Transaction<'_>, @@ -106,7 +159,7 @@ pub(super) fn apply( } let network = network(tx, wallet_id)?; for txid in affected { - repair_record(tx, wallet_id, &txid, network)?; + repair_best_effort(tx, &txid, || repair_record(tx, wallet_id, &txid, network))?; } Ok(()) } @@ -165,8 +218,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)?; @@ -295,20 +347,37 @@ fn repair_record( /// 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)?)?; + // Keys are collected first: savepoint rollbacks must not race an 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()? { + let key = (|| { + 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)?; + Ok::<_, WalletStorageError>((wallet_id, Txid::from_slice(&txid)?)) + })(); + match key { + Ok(key) => keys.push(key), + Err(error) if is_unreadable(&error) => { + tracing::warn!(%error, "skipping unreadable transaction key in history migration"); + } + Err(error) => return Err(error), + } } } + for (wallet_id, txid) in keys { + repair_best_effort(tx, &txid, || { + let Some(record) = prior_record(tx, &wallet_id, &txid)? else { + return Ok(()); + }; + index_record(tx, &wallet_id, &record)?; + repair_record(tx, &wallet_id, &txid, network(tx, &wallet_id)?) + })?; + } Ok(()) } 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..9e53e8867e3 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 @@ -1600,7 +1600,7 @@ mod tests { } #[test] - fn should_roll_back_history_migration_on_corrupt_record() { + fn should_skip_corrupt_record_during_history_migration() { 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(); @@ -1611,23 +1611,56 @@ mod tests { ) .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()); - let tables: i64 = conn + crate::sqlite::migrations::run(&mut conn) + .expect("one unreadable record must not block opening the database"); + let version: i64 = conn .query_row( - "SELECT count(*) FROM sqlite_master WHERE name = 'core_transaction_inputs'", + "SELECT max(version) FROM refinery_schema_history", [], |r| r.get(0), ) .unwrap(); - assert_eq!(tables, 0); - let version: i64 = conn + assert!(version >= 19); + let stored: Vec = conn .query_row( - "SELECT max(version) FROM refinery_schema_history", - [], + "SELECT record_blob FROM core_transactions WHERE wallet_id = ?1", + params![&wallet_id[..]], |r| r.get(0), ) .unwrap(); - assert_eq!(version, 18); + assert_eq!(stored, vec![0xff], "the unreadable row is left untouched"); + } + + #[test] + fn should_heal_corrupt_record_when_the_transaction_is_stored_again() { + 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(); + apply( + &tx, + &wallet_id, + &CoreChangeSet { + records: vec![record.clone()], + ..Default::default() + }, + ) + .expect("a corrupt stored copy must not wedge later writes"); + let healed = get_tx_record(&tx, &wallet_id, &record.txid, &LoadCtx::strict()) + .unwrap() + .unwrap(); + assert_eq!(healed.txid, record.txid); } #[test] From 740bbdd18af4159457721d395c750859483e1cc7 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:34:29 +0000 Subject: [PATCH 02/13] fix(platform-wallet-storage): keep history repairs reversible V019 and the per-round history repair rewrote stored transaction records in place and marked owned inputs spent without saying which transaction spent them, so a wrong repair could not be undone after the one-off pre-migration backup. Before its first rewrite, a record's original blob is now copied verbatim into the append-only core_transaction_record_originals table (V019 is unreleased and edited in place), and a repair's spent mark records its spender in spent_in_txid. Automatic release after a reorg is left as a TODO pending verification of upstream reorg handling. Co-Authored-By: Claude Opus 5.5 --- .../V019__core_transaction_accounting.rs | 12 ++++- .../src/sqlite/persister.rs | 4 ++ .../src/sqlite/schema/core_history.rs | 22 ++++++++- .../src/sqlite/schema/core_state.rs | 49 +++++++++++++++++-- .../tests/sqlite_schema_pinning.rs | 2 +- 5 files changed, 81 insertions(+), 8 deletions(-) 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/persister.rs b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs index 28623c535e2..31679e8b0fc 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs @@ -2418,6 +2418,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 +2462,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/schema/core_history.rs b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs index c74636b2117..b9cb1354098 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 @@ -242,11 +242,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() ], )?; } @@ -333,6 +343,14 @@ fn repair_record( 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![ 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 9e53e8867e3..be2acb222c5 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,13 +1627,26 @@ 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_skip_corrupt_record_during_history_migration() { 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)", @@ -1805,7 +1848,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_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 From 318d874745eb57130f35a104991578b0bbf0ebae Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:37:46 +0000 Subject: [PATCH 03/13] refactor(platform-wallet-storage): freeze the V019 history repair V019 called the live per-round history repair, so any later edit to that code (or to the queried tables) would silently change what V019 does for users upgrading from V018 or earlier. Move V019's repair into the migrations-local legacy_v019 module, following the legacy_v008 precedent, with its own record reader. Only the low-level blob codec stays shared; TransactionRecord's encoding is owned upstream. A pinning test fixes V019's result on a V018-shaped database: the repaired record bytes, the preserved original, the input index and the spent marks. Co-Authored-By: Claude Opus 5.5 --- .../src/sqlite/migrations.rs | 3 +- .../src/sqlite/migrations/legacy_v019.rs | 550 ++++++++++++++++++ .../src/sqlite/schema/core_history.rs | 37 -- 3 files changed, 552 insertions(+), 38 deletions(-) create mode 100644 packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs 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..9564fa8850f --- /dev/null +++ b/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs @@ -0,0 +1,550 @@ +//! 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 as repair evidence; an unreadable one is no evidence. +fn prior_record( + tx: &Transaction<'_>, + wallet_id: &WalletId, + txid: &Txid, +) -> Result, WalletStorageError> { + let read = (|| { + 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)?; + Ok(Some(record).filter(|record| record.txid == *txid)) + })(); + match read { + Err(error) if is_unreadable(&error) => { + tracing::warn!(%txid, %error, "skipping unreadable transaction record in history migration"); + Ok(None) + } + result => result, + } +} + +/// 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 { .. } + ) +} + +/// Run one record's repair atomically; unreadable stored data skips it instead of failing. +fn repair_best_effort( + tx: &Transaction<'_>, + txid: &Txid, + repair: impl FnOnce() -> Result<(), WalletStorageError>, +) -> Result<(), WalletStorageError> { + tx.execute_batch("SAVEPOINT core_history_repair")?; + let result = repair(); + if result.is_err() { + tx.execute_batch("ROLLBACK TO core_history_repair")?; + } + tx.execute_batch("RELEASE core_history_repair")?; + match result { + Err(error) if is_unreadable(&error) => { + tracing::warn!(%txid, %error, "skipping history repair of unreadable stored data"); + Ok(()) + } + result => result, + } +} + +/// 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, + txid: &Txid, + network: dashcore::Network, +) -> Result<(), WalletStorageError> { + let Some(mut record) = prior_record(tx, wallet_id, txid)? else { + return Ok(()); + }; + 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: savepoint rollbacks must not race an 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()? { + let key = (|| { + 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)?; + Ok::<_, WalletStorageError>((wallet_id, Txid::from_slice(&txid)?)) + })(); + match key { + Ok(key) => keys.push(key), + Err(error) if is_unreadable(&error) => { + tracing::warn!(%error, "skipping unreadable transaction key in history migration"); + } + Err(error) => return Err(error), + } + } + } + for (wallet_id, txid) in keys { + repair_best_effort(tx, &txid, || { + let Some(record) = prior_record(tx, &wallet_id, &txid)? else { + return Ok(()); + }; + index_record(tx, &wallet_id, &record)?; + repair_record(tx, &wallet_id, &txid, network(tx, &wallet_id)?) + })?; + } + 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/schema/core_history.rs b/packages/rs-platform-wallet-storage/src/sqlite/schema/core_history.rs index b9cb1354098..091afa8b073 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 @@ -362,40 +362,3 @@ fn repair_record( } Ok(()) } - -/// Backfill the input index and correct existing history in the migration transaction. -pub(crate) fn migrate(tx: &Transaction<'_>) -> Result<(), WalletStorageError> { - // Keys are collected first: savepoint rollbacks must not race an 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()? { - let key = (|| { - 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)?; - Ok::<_, WalletStorageError>((wallet_id, Txid::from_slice(&txid)?)) - })(); - match key { - Ok(key) => keys.push(key), - Err(error) if is_unreadable(&error) => { - tracing::warn!(%error, "skipping unreadable transaction key in history migration"); - } - Err(error) => return Err(error), - } - } - } - for (wallet_id, txid) in keys { - repair_best_effort(tx, &txid, || { - let Some(record) = prior_record(tx, &wallet_id, &txid)? else { - return Ok(()); - }; - index_record(tx, &wallet_id, &record)?; - repair_record(tx, &wallet_id, &txid, network(tx, &wallet_id)?) - })?; - } - Ok(()) -} From 283acf7012b89e36d8dcb252a69537a09b818fef Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:42:17 +0000 Subject: [PATCH 04/13] fix(platform-wallet-storage): restore unconfirmed spend reservations on load Load replayed only block-confirmed records, so a mempool or InstantSend spend never reached the account's spent set. A redelivered funding transaction (rescan, reorg re-connect, or the funding confirming after a restart) then re-credited the output that spend reserves, and the round wrote it back to SQLite as unspent. Replay unconfirmed records after the confirmed ones, parents before children, so the checker skips re-crediting reserved outputs. The persisted-unspent filter still keeps their own outputs from being credited unless persistence holds them. The regression tests now redeliver funding after reload, for confirmed and still-unconfirmed funding. Co-Authored-By: Claude Opus 5.5 --- .../src/sqlite/persister.rs | 2 +- .../src/sqlite/rehydrate.rs | 59 ++++++++++++++----- .../tests/sqlite_spent_rehydration.rs | 38 +++++++++++- 3 files changed, 81 insertions(+), 18 deletions(-) diff --git a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs index 31679e8b0fc..69967fa465a 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs @@ -1850,7 +1850,7 @@ fn load_one_wallet( })?; } let (wallet, wallet_info) = - super::rehydrate::restore_confirmed_transactions(wallet_info, wallet, core_state.records) + super::rehydrate::restore_recorded_transactions(wallet_info, wallet, core_state.records) .map_err(PersistenceError::from)?; Ok(platform_wallet::changeset::ClientWalletStartState { wallet, diff --git a/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs b/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs index db27b6106ce..fbf892b63be 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,8 +421,11 @@ 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( +/// 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( mut wallet_info: ManagedWalletInfo, mut wallet: Wallet, records: Vec, @@ -437,22 +441,14 @@ pub(crate) fn restore_confirmed_transactions( account.utxos.keys().map(move |outpoint| (*outpoint, owner)) }) .collect(); - let mut confirmed: Vec<_> = records - .into_iter() - .filter(|record| record.block_info().is_some()) - .collect(); - if confirmed.is_empty() { + if records.is_empty() { return Ok((wallet, wallet_info)); } - confirmed.sort_by_key(|record| { - record - .block_info() - .map(|block| (block.height(), block.position())) - }); + let replay = replay_order(records); // The checker only mutates in-memory state; no network requests or persistence. dash_async::block_on(async move { - for record in confirmed { + for record in replay { wallet_info .check_core_transaction( &record.transaction, @@ -505,6 +501,39 @@ pub(crate) fn restore_confirmed_transactions( .map_err(WalletStorageError::CoreHistoryReplay) } +/// 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())) + }); + // 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) + }) + }); + 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; + } + ordered.extend(ready); + pending = blocked; + } + ordered +} + /// Account identity of a funds account, stable across replay mutations. fn funds_account_type( account: &key_wallet::managed_account::ManagedCoreFundsAccount, 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] From 05ade4e463a84bab355b5a5ebb165603ac68793e Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:45:30 +0000 Subject: [PATCH 05/13] fix(platform-wallet-storage): replay load history without an async bridge The wallet checker is async only by trait shape and never awaits, yet load drove it through dash_async::block_on, which spawns a thread and runtime on current-thread runtimes, and surfaced bridge failures through a new public WalletStorageError::CoreHistoryReplay variant (a breaking change, as the enum is exhaustive). Poll the replay once instead. If a future upstream checker ever suspends, load logs an error and keeps the pre-replay projection rather than failing the wallet; a unit test pins first-poll completion so such a change fails in CI. This removes the variant and the dash-async dependency. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 1 - .../rs-platform-wallet-storage/Cargo.toml | 2 - .../src/sqlite/error.rs | 9 +- .../src/sqlite/persister.rs | 9 +- .../src/sqlite/rehydrate.rs | 182 +++++++++++++----- .../tests/sqlite_error_classification.rs | 2 - 6 files changed, 138 insertions(+), 67 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a0b69124ee6..324364168c1 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 d506d239da1..fe491fcc76b 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 } @@ -226,7 +225,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/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/persister.rs b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs index 69967fa465a..d7aa54e8134 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs @@ -1849,9 +1849,12 @@ fn load_one_wallet( )) })?; } - let (wallet, wallet_info) = - super::rehydrate::restore_recorded_transactions(wallet_info, wallet, core_state.records) - .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, diff --git a/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs b/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs index fbf892b63be..b2f578c74ef 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/rehydrate.rs @@ -426,10 +426,10 @@ pub fn apply_persisted_core_state( /// 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( - mut wallet_info: ManagedWalletInfo, - mut wallet: Wallet, + 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 @@ -442,63 +442,85 @@ pub(crate) fn restore_recorded_transactions( }) .collect(); if records.is_empty() { - return Ok((wallet, wallet_info)); + 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); - // The checker only mutates in-memory state; no network requests or persistence. - dash_async::block_on(async move { + // 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, - &mut wallet, - true, - false, - ) + .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() - .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(); - for account in wallet_info.accounts.all_funding_accounts_mut() { + 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); - 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(); - (wallet, wallet_info) - }) - .map_err(WalletStorageError::CoreHistoryReplay) + let placed = &placed; + account.utxos.keys().filter_map(move |outpoint| { + placed + .get(outpoint) + .filter(|parked| **parked != owner) + .map(|parked| (*outpoint, *parked)) + }) + }) + .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); + } + 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, + } } /// Order records as the chain would deliver them: confirmed by block position, @@ -3285,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/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"), } } From bc16c1760dc4391d2efc24b2f9f0307909ef6f05 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:46:42 +0000 Subject: [PATCH 06/13] fix(swift-sdk): classify repaired direction like the Rust repair The Swift accounting reconcile and the SQLite history repair disagreed on direction: Swift reported an asset lock as internal even when an output left the wallet, and a spend with no remaining outputs as internal. Swift now uses the Rust repair's rule (internal only when nothing leaves the wallet and something stays in it, or an asset lock burns into Platform). The Rust rule moves into a pure helper, and both sides test the same case table. The upstream rust-dashcore recompute stays out of scope. Co-Authored-By: Claude Opus 5.5 --- .../src/sqlite/schema/core_history.rs | 68 ++++++++++++++++--- .../Models/PersistentTransaction.swift | 5 +- .../TransactionAccountingTests.swift | 26 ++++++- 3 files changed, 88 insertions(+), 11 deletions(-) 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 091afa8b073..d7aeacba5ce 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 @@ -330,15 +330,12 @@ 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)?; @@ -362,3 +359,56 @@ fn repair_record( } Ok(()) } + +/// 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}" + ); + } + } +} diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift index 0fe3efe8f2e..c6d11db2591 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift @@ -287,10 +287,13 @@ 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 if allOutputsOwned && (!ownedOutputAmounts.isEmpty || isAssetLock) { direction = 2 } else { direction = 1 } return (received - spent, direction) } diff --git a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index ccde14865f4..cdc6e85de00 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 From e5fc327055e218ab1ff180be1138e5920830b3c1 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:47:17 +0000 Subject: [PATCH 07/13] test(swift-sdk): model asset-lock burns as OP_RETURN outputs The asset-lock fixtures used an empty-script output as the burn, which the aligned direction rule rightly treats as leaving the wallet. Use a real OP_RETURN so the fixtures match what an asset lock carries. Co-Authored-By: Claude Opus 5.5 --- .../SwiftDashSDKTests/TransactionAccountingTests.swift | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index cdc6e85de00..fa60ce5f72a 100644 --- a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift +++ b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift @@ -78,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) @@ -86,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 } @@ -266,7 +268,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) @@ -291,7 +293,7 @@ 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) From 4a30a5b3344aef2fad8444f3fe26c07e44c1c9d9 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:48:29 +0000 Subject: [PATCH 08/13] fix(swift-sdk): scope per-wallet transaction amounts netAmount(for:) returned the stored scalar through its single-wallet fast path before checking for unresolved inputs, and the accounting reconcile counted an address-matched output of any local wallet, so a transfer to another local wallet whose TXO was not linked yet showed the sender only part of what it paid. Unresolved inputs now make every amount provisional, and an unlinked address-matched output counts only when it belongs to a spending wallet. Co-Authored-By: Claude Opus 5.5 --- .../Models/PersistentTransaction.swift | 3 +- .../PlatformWalletPersistenceHandler.swift | 7 ++- .../TransactionAccountingTests.swift | 52 +++++++++++++++++++ 3 files changed, 60 insertions(+), 2 deletions(-) diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift index c6d11db2591..373123fcf0e 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 } diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift index 354e2cf32e8..09e8d0cce81 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift @@ -6837,6 +6837,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 } @@ -6846,7 +6850,8 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { let owner: PersistentCoreAddress? if let cached = roundIndex?.coreAddressesByAddress[address] { owner = cached } 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 } } diff --git a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index fa60ce5f72a..581100392cf 100644 --- a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift +++ b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift @@ -300,6 +300,58 @@ final class TransactionAccountingTests: XCTestCase { 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 From 8e5bcff11919388cc7238c0bebb4e65bf39a9501 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:48:57 +0000 Subject: [PATCH 09/13] fix(swift-sdk): keep asset-lock debits when no input is linked yet The guard that keeps a funded asset lock's Core debit and fee from being overwritten by a context-only zero update required an already-linked owned input. With every prevout still pending, the synthetic update erased the debit and fee, and reconcile could not restore them. A stored negative amount is itself the proof the wallet funded the lock, so the guard now keys on it. Co-Authored-By: Claude Opus 5.5 --- .../PlatformWalletPersistenceHandler.swift | 3 ++- .../TransactionAccountingTests.swift | 20 +++++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift index 09e8d0cce81..e3bebdd20be 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift @@ -2547,8 +2547,9 @@ 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. + // A stored debit is itself the proof we funded it: its inputs may not be linked yet. let preserveLockAccounting = tx.transaction_type_kind == 6 && tx.net_amount == 0 && !tx.has_fee - && record.netAmount != 0 && record.inputs.contains(where: Self.isWalletOwnedTxo) + && record.netAmount < 0 if !preserveLockAccounting { record.direction = tx.direction } if let typeName = tx.transaction_type { record.transactionType = String(cString: typeName) diff --git a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index 581100392cf..688c0f7c844 100644 --- a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift +++ b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift @@ -284,6 +284,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) From ec8b9458fa7e2983503b836ba1e894e4ea4e2453 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:49:46 +0000 Subject: [PATCH 10/13] fix(swift-sdk): never block wallet restore on load-time accounting loadWalletList's history accounting pass returned a load failure on any fetch, reconcile or save error, so a problem in display-only amounts left every wallet unrestored. It also fetched a PersistentCoreAddress per unlinked output. A failed pass now rolls back its own edits, logs a distinct event and lets the restore continue; addresses are read once into a lookup. A persisted completion marker is deferred (TODO) because it needs a SwiftData shape change. Co-Authored-By: Claude Opus 5.5 --- .../PlatformWalletPersistenceHandler.swift | 23 +++++++++++++++---- .../TransactionAccountingTests.swift | 12 ++++++++++ 2 files changed, 30 insertions(+), 5 deletions(-) diff --git a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift index e3bebdd20be..940fc664c30 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift @@ -6811,7 +6811,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 @@ -6850,6 +6851,7 @@ 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, spendingWallets.contains(account.wallet.walletId) { @@ -6918,6 +6920,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)) @@ -6929,15 +6935,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/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift index 688c0f7c844..26b690a985c 100644 --- a/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift +++ b/packages/swift-sdk/SwiftTests/SwiftDashSDKTests/TransactionAccountingTests.swift @@ -165,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 From b71f35a469c91ef2530bcb90ca6be5c99d5191cb Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:52:55 +0000 Subject: [PATCH 11/13] chore: address low-severity PR #5150 review findings - Transaction views treat an unresolved per-wallet amount as unavailable everywhere, so the fee no longer shows beside "Amount unavailable". - A failed accounting reconcile no longer fails the persistence round and is reported as persistence_transaction_accounting_failed, not save_failed. - Name the Swift transaction-kind and direction values instead of 1/3/6. - Rename the verdict test that never covered Uncredited and add one that does; V019's confirmed-spend marking is pinned by the V019 fixture test. - Drop the hand-written Unreleased section from the generated CHANGELOG. - Mark the divergent record-coalescing helpers with a TODO. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 6 -- .../src/sqlite/schema/core_state.rs | 56 ++++++++++++++++++- .../src/changeset/changeset.rs | 4 ++ .../Models/PersistentTransaction.swift | 23 ++++++-- .../PlatformWalletPersistenceHandler.swift | 17 +++++- .../Core/Views/TransactionDetailView.swift | 6 +- .../Core/Views/TransactionListView.swift | 6 +- 7 files changed, 99 insertions(+), 19 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 198da619b6e..d4ae8423857 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,9 +1,3 @@ -## Unreleased - -### Fixed - -- **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.4](https://github.com/dashpay/platform/compare/v4.2.0-beta.3...v4.2.0-beta.4) (2026-09-24) 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 be2acb222c5..54967106037 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 @@ -1707,7 +1707,7 @@ mod tests { } #[test] - fn should_keep_uncredited_outputs_spent() { + 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(); @@ -1745,6 +1745,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}; 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 373123fcf0e..f79cbbebdb2 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/Persistence/Models/PersistentTransaction.swift @@ -257,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. @@ -292,10 +293,11 @@ public final class PersistentTransaction { // 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 allOutputsOwned && (!ownedOutputAmounts.isEmpty || isAssetLock) { 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) } @@ -427,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 940fc664c30..2893e42c4a2 100644 --- a/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift +++ b/packages/swift-sdk/Sources/SwiftDashSDK/PlatformWallet/PlatformWalletPersistenceHandler.swift @@ -2548,7 +2548,8 @@ public final class PlatformWalletPersistenceHandler: @unchecked Sendable { record.blockHash = blockHashBytes.allSatisfy { $0 == 0 } ? nil : blockHashBytes // A context-only recovery record has zero accounting; a funded asset lock burns Core value. // A stored debit is itself the proof we funded it: its inputs may not be linked yet. - let preserveLockAccounting = tx.transaction_type_kind == 6 && tx.net_amount == 0 && !tx.has_fee + 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 { @@ -3428,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( @@ -6863,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 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) From 338067889e4c01651eb1742d5ba83e09955c1834 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:59:42 +0000 Subject: [PATCH 12/13] fix(platform-wallet-storage): fail on corrupt history unless a resync restores it Rework of 63a7f8e09d, which skipped undecodable records everywhere. Normal operation is strict again: the per-round repair and preserve_known_details treat a corrupt stored record as an error. V019 (which has no recovery mode) drops an undecodable record only when a Core resync re-delivers it: the row has a block height, so a filter rescan finds it again. It then deletes the row and its input-index rows and lowers that wallet's synced_height to just below its birth height; load hands that checkpoint to SPV, which rescans from birth and re-records the transaction. Spent marks and outputs it already produced are kept until the rescan confirms them. An unconfirmed corrupt record cannot be restored that way, so the migration still fails and rolls back. Co-Authored-By: Claude Opus 5.5 --- .../src/sqlite/migrations/legacy_v019.rs | 161 ++++++++++-------- .../src/sqlite/schema/core_history.rs | 58 +------ .../src/sqlite/schema/core_state.rs | 113 +++++++++--- 3 files changed, 184 insertions(+), 148 deletions(-) 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 index 9564fa8850f..e1bf7fcd5dc 100644 --- a/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs +++ b/packages/rs-platform-wallet-storage/src/sqlite/migrations/legacy_v019.rs @@ -21,40 +21,36 @@ 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 as repair evidence; an unreadable one is no evidence. -fn prior_record( +/// 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 read = (|| { - 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)?; - Ok(Some(record).filter(|record| record.txid == *txid)) - })(); - match read { - Err(error) if is_unreadable(&error) => { - tracing::warn!(%txid, %error, "skipping unreadable transaction record in history migration"); - Ok(None) - } - result => result, + 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. @@ -69,25 +65,56 @@ fn is_unreadable(error: &WalletStorageError) -> bool { ) } -/// Run one record's repair atomically; unreadable stored data skips it instead of failing. -fn repair_best_effort( +/// 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, - repair: impl FnOnce() -> Result<(), WalletStorageError>, + error: WalletStorageError, ) -> Result<(), WalletStorageError> { - tx.execute_batch("SAVEPOINT core_history_repair")?; - let result = repair(); - if result.is_err() { - tx.execute_batch("ROLLBACK TO core_history_repair")?; - } - tx.execute_batch("RELEASE core_history_repair")?; - match result { - Err(error) if is_unreadable(&error) => { - tracing::warn!(%txid, %error, "skipping history repair of unreadable stored data"); - Ok(()) - } - result => result, + 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. @@ -160,12 +187,11 @@ fn contact_only_script( fn repair_record( tx: &Transaction<'_>, wallet_id: &WalletId, - txid: &Txid, + mut record: TransactionRecord, network: dashcore::Network, ) -> Result<(), WalletStorageError> { - let Some(mut record) = prior_record(tx, wallet_id, txid)? else { - return Ok(()); - }; + 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(..) { @@ -305,37 +331,30 @@ fn repair_record( } pub(super) fn repair_history(tx: &Transaction<'_>) -> Result<(), WalletStorageError> { - // Keys are collected first: savepoint rollbacks must not race an open cursor. + // 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()? { - let key = (|| { - 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)?; - Ok::<_, WalletStorageError>((wallet_id, Txid::from_slice(&txid)?)) - })(); - match key { - Ok(key) => keys.push(key), - Err(error) if is_unreadable(&error) => { - tracing::warn!(%error, "skipping unreadable transaction key in history migration"); - } - Err(error) => return Err(error), - } + 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 { - repair_best_effort(tx, &txid, || { - let Some(record) = prior_record(tx, &wallet_id, &txid)? else { - return Ok(()); - }; - index_record(tx, &wallet_id, &record)?; - repair_record(tx, &wallet_id, &txid, network(tx, &wallet_id)?) - })?; + 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(()) } 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 d7aeacba5ce..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 @@ -28,13 +28,9 @@ pub(super) fn preserve_known_details( return Ok(merged); }; if previous.transaction != incoming.transaction { - // A txid commits to its body, so the stored copy is corrupt; the - // incoming record replaces it rather than wedging every later write. - tracing::warn!( - txid = %incoming.txid, - "stored transaction body disagrees with its txid; replacing it" - ); - return Ok(merged); + return Err(WalletStorageError::blob_decode( + "same transaction id has different raw transaction bodies", + )); } let mut inputs: BTreeMap<_, _> = previous .input_details @@ -63,55 +59,13 @@ pub(super) fn preserve_known_details( Ok(merged) } -/// Read a stored record as repair evidence; an unreadable one is no evidence. -/// -/// History repair is best-effort accounting and must never block opening the -/// database or storing new state, so undecodable or drifted rows are skipped. +/// Read a stored record strictly: corrupt history is an error, never skipped. fn prior_record( conn: &Connection, wallet_id: &WalletId, txid: &Txid, ) -> Result, WalletStorageError> { - match core_state::get_tx_record(conn, wallet_id, txid, &LoadCtx::recovery()) { - Err(error) if is_unreadable(&error) => { - tracing::warn!(%txid, %error, "skipping unreadable transaction record in history repair"); - Ok(None) - } - result => result, - } -} - -/// 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 { .. } - ) -} - -/// Run one record's repair atomically; unreadable stored data skips it instead of failing. -fn repair_best_effort( - tx: &Transaction<'_>, - txid: &Txid, - repair: impl FnOnce() -> Result<(), WalletStorageError>, -) -> Result<(), WalletStorageError> { - tx.execute_batch("SAVEPOINT core_history_repair")?; - let result = repair(); - if result.is_err() { - tx.execute_batch("ROLLBACK TO core_history_repair")?; - } - tx.execute_batch("RELEASE core_history_repair")?; - match result { - Err(error) if is_unreadable(&error) => { - tracing::warn!(%txid, %error, "skipping history repair of unreadable stored data"); - Ok(()) - } - result => result, - } + core_state::get_tx_record(conn, wallet_id, txid, &LoadCtx::strict()) } /// Index raw inputs independently of when their ownership becomes known. @@ -159,7 +113,7 @@ pub(super) fn apply( } let network = network(tx, wallet_id)?; for txid in affected { - repair_best_effort(tx, &txid, || repair_record(tx, wallet_id, &txid, network))?; + repair_record(tx, wallet_id, &txid, network)?; } Ok(()) } 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 54967106037..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 @@ -1642,20 +1642,45 @@ mod tests { ); } - #[test] - fn should_skip_corrupt_record_during_history_migration() { + /// 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; 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_transactions (wallet_id, txid, finalized, record_blob) VALUES (?1, ?2, 0, ?3)", params![&wallet_id[..], &[0u8;32][..], &[0xffu8][..]]).unwrap(); - crate::sqlite::migrations::run(&mut conn) - .expect("one unreadable record must not block opening the database"); + 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, 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'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(tables, 0); let version: i64 = conn .query_row( "SELECT max(version) FROM refinery_schema_history", @@ -1663,19 +1688,58 @@ mod tests { |r| r.get(0), ) .unwrap(); - assert!(version >= 19); - let stored: Vec = conn + assert_eq!(version, 18); + } + + #[test] + 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 record_blob FROM core_transactions WHERE wallet_id = ?1", + "SELECT last_processed_height, synced_height FROM core_sync_state WHERE wallet_id = ?1", params![&wallet_id[..]], - |r| r.get(0), + |r| Ok((r.get(0)?, r.get(1)?)), ) .unwrap(); - assert_eq!(stored, vec![0xff], "the unreadable row is left untouched"); + 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_heal_corrupt_record_when_the_transaction_is_stored_again() { + 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]; @@ -1691,19 +1755,18 @@ mod tests { ) .unwrap(); let tx = conn.transaction().unwrap(); - apply( - &tx, - &wallet_id, - &CoreChangeSet { - records: vec![record.clone()], - ..Default::default() - }, - ) - .expect("a corrupt stored copy must not wedge later writes"); - let healed = get_tx_record(&tx, &wallet_id, &record.txid, &LoadCtx::strict()) - .unwrap() - .unwrap(); - assert_eq!(healed.txid, record.txid); + assert!( + apply( + &tx, + &wallet_id, + &CoreChangeSet { + records: vec![record], + ..Default::default() + }, + ) + .is_err(), + "normal operation treats corrupt stored history as an error" + ); } #[test] From 10b15c9d3fe49a4d4467461950bb8abe7fb0d76b Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Tue, 29 Sep 2026 18:22:47 +0000 Subject: [PATCH 13/13] feat(wallet-storage): restore complete persisted Core wallet snapshots --- .../src/sqlite/persister.rs | 39 +- .../src/sqlite/schema/core_pool.rs | 187 +++++- .../tests/sqlite_persist_roundtrip.rs | 2 +- .../tests/sqlite_wallet_restore.rs | 613 ++++++++++++++++++ 4 files changed, 815 insertions(+), 26 deletions(-) create mode 100644 packages/rs-platform-wallet-storage/tests/sqlite_wallet_restore.rs diff --git a/packages/rs-platform-wallet-storage/src/sqlite/persister.rs b/packages/rs-platform-wallet-storage/src/sqlite/persister.rs index d7aa54e8134..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,25 +1831,17 @@ 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) - )) - })?; + .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, 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/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_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()); +}