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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ strum = { version = "0.26.1", features = ["derive"] }
bs58 = "0.5.0"
base64 = "0.22.1"
copypasta = "0.10.1"
dash-sdk = { git = "https://github.com/dashpay/platform", branch = "test/testWithoutSpan2" }
dash-sdk = { git = "https://github.com/dashpay/platform", rev = "c49af53698bad1070fcfc344ccaf6b3cc404e516" }
thiserror = "1"
serde = "1.0.197"
serde_json = "1.0.120"
Expand Down
262 changes: 148 additions & 114 deletions src/backend_task/contested_names/query_dpns_contested_resources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use crate::context::AppContext;
use crate::model::proof_log_item::{ProofLogItem, RequestType};
use dash_sdk::dpp::data_contract::accessors::v0::DataContractV0Getters;
use dash_sdk::dpp::data_contract::document_type::accessors::DocumentTypeV0Getters;
use dash_sdk::dpp::platform_value::Value;
use dash_sdk::drive::query::vote_polls_by_document_type_query::VotePollsByDocumentTypeQuery;
use dash_sdk::platform::FetchMany;
use dash_sdk::query_types::ContestedResource;
Expand All @@ -23,28 +24,32 @@ impl AppContext {
let Some(contested_index) = document_type.find_contested_index() else {
return Err("No contested index on dpns domains".to_string());
};
let query = VotePollsByDocumentTypeQuery {
contract_id: data_contract.id(),
document_type_name: document_type.name().to_string(),
index_name: contested_index.name.clone(),
start_at_value: None,
start_index_values: vec!["dash".into()], // hardcoded for dpns
end_index_values: vec![],
limit: None,
order_ascending: true,
};

let (contested_resources, metadata, proof) =
ContestedResource::fetch_many_with_metadata_and_proof(&sdk, query.clone(), None)
let mut start_at_value = None;
loop {
let query = VotePollsByDocumentTypeQuery {
contract_id: data_contract.id(),
document_type_name: document_type.name().to_string(),
index_name: contested_index.name.clone(),
start_at_value,
start_index_values: vec!["dash".into()], // hardcoded for dpns
end_index_values: vec![],
limit: Some(100),
order_ascending: true,
};

let (contested_resources) = ContestedResource::fetch_many(&sdk, query.clone())
.await
.map_err(|e| {
tracing::error!("error fetching contested resources: {}", e);
if let dash_sdk::Error::Proof(dash_sdk::ProofVerifierError::GroveDBError {
proof_bytes,
height,
time_ms,
error,
}) = &e
if let dash_sdk::Error::Proof(
dash_sdk::ProofVerifierError::GroveDBProofVerificationError {
proof_bytes,
path_query,
height,
time_ms,
error,
},
) = &e
{
// Encode the query using bincode
let encoded_query =
Expand All @@ -57,11 +62,23 @@ impl AppContext {
Err(e) => return e,
};

// Encode the path_query using bincode
let verification_path_query_bytes =
match bincode::encode_to_vec(&path_query, bincode::config::standard())
.map_err(|encode_err| {
tracing::error!("error encoding path_query: {}", encode_err);
format!("error encoding path_query: {}", encode_err)
}) {
Ok(encoded_path_query) => encoded_path_query,
Err(e) => return e,
};

if let Err(e) = self
.db
.insert_proof_log_item(ProofLogItem {
request_type: RequestType::GetContestedResources,
request_bytes: encoded_query,
verification_path_query_bytes,
height: *height,
time_ms: *time_ms,
proof_bytes: proof_bytes.clone(),
Expand All @@ -75,107 +92,124 @@ impl AppContext {
format!("error fetching contested resources: {}", e)
})?;

let contested_resources_as_strings: Vec<String> = contested_resources
.0
.into_iter()
.map(|contested_resource| {
contested_resource
.0
.as_str()
.expect("expected str")
.to_string()
})
.collect();

let names_to_be_updated = self
.db
.insert_name_contests_as_normalized_names(contested_resources_as_strings, &self)
.map_err(|e| e.to_string())?;

sender
.send(TaskResult::Refresh)
.await
.map_err(|e| e.to_string())?;

// Create a semaphore with 15 permits
let semaphore = Arc::new(Semaphore::new(24));

let mut handles = Vec::new();

let handle = {
let semaphore = semaphore.clone();
let sdk = sdk.clone();
let sender = sender.clone();
let self_ref = self.clone();

tokio::spawn(async move {
// Acquire a permit from the semaphore
let _permit: OwnedSemaphorePermit = semaphore.acquire_owned().await.unwrap();

match self_ref.query_dpns_ending_times(sdk, sender.clone()).await {
Ok(_) => {
// Send a refresh message if the query succeeded
sender
.send(TaskResult::Refresh)
.await
.expect("expected to send refresh");
}
Err(e) => {
tracing::error!("error querying dpns end times: {}", e);
sender
.send(TaskResult::Error(e))
.await
.expect("expected to send error");
}
}
})
};
let contested_resources_len = contested_resources.0.len();

handles.push(handle);

for name in names_to_be_updated {
// Clone the semaphore, sdk, and sender for each task
let semaphore = semaphore.clone();
let sdk = sdk.clone();
let sender = sender.clone();
let self_ref = self.clone(); // Assuming self is cloneable

// Spawn each task with a permit from the semaphore
let handle = tokio::spawn(async move {
// Acquire a permit from the semaphore
let _permit: OwnedSemaphorePermit = semaphore.acquire_owned().await.unwrap();

// Perform the query
match self_ref
.query_dpns_vote_contenders(&name, sdk, sender.clone())
.await
{
Ok(_) => {
// Send a refresh message if the query succeeded
sender
.send(TaskResult::Refresh)
.await
.expect("expected to send refresh");
}
Err(e) => {
tracing::error!("error querying dpns vote contenders for {}: {}", name, e);
sender
.send(TaskResult::Error(e))
.await
.expect("expected to send error");
if contested_resources_len == 0 {
break;
}

let contested_resources_as_strings: Vec<String> = contested_resources
.0
.into_iter()
.map(|contested_resource| {
contested_resource
.0
.as_str()
.expect("expected str")
.to_string()
})
.collect();

let last_found_name = contested_resources_as_strings.last().unwrap().clone();

let names_to_be_updated = self
.db
.insert_name_contests_as_normalized_names(contested_resources_as_strings, &self)
.map_err(|e| e.to_string())?;

sender
.send(TaskResult::Refresh)
.await
.map_err(|e| e.to_string())?;

// Create a semaphore with 15 permits
let semaphore = Arc::new(Semaphore::new(24));

let mut handles = Vec::new();

let handle = {
let semaphore = semaphore.clone();
let sdk = sdk.clone();
let sender = sender.clone();
let self_ref = self.clone();

tokio::spawn(async move {
// Acquire a permit from the semaphore
let _permit: OwnedSemaphorePermit = semaphore.acquire_owned().await.unwrap();

match self_ref.query_dpns_ending_times(sdk, sender.clone()).await {
Ok(_) => {
// Send a refresh message if the query succeeded
sender
.send(TaskResult::Refresh)
.await
.expect("expected to send refresh");
}
Err(e) => {
tracing::error!("error querying dpns end times: {}", e);
sender
.send(TaskResult::Error(e))
.await
.expect("expected to send error");
}
}
}
});
})
};

// Collect all task handles
handles.push(handle);
}

// Await all tasks
for handle in handles {
if let Err(e) = handle.await {
tracing::error!("Task failed: {:?}", e);
for name in names_to_be_updated {
// Clone the semaphore, sdk, and sender for each task
let semaphore = semaphore.clone();
let sdk = sdk.clone();
let sender = sender.clone();
let self_ref = self.clone(); // Assuming self is cloneable

// Spawn each task with a permit from the semaphore
let handle = tokio::spawn(async move {
// Acquire a permit from the semaphore
let _permit: OwnedSemaphorePermit = semaphore.acquire_owned().await.unwrap();

// Perform the query
match self_ref
.query_dpns_vote_contenders(&name, sdk, sender.clone())
.await
{
Ok(_) => {
// Send a refresh message if the query succeeded
sender
.send(TaskResult::Refresh)
.await
.expect("expected to send refresh");
}
Err(e) => {
tracing::error!(
"error querying dpns vote contenders for {}: {}",
name,
e
);
sender
.send(TaskResult::Error(e))
.await
.expect("expected to send error");
}
}
});

// Collect all task handles
handles.push(handle);
}

// Await all tasks
for handle in handles {
if let Err(e) = handle.await {
tracing::error!("Task failed: {:?}", e);
}
}
if contested_resources_len < 100 {
break;
}
start_at_value = Some((Value::Text(last_found_name), false))
}

Ok(())
Expand Down
9 changes: 5 additions & 4 deletions src/database/contested_names.rs
Original file line number Diff line number Diff line change
Expand Up @@ -699,10 +699,11 @@ impl Database {
for (name, new_ending_time) in name_contests {
// Check if the name exists in the database and retrieve the current ending time
let existing_ending_time: Option<TimestampMillis> =
select_stmt.query_row(params![network, name], |row| {
let ending_time: Result<Option<TimestampMillis>> = row.get(0);
ending_time
})?;
match select_stmt.query_row(params![network, name], |row| row.get(0)) {
Ok(ending_time) => ending_time,
Err(rusqlite::Error::QueryReturnedNoRows) => continue, // Handle no rows case gracefully
Err(e) => return Err(e.into()), // Propagate other errors
};

if let Some(existing_ending_time) = existing_ending_time {
// Update only if the new ending time is greater than the existing one
Expand Down
Loading