iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

commit 63999847d979e96a008c45b712ceb8d93b842e8c
parent 0e369580e4e99d9937c51e47328103d807013e30
Author: Joris Hartog <jorishartog@hotmail.com>
Date:   Mon, 24 Aug 2026 11:41:26 +0200

Make metrics projection incremental

Diffstat:
Msrc/adapters/http/api.rs | 8++++----
Msrc/adapters/http/metrics.rs | 44++++++++++++++++++++++++++++++++++++++++++--
Msrc/adapters/ui_data_store.rs | 1017+++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------------
Msrc/domain.rs | 1+
Msrc/domain/mine_policy.rs | 4++--
Msrc/main.rs | 30+++++++++++++++++++++++++++++-
6 files changed, 769 insertions(+), 335 deletions(-)

diff --git a/src/adapters/http/api.rs b/src/adapters/http/api.rs @@ -342,9 +342,6 @@ pub(super) async fn api_metrics( .ok() .and_then(Result::ok) .unwrap_or(false); - if !ready { - return Json(empty_metrics_response(enabled, true)); - } let store = state.ui_data_store.clone(); let rows = tokio::task::spawn_blocking(move || match query.limit { Some(limit) => store.load_recent_metrics(limit.clamp(1, DATASET_LIMIT)), @@ -361,7 +358,10 @@ pub(super) async fn api_metrics( .and_then(Result::ok) .map(metrics_leaderboards) .unwrap_or_default(); - Json(metrics_response(enabled, rows, leaderboards)) + if rows.is_empty() { + return Json(empty_metrics_response(enabled, !ready)); + } + Json(metrics_response(enabled, !ready, rows, leaderboards)) } fn empty_metrics_response(enabled: bool, preparing: bool) -> MetricsResponse { diff --git a/src/adapters/http/metrics.rs b/src/adapters/http/metrics.rs @@ -22,13 +22,14 @@ pub(super) fn network_health( pub(super) fn metrics_response( enabled: bool, + preparing: bool, rows: Vec<BlockMetricRow>, leaderboards: MetricsLeaderboards, ) -> MetricsResponse { let latest = rows.last().cloned(); MetricsResponse { enabled, - preparing: false, + preparing, latest, leaderboards, charts: vec![ @@ -304,7 +305,46 @@ fn median_peer_clock_offset(peers: &[PeerInfo], now_ms: u64) -> Option<i64> { #[cfg(test)] mod tests { - use super::{MempoolCounts, NetworkHealthLocalState, network_health_at}; + use crate::adapters::ui_data_store::BlockMetricRow; + + use super::{ + MempoolCounts, MetricsLeaderboards, NetworkHealthLocalState, metrics_response, + network_health_at, + }; + + #[test] + fn stale_metrics_remain_visible_while_the_latest_tip_is_preparing() { + let row = BlockMetricRow { + height: 7, + block_hash: "cached-tip".to_string(), + timestamp_ms: 1_000, + block_time_ms: Some(100), + mine_difficulty_bits: 10, + circulating_supply: 100, + known_wallet_addresses: 2, + transaction_count: 1, + transfer_count: 0, + burn_count: 1, + mine_count: 0, + burned_amount: 1, + total_burned_amount: 3, + fees_amount: 1, + reward_amount: 1, + vdf_rounds: 10, + finalizer_rank: 0, + }; + + let response = metrics_response( + true, + true, + vec![row.clone()], + MetricsLeaderboards::default(), + ); + + assert!(response.preparing); + assert_eq!(response.latest, Some(row)); + assert!(response.charts.iter().any(|chart| !chart.points.is_empty())); + } #[test] fn network_health_exposes_operator_chain_and_rejection_context() { diff --git a/src/adapters/ui_data_store.rs b/src/adapters/ui_data_store.rs @@ -13,8 +13,8 @@ use serde::Serialize; use crate::{ adapters::ui_index::{UiChainIndex, build_ui_chain_index}, domain::{ - Amount, Block, BurnLeaderRank, ChainSnapshot, Ledger, MINE_REWARD, OutPoint, Transaction, - TxInput, TxOutput, hex_hash, reward_outputs_for_block, + Amount, Block, BurnLeaderRank, ChainSnapshot, Ledger, MINE_RETARGET_WINDOW_BLOCKS, + MINE_REWARD, OutPoint, Transaction, TxInput, TxOutput, retarget_mine_difficulty_bits, }, }; @@ -39,6 +39,30 @@ CREATE TABLE IF NOT EXISTS block_metrics ( finalizer_rank INTEGER NOT NULL ); +CREATE TABLE IF NOT EXISTS metrics_cache_meta ( + id INTEGER PRIMARY KEY CHECK (id = 1), + schema_version INTEGER NOT NULL, + height INTEGER NOT NULL, + tip_hash TEXT NOT NULL, + updated_at_ms INTEGER NOT NULL +); + +CREATE TABLE IF NOT EXISTS metric_known_addresses ( + address TEXT PRIMARY KEY, + first_seen_height INTEGER NOT NULL +); + +CREATE TABLE IF NOT EXISTS ui_leaderboards ( + kind TEXT NOT NULL, + address TEXT NOT NULL, + amount INTEGER NOT NULL, + count INTEGER NOT NULL, + PRIMARY KEY (kind, address) +); + +CREATE INDEX IF NOT EXISTS idx_ui_leaderboards_rank +ON ui_leaderboards(kind, amount DESC, count DESC, address ASC); + CREATE TABLE IF NOT EXISTS ui_cache_meta ( id INTEGER PRIMARY KEY CHECK (id = 1), schema_version INTEGER NOT NULL, @@ -99,7 +123,8 @@ CREATE TABLE IF NOT EXISTS ui_burn_leader_rank_blocks ( ); "#; -const UI_CACHE_SCHEMA_VERSION: u32 = 2; +const UI_CACHE_SCHEMA_VERSION: u32 = 3; +const METRICS_CACHE_SCHEMA_VERSION: u32 = 1; #[derive(Clone, Debug, Eq, PartialEq, Serialize)] #[serde(rename_all = "camelCase")] @@ -201,8 +226,7 @@ impl SqliteUiDataStore { }) } - #[cfg(test)] - pub(crate) fn is_projected_to(&self, tip_hash: &str) -> Result<bool> { + pub fn is_projected_to(&self, tip_hash: &str) -> Result<bool> { self.with_connection(|connection| { let projected = connection .query_row( @@ -219,42 +243,48 @@ impl SqliteUiDataStore { }) } - pub(crate) fn metrics_are_projected_to(&self, tip_hash: &str) -> Result<bool> { + pub fn metrics_are_projected_to(&self, tip_hash: &str) -> Result<bool> { self.with_connection(|connection| { - connection + let projected = connection .query_row( - "SELECT EXISTS(SELECT 1 FROM block_metrics WHERE block_hash = ?1)", - [tip_hash], - |row| row.get(0), + "SELECT schema_version, tip_hash FROM metrics_cache_meta WHERE id = 1", + [], + |row| Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?)), ) - .context("failed to inspect metrics projection") + .optional() + .context("failed to inspect metrics projection")? + .is_some_and(|(schema_version, stored_tip_hash)| { + schema_version == METRICS_CACHE_SCHEMA_VERSION && stored_tip_hash == tip_hash + }); + Ok(projected) }) } 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")?; + + if keep_metrics { + self.project_metrics_snapshot(snapshot, updated_at_ms)?; + } else { + self.clear_metrics()?; + } + let ui_index = build_ui_chain_index(snapshot); - let utxos = Ledger::from_preverified_snapshot(snapshot.clone()) - .context("failed to rebuild ledger for UI UTXO projection")? - .all_utxos(); + let utxos = ledger.all_utxos(); let wallet_transactions = wallet_transactions_from_snapshot(snapshot); - let metrics = if keep_metrics { - Some(metrics_from_snapshot(snapshot)?) - } else { - None - }; + let leaderboards = build_ui_leaderboards(&utxos, &wallet_transactions)?; self.with_connection_mut(|connection| { let transaction = connection .transaction() .context("failed to start UI data projection transaction")?; - match metrics { - Some(metrics) => replace_metrics(&transaction, &metrics)?, - None => clear_metrics_in_transaction(&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_leaderboards(&transaction, &leaderboards)?; transaction .commit() .context("failed to commit UI data projection transaction")?; @@ -263,26 +293,12 @@ impl SqliteUiDataStore { } pub fn replace_metrics_for_snapshot(&self, snapshot: &ChainSnapshot) -> Result<()> { - let metrics = metrics_from_snapshot(snapshot)?; - self.with_connection_mut(|connection| { - let transaction = connection - .transaction() - .context("failed to start metrics transaction")?; - replace_metrics(&transaction, &metrics)?; - transaction - .commit() - .context("failed to commit metrics transaction")?; - Ok(()) - }) + validate_projection_snapshot_structure(snapshot)?; + self.project_metrics_snapshot(snapshot, unix_ms()) } pub fn clear_metrics(&self) -> Result<()> { - self.with_connection_mut(|connection| { - connection - .execute("DELETE FROM block_metrics", []) - .context("failed to delete block metrics")?; - Ok(()) - }) + self.with_connection_mut(clear_metrics_state) } pub fn clear_all(&self) -> Result<()> { @@ -299,6 +315,19 @@ impl SqliteUiDataStore { }) } + fn project_metrics_snapshot(&self, snapshot: &ChainSnapshot, updated_at_ms: u64) -> Result<()> { + self.with_connection_mut(|connection| { + let transaction = connection + .transaction() + .context("failed to start incremental metrics transaction")?; + project_metrics_in_transaction(&transaction, snapshot, updated_at_ms)?; + transaction + .commit() + .context("failed to commit incremental metrics transaction")?; + Ok(()) + }) + } + pub fn load_metrics(&self) -> Result<Vec<BlockMetricRow>> { self.with_connection(|connection| { let mut statement = connection @@ -464,103 +493,234 @@ fn load_balance_leaderboard( connection: &Connection, limit: usize, ) -> Result<Vec<UiLeaderboardEntry>> { + load_materialized_leaderboard(connection, "balances", limit, false) +} + +fn load_transaction_leaderboard( + connection: &Connection, + kind: &str, + limit: usize, +) -> Result<Vec<UiLeaderboardEntry>> { + load_materialized_leaderboard(connection, kind, limit, true) +} + +fn load_materialized_leaderboard( + connection: &Connection, + kind: &str, + limit: usize, + rank_by_count: bool, +) -> Result<Vec<UiLeaderboardEntry>> { + let order = if rank_by_count { + "amount DESC, count DESC, address ASC" + } else { + "amount DESC, address ASC" + }; let mut statement = connection - .prepare( + .prepare(&format!( r#" -SELECT address, SUM(amount) AS total_amount, COUNT(*) AS output_count -FROM ui_utxos -GROUP BY address -HAVING total_amount > 0 -ORDER BY total_amount DESC, address ASC -LIMIT ?1 -"#, - ) - .context("failed to prepare balance leaderboard query")?; +SELECT address, amount, count +FROM ui_leaderboards +WHERE kind = ?1 +ORDER BY {order} +LIMIT ?2 +"# + )) + .with_context(|| format!("failed to prepare {kind} leaderboard query"))?; let rows = statement - .query_map([limit as u64], |row| { + .query_map(params![kind, limit as u64], |row| { Ok(UiLeaderboardEntry { address: row.get(0)?, amount: row.get(1)?, count: row.get(2)?, }) }) - .context("failed to load balance leaderboard")?; + .with_context(|| format!("failed to load {kind} leaderboard"))?; rows.collect::<std::result::Result<Vec<_>, _>>() - .context("failed to read balance leaderboard rows") + .with_context(|| format!("failed to read {kind} leaderboard rows")) } -fn load_transaction_leaderboard( - connection: &Connection, - kind: &str, - limit: usize, -) -> Result<Vec<UiLeaderboardEntry>> { - let mut statement = connection - .prepare( +fn project_metrics_in_transaction( + transaction: &rusqlite::Transaction<'_>, + snapshot: &ChainSnapshot, + updated_at_ms: u64, +) -> Result<()> { + let tip = snapshot + .blocks + .last() + .context("cannot project metrics for empty chain snapshot")?; + let mut common_height = metrics_common_height(transaction, snapshot)?; + let mut previous = common_height + .map(|height| load_metric_at_height(transaction, height)) + .transpose()? + .flatten(); + if common_height.is_some() && previous.is_none() { + common_height = None; + } + + match common_height { + Some(height) => { + transaction + .execute("DELETE FROM block_metrics WHERE height > ?1", [height]) + .context("failed to truncate reorged block metrics")?; + transaction + .execute( + "DELETE FROM metric_known_addresses WHERE first_seen_height > ?1", + [height], + ) + .context("failed to truncate reorged metric addresses")?; + } + None => clear_metrics_in_transaction(transaction)?, + } + + let mut known_wallet_addresses = previous + .as_ref() + .map(|metric| metric.known_wallet_addresses) + .unwrap_or_default(); + if previous.is_none() { + for address in snapshot.genesis_allocations.keys() { + known_wallet_addresses = known_wallet_addresses + .checked_add(insert_metric_address(transaction, address, 0)?) + .context("known metric address count overflows")?; + } + } + + let start_height = common_height.map_or(0, |height| height.saturating_add(1)); + for block in snapshot + .blocks + .iter() + .filter(|block| block.height >= start_height) + { + known_wallet_addresses = known_wallet_addresses + .checked_add(index_metric_addresses(transaction, block)?) + .context("known metric address count overflows")?; + let mut metric = incremental_metric_for_block(snapshot, block, previous.as_ref())?; + metric.known_wallet_addresses = known_wallet_addresses; + insert_metric(transaction, &metric)?; + previous = Some(metric); + } + + transaction + .execute( r#" -SELECT address, transaction_json -FROM ui_wallet_transactions -WHERE kind = ?1 +INSERT INTO metrics_cache_meta (id, schema_version, height, tip_hash, updated_at_ms) +VALUES (1, ?1, ?2, ?3, ?4) +ON CONFLICT(id) DO UPDATE SET + schema_version = excluded.schema_version, + height = excluded.height, + tip_hash = excluded.tip_hash, + updated_at_ms = excluded.updated_at_ms "#, + params![ + METRICS_CACHE_SCHEMA_VERSION, + tip.height, + tip.hash, + updated_at_ms + ], ) - .with_context(|| format!("failed to prepare {kind} leaderboard query"))?; + .context("failed to update metrics cache metadata")?; + Ok(()) +} + +fn metrics_common_height( + transaction: &rusqlite::Transaction<'_>, + snapshot: &ChainSnapshot, +) -> Result<Option<u64>> { + let meta = transaction + .query_row( + "SELECT schema_version, height, tip_hash FROM metrics_cache_meta WHERE id = 1", + [], + |row| { + Ok(( + row.get::<_, u32>(0)?, + row.get::<_, u64>(1)?, + row.get::<_, String>(2)?, + )) + }, + ) + .optional() + .context("failed to load metrics cache metadata")?; + let Some((schema_version, stored_height, stored_hash)) = meta else { + return Ok(None); + }; + if schema_version != METRICS_CACHE_SCHEMA_VERSION { + return Ok(None); + } + + let candidate_height = stored_height.min(snapshot.blocks.len().saturating_sub(1) as u64); + if candidate_height == stored_height + && snapshot.blocks[candidate_height as usize].hash == stored_hash + { + return Ok(Some(candidate_height)); + } + + let mut statement = transaction + .prepare( + "SELECT height, block_hash FROM block_metrics WHERE height <= ?1 ORDER BY height DESC", + ) + .context("failed to prepare metrics common-ancestor query")?; let rows = statement - .query_map([kind], |row| { - Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?)) + .query_map([candidate_height], |row| { + Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)) }) - .with_context(|| format!("failed to load {kind} leaderboard"))?; - let mut entries = BTreeMap::<String, UiLeaderboardEntry>::new(); + .context("failed to query metrics common ancestor")?; for row in rows { - let (address, transaction_json) = - row.with_context(|| format!("failed to read {kind} leaderboard row"))?; - let transaction = - serde_json::from_slice::<Transaction>(&transaction_json).with_context(|| { - format!( - "failed to parse {kind} leaderboard transaction JSON with {} bytes", - transaction_json.len() - ) - })?; - let amount = match transaction { - Transaction::Mine { .. } => MINE_REWARD, - Transaction::Burn { amount, .. } => amount, - Transaction::Transfer { .. } => 0, - }; - let entry = entries - .entry(address.clone()) - .or_insert_with(|| UiLeaderboardEntry { - address, - amount: 0, - count: 0, - }); - entry.amount = entry - .amount - .checked_add(amount) - .with_context(|| format!("{kind} leaderboard amount overflow"))?; - entry.count = entry - .count - .checked_add(1) - .with_context(|| format!("{kind} leaderboard count overflow"))?; + let (height, hash) = row.context("failed to read metrics common ancestor")?; + if snapshot + .blocks + .get(height as usize) + .is_some_and(|block| block.hash == hash) + { + return Ok(Some(height)); + } } - let mut entries = entries.into_values().collect::<Vec<_>>(); - entries.sort_by(|left, right| { - right - .amount - .cmp(&left.amount) - .then_with(|| right.count.cmp(&left.count)) - .then_with(|| left.address.cmp(&right.address)) - }); - entries.truncate(limit); - Ok(entries) + Ok(None) } -fn replace_metrics( +fn load_metric_at_height( transaction: &rusqlite::Transaction<'_>, - metrics: &[BlockMetricRow], -) -> Result<()> { - clear_metrics_in_transaction(transaction)?; - for metric in metrics { - transaction - .execute( - r#" + height: u64, +) -> Result<Option<BlockMetricRow>> { + transaction + .query_row( + r#" +SELECT height, block_hash, timestamp_ms, block_time_ms, mine_difficulty_bits, + circulating_supply, known_wallet_addresses, transaction_count, transfer_count, burn_count, + mine_count, burned_amount, total_burned_amount, fees_amount, reward_amount, + vdf_rounds, finalizer_rank +FROM block_metrics +WHERE height = ?1 +"#, + [height], + |row| { + Ok(BlockMetricRow { + height: row.get(0)?, + block_hash: row.get(1)?, + timestamp_ms: row.get(2)?, + block_time_ms: row.get(3)?, + mine_difficulty_bits: row.get(4)?, + circulating_supply: row.get(5)?, + known_wallet_addresses: row.get(6)?, + transaction_count: row.get(7)?, + transfer_count: row.get(8)?, + burn_count: row.get(9)?, + mine_count: row.get(10)?, + burned_amount: row.get(11)?, + total_burned_amount: row.get(12)?, + fees_amount: row.get(13)?, + reward_amount: row.get(14)?, + vdf_rounds: row.get(15)?, + finalizer_rank: row.get(16)?, + }) + }, + ) + .optional() + .context("failed to load previous block metric") +} + +fn insert_metric(transaction: &rusqlite::Transaction<'_>, metric: &BlockMetricRow) -> Result<()> { + transaction + .execute( + r#" INSERT INTO block_metrics ( height, block_hash, timestamp_ms, block_time_ms, mine_difficulty_bits, circulating_supply, known_wallet_addresses, transaction_count, transfer_count, burn_count, @@ -568,28 +728,27 @@ INSERT INTO block_metrics ( finalizer_rank ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17) "#, - params![ - metric.height, - metric.block_hash, - metric.timestamp_ms, - metric.block_time_ms, - metric.mine_difficulty_bits, - metric.circulating_supply, - metric.known_wallet_addresses, - metric.transaction_count, - metric.transfer_count, - metric.burn_count, - metric.mine_count, - metric.burned_amount, - metric.total_burned_amount, - metric.fees_amount, - metric.reward_amount, - metric.vdf_rounds, - metric.finalizer_rank, - ], - ) - .with_context(|| format!("failed to insert metrics for block {}", metric.height))?; - } + params![ + metric.height, + metric.block_hash, + metric.timestamp_ms, + metric.block_time_ms, + metric.mine_difficulty_bits, + metric.circulating_supply, + metric.known_wallet_addresses, + metric.transaction_count, + metric.transfer_count, + metric.burn_count, + metric.mine_count, + metric.burned_amount, + metric.total_burned_amount, + metric.fees_amount, + metric.reward_amount, + metric.vdf_rounds, + metric.finalizer_rank, + ], + ) + .with_context(|| format!("failed to insert metrics for block {}", metric.height))?; Ok(()) } @@ -622,6 +781,23 @@ fn clear_metrics_in_transaction(transaction: &rusqlite::Transaction<'_>) -> Resu transaction .execute("DELETE FROM block_metrics", []) .context("failed to clear old block metrics")?; + transaction + .execute("DELETE FROM metrics_cache_meta", []) + .context("failed to clear old metrics cache metadata")?; + transaction + .execute("DELETE FROM metric_known_addresses", []) + .context("failed to clear old metric addresses")?; + Ok(()) +} + +fn clear_metrics_state(connection: &mut Connection) -> Result<()> { + let transaction = connection + .transaction() + .context("failed to start metrics cleanup transaction")?; + clear_metrics_in_transaction(&transaction)?; + transaction + .commit() + .context("failed to commit metrics cleanup transaction")?; Ok(()) } @@ -761,6 +937,32 @@ INSERT INTO ui_wallet_transactions ( Ok(()) } +fn replace_ui_leaderboards( + transaction: &rusqlite::Transaction<'_>, + rows: &[(String, UiLeaderboardEntry)], +) -> Result<()> { + transaction + .execute("DELETE FROM ui_leaderboards", []) + .context("failed to clear old UI leaderboards")?; + for (kind, row) in rows { + transaction + .execute( + r#" +INSERT INTO ui_leaderboards (kind, address, amount, count) +VALUES (?1, ?2, ?3, ?4) +"#, + params![kind, row.address, row.amount, row.count], + ) + .with_context(|| { + format!( + "failed to persist {kind} leaderboard row for {}", + row.address + ) + })?; + } + Ok(()) +} + fn clear_ui_chain_index_in_transaction(transaction: &rusqlite::Transaction<'_>) -> Result<()> { transaction .execute("DELETE FROM ui_cache_meta", []) @@ -775,6 +977,9 @@ fn clear_ui_chain_index_in_transaction(transaction: &rusqlite::Transaction<'_>) .execute("DELETE FROM ui_wallet_transactions", []) .context("failed to clear old UI wallet transaction index")?; transaction + .execute("DELETE FROM ui_leaderboards", []) + .context("failed to clear old UI leaderboards")?; + transaction .execute("DELETE FROM ui_burn_leader_ranks", []) .context("failed to clear old UI burn leader rank index")?; transaction @@ -1047,6 +1252,57 @@ ORDER BY block_hash, rank Ok(by_block_hash) } +fn build_ui_leaderboards( + utxos: &[(OutPoint, TxOutput)], + wallet_transactions: &[(String, WalletTransactionProjection)], +) -> Result<Vec<(String, UiLeaderboardEntry)>> { + let mut by_kind = BTreeMap::<(String, String), UiLeaderboardEntry>::new(); + for (_, output) in utxos { + if output.amount == 0 { + continue; + } + let key = ("balances".to_string(), output.address.clone()); + let entry = by_kind.entry(key).or_insert_with(|| UiLeaderboardEntry { + address: output.address.clone(), + amount: 0, + count: 0, + }); + entry.amount = entry + .amount + .checked_add(output.amount) + .context("balance leaderboard amount overflows")?; + entry.count = entry + .count + .checked_add(1) + .context("balance leaderboard count overflows")?; + } + for (address, projection) in wallet_transactions { + let (kind, amount) = match &projection.transaction { + Transaction::Mine { .. } => ("mine", MINE_REWARD), + Transaction::Burn { amount, .. } => ("burn", *amount), + Transaction::Transfer { .. } => continue, + }; + let key = (kind.to_string(), address.clone()); + let entry = by_kind.entry(key).or_insert_with(|| UiLeaderboardEntry { + address: address.clone(), + amount: 0, + count: 0, + }); + entry.amount = entry + .amount + .checked_add(amount) + .with_context(|| format!("{kind} leaderboard amount overflows"))?; + entry.count = entry + .count + .checked_add(1) + .with_context(|| format!("{kind} leaderboard count overflows"))?; + } + Ok(by_kind + .into_iter() + .map(|((kind, _), entry)| (kind, entry)) + .collect()) +} + fn wallet_transactions_from_snapshot( snapshot: &ChainSnapshot, ) -> Vec<(String, WalletTransactionProjection)> { @@ -1108,212 +1364,157 @@ fn transaction_kind(transaction: &Transaction) -> &'static str { } } -fn metrics_from_snapshot(snapshot: &ChainSnapshot) -> Result<Vec<BlockMetricRow>> { - let ledger = Ledger::from_persisted_snapshot(snapshot.clone()) - .context("failed to rebuild ledger for metrics")?; - let genesis = snapshot - .blocks - .first() - .cloned() - .context("cannot compute metrics for empty chain snapshot")?; - let mut running_ledger = Ledger::from_persisted_snapshot(ChainSnapshot { - genesis_allocations: snapshot.genesis_allocations.clone(), - vdf_rounds: snapshot.vdf_rounds, - launch_profile: snapshot.launch_profile.clone(), - blocks: vec![genesis], - }) - .context("failed to rebuild genesis ledger for metrics")?; - let mut known_wallet_addresses = snapshot - .genesis_allocations - .keys() - .cloned() - .collect::<BTreeSet<_>>(); - let mut total_burned_amount = 0_u64; - let mut rows = Vec::with_capacity(snapshot.blocks.len()); - let mut previous_timestamp_ms = None; - let mut metric_utxos = metric_genesis_utxos(snapshot); - - for block in &snapshot.blocks { - let reward_committee = if block.height == 0 { - Vec::new() - } else { - running_ledger.burn_committee_for_block(block) - }; - let mut transfer_count = 0_u64; - let mut burn_count = 0_u64; - let mut mine_count = 0_u64; - let mut burned_amount = 0_u64; - let burned_fee_amount = 0_u64; - let mut fees_amount = 0_u64; - - known_wallet_addresses.insert(block.miner.clone()); - for signature in &block.burn_bundle_section.signatures { - known_wallet_addresses.insert(signature.member.clone()); - } - for (_, output) in reward_outputs_for_block(block, &reward_committee) { - known_wallet_addresses.insert(output.address); - } - for transaction in &block.transactions { - collect_transaction_addresses(transaction, &mut known_wallet_addresses); - metric_apply_public_transaction(transaction, &mut metric_utxos)?; - fees_amount = fees_amount - .checked_add(transaction.fee()) - .context("block metric fees overflow")?; - match transaction { - Transaction::Transfer { .. } => transfer_count += 1, - Transaction::Burn { amount, .. } => { - burn_count += 1; - burned_amount = burned_amount - .checked_add(*amount) - .context("block metric burns overflow")?; - } - Transaction::Mine { .. } => { - mine_count += 1; - } - } +fn validate_projection_snapshot_structure(snapshot: &ChainSnapshot) -> Result<()> { + if snapshot.blocks.is_empty() { + anyhow::bail!("chain snapshot is empty"); + } + for (index, block) in snapshot.blocks.iter().enumerate() { + if block.height != index as u64 { + anyhow::bail!("chain snapshot block heights are not contiguous"); } - metric_index_block_reward(&mut metric_utxos, block, &reward_committee); - total_burned_amount = total_burned_amount - .checked_add(burned_amount) - .and_then(|amount| amount.checked_add(burned_fee_amount)) - .context("total burned metric overflows")?; - - if block.height > 0 { - running_ledger - .apply_preverified_block_at(block.clone(), u64::MAX) - .with_context(|| format!("failed to replay block {} for metrics", block.height))?; + if index > 0 && block.prev_hash != snapshot.blocks[index - 1].hash { + anyhow::bail!("chain snapshot parent hashes are not contiguous"); } - let circulating_supply = ledger_circulating_supply(&running_ledger)?; - let block_time_ms = - previous_timestamp_ms.map(|previous| block.timestamp_ms.saturating_sub(previous)); - previous_timestamp_ms = Some(block.timestamp_ms); - rows.push(BlockMetricRow { - height: block.height, - block_hash: block.hash.clone(), - timestamp_ms: block.timestamp_ms, - block_time_ms, - mine_difficulty_bits: ledger.mine_difficulty_bits_at_height(block.height), - circulating_supply, - known_wallet_addresses: known_wallet_addresses.len() as u64, - transaction_count: block.transactions.len() as u64, - transfer_count, - burn_count, - mine_count, - burned_amount, - total_burned_amount, - fees_amount, - reward_amount: block.reward, - vdf_rounds: block.vdf_rounds, - finalizer_rank: block.finalizer_rank, - }); } - Ok(rows) -} - -fn ledger_circulating_supply(ledger: &Ledger) -> Result<Amount> { - ledger - .status() - .balances - .values() - .try_fold(0_u64, |total, amount| { - total - .checked_add(*amount) - .context("circulating supply metric overflows") - }) -} - -fn metric_genesis_utxos(snapshot: &ChainSnapshot) -> BTreeMap<OutPoint, TxOutput> { - snapshot - .genesis_allocations - .iter() - .filter(|(_, amount)| **amount > 0) - .map(|(address, amount)| { - ( - metric_genesis_allocation_outpoint(address), - TxOutput { - address: address.clone(), - amount: *amount, - }, - ) - }) - .collect() -} - -fn metric_apply_public_transaction( - transaction: &Transaction, - utxos: &mut BTreeMap<OutPoint, TxOutput>, -) -> Result<()> { - metric_spend_transaction_inputs(transaction, utxos)?; - metric_index_transaction_outputs(utxos, transaction); Ok(()) } -fn metric_spend_transaction_inputs( - transaction: &Transaction, - utxos: &mut BTreeMap<OutPoint, TxOutput>, -) -> Result<Amount> { - let inputs = match transaction { - Transaction::Transfer { inputs, .. } | Transaction::Burn { inputs, .. } => inputs, - Transaction::Mine { .. } => return Ok(0), +fn incremental_metric_for_block( + snapshot: &ChainSnapshot, + block: &Block, + previous: Option<&BlockMetricRow>, +) -> Result<BlockMetricRow> { + let mut circulating_supply = match previous { + Some(previous) => previous.circulating_supply, + None => snapshot + .genesis_allocations + .values() + .try_fold(0_u64, |total, amount| { + total + .checked_add(*amount) + .context("genesis circulating supply overflows") + })?, }; - metric_spend_inputs(inputs, utxos) -} + let mut transfer_count = 0_u64; + let mut burn_count = 0_u64; + let mut mine_count = 0_u64; + let mut burned_amount = 0_u64; + let mut fees_amount = 0_u64; + + for transaction in &block.transactions { + fees_amount = fees_amount + .checked_add(transaction.fee()) + .context("block metric fees overflow")?; + match transaction { + Transaction::Transfer { fee, .. } => { + transfer_count += 1; + circulating_supply = circulating_supply + .checked_sub(*fee) + .context("transfer fee exceeds circulating supply")?; + } + Transaction::Burn { amount, fee, .. } => { + burn_count += 1; + burned_amount = burned_amount + .checked_add(*amount) + .context("block metric burns overflow")?; + circulating_supply = circulating_supply + .checked_sub(*amount) + .and_then(|supply| supply.checked_sub(*fee)) + .context("burn exceeds circulating supply")?; + } + Transaction::Mine { .. } => { + mine_count += 1; + circulating_supply = circulating_supply + .checked_add(MINE_REWARD) + .context("mine reward circulating supply overflows")?; + } + } + } + circulating_supply = circulating_supply + .checked_add(block.reward) + .context("block reward circulating supply overflows")?; + let total_burned_amount = previous + .map(|row| row.total_burned_amount) + .unwrap_or_default() + .checked_add(burned_amount) + .context("total burned metric overflows")?; -fn metric_spend_inputs( - inputs: &[TxInput], - utxos: &mut BTreeMap<OutPoint, TxOutput>, -) -> Result<Amount> { - inputs.iter().try_fold(0_u64, |total, input| { - let output = utxos.remove(&input.outpoint).with_context(|| { - format!( - "metric replay spends missing output {}:{}", - input.outpoint.txid, input.outpoint.index - ) - })?; - total - .checked_add(output.amount) - .context("metric replay input total overflows") + Ok(BlockMetricRow { + height: block.height, + block_hash: block.hash.clone(), + timestamp_ms: block.timestamp_ms, + block_time_ms: previous.map(|row| block.timestamp_ms.saturating_sub(row.timestamp_ms)), + mine_difficulty_bits: metric_difficulty_for_block(snapshot, block, previous), + circulating_supply, + known_wallet_addresses: 0, + transaction_count: block.transactions.len() as u64, + transfer_count, + burn_count, + mine_count, + burned_amount, + total_burned_amount, + fees_amount, + reward_amount: block.reward, + vdf_rounds: block.vdf_rounds, + finalizer_rank: block.finalizer_rank, }) } -fn metric_index_transaction_outputs( - utxos: &mut BTreeMap<OutPoint, TxOutput>, - transaction: &Transaction, -) { - let outputs = match transaction { - Transaction::Transfer { outputs, .. } => outputs.clone(), - Transaction::Burn { change, .. } => change.clone(), - Transaction::Mine { recipient, .. } => vec![TxOutput { - address: recipient.clone(), - amount: MINE_REWARD, - }], - }; - for (index, output) in outputs.iter().enumerate() { - utxos.insert( - OutPoint { - txid: transaction.signature().to_string(), - index: index as u32, - }, - output.clone(), - ); +fn metric_difficulty_for_block( + snapshot: &ChainSnapshot, + block: &Block, + previous: Option<&BlockMetricRow>, +) -> u32 { + let previous_difficulty = previous + .map(|row| row.mine_difficulty_bits) + .unwrap_or(snapshot.launch_profile.mine_difficulty_bits); + if block.height == 0 || block.height % MINE_RETARGET_WINDOW_BLOCKS != 0 { + return previous_difficulty; } + let window_start = block.height + 1 - MINE_RETARGET_WINDOW_BLOCKS; + let mine_actions = snapshot.blocks[window_start as usize..=block.height as usize] + .iter() + .flat_map(|candidate| candidate.transactions.iter()) + .filter(|transaction| matches!(transaction, Transaction::Mine { .. })) + .count() as u64; + retarget_mine_difficulty_bits(previous_difficulty, mine_actions) } -fn metric_index_block_reward( - utxos: &mut BTreeMap<OutPoint, TxOutput>, - block: &Block, - committee: &[crate::domain::BurnCommitteeMember], -) { - for (outpoint, output) in reward_outputs_for_block(block, committee) { - utxos.insert(outpoint, output); +fn index_metric_addresses(transaction: &rusqlite::Transaction<'_>, block: &Block) -> Result<u64> { + let mut inserted = insert_metric_address(transaction, &block.miner, block.height)?; + for signature in &block.burn_bundle_section.signatures { + inserted = inserted + .checked_add(insert_metric_address( + transaction, + &signature.member, + block.height, + )?) + .context("known metric address count overflows")?; + } + let mut addresses = BTreeSet::new(); + for public_transaction in &block.transactions { + collect_transaction_addresses(public_transaction, &mut addresses); } + for address in addresses { + inserted = inserted + .checked_add(insert_metric_address(transaction, &address, block.height)?) + .context("known metric address count overflows")?; + } + Ok(inserted) } -fn metric_genesis_allocation_outpoint(address: &str) -> OutPoint { - OutPoint { - txid: hex_hash(format!("iuna-genesis-allocation:{address}")), - index: 0, - } +fn insert_metric_address( + transaction: &rusqlite::Transaction<'_>, + address: &str, + first_seen_height: u64, +) -> Result<u64> { + let inserted = transaction + .execute( + "INSERT OR IGNORE INTO metric_known_addresses (address, first_seen_height) VALUES (?1, ?2)", + params![address, first_seen_height], + ) + .with_context(|| format!("failed to index metric address {address}"))?; + Ok(inserted as u64) } fn collect_transaction_addresses(transaction: &Transaction, addresses: &mut BTreeSet<String>) { @@ -1357,6 +1558,7 @@ fn unix_ms() -> u64 { mod tests { use std::collections::BTreeMap; + use rusqlite::Connection; use tempfile::tempdir; use crate::domain::{ChainSnapshot, GenesisBurn, Ledger, Wallet}; @@ -1372,6 +1574,26 @@ mod tests { .snapshot() } + fn test_ledger(seed: &str) -> (Ledger, Wallet) { + let wallet = Wallet::from_seed(seed); + let mut allocations = BTreeMap::new(); + allocations.insert(wallet.address().to_string(), 1_000); + let ledger = Ledger::new_with_genesis_burns( + allocations, + vec![GenesisBurn::new(wallet.address(), 10)], + 1, + ) + .unwrap(); + (ledger, wallet) + } + + fn append_test_block(ledger: &mut Ledger, wallet: &Wallet, timestamp_ms: u64) { + let burn = ledger.build_burn(wallet, 1, 1).unwrap(); + ledger.submit_transaction(burn).unwrap(); + let block = ledger.mine_next_block(wallet, timestamp_ms).unwrap(); + ledger.apply_locally_mined_block(block).unwrap(); + } + #[test] fn failed_ui_projection_keeps_last_committed_projection() { let dir = tempdir().unwrap(); @@ -1395,4 +1617,147 @@ mod tests { assert!(store.metrics_are_projected_to(&tip).unwrap()); assert!(!store.load_metrics().unwrap().is_empty()); } + + #[test] + fn metrics_projection_appends_without_rewriting_the_consistent_prefix() { + let dir = tempdir().unwrap(); + let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap(); + let (mut ledger, wallet) = test_ledger("incremental-metrics"); + append_test_block(&mut ledger, &wallet, 1_000); + append_test_block(&mut ledger, &wallet, 2_000); + store.project_snapshot(&ledger.snapshot(), true).unwrap(); + let prefix = store.load_metrics().unwrap(); + + let connection = Connection::open(store.path()).unwrap(); + connection + .execute_batch( + r#" +CREATE TRIGGER protect_metric_prefix +BEFORE DELETE ON block_metrics +WHEN OLD.height <= 2 +BEGIN + SELECT RAISE(FAIL, 'consistent metric prefix was rewritten'); +END; +"#, + ) + .unwrap(); + drop(connection); + + append_test_block(&mut ledger, &wallet, 3_000); + let snapshot = ledger.snapshot(); + store.project_snapshot(&snapshot, true).unwrap(); + let metrics = store.load_metrics().unwrap(); + + assert_eq!(&metrics[..prefix.len()], prefix.as_slice()); + assert_eq!(metrics.len(), snapshot.blocks.len()); + assert_eq!(metrics.last().unwrap().block_hash, ledger.tip_hash()); + assert_eq!( + metrics.last().unwrap().circulating_supply, + ledger + .all_utxos() + .iter() + .map(|(_, output)| output.amount) + .sum::<u64>() + ); + assert_eq!( + metrics.last().unwrap().mine_difficulty_bits, + ledger.mine_difficulty_bits_at_height(ledger.height()) + ); + } + + #[test] + fn metrics_projection_replaces_only_the_reorged_suffix() { + let dir = tempdir().unwrap(); + let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap(); + let (mut base, wallet) = test_ledger("reorg-metrics"); + 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); + store + .project_snapshot(&first_branch.snapshot(), true) + .unwrap(); + + let connection = Connection::open(store.path()).unwrap(); + connection + .execute_batch( + r#" +CREATE TRIGGER protect_metric_common_ancestor +BEFORE DELETE ON block_metrics +WHEN OLD.height <= 1 +BEGIN + SELECT RAISE(FAIL, 'metric common ancestor was rewritten'); +END; +"#, + ) + .unwrap(); + drop(connection); + + store + .project_snapshot(&second_branch.snapshot(), true) + .unwrap(); + let metrics = store.load_metrics().unwrap(); + + assert_eq!(metrics.len(), second_branch.chain().len()); + assert_eq!(metrics.last().unwrap().block_hash, second_branch.tip_hash()); + assert!( + store + .metrics_are_projected_to(second_branch.tip_hash()) + .unwrap() + ); + assert!( + !store + .metrics_are_projected_to(first_branch.tip_hash()) + .unwrap() + ); + } + + #[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(); + let (mut ledger, wallet) = test_ledger("retarget-metrics"); + + store.project_snapshot(&ledger.snapshot(), true).unwrap(); + for height in 1..=super::MINE_RETARGET_WINDOW_BLOCKS + 2 { + append_test_block(&mut ledger, &wallet, height * 1_000); + store.project_snapshot(&ledger.snapshot(), true).unwrap(); + let latest = store.load_metrics().unwrap().pop().unwrap(); + + assert_eq!(latest.height, height); + assert_eq!( + latest.mine_difficulty_bits, + ledger.mine_difficulty_bits_at_height(height) + ); + assert_eq!( + latest.circulating_supply, + ledger + .all_utxos() + .iter() + .map(|(_, output)| output.amount) + .sum::<u64>() + ); + } + } + + #[test] + fn leaderboards_are_served_from_their_materialized_projection() { + let dir = tempdir().unwrap(); + let store = SqliteUiDataStore::open(dir.path().join("ui_data.sqlite3")).unwrap(); + let (mut ledger, wallet) = test_ledger("materialized-leaderboards"); + append_test_block(&mut ledger, &wallet, 1_000); + store.project_snapshot(&ledger.snapshot(), true).unwrap(); + let expected = store.load_leaderboards(10).unwrap(); + assert!(!expected.balances.is_empty()); + assert!(!expected.burners.is_empty()); + + let connection = Connection::open(store.path()).unwrap(); + connection.execute("DELETE FROM ui_utxos", []).unwrap(); + connection + .execute("DELETE FROM ui_wallet_transactions", []) + .unwrap(); + + assert_eq!(store.load_leaderboards(10).unwrap(), expected); + } } diff --git a/src/domain.rs b/src/domain.rs @@ -50,6 +50,7 @@ use ledger_ops::{ pub(crate) use ledger_pending::MINE_ANCHOR_LIMIT_REACHED; pub use ledger_state::Ledger; use ledger_state::unix_now_ms; +pub(crate) use mine_policy::{MINE_RETARGET_WINDOW_BLOCKS, retarget_mine_difficulty_bits}; use mining::{mine_payload, mine_signature}; pub use profile::{GenesisBurn, LaunchProfile}; pub use protocol::{ diff --git a/src/domain/mine_policy.rs b/src/domain/mine_policy.rs @@ -4,14 +4,14 @@ use anyhow::{Result, bail}; use super::{Block, MINE_ACTIONS_PER_ANCHOR_LIMIT, Transaction}; -pub(super) const MINE_RETARGET_WINDOW_BLOCKS: u64 = 10; +pub(crate) const MINE_RETARGET_WINDOW_BLOCKS: u64 = 10; pub(super) const MINE_MAX_RETARGET_STEP_BITS: u32 = 2; pub(super) const MINE_MIN_DIFFICULTY_BITS: u32 = 10; pub(super) const MINE_MAX_ANCHOR_AGE_BLOCKS: u64 = MINE_RETARGET_WINDOW_BLOCKS; const MINE_TARGET_ACTIONS_PER_BLOCK: u64 = 1; -pub(super) fn retarget_mine_difficulty_bits(current: u32, mine_actions: u64) -> u32 { +pub(crate) fn retarget_mine_difficulty_bits(current: u32, mine_actions: u64) -> u32 { let target = MINE_RETARGET_WINDOW_BLOCKS.saturating_mul(MINE_TARGET_ACTIONS_PER_BLOCK); if target == 0 || mine_actions == target { return current.max(MINE_MIN_DIFFICULTY_BITS); diff --git a/src/main.rs b/src/main.rs @@ -1016,7 +1016,35 @@ async fn warm_ui_data_store( ) -> Result<()> { println!("warming UI data database..."); let started = Instant::now(); - project_ui_data_store(store, snapshot, keep_metrics).await?; + let tip_hash = snapshot + .blocks + .last() + .map(|block| block.hash.clone()) + .context("cannot warm UI data database from an empty chain")?; + let readiness_store = store.clone(); + let readiness_tip = tip_hash.clone(); + let (ui_ready, metrics_ready) = tokio::task::spawn_blocking(move || { + Ok::<_, anyhow::Error>(( + readiness_store.is_projected_to(&readiness_tip)?, + readiness_store.metrics_are_projected_to(&readiness_tip)?, + )) + }) + .await + .context("UI data readiness worker failed")??; + + if !ui_ready { + project_ui_data_store(store, snapshot, keep_metrics).await?; + } else if keep_metrics && !metrics_ready { + let metrics_store = store.clone(); + tokio::task::spawn_blocking(move || metrics_store.replace_metrics_for_snapshot(&snapshot)) + .await + .context("metrics warm-up worker failed")??; + } else if !keep_metrics { + let metrics_store = store.clone(); + tokio::task::spawn_blocking(move || metrics_store.clear_metrics()) + .await + .context("metrics cleanup worker failed")??; + } println!( "UI data database ready in {:.2}s", started.elapsed().as_secs_f64()