diff --git a/crates/rooch/src/commands/db/commands/mod.rs b/crates/rooch/src/commands/db/commands/mod.rs index fed0f3f2f5..d36dad69c3 100644 --- a/crates/rooch/src/commands/db/commands/mod.rs +++ b/crates/rooch/src/commands/db/commands/mod.rs @@ -37,6 +37,7 @@ pub mod rocksdb_stats; pub mod rollback; pub mod stat_changeset; pub mod state_prune; +pub mod tx_accumulator_compact; pub mod verify_order; fn open_rocks( diff --git a/crates/rooch/src/commands/db/commands/tx_accumulator_compact.rs b/crates/rooch/src/commands/db/commands/tx_accumulator_compact.rs new file mode 100644 index 0000000000..8e9a49a28a --- /dev/null +++ b/crates/rooch/src/commands/db/commands/tx_accumulator_compact.rs @@ -0,0 +1,363 @@ +// Copyright (c) RoochNetwork +// SPDX-License-Identifier: Apache-2.0 + +use crate::cli_types::CommandAction; +use crate::commands::db::commands::{load_accumulator, open_rocks}; +use crate::utils::open_rooch_db; +use accumulator::node_index::{FrozenSubTreeIterator, NodeIndex}; +use accumulator::{Accumulator, AccumulatorNode, AccumulatorTreeStore as _}; +use anyhow::{ensure, Result}; +use async_trait::async_trait; +use clap::Parser; +use moveos_types::h256::{ACCUMULATOR_PLACEHOLDER_HASH, H256}; +use rooch_config::R_OPT_NET_HELP; +use rooch_store::{RoochStore, TX_ACCUMULATOR_NODE_COLUMN_FAMILY_NAME}; +use rooch_types::error::RoochResult; +use rooch_types::rooch_network::RoochChainID; +use std::collections::HashMap; +use std::fmt::Write; +use std::path::PathBuf; +use std::time::Instant; + +/// Compact historical non-frozen transaction accumulator nodes. +/// +/// The command replays historical leaves and reconstructs the transient +/// non-frozen internal node hashes that older versions persisted. Frozen +/// subtrees and leaves are not selected for deletion. +#[derive(Debug, Parser)] +pub struct TxAccumulatorCompactCommand { + #[clap(long = "data-dir", short = 'd')] + pub base_data_dir: Option, + + #[clap(long, short = 'n', help = R_OPT_NET_HELP)] + pub chain_id: Option, + + /// First leaf index to include in deletion candidate collection. + /// Earlier leaves are still replayed to rebuild frozen roots. + #[clap(long, default_value_t = 0)] + pub start_index: u64, + + /// Stop before this leaf index. Defaults to the current accumulator leaf count. + #[clap(long)] + pub end_index: Option, + + /// Number of candidate hashes to check/delete per batch. + #[clap(long, default_value_t = 10000)] + pub batch_size: usize, + + /// Print progress after this many replayed leaves. 0 disables progress lines. + #[clap(long, default_value_t = 1_000_000)] + pub progress_interval: u64, + + /// Actually delete matched non-frozen nodes. Without this flag the command is a dry-run. + #[clap(long)] + pub execute: bool, + + /// Force RocksDB compaction on the transaction_acc_node column family after deletion. + #[clap(long)] + pub force_compaction: bool, +} + +#[derive(Debug, Default, Clone, Eq, PartialEq)] +struct CompactReport { + start_index: u64, + end_index: u64, + replayed_leaves: u64, + candidate_nodes: u64, + existing_nodes: u64, + deleted_nodes: u64, + dry_run: bool, +} + +#[async_trait] +impl CommandAction for TxAccumulatorCompactCommand { + async fn execute(self) -> RoochResult { + self.execute_impl().map_err(Into::into) + } +} + +impl TxAccumulatorCompactCommand { + fn execute_impl(self) -> Result { + ensure!(self.batch_size > 0, "batch-size must be greater than 0"); + + let started_at = Instant::now(); + let dry_run = !self.execute; + let (_root, rooch_db, _start_time) = + open_rooch_db(self.base_data_dir.clone(), self.chain_id.clone()); + let rooch_store = rooch_db.rooch_store.clone(); + let (tx_accumulator, _last_order) = load_accumulator(rooch_store.clone())?; + let leaf_count = tx_accumulator.num_leaves(); + let end_index = self.end_index.unwrap_or(leaf_count); + ensure!( + self.start_index <= end_index && end_index <= leaf_count, + "invalid range [{}, {}), current leaf count {}", + self.start_index, + end_index, + leaf_count + ); + + let mut replayer = NonFrozenNodeReplayer::new(); + let mut report = CompactReport { + start_index: self.start_index, + end_index, + dry_run, + ..Default::default() + }; + let mut pending_hashes = Vec::with_capacity(self.batch_size); + let mut out = String::new(); + + for leaf_index in 0..end_index { + let leaf = tx_accumulator.get_leaf(leaf_index)?.ok_or_else(|| { + anyhow::anyhow!("transaction accumulator leaf {} not found", leaf_index) + })?; + let non_frozen_hashes = replayer.append_one(leaf)?; + if leaf_index < self.start_index { + continue; + } + + report.replayed_leaves += 1; + report.candidate_nodes += non_frozen_hashes.len() as u64; + pending_hashes.extend(non_frozen_hashes); + if pending_hashes.len() >= self.batch_size { + flush_candidates(&rooch_store, &mut pending_hashes, dry_run, &mut report)?; + } + + if self.progress_interval > 0 && report.replayed_leaves % self.progress_interval == 0 { + writeln!( + out, + "progress: replayed={} candidates={} existing={} deleted={}", + report.replayed_leaves, + report.candidate_nodes, + report.existing_nodes, + report.deleted_nodes + )?; + } + } + flush_candidates(&rooch_store, &mut pending_hashes, dry_run, &mut report)?; + + drop(tx_accumulator); + drop(rooch_store); + drop(rooch_db); + + let compact_elapsed = if self.force_compaction && !dry_run { + Some(compact_tx_accumulator_cf( + self.base_data_dir, + self.chain_id, + )?) + } else { + None + }; + + writeln!(out, "=== Tx Accumulator Compact Result ===")?; + writeln!(out, "mode: {}", if dry_run { "dry-run" } else { "execute" })?; + writeln!(out, "range: [{}..{})", report.start_index, report.end_index)?; + writeln!(out, "replayed leaves: {}", report.replayed_leaves)?; + writeln!( + out, + "candidate non-frozen nodes: {}", + report.candidate_nodes + )?; + writeln!(out, "existing candidate nodes: {}", report.existing_nodes)?; + writeln!(out, "deleted nodes: {}", report.deleted_nodes)?; + if let Some(elapsed) = compact_elapsed { + writeln!(out, "rocksdb compaction: {:?}", elapsed)?; + } + writeln!(out, "elapsed: {:?}", started_at.elapsed())?; + Ok(out) + } +} + +fn flush_candidates( + rooch_store: &RoochStore, + pending_hashes: &mut Vec, + dry_run: bool, + report: &mut CompactReport, +) -> Result<()> { + if pending_hashes.is_empty() { + return Ok(()); + } + + pending_hashes.sort(); + pending_hashes.dedup(); + let existing_hashes = rooch_store + .transaction_accumulator_store + .multiple_get(pending_hashes.clone())? + .into_iter() + .zip(pending_hashes.iter()) + .filter_map(|(node, hash)| node.map(|_| *hash)) + .collect::>(); + + report.existing_nodes += existing_hashes.len() as u64; + if !dry_run && !existing_hashes.is_empty() { + let deleted = existing_hashes.len() as u64; + rooch_store + .transaction_accumulator_store + .delete_nodes(existing_hashes)?; + report.deleted_nodes += deleted; + } + + pending_hashes.clear(); + Ok(()) +} + +fn compact_tx_accumulator_cf( + base_data_dir: Option, + chain_id: Option, +) -> Result { + let db = open_rocks(base_data_dir, chain_id)?; + let raw = db.inner(); + let cf = raw + .cf_handle(TX_ACCUMULATOR_NODE_COLUMN_FAMILY_NAME) + .ok_or_else(|| anyhow::anyhow!("transaction accumulator column family not found"))?; + + raw.flush_wal(true)?; + raw.flush_cf(&cf)?; + use rocksdb::{BottommostLevelCompaction, CompactOptions}; + let mut copt = CompactOptions::default(); + copt.set_bottommost_level_compaction(BottommostLevelCompaction::Force); + copt.set_exclusive_manual_compaction(true); + + let start = Instant::now(); + raw.compact_range_cf_opt(&cf, None::<&[u8]>, None::<&[u8]>, &copt); + Ok(start.elapsed()) +} + +#[derive(Debug, Default)] +struct NonFrozenNodeReplayer { + num_leaves: u64, + frozen_roots: HashMap, +} + +impl NonFrozenNodeReplayer { + fn new() -> Self { + Self::default() + } + + fn append_one(&mut self, leaf: H256) -> Result> { + let leaf_pos = NodeIndex::from_leaf_index(self.num_leaves); + let last_new_leaf_count = self.num_leaves + 1; + let root_level = NodeIndex::root_level_from_leaf_count(last_new_leaf_count); + let mut new_frozen = HashMap::new(); + + let mut pos = leaf_pos; + let mut hash = leaf; + new_frozen.insert(pos, hash); + + while pos.is_right_child() { + let sibling = pos.sibling(); + let left_hash = new_frozen + .get(&sibling) + .or_else(|| self.frozen_roots.get(&sibling)) + .copied() + .ok_or_else(|| anyhow::anyhow!("missing frozen sibling {:?}", sibling))?; + let internal_node = AccumulatorNode::new_internal(pos.parent(), left_hash, hash); + hash = internal_node.hash(); + pos = pos.parent(); + new_frozen.insert(pos, hash); + } + + let mut non_frozen_hashes = Vec::new(); + for _ in pos.level()..root_level { + let internal_node = if pos.is_left_child() { + AccumulatorNode::new_internal(pos.parent(), hash, *ACCUMULATOR_PLACEHOLDER_HASH) + } else { + let sibling = pos.sibling(); + let left_hash = self + .frozen_roots + .get(&sibling) + .copied() + .unwrap_or(*ACCUMULATOR_PLACEHOLDER_HASH); + AccumulatorNode::new_internal(pos.parent(), left_hash, hash) + }; + hash = internal_node.hash(); + pos = pos.parent(); + non_frozen_hashes.push(hash); + } + + self.num_leaves = last_new_leaf_count; + self.frozen_roots = FrozenSubTreeIterator::new(self.num_leaves) + .map(|index| { + let hash = new_frozen + .get(&index) + .or_else(|| self.frozen_roots.get(&index)) + .copied() + .ok_or_else(|| anyhow::anyhow!("missing frozen root {:?}", index))?; + Ok((index, hash)) + }) + .collect::>>()?; + + Ok(non_frozen_hashes) + } +} + +#[cfg(test)] +fn replay_non_frozen_hashes(leaves: &[H256]) -> Result> { + let mut replayer = NonFrozenNodeReplayer::new(); + let mut hashes = Vec::new(); + for leaf in leaves { + hashes.extend(replayer.append_one(*leaf)?); + } + Ok(hashes) +} + +#[cfg(test)] +mod tests { + use super::*; + use accumulator::MerkleAccumulator; + + #[test] + fn test_replay_non_frozen_hashes_for_single_leaf_appends() { + let leaves = (0..8).map(|_| H256::random()).collect::>(); + let hashes = replay_non_frozen_hashes(&leaves).unwrap(); + + // Perfect tree sizes have no non-frozen root after the last append, but + // intermediate non-perfect prefixes still produce transient nodes. + assert!(!hashes.is_empty()); + assert_eq!(hashes.len(), 10); + } + + #[test] + fn test_flush_candidates_deletes_only_existing_nodes() { + let (rooch_store, _tmpdir) = RoochStore::mock_rooch_store().unwrap(); + let accumulator = + MerkleAccumulator::new_empty(rooch_store.get_transaction_accumulator_store()); + let leaves = (0..6).map(|_| H256::random()).collect::>(); + for leaf in &leaves { + accumulator.append(&[*leaf]).unwrap(); + accumulator.flush().unwrap(); + } + + let old_non_frozen_hashes = replay_non_frozen_hashes(&leaves).unwrap(); + let persisted_old_nodes = old_non_frozen_hashes + .iter() + .map(|hash| AccumulatorNode::new_leaf(NodeIndex::from_leaf_index(10_000), *hash)) + .collect::>(); + rooch_store + .transaction_accumulator_store + .save_nodes(persisted_old_nodes) + .unwrap(); + + let mut dry_run_report = CompactReport::default(); + let mut pending = old_non_frozen_hashes.clone(); + flush_candidates(&rooch_store, &mut pending, true, &mut dry_run_report).unwrap(); + assert_eq!( + dry_run_report.existing_nodes, + old_non_frozen_hashes.len() as u64 + ); + assert_eq!(dry_run_report.deleted_nodes, 0); + + let mut execute_report = CompactReport::default(); + let mut pending = old_non_frozen_hashes.clone(); + flush_candidates(&rooch_store, &mut pending, false, &mut execute_report).unwrap(); + assert_eq!( + execute_report.deleted_nodes, + old_non_frozen_hashes.len() as u64 + ); + + let remaining = rooch_store + .transaction_accumulator_store + .multiple_get(old_non_frozen_hashes) + .unwrap(); + assert!(remaining.into_iter().all(|node| node.is_none())); + } +} diff --git a/crates/rooch/src/commands/db/mod.rs b/crates/rooch/src/commands/db/mod.rs index b2ad9d8c2b..024764b65d 100644 --- a/crates/rooch/src/commands/db/mod.rs +++ b/crates/rooch/src/commands/db/mod.rs @@ -26,6 +26,7 @@ use crate::commands::db::commands::rocksdb_gc::RocksDBGcCommand; use crate::commands::db::commands::rocksdb_stats::RocksDBStatsCommand; use crate::commands::db::commands::stat_changeset::StatChangesetCommand; use crate::commands::db::commands::state_prune::StatePruneCommand; +use crate::commands::db::commands::tx_accumulator_compact::TxAccumulatorCompactCommand; use crate::commands::db::commands::verify_order::VerifyOrderCommand; use async_trait::async_trait; use clap::Parser; @@ -126,6 +127,9 @@ impl CommandAction for DB { DBCommand::GC(gc) => gc.execute().await, DBCommand::Recycle(recycle) => recycle.execute().await, DBCommand::StatePrune(state_prune) => state_prune.execute().await, + DBCommand::TxAccumulatorCompact(tx_accumulator_compact) => { + tx_accumulator_compact.execute().await + } } } } @@ -159,4 +163,5 @@ pub enum DBCommand { GC(GCCommand), Recycle(RecycleCommand), StatePrune(StatePruneCommand), + TxAccumulatorCompact(TxAccumulatorCompactCommand), } diff --git a/moveos/moveos-commons/accumulator/src/tests/test_accumulator.rs b/moveos/moveos-commons/accumulator/src/tests/test_accumulator.rs index adf032b470..dd88e61104 100644 --- a/moveos/moveos-commons/accumulator/src/tests/test_accumulator.rs +++ b/moveos/moveos-commons/accumulator/src/tests/test_accumulator.rs @@ -282,6 +282,20 @@ fn test_multiple_tree() { proof_verify(&accumulator2, root_hash2, &batch1, 0); } +#[test] +fn test_reopen_non_power_of_two_accumulator_without_non_frozen_nodes() { + let leaves = create_leaves(750..755); + let mock_store = Arc::new(MockAccumulatorStore::new()); + let accumulator = MerkleAccumulator::new_empty(mock_store.clone()); + let root_hash = accumulator.append(&leaves).unwrap(); + accumulator.flush().unwrap(); + + let accumulator_info = accumulator.get_info(); + let reopened = MerkleAccumulator::new_with_info(accumulator_info, mock_store); + assert_eq!(reopened.root_hash(), root_hash); + proof_verify(&reopened, root_hash, &leaves, 0); +} + #[test] fn test_update_left_leaf() { // construct a accumulator diff --git a/moveos/moveos-commons/accumulator/src/tree.rs b/moveos/moveos-commons/accumulator/src/tree.rs index 4b80e37e66..7f3751e1f0 100644 --- a/moveos/moveos-commons/accumulator/src/tree.rs +++ b/moveos/moveos-commons/accumulator/src/tree.rs @@ -180,10 +180,12 @@ impl AccumulatorTree { node.clone() }) .collect(); - //aggregator all nodes - not_frozen_nodes.extend_from_slice(&to_freeze); - self.update_temp_nodes(not_frozen_nodes.clone()); + // Only frozen nodes are durable. Non-frozen nodes can be reconstructed from + // frozen subtree roots and placeholders, so keep them in the index cache for + // this process but do not add them to the pending storage updates. + self.update_temp_nodes(to_freeze.clone()); // update to cache + not_frozen_nodes.extend_from_slice(&to_freeze); self.update_cache(not_frozen_nodes); // update self properties self.root_hash = hash; @@ -320,13 +322,24 @@ impl AccumulatorTree { if let Some(node_hash) = self.get_node_index(index_key) { return Ok(node_hash); } + if index.is_placeholder(self.rightmost_leaf_index()) { + return Ok(*ACCUMULATOR_PLACEHOLDER_HASH); + } + if let Some(node_hash) = self.compute_frontier_hash(index) { + return Ok(node_hash); + } // find parent hash,then get node by parent hash let root_index = NodeIndex::root_from_leaf_count(self.num_leaves); let level = root_index.level() + 1; let mut parent_hash = None; for _i in 0..level { index_key = temp_index.parent(); - if let Some(internal_parent_hash) = self.get_node_index(index_key) { + if let Some(internal_parent_hash) = self + .index_frozen_subtrees + .get(&index_key) + .copied() + .or_else(|| self.get_node_index(index_key)) + { parent_hash = Some(internal_parent_hash); break; } @@ -382,6 +395,26 @@ impl AccumulatorTree { bail!("node hash not found:{:?}", index) } + fn compute_frontier_hash(&self, index: NodeIndex) -> Option { + if let Some(node_hash) = self.index_frozen_subtrees.get(&index) { + return Some(*node_hash); + } + if index.is_placeholder(self.rightmost_leaf_index()) { + return Some(*ACCUMULATOR_PLACEHOLDER_HASH); + } + if index.is_leaf() { + return None; + } + + let left = self.compute_frontier_hash(index.left_child())?; + let right = self.compute_frontier_hash(index.right_child())?; + if left == *ACCUMULATOR_PLACEHOLDER_HASH && right == *ACCUMULATOR_PLACEHOLDER_HASH { + Some(*ACCUMULATOR_PLACEHOLDER_HASH) + } else { + Some(AccumulatorNode::new_internal(index, left, right).hash()) + } + } + fn save_node_indexes(&mut self, nodes: Vec) { let id = format!("{:p}", self); let cache = &mut self.index_cache;