commit dcf72db4cccf2efe53f8172f1363be0e14ffb974
parent dd478e9bbd684256d49b0ed36c6bd628588236e7
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Wed, 23 Sep 2026 22:15:21 +0200
perf(ui): incrementally warm derived data
Diffstat:
2 files changed, 364 insertions(+), 86 deletions(-)
diff --git a/src/adapters/ui_data_store.rs b/src/adapters/ui_data_store.rs
@@ -11,12 +11,12 @@ use rusqlite::{Connection, OptionalExtension, params, params_from_iter, types::V
use serde::Serialize;
use crate::{
- adapters::ui_index::{UiChainIndex, build_ui_chain_index},
+ adapters::ui_index::{UiChainIndex, build_ui_chain_index_for_blocks, projected_reward_outputs},
domain::{
AddressNetwork, Amount, Block, BurnLeaderRank, ChainSnapshot, Ledger,
MINE_RETARGET_WINDOW_BLOCKS, MINE_REWARD, OutPoint, Transaction, TransactionV2,
TransactionV2Domain, TxInput, TxOutput, decode_hex, encode_versioned_address, hex_encode,
- retarget_mine_difficulty_bits, reward_outputs_for_block,
+ retarget_mine_difficulty_bits,
},
};
@@ -305,8 +305,7 @@ impl SqliteUiDataStore {
pub fn project_snapshot(&self, snapshot: &ChainSnapshot, keep_metrics: bool) -> Result<()> {
let updated_at_ms = unix_ms();
validate_projection_snapshot_structure(snapshot)?;
- let ledger = Ledger::from_preverified_snapshot(snapshot.clone())
- .context("failed to rebuild ledger for UI UTXO projection")?;
+ let incremental_start = self.incremental_projection_start(snapshot)?;
if keep_metrics {
self.project_metrics_snapshot(snapshot, updated_at_ms)?;
@@ -314,21 +313,41 @@ impl SqliteUiDataStore {
self.clear_metrics()?;
}
- let ui_index = build_ui_chain_index(snapshot);
+ if incremental_start == Some(snapshot.blocks.len()) {
+ return Ok(());
+ }
+
+ let ledger = Ledger::from_preverified_snapshot(snapshot.clone())
+ .context("failed to rebuild ledger for UI UTXO projection")?;
+ let start = incremental_start.unwrap_or(0);
+ let blocks = &snapshot.blocks[start..];
+ let ui_index =
+ build_ui_chain_index_for_blocks(snapshot, &ledger, blocks, incremental_start.is_none());
let utxos = ledger.all_utxos();
- let wallet_transactions = wallet_transactions_from_snapshot(snapshot);
- let wallet_transactions_v2 =
- wallet_transactions_v2_from_snapshot(snapshot, &ledger.transaction_v2_domain()?);
- let leaderboards = build_ui_leaderboards(&utxos, &wallet_transactions)?;
+ let all_wallet_transactions = wallet_transactions_from_snapshot(snapshot);
+ let wallet_transactions = wallet_transactions_from_blocks(blocks);
+ let wallet_transactions_v2 = wallet_transactions_v2_from_blocks(
+ blocks,
+ &ledger.transaction_v2_domain()?,
+ AddressNetwork::from_profile_id(&snapshot.launch_profile.profile_id),
+ );
+ let leaderboards = build_ui_leaderboards(&utxos, &all_wallet_transactions)?;
self.with_connection_mut(|connection| {
let transaction = connection
.transaction()
.context("failed to start UI data projection transaction")?;
- replace_ui_chain_index(&transaction, &ui_index, updated_at_ms)?;
- replace_ui_utxos(&transaction, &utxos)?;
- replace_ui_wallet_transactions(&transaction, &wallet_transactions)?;
- replace_ui_wallet_transactions_v2(&transaction, &wallet_transactions_v2)?;
+ if incremental_start.is_some() {
+ append_ui_chain_index(&transaction, &ui_index, updated_at_ms)?;
+ replace_ui_utxos(&transaction, &utxos)?;
+ append_ui_wallet_transactions(&transaction, &wallet_transactions)?;
+ append_ui_wallet_transactions_v2(&transaction, &wallet_transactions_v2)?;
+ } else {
+ replace_ui_chain_index(&transaction, &ui_index, updated_at_ms)?;
+ replace_ui_utxos(&transaction, &utxos)?;
+ replace_ui_wallet_transactions(&transaction, &wallet_transactions)?;
+ replace_ui_wallet_transactions_v2(&transaction, &wallet_transactions_v2)?;
+ }
replace_ui_leaderboards(&transaction, &leaderboards)?;
transaction
.commit()
@@ -337,6 +356,30 @@ impl SqliteUiDataStore {
})
}
+ fn incremental_projection_start(&self, snapshot: &ChainSnapshot) -> Result<Option<usize>> {
+ self.with_connection(|connection| {
+ let meta = connection
+ .query_row(
+ "SELECT schema_version, tip_hash FROM ui_cache_meta WHERE id = 1",
+ [],
+ |row| Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?)),
+ )
+ .optional()
+ .context("failed to load UI projection metadata")?;
+ let Some((schema_version, tip_hash)) = meta else {
+ return Ok(None);
+ };
+ if schema_version != UI_CACHE_SCHEMA_VERSION {
+ return Ok(None);
+ }
+ Ok(snapshot
+ .blocks
+ .iter()
+ .position(|block| block.hash == tip_hash)
+ .map(|height| height + 1))
+ })
+ }
+
pub fn replace_metrics_for_snapshot(&self, snapshot: &ChainSnapshot) -> Result<()> {
validate_projection_snapshot_structure(snapshot)?;
self.project_metrics_snapshot(snapshot, unix_ms())
@@ -1010,6 +1053,76 @@ INSERT INTO ui_burn_leader_ranks (
Ok(())
}
+fn append_ui_chain_index(
+ transaction: &rusqlite::Transaction<'_>,
+ index: &UiChainIndex,
+ updated_at_ms: u64,
+) -> Result<()> {
+ let Some(tip_hash) = &index.tip_hash else {
+ return Ok(());
+ };
+ for (outpoint, output) in &index.outputs {
+ transaction
+ .execute(
+ r#"
+INSERT INTO ui_output_index (txid, output_index, address, amount)
+VALUES (?1, ?2, ?3, ?4)
+"#,
+ params![outpoint.txid, outpoint.index, output.address, output.amount],
+ )
+ .with_context(|| {
+ format!(
+ "failed to append UI output index row {}:{}",
+ outpoint.txid, outpoint.index
+ )
+ })?;
+ }
+ for (block_hash, ranks) in &index.burn_leader_ranks_by_hash {
+ transaction
+ .execute(
+ "INSERT INTO ui_burn_leader_rank_blocks (block_hash) VALUES (?1)",
+ params![block_hash],
+ )
+ .with_context(|| format!("failed to append UI burn leader rank block {block_hash}"))?;
+ for rank in ranks {
+ transaction
+ .execute(
+ r#"
+INSERT INTO ui_burn_leader_ranks (
+ block_hash, rank, ticket_id, owner, amount, eligible_from_height, eligible_until_height
+) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
+"#,
+ params![
+ block_hash,
+ rank.rank,
+ rank.ticket_id,
+ rank.owner,
+ rank.amount,
+ rank.eligible_from_height,
+ rank.eligible_until_height,
+ ],
+ )
+ .with_context(|| {
+ format!(
+ "failed to append UI burn leader rank {} for block {}",
+ rank.rank, block_hash
+ )
+ })?;
+ }
+ }
+ transaction
+ .execute(
+ r#"
+UPDATE ui_cache_meta
+SET schema_version = ?1, tip_hash = ?2, updated_at_ms = ?3
+WHERE id = 1
+"#,
+ params![UI_CACHE_SCHEMA_VERSION, tip_hash, updated_at_ms],
+ )
+ .context("failed to advance UI chain index metadata")?;
+ Ok(())
+}
+
fn replace_ui_utxos(
transaction: &rusqlite::Transaction<'_>,
utxos: &[(OutPoint, TxOutput)],
@@ -1076,6 +1189,43 @@ INSERT INTO ui_wallet_transactions (
Ok(())
}
+fn append_ui_wallet_transactions(
+ transaction: &rusqlite::Transaction<'_>,
+ rows: &[(String, WalletTransactionProjection)],
+) -> Result<()> {
+ for (address, row) in rows {
+ let transaction_json = serde_json::to_vec(&row.transaction)
+ .context("failed to serialize UI wallet transaction")?;
+ transaction
+ .execute(
+ r#"
+INSERT INTO ui_wallet_transactions (
+ address, sort_key, kind, signature, block_height, timestamp_ms, block_finalizer,
+ transaction_json
+) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
+"#,
+ params![
+ address,
+ row.sort_key,
+ row.kind,
+ row.transaction.signature(),
+ row.block_height,
+ row.timestamp_ms,
+ row.block_finalizer,
+ transaction_json,
+ ],
+ )
+ .with_context(|| {
+ format!(
+ "failed to append UI wallet transaction {} for {}",
+ row.transaction.signature(),
+ address
+ )
+ })?;
+ }
+ Ok(())
+}
+
fn replace_ui_wallet_transactions_v2(
transaction: &rusqlite::Transaction<'_>,
rows: &[(String, WalletTransactionV2Projection)],
@@ -1113,6 +1263,40 @@ INSERT INTO ui_wallet_transactions_v2 (
Ok(())
}
+fn append_ui_wallet_transactions_v2(
+ transaction: &rusqlite::Transaction<'_>,
+ rows: &[(String, WalletTransactionV2Projection)],
+) -> Result<()> {
+ for (address, row) in rows {
+ transaction
+ .execute(
+ r#"
+INSERT INTO ui_wallet_transactions_v2 (
+ address, sort_key, kind, transaction_id, block_height, timestamp_ms, block_finalizer,
+ envelope
+) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
+"#,
+ params![
+ address,
+ row.sort_key,
+ row.kind,
+ row.transaction_id,
+ row.block_height,
+ row.timestamp_ms,
+ row.block_finalizer,
+ row.envelope,
+ ],
+ )
+ .with_context(|| {
+ format!(
+ "failed to append UI wallet transaction v2 {} for {}",
+ row.transaction_id, address
+ )
+ })?;
+ }
+ Ok(())
+}
+
fn replace_ui_leaderboards(
transaction: &rusqlite::Transaction<'_>,
rows: &[(String, UiLeaderboardEntry)],
@@ -1606,17 +1790,12 @@ fn build_ui_leaderboards(
fn wallet_transactions_from_snapshot(
snapshot: &ChainSnapshot,
) -> Vec<(String, WalletTransactionProjection)> {
+ wallet_transactions_from_blocks(&snapshot.blocks)
+}
+
+fn wallet_transactions_from_blocks(blocks: &[Block]) -> Vec<(String, WalletTransactionProjection)> {
let mut rows = Vec::new();
- let mut running_ledger = snapshot.blocks.first().cloned().and_then(|genesis| {
- Ledger::from_preverified_snapshot(ChainSnapshot {
- genesis_allocations: snapshot.genesis_allocations.clone(),
- vdf_rounds: snapshot.vdf_rounds,
- launch_profile: snapshot.launch_profile.clone(),
- blocks: vec![genesis],
- })
- .ok()
- });
- for block in &snapshot.blocks {
+ for block in blocks {
for (index, transaction) in block.transactions.iter().rev().enumerate() {
push_wallet_transaction_projection(
&mut rows,
@@ -1625,16 +1804,8 @@ fn wallet_transactions_from_snapshot(
block.height as u128 * 10_000 + index as u128,
);
}
- let reward_committee = if block.height == 0 {
- Vec::new()
- } else {
- running_ledger
- .as_ref()
- .map(|ledger| ledger.burn_committee_for_block(block))
- .unwrap_or_default()
- };
let reward_outputs = if block.height > 0 {
- reward_outputs_for_block(block, &reward_committee)
+ projected_reward_outputs(block)
} else {
Vec::new()
};
@@ -1662,22 +1833,26 @@ fn wallet_transactions_from_snapshot(
},
));
}
- if block.height > 0 {
- if let Some(ledger) = running_ledger.as_mut() {
- let _ = ledger.apply_preverified_block_at(block.clone(), u64::MAX);
- }
- }
}
rows
}
+#[cfg(test)]
fn wallet_transactions_v2_from_snapshot(
snapshot: &ChainSnapshot,
expected_domain: &TransactionV2Domain,
) -> Vec<(String, WalletTransactionV2Projection)> {
let network = AddressNetwork::from_profile_id(&snapshot.launch_profile.profile_id);
+ wallet_transactions_v2_from_blocks(&snapshot.blocks, expected_domain, network)
+}
+
+fn wallet_transactions_v2_from_blocks(
+ blocks: &[Block],
+ expected_domain: &TransactionV2Domain,
+ network: AddressNetwork,
+) -> Vec<(String, WalletTransactionV2Projection)> {
let mut rows = Vec::new();
- for block in &snapshot.blocks {
+ for block in blocks {
for (index, envelope) in block.transactions_v2.iter().rev().enumerate() {
let Some(transaction) = decode_projected_transaction_v2(envelope, expected_domain)
else {
@@ -2339,6 +2514,102 @@ END;
}
#[test]
+ fn ui_projection_appends_without_rewriting_the_consistent_history() {
+ let dir = tempdir().unwrap();
+ let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap();
+ let (mut ledger, wallet) = test_ledger("incremental-ui-data");
+ append_test_block(&mut ledger, &wallet, 1_000);
+ store.project_snapshot(&ledger.snapshot(), true).unwrap();
+ let first_tip = ledger.tip_hash().to_string();
+
+ let connection = Connection::open(store.path()).unwrap();
+ connection
+ .execute_batch(
+ r#"
+CREATE TRIGGER protect_ui_output_prefix
+BEFORE DELETE ON ui_output_index
+BEGIN
+ SELECT RAISE(FAIL, 'consistent UI output history was rewritten');
+END;
+CREATE TRIGGER protect_ui_wallet_prefix
+BEFORE DELETE ON ui_wallet_transactions
+WHEN OLD.block_height <= 1
+BEGIN
+ SELECT RAISE(FAIL, 'consistent UI wallet history was rewritten');
+END;
+CREATE TRIGGER protect_ui_rank_prefix
+BEFORE DELETE ON ui_burn_leader_rank_blocks
+BEGIN
+ SELECT RAISE(FAIL, 'consistent UI rank history was rewritten');
+END;
+"#,
+ )
+ .unwrap();
+ drop(connection);
+
+ append_test_block(&mut ledger, &wallet, 2_000);
+ store.project_snapshot(&ledger.snapshot(), true).unwrap();
+
+ assert!(store.is_projected_to(ledger.tip_hash()).unwrap());
+ let connection = Connection::open(store.path()).unwrap();
+ let first_reward_is_retained = connection
+ .query_row(
+ "SELECT EXISTS(SELECT 1 FROM ui_output_index WHERE txid = ?1)",
+ [&first_tip],
+ |row| row.get::<_, bool>(0),
+ )
+ .unwrap();
+ let latest_reward_is_present = connection
+ .query_row(
+ "SELECT EXISTS(SELECT 1 FROM ui_output_index WHERE txid = ?1)",
+ [ledger.tip_hash()],
+ |row| row.get::<_, bool>(0),
+ )
+ .unwrap();
+ assert!(first_reward_is_retained);
+ assert!(latest_reward_is_present);
+ }
+
+ #[test]
+ fn ui_projection_falls_back_to_full_rebuild_for_a_reorg() {
+ let dir = tempdir().unwrap();
+ let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap();
+ let (mut base, wallet) = test_ledger("reorg-ui-data");
+ append_test_block(&mut base, &wallet, 1_000);
+ let mut first_branch = base.clone();
+ let mut second_branch = base;
+ append_test_block(&mut first_branch, &wallet, 2_000);
+ append_test_block(&mut second_branch, &wallet, 3_000);
+ let old_tip = first_branch.tip_hash().to_string();
+
+ store
+ .project_snapshot(&first_branch.snapshot(), true)
+ .unwrap();
+ store
+ .project_snapshot(&second_branch.snapshot(), true)
+ .unwrap();
+
+ assert!(store.is_projected_to(second_branch.tip_hash()).unwrap());
+ let connection = Connection::open(store.path()).unwrap();
+ let old_branch_output_remains = connection
+ .query_row(
+ "SELECT EXISTS(SELECT 1 FROM ui_output_index WHERE txid = ?1)",
+ [&old_tip],
+ |row| row.get::<_, bool>(0),
+ )
+ .unwrap();
+ let new_branch_output_is_present = connection
+ .query_row(
+ "SELECT EXISTS(SELECT 1 FROM ui_output_index WHERE txid = ?1)",
+ [second_branch.tip_hash()],
+ |row| row.get::<_, bool>(0),
+ )
+ .unwrap();
+ assert!(!old_branch_output_remains);
+ assert!(new_branch_output_is_present);
+ }
+
+ #[test]
fn incremental_metrics_match_the_ledger_across_a_difficulty_retarget() {
let dir = tempdir().unwrap();
let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap();
diff --git a/src/adapters/ui_index.rs b/src/adapters/ui_index.rs
@@ -1,9 +1,9 @@
use std::collections::BTreeMap;
use crate::domain::{
- AddressNetwork, Block, BurnLeaderRank, ChainSnapshot, Ledger, OutPoint, Transaction,
- TransactionV2, TxOutput, decode_hex, encode_versioned_address, genesis_allocation_outpoint,
- hex_encode, reward_outputs_for_block,
+ AddressNetwork, Block, BurnCommitteeMember, BurnLeaderRank, ChainSnapshot, Ledger, OutPoint,
+ Transaction, TransactionV2, TxOutput, decode_hex, encode_versioned_address,
+ genesis_allocation_outpoint, hex_encode, reward_outputs_for_block,
};
#[derive(Clone, Debug, Default, Eq, PartialEq)]
@@ -13,25 +13,33 @@ pub(crate) struct UiChainIndex {
pub(crate) burn_leader_ranks_by_hash: BTreeMap<String, Vec<BurnLeaderRank>>,
}
+#[cfg(test)]
pub(crate) fn build_ui_chain_index(snapshot: &ChainSnapshot) -> UiChainIndex {
+ let Ok(ledger) = Ledger::from_preverified_snapshot(snapshot.clone()) else {
+ return UiChainIndex::default();
+ };
+ build_ui_chain_index_for_blocks(snapshot, &ledger, &snapshot.blocks, true)
+}
+
+pub(crate) fn build_ui_chain_index_for_blocks(
+ snapshot: &ChainSnapshot,
+ ledger: &Ledger,
+ blocks: &[Block],
+ include_genesis_allocations: bool,
+) -> UiChainIndex {
UiChainIndex {
tip_hash: snapshot.blocks.last().map(|block| block.hash.clone()),
- outputs: known_chain_output_index(snapshot),
- burn_leader_ranks_by_hash: burn_leader_ranks_for_blocks(snapshot, &snapshot.blocks),
+ outputs: known_chain_output_index(snapshot, blocks, include_genesis_allocations),
+ burn_leader_ranks_by_hash: burn_leader_ranks_for_blocks(ledger, blocks),
}
}
pub(crate) fn burn_leader_ranks_for_blocks(
- snapshot: &ChainSnapshot,
+ ledger: &Ledger,
blocks: &[Block],
) -> BTreeMap<String, Vec<BurnLeaderRank>> {
- let Some(ranks_by_height) = Ledger::from_preverified_snapshot(snapshot.clone())
- .ok()
- .and_then(|ledger| {
- ledger
- .burn_leader_ranks_for_blocks(blocks.iter().map(|block| block.height))
- .ok()
- })
+ let Ok(ranks_by_height) =
+ ledger.burn_leader_ranks_for_blocks(blocks.iter().map(|block| block.height))
else {
return BTreeMap::new();
};
@@ -47,31 +55,28 @@ pub(crate) fn burn_leader_ranks_for_blocks(
.collect()
}
-fn known_chain_output_index(snapshot: &ChainSnapshot) -> BTreeMap<OutPoint, TxOutput> {
+fn known_chain_output_index(
+ snapshot: &ChainSnapshot,
+ blocks: &[Block],
+ include_genesis_allocations: bool,
+) -> BTreeMap<OutPoint, TxOutput> {
let mut outputs = BTreeMap::new();
let network = AddressNetwork::from_profile_id(&snapshot.launch_profile.profile_id);
- let mut running_ledger = snapshot.blocks.first().cloned().and_then(|genesis| {
- Ledger::from_preverified_snapshot(ChainSnapshot {
- genesis_allocations: snapshot.genesis_allocations.clone(),
- vdf_rounds: snapshot.vdf_rounds,
- launch_profile: snapshot.launch_profile.clone(),
- blocks: vec![genesis],
- })
- .ok()
- });
- for (address, amount) in &snapshot.genesis_allocations {
- if *amount == 0 {
- continue;
+ if include_genesis_allocations {
+ for (address, amount) in &snapshot.genesis_allocations {
+ if *amount == 0 {
+ continue;
+ }
+ outputs.insert(
+ genesis_allocation_outpoint(address),
+ TxOutput {
+ address: address.clone(),
+ amount: *amount,
+ },
+ );
}
- outputs.insert(
- genesis_allocation_outpoint(address),
- TxOutput {
- address: address.clone(),
- amount: *amount,
- },
- );
}
- for block in &snapshot.blocks {
+ for block in blocks {
for transaction in &block.transactions {
index_transaction_outputs(&mut outputs, transaction);
}
@@ -84,26 +89,28 @@ fn known_chain_output_index(snapshot: &ChainSnapshot) -> BTreeMap<OutPoint, TxOu
};
index_transaction_v2_outputs(&mut outputs, &transaction, &domain, network);
}
- let reward_committee = if block.height == 0 {
- Vec::new()
- } else {
- running_ledger
- .as_ref()
- .map(|ledger| ledger.burn_committee_for_block(block))
- .unwrap_or_default()
- };
- for (outpoint, output) in reward_outputs_for_block(block, &reward_committee) {
+ for (outpoint, output) in projected_reward_outputs(block) {
outputs.insert(outpoint, output);
}
- if block.height > 0 {
- if let Some(ledger) = running_ledger.as_mut() {
- let _ = ledger.apply_preverified_block_at(block.clone(), u64::MAX);
- }
- }
}
outputs
}
+pub(crate) fn projected_reward_outputs(block: &Block) -> Vec<(OutPoint, TxOutput)> {
+ let committee = block
+ .burn_bundle_section
+ .signatures
+ .iter()
+ .map(|signature| BurnCommitteeMember {
+ slot: signature.slot,
+ root: block.hash.clone(),
+ owner: signature.member.clone(),
+ weight: 0,
+ })
+ .collect::<Vec<_>>();
+ reward_outputs_for_block(block, &committee)
+}
+
fn index_transaction_v2_outputs(
outputs: &mut BTreeMap<OutPoint, TxOutput>,
transaction: &TransactionV2,