From 68944e83a3267a870af627e072f75125d2efddbd Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sun, 27 Sep 2026 14:39:25 +0300 Subject: [PATCH] volume: typed tier errors so a missing remote object answers NotFound (#11484) remote_storage/s3_tier.rs returned Result<_, String> from every transfer (upload_file, download_file, read_range[_blocking], delete_file[_blocking]) and from the tier runtime helpers. The tier move handlers could only wrap that in Status::internal, so a .dat whose remote object is gone was indistinguishable from an I/O failure to weed shell. Add TierError { NotFound, Io, RuntimeUnavailable, Aborted }. Each variant carries the existing message verbatim. NotFound follows the rules remote_storage/s3.rs already uses: raw 404 status on HEAD, NoSuchKey code on GET; a bare 404 on GET stays Io. A progress-callback Err becomes Aborted. VolumeError gains a transparent Tier variant and From for Status maps Tier(NotFound) to NotFound; the tier move handlers go through status_with_context, so their message text is unchanged. Every other tier failure is still Internal. The remote needle read path keeps io::Error::other, so its error kind and vacuum's handling of it do not change. Co-authored-by: Claude Opus 5.5 (1M context) --- seaweed-volume/src/remote_storage/s3.rs | 10 +- seaweed-volume/src/remote_storage/s3_tier.rs | 304 +++++++++++++------ seaweed-volume/src/server/grpc_server.rs | 118 ++++++- seaweed-volume/src/server/mod.rs | 26 +- seaweed-volume/src/storage/volume.rs | 3 + 5 files changed, 353 insertions(+), 108 deletions(-) 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> = std::sync::Mutex::new(None); -fn tier_handle() -> Result { +fn tier_handle() -> Result { let mut slot = TIER_RUNTIME .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); @@ -625,7 +666,12 @@ fn tier_handle() -> Result { .thread_name("tier-io") .enable_all() .build() - .map_err(|e| format!("failed to build the tier I/O tokio runtime: {}", e))?; + .map_err(|e| { + TierError::RuntimeUnavailable(format!( + "failed to build the tier I/O tokio runtime: {}", + e + )) + })?; *slot = Some(runtime); } Ok(slot.as_ref().expect("just initialised").handle().clone()) @@ -635,9 +681,9 @@ fn tier_handle() -> Result { /// finishes. The caller may be a worker of *another* tokio runtime, so this /// waits on a channel rather than `Handle::block_on`, which panics when /// called from inside any runtime context. -fn block_on_tier_future(future: F) -> Result +fn block_on_tier_future(future: F) -> Result where - F: Future> + Send + 'static, + F: Future> + Send + 'static, T: Send + 'static, { let handle = tier_handle()?; @@ -651,13 +697,15 @@ where match rx.recv() { Ok(Ok(result)) => result, Ok(Err(join_error)) => Err(describe_join_error(join_error)), - Err(_) => Err("tier I/O runtime dropped the task before it finished".to_string()), + Err(_) => Err(TierError::RuntimeUnavailable( + "tier I/O runtime dropped the task before it finished".to_string(), + )), } } /// Turn a `JoinError` into a message that keeps the panic payload, so an /// SDK panic surfaces as "boom" rather than a fixed "thread panicked". -fn describe_join_error(join_error: tokio::task::JoinError) -> String { +fn describe_join_error(join_error: tokio::task::JoinError) -> TierError { if join_error.is_panic() { let payload = join_error.into_panic(); let message = if let Some(s) = payload.downcast_ref::<&str>() { @@ -667,19 +715,20 @@ fn describe_join_error(join_error: tokio::task::JoinError) -> String { } else { "non-string panic payload".to_string() }; - format!("tier I/O task panicked: {}", message) + TierError::Io(format!("tier I/O task panicked: {}", message)) } else { - format!("tier I/O task failed: {}", join_error) + TierError::RuntimeUnavailable(format!("tier I/O task failed: {}", join_error)) } } #[cfg(test)] mod tests { use super::*; + use crate::remote_storage::s3::tests::{CannedResponse, NO_SUCH_KEY}; use std::collections::HashSet; use tokio::runtime::Handle; - fn probe() -> Result<(tokio::runtime::Id, Option), String> { + fn probe() -> Result<(tokio::runtime::Id, Option), TierError> { block_on_tier_future(async { Ok(( Handle::current().id(), @@ -709,10 +758,12 @@ mod tests { #[test] fn block_on_tier_future_returns_the_value_and_the_error() { - assert_eq!(block_on_tier_future(async { Ok(7u32) }), Ok(7)); - assert_eq!( - block_on_tier_future::<_, u32>(async { Err("nope".to_string()) }), - Err("nope".to_string()) + assert_eq!(block_on_tier_future(async { Ok(7u32) }).unwrap(), 7); + let err = block_on_tier_future::<_, u32>(async { Err(TierError::NotFound("nope".into())) }) + .unwrap_err(); + assert!( + matches!(&err, TierError::NotFound(m) if m == "nope"), + "{err:?}" ); } @@ -761,8 +812,9 @@ mod tests { Ok(()) }) .expect_err("a panicking future must be an error"); - assert!(err.contains("boom 42"), "got: {err}"); - assert!(err.contains("panicked"), "got: {err}"); + assert!(matches!(err, TierError::Io(_)), "got: {err:?}"); + assert!(err.to_string().contains("boom 42"), "got: {err}"); + assert!(err.to_string().contains("panicked"), "got: {err}"); } #[test] @@ -774,7 +826,83 @@ mod tests { Ok(()) }) .expect_err("a panicking future must be an error"); - assert!(err.contains("static boom"), "got: {err}"); + assert!(err.to_string().contains("static boom"), "got: {err}"); + } + + fn backend_answering(status: u16, body: &'static str) -> S3TierBackend { + let config = aws_sdk_s3::Config::builder() + .behavior_version(BehaviorVersion::latest()) + .region(Region::new("us-east-1")) + .credentials_provider(Credentials::new("AKIATEST", "secret", None, None, "test")) + .endpoint_url("http://127.0.0.1:1") + .force_path_style(true) + .http_client(CannedResponse { status, body }) + .retry_config(aws_sdk_s3::config::retry::RetryConfig::disabled()) + .build(); + S3TierBackend { + client: Client::from_conf(config), + bucket: "bucket".to_string(), + storage_class: "STANDARD".to_string(), + } + } + + #[tokio::test] + async fn download_head_404_is_not_found() { + let tmp = tempfile::tempdir().unwrap(); + let dest = tmp.path().join("1.dat"); + let err = backend_answering(404, "") + .download_file(dest.to_str().unwrap(), "missing", |_, _| Ok(())) + .await + .unwrap_err(); + assert!(matches!(err, TierError::NotFound(_)), "{err:?}"); + assert!( + err.to_string() + .starts_with("failed to head object missing: "), + "{err}" + ); + } + + #[tokio::test] + async fn download_head_403_is_io() { + let tmp = tempfile::tempdir().unwrap(); + let dest = tmp.path().join("1.dat"); + let err = backend_answering(403, "") + .download_file(dest.to_str().unwrap(), "denied", |_, _| Ok(())) + .await + .unwrap_err(); + assert!(matches!(err, TierError::Io(_)), "{err:?}"); + } + + #[tokio::test] + async fn read_range_no_such_key_is_not_found() { + let err = backend_answering(404, NO_SUCH_KEY) + .read_range("missing", 0, 8) + .await + .unwrap_err(); + assert!(matches!(err, TierError::NotFound(_)), "{err:?}"); + assert!( + err.to_string() + .starts_with("failed to get object missing range bytes=0-7: "), + "{err}" + ); + } + + #[tokio::test] + async fn read_range_bare_404_is_io() { + // As in Go, GET is not-found by the NoSuchKey code, not the status. + let err = backend_answering(404, "") + .read_range("missing", 0, 8) + .await + .unwrap_err(); + assert!(matches!(err, TierError::Io(_)), "{err:?}"); + } + + #[test] + fn read_range_blocking_no_such_key_is_not_found() { + let err = backend_answering(404, NO_SUCH_KEY) + .read_range_blocking("missing", 0, 8) + .unwrap_err(); + assert!(matches!(err, TierError::NotFound(_)), "{err:?}"); } #[test] diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index d4ca064a9..31c286a68 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -4442,10 +4442,10 @@ impl VolumeServer for VolumeGrpcService { }) .await .map_err(|e| { - Status::internal(format!( - "backend {} copy file {}: {}", - dest_backend_name, dat_path, e - )) + crate::server::status_with_context( + &format!("backend {} copy file {}", dest_backend_name, dat_path), + e.into(), + ) })?; // Deliberately no cancellation check here. Once the object is @@ -4664,10 +4664,10 @@ impl VolumeServer for VolumeGrpcService { rm ); } - Status::internal(format!( - "backend {} copy file {}: {}", - storage_name_clone, dat_path, e - )) + crate::server::status_with_context( + &format!("backend {} copy file {}", storage_name_clone, dat_path), + e.into(), + ) })?; // Restore the .dat mtime so a reload computes TTL from real data age, @@ -4796,10 +4796,13 @@ impl VolumeServer for VolumeGrpcService { // Only the last replica to download deletes the shared remote object. backend.delete_file(&storage_key).await.map_err(|e| { - Status::internal(format!( - "volume {} failed to delete remote file {}: {}", - vid, storage_key, e - )) + crate::server::status_with_context( + &format!( + "volume {} failed to delete remote file {}", + vid, storage_key + ), + e.into(), + ) })?; // Go does NOT send a final 100% progress message after download completion @@ -6196,7 +6199,9 @@ fn run_vacuum_compact( mod tests { use super::*; use crate::config::MinFreeSpace; - use crate::remote_storage::s3_tier::{S3TierBackend, S3TierConfig, global_s3_tier_registry}; + use crate::remote_storage::s3_tier::{ + S3TierBackend, S3TierConfig, TierError, global_s3_tier_registry, + }; use crate::security::{Guard, SigningKey}; use crate::server::grpc_client::GRPC_MAX_MESSAGE_SIZE; use crate::storage::needle_map::NeedleMapKind; @@ -6870,6 +6875,88 @@ mod tests { .remove("s3.tier_down_delete"); } + /// An S3 endpoint that answers every request 404 with no body. + fn spawn_s3_not_found_server() -> (String, tokio::sync::oneshot::Sender<()>) { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + listener.set_nonblocking(true).unwrap(); + let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>(); + let (ready_tx, ready_rx) = std::sync::mpsc::channel::<()>(); + std::thread::spawn(move || { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap(); + runtime.block_on(async move { + let app = axum::Router::new().fallback(axum::routing::any(|| async { + axum::http::StatusCode::NOT_FOUND + })); + let listener = tokio::net::TcpListener::from_std(listener).unwrap(); + let _ = ready_tx.send(()); + axum::serve(listener, app) + .with_graceful_shutdown(async move { + let _ = shutdown_rx.await; + }) + .await + .unwrap(); + }); + }); + ready_rx.recv().unwrap(); + (format!("http://{}", addr), shutdown_tx) + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_tier_move_from_remote_missing_object_is_not_found() { + let (service, tmp, shutdown_tx, _dat_bytes, _super_block_size, _delete_count) = + make_remote_only_service("tier_down_missing"); + let (endpoint, missing_shutdown_tx) = spawn_s3_not_found_server(); + global_s3_tier_registry().write().unwrap().register( + "s3.tier_down_missing".to_string(), + S3TierBackend::new(&S3TierConfig { + access_key: "access".to_string(), + secret_key: "secret".to_string(), + region: "us-east-1".to_string(), + bucket: "bucket-a".to_string(), + endpoint, + storage_class: "STANDARD".to_string(), + force_path_style: true, + }), + ); + let dat_path = format!("{}/1.dat", tmp.path().to_str().unwrap()); + + let mut stream = service + .volume_tier_move_dat_from_remote(Request::new( + volume_server_pb::VolumeTierMoveDatFromRemoteRequest { + volume_id: 1, + collection: String::new(), + keep_remote_dat_file: false, + }, + )) + .await + .unwrap() + .into_inner(); + let err = stream + .next() + .await + .expect("the stream must carry the failure") + .expect_err("a missing remote object must fail the tier-down"); + assert_eq!(err.code(), tonic::Code::NotFound, "{err:?}"); + assert!( + err.message().starts_with(&format!( + "backend s3.tier_down_missing copy file {dat_path}: " + )), + "the message must keep its context: {err:?}" + ); + assert!(!std::path::Path::new(&dat_path).exists()); + + let _ = shutdown_tx.send(()); + let _ = missing_shutdown_tx.send(()); + global_s3_tier_registry() + .write() + .unwrap() + .remove("s3.tier_down_missing"); + } + // The progress callback is the only thing a tier transfer polls, so it is // the only place a departing caller can be noticed. Go gives it an error // return for exactly this (s3_upload.go:99, s3_download.go:84); this checks @@ -6889,7 +6976,10 @@ mod tests { .download_file(&dest, "remote-key", |_, _| Err("caller gone".to_string())) .await .expect_err("a failing progress callback must abort the download"); - assert!(err.contains("caller gone"), "unexpected error: {}", err); + assert!( + matches!(&err, TierError::Aborted(m) if m == "caller gone"), + "unexpected error: {err:?}" + ); drop(service); let _ = shutdown_tx.send(()); diff --git a/seaweed-volume/src/server/mod.rs b/seaweed-volume/src/server/mod.rs index de34ebe71..f18fc7745 100644 --- a/seaweed-volume/src/server/mod.rs +++ b/seaweed-volume/src/server/mod.rs @@ -1,5 +1,6 @@ use tonic::Status; +use crate::remote_storage::s3_tier::TierError; use crate::storage::volume::VolumeError; #[cfg(unix)] @@ -23,7 +24,9 @@ impl From for Status { fn from(err: VolumeError) -> Self { let message = err.to_string(); match err { - VolumeError::NotFound | VolumeError::VolumeNotFound(_) => Status::not_found(message), + VolumeError::NotFound + | VolumeError::VolumeNotFound(_) + | VolumeError::Tier(TierError::NotFound(_)) => Status::not_found(message), VolumeError::ReadOnly | VolumeError::NotEmpty => Status::failed_precondition(message), VolumeError::InsufficientSpace { .. } => Status::resource_exhausted(message), VolumeError::AlreadyExists => Status::already_exists(message), @@ -77,6 +80,27 @@ mod tests { ); assert_eq!(code(VolumeError::AlreadyExists), Code::AlreadyExists); assert_eq!(code(VolumeError::NotInitialized), Code::Internal); + assert_eq!( + code(TierError::NotFound("gone".into()).into()), + Code::NotFound + ); + for tier in [ + TierError::Io("io".into()), + TierError::RuntimeUnavailable("rt".into()), + TierError::Aborted("bye".into()), + ] { + assert_eq!(code(tier.into()), Code::Internal); + } + + let status = status_with_context( + "backend s3.default copy file /data/1.dat", + TierError::NotFound("failed to head object k: NotFound".into()).into(), + ); + assert_eq!(status.code(), Code::NotFound); + assert_eq!( + status.message(), + "backend s3.default copy file /data/1.dat: failed to head object k: NotFound" + ); let status = status_with_context( "compact volume 7", diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index e3d837f6c..bb0df91be 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -90,6 +90,9 @@ pub enum VolumeError { #[error("streaming from remote-backed volume requires buffered fallback")] StreamingUnsupported, + + #[error(transparent)] + Tier(#[from] crate::remote_storage::s3_tier::TierError), } /// Returns true when a needle read failed because the on-disk bytes are