use tokio::sync::Semaphore;
use core::result::Result;
use std::{sync::Arc, time::Instant};
use crate::error::Error;
use super::types::BootParameters;
pub async fn get(
shasta_token: &str,
shasta_base_url: &str,
shasta_root_cert: &[u8],
xnames: &[String],
) -> Result<Vec<BootParameters>, Error> {
log::info!("Get BSS bootparameters");
let client;
let client_builder = reqwest::Client::builder()
.add_root_certificate(reqwest::Certificate::from_pem(shasta_root_cert)?);
if std::env::var("SOCKS5").is_ok() {
log::debug!("SOCKS5 enabled");
let socks5proxy = reqwest::Proxy::all(std::env::var("SOCKS5").unwrap())?;
client = client_builder.proxy(socks5proxy).build()?;
} else {
client = client_builder.build()?;
}
let url_api =
format!("{}/bss/boot/v1/bootparameters", shasta_base_url.to_string());
let params: Vec<_> = xnames.iter().map(|xname| ("name", xname)).collect();
let response = client
.get(url_api)
.query(¶ms)
.bearer_auth(shasta_token)
.send()
.await
.map_err(|error| Error::NetError(error))?;
if response.status().is_success() {
response
.json::<Vec<BootParameters>>()
.await
.map_err(|error| Error::NetError(error))
} else {
let payload = response
.text()
.await
.map_err(|error| Error::NetError(error))?;
Err(Error::Message(payload))
}
}
pub async fn get_all(
shasta_token: &str,
shasta_base_url: &str,
shasta_root_cert: &[u8],
) -> Result<Vec<BootParameters>, Error> {
get(shasta_token, shasta_base_url, shasta_root_cert, &[]).await
}
pub async fn get_multiple(
shasta_token: &str,
shasta_base_url: &str,
shasta_root_cert: &[u8],
xnames: &[String],
) -> Result<Vec<BootParameters>, Error> {
let start = Instant::now();
let chunk_size = 30;
let mut boot_params_vec = Vec::new();
let mut tasks = tokio::task::JoinSet::new();
let sem = Arc::new(Semaphore::new(10));
for sub_node_list in xnames.chunks(chunk_size) {
let shasta_token_string = shasta_token.to_string();
let shasta_base_url_string = shasta_base_url.to_string();
let shasta_root_cert_vec = shasta_root_cert.to_vec();
let permit = Arc::clone(&sem).acquire_owned().await;
let node_vec = sub_node_list.to_vec();
tasks.spawn(async move {
let _permit = permit;
get(
&shasta_token_string,
&shasta_base_url_string,
&shasta_root_cert_vec,
&node_vec,
)
.await
.unwrap()
});
}
while let Some(message) = tasks.join_next().await {
if let Ok(mut node_status_vec) = message {
boot_params_vec.append(&mut node_status_vec);
}
}
let duration = start.elapsed();
log::info!("Time elapsed to get BSS bootparameters is: {:?}", duration);
Ok(boot_params_vec)
}
pub fn post(
base_url: &str,
auth_token: &str,
root_cert: &[u8],
boot_parameters: BootParameters,
) -> Result<(), Error> {
let client_builder = reqwest::blocking::Client::builder()
.add_root_certificate(reqwest::Certificate::from_pem(root_cert)?);
let client = if let Ok(socks5_env) = std::env::var("SOCKS5") {
log::debug!("SOCKS5 enabled");
let socks5proxy = reqwest::Proxy::all(socks5_env)?;
client_builder.proxy(socks5proxy).build()?
} else {
client_builder.build()?
};
let api_url = format!("{}/boot/v1/bootparameters", base_url);
let response = client
.post(api_url)
.bearer_auth(auth_token)
.json(&boot_parameters)
.send()
.map_err(|error| Error::NetError(error))?;
if response.status().is_success() {
Ok(())
} else {
Err(Error::Message(response.text()?))
}
}
pub async fn put(
shasta_base_url: &str,
shasta_token: &str,
shasta_root_cert: &[u8],
boot_parameters: BootParameters,
) -> Result<BootParameters, Error> {
let client;
let client_builder = reqwest::Client::builder()
.add_root_certificate(reqwest::Certificate::from_pem(shasta_root_cert)?);
if std::env::var("SOCKS5").is_ok() {
log::debug!("SOCKS5 enabled");
let socks5proxy = reqwest::Proxy::all(std::env::var("SOCKS5").unwrap())?;
client = client_builder.proxy(socks5proxy).build()?;
} else {
client = client_builder.build()?;
}
let api_url = format!("{}/bss/boot/v1/bootparameters", shasta_base_url);
log::debug!(
"request payload:\n{}",
serde_json::to_string_pretty(&boot_parameters).unwrap()
);
let response = client
.put(api_url)
.json(&boot_parameters)
.bearer_auth(shasta_token)
.send()
.await
.map_err(|error| Error::NetError(error))?;
if response.status().is_success() {
Ok(response.json().await?)
} else {
Err(Error::Message(response.text().await?))
}
}
pub async fn patch(
shasta_base_url: &str,
shasta_token: &str,
shasta_root_cert: &[u8],
boot_parameters: &BootParameters,
) -> Result<(), Error> {
let client;
let client_builder = reqwest::Client::builder()
.add_root_certificate(reqwest::Certificate::from_pem(shasta_root_cert)?);
if std::env::var("SOCKS5").is_ok() {
log::debug!("SOCKS5 enabled");
let socks5proxy = reqwest::Proxy::all(std::env::var("SOCKS5").unwrap())?;
client = client_builder.proxy(socks5proxy).build()?;
} else {
client = client_builder.build()?;
}
let api_url = format!("{}/bss/boot/v1/bootparameters", shasta_base_url);
let response = client
.patch(api_url)
.json(&boot_parameters)
.bearer_auth(shasta_token)
.send()
.await
.map_err(|error| Error::NetError(error))?;
if response.status().is_success() {
Ok(())
} else {
Err(Error::Message(response.text().await?))
}
}