use std::time::Duration;
use reqwest::{Client, Response, StatusCode};
use serde::de::DeserializeOwned;
use serde_json::Value;
use crate::collect::env_expand::expand_env_var;
use crate::collect::errors::{CollectError, Result};
pub type Credentials = (String, String);
const MAX_HONOURED_RETRY_AFTER: Duration = Duration::from_secs(30);
pub async fn post_json<T: DeserializeOwned>(
client: &Client,
credentials: Option<&Credentials>,
url: &str,
body: &Value,
) -> Result<T> {
let mut req = client.post(url).json(body);
if let Some((user, token)) = credentials {
req = req.basic_auth(user, Some(token));
}
decode(req.send().await?).await
}
pub async fn get_json<T: DeserializeOwned>(
client: &Client,
credentials: Option<&Credentials>,
url: &str,
) -> Result<T> {
let mut req = client.get(url);
if let Some((user, token)) = credentials {
req = req.basic_auth(user, Some(token));
}
decode(req.send().await?).await
}
async fn decode<T: DeserializeOwned>(resp: Response) -> Result<T> {
let status = resp.status();
if status == StatusCode::TOO_MANY_REQUESTS || status == StatusCode::SERVICE_UNAVAILABLE {
return Err(CollectError::Throttled {
status: status.as_u16(),
retry_after: retry_after(&resp),
});
}
let resp = resp.error_for_status()?;
Ok(resp.json().await?)
}
fn retry_after(resp: &Response) -> Option<Duration> {
resp.headers()
.get(reqwest::header::RETRY_AFTER)?
.to_str()
.ok()?
.trim()
.parse::<u64>()
.ok()
.map(|secs| Duration::from_secs(secs).min(MAX_HONOURED_RETRY_AFTER))
}
pub fn expand_credential(field: &str, raw: &str) -> Result<String> {
let expanded = expand_env_var(raw);
if expanded.is_empty() {
return Err(CollectError::Config(format!(
"{field} is empty after expansion (config value: `{raw}`) — set the \
referenced environment variable, or remove the field to run \
unauthenticated"
)));
}
Ok(expanded)
}