use crate::traits::{StorageBackend, StorageError, StorageItem};
use crate::providers::IPFSCredentials;
use crate::providers::utils::translate_http_error;
use async_trait::async_trait;
use std::path::Path;
use std::time::SystemTime;
use tokio::fs;
use tracing::info;
use serde::Deserialize;
#[derive(Deserialize, Debug)]
struct PinataMetadata {
name: String,
}
#[derive(Deserialize, Debug)]
struct PinataRow {
ipfs_pin_hash: String,
metadata: PinataMetadata,
size: u64,
date_pinned: Option<String>,
}
#[derive(Deserialize, Debug)]
struct PinataPinListResponse {
rows: Vec<PinataRow>,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct PinataPinFileResponse {
#[serde(rename = "IpfsHash")]
ipfs_hash: String,
}
pub struct IPFSProvider {
client: reqwest::Client,
credentials: IPFSCredentials,
api_url: String,
gateway_url: String,
}
crate::impl_provider_builder!(IPFSProvider, IPFSProviderBuilder, IPFSCredentials, relative);
impl IPFSProvider {
pub fn with_client_options(
credentials: IPFSCredentials,
timeout: Option<std::time::Duration>,
custom_headers: Option<reqwest::header::HeaderMap>,
) -> Self {
let api_url = if let Some(ref ep) = credentials.endpoint {
ep.trim_end_matches('/').to_string()
} else {
"https://api.pinata.cloud".to_string()
};
let gateway_url = if let Some(ref gw) = credentials.gateway_url {
gw.clone()
} else {
"https://gateway.pinata.cloud/ipfs/".to_string()
};
Self {
client: super::utils::build_http_client(timeout, custom_headers),
credentials,
api_url,
gateway_url,
}
}
async fn resolve_cid(&self, remote_path: &str) -> Result<String, StorageError> {
let query_url = format!("{}/data/pinList", self.api_url);
let res = super::utils::apply_bearer_auth(self.client.get(&query_url), &self.credentials.jwt_token)
.query(&[
("status", "pinned"),
("metadata[name]", remote_path),
])
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "resolve_cid").await);
}
let body: PinataPinListResponse = res.json().await?;
let row = body.rows.first()
.ok_or_else(|| StorageError::NotFound(format!("File '{}' not found in IPFS pinned index", remote_path)))?;
Ok(row.ipfs_pin_hash.clone())
}
}
#[async_trait]
impl StorageBackend for IPFSProvider {
fn name(&self) -> &str {
"IPFS"
}
async fn upload(&self, local_path: &Path, remote_path: &str) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "upload", || async {
let clean_path = self.format_path(remote_path);
info!("[{}] Real upload starting for '{}'", self.name(), clean_path);
let file_content = fs::read(local_path).await?;
let file_name = Path::new(clean_path.as_ref())
.file_name()
.and_then(|n| n.to_str())
.unwrap_or(clean_path.as_ref())
.to_string();
let upload_url = format!("{}/pinning/pinFileToIPFS", self.api_url);
let metadata_json = serde_json::json!({
"name": clean_path
}).to_string();
let form = reqwest::multipart::Form::new()
.part("file", reqwest::multipart::Part::bytes(file_content).file_name(file_name))
.text("pinataMetadata", metadata_json);
let res = super::utils::apply_bearer_auth(self.client.post(&upload_url), &self.credentials.jwt_token)
.multipart(form)
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "upload").await);
}
let _body: PinataPinFileResponse = res.json().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 clean_path = self.format_path(remote_path);
let cid = self.resolve_cid(&clean_path).await?;
let download_url = if self.gateway_url.ends_with('/') {
format!("{}{}", self.gateway_url, cid)
} else {
format!("{}/{}", self.gateway_url, cid)
};
let dl_res = self.client.get(&download_url).send().await?;
if !dl_res.status().is_success() {
return Err(translate_http_error(dl_res, self.name(), "download").await);
}
if let Some(parent) = local_path.parent() {
fs::create_dir_all(parent).await?;
}
let bytes = dl_res.bytes().await?;
fs::write(local_path, bytes).await?;
Ok(())
}).await
}
async fn delete(&self, remote_path: &str) -> Result<(), StorageError> {
super::utils::execute_with_retry(self.name(), "delete", || async {
let clean_path = self.format_path(remote_path);
let cid = self.resolve_cid(&clean_path).await?;
let unpin_url = format!("{}/pinning/unpin/{}", self.api_url, cid);
let res = super::utils::apply_bearer_auth(self.client.delete(&unpin_url), &self.credentials.jwt_token)
.send()
.await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "delete").await);
}
Ok(())
}).await
}
async fn list(&self, remote_path: &str) -> Result<Vec<StorageItem>, StorageError> {
super::utils::execute_with_retry(self.name(), "list", || async {
let clean_path = self.format_path(remote_path);
let list_url = format!("{}/data/pinList", self.api_url);
let mut req = super::utils::apply_bearer_auth(self.client.get(&list_url), &self.credentials.jwt_token)
.query(&[("status", "pinned")]);
if !clean_path.is_empty() {
req = req.query(&[("metadata[name]", &clean_path)]);
}
let res = req.send().await?;
if !res.status().is_success() {
return Err(translate_http_error(res, self.name(), "list").await);
}
let list_response: PinataPinListResponse = res.json().await?;
let mut items = Vec::new();
for row in list_response.rows {
let item_path = super::utils::strip_destination_prefix(Path::new(&row.metadata.name), self.credentials.common.destination_folder.as_deref());
let modified = row.date_pinned
.as_ref()
.and_then(|t| time::OffsetDateTime::parse(t, &time::format_description::well_known::Rfc3339).ok())
.map(SystemTime::from)
.unwrap_or(SystemTime::now());
items.push(StorageItem {
path: item_path,
size: row.size,
modified,
is_dir: false,
checksum: None,
permissions: None,
});
}
Ok(items)
}).await
}
}
pub struct IPFSProviderBuilder {
pub credentials: IPFSCredentials,
pub timeout: Option<std::time::Duration>,
pub custom_headers: Option<reqwest::header::HeaderMap>,
}
impl IPFSProviderBuilder {
pub fn new(credentials: IPFSCredentials) -> Self {
Self {
credentials,
timeout: None,
custom_headers: None,
}
}
pub fn build(self) -> IPFSProvider {
IPFSProvider::with_client_options(self.credentials, self.timeout, self.custom_headers)
}
}