use anyhow::Result;
use futures::stream::{self, StreamExt, TryStreamExt};
use std::collections::HashSet;
use vrchatapi::apis;
pub async fn fetch_pages_parallel(
api_config: &vrchatapi::apis::configuration::Configuration,
offline: Option<bool>,
limit: Option<i32>,
) -> Result<Vec<vrchatapi::models::LimitedUserFriend>> {
let page_size = 60i32;
let max_pages = limit
.map(|lim| {
let min_pages = ((lim * 2) as f32 / page_size as f32).ceil() as usize;
min_pages.max(3) })
.unwrap_or(20);
let offsets = (0..max_pages)
.map(move |i| i as i32 * page_size)
.collect::<Vec<_>>();
let friends_batches: Vec<Vec<vrchatapi::models::LimitedUserFriend>> = stream::iter(offsets)
.map(|offset| {
let cfg = api_config.clone();
async move {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
apis::friends_api::get_friends(&cfg, Some(offset), Some(page_size), offline).await
}
})
.buffered(5) .take_while(|res| {
futures::future::ready(match res {
Ok(batch) => !batch.is_empty(),
Err(_) => true, })
})
.try_collect() .await?;
let result = friends_batches.into_iter().flatten().collect::<Vec<_>>();
Ok(result)
}
pub async fn fetch_all_friends_parallel(
api_config: &vrchatapi::apis::configuration::Configuration,
limit: Option<i32>,
) -> Result<Vec<vrchatapi::models::LimitedUserFriend>> {
let api_config_clone = api_config.clone();
let online_task =
tokio::spawn(
async move { fetch_pages_parallel(&api_config_clone, Some(false), limit).await },
);
let api_config_clone = api_config.clone();
let offline_task =
tokio::spawn(
async move { fetch_pages_parallel(&api_config_clone, Some(true), limit).await },
);
let (online_result, offline_result) = tokio::try_join!(online_task, offline_task)?;
let online = online_result?;
let offline = offline_result?;
let mut seen = HashSet::new();
let mut merged = Vec::new();
for friend in online.into_iter().chain(offline) {
if seen.insert(friend.id.clone()) {
merged.push(friend);
}
}
Ok(merged)
}