use std::collections::VecDeque;
use std::io::Read;
use std::sync::Arc;
use ibmcloud_iam::token::TokenManager;
use quick_xml::de::from_str;
use reqwest;
use serde;
use serde::{Deserialize, Serialize};
use tracing::error;
pub type Error = Box<dyn std::error::Error>;
#[derive(Deserialize, Serialize, Debug)]
pub struct ListAllMyBucketsResult {
#[serde(rename = "Owner")]
owner: Owner,
#[serde(rename = "Buckets")]
buckets: Buckets,
}
#[derive(Deserialize, Serialize, Debug)]
pub struct Buckets {
#[serde(rename = "Bucket")]
list: Vec<Bucket>,
}
#[derive(Deserialize, Serialize, Debug)]
pub struct Owner {
#[serde(rename = "$unflatten=ID")]
id: String,
#[serde(rename = "$unflatten=DisplayName")]
display_name: String,
}
#[derive(Deserialize, Serialize, Debug)]
pub struct Bucket {
#[serde(rename = "$unflatten=Name")]
pub name: String,
#[serde(rename = "$unflatten=CreationDate")]
pub creation_date: String,
}
fn default_contents() -> Vec<Contents> {
Vec::new()
}
#[derive(Deserialize, Serialize, Debug, PartialEq)]
pub struct ListBucketResult {
#[serde(rename = "Contents", default = "default_contents")]
contents: Vec<Contents>,
#[serde(rename = "$unflatten=KeyCount")]
key_count: u64,
#[serde(rename = "$unflatten=MaxKeys")]
max_keys: u64,
#[serde(rename = "$unflatten=NextContinuationToken")]
next_token: Option<String>,
}
#[derive(Deserialize, Serialize, Debug, PartialEq)]
pub struct Contents {
#[serde(rename = "$unflatten=Key")]
pub key: String,
#[serde(rename = "$unflatten=LastModified")]
pub last_modified: String,
#[serde(rename = "$unflatten=ETag")]
pub etag: String,
#[serde(rename = "$unflatten=Size")]
pub size: u64,
#[serde(rename = "$unflatten=StorageClass")]
pub storage_class: String,
}
pub struct Client {
pub(crate) tm: Arc<TokenManager>,
pub(crate) endpoint: String,
pub(crate) client: reqwest::blocking::Client,
}
impl Client {
pub fn new(tm: Arc<TokenManager>, endpoint: &str) -> Self {
Self {
tm: tm,
endpoint: endpoint.to_string(),
client: reqwest::blocking::Client::new(),
}
}
pub fn list_buckets(&self, instance_id: &str) -> Result<Vec<Bucket>, Error> {
let c = &self.client;
let url = format!("https://{}/", self.endpoint);
let response = c
.get(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.header("ibm-service-instance-id", instance_id.to_string())
.send()?;
let text: String = check_response(response)?.text()?;
let bucket_resp: ListAllMyBucketsResult = from_str(&text)?;
Ok(bucket_resp.buckets.list)
}
pub fn list_objects(
&self,
bucket: &str,
prefix: Option<String>,
start_after: Option<String>,
) -> ObjectIterator {
ObjectIterator::new(self, bucket, prefix.clone(), start_after.clone())
}
fn _list_objects(
&self,
bucket: &str,
prefix: &Option<String>,
continuation_token: &Option<String>,
start_after: &Option<String>,
) -> Result<ListBucketResult, Error> {
let c = &self.client;
let url = build_list_objects_url(
&self.endpoint,
bucket,
prefix,
continuation_token,
start_after,
)?;
let response = c
.get(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.send()?;
let text: String = check_response(response)?.text()?;
let objlist: ListBucketResult = from_str(&text)?;
Ok(objlist)
}
pub fn get_object_at_range(
&self,
bucket: &str,
key: &str,
start: u64,
end: Option<u64>,
) -> Result<Box<dyn Read>, Error> {
let c = &self.client;
let url = format!("https://{}.{}/{}", bucket, self.endpoint, key);
let mut end_str = "".to_string();
if let Some(e) = end {
end_str = format!("{}", e);
}
let response = c
.get(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.header("Range", format!("bytes={}-{}", start, end_str))
.send()?;
let r = check_response(response)?;
Ok(Box::new(r))
}
pub fn get_object(&self, bucket: &str, key: &str) -> Result<Box<dyn Read>, Error> {
let c = &self.client;
let url = format!("https://{}.{}/{}", bucket, self.endpoint, key);
let response = c
.get(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.send()?;
let r = check_response(response)?;
Ok(Box::new(r))
}
pub fn put_object<B: Into<reqwest::blocking::Body>>(
&self,
bucket: &str,
key: &str,
body: B,
) -> Result<(), Error> {
let c = &self.client;
let url = format!("https://{}.{}/{}", bucket, self.endpoint, key);
let response = c
.put(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.body(body)
.send()?;
let _r = check_response(response)?;
Ok(())
}
pub fn delete_object(&self, bucket: &str, key: &str) -> Result<(), Error> {
let c = &self.client;
let url = format!("https://{}.{}/{}", bucket, self.endpoint, key);
let response = c
.delete(url)
.header(
"Authorization",
format!("Bearer {}", self.tm.token()?.access_token),
)
.send()?;
check_response(response)?;
Ok(())
}
}
pub(crate) fn check_response(
response: reqwest::blocking::Response,
) -> Result<reqwest::blocking::Response, Error> {
if !response.status().is_success() {
return Err(format!(
"request failed: code='{}' body='{:?}'",
response.status(),
response.text().unwrap()
)
.into());
}
Ok(response)
}
pub struct ObjectIterator<'a> {
client: &'a Client,
bucket: String,
prefix: Option<String>,
continuation_token: Option<String>,
start_after: Option<String>,
results: VecDeque<Contents>,
complete: bool,
}
impl<'a> ObjectIterator<'a> {
pub fn new(
client: &'a Client,
bucket: &str,
prefix: Option<String>,
start_after: Option<String>,
) -> Self {
Self {
client,
bucket: bucket.to_string(),
prefix: prefix,
continuation_token: None,
start_after: start_after,
results: VecDeque::new(),
complete: false,
}
}
}
impl Iterator for ObjectIterator<'_> {
type Item = Contents;
fn next(&mut self) -> Option<Self::Item> {
if self.results.len() < 1 {
if self.complete {
return None;
}
match self.client._list_objects(
&self.bucket,
&self.prefix,
&self.continuation_token,
&self.start_after,
) {
Ok(mut v) => {
if v.contents.len() < 1 {
self.complete = true;
return None;
}
for o in v.contents.drain(..) {
self.results.push_back(o);
}
if v.next_token.is_some() {
self.continuation_token = v.next_token;
} else {
self.complete = true;
}
}
Err(e) => {
error!(e);
return None;
}
}
}
Some(self.results.pop_front().unwrap())
}
}
fn build_list_objects_url(
endpoint: &str,
bucket: &str,
prefix: &Option<String>,
continuation_token: &Option<String>,
start_after: &Option<String>,
) -> Result<reqwest::Url, Error> {
let mut url = reqwest::Url::parse(&format!("https://{}.{}/?list-type=2", bucket, endpoint))?;
if let Some(tok) = continuation_token {
url.query_pairs_mut().append_pair("continuation-token", tok);
}
if let Some(pre) = prefix {
url.query_pairs_mut().append_pair("prefix", pre);
}
if let Some(after) = start_after {
url.query_pairs_mut().append_pair("start-after", after);
}
Ok(url)
}
#[cfg(test)]
mod tests {
use super::*;
use quick_xml::se::to_string;
#[test]
fn test_bucket_list_response() {
let res = ListAllMyBucketsResult {
owner: Owner {
id: "asdfasdfa".to_string(),
display_name: "12315123".to_string(),
},
buckets: Buckets {
list: vec![
Bucket {
name: "asdasdfasdfadfadf".to_string(),
creation_date: "1238218238902389023890".to_string(),
},
Bucket {
name: "asdasdfasdfadfadf".to_string(),
creation_date: "1238218238902389023890".to_string(),
},
Bucket {
name: "asdasdfasdfadfadf".to_string(),
creation_date: "1238218238902389023890".to_string(),
},
],
},
};
let exp = "<ListAllMyBucketsResult><Owner><ID>asdfasdfa</ID><DisplayName>12315123</DisplayName></Owner><Buckets><Bucket><Name>asdasdfasdfadfadf</Name><CreationDate>1238218238902389023890</CreationDate></Bucket><Bucket><Name>asdasdfasdfadfadf</Name><CreationDate>1238218238902389023890</CreationDate></Bucket><Bucket><Name>asdasdfasdfadfadf</Name><CreationDate>1238218238902389023890</CreationDate></Bucket></Buckets></ListAllMyBucketsResult>";
let out = to_string(&res).unwrap();
assert_eq!(out, exp);
}
#[test]
fn test_list_objects_empty_bucket() {
let input = r#"<?xml version="1.0" encoding="UTF-8" standalone="yes"?><ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Name>logbase</Name><Prefix></Prefix><KeyCount>0</KeyCount><MaxKeys>1000</MaxKeys><Delimiter></Delimiter><IsTruncated>false</IsTruncated></ListBucketResult>"#;
let exp = ListBucketResult {
contents: vec![],
key_count: 0,
max_keys: 1000,
next_token: None,
};
let objs: ListBucketResult = from_str(&input).unwrap();
assert_eq!(objs, exp);
}
#[test]
fn test_build_list_objects_url() {
let res = build_list_objects_url(
"cos.cloud.ibm.com",
"test-bucket-123",
&None,
&None,
&Some("object-key/with/special=characters+001.stuff".to_string()),
);
let mut url = reqwest::Url::parse("https://test-bucket-123.cos.cloud.ibm.com/").unwrap();
url.query_pairs_mut()
.append_pair("list-type", "2")
.append_pair(
"start-after",
"object-key/with/special=characters+001.stuff",
);
assert_eq!(res.unwrap(), url);
}
}