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<VolumeError> 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) <noreply@anthropic.com>
This commit is contained in:
Eliah Rusin
2026-09-27 19:39:25 +08:00
committed by GitHub
co-authored by Claude Opus 5.5
parent 00310f6588
commit 68944e83a3
5 changed files with 353 additions and 108 deletions
+5 -5
View File
@@ -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#"<?xml version="1.0" encoding="UTF-8"?>
pub(crate) const NO_SUCH_KEY: &str = r#"<?xml version="1.0" encoding="UTF-8"?>
<Error><Code>NoSuchKey</Code><Message>The specified key does not exist.</Message><Key>dir/missing</Key></Error>"#;
const NOT_FOUND_BODY: &str = r#"<?xml version="1.0" encoding="UTF-8"?>
+216 -88
View File
@@ -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<HeadObjectError, HttpResponse>) -> 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<GetObjectError, HttpResponse>,
) -> 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<u64, String>
) -> Result<u64, TierError>
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<Vec<u8>, String> {
pub async fn read_range(
&self,
key: &str,
offset: u64,
size: usize,
) -> Result<Vec<u8>, 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<Vec<u8>, String> {
) -> Result<Vec<u8>, 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<S3TierRegistry> {
static TIER_RUNTIME: std::sync::Mutex<Option<tokio::runtime::Runtime>> =
std::sync::Mutex::new(None);
fn tier_handle() -> Result<tokio::runtime::Handle, String> {
fn tier_handle() -> Result<tokio::runtime::Handle, TierError> {
let mut slot = TIER_RUNTIME
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
@@ -625,7 +666,12 @@ fn tier_handle() -> Result<tokio::runtime::Handle, String> {
.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<tokio::runtime::Handle, String> {
/// 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<F, T>(future: F) -> Result<T, String>
fn block_on_tier_future<F, T>(future: F) -> Result<T, TierError>
where
F: Future<Output = Result<T, String>> + Send + 'static,
F: Future<Output = Result<T, TierError>> + 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>), String> {
fn probe() -> Result<(tokio::runtime::Id, Option<String>), 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]
+104 -14
View File
@@ -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(());
+25 -1
View File
@@ -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<VolumeError> 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",
+3
View File
@@ -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