diff --git a/seaweed-volume/src/remote_storage/s3.rs b/seaweed-volume/src/remote_storage/s3.rs
index 03731c03f..400098433 100644
--- a/seaweed-volume/src/remote_storage/s3.rs
+++ b/seaweed-volume/src/remote_storage/s3.rs
@@ -209,7 +209,7 @@ impl RemoteStorageClient for S3RemoteStorageClient {
}
#[cfg(test)]
-mod tests {
+pub(crate) mod tests {
use super::*;
use aws_sdk_s3::config::http::{HttpRequest, HttpResponse};
use aws_sdk_s3::config::retry::RetryConfig;
@@ -224,9 +224,9 @@ mod tests {
/// so the error-mapping paths can be exercised without a network or a
/// running S3 server.
#[derive(Debug, Clone)]
- struct CannedResponse {
- status: u16,
- body: &'static str,
+ pub(crate) struct CannedResponse {
+ pub(crate) status: u16,
+ pub(crate) body: &'static str,
}
impl HttpConnector for CannedResponse {
@@ -272,7 +272,7 @@ mod tests {
}
}
- const NO_SUCH_KEY: &str = r#"
+ pub(crate) const NO_SUCH_KEY: &str = r#"
NoSuchKeyThe specified key does not exist.dir/missing"#;
const NOT_FOUND_BODY: &str = r#"
diff --git a/seaweed-volume/src/remote_storage/s3_tier.rs b/seaweed-volume/src/remote_storage/s3_tier.rs
index 82d2f1102..1bcea50bb 100644
--- a/seaweed-volume/src/remote_storage/s3_tier.rs
+++ b/seaweed-volume/src/remote_storage/s3_tier.rs
@@ -8,8 +8,11 @@ use std::future::Future;
use std::sync::{Arc, OnceLock, RwLock};
use aws_sdk_s3::Client;
+use aws_sdk_s3::config::http::HttpResponse;
use aws_sdk_s3::config::{BehaviorVersion, Credentials, Region};
-use aws_sdk_s3::error::DisplayErrorContext;
+use aws_sdk_s3::error::{DisplayErrorContext, SdkError};
+use aws_sdk_s3::operation::get_object::GetObjectError;
+use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart};
use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
use tokio::sync::Semaphore;
@@ -17,6 +20,53 @@ use tokio::sync::Semaphore;
/// Concurrency limit for multipart upload/download (matches Go's s3manager).
const CONCURRENCY: usize = 5;
+/// A tier transfer failure. The variant is what callers match on; the
+/// message is the operator-facing text.
+#[derive(Debug, thiserror::Error)]
+pub enum TierError {
+ /// The remote object does not exist.
+ #[error("{0}")]
+ NotFound(String),
+ /// An S3 request or a local file operation failed.
+ #[error("{0}")]
+ Io(String),
+ /// The tier I/O runtime could not be built or dropped the task.
+ #[error("{0}")]
+ RuntimeUnavailable(String),
+ /// The progress callback asked to stop.
+ #[error("{0}")]
+ Aborted(String),
+}
+
+// Not-found rules as in remote_storage/s3.rs: HEAD by the raw 404 status,
+// GET by the NoSuchKey code only.
+fn head_object_error(key: &str, e: SdkError) -> TierError {
+ let message = format!("failed to head object {}: {}", key, DisplayErrorContext(&e));
+ match e {
+ SdkError::ServiceError(ref se) if se.raw().status().as_u16() == 404 => {
+ TierError::NotFound(message)
+ }
+ _ => TierError::Io(message),
+ }
+}
+
+fn get_object_error(
+ key: &str,
+ range: &str,
+ e: SdkError,
+) -> TierError {
+ let message = format!(
+ "failed to get object {} range {}: {}",
+ key,
+ range,
+ DisplayErrorContext(&e)
+ );
+ match e {
+ SdkError::ServiceError(ref se) if se.err().is_no_such_key() => TierError::NotFound(message),
+ _ => TierError::Io(message),
+ }
+}
+
/// Configuration for an S3 tier backend.
#[derive(Debug, Clone)]
pub struct S3TierConfig {
@@ -90,7 +140,7 @@ impl S3TierBackend {
&self,
file_path: &str,
progress_fn: F,
- ) -> Result<(String, u64), String>
+ ) -> Result<(String, u64), TierError>
where
F: FnMut(i64, f32) -> Result<(), String> + Send + Sync + 'static,
{
@@ -98,7 +148,7 @@ impl S3TierBackend {
let metadata = tokio::fs::metadata(file_path)
.await
- .map_err(|e| format!("failed to stat file {}: {}", file_path, e))?;
+ .map_err(|e| TierError::Io(format!("failed to stat file {}: {}", file_path, e)))?;
let file_size = metadata.len();
// Calculate part size: start at 64MB, scale up for very large files (matches Go)
@@ -121,15 +171,15 @@ impl S3TierBackend {
.send()
.await
.map_err(|e| {
- format!(
+ TierError::Io(format!(
"failed to create multipart upload: {}",
DisplayErrorContext(&e)
- )
+ ))
})?;
let upload_id = create_resp
.upload_id()
- .ok_or_else(|| "no upload_id in multipart upload response".to_string())?
+ .ok_or_else(|| TierError::Io("no upload_id in multipart upload response".to_string()))?
.to_string();
// Build list of (part_number, offset, size) for all parts
@@ -165,19 +215,21 @@ impl S3TierBackend {
let _permit = sem
.acquire()
.await
- .map_err(|e| format!("semaphore error: {}", e))?;
+ .map_err(|e| TierError::Io(format!("semaphore error: {}", e)))?;
// Read this part's data from the file at the correct offset
let mut file = tokio::fs::File::open(&fp)
.await
- .map_err(|e| format!("failed to open file {}: {}", fp, e))?;
+ .map_err(|e| TierError::Io(format!("failed to open file {}: {}", fp, e)))?;
file.seek(std::io::SeekFrom::Start(off))
.await
- .map_err(|e| format!("failed to seek to offset {}: {}", off, e))?;
+ .map_err(|e| {
+ TierError::Io(format!("failed to seek to offset {}: {}", off, e))
+ })?;
let mut buf = vec![0u8; size];
- file.read_exact(&mut buf)
- .await
- .map_err(|e| format!("failed to read file at offset {}: {}", off, e))?;
+ file.read_exact(&mut buf).await.map_err(|e| {
+ TierError::Io(format!("failed to read file at offset {}: {}", off, e))
+ })?;
let upload_part_resp = client
.upload_part()
@@ -189,12 +241,12 @@ impl S3TierBackend {
.send()
.await
.map_err(|e| {
- format!(
+ TierError::Io(format!(
"failed to upload part {} at offset {}: {}",
pn,
off,
DisplayErrorContext(&e)
- )
+ ))
})?;
let e_tag = upload_part_resp.e_tag().unwrap_or_default().to_string();
@@ -213,9 +265,9 @@ impl S3TierBackend {
};
(guard.1)(uploaded as i64, pct)
};
- progress_result?;
+ progress_result.map_err(TierError::Aborted)?;
- Ok::<_, String>(
+ Ok::<_, TierError>(
CompletedPart::builder()
.e_tag(e_tag)
.part_number(pn)
@@ -230,7 +282,7 @@ impl S3TierBackend {
for handle in handles {
let part = handle
.await
- .map_err(|e| format!("upload task panicked: {}", e))??;
+ .map_err(|e| TierError::Io(format!("upload task panicked: {}", e)))??;
completed_parts.push(part);
}
@@ -248,13 +300,13 @@ impl S3TierBackend {
.send()
.await
.map_err(|e| {
- format!(
+ TierError::Io(format!(
"failed to complete multipart upload: {}",
DisplayErrorContext(&e)
- )
+ ))
})?;
- Ok::<(), String>(())
+ Ok::<(), TierError>(())
}
.await;
@@ -297,7 +349,7 @@ impl S3TierBackend {
dest_path: &str,
key: &str,
progress_fn: F,
- ) -> Result
+ ) -> Result
where
F: FnMut(i64, f32) -> Result<(), String> + Send + Sync + 'static,
{
@@ -309,7 +361,7 @@ impl S3TierBackend {
.key(key)
.send()
.await
- .map_err(|e| format!("failed to head object {}: {}", key, DisplayErrorContext(&e)))?;
+ .map_err(|e| head_object_error(key, e))?;
let file_size = head_resp.content_length().unwrap_or(0) as u64;
@@ -321,10 +373,12 @@ impl S3TierBackend {
.truncate(true)
.open(dest_path)
.await
- .map_err(|e| format!("failed to open dest file {}: {}", dest_path, e))?;
+ .map_err(|e| {
+ TierError::Io(format!("failed to open dest file {}: {}", dest_path, e))
+ })?;
file.set_len(file_size)
.await
- .map_err(|e| format!("failed to set file length: {}", e))?;
+ .map_err(|e| TierError::Io(format!("failed to set file length: {}", e)))?;
}
let part_size: u64 = 64 * 1024 * 1024;
@@ -360,7 +414,7 @@ impl S3TierBackend {
let _permit = sem
.acquire()
.await
- .map_err(|e| format!("semaphore error: {}", e))?;
+ .map_err(|e| TierError::Io(format!("semaphore error: {}", e)))?;
let end = off + size - 1;
let range = format!("bytes={}-{}", off, end);
@@ -372,20 +426,13 @@ impl S3TierBackend {
.range(&range)
.send()
.await
- .map_err(|e| {
- format!(
- "failed to get object {} range {}: {}",
- key,
- range,
- DisplayErrorContext(&e)
- )
- })?;
+ .map_err(|e| get_object_error(&key, &range, e))?;
let body = get_resp
.body
.collect()
.await
- .map_err(|e| format!("failed to read body: {}", e))?;
+ .map_err(|e| TierError::Io(format!("failed to read body: {}", e)))?;
let bytes = body.into_bytes();
// Write at the correct offset (like Go's WriteAt)
@@ -393,13 +440,17 @@ impl S3TierBackend {
.write(true)
.open(&dp)
.await
- .map_err(|e| format!("failed to open dest file {}: {}", dp, e))?;
+ .map_err(|e| {
+ TierError::Io(format!("failed to open dest file {}: {}", dp, e))
+ })?;
file.seek(std::io::SeekFrom::Start(off))
.await
- .map_err(|e| format!("failed to seek to offset {}: {}", off, e))?;
+ .map_err(|e| {
+ TierError::Io(format!("failed to seek to offset {}: {}", off, e))
+ })?;
file.write_all(&bytes)
.await
- .map_err(|e| format!("failed to write to {}: {}", dp, e))?;
+ .map_err(|e| TierError::Io(format!("failed to write to {}: {}", dp, e)))?;
// Report progress. The lock is released before the result is
// propagated so an aborting callback cannot poison the mutex
@@ -415,9 +466,9 @@ impl S3TierBackend {
};
(guard.1)(downloaded as i64, pct)
};
- progress_result?;
+ progress_result.map_err(TierError::Aborted)?;
- Ok::<_, String>(())
+ Ok::<_, TierError>(())
}));
}
@@ -425,7 +476,7 @@ impl S3TierBackend {
for handle in handles {
handle
.await
- .map_err(|e| format!("download task panicked: {}", e))??;
+ .map_err(|e| TierError::Io(format!("download task panicked: {}", e)))??;
}
// fsync the file so its content is durable before the caller trims the .vif
@@ -434,16 +485,21 @@ impl S3TierBackend {
.write(true)
.open(dest_path)
.await
- .map_err(|e| format!("failed to open {} for fsync: {}", dest_path, e))?;
+ .map_err(|e| TierError::Io(format!("failed to open {} for fsync: {}", dest_path, e)))?;
synced
.sync_all()
.await
- .map_err(|e| format!("failed to fsync {}: {}", dest_path, e))?;
+ .map_err(|e| TierError::Io(format!("failed to fsync {}: {}", dest_path, e)))?;
Ok(file_size)
}
- pub async fn read_range(&self, key: &str, offset: u64, size: usize) -> Result, String> {
+ pub async fn read_range(
+ &self,
+ key: &str,
+ offset: u64,
+ size: usize,
+ ) -> Result, TierError> {
let end = offset + (size as u64).saturating_sub(1);
let range = format!("bytes={}-{}", offset, end);
let resp = self
@@ -454,25 +510,18 @@ impl S3TierBackend {
.range(&range)
.send()
.await
- .map_err(|e| {
- format!(
- "failed to get object {} range {}: {}",
- key,
- range,
- DisplayErrorContext(&e)
- )
- })?;
+ .map_err(|e| get_object_error(key, &range, e))?;
let body = resp
.body
.collect()
.await
- .map_err(|e| format!("failed to read object {} body: {}", key, e))?;
+ .map_err(|e| TierError::Io(format!("failed to read object {} body: {}", key, e)))?;
Ok(body.into_bytes().to_vec())
}
/// Delete a file from S3.
- pub async fn delete_file(&self, key: &str) -> Result<(), String> {
+ pub async fn delete_file(&self, key: &str) -> Result<(), TierError> {
self.client
.delete_object()
.bucket(&self.bucket)
@@ -480,16 +529,16 @@ impl S3TierBackend {
.send()
.await
.map_err(|e| {
- format!(
+ TierError::Io(format!(
"failed to delete object {}: {}",
key,
DisplayErrorContext(&e)
- )
+ ))
})?;
Ok(())
}
- pub fn delete_file_blocking(&self, key: &str) -> Result<(), String> {
+ pub fn delete_file_blocking(&self, key: &str) -> Result<(), TierError> {
let client = self.client.clone();
let bucket = self.bucket.clone();
let key = key.to_string();
@@ -501,11 +550,11 @@ impl S3TierBackend {
.send()
.await
.map_err(|e| {
- format!(
+ TierError::Io(format!(
"failed to delete object {}: {}",
key,
DisplayErrorContext(&e)
- )
+ ))
})?;
Ok(())
})
@@ -516,7 +565,7 @@ impl S3TierBackend {
key: &str,
offset: u64,
size: usize,
- ) -> Result, String> {
+ ) -> Result, TierError> {
let client = self.client.clone();
let bucket = self.bucket.clone();
let key = key.to_string();
@@ -530,20 +579,12 @@ impl S3TierBackend {
.range(&range)
.send()
.await
- .map_err(|e| {
- format!(
- "failed to get object {} range {}: {}",
- key,
- range,
- DisplayErrorContext(&e)
- )
- })?;
+ .map_err(|e| get_object_error(&key, &range, e))?;
- let body = resp
- .body
- .collect()
- .await
- .map_err(|e| format!("failed to read object {} body: {}", key, e))?;
+ let body =
+ resp.body.collect().await.map_err(|e| {
+ TierError::Io(format!("failed to read object {} body: {}", key, e))
+ })?;
Ok(body.into_bytes().to_vec())
})
}
@@ -615,7 +656,7 @@ pub fn global_s3_tier_registry() -> &'static RwLock {
static TIER_RUNTIME: std::sync::Mutex