Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
108 changes: 54 additions & 54 deletions Cargo.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
10 changes: 10 additions & 0 deletions applications/tari_app_utilities/src/seed_peer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ impl SeedPeer {
pub fn to_peer_id(&self) -> Option<PeerId> {
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())
Expand All @@ -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 {
Expand Down Expand Up @@ -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";
Expand Down
9 changes: 2 additions & 7 deletions applications/tari_indexer/src/bootstrap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
30 changes: 10 additions & 20 deletions applications/tari_indexer/src/event_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand All @@ -42,27 +43,16 @@ impl EventManager {

pub async fn get_events_from_db(
&self,
topic: Option<String>,
substate_id: Option<SubstateId>,
topic: Option<&str>,
substate_id: Option<&SubstateId>,
offset: u32,
limit: u32,
) -> Result<Vec<Event>, anyhow::Error> {
let rows = self
) -> Result<Vec<(TransactionId, Event)>, 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::<BTreeMap<String, String>>(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)
}
}
18 changes: 13 additions & 5 deletions applications/tari_indexer/src/graphql/model/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -43,12 +44,15 @@ pub struct Event {
}

impl Event {
fn from_engine_event(event: tari_engine_types::events::Event) -> Result<Self, anyhow::Error> {
fn from_engine_event(
transaction_id: TransactionId,
event: tari_engine_types::events::Event,
) -> Result<Self, anyhow::Error> {
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(),
})
}
Expand All @@ -75,14 +79,18 @@ impl EventQuery {
let substate_id = substate_id.map(|str| SubstateId::from_str(&str)).transpose()?;
let event_manager = ctx.data_unchecked::<EventManager>();
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()
}
}
1 change: 1 addition & 0 deletions applications/tari_indexer/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,11 +40,17 @@ impl ValidatorCommitteeRpcPool {

pub async fn new_session(&mut self) -> Result<ValidatorRpcSession, ValidatorCommitteeClientError> {
let epoch = self.epoch_manager.get_current_epoch();
let mut last_error = None::<ValidatorCommitteeClientError>;
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),
Expand All @@ -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);
},
}
Expand Down Expand Up @@ -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;
};
Expand Down
27 changes: 17 additions & 10 deletions applications/tari_indexer/src/network_state_sync/event_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,32 +7,39 @@ use tari_template_lib::types::{EntityId, TemplateAddress};

#[derive(Default, Debug, Serialize, Deserialize, Clone)]
pub struct EventFilter {
pub topic: Option<String>,
pub topic: Option<Box<str>>,
pub entity_id: Option<EntityId>,
pub substate_id: Option<SubstateId>,
pub template_address: Option<TemplateAddress>,
}

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
})
}
}
30 changes: 11 additions & 19 deletions applications/tari_indexer/src/network_state_sync/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -333,7 +331,7 @@ impl NetworkWideStateSync {
sync_plan_mut: &mut SyncPlan,
update_buf: &mut Vec<(Epoch, SubstateUpdateProof)>,
utxos_buf: &mut Vec<UtxoUpdateRecord>,
transactions_buf: &mut Vec<TransactionReceipt>,
transactions_buf: &mut Vec<(TransactionReceiptAddress, TransactionReceipt)>,
validator_fee_pools_buf: &mut Vec<SubstateData>,
shard_group: ShardGroup,
session: &mut ValidatorRpcSession,
Expand Down Expand Up @@ -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(..))?;
Expand Down Expand Up @@ -439,21 +438,12 @@ impl NetworkWideStateSync {
Ok(())
}

fn persist_transaction_receipts<I: IntoIterator<Item = TransactionReceipt>>(
fn persist_transaction_receipts<I: IntoIterator<Item = (TransactionReceiptAddress, TransactionReceipt)>>(
&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(())
}
}
Expand All @@ -466,7 +456,7 @@ fn extend_bufs_from_substate_update(
update_buf: &mut Vec<(Epoch, SubstateUpdateProof)>,
templates_buf: &mut Vec<TemplateChange>,
utxos_buf: &mut Vec<UtxoUpdateRecord>,
transactions_buf: &mut Vec<TransactionReceipt>,
transactions_buf: &mut Vec<(TransactionReceiptAddress, TransactionReceipt)>,
validator_fee_pools_buf: &mut Vec<SubstateData>,
) -> Result<(), NetworkStateSyncError> {
match &update {
Expand All @@ -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() {
Expand Down
7 changes: 7 additions & 0 deletions applications/tari_indexer/src/rest_api/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand All @@ -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(),
Expand Down Expand Up @@ -80,6 +82,10 @@ impl HandlerContext {
&self.inner.transaction_manager
}

pub fn read_only_store(&self) -> &ReadOnlyStore<SqliteIndexerStore> {
&self.inner.read_only_store
}

pub fn template_manager(&self) -> &TemplateManager<PeerAddress> {
&self.inner.template_manager
}
Expand Down Expand Up @@ -111,6 +117,7 @@ struct InnerContext {
substate_manager: SubstateManager,
transaction_manager:
TransactionManager<EpochManagerHandle<PeerAddress>, TariValidatorNodeRpcClientFactory, SqliteIndexerStore>,
read_only_store: ReadOnlyStore<SqliteIndexerStore>,
template_manager: TemplateManager<PeerAddress>,
dry_run_transaction_processor: DryRunTransactionProcessor,
}
Loading
Loading