diff --git a/Cargo.lock b/Cargo.lock index dfe9f07310..688aee69ad 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1840,7 +1840,7 @@ dependencies = [ [[package]] name = "consensus_tests" -version = "0.12.0" +version = "0.13.0" dependencies = [ "fern", "futures 0.3.31", @@ -2607,7 +2607,7 @@ dependencies = [ [[package]] name = "db_inspector" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "axum 0.8.4", @@ -3746,7 +3746,7 @@ checksum = "8f5f3913fa0bfe7ee1fd8248b6b9f42a5af4b9d65ec2dd2c3c26132b950ecfc2" [[package]] name = "generate_ristretto_value_lookup" -version = "0.12.0" +version = "0.13.0" dependencies = [ "clap 3.2.25", "human_bytes", @@ -4993,7 +4993,7 @@ checksum = "8bb03732005da905c88227371639bf1ad885cc712789c011c31c5fb3ab3ccf02" [[package]] name = "integration_tests" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "config", @@ -5761,7 +5761,7 @@ dependencies = [ [[package]] name = "libp2p-messaging" -version = "0.12.0" +version = "0.13.0" dependencies = [ "async-trait", "futures-bounded", @@ -5940,7 +5940,7 @@ dependencies = [ [[package]] name = "libp2p-substream" -version = "0.12.0" +version = "0.13.0" dependencies = [ "libp2p", "prometheus-client", @@ -8466,7 +8466,7 @@ dependencies = [ [[package]] name = "proto_builder" -version = "0.12.0" +version = "0.13.0" dependencies = [ "prost-build 0.14.1", "sha2", @@ -10300,7 +10300,7 @@ dependencies = [ [[package]] name = "sqlite_message_logger" -version = "0.12.0" +version = "0.13.0" dependencies = [ "chrono", "diesel", @@ -10341,7 +10341,7 @@ checksum = "e7386b49cb287f6fafbfd3bd604914bccb99fb8d53483f40e1ecfda5d45f3370" [[package]] name = "state_store_tests" -version = "0.12.0" +version = "0.13.0" dependencies = [ "env_logger 0.11.8", "indexmap 2.11.4", @@ -10730,7 +10730,7 @@ dependencies = [ [[package]] name = "tari_base_node_client" -version = "0.12.0" +version = "0.13.0" dependencies = [ "log", "minotari_app_grpc", @@ -10954,7 +10954,7 @@ dependencies = [ [[package]] name = "tari_consensus" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "indexmap 2.11.4", @@ -10978,7 +10978,7 @@ dependencies = [ [[package]] name = "tari_consensus_types" -version = "0.12.0" +version = "0.13.0" dependencies = [ "borsh", "serde", @@ -11080,7 +11080,7 @@ dependencies = [ [[package]] name = "tari_engine" -version = "0.12.0" +version = "0.13.0" dependencies = [ "blake2", "cargo_toml 0.22.3", @@ -11109,7 +11109,7 @@ dependencies = [ [[package]] name = "tari_engine_types" -version = "0.12.0" +version = "0.13.0" dependencies = [ "base64 0.21.7", "bincode 2.0.1", @@ -11136,7 +11136,7 @@ dependencies = [ [[package]] name = "tari_epoch_manager" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "log", @@ -11158,7 +11158,7 @@ dependencies = [ [[package]] name = "tari_epoch_oracles" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "blake2", @@ -11199,7 +11199,7 @@ dependencies = [ [[package]] name = "tari_indexer" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "async-graphql", @@ -11263,7 +11263,7 @@ dependencies = [ [[package]] name = "tari_indexer_client" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bounded-vec", @@ -11293,7 +11293,7 @@ dependencies = [ [[package]] name = "tari_indexer_lib" -version = "0.12.0" +version = "0.13.0" dependencies = [ "log", "serde", @@ -11369,7 +11369,7 @@ dependencies = [ [[package]] name = "tari_networking" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "async-trait", @@ -11406,7 +11406,7 @@ dependencies = [ [[package]] name = "tari_ootle_address" -version = "0.12.0" +version = "0.13.0" dependencies = [ "bech32", "bincode 2.0.1", @@ -11422,7 +11422,7 @@ dependencies = [ [[package]] name = "tari_ootle_app_utilities" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bincode 2.0.1", @@ -11461,7 +11461,7 @@ dependencies = [ [[package]] name = "tari_ootle_common_types" -version = "0.12.0" +version = "0.13.0" dependencies = [ "blake2", "borsh", @@ -11489,7 +11489,7 @@ dependencies = [ [[package]] name = "tari_ootle_p2p" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "prost 0.14.1", @@ -11511,7 +11511,7 @@ dependencies = [ [[package]] name = "tari_ootle_storage" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bitflags 2.9.2", @@ -11536,7 +11536,7 @@ dependencies = [ [[package]] name = "tari_ootle_storage_sqlite" -version = "0.12.0" +version = "0.13.0" dependencies = [ "diesel", "diesel_migrations", @@ -11556,7 +11556,7 @@ dependencies = [ [[package]] name = "tari_ootle_wallet_cli" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "base64 0.22.1", @@ -11585,7 +11585,7 @@ dependencies = [ [[package]] name = "tari_ootle_wallet_crypto" -version = "0.12.0" +version = "0.13.0" dependencies = [ "argon2 0.5.3", "blake2", @@ -11609,7 +11609,7 @@ dependencies = [ [[package]] name = "tari_ootle_wallet_sdk" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "blake2", @@ -11643,7 +11643,7 @@ dependencies = [ [[package]] name = "tari_ootle_wallet_sdk_services" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "futures 0.3.31", @@ -11668,7 +11668,7 @@ dependencies = [ [[package]] name = "tari_ootle_wallet_storage_sqlite" -version = "0.12.0" +version = "0.13.0" dependencies = [ "bigdecimal", "diesel", @@ -11692,7 +11692,7 @@ dependencies = [ [[package]] name = "tari_ootle_walletd" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "async-trait", @@ -11777,7 +11777,7 @@ dependencies = [ [[package]] name = "tari_rpc_framework" -version = "0.12.0" +version = "0.13.0" dependencies = [ "async-trait", "bitflags 2.9.2", @@ -11801,7 +11801,7 @@ dependencies = [ [[package]] name = "tari_rpc_macros" -version = "0.12.0" +version = "0.13.0" dependencies = [ "proc-macro2", "quote", @@ -11810,7 +11810,7 @@ dependencies = [ [[package]] name = "tari_rpc_state_sync" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "futures 0.3.31", @@ -11830,7 +11830,7 @@ dependencies = [ [[package]] name = "tari_scaffolder" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "clap 3.2.25", @@ -11902,7 +11902,7 @@ dependencies = [ [[package]] name = "tari_signaling_server" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "axum 0.8.4", @@ -11928,7 +11928,7 @@ dependencies = [ [[package]] name = "tari_state_store_rocksdb" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bincode 2.0.1", @@ -11954,7 +11954,7 @@ dependencies = [ [[package]] name = "tari_state_tree" -version = "0.12.0" +version = "0.13.0" dependencies = [ "indexmap 2.11.4", "log", @@ -11981,7 +11981,7 @@ dependencies = [ [[package]] name = "tari_swarm" -version = "0.12.0" +version = "0.13.0" dependencies = [ "libp2p", "libp2p-messaging", @@ -11991,7 +11991,7 @@ dependencies = [ [[package]] name = "tari_swarm_daemon" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "async-trait", @@ -12047,7 +12047,7 @@ dependencies = [ [[package]] name = "tari_template_builtin" -version = "0.12.0" +version = "0.13.0" dependencies = [ "tari_engine_types", "tari_template_lib", @@ -12099,7 +12099,7 @@ dependencies = [ [[package]] name = "tari_template_manager" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bytes 1.10.1", @@ -12165,7 +12165,7 @@ dependencies = [ [[package]] name = "tari_transaction" -version = "0.12.0" +version = "0.13.0" dependencies = [ "borsh", "hex", @@ -12253,7 +12253,7 @@ dependencies = [ [[package]] name = "tari_transaction_manifest" -version = "0.12.0" +version = "0.13.0" dependencies = [ "proc-macro2", "syn 2.0.103", @@ -12286,7 +12286,7 @@ dependencies = [ [[package]] name = "tari_validator_node" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "axum 0.8.4", @@ -12348,7 +12348,7 @@ dependencies = [ [[package]] name = "tari_validator_node_cli" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "clap 3.2.25", @@ -12376,7 +12376,7 @@ dependencies = [ [[package]] name = "tari_validator_node_client" -version = "0.12.0" +version = "0.13.0" dependencies = [ "indexmap 2.11.4", "multiaddr 0.18.1", @@ -12399,7 +12399,7 @@ dependencies = [ [[package]] name = "tari_validator_node_rpc" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "prost 0.14.1", @@ -12421,7 +12421,7 @@ dependencies = [ [[package]] name = "tari_wallet_daemon_client" -version = "0.12.0" +version = "0.13.0" dependencies = [ "reqwest 0.11.27", "serde", @@ -12442,7 +12442,7 @@ dependencies = [ [[package]] name = "tari_watcher" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "clap 3.2.25", @@ -12470,7 +12470,7 @@ dependencies = [ [[package]] name = "tariswap_bench" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "clap 4.5.48", @@ -13129,7 +13129,7 @@ dependencies = [ [[package]] name = "transaction_generator" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "bincode 2.0.1", @@ -13149,7 +13149,7 @@ dependencies = [ [[package]] name = "transaction_submitter" -version = "0.12.0" +version = "0.13.0" dependencies = [ "anyhow", "clap 4.5.48", diff --git a/Cargo.toml b/Cargo.toml index bcc67ca6db..fd7817d7cf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ # NOTE: When editing this version, also edit the versions in template_built_in/templates/account and account_nft [workspace.package] -version = "0.12.0" +version = "0.13.0" edition = "2021" authors = ["The Tari Development Community"] repository = "https://github.com/tari-project/tari-ootle" diff --git a/applications/tari_app_utilities/src/seed_peer.rs b/applications/tari_app_utilities/src/seed_peer.rs index 29002f841e..5e772737de 100644 --- a/applications/tari_app_utilities/src/seed_peer.rs +++ b/applications/tari_app_utilities/src/seed_peer.rs @@ -38,6 +38,7 @@ impl SeedPeer { pub fn to_peer_id(&self) -> Option { let pk = self.public_key.as_ref()?; let pk = identity::PublicKey::from( + // invariant: we only construct SeedPeer with valid public keys identity::sr25519::PublicKey::try_from_bytes(pk.as_bytes()).expect("invariant: valid public key"), ); Some(pk.to_peer_id()) @@ -59,6 +60,7 @@ impl FromStr for SeedPeer { } let (pk, address) = s.split_once("::").ok_or_else(|| anyhow!("Invalid seed peer format"))?; + // Ensure that the public key is valid let public_key = RistrettoPublicKey::from_hex(pk).map_err(|err| anyhow!("Invalid public key {err}"))?; let address = address.parse().map_err(|err| anyhow!("Invalid address {err}"))?; Ok(SeedPeer { @@ -101,6 +103,14 @@ mod tests { assert!(seed_peer.to_peer_id().is_some()); } + #[test] + fn it_fails_to_parse_with_invalid_public_key() { + let s = "deadbeaf00000000000deadbeaf0000000000000000000000000000000000000::/ip4/127.0.0.1/tcp/8080"; + + let err = SeedPeer::from_str(s).unwrap_err(); + assert!(err.to_string().contains("Invalid public key")); + } + #[test] fn it_parses_without_public_key() { let s = "/ip4/127.0.0.1/tcp/8080"; diff --git a/applications/tari_indexer/src/bootstrap.rs b/applications/tari_indexer/src/bootstrap.rs index 0dd486e9a9..27b63adb3e 100644 --- a/applications/tari_indexer/src/bootstrap.rs +++ b/applications/tari_indexer/src/bootstrap.rs @@ -68,13 +68,8 @@ use crate::{ network_client::TariNetworkClient, network_state_sync, network_state_sync::NetworkWideStateSyncConfig, - storage_sqlite::{ - models::Key, - IndexerStore, - IndexerStoreReadTransaction, - IndexerStoreWriteTransaction, - SqliteIndexerStore, - }, + storage_sqlite::{models::Key, SqliteIndexerStore}, + store::{IndexerStore, IndexerStoreReadTransaction, IndexerStoreWriteTransaction}, substate_file_cache::SubstateFileCache, substate_manager::SubstateManager, transaction_manager::TransactionManager, diff --git a/applications/tari_indexer/src/event_manager.rs b/applications/tari_indexer/src/event_manager.rs index 1efbf0f323..25c22f4e30 100644 --- a/applications/tari_indexer/src/event_manager.rs +++ b/applications/tari_indexer/src/event_manager.rs @@ -20,13 +20,14 @@ // WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE // USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. -use std::{collections::BTreeMap, str::FromStr}; - use log::*; use tari_engine_types::{events::Event, substate::SubstateId}; -use tari_template_lib::{models::Metadata, types::Hash}; +use tari_transaction::TransactionId; -use crate::storage_sqlite::{IndexerStore, IndexerStoreReadTransaction, SqliteIndexerStore}; +use crate::{ + storage_sqlite::SqliteIndexerStore, + store::{IndexerStoreReadTransaction, IndexerStoreReader}, +}; const LOG_TARGET: &str = "tari::indexer::event_manager"; @@ -42,27 +43,16 @@ impl EventManager { pub async fn get_events_from_db( &self, - topic: Option, - substate_id: Option, + topic: Option<&str>, + substate_id: Option<&SubstateId>, offset: u32, limit: u32, - ) -> Result, anyhow::Error> { - let rows = self + ) -> Result, anyhow::Error> { + let events = self .substate_store .with_read_tx(|tx| tx.get_events(substate_id, topic, offset, limit))?; - debug!(target: LOG_TARGET, "Found {} events", rows.len()); - - let mut events = Vec::with_capacity(rows.len()); - for row in rows { - let substate_id = row.substate_id.map(|str| SubstateId::from_str(&str)).transpose()?; - let template_address = Hash::from_hex(&row.template_address)?; - let tx_hash = Hash::from_hex(&row.tx_hash)?; - let topic = row.topic; - let payload = Metadata::from(serde_json::from_str::>(row.payload.as_str())?); - events.push(Event::new(substate_id, template_address, tx_hash, topic, payload)); - } - + debug!(target: LOG_TARGET, "Found {} events", events.len()); Ok(events) } } diff --git a/applications/tari_indexer/src/graphql/model/events.rs b/applications/tari_indexer/src/graphql/model/events.rs index ddb169bd05..634bc1ab65 100644 --- a/applications/tari_indexer/src/graphql/model/events.rs +++ b/applications/tari_indexer/src/graphql/model/events.rs @@ -27,6 +27,7 @@ use log::*; use serde::{Deserialize, Serialize}; use tari_engine_types::substate::SubstateId; use tari_ootle_common_types::displayable::Displayable; +use tari_transaction::TransactionId; use crate::event_manager::EventManager; @@ -43,12 +44,15 @@ pub struct Event { } impl Event { - fn from_engine_event(event: tari_engine_types::events::Event) -> Result { + fn from_engine_event( + transaction_id: TransactionId, + event: tari_engine_types::events::Event, + ) -> Result { Ok(Self { substate_id: event.substate_id().map(|sub_id| sub_id.to_string()), template_address: event.template_address().into_array(), - tx_hash: event.tx_hash().into_array(), - topic: event.topic(), + tx_hash: transaction_id.into_array(), + topic: event.topic().to_string(), payload: event.into_payload().into_iter().collect(), }) } @@ -75,14 +79,18 @@ impl EventQuery { let substate_id = substate_id.map(|str| SubstateId::from_str(&str)).transpose()?; let event_manager = ctx.data_unchecked::(); let limit = limit.unwrap_or(100); + if limit == 0 { + return Ok(vec![]); + } + if limit > 1000 { return Err(anyhow::anyhow!("Limit cannot be greater than 1000")); } event_manager - .get_events_from_db(topic, substate_id, offset.unwrap_or(0), limit) + .get_events_from_db(topic.as_deref(), substate_id.as_ref(), offset.unwrap_or(0), limit) .await? .into_iter() - .map(Event::from_engine_event) + .map(|(id, ev)| Event::from_engine_event(id, ev)) .collect() } } diff --git a/applications/tari_indexer/src/lib.rs b/applications/tari_indexer/src/lib.rs index 7509f0a1d8..8c27970aa4 100644 --- a/applications/tari_indexer/src/lib.rs +++ b/applications/tari_indexer/src/lib.rs @@ -39,6 +39,7 @@ mod event_manager; mod network_client; mod network_state_sync; mod storage_sqlite; +mod store; mod substate_file_cache; mod substate_manager; mod transaction_manager; diff --git a/applications/tari_indexer/src/network_state_sync/block_scanner.rs b/applications/tari_indexer/src/network_state_sync/block_scanner.rs index 1950e71bbe..965545c0d3 100644 --- a/applications/tari_indexer/src/network_state_sync/block_scanner.rs +++ b/applications/tari_indexer/src/network_state_sync/block_scanner.rs @@ -12,13 +12,8 @@ use tari_validator_node_rpc::client::{TariValidatorNodeRpcClientFactory, Validat use crate::{ block_data::BlockData, - storage_sqlite::{ - models::NewScannedBlockId, - IndexerStore, - IndexerStoreReadTransaction, - IndexerStoreWriteTransaction, - SqliteIndexerStore, - }, + storage_sqlite::{models::NewScannedBlockId, SqliteIndexerStore}, + store::{IndexerStore, IndexerStoreReadTransaction, IndexerStoreReader, IndexerStoreWriteTransaction}, }; const LOG_TARGET: &str = "tari::indexer::block_scanner"; diff --git a/applications/tari_indexer/src/network_state_sync/committee_client.rs b/applications/tari_indexer/src/network_state_sync/committee_client.rs index 83cf0a581c..bbbe0a3b7f 100644 --- a/applications/tari_indexer/src/network_state_sync/committee_client.rs +++ b/applications/tari_indexer/src/network_state_sync/committee_client.rs @@ -40,11 +40,17 @@ impl ValidatorCommitteeRpcPool { pub async fn new_session(&mut self) -> Result { let epoch = self.epoch_manager.get_current_epoch(); + let mut last_error = None::; loop { let member = self .epoch_manager .get_random_committee_member(epoch, Some(self.shard_group), self.past_failed_nodes.clone()) - .await?; + .await + .optional()? + .ok_or_else(|| ValidatorCommitteeClientError::AllValidatorsFailed { + committee_size: self.past_failed_nodes.len(), + last_error: last_error.as_ref().map(|e| e.to_string()), + })?; let result = self.session_for_peer(member.address).await; match result { Ok(session) => return Ok(session), @@ -53,6 +59,7 @@ impl ValidatorCommitteeRpcPool { target: LOG_TARGET, "Failed to create session for validator '{}': {}", member, err ); + last_error = Some(err); self.past_failed_nodes.push(member.address); }, } @@ -94,6 +101,10 @@ impl ValidatorCommitteeRpcPool { .await .optional()? else { + warn!( + target: LOG_TARGET, + "All {} committee members have been attempted and failed.", attempted.len() + ); // No more committee members to try break; }; diff --git a/applications/tari_indexer/src/network_state_sync/event_filter.rs b/applications/tari_indexer/src/network_state_sync/event_filter.rs index dfc33fe603..45e74e7aba 100644 --- a/applications/tari_indexer/src/network_state_sync/event_filter.rs +++ b/applications/tari_indexer/src/network_state_sync/event_filter.rs @@ -7,7 +7,7 @@ use tari_template_lib::types::{EntityId, TemplateAddress}; #[derive(Default, Debug, Serialize, Deserialize, Clone)] pub struct EventFilter { - pub topic: Option, + pub topic: Option>, pub entity_id: Option, pub substate_id: Option, pub template_address: Option, @@ -15,24 +15,31 @@ pub struct EventFilter { impl EventFilter { pub fn matches(&self, event: &Event) -> bool { - let matches_topic = self.topic.as_ref().is_none_or(|t| *t == event.topic()); - let matches_template = self + if self.topic.as_ref().is_some_and(|t| t.as_ref() != event.topic()) { + return false; + } + + if self .template_address .as_ref() - .is_none_or(|t| *t == event.template_address()); + .is_some_and(|t| *t != event.template_address()) + { + return false; + } - let matches_substate_id = self + if self .substate_id .as_ref() - .is_none_or(|substate_id| event.substate_id().map(|s| s == substate_id).unwrap_or(false)); + .is_some_and(|substate_id| event.substate_id().map(|s| s != substate_id).unwrap_or(true)) + { + return false; + } - let matches_entity_id = self.entity_id.as_ref().is_none_or(|entity_id| { + self.entity_id.as_ref().is_none_or(|entity_id| { event .substate_id() .map(|s| s.to_object_key().as_entity_id() == *entity_id) .unwrap_or(false) - }); - - matches_topic && matches_template && matches_substate_id && matches_entity_id + }) } } diff --git a/applications/tari_indexer/src/network_state_sync/worker.rs b/applications/tari_indexer/src/network_state_sync/worker.rs index 53e27b16d6..601ce8a55e 100644 --- a/applications/tari_indexer/src/network_state_sync/worker.rs +++ b/applications/tari_indexer/src/network_state_sync/worker.rs @@ -7,7 +7,7 @@ use futures::StreamExt; use log::*; use tari_engine_types::{ substate::{SubstateId, SubstateValue}, - transaction_receipt::TransactionReceipt, + transaction_receipt::{TransactionReceipt, TransactionReceiptAddress}, }; use tari_epoch_manager::{service::EpochManagerHandle, EpochManagerEvent, EpochManagerReader}; use tari_networking::NetworkingHandle; @@ -41,12 +41,10 @@ use crate::{ }, storage_sqlite::{ models::{Key, UtxoSpent, UtxoUnspent, UtxoUpdateRecord}, - IndexerStore, - IndexerStoreReadTransaction, - IndexerStoreWriteTransaction, SqliteIndexerStore, SqliteStoreWriteTransaction, }, + store::{IndexerStore, IndexerStoreReadTransaction, IndexerStoreReader, IndexerStoreWriteTransaction}, }; const LOG_TARGET: &str = "tari::indexer::network_state_sync::worker"; @@ -233,7 +231,7 @@ impl NetworkWideStateSync { // TODO: continue on failure for checkpoint in checkpoints { - info!(target: LOG_TARGET, "🌍️ Validating checkpoint for shard group {shard_group}: {:?}", checkpoint); + info!(target: LOG_TARGET, "🌍️ Validating checkpoint for shard group {shard_group}: {}", checkpoint.header().calculate_hash()); // TODO: we require historical committees to validate older checkpoints. Figure out the best way to // avoid needing the data (e.g. VN merkle inclusion proof + historic L1 block MR), or, // decide it is ok to require this data to be locally stored by all indexers. For now, to avoid @@ -333,7 +331,7 @@ impl NetworkWideStateSync { sync_plan_mut: &mut SyncPlan, update_buf: &mut Vec<(Epoch, SubstateUpdateProof)>, utxos_buf: &mut Vec, - transactions_buf: &mut Vec, + transactions_buf: &mut Vec<(TransactionReceiptAddress, TransactionReceipt)>, validator_fee_pools_buf: &mut Vec, shard_group: ShardGroup, session: &mut ValidatorRpcSession, @@ -411,6 +409,7 @@ impl NetworkWideStateSync { self.store.clone().with_write_tx(|tx| { debug!(target: LOG_TARGET, "✅ Committing {} updates for shard {shard} (epoch: {msg_epoch}, state version: {state_version})", update_buf.len()); + // TODO: this is not currently used. Consider removing. tx.batch_insert_substate_transitions(shard, state_version, update_buf.drain(..))?; debug!(target: LOG_TARGET, "✅ Committing {} UTXOs for shard {shard} (epoch: {msg_epoch})", utxos_buf.len()); tx.batch_insert_utxo_updates(utxos_buf.drain(..))?; @@ -439,21 +438,12 @@ impl NetworkWideStateSync { Ok(()) } - fn persist_transaction_receipts>( + fn persist_transaction_receipts>( &self, tx: &mut SqliteStoreWriteTransaction<'_>, receipts: I, ) -> Result<(), StorageError> { - let events = receipts - .into_iter() - .flat_map(|receipt| receipt.events) - // only keep the events specified by the indexer filter - .filter(|event| { - self.config.event_filters.is_empty() || self.config.event_filters.iter().any(|filter| filter.matches(event)) - }); - - tx.batch_insert_events(events)?; - + tx.batch_insert_transaction_receipts(receipts, &self.config.event_filters)?; Ok(()) } } @@ -466,7 +456,7 @@ fn extend_bufs_from_substate_update( update_buf: &mut Vec<(Epoch, SubstateUpdateProof)>, templates_buf: &mut Vec, utxos_buf: &mut Vec, - transactions_buf: &mut Vec, + transactions_buf: &mut Vec<(TransactionReceiptAddress, TransactionReceipt)>, validator_fee_pools_buf: &mut Vec, ) -> Result<(), NetworkStateSyncError> { match &update { @@ -489,7 +479,9 @@ fn extend_bufs_from_substate_update( }; }, Some(SubstateValue::TransactionReceipt(receipt)) => { - transactions_buf.push(receipt.clone()); + if let Some(address) = update.substate_id().as_transaction_receipt_address() { + transactions_buf.push((address, receipt.clone())); + } }, Some(SubstateValue::Template(template)) => { if let Some(address) = create.substate.substate_id().as_template() { diff --git a/applications/tari_indexer/src/rest_api/context.rs b/applications/tari_indexer/src/rest_api/context.rs index e5ec8546d9..098aab2580 100644 --- a/applications/tari_indexer/src/rest_api/context.rs +++ b/applications/tari_indexer/src/rest_api/context.rs @@ -23,6 +23,7 @@ use crate::{ dry_run::processor::DryRunTransactionProcessor, rest_api::cache::HttpCacheConfig, storage_sqlite::SqliteIndexerStore, + store::ReadOnlyStore, substate_manager::SubstateManager, transaction_manager::TransactionManager, }; @@ -37,6 +38,7 @@ impl HandlerContext { Self { inner: Arc::new(InnerContext { cache_control_enabled: true, + read_only_store: ReadOnlyStore::new(services.store.clone()), global_db: services.global_db.clone(), epoch_manager: services.epoch_manager.clone(), networking: services.networking.clone(), @@ -80,6 +82,10 @@ impl HandlerContext { &self.inner.transaction_manager } + pub fn read_only_store(&self) -> &ReadOnlyStore { + &self.inner.read_only_store + } + pub fn template_manager(&self) -> &TemplateManager { &self.inner.template_manager } @@ -111,6 +117,7 @@ struct InnerContext { substate_manager: SubstateManager, transaction_manager: TransactionManager, TariValidatorNodeRpcClientFactory, SqliteIndexerStore>, + read_only_store: ReadOnlyStore, template_manager: TemplateManager, dry_run_transaction_processor: DryRunTransactionProcessor, } diff --git a/applications/tari_indexer/src/rest_api/handlers/mod.rs b/applications/tari_indexer/src/rest_api/handlers/mod.rs index 21a22a345d..baf6767ae2 100644 --- a/applications/tari_indexer/src/rest_api/handlers/mod.rs +++ b/applications/tari_indexer/src/rest_api/handlers/mod.rs @@ -6,6 +6,7 @@ pub mod network; pub mod nfts; pub mod substates; pub mod templates; +pub mod transaction_receipts; pub mod transactions; pub mod utxos; diff --git a/applications/tari_indexer/src/rest_api/handlers/transaction_receipts.rs b/applications/tari_indexer/src/rest_api/handlers/transaction_receipts.rs new file mode 100644 index 0000000000..a025c29934 --- /dev/null +++ b/applications/tari_indexer/src/rest_api/handlers/transaction_receipts.rs @@ -0,0 +1,57 @@ +// Copyright 2025 The Tari Project +// SPDX-License-Identifier: BSD-3-Clause + +use axum::{ + extract::{Path, Query}, + response::Response, + Extension, + Json, +}; +use tari_engine_types::transaction_receipt::TransactionReceiptAddress; +use tari_indexer_client::types::{ + GetTransactionReceiptResponse, + ListTransactionReceiptsRequest, + ListTransactionReceiptsResponse, +}; +use tari_ootle_common_types::optional::Optional; + +use crate::rest_api::{context::HandlerContext, error::ErrorResponse, handlers::HandlerResult}; + +#[utoipa::path(get, path = "/transaction-receipts", description = "List transaction receipts")] +pub async fn list_transaction_receipts( + Extension(context): Extension, + Query(req): Query, +) -> HandlerResult { + let limit = req.limit.unwrap_or(100); + if limit > 100 { + return Err(ErrorResponse::bad_request( + "Limit cannot be greater than 100".to_string(), + )); + } + + let receipts = context + .read_only_store() + .list_transaction_receipts(req.last_id, u64::from(limit), req.ordering) + .map_err(ErrorResponse::anyhow)?; + + Ok(context.apply_cache_control(Json(ListTransactionReceiptsResponse { receipts }), 10)) +} + +#[utoipa::path( + get, + path = "/transaction-receipt/{transaction_id}", + description = "Get the transaction receipt by transaction ID" +)] +pub async fn get_transaction_receipt( + Extension(context): Extension, + Path(receipt_addr): Path, +) -> HandlerResult> { + let receipt = context + .read_only_store() + .get_transaction_receipt(&receipt_addr) + .optional() + .map_err(ErrorResponse::anyhow)? + .ok_or_else(|| ErrorResponse::not_found(format!("Transaction receipt {receipt_addr} not found")))?; + + Ok(Json(GetTransactionReceiptResponse { receipt })) +} diff --git a/applications/tari_indexer/src/rest_api/handlers/transactions.rs b/applications/tari_indexer/src/rest_api/handlers/transactions.rs index 54104c90a7..cd2f46eafc 100644 --- a/applications/tari_indexer/src/rest_api/handlers/transactions.rs +++ b/applications/tari_indexer/src/rest_api/handlers/transactions.rs @@ -112,7 +112,7 @@ pub async fn list_recent_transactions( } #[utoipa::path( - post, + get, path = "/transactions/{transaction_id}/result", description = "Get the result of a submitted transaction (by transaction ID)" )] diff --git a/applications/tari_indexer/src/rest_api/server.rs b/applications/tari_indexer/src/rest_api/server.rs index 79c1d6be9e..061c5f001a 100644 --- a/applications/tari_indexer/src/rest_api/server.rs +++ b/applications/tari_indexer/src/rest_api/server.rs @@ -43,6 +43,8 @@ const REQUEST_BODY_LIMIT: usize = 4 * 1024 * 1024; // 4 MB handlers::templates::list_templates, handlers::utxos::fetch_utxos, handlers::utxos::stream_utxo_updates, + handlers::transaction_receipts::list_transaction_receipts, + handlers::transaction_receipts::get_transaction_receipt ))] pub struct ApiDoc; @@ -93,6 +95,15 @@ impl Server { .route("/fetch", post(handlers::utxos::fetch_utxos)) .route("/stream", post(handlers::utxos::stream_utxo_updates)) ) + .nest( + "/transaction-receipts", + Router::new() + .route("/", get(handlers::transaction_receipts::list_transaction_receipts)) + .route( + "/{address}", + get(handlers::transaction_receipts::get_transaction_receipt), + ), + ) .layer(CorsLayer::permissive()) .layer(RequestBodyLimitLayer::new(REQUEST_BODY_LIMIT)) .merge(SwaggerUi::new("/swagger-ui").url("/openapi.json", ApiDoc::openapi())) diff --git a/applications/tari_indexer/src/storage_sqlite/migrations/2023-02-16-145719_initial/up.sql b/applications/tari_indexer/src/storage_sqlite/migrations/2023-02-16-145719_initial/up.sql index 07a821ef26..731837ec52 100644 --- a/applications/tari_indexer/src/storage_sqlite/migrations/2023-02-16-145719_initial/up.sql +++ b/applications/tari_indexer/src/storage_sqlite/migrations/2023-02-16-145719_initial/up.sql @@ -62,6 +62,17 @@ create table events -- DB index for faster collection scan queries create index events_indexer on events (template_address, tx_hash); +-- Transaction receipts +create table transaction_receipts +( + id integer not NULL primary key AUTOINCREMENT, + address text not NULL, + data text not NULL, + created_at timestamp not null default current_timestamp +); + +create unique index transaction_receipts_address_uniq on transaction_receipts (address); + -- Latest scanned blocks, separately by committee (epoch + shard) -- Used mostly for efficient scanning of events in the whole network create table scanned_block_ids diff --git a/applications/tari_indexer/src/storage_sqlite/models/events.rs b/applications/tari_indexer/src/storage_sqlite/models/events.rs index a0c2474f49..6586637bd5 100644 --- a/applications/tari_indexer/src/storage_sqlite/models/events.rs +++ b/applications/tari_indexer/src/storage_sqlite/models/events.rs @@ -21,11 +21,10 @@ // USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. // -use std::{convert::TryFrom, str::FromStr}; +use std::convert::TryFrom; use diesel::sql_types::{Nullable, Text}; use serde::{Deserialize, Serialize}; -use tari_engine_types::substate::SubstateId; use tari_ootle_storage::time::PrimitiveDateTime; use tari_template_lib::types::Hash; @@ -43,13 +42,13 @@ pub struct EventRecord { pub created_at: PrimitiveDateTime, } -#[derive(Debug, Clone, Insertable, AsChangeset)] +#[derive(Debug, Clone, Insertable)] #[diesel(table_name = events)] #[diesel(treat_none_as_null = true)] -pub struct NewEvent { +pub struct NewEvent<'a> { pub template_address: String, - pub tx_hash: String, - pub topic: String, + pub tx_hash: &'a str, + pub topic: &'a str, pub payload: String, pub substate_id: Option, } @@ -90,28 +89,26 @@ impl TryFrom for crate::graphql::model::events::Event { } } -impl TryFrom for tari_engine_types::events::Event { - type Error = anyhow::Error; - - fn try_from(event_data: EventData) -> Result { - let substate_id = event_data - .substate_id - .clone() - .map(|sub_id| SubstateId::from_str(&sub_id)) - .transpose()?; - let template_address = Hash::from_hex(&event_data.template_address)?; - let tx_hash = Hash::from_hex(&event_data.tx_hash)?; - let payload = serde_json::from_str(event_data.payload.as_str())?; - - Ok(Self::new( - substate_id, - template_address, - tx_hash, - event_data.topic, - payload, - )) - } -} +// impl TryFrom for tari_engine_types::events::Event { +// type Error = anyhow::Error; +// +// fn try_from(event_data: EventData) -> Result { +// let substate_id = event_data +// .substate_id +// .clone() +// .map(|sub_id| SubstateId::from_str(&sub_id)) +// .transpose()?; +// let template_address = Hash::from_hex(&event_data.template_address)?; +// let payload = serde_json::from_str(event_data.payload.as_str())?; +// +// Ok(Self::new( +// substate_id, +// template_address, +// event_data.topic, +// payload, +// )) +// } +// } // To keep track of the latest blocks that we scanned for events diff --git a/applications/tari_indexer/src/storage_sqlite/reader.rs b/applications/tari_indexer/src/storage_sqlite/reader.rs index 7c6ac1c15b..aa55a4db98 100644 --- a/applications/tari_indexer/src/storage_sqlite/reader.rs +++ b/applications/tari_indexer/src/storage_sqlite/reader.rs @@ -8,6 +8,7 @@ use diesel::{ sql_types, BoolExpressionMethods, ExpressionMethods, + NullableExpressionMethods, OptionalExtension, QueryDsl, RunQueryDsl, @@ -18,18 +19,27 @@ use log::info; use serde::de::DeserializeOwned; use tari_consensus_types::BlockId; use tari_engine_types::{ + events::Event, substate::{SubstateId, SubstateValue}, + transaction_receipt::{TransactionReceipt, TransactionReceiptAddress}, Utxo, }; use tari_indexer_client::types::{ListSubstateItem, NonFungibleSubstate, TransactionEntry}; -use tari_ootle_common_types::{shard::Shard, substate_type::SubstateType, Epoch, ShardGroup, StateVersion}; -use tari_ootle_storage::{time::PrimitiveDateTime, StorageError}; +use tari_ootle_common_types::{ + displayable::Displayable, + shard::Shard, + substate_type::SubstateType, + Epoch, + ShardGroup, + StateVersion, +}; +use tari_ootle_storage::{time::PrimitiveDateTime, Ordering, StorageError}; use tari_ootle_storage_sqlite::SqliteTransaction; use tari_ootle_wallet_sdk::models::WalletUtxoUpdate; use tari_template_lib::{ models::{ResourceAddress, UtxoId}, prelude::{RistrettoPublicKeyBytes, TemplateAddress}, - types::crypto::UtxoTag, + types::{crypto::UtxoTag, Hash}, }; use tari_transaction::{Transaction, TransactionId}; @@ -38,8 +48,8 @@ use crate::{ models, models::{EventRecord, KeyValue, ScannedBlockId, SubstateRecord}, serialization::{deserialize_hex_try_from, deserialize_json, serialize_hex}, - IndexerStoreReadTransaction, }, + store::IndexerStoreReadTransaction, substate_manager::SubstateResponse, }; @@ -87,17 +97,22 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { query = query.offset(offset as i64); } - let substates: Vec = query + let substates = query .order_by(substates::id.desc()) - .get_results(self.connection()) + .load_iter::(self.connection()) .map_err(|e| StorageError::QueryError { reason: format!("list_substates: {}", e), })?; let items = substates .into_iter() - .map(|s| { - let substate_id = SubstateId::from_str(&s.address)?; + .map(|res| { + let s = res.map_err(|e| StorageError::QueryError { + reason: format!("list_substates: {}", e), + })?; + let substate_id = SubstateId::from_str(&s.address).map_err(|e| StorageError::DataInconsistency { + details: format!("Invalid substate address {}: {}", s.address, e), + })?; let version = s.version as u32; let template_address = s.template_address.map(|h| deserialize_hex_try_from(&h)).transpose()?; let timestamp = s.updated_at; @@ -109,7 +124,7 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { timestamp, }) }) - .collect::, anyhow::Error>>() + .collect::, StorageError>>() .map_err(|e| StorageError::QueryError { reason: format!("list_substates: invalid substate items: {}", e), })?; @@ -150,12 +165,19 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { let rows = substates::table .select(substates::all_columns) .filter(substates::address.eq_any(str_ids)) - .get_results::(self.connection()) + .load_iter::(self.connection()) .map_err(|e| StorageError::QueryError { - reason: format!("get_substate: {}", e), + reason: format!("get_substates: {}", e), })?; - rows.into_iter().map(TryInto::try_into).collect() + rows.into_iter() + .map(|res| { + res.map_err(|e| StorageError::QueryError { + reason: format!("get_substates: {e}"), + }) + .and_then(TryInto::try_into) + }) + .collect() } fn get_non_fungible_count(&mut self, resource_address: String) -> Result { @@ -184,13 +206,16 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { .filter(substates::address.like(format!("nft_{}_%", resource_address.as_object_key()))) .limit(limit as i64) .offset(offset as i64) - .get_results::(self.connection()) + .load_iter::(self.connection()) .map_err(|e| StorageError::QueryError { reason: format!("get_non_fungibles_by_resource_address: {}", e), })?; res.into_iter() - .map(|row| { + .map(|res| { + let row = res.map_err(|e| StorageError::QueryError { + reason: format!("get_non_fungibles_by_resource_address: {}", e), + })?; let value: SubstateValue = serde_json::from_str(&row.data).map_err(|e| StorageError::DataInconsistency { details: format!("Failed to parse substate data: {}", e), @@ -211,11 +236,11 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { fn get_events( &mut self, - substate_id_filter: Option, - topic_filter: Option, + substate_id_filter: Option<&SubstateId>, + topic_filter: Option<&str>, offset: u32, limit: u32, - ) -> Result, StorageError> { + ) -> Result, StorageError> { // TODO: allow to query by payload as well, unifying all event methods into one info!( target: LOG_TARGET, @@ -236,19 +261,52 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { query = query.filter(events::topic.eq(topic)); } - query = query.offset(offset.into()); - if limit > 0 { - query = query.limit(limit.into()); - } - - let events = query + let event_rows = query + .offset(offset.into()) + .limit(limit.into()) .order(events::id.desc()) - .get_results::(self.connection()) + .load_iter::(self.connection()) .map_err(|e| StorageError::QueryError { reason: format!("get_events: {}", e), })?; - Ok(events) + event_rows + .map(|res| { + res.map_err(|e| StorageError::QueryError { + reason: format!("get_events: {}", e), + }) + .and_then(|row| { + let substate_id = row + .substate_id + .as_ref() + .map(|str| SubstateId::from_str(str)) + .transpose() + .map_err(|e| StorageError::DataInconsistency { + details: format!( + "Invalid substate_id {} in events table: {}", + row.substate_id.display(), + e + ), + })?; + let template_address = + Hash::from_hex(&row.template_address).map_err(|e| StorageError::DataInconsistency { + details: format!( + "Invalid template_address {} in events table: {}", + row.template_address, e + ), + })?; + let tx_hash = Hash::from_hex(&row.tx_hash).map_err(|e| StorageError::DataInconsistency { + details: format!("Invalid tx_hash {} in events table: {}", row.tx_hash, e), + })?; + let topic = row.topic; + let payload = deserialize_json(&row.payload)?; + Ok(( + TransactionId::from(tx_hash), + Event::new(substate_id, template_address, topic, payload), + )) + }) + }) + .collect() } fn get_oldest_scanned_epoch(&mut self) -> Result, StorageError> { @@ -340,6 +398,80 @@ impl IndexerStoreReadTransaction for SqliteStoreReadTransaction<'_> { .collect() } + fn list_transaction_receipts( + &mut self, + last_id: Option, + limit: u64, + ordering: Ordering, + ) -> Result, StorageError> { + const OPERATION: &str = "list_transaction_receipts"; + use crate::storage_sqlite::schema::transaction_receipts; + + let mut query = transaction_receipts::table + .select((transaction_receipts::address, transaction_receipts::data)) + .into_boxed(); + if let Some(last_id) = last_id { + let tr = alias!(transaction_receipts as tr); + let subquery = tr + .select(tr.field(transaction_receipts::id)) + .filter(transaction_receipts::address.eq(last_id.to_string())) + .limit(1) + .single_value() + .assume_not_null(); + + query = match ordering { + Ordering::Ascending => query.filter(transaction_receipts::id.gt(subquery)), + Ordering::Descending => query.filter(transaction_receipts::id.lt(subquery)), + } + } + + query = match ordering { + Ordering::Ascending => query.order_by(transaction_receipts::id.asc()), + Ordering::Descending => query.order_by(transaction_receipts::id.desc()), + }; + + let rows = query + .limit(limit as i64) + .load_iter::<(String, String), _>(self.connection()) + .map_err(|e| StorageError::QueryError { + reason: format!("{OPERATION}: {}", e), + })?; + + rows.into_iter() + .map(|row| { + row.map_err(|e| StorageError::QueryError { + reason: format!("{OPERATION}: {}", e), + }) + .and_then(|(addr, data)| { + Ok(( + TransactionReceiptAddress::from_str(&addr).map_err(|e| StorageError::DataInconsistency { + details: format!("Invalid transaction receipt address {}: {}", addr, e), + })?, + deserialize_json(data)?, + )) + }) + }) + .collect() + } + + fn get_transaction_receipt( + &mut self, + address: &TransactionReceiptAddress, + ) -> Result { + const OPERATION: &str = "get_transaction_receipt"; + use crate::storage_sqlite::schema::transaction_receipts; + + let receipt_entry = transaction_receipts::table + .select(transaction_receipts::data) + .filter(transaction_receipts::address.eq(address.to_string())) + .first::(self.connection()) + .map_err(|e| StorageError::QueryError { + reason: format!("{OPERATION}: {}", e), + })?; + + deserialize_json(&receipt_entry) + } + // -------------------------------- KeyValues -------------------------------- // fn key_value_get_value, T: DeserializeOwned>(&mut self, key: K) -> Result { let key_value = self.key_value_get_raw(key)?; diff --git a/applications/tari_indexer/src/storage_sqlite/schema.rs b/applications/tari_indexer/src/storage_sqlite/schema.rs index 20972737ef..b628c21c38 100644 --- a/applications/tari_indexer/src/storage_sqlite/schema.rs +++ b/applications/tari_indexer/src/storage_sqlite/schema.rs @@ -79,6 +79,15 @@ diesel::table! { } } +diesel::table! { + transaction_receipts (id) { + id -> Integer, + address -> Text, + data -> Text, + created_at -> Timestamp, + } +} + diesel::table! { transactions (id) { id -> Integer, @@ -114,6 +123,7 @@ diesel::allow_tables_to_appear_in_same_query!( scanned_block_ids, substate_transitions, substates, + transaction_receipts, transactions, utxos, ); diff --git a/applications/tari_indexer/src/storage_sqlite/store_factory.rs b/applications/tari_indexer/src/storage_sqlite/store_factory.rs index d4d32ae599..957dba852c 100644 --- a/applications/tari_indexer/src/storage_sqlite/store_factory.rs +++ b/applications/tari_indexer/src/storage_sqlite/store_factory.rs @@ -4,40 +4,18 @@ use std::{ fmt::Debug, fs::create_dir_all, - ops::{Deref, DerefMut}, path::PathBuf, sync::{Arc, Mutex}, }; use diesel::{sql_query, Connection, RunQueryDsl, SqliteConnection}; use diesel_migrations::{EmbeddedMigrations, MigrationHarness}; -use serde::{de::DeserializeOwned, Serialize}; -use tari_consensus_types::BlockId; -use tari_engine_types::{events::Event, substate::SubstateId, Utxo}; -use tari_indexer_client::types::{ListSubstateItem, NonFungibleSubstate, TransactionEntry}; -use tari_ootle_common_types::{shard::Shard, substate_type::SubstateType, Epoch, ShardGroup, StateVersion}; -use tari_ootle_storage::{ - consensus_models::{EpochCheckpoint, SubstateData, SubstateUpdateProof}, - StorageError, -}; +use tari_ootle_storage::StorageError; use tari_ootle_storage_sqlite::{error::SqliteStorageError, SqliteTransaction}; -use tari_ootle_wallet_sdk::models::WalletUtxoUpdate; -use tari_template_lib::{ - models::{ResourceAddress, UtxoId}, - types::{ - crypto::{RistrettoPublicKeyBytes, UtxoTag}, - TemplateAddress, - }, -}; -use tari_transaction::{Transaction, TransactionId}; use crate::{ - storage_sqlite::{ - models::{EventRecord, KeyValue, NewScannedBlockId, SubstateRecord, UtxoUpdateRecord}, - reader::SqliteStoreReadTransaction, - writer::SqliteStoreWriteTransaction, - }, - substate_manager::SubstateResponse, + storage_sqlite::{reader::SqliteStoreReadTransaction, writer::SqliteStoreWriteTransaction}, + store::{IndexerStore, IndexerStoreReader}, }; const LOG_TARGET: &str = "tari::indexer::storage_sqlite"; @@ -62,7 +40,7 @@ impl SqliteIndexerStore { .execute(&mut connection) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "set pragma".to_string(), + operation: "set pragma", })?; Ok(Self { @@ -75,145 +53,21 @@ impl Debug for SqliteIndexerStore { write!(f, "SqliteIndexerStore {{ connection: ... }}") } } -pub trait IndexerStore { - type ReadTransaction<'a>: IndexerStoreReadTransaction - where Self: 'a; - type WriteTransaction<'a>: IndexerStoreWriteTransaction + Deref> + DerefMut - where Self: 'a; - fn create_read_tx(&self) -> Result, StorageError>; - fn create_write_tx(&self) -> Result, StorageError>; - - fn with_write_tx) -> Result, R, E>(&self, f: F) -> Result - where E: From { - let mut tx = self.create_write_tx()?; - match f(&mut tx) { - Ok(r) => { - tx.commit()?; - Ok(r) - }, - Err(e) => { - if let Err(err) = tx.rollback() { - log::error!(target: LOG_TARGET, "Failed to rollback transaction: {}", err); - } - Err(e) - }, - } - } - - fn with_read_tx) -> Result, R, E>(&self, f: F) -> Result - where E: From { - let mut tx = self.create_read_tx()?; - let ret = f(&mut tx)?; - Ok(ret) - } -} - -impl IndexerStore for SqliteIndexerStore { +impl IndexerStoreReader for SqliteIndexerStore { type ReadTransaction<'a> = SqliteStoreReadTransaction<'a>; - type WriteTransaction<'a> = SqliteStoreWriteTransaction<'a>; fn create_read_tx(&self) -> Result, StorageError> { let tx = SqliteTransaction::begin(self.connection.lock().unwrap())?; Ok(SqliteStoreReadTransaction::new(tx)) } +} + +impl IndexerStore for SqliteIndexerStore { + type WriteTransaction<'a> = SqliteStoreWriteTransaction<'a>; fn create_write_tx(&self) -> Result, StorageError> { let tx = SqliteTransaction::begin(self.connection.lock().unwrap())?; Ok(SqliteStoreWriteTransaction::new(tx)) } } - -pub trait IndexerStoreReadTransaction { - fn list_substates( - &mut self, - filter_by_type: Option, - filter_by_template: Option, - limit: Option, - offset: Option, - ) -> Result, StorageError>; - fn get_substate( - &mut self, - address: &SubstateId, - version: Option, - ) -> Result, StorageError>; - - fn get_substates(&mut self, ids: &[SubstateId]) -> Result, StorageError>; - fn get_non_fungible_count(&mut self, resource_address: String) -> Result; - fn get_non_fungibles_by_resource_address( - &mut self, - resource_address: ResourceAddress, - limit: usize, - offset: usize, - ) -> Result, StorageError>; - - fn get_events( - &mut self, - substate_id_filter: Option, - topic_filter: Option, - offset: u32, - limit: u32, - ) -> Result, StorageError>; - fn get_oldest_scanned_epoch(&mut self) -> Result, StorageError>; - fn get_last_scanned_block_id( - &mut self, - epoch: Epoch, - shard_group: ShardGroup, - ) -> Result, StorageError>; - fn list_recent_transactions( - &mut self, - last_transaction_id: Option, - limit: usize, - ) -> Result, StorageError>; - - // -------------------------------- KeyValues -------------------------------- // - fn key_value_get_value, T: DeserializeOwned>(&mut self, key: K) -> Result; - fn key_value_get_raw>(&mut self, key: K) -> Result, StorageError>; - - fn utxos_get_max_state_version( - &mut self, - resource_address: ResourceAddress, - shard: Shard, - ) -> Result; - - /// Get UTXO updates for a given resource address and shard, starting from a specific state version. - /// - /// Returns a tuple containing the maximum returned state version, and a vector of UTXO - /// updates. - fn utxos_get_updates( - &mut self, - resource_address: ResourceAddress, - shard: Shard, - from_state_version: StateVersion, - unspents_only: bool, - limit: u32, - ) -> Result<(StateVersion, Vec), StorageError>; - - fn utxos_get_unspent_by_public_nonce_and_tag( - &mut self, - resource_address: &ResourceAddress, - public_nonce_and_tag: &[(UtxoTag, RistrettoPublicKeyBytes)], - ) -> Result, StorageError>; -} - -pub trait IndexerStoreWriteTransaction { - fn commit(self) -> Result<(), StorageError>; - fn rollback(self) -> Result<(), StorageError>; - fn key_value_set, V: Serialize>(&mut self, key: K, value: V) -> Result<(), StorageError>; - fn batch_insert_substate_transitions>( - &mut self, - shard: Shard, - state_version: StateVersion, - updates: I, - ) -> Result<(), StorageError>; - fn batch_insert_utxo_updates>( - &mut self, - updates: I, - ) -> Result<(), StorageError>; - fn upsert_substate(&mut self, substate: &SubstateData) -> Result<(), StorageError>; - fn batch_insert_events>(&mut self, events: I) -> Result<(), StorageError>; - fn save_scanned_block_id(&mut self, new_scanned_block_id: NewScannedBlockId) -> Result<(), StorageError>; - fn delete_scanned_epochs_older_than(&mut self, epoch: Epoch) -> Result<(), StorageError>; - fn insert_or_ignore_transaction(&mut self, transaction: &Transaction) -> Result<(), StorageError>; - fn insert_or_ignore_epoch_checkpoint(&mut self, epoch_checkpoint: &EpochCheckpoint) -> Result<(), StorageError>; -} diff --git a/applications/tari_indexer/src/storage_sqlite/writer.rs b/applications/tari_indexer/src/storage_sqlite/writer.rs index 4f932cb5d5..a57ca68adc 100644 --- a/applications/tari_indexer/src/storage_sqlite/writer.rs +++ b/applications/tari_indexer/src/storage_sqlite/writer.rs @@ -6,7 +6,7 @@ use std::ops::{Deref, DerefMut}; use diesel::{OptionalExtension, QueryDsl, RunQueryDsl, SqliteConnection}; use log::{debug, info, warn}; use serde::Serialize; -use tari_engine_types::events::Event; +use tari_engine_types::transaction_receipt::{TransactionReceipt, TransactionReceiptAddress}; use tari_ootle_common_types::{shard::Shard, substate_type::SubstateType, Epoch, StateVersion}; use tari_ootle_storage::{ consensus_models::{EpochCheckpoint, SubstateData, SubstateUpdateProof}, @@ -17,6 +17,7 @@ use tari_transaction::Transaction; use crate::{ diesel::ExpressionMethods, + network_state_sync::EventFilter, storage_sqlite::{ models::{ NewEvent, @@ -29,8 +30,8 @@ use crate::{ }, reader::SqliteStoreReadTransaction, serialization::{serialize_bincode, serialize_hex, serialize_json}, - IndexerStoreWriteTransaction, }, + store::IndexerStoreWriteTransaction, }; const LOG_TARGET: &str = "tari::indexer::storage_sqlite::writer"; @@ -236,29 +237,51 @@ impl IndexerStoreWriteTransaction for SqliteStoreWriteTransaction<'_> { Ok(()) } - fn batch_insert_events>(&mut self, events: I) -> Result<(), StorageError> { - const OPERATION: &str = "batch_insert_events"; - use crate::storage_sqlite::schema::events; - - let events = events.into_iter().map(|event| { - Ok::<_, StorageError>(NewEvent { - template_address: event.template_address().to_string(), - tx_hash: event.tx_hash().to_string(), - topic: event.topic(), - payload: serialize_json(event.payload())?, - substate_id: event.substate_id().map(|s| s.to_string()), - }) - }); - - for result in events { - let event = result?; - - diesel::insert_into(events::table) - .values(event) + fn batch_insert_transaction_receipts>( + &mut self, + receipts: I, + event_filters: &[EventFilter], + ) -> Result<(), StorageError> { + const OPERATION: &str = "batch_insert_transaction_receipts"; + use crate::storage_sqlite::schema::{events, transaction_receipts}; + + for (receipt_addr, receipt) in receipts { + let receipt_addr_hex = serialize_hex(receipt_addr.as_object_key()); + + diesel::insert_into(transaction_receipts::table) + .values(( + transaction_receipts::address.eq(&receipt_addr_hex), + transaction_receipts::data.eq(serialize_json(&receipt)?), + )) .execute(self.connection()) - .map_err(|e| StorageError::QueryError { - reason: format!("{OPERATION}: {}", e), - })?; + .map_err(|e| StorageError::general(OPERATION, e))?; + + // Insert events + let events = receipt.events.iter() + // only keep the events specified by the indexer filter + .filter(|event| { + event_filters.is_empty() || event_filters.iter().any(|filter| filter.matches(event)) + }) + .map(|event| { + Ok::<_, StorageError>(NewEvent { + template_address: event.template_address().to_string(), + tx_hash: &receipt_addr_hex, + topic: event.topic(), + payload: serialize_json(event.payload())?, + substate_id: event.substate_id().map(|s| s.to_string()), + }) + }); + + for result in events { + let event = result?; + + diesel::insert_into(events::table) + .values(event) + .execute(self.connection()) + .map_err(|e| StorageError::QueryError { + reason: format!("{OPERATION}: {}", e), + })?; + } } Ok(()) diff --git a/applications/tari_indexer/src/store.rs b/applications/tari_indexer/src/store.rs new file mode 100644 index 0000000000..a2885855d5 --- /dev/null +++ b/applications/tari_indexer/src/store.rs @@ -0,0 +1,211 @@ +// Copyright 2025 The Tari Project +// SPDX-License-Identifier: BSD-3-Clause + +use std::ops::{Deref, DerefMut}; + +use serde::{de::DeserializeOwned, Serialize}; +use tari_consensus_types::BlockId; +use tari_engine_types::{ + events::Event, + substate::SubstateId, + transaction_receipt::{TransactionReceipt, TransactionReceiptAddress}, + Utxo, +}; +use tari_indexer_client::types::{ListSubstateItem, NonFungibleSubstate, TransactionEntry}; +use tari_ootle_common_types::{shard::Shard, substate_type::SubstateType, Epoch, ShardGroup, StateVersion}; +use tari_ootle_storage::{ + consensus_models::{EpochCheckpoint, SubstateData, SubstateUpdateProof}, + Ordering, + StorageError, +}; +use tari_ootle_wallet_sdk::models::WalletUtxoUpdate; +use tari_template_lib::{ + models::{ResourceAddress, UtxoId}, + prelude::RistrettoPublicKeyBytes, + types::{crypto::UtxoTag, TemplateAddress}, +}; +use tari_transaction::{Transaction, TransactionId}; + +use crate::{ + network_state_sync::EventFilter, + storage_sqlite::models::{KeyValue, NewScannedBlockId, SubstateRecord, UtxoUpdateRecord}, + substate_manager::SubstateResponse, +}; + +const LOG_TARGET: &str = "tari::indexer::store"; + +pub trait IndexerStore: IndexerStoreReader { + type WriteTransaction<'a>: IndexerStoreWriteTransaction + Deref> + DerefMut + where Self: 'a; + + fn create_write_tx(&self) -> Result, StorageError>; + + fn with_write_tx) -> Result, R, E>(&self, f: F) -> Result + where E: From { + let mut tx = self.create_write_tx()?; + match f(&mut tx) { + Ok(r) => { + tx.commit()?; + Ok(r) + }, + Err(e) => { + if let Err(err) = tx.rollback() { + log::error!(target: LOG_TARGET, "Failed to rollback transaction: {}", err); + } + Err(e) + }, + } + } +} + +pub trait IndexerStoreReader { + type ReadTransaction<'a>: IndexerStoreReadTransaction + where Self: 'a; + + fn create_read_tx(&self) -> Result, StorageError>; + + fn with_read_tx) -> Result, R, E>(&self, f: F) -> Result + where E: From { + let mut tx = self.create_read_tx()?; + let ret = f(&mut tx)?; + Ok(ret) + } +} + +pub trait IndexerStoreReadTransaction { + fn list_substates( + &mut self, + filter_by_type: Option, + filter_by_template: Option, + limit: Option, + offset: Option, + ) -> Result, StorageError>; + fn get_substate( + &mut self, + address: &SubstateId, + version: Option, + ) -> Result, StorageError>; + + fn get_substates(&mut self, ids: &[SubstateId]) -> Result, StorageError>; + fn get_non_fungible_count(&mut self, resource_address: String) -> Result; + fn get_non_fungibles_by_resource_address( + &mut self, + resource_address: ResourceAddress, + limit: usize, + offset: usize, + ) -> Result, StorageError>; + + fn get_events( + &mut self, + substate_id_filter: Option<&SubstateId>, + topic_filter: Option<&str>, + offset: u32, + limit: u32, + ) -> Result, StorageError>; + fn get_oldest_scanned_epoch(&mut self) -> Result, StorageError>; + fn get_last_scanned_block_id( + &mut self, + epoch: Epoch, + shard_group: ShardGroup, + ) -> Result, StorageError>; + fn list_recent_transactions( + &mut self, + last_transaction_id: Option, + limit: usize, + ) -> Result, StorageError>; + + // -------------------------------- Transaction Receipts -------------------------------- // + fn list_transaction_receipts( + &mut self, + last_id: Option, + limit: u64, + ordering: Ordering, + ) -> Result, StorageError>; + + fn get_transaction_receipt( + &mut self, + address: &TransactionReceiptAddress, + ) -> Result; + + // -------------------------------- KeyValues -------------------------------- // + fn key_value_get_value, T: DeserializeOwned>(&mut self, key: K) -> Result; + fn key_value_get_raw>(&mut self, key: K) -> Result, StorageError>; + + fn utxos_get_max_state_version( + &mut self, + resource_address: ResourceAddress, + shard: Shard, + ) -> Result; + + /// Get UTXO updates for a given resource address and shard, starting from a specific state version. + /// + /// Returns a tuple containing the maximum returned state version, and a vector of UTXO + /// updates. + fn utxos_get_updates( + &mut self, + resource_address: ResourceAddress, + shard: Shard, + from_state_version: StateVersion, + unspents_only: bool, + limit: u32, + ) -> Result<(StateVersion, Vec), StorageError>; + + fn utxos_get_unspent_by_public_nonce_and_tag( + &mut self, + resource_address: &ResourceAddress, + public_nonce_and_tag: &[(UtxoTag, RistrettoPublicKeyBytes)], + ) -> Result, StorageError>; +} + +pub trait IndexerStoreWriteTransaction { + fn commit(self) -> Result<(), StorageError>; + fn rollback(self) -> Result<(), StorageError>; + fn key_value_set, V: Serialize>(&mut self, key: K, value: V) -> Result<(), StorageError>; + fn batch_insert_substate_transitions>( + &mut self, + shard: Shard, + state_version: StateVersion, + updates: I, + ) -> Result<(), StorageError>; + fn batch_insert_utxo_updates>( + &mut self, + updates: I, + ) -> Result<(), StorageError>; + fn upsert_substate(&mut self, substate: &SubstateData) -> Result<(), StorageError>; + fn batch_insert_transaction_receipts>( + &mut self, + receipts: I, + event_filters: &[EventFilter], + ) -> Result<(), StorageError>; + fn save_scanned_block_id(&mut self, new_scanned_block_id: NewScannedBlockId) -> Result<(), StorageError>; + fn delete_scanned_epochs_older_than(&mut self, epoch: Epoch) -> Result<(), StorageError>; + fn insert_or_ignore_transaction(&mut self, transaction: &Transaction) -> Result<(), StorageError>; + fn insert_or_ignore_epoch_checkpoint(&mut self, epoch_checkpoint: &EpochCheckpoint) -> Result<(), StorageError>; +} + +pub struct ReadOnlyStore { + inner: T, +} + +impl ReadOnlyStore { + pub fn new(inner: T) -> Self { + Self { inner } + } + + pub fn list_transaction_receipts( + &self, + last_id: Option, + limit: u64, + ordering: Ordering, + ) -> Result, StorageError> { + self.inner + .with_read_tx(|tx| tx.list_transaction_receipts(last_id, limit, ordering)) + } + + pub fn get_transaction_receipt( + &self, + address: &TransactionReceiptAddress, + ) -> Result { + self.inner.with_read_tx(|tx| tx.get_transaction_receipt(address)) + } +} diff --git a/applications/tari_indexer/src/substate_manager.rs b/applications/tari_indexer/src/substate_manager.rs index 303c873fbc..154133bbe5 100644 --- a/applications/tari_indexer/src/substate_manager.rs +++ b/applications/tari_indexer/src/substate_manager.rs @@ -50,7 +50,8 @@ use tari_validator_node_rpc::client::{SubstateResult, TariValidatorNodeRpcClient use crate::{ network_state_sync::SyncProgress, - storage_sqlite::{models::Key, IndexerStore, IndexerStoreReadTransaction, SqliteIndexerStore}, + storage_sqlite::{models::Key, SqliteIndexerStore}, + store::{IndexerStoreReadTransaction, IndexerStoreReader}, substate_file_cache::SubstateFileCache, }; diff --git a/applications/tari_indexer/src/transaction_manager/mod.rs b/applications/tari_indexer/src/transaction_manager/mod.rs index a50861c1c7..e25c0efa92 100644 --- a/applications/tari_indexer/src/transaction_manager/mod.rs +++ b/applications/tari_indexer/src/transaction_manager/mod.rs @@ -40,7 +40,7 @@ use tari_validator_node_rpc::client::{ use crate::{ network_client::TariNetworkClient, - storage_sqlite::{IndexerStore, IndexerStoreReadTransaction, IndexerStoreWriteTransaction}, + store::{IndexerStore, IndexerStoreReadTransaction, IndexerStoreWriteTransaction}, transaction_manager::error::TransactionManagerError, }; diff --git a/applications/tari_indexer/web_ui/src/routes/Transaction/components/Events.tsx b/applications/tari_indexer/web_ui/src/routes/Transaction/components/Events.tsx index f17f572123..818d7dd3a4 100644 --- a/applications/tari_indexer/web_ui/src/routes/Transaction/components/Events.tsx +++ b/applications/tari_indexer/web_ui/src/routes/Transaction/components/Events.tsx @@ -47,11 +47,11 @@ interface EventsProps { } function Events({ - events, - expandAllTrigger = 0, - collapseAllTrigger = 0, - onExpandedChange, -}: EventsProps) { + events, + expandAllTrigger = 0, + collapseAllTrigger = 0, + onExpandedChange, + }: EventsProps) { const [expanded, setExpanded] = useState(false); if (!events || events.length === 0) { @@ -116,10 +116,6 @@ function Events({ Template Address {event.template_address} - - Transaction Hash - {event.tx_hash} - Payload diff --git a/applications/tari_validator_node/web_ui/src/routes/Transactions/Events.tsx b/applications/tari_validator_node/web_ui/src/routes/Transactions/Events.tsx index 3c0bfb1311..8534dd1a98 100644 --- a/applications/tari_validator_node/web_ui/src/routes/Transactions/Events.tsx +++ b/applications/tari_validator_node/web_ui/src/routes/Transactions/Events.tsx @@ -30,11 +30,11 @@ import KeyboardArrowUpIcon from "@mui/icons-material/KeyboardArrowUp"; import CodeBlockDialog from "../../Components/CodeBlock"; import { Event, shortenString, shortenSubstateId, substateIdToString } from "@tari-project/typescript-bindings"; -function RowData({ substate_id, template_address, topic, tx_hash, payload }: Event, index: number) { +function RowData({ substate_id, template_address, topic, payload }: Event) { const [open, setOpen] = useState(false); return ( <> - + - - {shortenString(tx_hash)} - - @@ -84,17 +80,15 @@ export default function Events({ data }: { data: Event[] }) { Topic Substate Id Template Address - Transaction Hash - {data.map(({ substate_id, template_address, topic, tx_hash, payload }: Event, index: number) => { + {data.map(({ substate_id, template_address, topic, payload }: Event, index: number) => { return ( diff --git a/applications/tari_walletd/src/main.rs b/applications/tari_walletd/src/main.rs index 73bca6b3fb..9c6f5019cd 100644 --- a/applications/tari_walletd/src/main.rs +++ b/applications/tari_walletd/src/main.rs @@ -41,6 +41,7 @@ use tari_shutdown::Shutdown; const LOG_TARGET: &str = "tari::wallet_daemon"; +#[allow(clippy::too_many_lines)] #[tokio::main] async fn main() -> Result<(), anyhow::Error> { // Set up a panic hook which prints the default rust panic message but also exits the process. This makes a panic in @@ -62,6 +63,10 @@ async fn main() -> Result<(), anyhow::Error> { if let Some(password) = cli.override_keyring_password.take() { config.ootle_wallet_daemon.override_keyring_password = Some(password); } + if let Some(ref url) = cli.indexer_api_url { + // TODO: not sure why the normal load_configuration override doesnt work + config.ootle_wallet_daemon.indexer_api_url = url.clone(); + } match &cli.command { Some(Subcommand::Run) | None => run(cli, config).await?, diff --git a/applications/tari_walletd/web_ui/src/routes/Transactions/Events.tsx b/applications/tari_walletd/web_ui/src/routes/Transactions/Events.tsx index a104bc1e07..ec05080e14 100644 --- a/applications/tari_walletd/web_ui/src/routes/Transactions/Events.tsx +++ b/applications/tari_walletd/web_ui/src/routes/Transactions/Events.tsx @@ -21,7 +21,18 @@ // USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. import { useState } from "react"; -import { TableContainer, Table, TableHead, TableRow, TableCell, TableBody, Collapse, Box, Chip, Typography } from "@mui/material"; +import { + TableContainer, + Table, + TableHead, + TableRow, + TableCell, + TableBody, + Collapse, + Box, + Chip, + Typography, +} from "@mui/material"; import { DataTableCell, AccordionIconButton } from "@components/StyledComponents"; import KeyboardArrowDownIcon from "@mui/icons-material/KeyboardArrowDown"; import KeyboardArrowUpIcon from "@mui/icons-material/KeyboardArrowUp"; @@ -34,42 +45,52 @@ function renderPayloadField(key: string, value: any) { if (key === "amount" && typeof value === "string") { return ( - Amount: + + Amount: + ); } - + if (key === "resource_type" && typeof value === "string") { return ( - Resource Type: + + Resource Type: + ); } - + if (key === "resource" || key === "resource_address") { return ( - Resource: + + Resource: + ); } - + if (key === "module_name" && typeof value === "string") { return ( - Module: + + Module: + ); } - + return ( - {key}: + + {key}: + {String(value)} ); @@ -79,28 +100,29 @@ function renderPayload(payload: any) { if (!payload || typeof payload !== "object") { return {JSON.stringify(payload)}; } - + return ( {Object.entries(payload).map(([key, value]) => ( - - {renderPayloadField(key, value)} - + {renderPayloadField(key, value)} ))} ); } -function RowData({ substate_id, template_address, topic, tx_hash, payload }: Event, index: number) { +function RowData({ substate_id, template_address, topic, payload }: Event, index: number) { const [open, setOpen] = useState(false); const theme = useTheme(); return ( <> {topic} - {substate_id ? : "--"} - {template_address ? : "--"} - {tx_hash ? : "--"} + + {substate_id ? : "--"} + + + {template_address ? : "--"} + - Payload Details + + Payload Details + {renderPayload(payload)} @@ -151,13 +175,12 @@ export default function Events({ data }: { data: Event[] }) { - {data.map(({ substate_id, template_address, topic, tx_hash, payload }: Event, index: number) => { + {data.map(({ substate_id, template_address, topic, payload }: Event, index: number) => { return ( diff --git a/applications/tari_walletd/web_ui/src/routes/Wallet/Components/Accounts.tsx b/applications/tari_walletd/web_ui/src/routes/Wallet/Components/Accounts.tsx index a3b063a3d9..b26dee02e1 100644 --- a/applications/tari_walletd/web_ui/src/routes/Wallet/Components/Accounts.tsx +++ b/applications/tari_walletd/web_ui/src/routes/Wallet/Components/Accounts.tsx @@ -41,13 +41,13 @@ import queryClient from "@api/queryClient"; import { AccountInfo, substateIdToString, shortenSubstateId } from "@tari-project/typescript-bindings"; import CopyAddress from "@components/CopyAddress"; -function Account(account: AccountInfo, index: number) { +function Account(account: AccountInfo) { const { account: { name, component_address }, address, } = account; return ( - + }; diff --git a/bindings/src/types/Event.ts b/bindings/src/types/Event.ts index 4258aabb9e..130dd0fd8a 100644 --- a/bindings/src/types/Event.ts +++ b/bindings/src/types/Event.ts @@ -3,10 +3,4 @@ import type { Hash } from "./Hash"; import type { Metadata } from "./Metadata"; import type { SubstateId } from "./SubstateId"; -export type Event = { - substate_id: SubstateId | null; - template_address: Hash; - tx_hash: Hash; - topic: string; - payload: Metadata; -}; +export type Event = { substate_id: SubstateId | null; template_address: Hash; topic: string; payload: Metadata }; diff --git a/bindings/src/types/Hash64.ts b/bindings/src/types/Hash64.ts index 7caeb25d71..ffadd7872f 100644 --- a/bindings/src/types/Hash64.ts +++ b/bindings/src/types/Hash64.ts @@ -1,6 +1,6 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. /** - * Representation of a 32-byte hash value + * Representation of a 64-byte hash value */ export type Hash64 = string; diff --git a/bindings/src/types/TransactionReceipt.ts b/bindings/src/types/TransactionReceipt.ts index f3418ecbc8..cf14c8fa8f 100644 --- a/bindings/src/types/TransactionReceipt.ts +++ b/bindings/src/types/TransactionReceipt.ts @@ -1,11 +1,13 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DiffSummary } from "./DiffSummary"; import type { Event } from "./Event"; import type { FeeReceipt } from "./FeeReceipt"; -import type { Hash } from "./Hash"; import type { LogEntry } from "./LogEntry"; +import type { ValidatorFeeWithdrawal } from "./ValidatorFeeWithdrawal"; export type TransactionReceipt = { - transaction_hash: Hash; + diff_summary: DiffSummary; + fee_withdrawals: Array; events: Array; logs: Array; fee_receipt: FeeReceipt; diff --git a/bindings/src/types/UpSubstate.ts b/bindings/src/types/UpSubstate.ts new file mode 100644 index 0000000000..26ae9366c9 --- /dev/null +++ b/bindings/src/types/UpSubstate.ts @@ -0,0 +1,4 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { SubstateId } from "./SubstateId"; + +export type UpSubstate = { substate_id: SubstateId; version: number; value_hash: string }; diff --git a/bindings/src/types/tari-indexer-client/GetTransactionReceiptResponse.ts b/bindings/src/types/tari-indexer-client/GetTransactionReceiptResponse.ts new file mode 100644 index 0000000000..613d06bd16 --- /dev/null +++ b/bindings/src/types/tari-indexer-client/GetTransactionReceiptResponse.ts @@ -0,0 +1,4 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { TransactionReceipt } from "../TransactionReceipt"; + +export type GetTransactionReceiptResponse = { receipt: TransactionReceipt }; diff --git a/bindings/src/types/tari-indexer-client/ListTransactionReceiptsRequest.ts b/bindings/src/types/tari-indexer-client/ListTransactionReceiptsRequest.ts new file mode 100644 index 0000000000..a6e1e3f84b --- /dev/null +++ b/bindings/src/types/tari-indexer-client/ListTransactionReceiptsRequest.ts @@ -0,0 +1,9 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { Ordering } from "../Ordering"; +import type { TransactionReceiptAddress } from "../TransactionReceiptAddress"; + +export type ListTransactionReceiptsRequest = { + limit?: number | null; + last_id?: TransactionReceiptAddress | null; + ordering: Ordering; +}; diff --git a/bindings/src/types/tari-indexer-client/ListTransactionReceiptsResponse.ts b/bindings/src/types/tari-indexer-client/ListTransactionReceiptsResponse.ts new file mode 100644 index 0000000000..8f02d92932 --- /dev/null +++ b/bindings/src/types/tari-indexer-client/ListTransactionReceiptsResponse.ts @@ -0,0 +1,5 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { TransactionReceipt } from "../TransactionReceipt"; +import type { TransactionReceiptAddress } from "../TransactionReceiptAddress"; + +export type ListTransactionReceiptsResponse = { receipts: Array<[TransactionReceiptAddress, TransactionReceipt]> }; diff --git a/clients/javascript/indexer_client/package.json b/clients/javascript/indexer_client/package.json index b700edb255..04a31b4d96 100644 --- a/clients/javascript/indexer_client/package.json +++ b/clients/javascript/indexer_client/package.json @@ -1,6 +1,6 @@ { "name": "@tari-project/indexer-client", - "version": "1.0.0", + "version": "1.0.1", "description": "Tari indexer REST API client library", "homepage": "https://github.com/tari-project/tari-ootle#readme", "bugs": { diff --git a/clients/javascript/indexer_client/src/index.ts b/clients/javascript/indexer_client/src/index.ts index 9ee10ebb60..a09993ce96 100644 --- a/clients/javascript/indexer_client/src/index.ts +++ b/clients/javascript/indexer_client/src/index.ts @@ -9,7 +9,7 @@ import type { GetRecentTransactionsRequest, GetRecentTransactionsResponse, GetSubstatesRequest, - GetSubstatesResponse, + GetSubstatesResponse, GetTransactionReceiptResponse, IndexerAddPeerRequest, IndexerAddPeerResponse, IndexerGetConnectionsResponse, @@ -19,7 +19,7 @@ import type { IndexerGetTransactionResultResponse, IndexerReadyResponse, ListRecentTransactionsRequest, ListRecentTransactionsResponse, ListSubstatesRequest, - ListSubstatesResponse, ListTemplatesResponse, + ListSubstatesResponse, ListTemplatesResponse, ListTransactionReceiptsRequest, ListTransactionReceiptsResponse, rejectReasonToString, stringToSubstateId, SubstateId, @@ -27,7 +27,7 @@ import type { TemplatesGetResponse, TemplatesListAuthoredRequest, TemplatesListAuthoredResponse, - TransactionId, + TransactionId, TransactionReceiptAddress, TransactionSubmitRequest, TransactionSubmitResponse, } from "@tari-project/typescript-bindings"; @@ -108,6 +108,14 @@ export class IndexerClient { return this.transport.sendGet(`transactions/recent`, params); } + public listTransactionReceipts(params: ListTransactionReceiptsRequest): Promise { + return this.transport.sendGet(`transaction-receipts`, params); + } + + public getTransactionReceipt(address: TransactionReceiptAddress): Promise { + return this.transport.sendGet(`transaction-receipts/${address}`, {}); + } + public templatesGet(template_address: string): Promise { return this.transport.sendGet(`templates/${encodeURIComponent(template_address)}`, {}); } diff --git a/clients/javascript/indexer_client/src/transports/fetch.ts b/clients/javascript/indexer_client/src/transports/fetch.ts index f88429df50..0a602288e0 100644 --- a/clients/javascript/indexer_client/src/transports/fetch.ts +++ b/clients/javascript/indexer_client/src/transports/fetch.ts @@ -67,6 +67,11 @@ export default class FetchTransport implements HttpTransport { query = urlParams; } if (typeof request.body === "object" && request.body !== null) { + // Add content-type header + if (!request.headers) { + request.headers = {}; + } + request.headers["Content-Type"] = "application/json"; request.body = JSON.stringify(request.body); } diff --git a/clients/javascript/wallet_daemon_client/package.json b/clients/javascript/wallet_daemon_client/package.json index 6be68d6bea..32583197fb 100644 --- a/clients/javascript/wallet_daemon_client/package.json +++ b/clients/javascript/wallet_daemon_client/package.json @@ -1,6 +1,6 @@ { "name": "@tari-project/wallet_jrpc_client", - "version": "1.9.2", + "version": "1.9.3", "description": "Tari wallet JSON-RPC client library", "homepage": "https://github.com/tari-project/tari-ootle#readme", "bugs": { diff --git a/clients/tari_indexer_client/src/rest_api_client.rs b/clients/tari_indexer_client/src/rest_api_client.rs index 78dfb30684..bdb0e4f199 100644 --- a/clients/tari_indexer_client/src/rest_api_client.rs +++ b/clients/tari_indexer_client/src/rest_api_client.rs @@ -3,7 +3,7 @@ use reqwest::{header, header::HeaderMap, IntoUrl, Url}; use serde::{de::DeserializeOwned, Serialize}; -use tari_engine_types::substate::SubstateId; +use tari_engine_types::{substate::SubstateId, transaction_receipt::TransactionReceiptAddress}; use crate::{ error::IndexerRestClientError, @@ -23,6 +23,7 @@ use crate::{ GetSubstatesResponse, GetTemplateDefinitionRequest, GetTemplateDefinitionResponse, + GetTransactionReceiptResponse, GetTransactionResultRequest, GetTransactionResultResponse, GetUnspentUtxosRequest, @@ -33,6 +34,8 @@ use crate::{ ListRecentTransactionsResponse, ListSubstatesRequest, ListSubstatesResponse, + ListTransactionReceiptsRequest, + ListTransactionReceiptsResponse, SubmitTransactionRequest, SubmitTransactionResponse, }, @@ -169,6 +172,20 @@ impl IndexerRestApiClient { self.send_post("utxos/fetch", req).await } + pub async fn list_transaction_receipts( + &mut self, + req: ListTransactionReceiptsRequest, + ) -> Result { + self.send_get("transaction-receipts", req).await + } + + pub async fn get_transaction_receipt( + &mut self, + address: TransactionReceiptAddress, + ) -> Result { + self.send_get(format!("transaction-receipts/{}", address), ()).await + } + pub async fn get_network_sync_state(&mut self) -> Result { self.send_get("network/stats", ()).await } diff --git a/clients/tari_indexer_client/src/types.rs b/clients/tari_indexer_client/src/types.rs index 2945a69522..a0f105d25b 100644 --- a/clients/tari_indexer_client/src/types.rs +++ b/clients/tari_indexer_client/src/types.rs @@ -12,6 +12,7 @@ use tari_engine_types::{ commit_result::ExecuteResult, substate::{Substate, SubstateId, SubstateValue}, template_lib_models::{NonFungibleAddress, ResourceAddress, UtxoId}, + transaction_receipt::{TransactionReceipt, TransactionReceiptAddress}, Utxo, }; use tari_ootle_common_types::{ @@ -22,7 +23,7 @@ use tari_ootle_common_types::{ ShardGroup, StateVersion, }; -use tari_ootle_storage::time::PrimitiveDateTime; +use tari_ootle_storage::{time::PrimitiveDateTime, Ordering}; use tari_ootle_wallet_sdk::models::UtxoUpdateSet; use tari_template_abi::TemplateDef; use tari_template_lib_types::{ @@ -421,3 +422,26 @@ pub struct SyncProgress { pub checkpoint_progress: Vec<(ShardGroup, Epoch)>, pub last_state_versions: Vec<(Shard, (StateVersion, Epoch))>, } + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export, export_to = "tari-indexer-client/"))] +pub struct ListTransactionReceiptsRequest { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub limit: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub last_id: Option, + #[serde(default)] + pub ordering: Ordering, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export, export_to = "tari-indexer-client/"))] +pub struct ListTransactionReceiptsResponse { + pub receipts: Vec<(TransactionReceiptAddress, TransactionReceipt)>, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export, export_to = "tari-indexer-client/"))] +pub struct GetTransactionReceiptResponse { + pub receipt: TransactionReceipt, +} diff --git a/crates/consensus_tests/src/support/transaction.rs b/crates/consensus_tests/src/support/transaction.rs index 216bf37b7b..8c2b92fa15 100644 --- a/crates/consensus_tests/src/support/transaction.rs +++ b/crates/consensus_tests/src/support/transaction.rs @@ -26,6 +26,7 @@ pub fn build_transaction_from(tx: Transaction) -> TransactionRecord { TransactionRecord::new(tx) } +#[allow(clippy::too_many_lines)] pub fn create_execution_result_for_transaction( transaction: &Transaction, decision: Decision, @@ -105,9 +106,10 @@ pub fn create_execution_result_for_transaction( diff.up( SubstateId::TransactionReceipt(TransactionReceiptAddress::from(transaction.calculate_id())), Substate::new(0, TransactionReceipt { - transaction_hash: transaction.calculate_id().into_array().into(), - events: vec![], - logs: vec![], + diff_summary: Default::default(), + fee_withdrawals: Default::default(), + events: Default::default(), + logs: Default::default(), fee_receipt: FeeReceipt { total_fee_payment: fee, total_fees_paid: fee, diff --git a/crates/engine/src/runtime/impl.rs b/crates/engine/src/runtime/impl.rs index bbb8ad2aa2..30bdad6a72 100644 --- a/crates/engine/src/runtime/impl.rs +++ b/crates/engine/src/runtime/impl.rs @@ -248,28 +248,18 @@ impl> RuntimeInte }) } - fn emit_std_event( + fn emit_std_event>( &self, object_name: &str, action: &str, - substate_id: SubstateId, + substate_id: T, payload: Metadata, - state: &mut WorkingState, + state_mut: &mut WorkingState, ) -> Result<(), RuntimeError> { - let tx_hash = self.entity_id_provider.transaction_hash(); - let (&template_address, _) = state.current_template()?; - - let event = Event::std( - Some(substate_id), - template_address, - tx_hash, - object_name, - action, - payload, - ); + let (&template_address, _) = state_mut.current_template()?; + let event = Event::std(Some(substate_id.into()), template_address, object_name, action, payload); debug!(target: LOG_TARGET, "Emitted event {}", event); - state.push_event(event)?; - + state_mut.push_event(event)?; Ok(()) } @@ -457,13 +447,12 @@ impl> RuntimeInte ) })?; let substate_id = component_address_option.map(SubstateId::Component); - let tx_hash = self.entity_id_provider.transaction_hash(); let template_address = self.tracker.get_template_address()?; let module = self.tracker.get_template_module_name()?; let topic = format!("{module}.{topic}"); self.tracker - .add_event(Event::custom(substate_id, template_address, tx_hash, topic, payload))?; + .add_event(Event::custom(substate_id, template_address, topic, payload))?; Ok(()) } @@ -837,7 +826,7 @@ impl> RuntimeInte payload.insert(IMAGE_URL, image_url); } - self.emit_std_event("resource", "create", resource_address.into(), payload, state_mut)?; + self.emit_std_event("resource", "create", resource_address, payload, state_mut)?; state_mut.new_substate(resource_address, resource)?; let resource_lock = state_mut.write_lock_substate(&SubstateId::Resource(resource_address))?; @@ -931,7 +920,7 @@ impl> RuntimeInte ("resource_type", resource.resource_type().to_string()), ("amount", resource.amount().to_string()), ]); - self.emit_std_event("resource", "mint", resource_address.into(), payload, state_mut)?; + self.emit_std_event("resource", "mint", resource_address, payload, state_mut)?; state_mut.new_bucket(bucket_id, resource)?; let bucket = tari_template_lib::models::Bucket::from_id(bucket_id); @@ -982,7 +971,7 @@ impl> RuntimeInte ("vault_id", arg.vault_id.to_string()), ("recall_desc", arg.resource.to_string()), ]); - self.emit_std_event("resource", "recall", resource_address.into(), payload, state_mut)?; + self.emit_std_event("resource", "recall", resource_address, payload, state_mut)?; let bucket_id = state_mut.id_provider()?.new_bucket_id(); state_mut.new_bucket(bucket_id, resource)?; @@ -1078,7 +1067,7 @@ impl> RuntimeInte self.emit_std_event( "resource", "update_nonfungible_data", - resource_address.into(), + resource_address, payload, state_mut, )?; @@ -1121,13 +1110,7 @@ impl> RuntimeInte let resource_mut = state_mut.get_resource_mut(&resource_lock)?; resource_mut.set_access_rules(access_rules); let payload = Metadata::from_iter([("resource_type", resource_mut.resource_type().to_string())]); - self.emit_std_event( - "resource", - "update_access_rules", - resource_address.into(), - payload, - state_mut, - )?; + self.emit_std_event("resource", "update_access_rules", resource_address, payload, state_mut)?; state_mut.unlock_substate(resource_lock)?; @@ -1178,8 +1161,8 @@ impl> RuntimeInte state_mut.set_vault_freeze(&vault_lock, arg.flags)?; let payload = Metadata::from_iter([("vault_id", arg.vault_id.to_string()), ("flags", arg.flags.to_string())]); - let action = if arg.flags.is_empty() { "freeze" } else { "unfreeze" }; - self.emit_std_event("resource", action, resource_address.into(), payload, state_mut)?; + let action = if arg.flags.is_empty() { "unfreeze" } else { "freeze" }; + self.emit_std_event("resource", action, resource_address, payload, state_mut)?; state_mut.unlock_substate(vault_lock)?; @@ -1478,7 +1461,7 @@ impl> RuntimeInte ("amount", bucket.amount().to_string()), ]); - self.emit_std_event("vault", "deposit", vault_id.into(), payload, state_mut)?; + self.emit_std_event("vault", "deposit", vault_id, payload, state_mut)?; let vault_mut = state_mut.get_vault_mut(&vault_lock)?; @@ -1567,7 +1550,7 @@ impl> RuntimeInte ("amount", public_amount.to_string()), ]); - self.emit_std_event("vault", "withdraw", vault_id.into(), payload, state)?; + self.emit_std_event("vault", "withdraw", vault_id, payload, state)?; let bucket_id = state.id_provider()?.new_bucket_id(); state.new_bucket(bucket_id, resource_container)?; @@ -1715,6 +1698,14 @@ impl> RuntimeInte }); } + self.emit_std_event( + "vault", + "pay_fee", + vault_id, + Metadata::from_iter([("amount", container.amount().to_string())]), + state_mut, + )?; + state_mut.pay_fee(container, Some(vault_id))?; state_mut.unlock_substate(resource_lock)?; @@ -2672,12 +2663,12 @@ impl> RuntimeInte fn publish_template(&self, template: Vec) -> Result<(), RuntimeError> { self.invoke_modules_on_runtime_call("publish_template")?; - self.tracker.write_with(|state| { + self.tracker.write_with(|state_mut| { let template_byte_size = template.len(); let binary_hash = hash_template_code(&template); let template_address = PublishedTemplateAddress::from_author_and_binary_hash(&self.seal_signer_public_key, &binary_hash); - state.new_substate( + state_mut.new_substate( template_address, SubstateValue::Template(PublishedTemplate { // We essentially store the pre-image of the template address in the substate @@ -2686,15 +2677,14 @@ impl> RuntimeInte }), )?; // Mark template substate as owned by current call stack - let scope_mut = state.current_call_scope_mut()?; + let scope_mut = state_mut.current_call_scope_mut()?; scope_mut.move_node_to_owned(&template_address.into())?; // Publish template event let mut metadata = Metadata::new(); metadata.insert("template_byte_size".to_string(), template_byte_size.to_string()); - state.push_event(Event::std( + state_mut.push_event(Event::std( Some(template_address.into()), template_address.as_hash(), - state.transaction_hash(), "template", "publish", metadata, diff --git a/crates/engine/src/runtime/tracker.rs b/crates/engine/src/runtime/tracker.rs index 866818a367..9f32e2c3ce 100644 --- a/crates/engine/src/runtime/tracker.rs +++ b/crates/engine/src/runtime/tracker.rs @@ -32,7 +32,7 @@ use tari_engine_types::{ indexed_value::{IndexedValue, IndexedWellKnownTypes}, lock::LockFlag, logs::LogEntry, - substate::{SubstateId, SubstateValue}, + substate::{Substate, SubstateId, SubstateValue}, virtual_substate::VirtualSubstates, }; use tari_ootle_common_types::Epoch; @@ -181,7 +181,6 @@ impl StateTracker { state.push_event(Event::std( Some(substate_id), template_address, - state.transaction_hash(), "component", "created", Metadata::from([("module_name".to_string(), module_name)]), @@ -261,25 +260,23 @@ impl StateTracker { remaining: state.call_frame_depth(), }); } - // Resolve the transfers to the fee pool resource and vault refunds - let transaction_receipt = state.finalize_fees(&mut substates_to_persist)?; - - let fee_receipt = transaction_receipt.fee_receipt.clone(); let fee_withdrawals = state.take_validator_fee_withdrawals(); - let result = - state.generate_substate_diff(transaction_receipt, substates_to_persist, downed_utxos, fee_withdrawals); - let result = match result { - Ok(substate_diff) => TransactionResult::Accept(substate_diff), - Err(err) => TransactionResult::Reject(err.to_reject_reason()), - }; + // Resolve the transfers to the fee pool resource and vault refunds + let fee_receipt = state.finalize_fee_receipt(&mut substates_to_persist)?; + let mut diff = state.generate_substate_diff(substates_to_persist, downed_utxos, fee_withdrawals)?; + let transaction_receipt = state.finalize_transaction_receipt(&diff, fee_receipt.clone())?; + diff.up( + SubstateId::TransactionReceipt(state.transaction_hash().into()), + Substate::new(0, transaction_receipt), + ); let finalized = FinalizeResult::new( state.transaction_hash(), state.take_logs(), state.take_events(), - result, + TransactionResult::Accept(diff), fee_receipt, ); diff --git a/crates/engine/src/runtime/working_state.rs b/crates/engine/src/runtime/working_state.rs index 7a12ce2e2c..89707dc616 100644 --- a/crates/engine/src/runtime/working_state.rs +++ b/crates/engine/src/runtime/working_state.rs @@ -244,7 +244,6 @@ impl WorkingState { self.push_event(Event::std( Some(locked.substate_id().clone()), template_address, - self.transaction_hash(), "component", "updated", tari_template_lib::models::Metadata::from([("module_name".to_string(), module_name)]), @@ -681,7 +680,6 @@ impl WorkingState { Event::std( Some(vault_lock.substate_id().clone()), template_address, - self.transaction_hash(), "vault", "unfrozen", tari_template_lib::models::Metadata::default(), @@ -690,7 +688,6 @@ impl WorkingState { Event::std( Some(vault_lock.substate_id().clone()), template_address, - self.transaction_hash(), "vault", "set_freeze", flags.iter().map(|f| (f.to_string(), "true".to_string())).collect(), @@ -995,93 +992,17 @@ impl WorkingState { self.last_instruction_output = Some(output); } - pub fn finalize_fees( + pub fn finalize_transaction_receipt( &mut self, - substates_to_persist: &mut IndexMap, + diff: &SubstateDiff, + fee_receipt: FeeReceipt, ) -> Result { - let total_fees = self.fee_state.total_charges(); - - let total_fee_payment = self.fee_state.total_payments(); - - let mut fee_resource = ResourceContainer::stealth(STEALTH_TARI_RESOURCE_ADDRESS, Amount::zero()); - - // Collect the fee - let mut remaining_fees = total_fees; - let mut total_fee_overcharge = 0; - // First collect fees that cannot be refunded (we have to take all fees even if they exceed the required amount) - for resx in self.fee_state.non_refundable_fee_payments_mut_iter() { - // PANIC: this is checked by FeeState - let paid_amount = resx.amount().to_u64_checked().expect("invalid fee entry in fee state"); - - debug!( - target: LOG_TARGET, - "Collecting {} of non-refundable fees", resx.amount() - ); - - // If there is no refund vault, we must take the entire amount to avoid destroying funds - fee_resource.deposit(resx.withdraw(resx.amount())?)?; - if remaining_fees < paid_amount { - total_fee_overcharge += paid_amount - remaining_fees; - } - - remaining_fees = remaining_fees.saturating_sub(paid_amount); - } - - if remaining_fees > 0 { - for (resx, _) in self.fee_state.refundable_fee_payments_iter_mut() { - if remaining_fees == 0 { - break; - } - - debug!( - target: LOG_TARGET, - "Collecting {} of refundable fees", resx.amount() - ); - - // PANIC: this is checked by FeeState - let paid_amount = resx.amount().to_u64_checked().expect("invalid fee entry in fee state"); - - // Withdraw only what is needed - let amount_to_withdraw = cmp::min(paid_amount, remaining_fees); - fee_resource.deposit(resx.withdraw(amount_to_withdraw.into())?)?; - remaining_fees = remaining_fees.saturating_sub(amount_to_withdraw); - } - } - - // Refund the remaining refundable payments if any - for (mut resx, refund_vault) in self.fee_state.drain_refundable_fee_payments() { - if resx.amount().is_zero() { - debug_assert!(!resx.amount().is_negative()); - continue; - } - - debug!( - target: LOG_TARGET, - "Refunding {} of fees to vault {}", resx.amount(), refund_vault - ); - let vault_mut = substates_to_persist - .get_mut(&SubstateId::Vault(refund_vault)) - .expect("invariant: vault that made fee payment not in changeset") - .as_vault_mut() - .expect("invariant: substate substate_id for fee refund is not a vault"); - vault_mut.resource_container_mut().deposit(resx.withdraw_all()?)?; - } - - let total_fees_paid = fee_resource - .amount() - .to_u64_checked() - .expect("FeeState guarantees that the total fee payments fit in an u64"); - Ok(TransactionReceipt { - transaction_hash: self.transaction_hash, - events: self.events.clone(), - logs: self.logs.clone(), - fee_receipt: FeeReceipt { - total_fee_payment, - total_fees_paid, - total_fee_overcharge, - cost_breakdown: self.fee_state.take_fee_charges(), - }, + diff_summary: diff.into(), + fee_withdrawals: diff.validator_fee_withdrawals().to_vec().into_boxed_slice(), + events: self.events.clone().into_boxed_slice(), + logs: self.logs.clone().into_boxed_slice(), + fee_receipt, }) } @@ -1396,9 +1317,93 @@ impl WorkingState { &self.logs } + pub fn finalize_fee_receipt( + &mut self, + substates_to_persist: &mut IndexMap, + ) -> Result { + let total_fees = self.fee_state.total_charges(); + + let total_fee_payment = self.fee_state.total_payments(); + + let mut fee_resource = ResourceContainer::stealth(STEALTH_TARI_RESOURCE_ADDRESS, Amount::zero()); + + // Collect the fee + let mut remaining_fees = total_fees; + let mut total_fee_overcharge = 0; + // First collect fees that cannot be refunded (we have to take all fees even if they exceed the required amount) + for resx in self.fee_state.non_refundable_fee_payments_mut_iter() { + // PANIC: this is checked by FeeState + let paid_amount = resx.amount().to_u64_checked().expect("invalid fee entry in fee state"); + + debug!( + target: LOG_TARGET, + "Collecting {} of non-refundable fees", resx.amount() + ); + + // If there is no refund vault, we must take the entire amount to avoid destroying funds + fee_resource.deposit(resx.withdraw(resx.amount())?)?; + if remaining_fees < paid_amount { + total_fee_overcharge += paid_amount - remaining_fees; + } + + remaining_fees = remaining_fees.saturating_sub(paid_amount); + } + + if remaining_fees > 0 { + for (resx, _) in self.fee_state.refundable_fee_payments_iter_mut() { + if remaining_fees == 0 { + break; + } + + debug!( + target: LOG_TARGET, + "Collecting {} of refundable fees", resx.amount() + ); + + // PANIC: this is checked by FeeState + let paid_amount = resx.amount().to_u64_checked().expect("invalid fee entry in fee state"); + + // Withdraw only what is needed + let amount_to_withdraw = cmp::min(paid_amount, remaining_fees); + fee_resource.deposit(resx.withdraw(amount_to_withdraw.into())?)?; + remaining_fees = remaining_fees.saturating_sub(amount_to_withdraw); + } + } + + // Refund the remaining refundable payments if any + for (mut resx, refund_vault) in self.fee_state.drain_refundable_fee_payments() { + if resx.amount().is_zero() { + debug_assert!(!resx.amount().is_negative()); + continue; + } + + debug!( + target: LOG_TARGET, + "Refunding {} of fees to vault {}", resx.amount(), refund_vault + ); + let vault_mut = substates_to_persist + .get_mut(&SubstateId::Vault(refund_vault)) + .expect("invariant: vault that made fee payment not in changeset") + .as_vault_mut() + .expect("invariant: substate substate_id for fee refund is not a vault"); + vault_mut.resource_container_mut().deposit(resx.withdraw_all()?)?; + } + + let total_fees_paid = fee_resource + .amount() + .to_u64_checked() + .expect("FeeState guarantees that the total fee payments fit in an u64"); + + Ok(FeeReceipt { + total_fee_payment, + total_fees_paid, + total_fee_overcharge, + cost_breakdown: self.fee_state.take_fee_charges(), + }) + } + pub fn generate_substate_diff( &self, - transaction_receipt: TransactionReceipt, substates_to_persist: IndexMap, downed_utxos: IndexSet, fee_withdrawals: Vec, @@ -1427,11 +1432,6 @@ impl WorkingState { substate_diff.down(SubstateId::Utxo(downed_utxo), spent_utxo.version()); } - substate_diff.up( - SubstateId::TransactionReceipt(transaction_receipt.transaction_hash.into()), - Substate::new(0, SubstateValue::TransactionReceipt(transaction_receipt)), - ); - Ok(substate_diff) } diff --git a/crates/engine/tests/freeze.rs b/crates/engine/tests/freeze.rs index 69502d2fe9..5512040408 100644 --- a/crates/engine/tests/freeze.rs +++ b/crates/engine/tests/freeze.rs @@ -1,4 +1,4 @@ -// Copyright 2023 The Tari Project +// Copyright 2025 The Tari Project // SPDX-License-Identifier: BSD-3-Clause use std::collections::BTreeMap; diff --git a/crates/engine_types/src/commit_result.rs b/crates/engine_types/src/commit_result.rs index 2e9ce5ec64..052089ac2a 100644 --- a/crates/engine_types/src/commit_result.rs +++ b/crates/engine_types/src/commit_result.rs @@ -25,7 +25,6 @@ use std::{ time::Duration, }; -use borsh::BorshSerialize; use serde::{de::DeserializeOwned, Deserialize, Serialize}; use tari_template_lib::types::Hash; @@ -35,6 +34,7 @@ use crate::{ instruction_result::InstructionResult, logs::LogEntry, substate::SubstateDiff, + transaction_receipt::TransactionReceipt, }; #[derive(Debug, Clone, Serialize, Deserialize)] @@ -216,6 +216,13 @@ impl FinalizeResult { pub fn is_reject(&self) -> bool { matches!(self.result, TransactionResult::Reject(_)) } + + pub fn get_transaction_receipt(&self) -> Option<&TransactionReceipt> { + self.result.any_accept().and_then(|diff| { + diff.up_iter() + .find_map(|(_, s)| s.substate_value().as_transaction_receipt()) + }) + } } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -353,7 +360,7 @@ impl Display for RejectReason { } } -#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash, Deserialize, Serialize, BorshSerialize)] +#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash, Deserialize, Serialize, borsh::BorshSerialize)] #[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export))] pub enum AbortReason { ForeignPledgeInputConflict, diff --git a/crates/engine_types/src/events.rs b/crates/engine_types/src/events.rs index 8cdca7c7ce..b666dd4bbb 100644 --- a/crates/engine_types/src/events.rs +++ b/crates/engine_types/src/events.rs @@ -23,10 +23,7 @@ use std::fmt::Display; use serde::{Deserialize, Serialize}; -use tari_template_lib::{ - models::Metadata, - types::{Hash, TemplateAddress}, -}; +use tari_template_lib::{models::Metadata, types::TemplateAddress}; use crate::substate::SubstateId; @@ -42,7 +39,6 @@ fn std_event(object_name: &str, action_name: &str) -> String { pub struct Event { substate_id: Option, template_address: TemplateAddress, - tx_hash: Hash, topic: String, payload: Metadata, } @@ -51,14 +47,12 @@ impl Event { pub fn new( substate_id: Option, template_address: TemplateAddress, - tx_hash: Hash, topic: String, payload: Metadata, ) -> Self { Self { substate_id, template_address, - tx_hash, topic, payload, } @@ -67,17 +61,15 @@ impl Event { pub fn custom( substate_id: Option, template_address: TemplateAddress, - tx_hash: Hash, topic: String, payload: Metadata, ) -> Self { - Self::new(substate_id, template_address, tx_hash, topic, payload) + Self::new(substate_id, template_address, topic, payload) } pub fn std( substate_id: Option, template_address: TemplateAddress, - tx_hash: Hash, object_name: &str, action_name: &str, payload: Metadata, @@ -85,7 +77,6 @@ impl Event { Self::new( substate_id, template_address, - tx_hash, std_event(object_name, action_name), payload, ) @@ -117,20 +108,12 @@ impl Event { self.template_address } - pub fn tx_hash(&self) -> Hash { - self.tx_hash + pub fn topic(&self) -> &str { + &self.topic } - pub fn topic(&self) -> String { - self.topic.clone() - } - - pub fn add_payload(&mut self, key: String, value: String) { - self.payload.insert(key, value); - } - - pub fn get_payload(&self, key: &str) -> Option { - self.payload.get(key).cloned() + pub fn get_payload(&self, key: &str) -> Option<&str> { + self.payload.get(key).map(|s| s.as_str()) } pub fn payload(&self) -> &Metadata { @@ -146,10 +129,9 @@ impl Display for Event { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!( f, - "event: substate_id {:?}, template_address {}, tx_hash {}, topic {} and payload {}", + "event: substate_id {:?}, template_address {}, topic {} and payload {}", self.substate_id.as_ref().map(|e| e.to_string()), self.template_address, - self.tx_hash, self.topic, self.payload ) diff --git a/crates/engine_types/src/transaction_receipt.rs b/crates/engine_types/src/transaction_receipt.rs index 45618a1b90..1d61ff3f73 100644 --- a/crates/engine_types/src/transaction_receipt.rs +++ b/crates/engine_types/src/transaction_receipt.rs @@ -9,12 +9,20 @@ use std::{ use serde::{Deserialize, Serialize}; use tari_bor::BorTag; +use tari_common_types::types::FixedHash; use tari_template_lib::{ models::{address_prefixes, BinaryTag}, types::{Hash, KeyParseError, ObjectKey}, }; -use crate::{events::Event, fees::FeeReceipt, logs::LogEntry}; +use crate::{ + events::Event, + fees::FeeReceipt, + logs::LogEntry, + serde_with, + substate::{hash_substate, SubstateDiff, SubstateId}, + ValidatorFeeWithdrawal, +}; const TAG: u64 = BinaryTag::TransactionReceipt.as_u64(); @@ -78,8 +86,39 @@ impl FromStr for TransactionReceiptAddress { #[derive(Debug, Clone, Serialize, Deserialize, borsh::BorshSerialize)] #[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export))] pub struct TransactionReceipt { - pub transaction_hash: Hash, - pub events: Vec, - pub logs: Vec, + pub diff_summary: DiffSummary, + pub fee_withdrawals: Box<[ValidatorFeeWithdrawal]>, + pub events: Box<[Event]>, + pub logs: Box<[LogEntry]>, pub fee_receipt: FeeReceipt, } +#[derive(Debug, Clone, Default, Serialize, Deserialize, borsh::BorshSerialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export))] +pub struct DiffSummary { + pub upped: Box<[UpSubstate]>, +} + +impl From<&SubstateDiff> for DiffSummary { + fn from(diff: &SubstateDiff) -> Self { + Self { + upped: diff + .up_iter() + .map(|(id, s)| UpSubstate { + substate_id: id.clone(), + version: s.version(), + value_hash: hash_substate(s.substate_value(), s.version()), + }) + .collect(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, borsh::BorshSerialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export))] +pub struct UpSubstate { + pub substate_id: SubstateId, + pub version: u32, + #[cfg_attr(feature = "ts", ts(type = "string"))] + #[serde(with = "serde_with::hex")] + pub value_hash: FixedHash, +} diff --git a/crates/engine_types/src/validator_fee.rs b/crates/engine_types/src/validator_fee.rs index 9ff1c9a4e2..11fb7df3dd 100644 --- a/crates/engine_types/src/validator_fee.rs +++ b/crates/engine_types/src/validator_fee.rs @@ -176,7 +176,7 @@ impl ValidatorFeePool { } } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize, Deserialize, borsh::BorshSerialize)] #[cfg_attr(feature = "ts", derive(ts_rs::TS), ts(export))] pub struct ValidatorFeeWithdrawal { pub address: ValidatorFeePoolAddress, diff --git a/crates/storage_sqlite/src/error.rs b/crates/storage_sqlite/src/error.rs index 5933dbb05a..978e8f06ef 100644 --- a/crates/storage_sqlite/src/error.rs +++ b/crates/storage_sqlite/src/error.rs @@ -43,7 +43,7 @@ pub enum SqliteStorageError { #[error("General diesel error during operation {operation}: {source}")] DieselError { source: diesel::result::Error, - operation: String, + operation: &'static str, }, #[error("Could not migrate the database: {source}")] MigrationError { @@ -86,6 +86,12 @@ impl From for StorageError { SqliteStorageError::ConnectionError { .. } => StorageError::ConnectionError { reason: source.to_string(), }, + SqliteStorageError::DieselError { source, operation } if source == diesel::result::Error::NotFound => { + StorageError::NotFoundDbAdapter { + operation, + source: source.into(), + } + }, SqliteStorageError::DieselError { .. } => StorageError::QueryError { reason: source.to_string(), }, diff --git a/crates/storage_sqlite/src/global/backend_adapter.rs b/crates/storage_sqlite/src/global/backend_adapter.rs index 01888be8dd..421d481ae8 100644 --- a/crates/storage_sqlite/src/global/backend_adapter.rs +++ b/crates/storage_sqlite/src/global/backend_adapter.rs @@ -111,7 +111,7 @@ impl SqliteGlobalDbAdapter { .get_result::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "exists::metadata".to_string(), + operation: "exists::metadata", })?; Ok(result > 0) } @@ -162,7 +162,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .optional() .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::metadata_key".to_string(), + operation: "get::metadata_key", })?; let v = row.map(|r| serde_json::from_slice(&r.value)).transpose()?; @@ -184,14 +184,14 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "update::metadata".to_string(), + operation: "update::metadata", })?, Ok(false) => diesel::insert_into(metadata::table) .values((metadata::key_name.eq(key), metadata::value.eq(value))) .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::metadata".to_string(), + operation: "insert::metadata", })?, Err(e) => return Err(e), }; @@ -220,7 +220,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_result::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "exists::metadata".to_string(), + operation: "exists::metadata", })?; Ok(result > 0) } @@ -238,7 +238,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "set_status".to_string(), + operation: "set_status", })?; if num_affected == 0 { return Err(SqliteStorageError::NotFound { @@ -257,7 +257,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .optional() .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get_template".to_string(), + operation: "get_template", })?; match template { @@ -293,7 +293,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get_templates".to_string(), + operation: "get_templates", })?; templates @@ -315,7 +315,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get_templates_by_addresses".to_string(), + operation: "get_templates_by_addresses", })? .into_iter() .map(|t| t.try_into().map_err(SqliteStorageError::TemplateConversion)) @@ -334,7 +334,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get_pending_template".to_string(), + operation: "get_pending_template", })?; templates @@ -375,7 +375,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert_template".to_string(), + operation: "insert_template", })?; Ok(()) @@ -404,7 +404,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "update_template".to_string(), + operation: "update_template", })?; Ok(()) @@ -435,7 +435,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::validator_nodes".to_string(), + operation: "insert::validator_nodes", })?; Ok(()) @@ -455,7 +455,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "remove::validator_nodes".to_string(), + operation: "remove::validator_nodes", })?; Ok(()) @@ -478,7 +478,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: format!("get::get_validator_nodes_within_epochs({})", epoch), + operation: "get::get_validator_nodes_within_epochs", })?; distinct_validators_sorted(sqlite_vns) @@ -499,7 +499,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: format!("get::get_validator_nodes_within_epochs({})", epoch), + operation: "get::get_validator_nodes_within_epochs", })?; sqlite_vns.into_iter().map(TryInto::try_into).collect() @@ -522,7 +522,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .first::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::validator_node".to_string(), + operation: "get::validator_node", })?; let vn = vn.try_into()?; @@ -546,7 +546,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .first::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::validator_node".to_string(), + operation: "get::validator_node", })?; let vn = vn.try_into()?; @@ -563,7 +563,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_result::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "count_validator_nodes".to_string(), + operation: "count_validator_nodes", })?; Ok(count.cnt as u64) @@ -585,7 +585,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_result::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "count_validator_nodes".to_string(), + operation: "count_validator_nodes", })?; Ok(count as u64) @@ -609,7 +609,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .first::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "validator_nodes_set_committee_bucket".to_string(), + operation: "validator_nodes_set_committee_bucket", })?; diesel::insert_into(committees::table) @@ -622,7 +622,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::committee_bucket".to_string(), + operation: "insert::committee_bucket", })?; Ok(()) } @@ -654,7 +654,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "validator_nodes_get_for_shard_group".to_string(), + operation: "validator_nodes_get_for_shard_group", })?; debug!(target: LOG_TARGET, "Found {} validators", validators.len()); @@ -696,7 +696,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .get_results::<(DbValidatorNode, DbCommittee)>(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "validator_nodes_get_overlapping_shard_group".to_string(), + operation: "validator_nodes_get_overlapping_shard_group", })?; debug!(target: LOG_TARGET, "Found {} validators", validators.len()); @@ -750,7 +750,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .first::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::validator_node".to_string(), + operation: "get::validator_node", })?; let vn = vn.try_into()?; @@ -782,7 +782,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .load::<(i32, i32, String, Vec, i64)>(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "validator_nodes_get_committees".to_string(), + operation: "validator_nodes_get_committees", })?; let mut committees = HashMap::new(); @@ -814,7 +814,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::epoch".to_string(), + operation: "insert::epoch", })?; Ok(()) @@ -829,7 +829,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .optional() .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::epoch".to_string(), + operation: "get::epoch", })?; match query_res { @@ -854,7 +854,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::bmt".to_string(), + operation: "insert::bmt", })?; Ok(()) @@ -873,7 +873,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .optional() .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::bmt".to_string(), + operation: "get::bmt", })?; match query_res { Some(bmt) => Ok(Some(serde_json::from_slice(&bmt.bmt)?)), @@ -897,7 +897,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::layer_one_transaction".to_string(), + operation: "insert::layer_one_transaction", })?; Ok(()) @@ -921,7 +921,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .execute(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "insert::block_header".to_string(), + operation: "insert::block_header", })?; Ok(()) @@ -941,7 +941,7 @@ impl GlobalDbAdapter for SqliteGlobalDbAdapter { .first::(tx.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "get::block_header_by_hash".to_string(), + operation: "get::block_header_by_hash", })?; header.try_into() diff --git a/crates/storage_sqlite/src/sqlite_db_factory.rs b/crates/storage_sqlite/src/sqlite_db_factory.rs index dc6a5a91fb..4df72b97b4 100644 --- a/crates/storage_sqlite/src/sqlite_db_factory.rs +++ b/crates/storage_sqlite/src/sqlite_db_factory.rs @@ -65,7 +65,7 @@ impl DbFactory for SqliteDbFactory { .execute(&mut connection) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "set pragma".to_string(), + operation: "set pragma", })?; Ok(GlobalDb::new(SqliteGlobalDbAdapter::new(connection))) } diff --git a/crates/storage_sqlite/src/sqlite_transaction.rs b/crates/storage_sqlite/src/sqlite_transaction.rs index 7c413a890a..f5787fc334 100644 --- a/crates/storage_sqlite/src/sqlite_transaction.rs +++ b/crates/storage_sqlite/src/sqlite_transaction.rs @@ -62,7 +62,7 @@ impl<'a> SqliteTransaction<'a> { .execute(self.connection()) .map_err(|source| SqliteStorageError::DieselError { source, - operation: "execute sql".to_string(), + operation: "execute sql", })?; Ok(()) diff --git a/crates/template_builtin/templates/account/src/lib.rs b/crates/template_builtin/templates/account/src/lib.rs index eccb063386..3e14fd9713 100644 --- a/crates/template_builtin/templates/account/src/lib.rs +++ b/crates/template_builtin/templates/account/src/lib.rs @@ -85,31 +85,19 @@ mod account_template { } pub fn withdraw(&mut self, resource: ResourceAddress, amount: Amount) -> Bucket { - emit_event("withdraw", [ - ("amount", amount.to_string()), - ("resource", resource.to_string()), - ]); + // An event is emitted by the vault.withdraw method let v = self.get_vault_mut(resource); v.withdraw(amount) } pub fn withdraw_non_fungible(&mut self, resource: ResourceAddress, nf_id: NonFungibleId) -> Bucket { - emit_event("withdraw_non_fungible", [ - ("id", nf_id.to_string()), - ("resource", resource.to_string()), - ]); + // An event is emitted by the vault.withdraw_non_fungibles method let v = self.get_vault_mut(resource); v.withdraw_non_fungibles([nf_id]) } pub fn withdraw_many_non_fungibles(&mut self, resource: ResourceAddress, nf_ids: Vec) -> Bucket { - emit_event("withdraw_many_non_fungibles", [ - ("resource", resource.to_string()), - ( - "ids", - nf_ids.iter().map(ToString::to_string).collect::>().join(","), - ), - ]); + // An event is emitted by the vault.withdraw_non_fungibles method let v = self.get_vault_mut(resource); v.withdraw_non_fungibles(nf_ids) } @@ -119,21 +107,13 @@ mod account_template { resource: ResourceAddress, withdraw_proof: ConfidentialWithdrawProof, ) -> Bucket { - emit_event("withdraw_confidential", [ - ("num_inputs", withdraw_proof.inputs.len().to_string()), - ("resource", resource.to_string()), - ]); - + // An event is emitted by the vault.withdraw_confidential method let v = self.get_vault_mut(resource); v.withdraw_confidential(withdraw_proof) } pub fn deposit(&mut self, bucket: Bucket) { - emit_event("deposit", [ - ("amount", bucket.amount().to_string()), - ("resource", bucket.resource_address().to_string()), - ("resource_type", bucket.resource_type().to_string()), - ]); + // An event is emitted by the vault.deposit method let resource_address = bucket.resource_address(); let vault_mut = self .vaults @@ -168,10 +148,7 @@ mod account_template { /// vault. It will panic if the proof is invalid or the resource type contained in the vault is not /// confidential. This is useful for converting confidential tokens into revealed tokens and vice versa. pub fn join_confidential(&mut self, resource: ResourceAddress, proof: ConfidentialWithdrawProof) { - emit_event("join_confidential", [ - ("num_inputs", proof.inputs.len().to_string()), - ("resource", resource.to_string()), - ]); + // An event is emitted by the vault.withdraw_confidential and vault.deposit methods let vault_mut = self.get_vault_mut(resource); let bucket = vault_mut.withdraw_confidential(proof); vault_mut.deposit(bucket); diff --git a/crates/template_lib/src/args/freeze_flags.rs b/crates/template_lib/src/args/freeze_flags.rs index 1b16eb6bba..c4994952cc 100644 --- a/crates/template_lib/src/args/freeze_flags.rs +++ b/crates/template_lib/src/args/freeze_flags.rs @@ -4,7 +4,7 @@ use tari_bor::{Deserialize, Serialize}; use tari_template_abi::rust::{fmt::Display, iter, ops}; -const ALL_FLAGS: u8 = 0b0011; +const ALL_FLAGS: u8 = VaultFreezeFlag::Deposits as u8 | VaultFreezeFlag::Withdrawals as u8; #[derive(Clone, Debug, Copy, Default, Serialize, Deserialize)] #[serde(transparent)] diff --git a/crates/transaction/src/transaction_id.rs b/crates/transaction/src/transaction_id.rs index 4e20a8ef57..64c43794f3 100644 --- a/crates/transaction/src/transaction_id.rs +++ b/crates/transaction/src/transaction_id.rs @@ -118,3 +118,9 @@ impl From for TransactionId { Self::new(hash.into_array()) } } + +impl From for TransactionId { + fn from(address: TransactionReceiptAddress) -> Self { + Self::new(address.as_object_key().into_array()) + } +} diff --git a/crates/wallet/sdk/src/apis/stealth_transfer.rs b/crates/wallet/sdk/src/apis/stealth_transfer.rs index 52717d8062..cd5bca3db4 100644 --- a/crates/wallet/sdk/src/apis/stealth_transfer.rs +++ b/crates/wallet/sdk/src/apis/stealth_transfer.rs @@ -181,13 +181,8 @@ where src_vault.id ); - self.outputs_api.lock_funds_in_vault(lock_id, &src_vault.id, revealed_to_spend) - .inspect_err(|_| { - // TODO: atomic rollback will help with this - if let Err(err) = self.outputs_api.release_lock(lock_id) { - error!(target: LOG_TARGET, "Failed to release lock outputs for resource {}: {}", resource_address, err); - } - })?; + self.outputs_api + .lock_funds_in_vault(lock_id, &src_vault.id, revealed_to_spend)?; return Ok(InputsToSpend { inputs: vec![], @@ -212,26 +207,15 @@ where &resource_address, lock_id, utxo_amount_to_spend, - ) - .inspect_err(|_| { - // TODO: atomic rollback will help with this - if let Err(err) = self.outputs_api.release_lock(lock_id) { - error!(target: LOG_TARGET, "Failed to release lock outputs for resource {}: {}", resource_address, err); - } - })?; + )?; let total_confidential_spent = Amount::sum_from_positive(inputs.iter().map(|i| i.value)) // The wallet has somehow stored a negative amount, which should not happen. .expect("BUG: an unblinded input amount was negative"); if let Some(ref src_vault) = maybe_src_vault { - self.outputs_api.lock_revealed_funds(lock_id, &src_vault.id, revealed_to_spend) - .inspect_err(|_| { - // TODO: atomic rollback will help with this - if let Err(err) = self.outputs_api.release_lock(lock_id) { - error!(target: LOG_TARGET, "Failed to release lock outputs for resource {}: {}", resource_address, err); - } - })?; + self.outputs_api + .lock_revealed_funds(lock_id, &src_vault.id, revealed_to_spend)?; } info!( @@ -264,29 +248,23 @@ where .unwrap_or_else(Amount::zero); if available_revealed_funds < revealed_to_spend { - self.outputs_api.release_lock(lock_id)?; return Err(StealthTransferApiError::InsufficientFunds); } if revealed_to_spend.is_positive() { - match maybe_src_vault { - Some(vault) => { - self.outputs_api - .lock_revealed_funds(lock_id, &vault.id, revealed_to_spend)?; - }, - None => { - if let Err(err) = self.outputs_api.release_lock(lock_id) { - error!(target: LOG_TARGET, "🚨 Failed to release lock outputs for resource {}: {}", resource_address, err); - } - return Err(StealthTransferApiError::InsufficientRevealedFunds { + let vault = + maybe_src_vault + .as_ref() + .ok_or_else(|| StealthTransferApiError::InsufficientRevealedFunds { details: format!( "PreferConfidential: No vault found for resource {} in account {}. Need to spend \ {} revealed funds", resource_address, owner_account_component_address, revealed_to_spend ), - }); - }, - } + })?; + + self.outputs_api + .lock_revealed_funds(lock_id, &vault.id, revealed_to_spend)?; } Ok(InputsToSpend { @@ -431,33 +409,19 @@ where let lock_id = self.outputs_api.create_lock()?; // Lock up funds for fees and transfer - let fee_inputs_to_spend = self.lock_fee_inputs(lock_id, &owner_account, ¶ms)?; + let fee_inputs_to_spend = + self.unlock_on_failure(lock_id, self.lock_fee_inputs(lock_id, &owner_account, ¶ms))?; - let inputs_to_spend = match self.lock_inputs_for_transfer( + let inputs_to_spend = self.unlock_on_failure( lock_id, - owner_account.account().component_address(), - params.resource_address, - params.total_output_amount(), - params.input_selection, - ) { - Ok(inputs) => inputs, - Err(e) => { - warn!(target: LOG_TARGET, "Unlocking fee fund locks after error: {}", e); - // This is a hack that addresses the case where input locking fails after the fee transaction. - // However, any error after this point do not undo locking. This is a limitation - // of the current design - the db transaction should be passed in and - // automatically rolled back on error. - if let Err(err) = self.outputs_api.release_lock(lock_id) { - error!( - target: LOG_TARGET, - "Failed to release fee inputs for transfer: {}", - err - ); - } - - return Err(e); - }, - }; + self.lock_inputs_for_transfer( + lock_id, + owner_account.account().component_address(), + params.resource_address, + params.total_output_amount(), + params.input_selection, + ), + )?; // TODO: use single db transaction across calls // --- Any error from here can result in funds staying locked --- @@ -481,11 +445,9 @@ where let (signing_key_branch, signing_key_id) = if must_sign_with_account_key { (KeyBranch::Account, owner_key_id) } else { - ( - KeyBranch::Nonce, - // Only usage of key manager - too bad :/ - KeyId::derived(self.key_manager_api.next_derived_key_index(KeyBranch::Nonce)?), - ) + let next_index = + self.unlock_on_failure(lock_id, self.key_manager_api.next_derived_key_index(KeyBranch::Nonce))?; + (KeyBranch::Nonce, KeyId::derived(next_index)) }; let required_signer = self .key_manager_api @@ -493,28 +455,34 @@ where let required_signer_pk = required_signer.public_key.to_byte_type(); // Generate fee transfer statement - let fee_transfer_statement = self.outputs_api.generate_transfer_statement(TransferStatementParams { - spend_key_branch: KeyBranch::Account, - spend_key_id: owner_key_id, - view_only_key_id: owner_account.view_only_key_id(), - resource_address: ¶ms.resource_address, - resource_view_key: resource_view_key.clone(), - inputs: &fee_inputs_to_spend.inputs, - input_revealed_amount: fee_inputs_to_spend.revealed, - outputs: fee_change_output.into_iter(), - output_revealed_amount: Amount::from(params.max_fee), - required_signer: required_signer_pk, - })?; + let fee_transfer_statement = self.unlock_on_failure( + lock_id, + self.outputs_api.generate_transfer_statement(TransferStatementParams { + spend_key_branch: KeyBranch::Account, + spend_key_id: owner_key_id, + view_only_key_id: owner_account.view_only_key_id(), + resource_address: ¶ms.resource_address, + resource_view_key: resource_view_key.clone(), + inputs: &fee_inputs_to_spend.inputs, + input_revealed_amount: fee_inputs_to_spend.revealed, + outputs: fee_change_output, + output_revealed_amount: Amount::from(params.max_fee), + required_signer: required_signer_pk, + }), + )?; // Add the unconfirmed fee change output to the wallet store if let Some(output) = fee_transfer_statement.outputs_statement.outputs.first() { - self.add_unconfirmed_output_from_statement( + self.unlock_on_failure( lock_id, - &owner_account, - params.resource_address, - output, - fee_stealth_change_amt, - None, + self.add_unconfirmed_output_from_statement( + lock_id, + &owner_account, + params.resource_address, + output, + fee_stealth_change_amt, + None, + ), )?; } @@ -548,8 +516,7 @@ where .total_amount() .checked_sub_positive(params.total_output_amount()) .unwrap_or_else(|| { - // This is a bug because the wallet chooses inputs based on the required outputs. This function should - // not have been called if there are insufficient funds. + // This is a bug because the wallet chooses inputs based on the required outputs. error!( target: LOG_TARGET, "BUG: total_stealth_input_amount or params.total_amount() are negative after validation" @@ -561,27 +528,30 @@ where owner_address: &owner_address, amount: change_amount, memo: None, - }) - .filter(|o| o.amount.is_positive()); + }); - let transfer_statement = self.outputs_api.generate_transfer_statement(TransferStatementParams { - spend_key_branch: KeyBranch::Account, - spend_key_id: owner_key_id, - view_only_key_id: owner_account.view_only_key_id(), - resource_address: ¶ms.resource_address, - resource_view_key, - inputs: &inputs_to_spend.inputs, - input_revealed_amount: inputs_to_spend.revealed, - outputs: Some(OutputToCreate { - amount: params.blinded_output_amount, - owner_address: &destination_address, - memo: params.output_memo.as_ref(), - }) - .into_iter() - .chain(change_output), - output_revealed_amount: params.revealed_output_amount, - required_signer: required_signer_pk, - })?; + let transfer_statement = self.unlock_on_failure( + lock_id, + self.outputs_api.generate_transfer_statement(TransferStatementParams { + spend_key_branch: KeyBranch::Account, + spend_key_id: owner_key_id, + view_only_key_id: owner_account.view_only_key_id(), + resource_address: ¶ms.resource_address, + resource_view_key, + inputs: &inputs_to_spend.inputs, + input_revealed_amount: inputs_to_spend.revealed, + outputs: Some(OutputToCreate { + amount: params.blinded_output_amount, + owner_address: &destination_address, + memo: params.output_memo.as_ref(), + }) + .into_iter() + .chain(change_output) + .filter(|o| o.amount.is_positive()), + output_revealed_amount: params.revealed_output_amount, + required_signer: required_signer_pk, + }), + )?; // Add all input UTXO substates to transaction inputs substate_inputs.extend( @@ -605,30 +575,36 @@ where .map(SubstateRequirement::unversioned), ); - let result = self.generate_transfer_transaction( - &owner_account, - params, - substate_inputs, - fee_transfer_statement, - transfer_statement, - need_to_create_dest_account, - ); + let transaction = self.unlock_on_failure( + lock_id, + self.generate_transfer_transaction( + &owner_account, + params, + substate_inputs, + fee_transfer_statement, + transfer_statement, + need_to_create_dest_account, + ), + )?; + + Ok(TransferOutput { + transaction, + lock_id, + fee_inputs: fee_inputs_to_spend, + transfer_inputs: inputs_to_spend, + signing_key_branch, + signing_key_id, + }) + } + fn unlock_on_failure(&self, lock_id: WalletLockId, result: Result) -> Result { match result { - Ok(transaction) => Ok(TransferOutput { - transaction, - lock_id, - fee_inputs: fee_inputs_to_spend, - transfer_inputs: inputs_to_spend, - signing_key_branch, - signing_key_id, - }), - Err(err) => { - // Unlock inputs - if let Err(e) = self.outputs_api.release_lock(lock_id) { - error!(target: LOG_TARGET, "Failed to release inputs lock after error: {}", e); + Ok(value) => Ok(value), + Err(e) => { + if let Err(err) = self.outputs_api.release_lock(lock_id) { + error!(target: LOG_TARGET, "Failed to release inputs lock after error: {}", err); } - Err(err) + Err(e) }, } } diff --git a/crates/wallet/sdk/src/apis/substate.rs b/crates/wallet/sdk/src/apis/substate.rs index 5e75a38b83..f6afb6d678 100644 --- a/crates/wallet/sdk/src/apis/substate.rs +++ b/crates/wallet/sdk/src/apis/substate.rs @@ -8,7 +8,6 @@ use tari_engine_types::{ indexed_value::{IndexedValueError, IndexedWellKnownTypes}, resource::Resource, substate::{Substate, SubstateId, SubstateValue}, - transaction_receipt::TransactionReceiptAddress, }; use tari_ootle_common_types::{ displayable::Displayable, @@ -141,10 +140,17 @@ where } }, SubstateValue::Resource(_) => {}, - SubstateValue::TransactionReceipt(tx_receipt) => { - let tx_receipt_addr = SubstateId::TransactionReceipt(TransactionReceiptAddress::from_hash( - tx_receipt.transaction_hash, - )); + SubstateValue::TransactionReceipt(_) => { + let addr = substate_id + .substate_id() + .as_transaction_receipt_address() + .ok_or_else(|| { + SubstateApiError::InvalidValidatorNodeResponse(format!( + "Transaction receipt substate and substate ID mismatch! Got {}", + substate_id + )) + })?; + let tx_receipt_addr = SubstateId::TransactionReceipt(addr); if substate_ids.contains(&tx_receipt_addr) { continue; } diff --git a/integration_tests/tests/features/indexer.feature b/integration_tests/tests/features/indexer.feature index 063190b462..eff18384f8 100644 --- a/integration_tests/tests/features/indexer.feature +++ b/integration_tests/tests/features/indexer.feature @@ -54,7 +54,7 @@ Feature: Indexer node Then I wait for the indexer INDEXER to sync with the network # Scan the network for the event emitted on ACC creation - When indexer INDEXER scans the network events for account ACC with topics std.component.created,Account.pay_fee,Account.deposit,std.component.updated + When indexer INDEXER scans the network events for account ACC with topics std.component.created,std.vault.pay_fee,std.component.updated Scenario: Indexer GraphQL requests work # Initialize a base node, wallet, miner and VN @@ -75,10 +75,10 @@ Feature: Indexer node Then I wait for the indexer INDEXER to sync with the network ##### Scenario # Scan the network for the event emitted on ACC_1 creation - When indexer INDEXER scans the network events for account ACC_1 with topics std.component.created,Account.pay_fee + When indexer INDEXER scans the network events for account ACC_1 with topics std.component.created,std.vault.pay_fee # Scan the network for the event emitted on ACC_2 creation - When indexer INDEXER scans the network events for account ACC_2 with topics std.component.created,Account.pay_fee + When indexer INDEXER scans the network events for account ACC_2 with topics std.component.created,std.vault.pay_fee Scenario: Indexer GraphQL filtering and pagination of events Given a network with registered validator VN and wallet daemon WALLET_D