use crate::traits::{StorageBackend, StorageError, StorageItem};
use crate::providers::OAuthCredentials;
use crate::providers::utils::translate_http_error;
use async_trait::async_trait;
use std::path::{Path, PathBuf};
use tracing::info;
pub struct GoogleDriveProvider {
client: reqwest::Client,
credentials: OAuthCredentials,
api_url: String,
upload_url: String,
token_manager: std::sync::Arc<super::utils::OAuthTokenManager>,
upload_limiter: Option<crate::rate_limit::TokenBucket>,
download_limiter: Option<crate::rate_limit::TokenBucket>,
}
crate::impl_provider_builder!(GoogleDriveProvider, GoogleDriveProviderBuilder, OAuthCredentials);
crate::impl_oauth_token_helper!(GoogleDriveProvider);
impl GoogleDriveProvider {
pub fn with_client_options(
credentials: OAuthCredentials,
timeout: Option<std::time::Duration>,
custom_headers: Option<reqwest::header::HeaderMap>,
) -> Self {
let client = super::utils::build_http_client(timeout, custom_headers);
let auth_url = "https://oauth2.googleapis.com/token".to_string();
let token_manager = std::sync::Arc::new(super::utils::OAuthTokenManager::new(
client.clone(),
&auth_url,
&credentials.client_id,
&credentials.client_secret,
&credentials.refresh_token,
"Google Drive",
));
Self {
client,
credentials,
api_url: "https://www.googleapis.com/drive/v3/files".to_string(),
upload_url: "https://www.googleapis.com/upload/drive/v3/files".to_string(),
token_manager,
upload_limiter: None,
download_limiter: None,
}
}
pub fn with_limiters(
mut self,
upload_limiter: Option<crate::rate_limit::TokenBucket>,
download_limiter: Option<crate::rate_limit::TokenBucket>,
) -> Self {
if self.upload_limiter.is_none() {
self.upload_limiter = upload_limiter;
}
if self.download_limiter.is_none() {
self.download_limiter = download_limiter;
}
self
}
#[cfg(test)]
pub fn with_endpoints(mut self, auth_url: String, api_url: String, upload_url: String) -> Self {
self.api_url = api_url;
self.upload_url = upload_url;
self.token_manager = std::sync::Arc::new(super::utils::OAuthTokenManager::new(
self.client.clone(),
&auth_url,
&self.credentials.client_id,
&self.credentials.client_secret,
&self.credentials.refresh_token,
"Google Drive",
));
self
}
async fn get_or_create_folder_id(&self, token: &str, parent_id: &str, name: &str) -> Result<String, StorageError> {
let query = format!(
"name = '{}' and '{}' in parents and mimeType = 'application/vnd.google-apps.folder' and trashed = false",
name.replace('\'', "\\'"),
parent_id
);
let res = super::utils::apply_bearer_auth(self.client.get(&self.api_url), token)
.query(&[("q", &query), ("fields", &"files(id)".to_string())])
.send()
.await?
.json::<serde_json::Value>()
.await?;
if let Some(files) = res["files"].as_array() {
if !files.is_empty() {
return Ok(files[0]["id"].as_str().unwrap().to_string());
}
}
let body = serde_json::json!({
"name": name,
"parents": [parent_id],
"mimeType": "application/vnd.google-apps.folder"
});
let create_res = super::utils::apply_bearer_auth(self.client.post(&self.api_url), token)
.json(&body)
.send()
.await?
.json::<serde_json::Value>()
.await?;
let id = create_res["id"].as_str()
.ok_or_else(|| StorageError::Provider { message: format!("Failed to create folder '{}' in Google Drive: {:?}", name, create_res), status: None })?
.to_string();
Ok(id)
}
async fn get_or_create_file_id(&self, token: &str, path: &str, is_folder: bool) -> Result<String, StorageError> {
let normalized = super::utils::normalize_remote_path(path);
let parts: Vec<&str> = normalized.split('/').filter(|s| !s.is_empty()).collect();
let mut parent_id = "root".to_string();
if let Some(ref dest_folder) = self.credentials.common.destination_folder {
let normalized_dest = super::utils::normalize_remote_path(dest_folder);
if !normalized_dest.is_empty() {
for seg in normalized_dest.split('/').filter(|s| !s.is_empty()) {
parent_id = self.get_or_create_folder_id(token, &parent_id, seg).await?;
}
}
}
for (i, part) in parts.iter().enumerate() {
let is_last = i == parts.len() - 1;
let current_is_folder = !is_last || is_folder;
if current_is_folder {
parent_id = self.get_or_create_folder_id(token, &parent_id, part).await?;
} else {
let query = format!(
"name = '{}' and '{}' in parents and mimeType != 'application/vnd.google-apps.folder' and trashed = false",
part.replace('\'', "\\'"),
parent_id
);
let res = super::utils::apply_bearer_auth(self.client.get(&self.api_url), token)
.query(&[("q", &query), ("fields", &"files(id)".to_string())])
.send()
.await?
.json::<serde_json::Value>()
.await?;
if let Some(files) = res["files"].as_array() {
if !files.is_empty() {
parent_id = files[0]["id"].as_str().unwrap().to_string();
continue;
}
}
let body = serde_json::json!({
"name": part,
"parents": [parent_id],
"mimeType": "application/octet-stream"
});
let create_res = super::utils::apply_bearer_auth(self.client.post(&self.api_url), token)
.json(&body)
.send()
.await?
.json::<serde_json::Value>()
.await?;
parent_id = create_res["id"].as_str()
.ok_or_else(|| StorageError::Provider { message: format!("Failed to create file '{}' in Google Drive: {:?}", part, create_res), status: None })?
.to_string();
}
}
Ok(parent_id)
}
}
#[async_trait]
impl StorageBackend for GoogleDriveProvider {
fn name(&self) -> &str {
"Google Drive"
}
async fn upload(&self, local_path: &Path, remote_path: &str) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "upload", || async {
let token = self.get_access_token().await?;
let file_id = self.get_or_create_file_id(&token, remote_path, false).await?;
info!("[{}] Real upload starting for '{}' (ID: {})", self.name(), remote_path, file_id);
let (body, size) = super::utils::get_upload_body(local_path, self.upload_limiter.clone()).await?;
let upload_url = format!("{}/{}?uploadType=media", self.upload_url, file_id);
let res = super::utils::apply_bearer_auth(self.client.patch(&upload_url), &token)
.header("Content-Type", "application/octet-stream")
.header("Content-Length", size.to_string())
.body(body)
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "upload").await);
}
Ok(())
}).await
}
async fn download(&self, remote_path: &str, local_path: &Path) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "download", || async {
let token = self.get_access_token().await?;
let file_id = self.get_or_create_file_id(&token, remote_path, false).await?;
let download_url = format!("{}/{}?alt=media", self.api_url, file_id);
let res = super::utils::apply_bearer_auth(self.client.get(&download_url), &token)
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "download").await);
}
super::utils::download_rate_limited(res, local_path, self.download_limiter.clone()).await?;
Ok(())
}).await
}
async fn delete(&self, remote_path: &str) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "delete", || async {
let token = self.get_access_token().await?;
let file_id = self.get_or_create_file_id(&token, remote_path, false).await?;
let delete_url = format!("{}/{}", self.api_url, file_id);
let res = super::utils::apply_bearer_auth(self.client.delete(&delete_url), &token)
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "delete").await);
}
Ok(())
}).await
}
async fn create_folder(&self, remote_path: &str) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "create_folder", || async {
let token = self.get_access_token().await?;
let _ = self.get_or_create_file_id(&token, remote_path, true).await?;
Ok(())
}).await
}
async fn list(&self, remote_path: &str) -> Result<Vec<StorageItem>, StorageError> {
super::utils::execute_with_retry(self.name(), "list", || async {
let token = self.get_access_token().await?;
let folder_id = self.get_or_create_file_id(&token, remote_path, true).await?;
let query = format!("'{}' in parents and trashed = false", folder_id);
let mut items = Vec::new();
let mut next_page_token: Option<String> = None;
loop {
let req = super::utils::apply_bearer_auth(self.client.get(&self.api_url), &token);
let fields = "nextPageToken, files(id, name, size, mimeType, modifiedTime, md5Checksum)".to_string();
let mut query_params = vec![
("q", query.clone()),
("fields", fields),
];
let page_token_str;
if let Some(ref page_token) = next_page_token {
page_token_str = page_token.clone();
query_params.push(("pageToken", page_token_str));
}
let res = req.query(&query_params)
.send()
.await?
.json::<serde_json::Value>()
.await?;
if let Some(files) = res["files"].as_array() {
for file in files {
let name = file["name"].as_str().unwrap_or("").to_string();
let size = file["size"].as_str().unwrap_or("0").parse::<u64>().unwrap_or(0);
let mime_type = file["mimeType"].as_str().unwrap_or("");
let is_dir = mime_type == "application/vnd.google-apps.folder";
if mime_type.starts_with("application/vnd.google-apps.") && !is_dir {
continue;
}
let checksum = file["md5Checksum"].as_str().map(|s| s.to_string());
let modified = file["modifiedTime"].as_str()
.and_then(|t| time::OffsetDateTime::parse(t, &time::format_description::well_known::Rfc3339).ok())
.map(std::time::SystemTime::from)
.unwrap_or_else(std::time::SystemTime::now);
let rel_path = if remote_path.is_empty() {
name
} else {
format!("{}/{}", remote_path, name)
};
items.push(StorageItem {
path: PathBuf::from(rel_path),
size,
modified,
is_dir,
checksum,
permissions: None,
});
}
}
next_page_token = res["nextPageToken"].as_str().map(|s| s.to_string());
if next_page_token.is_none() {
break;
}
}
Ok(items)
}).await
}
async fn compute_local_checksum(&self, local_path: &Path) -> Result<Option<String>, StorageError> {
Ok(crate::checksum::compute_md5(local_path).await.ok())
}
}
pub struct GoogleDriveProviderBuilder {
pub credentials: OAuthCredentials,
pub timeout: Option<std::time::Duration>,
pub custom_headers: Option<reqwest::header::HeaderMap>,
}
impl GoogleDriveProviderBuilder {
pub fn new(credentials: OAuthCredentials) -> Self {
Self {
credentials,
timeout: None,
custom_headers: None,
}
}
pub fn build(self) -> GoogleDriveProvider {
GoogleDriveProvider::with_client_options(self.credentials, self.timeout, self.custom_headers)
}
}