use serde::{Deserialize, Serialize};
use tokio::time::{Duration, sleep};
use tracing::{instrument, trace};
use super::{Client, validate_limit, validate_market_id, validate_min_balance};
use crate::error::Result;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Holder {
#[serde(rename = "proxyWallet")]
pub proxy_wallet: String,
pub bio: String,
pub asset: String,
pub pseudonym: String,
pub amount: f64,
#[serde(rename = "displayUsernamePublic")]
pub display_username_public: bool,
#[serde(rename = "outcomeIndex")]
pub outcome_index: i32,
pub name: String,
#[serde(rename = "profileImage")]
pub profile_image: String,
#[serde(rename = "profileImageOptimized")]
pub profile_image_optimized: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MarketTopHolders {
pub token: String,
pub holders: Vec<Holder>,
}
const MAX_RETRIES: usize = 2;
const RETRY_DELAY_MS: u64 = 300;
fn is_retryable_status(status: u16) -> bool {
status == 408 || status == 429 || (500..=504).contains(&status)
}
impl Client {
#[instrument(skip(self, markets), level = "trace")]
pub async fn get_market_top_holders(
&self,
markets: &[&str],
limit: Option<i32>,
min_balance: Option<i32>,
) -> Result<Vec<MarketTopHolders>> {
for market_id in markets {
validate_market_id(market_id)?;
}
validate_limit(limit)?;
validate_min_balance(min_balance)?;
let mut url = self.build_url("holders");
if !markets.is_empty() {
let market_value = markets.join(",");
url.query_pairs_mut().append_pair("market", &market_value);
}
if let Some(l) = limit {
url.query_pairs_mut().append_pair("limit", &l.to_string());
}
if let Some(mb) = min_balance {
url.query_pairs_mut()
.append_pair("minBalance", &mb.to_string());
}
trace!(url = %url, method = "GET", market_count = markets.len(), "sending HTTP request");
for attempt in 0..=MAX_RETRIES {
match self.http_client.get(url.clone()).send().await {
Ok(resp) => {
let status = resp.status();
trace!(status = %status, attempt = attempt, "received HTTP response");
if status.is_success() {
let holders_response: Vec<MarketTopHolders> = resp.json().await?;
trace!(count = holders_response.len(), "received holders data");
return Ok(holders_response);
}
let status_code = status.as_u16();
if attempt < MAX_RETRIES && is_retryable_status(status_code) {
trace!(status = status_code, attempt = attempt, "retrying request");
sleep(Duration::from_millis(RETRY_DELAY_MS * (attempt as u64 + 1))).await;
continue;
}
return Err(self.check_response(resp).await.unwrap_err());
}
Err(err) => {
trace!(error = %err, attempt = attempt, "HTTP request error");
if attempt < MAX_RETRIES && err.is_timeout() {
trace!(attempt = attempt, "retrying after timeout");
sleep(Duration::from_millis(RETRY_DELAY_MS * (attempt as u64 + 1))).await;
continue;
}
return Err(err.into());
}
}
}
unreachable!("retry loop should return success or error")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_is_retryable_status() {
assert!(is_retryable_status(408)); assert!(is_retryable_status(429)); assert!(is_retryable_status(500)); assert!(is_retryable_status(501)); assert!(is_retryable_status(502)); assert!(is_retryable_status(503)); assert!(is_retryable_status(504));
assert!(!is_retryable_status(200)); assert!(!is_retryable_status(400)); assert!(!is_retryable_status(401)); assert!(!is_retryable_status(403)); assert!(!is_retryable_status(404)); assert!(!is_retryable_status(499)); assert!(!is_retryable_status(505)); }
}