use super::connection::Connection;
use crate::block_tlb::{BlockIdExt, MaybeAccount};
use crate::clients::client_types::MasterchainInfo;
use crate::clients::lite_client::config::{LiteClientConfig, LiteReqParams};
use crate::error::TLError;
use crate::libs_dict::LibsDict;
use crate::unwrap_lite_response;
use auto_pool::config::{AutoPoolConfig, PickStrategy};
use auto_pool::pool::AutoPool;
use std::cmp::max;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering::Relaxed;
use std::sync::Arc;
use std::time::Duration;
use tokio_retry::strategy::FixedInterval;
use tokio_retry::RetryIf;
use ton_lib_core::cell::{TonCellRef, TonHash};
use ton_lib_core::constants::{TON_MASTERCHAIN, TON_SHARD_FULL};
use ton_lib_core::error::TLCoreError;
use ton_lib_core::traits::tlb::TLB;
use ton_lib_core::types::TonAddress;
use ton_liteapi::tl::common::{AccountId, Int256};
use ton_liteapi::tl::request::{
GetAccountState, GetBlock, GetLibraries, LookupBlock, Request, WaitMasterchainSeqno, WrappedRequest,
};
use ton_liteapi::tl::response::{BlockData, Response};
const WAIT_MC_SEQNO_MS: u32 = 5000;
const WAIT_CONNECTION_MS: u64 = 5;
#[derive(Clone)]
pub struct LiteClient {
inner: Arc<Inner>,
}
impl LiteClient {
pub fn new(config: LiteClientConfig) -> Result<Self, TLError> {
Ok(Self {
inner: Arc::new(Inner::new(config)?),
})
}
pub async fn get_mc_info(&self) -> Result<MasterchainInfo, TLError> {
let rsp = self.exec(Request::GetMasterchainInfo, None, None).await?;
let mc_info = unwrap_lite_response!(rsp, MasterchainInfo)?;
Ok(mc_info.into())
}
pub async fn lookup_mc_block(&self, seqno: u32) -> Result<BlockIdExt, TLError> {
self.lookup_block(TON_MASTERCHAIN, TON_SHARD_FULL, seqno).await
}
pub async fn lookup_block(&self, wc: i32, shard: u64, seqno: u32) -> Result<BlockIdExt, TLError> {
let req = Request::LookupBlock(LookupBlock {
mode: (),
id: ton_liteapi::tl::common::BlockId {
workchain: wc,
shard,
seqno,
},
seqno: Some(()),
lt: None,
utime: None,
with_state_update: None,
with_value_flow: None,
with_extra: None,
with_shard_hashes: None,
with_prev_blk_signatures: None,
});
let rsp = self.exec(req, Some(seqno), None).await?;
let lite_id = unwrap_lite_response!(rsp, BlockHeader)?.id;
Ok(lite_id.into())
}
pub async fn get_block(&self, block_id: BlockIdExt, params: Option<LiteReqParams>) -> Result<BlockData, TLError> {
let seqno = block_id.seqno;
let req = Request::GetBlock(GetBlock { id: block_id.into() });
let rsp = self.exec(req, Some(seqno), params).await?;
unwrap_lite_response!(rsp, BlockData)
}
pub async fn get_account_state(
&self,
address: &TonAddress,
mc_seqno: u32,
params: Option<LiteReqParams>,
) -> Result<MaybeAccount, TLError> {
let req = Request::GetAccountState(GetAccountState {
id: self.lookup_mc_block(mc_seqno).await?.into(),
account: AccountId {
workchain: address.workchain,
id: Int256(*address.hash.as_slice_sized()),
},
});
let rsp = self.exec_with_timeout(req, Some(mc_seqno), params).await?;
let account_state_rsp = unwrap_lite_response!(rsp, AccountState)?;
Ok(MaybeAccount::from_boc(&account_state_rsp.state)?)
}
pub async fn get_libs(&self, lib_ids: &[TonHash], params: Option<LiteReqParams>) -> Result<LibsDict, TLError> {
self.inner.get_libs_impl(lib_ids, params).await
}
pub async fn exec(
&self,
req: Request,
wait_mc_seqno: Option<u32>,
params: Option<LiteReqParams>,
) -> Result<Response, TLError> {
self.exec_with_timeout(req, wait_mc_seqno, params).await
}
pub async fn exec_with_timeout(
&self,
request: Request,
wait_mc_seqno: Option<u32>,
params: Option<LiteReqParams>,
) -> Result<Response, TLError> {
self.inner.exec_with_retries(request, wait_mc_seqno, params).await
}
}
struct Inner {
config: LiteClientConfig,
conn_pool: AutoPool<Connection>,
global_req_id: AtomicU64,
}
impl Inner {
fn new(config: LiteClientConfig) -> Result<Self, TLError> {
let conn_per_node = max(1, config.connections_per_node);
log::info!(
"Creating LiteClient with {} conns per node; nodes_cnt: {}, default_req_params: {:?}",
conn_per_node,
config.net_config.lite_endpoints.len(),
config.default_req_params,
);
let mut connections = Vec::new();
for _ in 0..conn_per_node {
for endpoint in &config.net_config.lite_endpoints {
let conn = Connection::new(endpoint.clone(), config.conn_timeout)?;
connections.push(conn);
}
}
let ap_config = AutoPoolConfig {
wait_duration: Duration::MAX,
lock_duration: Duration::from_millis(2),
sleep_duration: Duration::from_millis(WAIT_CONNECTION_MS),
pick_strategy: PickStrategy::RANDOM,
};
let connection_pool = AutoPool::new_with_config(ap_config, connections);
Ok(Self {
config,
conn_pool: connection_pool,
global_req_id: AtomicU64::new(0),
})
}
async fn get_libs_impl(&self, lib_ids: &[TonHash], params: Option<LiteReqParams>) -> Result<LibsDict, TLError> {
let mut libs_dict = LibsDict::default();
for chunk in lib_ids.chunks(16) {
let request = Request::GetLibraries(GetLibraries {
library_list: chunk.iter().map(|x| Int256(*x.as_slice_sized())).collect(),
});
let rsp = self.exec_with_retries(request, None, params).await?;
let result = unwrap_lite_response!(rsp, LibraryResult)?;
let dict_items = result
.result
.into_iter()
.map(|x| {
let hash = TonHash::from_slice_sized(&x.hash.0);
let lib = TonCellRef::from_boc(x.data.as_slice())?;
Ok::<_, TLCoreError>((hash, lib))
})
.collect::<Result<Vec<_>, TLCoreError>>()?;
let req_cnt = chunk.len();
let rsp_cnt = dict_items.len();
if req_cnt != rsp_cnt {
let got_hashes: Vec<_> = dict_items.iter().map(|x| &x.0).collect();
log::warn!(
"[get_libs_impl] expected {req_cnt} libs, got {rsp_cnt}:\n\
requested: {chunk:?}\n\
got: {got_hashes:?}",
);
}
for item in dict_items {
libs_dict.insert(item.0, item.1);
}
}
Ok(libs_dict)
}
async fn exec_with_retries(
&self,
req: Request,
wait_seqno: Option<u32>,
params: Option<LiteReqParams>,
) -> Result<Response, TLError> {
let wrap_req = WrappedRequest {
wait_masterchain_seqno: wait_seqno.map(|seqno| WaitMasterchainSeqno {
seqno,
timeout_ms: WAIT_MC_SEQNO_MS,
}),
request: req,
};
let req_params = params.as_ref().unwrap_or(&self.config.default_req_params);
let req_id = self.global_req_id.fetch_add(1, Relaxed);
let fi = FixedInterval::new(req_params.retry_waiting);
let strategy = fi.take(req_params.retries_count as usize);
let exec_request = || async { self.exec_impl(req_id, &wrap_req, req_params.query_timeout).await };
RetryIf::spawn(strategy, exec_request, retry_condition).await
}
async fn exec_impl(&self, req_id: u64, req: &WrappedRequest, req_timeout: Duration) -> Result<Response, TLError> {
log::trace!("LiteClient exec_impl: req_id={req_id}, req={req:?}");
let mut conn = self.conn_pool.get_async().await.unwrap();
conn.exec(req.clone(), req_timeout).await
}
}
fn retry_condition(error: &TLError) -> bool { !matches!(error, TLError::LiteClientWrongResponse(..)) }