use super::Client;
use crate::client::utils::retry;
use crate::client::{connections::QueryResult, errors::Error};
use crate::messaging::{
data::{DataQuery, ServiceMsg},
ServiceAuth, WireMsg,
};
use crate::types::{PublicKey, Signature};
use bytes::Bytes;
use tracing::{debug, info_span, Instrument};
impl Client {
#[instrument(skip(self), level = "debug")]
pub(crate) async fn send_query(&self, query: DataQuery) -> Result<QueryResult, Error> {
let client_pk = self.public_key();
let msg = ServiceMsg::Query(query.clone());
let serialised_query = WireMsg::serialize_msg_payload(&msg)?;
let signature = self.keypair.sign(&serialised_query);
let max_retry_count = 11.0;
let starting_query_timeout = self.query_timeout.div_f32(max_retry_count);
trace!(
"Setting up query retry, initial interval is: {:?}",
starting_query_timeout
);
retry(
|| {
async {
debug!(
"Attempting {:?} with a query timeout of {:?}",
query, starting_query_timeout
);
let res = tokio::time::timeout(
starting_query_timeout,
self.send_signed_query(
query.clone(),
client_pk,
serialised_query.clone(),
signature.clone(),
),
)
.await;
match res {
Ok(inner_result) => match inner_result {
Ok(query_result) => Ok(Ok(query_result)),
Err(error) => match error {
Error::InsufficientElderConnections { .. } => {
Err(error).map_err(backoff::Error::Transient)
}
_ => Err(error).map_err(backoff::Error::Permanent),
},
},
Err(_elapsed) => {
Err(Error::QueryTimedOut).map_err(backoff::Error::Transient)
}
}
}
.instrument(info_span!("Attempting a query"))
},
starting_query_timeout,
self.query_timeout,
)
.await
.map_err(|_| {
debug!("retries all failed for {:?}, returning no response", query);
Error::NoResponse
})?
}
pub(crate) async fn send_signed_query(
&self,
query: DataQuery,
client_pk: PublicKey,
serialised_query: Bytes,
signature: Signature,
) -> Result<QueryResult, Error> {
debug!("Sending Query: {:?}", query);
let auth = ServiceAuth {
public_key: client_pk,
signature,
};
self.session.send_query(query, auth, serialised_query).await
}
}