use reqwest::Method;
use crate::{
Client, Connection, ConnectionAuditEntry, ConnectionLabelStat, ConnectionLabelsUpdate,
ConnectionListOptions, ConnectionListPage, ConnectionPagination, ConnectionRequest,
ConnectionStatusUpdate, Error, RequestOptions, Result,
};
use super::path_segment;
#[derive(Debug, Clone)]
pub struct Connections {
client: Client,
}
impl Connections {
pub(crate) fn new(client: Client) -> Self {
Self { client }
}
pub async fn create(
&self,
input: &ConnectionRequest,
options: RequestOptions,
) -> Result<Connection> {
let request = self
.client
.request(Method::POST, "/v1/connections")?
.json(input);
self.client.send(request, options).await
}
pub async fn list(&self, options: &ConnectionListOptions) -> Result<Vec<Connection>> {
Ok(self.list_page(options).await?.connections)
}
pub async fn list_page(&self, options: &ConnectionListOptions) -> Result<ConnectionListPage> {
let request = self
.client
.request(Method::GET, "/v1/connections")?
.query(options);
let value: serde_json::Value = self.client.send(request, RequestOptions::default()).await?;
match value {
serde_json::Value::Array(items) => {
let connections: Vec<Connection> =
serde_json::from_value(serde_json::Value::Array(items))
.map_err(|error| Error::UnexpectedResponse(error.to_string()))?;
let total = connections.len() as u64;
Ok(ConnectionListPage {
connections,
pagination: ConnectionPagination {
total,
limit: options.limit.unwrap_or(total as u32),
offset: options.offset.unwrap_or(0),
},
})
}
serde_json::Value::Object(_) => serde_json::from_value(value)
.map_err(|error| Error::UnexpectedResponse(error.to_string())),
_ => Err(Error::UnexpectedResponse(
"expected a connections response object".into(),
)),
}
}
pub async fn list_all(&self, options: &ConnectionListOptions) -> Result<Vec<Connection>> {
let mut filters = options.clone();
let mut connections = Vec::new();
let mut offset = 0_u32;
loop {
filters.limit = Some(200);
filters.offset = Some(offset);
let page = self.list_page(&filters).await?;
let batch_len = page.connections.len() as u32;
connections.extend(page.connections);
let next_offset = offset.saturating_add(batch_len);
if batch_len == 0 || u64::from(next_offset) >= page.pagination.total {
return Ok(connections);
}
offset = next_offset;
}
}
pub async fn get(&self, id: &str) -> Result<Connection> {
self.get_path(&format!("/v1/connections/{}", path_segment(id)))
.await
}
pub async fn delete(&self, id: &str, options: RequestOptions) -> Result<()> {
let request = self.client.request(
Method::DELETE,
&format!("/v1/connections/{}", path_segment(id)),
)?;
self.client.send(request, options).await
}
pub async fn validate(&self, id: &str, options: RequestOptions) -> Result<Connection> {
self.client
.post_empty(
&format!("/v1/connections/{}/validate", path_segment(id)),
options,
)
.await
}
pub async fn list_audit(
&self,
id: &str,
options: &ConnectionListOptions,
) -> Result<Vec<ConnectionAuditEntry>> {
let request = self
.client
.request(
Method::GET,
&format!("/v1/connections/{}/audit", path_segment(id)),
)?
.query(options);
self.client.send_collection(request, "audit").await
}
pub async fn list_labels(&self) -> Result<Vec<ConnectionLabelStat>> {
let request = self.client.request(Method::GET, "/v1/connections/labels")?;
self.client.send_collection(request, "labels").await
}
pub async fn update_labels(
&self,
id: &str,
input: &ConnectionLabelsUpdate,
options: RequestOptions,
) -> Result<Connection> {
self.patch(
&format!("/v1/connections/{}/labels", path_segment(id)),
input,
options,
)
.await
}
pub async fn update_status(
&self,
id: &str,
input: &ConnectionStatusUpdate,
options: RequestOptions,
) -> Result<Connection> {
self.patch(
&format!("/v1/connections/{}/status", path_segment(id)),
input,
options,
)
.await
}
pub async fn test(&self, input: &ConnectionRequest, options: RequestOptions) -> Result<bool> {
#[derive(serde::Deserialize)]
struct TestResult {
#[serde(alias = "ok")]
success: bool,
}
let request = self
.client
.request(Method::POST, "/v1/connections/test")?
.json(input);
let result: TestResult = self.client.send(request, options).await?;
Ok(result.success)
}
async fn get_path(&self, path: &str) -> Result<Connection> {
let request = self.client.request(Method::GET, path)?;
self.client.send(request, RequestOptions::default()).await
}
async fn patch<T: serde::Serialize + ?Sized>(
&self,
path: &str,
input: &T,
options: RequestOptions,
) -> Result<Connection> {
let request = self.client.request(Method::PATCH, path)?.json(input);
self.client.send(request, options).await
}
}