use reqwest::Client;
use reqwest::multipart::{Form, Part};
use serde::{Serialize, Deserialize};
use thiserror::Error;
use std::collections::HashMap;
use url::Url;
use bytes::Bytes;
use std::path::Path;
use tokio::fs::File;
use tokio::io::AsyncReadExt;
#[cfg(test)]
use wiremock::{MockServer, Mock, ResponseTemplate};
#[cfg(test)]
use wiremock::matchers::{method, path, query_param};
#[cfg(test)]
use serde_json::json;
pub type Result<T> = std::result::Result<T, StorageError>;
#[derive(Error, Debug)]
pub enum StorageError {
#[error("API error: {0}")]
ApiError(String),
#[error("Network error: {0}")]
NetworkError(#[from] reqwest::Error),
#[error("JSON serialization error: {0}")]
SerializationError(#[from] serde_json::Error),
#[error("URL parse error: {0}")]
UrlParseError(#[from] url::ParseError),
#[error("Storage error: {0}")]
StorageError(String),
#[error("File not found: {0}")]
FileNotFound(String),
#[error("IO error: {0}")]
IoError(#[from] std::io::Error),
#[error("Request error: {0}")]
RequestError(String),
#[error("Deserialization error: {0}")]
DeserializationError(String),
}
impl StorageError {
pub fn new(message: String) -> Self {
Self::StorageError(message)
}
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct FileOptions {
pub cache_control: Option<String>,
pub content_type: Option<String>,
pub upsert: Option<bool>,
}
impl FileOptions {
pub fn new() -> Self {
Self::default()
}
pub fn with_cache_control(mut self, cache_control: &str) -> Self {
self.cache_control = Some(cache_control.to_string());
self
}
pub fn with_content_type(mut self, content_type: &str) -> Self {
self.content_type = Some(content_type.to_string());
self
}
pub fn with_upsert(mut self, upsert: bool) -> Self {
self.upsert = Some(upsert);
self
}
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct ListOptions {
pub limit: Option<i32>,
pub offset: Option<i32>,
pub sort_by: Option<SortBy>,
pub search: Option<String>,
}
impl ListOptions {
pub fn new() -> Self {
Self::default()
}
pub fn limit(mut self, limit: i32) -> Self {
self.limit = Some(limit);
self
}
pub fn offset(mut self, offset: i32) -> Self {
self.offset = Some(offset);
self
}
pub fn sort_by(mut self, column: &str, order: SortOrder) -> Self {
self.sort_by = Some(SortBy {
column: column.to_string(),
order,
});
self
}
pub fn search(mut self, search: &str) -> Self {
self.search = Some(search.to_string());
self
}
}
#[derive(Debug, Clone, Serialize)]
pub struct SortBy {
pub column: String,
pub order: SortOrder,
}
impl ToString for SortBy {
fn to_string(&self) -> String {
format!("{}:{:?}", self.column, self.order).to_lowercase()
}
}
#[derive(Debug, Clone, Copy, Serialize)]
#[serde(rename_all = "lowercase")]
pub enum SortOrder {
Asc,
Desc,
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct ImageTransformOptions {
pub width: Option<u32>,
pub height: Option<u32>,
pub resize: Option<String>,
pub format: Option<String>,
pub quality: Option<u32>,
}
impl ImageTransformOptions {
pub fn new() -> Self {
Self::default()
}
pub fn with_width(mut self, width: u32) -> Self {
self.width = Some(width);
self
}
pub fn with_height(mut self, height: u32) -> Self {
self.height = Some(height);
self
}
pub fn with_resize(mut self, resize: &str) -> Self {
self.resize = Some(resize.to_string());
self
}
pub fn with_format(mut self, format: &str) -> Self {
self.format = Some(format.to_string());
self
}
pub fn with_quality(mut self, quality: u32) -> Self {
self.quality = Some(quality.min(100));
self
}
fn to_query_params(&self) -> String {
let mut params = Vec::new();
if let Some(width) = self.width {
params.push(format!("width={}", width));
}
if let Some(height) = self.height {
params.push(format!("height={}", height));
}
if let Some(resize) = &self.resize {
params.push(format!("resize={}", resize));
}
if let Some(format) = &self.format {
params.push(format!("format={}", format));
}
if let Some(quality) = self.quality {
params.push(format!("quality={}", quality));
}
if params.is_empty() {
String::new()
} else {
format!("?{}", params.join("&"))
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FileObject {
pub name: String,
pub bucket_id: String,
pub owner: String,
pub id: String,
pub updated_at: String,
pub created_at: String,
pub last_accessed_at: String,
pub metadata: Option<serde_json::Value>,
pub mime_type: Option<String>,
pub size: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Bucket {
pub id: String,
pub name: String,
pub owner: String,
pub public: bool,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Deserialize)]
pub struct InitiateMultipartUploadResponse {
pub id: String,
#[serde(rename = "uploadId")]
pub upload_id: String,
pub key: String,
pub bucket: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct UploadedPartInfo {
#[serde(rename = "partNumber")]
pub part_number: u32,
pub etag: String,
}
#[derive(Debug, Clone, Serialize)]
struct CompleteMultipartUploadRequest {
#[serde(rename = "uploadId")]
pub upload_id: String,
pub parts: Vec<UploadedPartInfo>,
}
pub struct StorageBucketClient<'a> {
parent: &'a StorageClient,
bucket_id: String,
}
pub struct StorageClient {
base_url: String,
api_key: String,
http_client: Client,
}
impl StorageClient {
pub fn new(base_url: &str, api_key: &str, http_client: Client) -> Self {
Self {
base_url: base_url.to_string(),
api_key: api_key.to_string(),
http_client,
}
}
pub fn from<'a>(&'a self, bucket_id: &str) -> StorageBucketClient<'a> {
StorageBucketClient {
parent: self,
bucket_id: bucket_id.to_string(),
}
}
pub async fn list_buckets(&self) -> Result<Vec<Bucket>> {
let url = format!("{}/storage/v1/bucket", self.base_url);
let response = self.http_client.get(&url)
.header("apikey", &self.api_key)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let buckets = response.json::<Vec<Bucket>>().await?;
Ok(buckets)
}
pub async fn create_bucket(&self, bucket_id: &str, is_public: bool) -> Result<Bucket> {
let url = format!("{}/storage/v1/bucket", self.base_url);
let payload = serde_json::json!({
"id": bucket_id,
"name": bucket_id,
"public": is_public
});
let response = self.http_client.post(&url)
.header("apikey", &self.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let bucket = response.json::<Bucket>().await?;
Ok(bucket)
}
pub async fn delete_bucket(&self, bucket_id: &str) -> Result<()> {
let url = format!("{}/storage/v1/bucket/{}", self.base_url, bucket_id);
let response = self.http_client.delete(&url)
.header("apikey", &self.api_key)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn update_bucket(&self, bucket_id: &str, is_public: bool) -> Result<Bucket> {
let url = format!("{}/storage/v1/bucket/{}", self.base_url, bucket_id);
let payload = serde_json::json!({
"id": bucket_id,
"public": is_public
});
let response = self.http_client.put(&url)
.header("apikey", &self.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let bucket = response.json::<Bucket>().await?;
Ok(bucket)
}
}
impl<'a> StorageBucketClient<'a> {
pub async fn upload(&self, path: &str, file_path: &Path, options: Option<FileOptions>) -> Result<FileObject> {
let mut url = Url::parse(&self.parent.base_url)?;
url.set_path(&format!("/storage/v1/object/{}/{}", self.bucket_id, path));
if let Some(opts) = &options {
let mut query_pairs = url.query_pairs_mut();
if let Some(cache_control) = &opts.cache_control {
query_pairs.append_pair("cache_control", cache_control);
}
if let Some(upsert) = &opts.upsert {
query_pairs.append_pair("upsert", &upsert.to_string());
}
}
let mut file = File::open(file_path).await?;
let mut contents = Vec::new();
file.read_to_end(&mut contents).await?;
let part = Part::bytes(contents)
.file_name(file_path.file_name().unwrap().to_string_lossy().to_string());
let form = Form::new().part("file", part);
let response = self.parent.http_client
.post(url)
.header("apikey", &self.parent.api_key)
.header("Authorization", format!("Bearer {}", &self.parent.api_key))
.multipart(form)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let file_object = response.json::<FileObject>().await?;
Ok(file_object)
}
pub async fn download(&self, path: &str) -> Result<Bytes> {
let mut url = Url::parse(&self.parent.base_url)?;
url.set_path(&format!("/storage/v1/object/{}/{}", self.bucket_id, path));
let response = self.parent.http_client
.get(url)
.header("apikey", &self.parent.api_key)
.header("Authorization", format!("Bearer {}", &self.parent.api_key))
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let bytes = response.bytes().await?;
Ok(bytes)
}
pub async fn list(&self, prefix: &str, options: Option<ListOptions>) -> Result<Vec<FileObject>> {
let mut url = Url::parse(&self.parent.base_url)?;
url.set_path(&format!("/storage/v1/object/list/{}", self.bucket_id));
{
let mut query_pairs = url.query_pairs_mut();
query_pairs.append_pair("prefix", prefix);
if let Some(opts) = &options {
if let Some(limit) = opts.limit {
query_pairs.append_pair("limit", &limit.to_string());
}
if let Some(offset) = opts.offset {
query_pairs.append_pair("offset", &offset.to_string());
}
if let Some(sort_by) = &opts.sort_by {
query_pairs.append_pair("sortBy", &sort_by.to_string());
}
if let Some(search) = &opts.search {
query_pairs.append_pair("search", search);
}
}
}
let response = self.parent.http_client
.get(url)
.header("apikey", &self.parent.api_key)
.header("Authorization", format!("Bearer {}", &self.parent.api_key))
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let files = response.json::<Vec<FileObject>>().await?;
Ok(files)
}
pub async fn remove(&self, paths: Vec<&str>) -> Result<()> {
let url = format!("{}/storage/v1/object/{}", self.parent.base_url, self.bucket_id);
let payload = serde_json::json!({
"prefixes": paths
});
let response = self.parent.http_client.delete(&url)
.header("apikey", &self.parent.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub fn get_public_url(&self, path: &str) -> String {
format!("{}/storage/v1/object/public/{}/{}", self.parent.base_url, self.bucket_id, path)
}
pub async fn create_signed_url(&self, path: &str, expires_in: i32) -> Result<String> {
let url = format!("{}/storage/v1/object/sign/{}/{}", self.parent.base_url, self.bucket_id, path);
let payload = serde_json::json!({
"expiresIn": expires_in
});
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
#[derive(Deserialize)]
struct SignedUrlResponse {
signed_url: String,
}
let signed_url = response.json::<SignedUrlResponse>().await?;
Ok(signed_url.signed_url)
}
pub async fn initiate_multipart_upload(
&self,
path: &str,
options: Option<FileOptions>
) -> Result<InitiateMultipartUploadResponse> {
let url = format!(
"{}/storage/v1/upload/initiate",
self.parent.base_url
);
let options = options.unwrap_or_default();
let cache_control = options.cache_control.unwrap_or_else(|| "max-age=3600".to_string());
let content_type = options.content_type.unwrap_or_else(|| "application/octet-stream".to_string());
let upsert = options.upsert.unwrap_or(false);
let payload = serde_json::json!({
"bucket": self.bucket_id,
"name": path,
"cacheControl": cache_control,
"contentType": content_type,
"upsert": upsert,
});
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let initiate_response: InitiateMultipartUploadResponse = response.json().await?;
Ok(initiate_response)
}
pub async fn upload_part(
&self,
upload_id: &str,
part_number: u32,
data: Bytes
) -> Result<UploadedPartInfo> {
let url = format!(
"{}/storage/v1/upload/part",
self.parent.base_url
);
let body = reqwest::Body::from(data);
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.query(&[
("uploadId", upload_id),
("partNumber", &part_number.to_string()),
("bucket", &self.bucket_id),
])
.body(body)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let etag = response.headers()
.get("etag")
.ok_or_else(|| StorageError::new("ETag header not found in response".to_string()))?
.to_str()
.map_err(|e| StorageError::new(format!("Invalid ETag header: {}", e)))?
.to_string();
let part_info = UploadedPartInfo {
part_number,
etag,
};
Ok(part_info)
}
pub async fn complete_multipart_upload(
&self,
upload_id: &str,
path: &str,
parts: Vec<UploadedPartInfo>
) -> Result<FileObject> {
let url = format!(
"{}/storage/v1/upload/complete",
self.parent.base_url
);
let request_data = CompleteMultipartUploadRequest {
upload_id: upload_id.to_string(),
parts,
};
let params = [
("bucket", self.bucket_id.to_string()),
("key", path.to_string()), ];
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.header("Content-Type", "application/json")
.json(&request_data)
.query(¶ms)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let file_object: FileObject = response.json().await?;
Ok(file_object)
}
pub async fn abort_multipart_upload(
&self,
upload_id: &str,
path: &str
) -> Result<()> {
let url = format!(
"{}/storage/v1/upload/abort",
self.parent.base_url
);
let payload = serde_json::json!({
"uploadId": upload_id,
"bucket": self.bucket_id,
"key": path,
});
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.header("Content-Type", "application/json")
.json(&payload)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn upload_large_file(
&self,
path: &str,
file_path: &Path,
chunk_size: usize,
options: Option<FileOptions>
) -> Result<FileObject> {
let mut file = File::open(file_path).await?;
let file_size = file.metadata().await?.len() as usize;
let chunk_count = (file_size + chunk_size - 1) / chunk_size;
if chunk_count == 0 {
return Err(StorageError::new("File is empty".to_string()));
}
let init_response = self.initiate_multipart_upload(path, options).await?;
let mut uploaded_parts = Vec::with_capacity(chunk_count);
let mut buffer = vec![0u8; chunk_size];
for part_number in 1..=chunk_count as u32 {
let n = file.read(&mut buffer).await?;
if n == 0 {
break;
}
let chunk_data = Bytes::from(buffer[0..n].to_vec());
let part_info = self.upload_part(&init_response.upload_id, part_number, chunk_data).await?;
uploaded_parts.push(part_info);
}
let file_object = self.complete_multipart_upload(
&init_response.upload_id,
path,
uploaded_parts
).await?;
Ok(file_object)
}
pub async fn transform_image(&self, path: &str, options: ImageTransformOptions) -> Result<Bytes> {
let url = format!(
"{}/storage/v1/object/transform/{}/{}{}",
self.parent.base_url,
self.bucket_id,
path,
options.to_query_params()
);
let response = self.parent.http_client.get(&url)
.header("apikey", &self.parent.api_key)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
let data = response.bytes().await?;
Ok(data)
}
pub fn get_public_transform_url(&self, path: &str, options: ImageTransformOptions) -> String {
format!(
"{}/storage/v1/object/public/transform/{}/{}{}",
self.parent.base_url,
self.bucket_id,
path,
options.to_query_params()
)
}
pub async fn create_signed_transform_url(
&self,
path: &str,
options: ImageTransformOptions,
expires_in: i32
) -> Result<String> {
let transform_path = format!("transform/{}/{}{}", self.bucket_id, path, options.to_query_params());
let url = format!(
"{}/storage/v1/object/sign/{}",
self.parent.base_url,
transform_path,
);
let params = serde_json::json!({
"expiresIn": expires_in
});
let response = self.parent.http_client.post(&url)
.header("apikey", &self.parent.api_key)
.json(¶ms)
.send()
.await?;
if !response.status().is_success() {
let error_text = response.text().await?;
return Err(StorageError::ApiError(error_text));
}
#[derive(Deserialize)]
struct SignedUrlResponse {
signed_url: String,
}
let result = response.json::<SignedUrlResponse>().await?;
Ok(result.signed_url)
}
pub fn s3_compatible(&self, options: s3::S3Options) -> s3::S3BucketClient {
s3::S3BucketClient::new(
&self.parent.base_url,
&self.parent.api_key,
&self.bucket_id,
self.parent.http_client.clone(),
options
)
}
}
pub mod s3 {
use crate::StorageError;
use crate::Result;
use reqwest::Client;
use serde_json::Value;
use std::collections::HashMap;
use bytes::Bytes;
use serde::{Serialize, Deserialize};
use serde_json;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct S3Options {
pub access_key_id: String,
pub secret_access_key: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub region: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub endpoint: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub force_path_style: Option<bool>,
}
impl Default for S3Options {
fn default() -> Self {
Self {
access_key_id: String::new(),
secret_access_key: String::new(),
region: Some("auto".to_string()),
endpoint: None,
force_path_style: Some(true),
}
}
}
pub struct S3Client {
pub options: S3Options,
pub base_url: String,
pub api_key: String,
pub http_client: Client,
}
impl S3Client {
pub fn new(base_url: &str, api_key: &str, http_client: Client, options: S3Options) -> Self {
Self {
options,
base_url: base_url.to_string(),
api_key: api_key.to_string(),
http_client,
}
}
pub async fn create_bucket(&self, bucket_name: &str, is_public: bool) -> Result<()> {
let url = format!("{}/storage/v1/bucket", self.base_url);
let payload = serde_json::json!({
"name": bucket_name,
"public": is_public,
"file_size_limit": null,
"allowed_mime_types": null
});
let response = self.http_client
.post(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.json(&payload)
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn delete_bucket(&self, bucket_name: &str) -> Result<()> {
let url = format!("{}/storage/v1/bucket/{}", self.base_url, bucket_name);
let response = self.http_client
.delete(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn list_buckets(&self) -> Result<Vec<serde_json::Value>> {
let url = format!("{}/storage/v1/bucket", self.base_url);
let response = self.http_client
.get(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
let buckets = response.json::<Vec<serde_json::Value>>().await
.map_err(|e| StorageError::DeserializationError(e.to_string()))?;
Ok(buckets)
}
pub fn bucket(&self, bucket_name: &str) -> S3BucketClient {
S3BucketClient::new(
&self.base_url,
&self.api_key,
bucket_name,
self.http_client.clone(),
self.options.clone()
)
}
}
pub struct S3BucketClient {
pub base_url: String,
pub api_key: String,
pub bucket_name: String,
pub http_client: Client,
pub options: S3Options,
}
impl S3BucketClient {
pub fn new(
base_url: &str,
api_key: &str,
bucket_name: &str,
http_client: Client,
options: S3Options
) -> Self {
Self {
base_url: base_url.to_string(),
api_key: api_key.to_string(),
bucket_name: bucket_name.to_string(),
http_client,
options,
}
}
pub async fn put_object(&self, path: &str, data: Bytes, content_type: Option<String>, metadata: Option<HashMap<String, String>>) -> Result<()> {
let url = format!(
"{}/storage/v1/object/{}/{}",
self.base_url,
self.bucket_name,
path.trim_start_matches('/')
);
let content_type = content_type.unwrap_or_else(|| "application/octet-stream".to_string());
let mut request = self.http_client
.put(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.header("Content-Type", content_type)
.body(data);
if let Some(metadata) = metadata {
for (key, value) in metadata {
request = request.header(&format!("x-amz-meta-{}", key), value);
}
}
let response = request
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn get_object(&self, path: &str) -> Result<Bytes> {
let url = format!(
"{}/storage/v1/object/{}/{}",
self.base_url,
self.bucket_name,
path.trim_start_matches('/')
);
let response = self.http_client
.get(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
let data = response.bytes().await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
Ok(data)
}
pub async fn head_object(&self, path: &str) -> Result<HashMap<String, String>> {
let url = format!(
"{}/storage/v1/object/{}/{}",
self.base_url,
self.bucket_name,
path.trim_start_matches('/')
);
let response = self.http_client
.head(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
return Err(StorageError::ApiError("Object not found".to_string()));
}
let mut metadata = HashMap::new();
for (key, value) in response.headers() {
let key_str = key.to_string();
if key_str.starts_with("x-amz-meta-") {
let meta_key = key_str.trim_start_matches("x-amz-meta-").to_string();
metadata.insert(meta_key, value.to_str().unwrap_or_default().to_string());
}
}
Ok(metadata)
}
pub async fn delete_object(&self, path: &str) -> Result<()> {
let url = format!(
"{}/storage/v1/object/{}/{}",
self.base_url,
self.bucket_name,
path.trim_start_matches('/')
);
let response = self.http_client
.delete(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
pub async fn list_objects(&self, prefix: Option<&str>, delimiter: Option<&str>, max_keys: Option<i32>) -> Result<serde_json::Value> {
let mut url = format!(
"{}/storage/v1/object/list/{}",
self.base_url,
self.bucket_name
);
let mut query_params = Vec::new();
if let Some(prefix) = prefix {
query_params.push(format!("prefix={}", prefix));
}
if let Some(delimiter) = delimiter {
query_params.push(format!("delimiter={}", delimiter));
}
if let Some(max_keys) = max_keys {
query_params.push(format!("max-keys={}", max_keys));
}
if !query_params.is_empty() {
url = format!("{}?{}", url, query_params.join("&"));
}
let response = self.http_client
.get(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
let objects = response.json::<serde_json::Value>().await
.map_err(|e| StorageError::DeserializationError(e.to_string()))?;
Ok(objects)
}
pub async fn copy_object(&self, source_path: &str, destination_path: &str) -> Result<()> {
let url = format!(
"{}/storage/v1/object/copy",
self.base_url
);
let payload = serde_json::json!({
"bucketId": self.bucket_name,
"sourceKey": source_path,
"destinationKey": destination_path
});
let response = self.http_client
.post(&url)
.header("apikey", &self.api_key)
.header("Authorization", format!("Bearer {}", &self.api_key))
.json(&payload)
.send()
.await
.map_err(|e| StorageError::RequestError(e.to_string()))?;
if !response.status().is_success() {
let error_text = response.text().await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(StorageError::ApiError(error_text));
}
Ok(())
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_list_buckets() {
}
#[tokio::test]
async fn test_multipart_upload() {
}
#[tokio::test]
async fn test_transform_image() {
let mock_server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/storage/v1/object/transform/test-bucket/image.jpg"))
.and(query_param("width", "300"))
.and(query_param("height", "200"))
.and(query_param("format", "webp"))
.respond_with(ResponseTemplate::new(200)
.set_body_bytes(vec![0, 1, 2, 3, 4]) .append_header("Content-Type", "image/webp")
)
.mount(&mock_server)
.await;
let http_client = reqwest::Client::new();
let storage_client = StorageClient::new(&mock_server.uri(), "fake-key", http_client);
let bucket_client = storage_client.from("test-bucket");
let options = ImageTransformOptions::new()
.with_width(300)
.with_height(200)
.with_format("webp");
let result = bucket_client.transform_image("image.jpg", options).await;
assert!(result.is_ok());
let image_data = result.unwrap();
assert_eq!(image_data.len(), 5);
}
#[tokio::test]
async fn test_get_public_transform_url() {
let http_client = reqwest::Client::new();
let storage_client = StorageClient::new("https://example.com", "fake-key", http_client);
let bucket_client = storage_client.from("test-bucket");
let options = ImageTransformOptions::new()
.with_width(300)
.with_height(200)
.with_format("webp");
let url = bucket_client.get_public_transform_url("image.jpg", options);
assert!(url.contains("https://example.com/storage/v1/object/public/transform/test-bucket/image.jpg"));
assert!(url.contains("width=300"));
assert!(url.contains("height=200"));
assert!(url.contains("format=webp"));
}
#[tokio::test]
async fn test_create_signed_transform_url() {
let mock_server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/storage/v1/object/sign/transform/test-bucket/image.jpg"))
.respond_with(ResponseTemplate::new(200)
.set_body_json(json!({
"signed_url": "https://example.com/storage/v1/object/signed/transform/test-bucket/image.jpg?width=300&height=200&format=webp&token=abc123"
}))
)
.mount(&mock_server)
.await;
let http_client = reqwest::Client::new();
let storage_client = StorageClient::new(&mock_server.uri(), "fake-key", http_client);
let bucket_client = storage_client.from("test-bucket");
let options = ImageTransformOptions::new()
.with_width(300)
.with_height(200)
.with_format("webp");
let result = bucket_client.create_signed_transform_url("image.jpg", options, 60).await;
assert!(result.is_ok());
let url = result.unwrap();
assert!(url.contains("https://example.com/storage/v1/object/signed/transform/test-bucket/image.jpg"));
assert!(url.contains("token=abc123"));
}
}