From e78b0724c63c5058d88b632920485d9c1d655e98 Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Wed, 22 Jul 2026 16:19:40 +0300 Subject: [PATCH 1/6] feat(wallet): parallel actualization owned --- lez/wallet/src/multi_client.rs | 34 ++++++++++++++++++++++++++++++++-- 1 file changed, 32 insertions(+), 2 deletions(-) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index d987c55c..a2c6ca17 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -13,7 +13,7 @@ use anyhow::{Context as _, Result}; use lee_core::BlockId; use sequencer_service_rpc::{RpcClient as _, SequencerClient, SequencerClientBuilder}; use serde::{Deserialize, Serialize}; -use tokio::sync::RwLock; +use tokio::{sync::RwLock, task::JoinSet}; use url::Url; use crate::config::{MultiSequencerClientConfig, SequencerConnectionData}; @@ -450,7 +450,7 @@ pub async fn calibrate_client( }) } -pub async fn actualize_client(client: &SequencerClient) -> StatisticsUpdate { +async fn actualize_client(client: &SequencerClient) -> StatisticsUpdate { let (latency, block_id) = measure_request_duration(client).await; #[expect(clippy::as_conversions, reason = "int to float conversion is safe")] @@ -464,6 +464,36 @@ pub async fn actualize_client(client: &SequencerClient) -> StatisticsUpdate { }) } +async fn actualize_client_owned(client: SequencerClient) -> StatisticsUpdate { + let (latency, block_id) = measure_request_duration(&client).await; + + #[expect(clippy::as_conversions, reason = "int to float conversion is safe")] + let latency = latency as f32; + + block_id.map_or(StatisticsUpdate::Failure, |new_latest_block_id| { + StatisticsUpdate::Success { + latency, + new_latest_block_id, + } + }) +} + +pub async fn multi_actualize_clients(clients: &[SequencerClient]) -> Vec { + let mut join_set = JoinSet::new(); + + for client in clients { + join_set.spawn(actualize_client_owned(client.clone())); + } + + let mut statistic_updates = vec![]; + + while let Some(statistic_resp) = join_set.join_next().await { + statistic_updates.push(statistic_resp.unwrap_or(StatisticsUpdate::Failure)); + } + + statistic_updates +} + #[must_use] pub fn choose_leaders( client_list: &HashMap, From d34d269350fa9587acf87b63167b838b28209be4 Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Thu, 23 Jul 2026 14:56:18 +0300 Subject: [PATCH 2/6] fix(wallet): more parallelization --- lez/wallet/src/multi_client.rs | 252 ++++++++++++++++++++------------- 1 file changed, 150 insertions(+), 102 deletions(-) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index fafbbafb..902f318a 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -13,7 +13,7 @@ use anyhow::{Context as _, Result}; use lee_core::BlockId; use sequencer_service_rpc::{RpcClient as _, SequencerClient, SequencerClientBuilder}; use serde::{Deserialize, Serialize}; -use tokio::{sync::RwLock, task::JoinSet}; +use tokio::sync::RwLock; use url::Url; use crate::config::{MultiSequencerClientConfig, SequencerConnectionData}; @@ -93,7 +93,9 @@ impl MultiSequencerClient { statistics: &mut HashMap, multi_sequencer_client_config: &MultiSequencerClientConfig, ) -> Result> { - let mut client_list = HashMap::new(); + let mut actualization_list = vec![]; + let mut calibration_list = vec![]; + let mut client_list = vec![]; for SequencerConnectionData { sequencer_addr, @@ -118,30 +120,35 @@ impl MultiSequencerClient { .context("Failed to create sequencer client")? }; - // If there is statistics for client, actualize it - if let Some(statistic_mut) = statistics.get_mut(sequencer_addr) { - let statistic_updates = actualize_client(&sequencer_client).await; - - log::debug!( - "Metered call for {sequencer_addr:?}, statistic updates is {statistic_updates:?}" - ); - - statistic_mut.apply_updates(&[statistic_updates]); - - client_list.insert(sequencer_addr.clone(), sequencer_client); - // Otherwise calibrate client data - } else if let Some(client_statistics) = calibrate_client( - &sequencer_client, - multi_sequencer_client_config.calibration_limit, - ) - .await - { - statistics.insert(sequencer_addr.clone(), client_statistics); - client_list.insert(sequencer_addr.clone(), sequencer_client); - // There is no point in adding uncalibrated client + if statistics.contains_key(sequencer_addr) { + actualization_list.push((sequencer_addr.clone(), sequencer_client.clone())); } else { - log::warn!("Client {sequencer_addr:?} failed all {} calibration attempts, it may be unhealthy. - \n Consider bumping calibration_limit or remove this client altogether", multi_sequencer_client_config.calibration_limit); + calibration_list.push((sequencer_addr.clone(), sequencer_client.clone())); + } + + client_list.push((sequencer_addr.clone(), sequencer_client)); + } + + // This actually runs in one thread, but eah exact future is parallelized. + let (actualization_res, callibration_res) = tokio::join!( + multi_actualize_clients(&actualization_list), + multi_calibrate_clients( + &calibration_list, + multi_sequencer_client_config.calibration_limit + ) + ); + + for (addr, statistic_update_opt) in actualization_res { + if let Some(statistic_mut) = statistics.get_mut(&addr) + && let Some(statistic_update) = statistic_update_opt + { + statistic_mut.apply_updates(&[statistic_update]); + } + } + + for (addr, statistic_opt) in callibration_res { + if let Some(statistic) = statistic_opt { + statistics.insert(addr, statistic); } } @@ -215,7 +222,8 @@ impl MultiSequencerClient { leader_url: &Url, statistic_map: &mut HashMap>, ) -> Result { - let (resp, statistics_update) = tokio::join!(call(leader), actualize_client(leader)); + let (resp, statistics_update) = + tokio::join!(call(leader), actualize_client(leader.clone())); log::debug!("Metered call for {leader_url:?}, statistic updates is {statistics_update:?}",); @@ -290,7 +298,8 @@ impl MultiSequencerClient { let mut results = vec![]; for (leader, leader_url) in leaders { - let (resp, statistics_update) = tokio::join!(call(leader), actualize_client(leader)); + let (resp, statistics_update) = + tokio::join!(call(leader), actualize_client(leader.clone())); log::debug!( "Metered call for {leader_url:?}, statistic updates is {statistics_update:?}", @@ -414,17 +423,16 @@ async fn measure_request_duration(client: &SequencerClient) -> (u128, Option Option { +/// Calibrate statistics for one client. Takes `client` by value deliberately, cloning +/// `SequencerClient` is cheap. +async fn calibrate_client(client: SequencerClient, calibration_limit: usize) -> Option { let mut latencies = vec![]; let mut latest_block_id = 0; let mut errors: u64 = 0; // ToDo: Add some DDoS adaptation for _ in 0..calibration_limit { - let (latency, block_id) = measure_request_duration(client).await; + let (latency, block_id) = measure_request_duration(&client).await; let Some(block_id) = block_id else { errors = errors.saturating_add(1); @@ -458,21 +466,9 @@ pub async fn calibrate_client( }) } -async fn actualize_client(client: &SequencerClient) -> StatisticsUpdate { - let (latency, block_id) = measure_request_duration(client).await; - - #[expect(clippy::as_conversions, reason = "int to float conversion is safe")] - let latency = latency as f32; - - block_id.map_or(StatisticsUpdate::Failure, |new_latest_block_id| { - StatisticsUpdate::Success { - latency, - new_latest_block_id, - } - }) -} - -async fn actualize_client_owned(client: SequencerClient) -> StatisticsUpdate { +/// Actualize statistics for one client. Takes `client` by value deliberately, cloning +/// `SequencerClient` is cheap. +async fn actualize_client(client: SequencerClient) -> StatisticsUpdate { let (latency, block_id) = measure_request_duration(&client).await; #[expect(clippy::as_conversions, reason = "int to float conversion is safe")] @@ -486,32 +482,85 @@ async fn actualize_client_owned(client: SequencerClient) -> StatisticsUpdate { }) } -pub async fn multi_actualize_clients(clients: &[SequencerClient]) -> Vec { - let mut join_set = JoinSet::new(); +// Next 2 functions can be done in uniform way right now, but it will be incredibly cursed. +// ToDo: Return to it when "return-type notation" is stable - for client in clients { - join_set.spawn(actualize_client_owned(client.clone())); +pub async fn multi_actualize_clients( + clients: &[(Url, SequencerClient)], +) -> Vec<(Url, Option)> { + let mut handle_map = HashMap::new(); + + for (url, client) in clients { + let actualization_task = tokio::task::spawn(actualize_client(client.clone())); + handle_map.insert(url.clone(), actualization_task); } let mut statistic_updates = vec![]; - while let Some(statistic_resp) = join_set.join_next().await { - statistic_updates.push(statistic_resp.unwrap_or(StatisticsUpdate::Failure)); + // We anyhow need results of each one + #[expect( + clippy::iter_over_hash_type, + reason = "Ordering of map updates is not important" + )] + for (url, task) in handle_map { + let task_opt = task.await.ok(); + statistic_updates.push((url, task_opt)); } statistic_updates } +pub async fn multi_calibrate_clients( + clients: &[(Url, SequencerClient)], + calibration_limit: usize, +) -> Vec<(Url, Option)> { + let mut handle_map = HashMap::new(); + + for (url, client) in clients { + let calibration_task = + tokio::task::spawn(calibrate_client(client.clone(), calibration_limit)); + handle_map.insert(url.clone(), calibration_task); + } + + let mut statistics = vec![]; + + // We anyhow need results of each one + #[expect( + clippy::iter_over_hash_type, + reason = "Ordering of map updates is not important" + )] + for (url, task) in handle_map { + let task_opt = task + .await + .ok() + .inspect(|val| { + if val.is_none() { + log::warn!( + "Client {url:?} failed all {calibration_limit} calibration attempts, it may be unhealthy. + \n Consider bumping calibration_limit or remove this client altogether" + ); + } + }) + .flatten(); + statistics.push((url, task_opt)); + } + + statistics +} + #[must_use] -pub fn choose_leaders( - client_list: &HashMap, +/// Choosing leaders up to a `distribution_limit`. +/// +/// Assumes that all clients have their statistics in `statistics`. +fn choose_leaders( + client_vec: &[(Url, SequencerClient)], statistics: &HashMap, distribution_limit: usize, ) -> Option> { // Sort out all unmetered clients - let mut client_vec: Vec<_> = client_list - .keys() - .filter(|item| statistics.contains_key(*item)) + let client_vec: Vec<_> = client_vec + .iter() + .filter(|item| statistics.contains_key(&item.0)) .collect(); if client_vec.is_empty() { @@ -519,23 +568,28 @@ pub fn choose_leaders( } // Considering the nature of our requests, the latest_block_id is the dominant characteristic - let max_block_id_addr = client_vec.iter().fold(client_vec[0], |acc, x| { - let old_latest_block_id = statistics.get(acc).unwrap().latest_block_id; - let new_latest_block_id = statistics.get(*x).unwrap().latest_block_id; - if new_latest_block_id > old_latest_block_id { - *x - } else { - acc - } - }); + let max_block_id_addr = client_vec + .iter() + .fold(client_vec.first().unwrap(), |acc, x| { + let old_latest_block_id = statistics.get(&acc.0).unwrap().latest_block_id; + let new_latest_block_id = statistics.get(&x.0).unwrap().latest_block_id; + if new_latest_block_id > old_latest_block_id { + x + } else { + acc + } + }); - let max_block_id = statistics.get(max_block_id_addr).unwrap().latest_block_id; + let max_block_id = statistics + .get(&max_block_id_addr.0) + .unwrap() + .latest_block_id; // Sort out all clients running late - client_vec = client_vec + let mut res_vec: Vec<_> = client_vec .iter() .filter_map(|x| { - let latest_block_id = statistics.get(*x).unwrap().latest_block_id; + let latest_block_id = statistics.get(&x.0).unwrap().latest_block_id; (latest_block_id == max_block_id).then_some(*x) }) @@ -543,21 +597,21 @@ pub fn choose_leaders( // Get the clients with lesser or equal to average error ratio #[expect(clippy::as_conversions, reason = "int to float conversion is safe")] - let avg_err_ratio = client_vec.iter().fold(0_f32, |acc, x| { + let avg_err_ratio = res_vec.iter().fold(0_f32, |acc, x| { acc + error_ratio( - statistics.get(*x).unwrap().errors, - statistics.get(*x).unwrap().sample_size, + statistics.get(&x.0).unwrap().errors, + statistics.get(&x.0).unwrap().sample_size, ) - }) / (client_vec.len() as f32); + }) / (res_vec.len() as f32); - client_vec.sort_by(|a, b| { + res_vec.sort_by(|a, b| { let err_ratio_a = error_ratio( - statistics.get(*a).unwrap().errors, - statistics.get(*a).unwrap().sample_size, + statistics.get(&a.0).unwrap().errors, + statistics.get(&a.0).unwrap().sample_size, ); let err_ratio_b = error_ratio( - statistics.get(*b).unwrap().errors, - statistics.get(*b).unwrap().sample_size, + statistics.get(&b.0).unwrap().errors, + statistics.get(&b.0).unwrap().sample_size, ); err_ratio_a @@ -565,22 +619,22 @@ pub fn choose_leaders( .expect("Ratios must be a valid numbers") }); - let mut client_vec = client_vec[..(client_vec + res_vec = res_vec[..(res_vec .iter() .position(|item| { error_ratio( - statistics.get(*item).unwrap().errors, - statistics.get(*item).unwrap().sample_size, + statistics.get(&item.0).unwrap().errors, + statistics.get(&item.0).unwrap().sample_size, ) > avg_err_ratio }) - .unwrap_or(client_vec.len()))] + .unwrap_or(res_vec.len()))] .to_vec(); // Choose clients with least latency and variance - client_vec.sort_by(|a, b| { - let left = statistics.get(*a).unwrap(); + res_vec.sort_by(|a, b| { + let left = statistics.get(&a.0).unwrap(); let (left_lat, left_var) = (left.latency_avg, left.latency_var); - let right = statistics.get(*b).unwrap(); + let right = statistics.get(&b.0).unwrap(); let (right_lat, right_var) = (right.latency_avg, right.latency_var); let right_std = right_var.sqrt(); @@ -605,16 +659,10 @@ pub fn choose_leaders( }); Some( - client_vec - .iter() + res_vec + .into_iter() + .map(|(a, b)| (b.clone(), a.clone())) .take(distribution_limit) - .map(|addr| { - let client = client_list - .get(*addr) - .expect("Missing clients already sorted out"); - - (client.clone(), (*addr).clone()) - }) .collect(), ) } @@ -721,7 +769,7 @@ mod tests { builder.build(url).unwrap() } - fn four_client_list() -> (HashMap, [Url; 4]) { + fn four_client_list() -> (Vec<(Url, SequencerClient)>, [Url; 4]) { let addr_leader = Url::parse("http://127.0.0.1:3040").unwrap(); let addr_1 = Url::parse("http://127.0.0.1:3041").unwrap(); let addr_2 = Url::parse("http://127.0.0.1:3042").unwrap(); @@ -732,12 +780,12 @@ mod tests { let client_2 = client_from_url_unchecked(&addr_2); let client_3 = client_from_url_unchecked(&addr_3); - let mut client_list = HashMap::new(); - - client_list.insert(addr_leader.clone(), leader); - client_list.insert(addr_1.clone(), client_1); - client_list.insert(addr_2.clone(), client_2); - client_list.insert(addr_3.clone(), client_3); + let client_list = vec![ + (addr_leader.clone(), leader), + (addr_1.clone(), client_1), + (addr_2.clone(), client_2), + (addr_3.clone(), client_3), + ]; (client_list, [addr_leader, addr_1, addr_2, addr_3]) } From a3c2f08ea58e992faefeeaaa75c6d5b61a682ef6 Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Fri, 24 Jul 2026 13:46:45 +0300 Subject: [PATCH 3/6] fix(wallet_ffi): simplification --- lez/wallet/src/lib.rs | 25 ++---- lez/wallet/src/multi_client.rs | 136 ++++++++++++++++++++------------- 2 files changed, 92 insertions(+), 69 deletions(-) diff --git a/lez/wallet/src/lib.rs b/lez/wallet/src/lib.rs index ebe5028a..8fa65d78 100644 --- a/lez/wallet/src/lib.rs +++ b/lez/wallet/src/lib.rs @@ -87,6 +87,8 @@ pub enum ExecutionFailureKind { SignError(anyhow::Error), #[error("Sending transaction failed for each client")] MultiSequencerTransactionSendError, + #[error("Failed to join a task: {0}")] + JoinError(#[from] tokio::task::JoinError), } pub struct WalletCore { @@ -808,11 +810,7 @@ impl WalletCore { let call_res = self .multi_sequencer_client - .metered_send(async |client: &SequencerClient| { - client - .send_transaction(LeeTransaction::PrivacyPreserving(tx.clone())) - .await - }) + .metered_send_transaction(LeeTransaction::PrivacyPreserving(tx)) .await .into_iter() .find(std::result::Result::is_ok) @@ -877,17 +875,12 @@ impl WalletCore { let tx = lee::public_transaction::PublicTransaction::new(message, witness_set); - Ok(self - .multi_sequencer_client - .metered_send(async |client: &SequencerClient| { - client - .send_transaction(LeeTransaction::Public(tx.clone())) - .await - }) + self.multi_sequencer_client + .metered_send_transaction(LeeTransaction::Public(tx)) .await .into_iter() .find(std::result::Result::is_ok) - .ok_or(ExecutionFailureKind::MultiSequencerTransactionSendError)??) + .ok_or(ExecutionFailureKind::MultiSequencerTransactionSendError)? } pub async fn send_program_deployment_transaction(&self, bytecode: Vec) -> Result { @@ -896,11 +889,7 @@ impl WalletCore { Ok(self .multi_sequencer_client - .metered_send(async |client: &SequencerClient| { - client - .send_transaction(LeeTransaction::ProgramDeployment(transaction.clone())) - .await - }) + .metered_send_transaction(LeeTransaction::ProgramDeployment(transaction)) .await .into_iter() .find(std::result::Result::is_ok) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index 0e7e2fef..3bb8a215 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -10,13 +10,17 @@ use std::{collections::HashMap, path::Path, sync::Arc}; use anyhow::{Context as _, Result}; +use common::{HashType, transaction::LeeTransaction}; use lee_core::BlockId; use sequencer_service_rpc::{RpcClient as _, SequencerClient, SequencerClientBuilder}; use serde::{Deserialize, Serialize}; -use tokio::sync::RwLock; +use tokio::{sync::RwLock, task::JoinSet}; use url::Url; -use crate::config::{MultiSequencerClientConfig, SequencerConnectionData}; +use crate::{ + ExecutionFailureKind, + config::{MultiSequencerClientConfig, SequencerConnectionData}, +}; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Statistics { @@ -131,13 +135,19 @@ impl MultiSequencerClient { // This actually runs in one thread, but eah exact future is parallelized. let (actualization_res, callibration_res) = tokio::join!( - multi_actualize_clients(&actualization_list), + multi_actualize_clients(actualization_list), multi_calibrate_clients( - &calibration_list, + calibration_list, multi_sequencer_client_config.calibration_limit ) ); + for (addr, statistic_opt) in callibration_res { + if let Some(statistic) = statistic_opt { + statistics.insert(addr, statistic); + } + } + for (addr, statistic_update_opt) in actualization_res { if let Some(statistic_mut) = statistics.get_mut(&addr) && let Some(statistic_update) = statistic_update_opt @@ -146,14 +156,8 @@ impl MultiSequencerClient { } } - for (addr, statistic_opt) in callibration_res { - if let Some(statistic) = statistic_opt { - statistics.insert(addr, statistic); - } - } - let leader_list = choose_leaders( - &client_list, + client_list, statistics, multi_sequencer_client_config.distribution_limit, ) @@ -165,7 +169,10 @@ impl MultiSequencerClient { ); } - log::info!("Chosen leaders is {leader_list:?}"); + log::info!( + "Chosen leaders is {:?}", + leader_list.iter().map(|(_, addr)| addr).collect::>() + ); Ok(leader_list) } @@ -194,7 +201,10 @@ impl MultiSequencerClient { ) -> Result<()> { let leader_list = Self::setup(conn_data, statistics, multi_sequencer_client_config).await?; - log::info!("Chosen leaders is {leader_list:#?}"); + log::info!( + "Chosen leaders is {:?}", + leader_list.iter().map(|(_, addr)| addr).collect::>() + ); self.leader_list = leader_list; @@ -288,32 +298,55 @@ impl MultiSequencerClient { resp } - /// Metered call for `distribution_limit` amount of leaders for sending data, usually - /// transaction. - pub async fn metered_send Result>( + /// Metered `send_transaction` for `distribution_limit` amount of leaders. + /// + /// Less abstract that it could be, clean way to implement it in a more general way is + /// "return-type notation". + /// + /// `ToDo`: Return to it, when "return-type notation" is stable. + pub async fn metered_send_transaction( &self, - call: I, - ) -> Vec> { - let leaders = self.leaders().iter().take(self.config().distribution_limit); - + tx: LeeTransaction, + ) -> Vec> { // Collecting all statistics into one map to lock updates only once let mut statistic_map: HashMap> = HashMap::new(); let mut results = vec![]; + let mut join_set = JoinSet::new(); - for (leader, leader_url) in leaders { - let (resp, statistics_update) = - tokio::join!(call(leader), actualize_client(leader.clone())); + for (leader, leader_url) in self.leaders() { + let curr_leader = leader.clone(); + let curr_tx = tx.clone(); + let curr_url = leader_url.clone(); - log::debug!( - "Metered call for {leader_url:?}, statistic updates is {statistics_update:?}", - ); + join_set.spawn(async move { + ( + tokio::join!( + curr_leader.send_transaction(curr_tx), + actualize_client(curr_leader.clone()) + ), + curr_url, + ) + }); + } - statistic_map - .entry(leader_url.clone()) - .or_default() - .push(statistics_update); - results.push(resp); + while let Some(resp) = join_set.join_next().await { + let res = resp + .map_err(Into::into) + .and_then(|((resp, statistics_update), leader_url)| { + log::debug!( + "Metered call for {leader_url:?}, statistic updates is {statistics_update:?}", + ); + + statistic_map + .entry(leader_url) + .or_default() + .push(statistics_update); + + resp.map_err(Into::into) + }); + + results.push(res); } { @@ -492,13 +525,14 @@ async fn actualize_client(client: SequencerClient) -> StatisticsUpdate { // ToDo: Return to it when "return-type notation" is stable pub async fn multi_actualize_clients( - clients: &[(Url, SequencerClient)], + clients: Vec<(Url, SequencerClient)>, ) -> Vec<(Url, Option)> { let mut handle_map = HashMap::new(); for (url, client) in clients { - let actualization_task = tokio::task::spawn(actualize_client(client.clone())); - handle_map.insert(url.clone(), actualization_task); + // `client` here must have 'static lifetime, so we can not use reference + let actualization_task = tokio::task::spawn(actualize_client(client)); + handle_map.insert(url, actualization_task); } let mut statistic_updates = vec![]; @@ -517,15 +551,15 @@ pub async fn multi_actualize_clients( } pub async fn multi_calibrate_clients( - clients: &[(Url, SequencerClient)], + clients: Vec<(Url, SequencerClient)>, calibration_limit: usize, ) -> Vec<(Url, Option)> { let mut handle_map = HashMap::new(); for (url, client) in clients { - let calibration_task = - tokio::task::spawn(calibrate_client(client.clone(), calibration_limit)); - handle_map.insert(url.clone(), calibration_task); + // `client` here must have 'static lifetime, so we can not use reference + let calibration_task = tokio::task::spawn(calibrate_client(client, calibration_limit)); + handle_map.insert(url, calibration_task); } let mut statistics = vec![]; @@ -559,13 +593,13 @@ pub async fn multi_calibrate_clients( /// /// Assumes that all clients have their statistics in `statistics`. fn choose_leaders( - client_vec: &[(Url, SequencerClient)], + client_vec: Vec<(Url, SequencerClient)>, statistics: &HashMap, distribution_limit: usize, ) -> Option> { // Sort out all unmetered clients let client_vec: Vec<_> = client_vec - .iter() + .into_iter() .filter(|item| statistics.contains_key(&item.0)) .collect(); @@ -594,10 +628,10 @@ fn choose_leaders( // Sort out all clients running late let mut res_vec: Vec<_> = client_vec .iter() - .filter_map(|x| { + .filter(|x| { let latest_block_id = statistics.get(&x.0).unwrap().latest_block_id; - (latest_block_id == max_block_id).then_some(*x) + latest_block_id == max_block_id }) .collect(); @@ -1026,7 +1060,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 1).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 1).unwrap(); let (_, leader_url) = leaders.first().unwrap(); @@ -1083,7 +1117,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 1).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 1).unwrap(); let (_, leader_url) = leaders.first().unwrap(); @@ -1140,7 +1174,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 1).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 1).unwrap(); let (_, leader_url) = leaders.first().unwrap(); @@ -1197,7 +1231,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 1).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 1).unwrap(); let (_, leader_url) = leaders.first().unwrap(); @@ -1254,7 +1288,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 2).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 2).unwrap(); let mut url_set_origin = HashSet::new(); let mut url_set_res = HashSet::new(); @@ -1322,7 +1356,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 2).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 2).unwrap(); let (_, helm) = leaders.first().unwrap(); @@ -1380,7 +1414,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 2).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 2).unwrap(); let (_, leader_url_first) = leaders[0].clone(); let (_, leader_url_second) = leaders[1].clone(); @@ -1440,7 +1474,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 2).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 2).unwrap(); let (_, leader_url_first) = leaders[0].clone(); let (_, leader_url_second) = leaders[1].clone(); @@ -1500,7 +1534,7 @@ mod tests { }, ); - let leaders = choose_leaders(&client_list, &statistics, 2).unwrap(); + let leaders = choose_leaders(client_list, &statistics, 2).unwrap(); let (_, leader_url_first) = leaders[0].clone(); let (_, leader_url_second) = leaders[1].clone(); From 60a78aec08212eb0fd71bb0e007a0cf4be273834 Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Mon, 27 Jul 2026 13:24:56 +0300 Subject: [PATCH 4/6] fix(wallet): suggestions fix --- lez/wallet/src/multi_client.rs | 42 ++++++++++++++++------------------ 1 file changed, 20 insertions(+), 22 deletions(-) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index 3bb8a215..b248cc53 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -11,7 +11,9 @@ use std::{collections::HashMap, path::Path, sync::Arc}; use anyhow::{Context as _, Result}; use common::{HashType, transaction::LeeTransaction}; +use itertools::Itertools as _; use lee_core::BlockId; +use log::warn; use sequencer_service_rpc::{RpcClient as _, SequencerClient, SequencerClientBuilder}; use serde::{Deserialize, Serialize}; use tokio::{sync::RwLock, task::JoinSet}; @@ -97,6 +99,14 @@ impl MultiSequencerClient { statistics: &mut HashMap, multi_sequencer_client_config: &MultiSequencerClientConfig, ) -> Result> { + if !conn_data + .iter() + .map(|conn| &conn.sequencer_addr) + .all_unique() + { + anyhow::bail!("All adresses must be unique"); + } + let mut actualization_list = vec![]; let mut calibration_list = vec![]; let mut client_list = vec![]; @@ -133,7 +143,7 @@ impl MultiSequencerClient { client_list.push((sequencer_addr.clone(), sequencer_client)); } - // This actually runs in one thread, but eah exact future is parallelized. + // This actually starts in one thread, but each exact future produces more tasks. let (actualization_res, callibration_res) = tokio::join!( multi_actualize_clients(actualization_list), multi_calibrate_clients( @@ -300,7 +310,7 @@ impl MultiSequencerClient { /// Metered `send_transaction` for `distribution_limit` amount of leaders. /// - /// Less abstract that it could be, clean way to implement it in a more general way is + /// Less abstract than it could be, clean way to implement it in a more general way is /// "return-type notation". /// /// `ToDo`: Return to it, when "return-type notation" is stable. @@ -332,6 +342,7 @@ impl MultiSequencerClient { while let Some(resp) = join_set.join_next().await { let res = resp + .inspect_err(|j_err| warn!("Task failed with join error: {j_err:?}")) .map_err(Into::into) .and_then(|((resp, statistics_update), leader_url)| { log::debug!( @@ -524,7 +535,7 @@ async fn actualize_client(client: SequencerClient) -> StatisticsUpdate { // Next 2 functions can be done in uniform way right now, but it will be incredibly cursed. // ToDo: Return to it when "return-type notation" is stable -pub async fn multi_actualize_clients( +async fn multi_actualize_clients( clients: Vec<(Url, SequencerClient)>, ) -> Vec<(Url, Option)> { let mut handle_map = HashMap::new(); @@ -550,7 +561,7 @@ pub async fn multi_actualize_clients( statistic_updates } -pub async fn multi_calibrate_clients( +async fn multi_calibrate_clients( clients: Vec<(Url, SequencerClient)>, calibration_limit: usize, ) -> Vec<(Url, Option)> { @@ -627,7 +638,7 @@ fn choose_leaders( // Sort out all clients running late let mut res_vec: Vec<_> = client_vec - .iter() + .into_iter() .filter(|x| { let latest_block_id = statistics.get(&x.0).unwrap().latest_block_id; @@ -680,28 +691,15 @@ fn choose_leaders( let right_std = right_var.sqrt(); let left_std = left_var.sqrt(); - // Client is better if its average is better and variance does not make it worse - // So basically we want this: - // [-right_std < left_lat < right_lat < +left_std < +right_std] - // - // However one can argue that this: - // - // [-right_std < right_lat < left_lat < +left_std < +right_std] - // - // is still better, but it is up to discussion - let first_ordering = left_lat.total_cmp(&right_lat); - match first_ordering { - std::cmp::Ordering::Greater => first_ordering, - std::cmp::Ordering::Less | std::cmp::Ordering::Equal => { - (left_lat + left_std).total_cmp(&(right_lat + right_std)) - } - } + // Client is better if its average + std is lesser: + // [-right_std < right_lat < +left_std < +right_std] + (left_lat + left_std).total_cmp(&(right_lat + right_std)) }); Some( res_vec .into_iter() - .map(|(a, b)| (b.clone(), a.clone())) + .map(|(a, b)| (b, a)) .take(distribution_limit) .collect(), ) From 8389f060c46cd05074731049d80fd5b3d59d92c6 Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Mon, 27 Jul 2026 16:14:03 +0300 Subject: [PATCH 5/6] fix(wallet): suggestions 2 --- lez/wallet/src/multi_client.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index b248cc53..84cff2ee 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -13,7 +13,6 @@ use anyhow::{Context as _, Result}; use common::{HashType, transaction::LeeTransaction}; use itertools::Itertools as _; use lee_core::BlockId; -use log::warn; use sequencer_service_rpc::{RpcClient as _, SequencerClient, SequencerClientBuilder}; use serde::{Deserialize, Serialize}; use tokio::{sync::RwLock, task::JoinSet}; @@ -342,7 +341,7 @@ impl MultiSequencerClient { while let Some(resp) = join_set.join_next().await { let res = resp - .inspect_err(|j_err| warn!("Task failed with join error: {j_err:?}")) + .inspect_err(|j_err| log::warn!("Task failed with join error: {j_err:?}")) .map_err(Into::into) .and_then(|((resp, statistics_update), leader_url)| { log::debug!( @@ -599,10 +598,10 @@ async fn multi_calibrate_clients( statistics } -#[must_use] /// Choosing leaders up to a `distribution_limit`. /// /// Assumes that all clients have their statistics in `statistics`. +#[must_use] fn choose_leaders( client_vec: Vec<(Url, SequencerClient)>, statistics: &HashMap, From 277f64147261a2db464cd5285ad26110a0939b7b Mon Sep 17 00:00:00 2001 From: Pravdyvy Date: Tue, 28 Jul 2026 07:55:43 +0300 Subject: [PATCH 6/6] fix(wallet): suggestions 3 --- lez/wallet/src/multi_client.rs | 23 ++++++++++++----------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/lez/wallet/src/multi_client.rs b/lez/wallet/src/multi_client.rs index 84cff2ee..49c6bbda 100644 --- a/lez/wallet/src/multi_client.rs +++ b/lez/wallet/src/multi_client.rs @@ -103,7 +103,7 @@ impl MultiSequencerClient { .map(|conn| &conn.sequencer_addr) .all_unique() { - anyhow::bail!("All adresses must be unique"); + anyhow::bail!("All addresses must be unique"); } let mut actualization_list = vec![]; @@ -669,16 +669,17 @@ fn choose_leaders( .expect("Ratios must be a valid numbers") }); - res_vec = res_vec[..(res_vec - .iter() - .position(|item| { - error_ratio( - statistics.get(&item.0).unwrap().errors, - statistics.get(&item.0).unwrap().sample_size, - ) > avg_err_ratio - }) - .unwrap_or(res_vec.len()))] - .to_vec(); + res_vec.truncate( + res_vec + .iter() + .position(|item| { + error_ratio( + statistics.get(&item.0).unwrap().errors, + statistics.get(&item.0).unwrap().sample_size, + ) > avg_err_ratio + }) + .unwrap_or(res_vec.len()), + ); // Choose clients with least latency and variance res_vec.sort_by(|a, b| {