diff --git a/Cargo.lock b/Cargo.lock index 1f691dfc..5cabfbde 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -208,7 +208,7 @@ dependencies = [ "argh_shared", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -251,7 +251,7 @@ checksum = "c7c24de15d275a1ecfd47a380fb4d5ec9bfe0933f309ed5e705b775596a3574d" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -262,7 +262,7 @@ checksum = "9035ad2d096bed7955a320ee7e2230574d28fd3c3a0f186cbea1ff3c7eed5dbb" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -757,7 +757,7 @@ checksum = "f46882e17999c6cc590af592290432be3bce0428cb0d5f8b6715e4dc7b383eb3" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -781,7 +781,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn", + "syn 2.0.117", ] [[package]] @@ -792,7 +792,7 @@ checksum = "d38308df82d1080de0afee5d069fa14b0326a88c14f15c5ccda35b4a6c414c81" dependencies = [ "darling_core", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -832,6 +832,17 @@ dependencies = [ "serde_core", ] +[[package]] +name = "derive-error-kind" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3a422c2552b72ab3fc01435089a5d9f5337e42383543e885275699ef0f539ba" +dependencies = [ + "proc-macro2", + "quote", + "syn 1.0.109", +] + [[package]] name = "derive_more" version = "2.1.1" @@ -851,7 +862,7 @@ dependencies = [ "proc-macro2", "quote", "rustc_version", - "syn", + "syn 2.0.117", "unicode-xid", ] @@ -885,7 +896,7 @@ checksum = "97369cbbc041bc366949bc74d34658d6cda5621039731c6310521892a3a20ae0" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1201,7 +1212,7 @@ checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1986,7 +1997,7 @@ dependencies = [ "quote", "rustc_version", "simd_cesu8", - "syn", + "syn 2.0.117", ] [[package]] @@ -2011,7 +2022,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "38c0b942f458fe50cdac086d2f946512305e5631e720728f2a61aabcd47a6264" dependencies = [ "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2738,6 +2749,7 @@ dependencies = [ "bigtable_rs", "bytes", "chrono", + "derive-error-kind", "futures", "futures-util", "gcp_auth", @@ -2819,7 +2831,7 @@ checksum = "a948666b637a0f465e8564c73e89d4dde00d72d4d473cc972f390fc3dcee7d9c" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2948,7 +2960,7 @@ dependencies = [ "proc-macro2", "proc-macro2-diagnostics", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -3004,7 +3016,7 @@ checksum = "d9b20ed30f105399776b9c883e68e536ef602a16ae6f596d2c473591d6ad64c6" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -3107,7 +3119,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "479ca8adacdd7ce8f1fb39ce9ecccbfe93a3f1344b3d0d97f20bc0196208f62b" dependencies = [ "proc-macro2", - "syn", + "syn 2.0.117", ] [[package]] @@ -3136,7 +3148,7 @@ checksum = "af066a9c399a26e020ada66a034357a868728e72cd426f3adcd35f80d88d88c8" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", "version_check", "yansi", ] @@ -3168,7 +3180,7 @@ dependencies = [ "pulldown-cmark", "pulldown-cmark-to-cmark", "regex", - "syn", + "syn 2.0.117", "tempfile", ] @@ -3182,7 +3194,7 @@ dependencies = [ "itertools", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -3511,7 +3523,7 @@ checksum = "b7186006dcb21920990093f30e3dea63b7d6e977bf1256be20c3563a5db070da" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4060,7 +4072,7 @@ checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4127,7 +4139,7 @@ dependencies = [ "darling", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4325,6 +4337,17 @@ version = "2.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" +[[package]] +name = "syn" +version = "1.0.109" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72b64191b275b66ffe2469e8af2c1cfe3bafa67b529ead792a6d0160888b4237" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + [[package]] name = "syn" version = "2.0.117" @@ -4353,7 +4376,7 @@ checksum = "728a70f3dbaf5bab7f0c4b1ac8d7ae5ea60a4b5549c8a5914361c99147a709d2" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4422,7 +4445,7 @@ checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4433,7 +4456,7 @@ checksum = "ebc4ee7f67670e9b64d05fa4253e753e016c6c95ff35b89b7941d6b856dec1d5" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4557,7 +4580,7 @@ checksum = "385a6cb71ab9ab790c5fe8d67f1645e6c450a7ce006a33de03daa956cf70a496" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4644,7 +4667,7 @@ dependencies = [ "prettyplease", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4669,7 +4692,7 @@ dependencies = [ "prost-build", "prost-types", "quote", - "syn", + "syn 2.0.117", "tempfile", "tonic-build", ] @@ -4763,7 +4786,7 @@ checksum = "7490cfa5ec963746568740651ac6781f701c9c5ea257c58e057f3ba8cf69e8da" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -4867,7 +4890,7 @@ checksum = "27a7a9b72ba121f6f1f6c3632b85604cac41aedb5ddc70accbebb6cac83de846" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -5109,7 +5132,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn", + "syn 2.0.117", "wasm-bindgen-shared", ] @@ -5256,7 +5279,7 @@ checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -5267,7 +5290,7 @@ checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -5633,7 +5656,7 @@ dependencies = [ "heck", "indexmap 2.13.0", "prettyplease", - "syn", + "syn 2.0.117", "wasm-metadata", "wit-bindgen-core", "wit-component", @@ -5649,7 +5672,7 @@ dependencies = [ "prettyplease", "proc-macro2", "quote", - "syn", + "syn 2.0.117", "wit-bindgen-core", "wit-bindgen-rust", ] @@ -5722,7 +5745,7 @@ checksum = "b659052874eb698efe5b9e8cf382204678a0086ebf46982b79d6ca3182927e5d" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", "synstructure", ] @@ -5743,7 +5766,7 @@ checksum = "f65c489a7071a749c849713807783f70672b28094011623e200cb86dcb835953" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -5763,7 +5786,7 @@ checksum = "d71e5d6e06ab090c67b5e44993ec16b72dcbaabc526db883a360057678b48502" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", "synstructure", ] @@ -5784,7 +5807,7 @@ checksum = "3c50655cbb0fe3fc43170059e702f1ce5e19b84cec58dc87b037a09935c2f328" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -5817,7 +5840,7 @@ checksum = "eadce39539ca5cb3985590102671f2567e659fca9666581ad3411d59207951f3" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 9413dd9d..2f487e98 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -40,6 +40,7 @@ base64 = "0.22.1" bigtable_rs = { git = "https://github.com/getsentry/bigtable_rs.git", rev = "4cb75bc5e5f87204363973f6302107768e64972e" } bytes = "1.12.0" bytesize = "2.4.0" +derive-error-kind = "0.1.0" chrono = "0.4.45" console = "0.16.3" elegant-departure = "0.3.2" diff --git a/objectstore-server/docs/architecture.md b/objectstore-server/docs/architecture.md index 290de51b..21e4b0cb 100644 --- a/objectstore-server/docs/architecture.md +++ b/objectstore-server/docs/architecture.md @@ -198,7 +198,7 @@ Beyond rate limiting and the web concurrency limit, the [`StorageService`](objectstore_service::StorageService) enforces a second layer of backpressure through a concurrency limit on in-flight backend operations, configured via `service.max_concurrency`. When exceeded, requests -receive HTTP 429. See the [service architecture docs](objectstore_service) for +receive HTTP 503. See the [service architecture docs](objectstore_service) for details. ## KEDA Metrics diff --git a/objectstore-server/src/endpoints/common.rs b/objectstore-server/src/endpoints/common.rs index 06b6014e..cda8de16 100644 --- a/objectstore-server/src/endpoints/common.rs +++ b/objectstore-server/src/endpoints/common.rs @@ -6,7 +6,7 @@ use axum::Json; use axum::http::StatusCode; use axum::response::{IntoResponse, Response}; use http::HeaderValue; -use objectstore_service::error::Error as ServiceError; +use objectstore_service::error::{Error as ServiceError, ErrorKind as ServiceErrorKind}; use serde::{Deserialize, Serialize}; use thiserror::Error; @@ -92,18 +92,23 @@ impl ApiError { StatusCode::INTERNAL_SERVER_ERROR } - ApiError::Service(ServiceError::Client(_)) => StatusCode::BAD_REQUEST, - ApiError::Service(ServiceError::Metadata(_)) => StatusCode::BAD_REQUEST, - ApiError::Service(ServiceError::RangeNotSatisfiable { .. }) => { - StatusCode::RANGE_NOT_SATISFIABLE - } - ApiError::Service(ServiceError::InvalidUploadId(_)) => StatusCode::BAD_REQUEST, - ApiError::Service(ServiceError::AtCapacity) => StatusCode::TOO_MANY_REQUESTS, - ApiError::Service(ServiceError::NotImplemented) => StatusCode::NOT_IMPLEMENTED, - ApiError::Service(_) => { - objectstore_log::error!(!!self, "error handling request"); - StatusCode::INTERNAL_SERVER_ERROR - } + ApiError::Service(error) => match error.kind() { + ServiceErrorKind::ClientStream | ServiceErrorKind::InvalidInput => { + StatusCode::BAD_REQUEST + } + ServiceErrorKind::RangeNotSatisfiable => StatusCode::RANGE_NOT_SATISFIABLE, + ServiceErrorKind::BackendRateLimited => StatusCode::TOO_MANY_REQUESTS, + // Local load-shedding and backend unavailability are both surfaced as a + // temporary 503 so callers back off and retry. + ServiceErrorKind::AtCapacity + | ServiceErrorKind::BackendTimeout + | ServiceErrorKind::BackendUnavailable => StatusCode::SERVICE_UNAVAILABLE, + ServiceErrorKind::NotImplemented => StatusCode::NOT_IMPLEMENTED, + ServiceErrorKind::Internal => { + objectstore_log::error!(!!self, "error handling request"); + StatusCode::INTERNAL_SERVER_ERROR + } + }, ApiError::Internal(_) => { objectstore_log::error!(!!self, "internal error"); diff --git a/objectstore-server/src/endpoints/multipart.rs b/objectstore-server/src/endpoints/multipart.rs index f1e231d2..5d81aaa5 100644 --- a/objectstore-server/src/endpoints/multipart.rs +++ b/objectstore-server/src/endpoints/multipart.rs @@ -96,7 +96,8 @@ async fn initiate_inner( headers: HeaderMap, ) -> ApiResult { // TODO: Update time_created in `complete`, when we have a Service API to mutate metadata. - let metadata = Metadata::from_insert_headers(&headers, "").map_err(ServiceError::from)?; + let metadata = Metadata::from_insert_headers(&headers, "") + .map_err(|cause| ServiceError::metadata_client("invalid object metadata headers", cause))?; state .config diff --git a/objectstore-server/src/endpoints/objects.rs b/objectstore-server/src/endpoints/objects.rs index a08b6620..6b9b0cbe 100644 --- a/objectstore-server/src/endpoints/objects.rs +++ b/objectstore-server/src/endpoints/objects.rs @@ -43,7 +43,8 @@ async fn objects_post( headers: HeaderMap, MeteredBody(body): MeteredBody, ) -> ApiResult { - let metadata = Metadata::from_insert_headers(&headers, "").map_err(ServiceError::from)?; + let metadata = Metadata::from_insert_headers(&headers, "") + .map_err(|cause| ServiceError::metadata_client("invalid object metadata headers", cause))?; state .config @@ -88,7 +89,9 @@ async fn object_get( }; let stream = state.meter_stream(stream, &context); - let metadata_headers = metadata.to_headers("").map_err(ServiceError::from)?; + let metadata_headers = metadata + .to_headers("") + .map_err(|cause| ServiceError::metadata("failed to serialize object metadata", cause))?; let mut response = match content_range { Some(ref content_range) => { @@ -120,7 +123,9 @@ async fn object_head(service: AuthAwareService, Xt(id): Xt) -> ApiResu return Ok(StatusCode::NOT_FOUND.into_response()); }; - let headers = metadata.to_headers("").map_err(ServiceError::from)?; + let headers = metadata + .to_headers("") + .map_err(|cause| ServiceError::metadata("failed to serialize object metadata", cause))?; let mut response = (StatusCode::NO_CONTENT, headers).into_response(); insert_content_disposition(&mut response, &metadata); @@ -175,7 +180,8 @@ async fn object_put( headers: HeaderMap, MeteredBody(body): MeteredBody, ) -> ApiResult { - let metadata = Metadata::from_insert_headers(&headers, "").map_err(ServiceError::from)?; + let metadata = Metadata::from_insert_headers(&headers, "") + .map_err(|cause| ServiceError::metadata_client("invalid object metadata headers", cause))?; let ObjectId { context, key } = id; diff --git a/objectstore-server/src/extractors/body.rs b/objectstore-server/src/extractors/body.rs index b4d3d772..6411a1dc 100644 --- a/objectstore-server/src/extractors/body.rs +++ b/objectstore-server/src/extractors/body.rs @@ -5,7 +5,7 @@ use std::convert::Infallible; use axum::extract::{FromRequest, FromRequestParts, Path, Request}; use futures_util::{StreamExt, TryStreamExt}; use objectstore_service::id::ObjectContext; -use objectstore_service::stream::{ClientError, ClientStream}; +use objectstore_service::stream::{ClientStream, ClientStreamError}; use super::id::ContextParams; use crate::state::ServiceState; @@ -37,7 +37,10 @@ impl FromRequest for MeteredBody { usecase: params.usecase, scopes: params.scopes, }; - let stream = body.into_data_stream().map_err(ClientError::new).boxed(); + let stream = body + .into_data_stream() + .map_err(ClientStreamError::new) + .boxed(); let stream = state.meter_stream(stream, &context).boxed(); Ok(Self(stream)) } diff --git a/objectstore-server/tests/limits.rs b/objectstore-server/tests/limits.rs index 4ed1fc06..4aee966e 100644 --- a/objectstore-server/tests/limits.rs +++ b/objectstore-server/tests/limits.rs @@ -534,9 +534,9 @@ async fn test_bandwidth_scope_pct_limit() -> Result<()> { } #[tokio::test] -async fn test_batch_at_capacity_returns_429() -> Result<()> { +async fn test_batch_at_capacity_returns_503() -> Result<()> { // With max_concurrency=0 the service has no permits available, so - // BatchExecutor::new() returns AtCapacity and the endpoint responds 429. + // BatchExecutor::new() returns AtCapacity and the endpoint responds 503. let server = TestServer::with_config(Config { service: Service { max_concurrency: 0 }, auth: AuthZ { @@ -566,8 +566,8 @@ async fn test_batch_at_capacity_returns_429() -> Result<()> { assert_eq!( response.status(), - reqwest::StatusCode::TOO_MANY_REQUESTS, - "expected 429 when service has no available permits" + reqwest::StatusCode::SERVICE_UNAVAILABLE, + "expected 503 when service has no available permits" ); Ok(()) diff --git a/objectstore-service/Cargo.toml b/objectstore-service/Cargo.toml index df6b771a..0132a1c6 100644 --- a/objectstore-service/Cargo.toml +++ b/objectstore-service/Cargo.toml @@ -16,6 +16,7 @@ base64 = { workspace = true } bigtable_rs = { workspace = true } bytes = { workspace = true } chrono = { workspace = true } +derive-error-kind = { workspace = true } futures-util = { workspace = true } gcp_auth = { workspace = true } humantime = { workspace = true } diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 73ad50b3..f4a8e233 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -179,6 +179,13 @@ The service applies backpressure to protect backends from overload. Rather than queueing work when capacity is exhausted, the service rejects operations immediately so the caller can shed load or retry. +# Error Classification + +The service error type exposes a fine-grained [`ErrorKind`](error::ErrorKind) +that describes the *cause* of a failure independent of HTTP semantics, so API +layers can map failures without matching every specific variant and remain free +to decide how each kind maps onto a status code. + ## Concurrency Limit A semaphore caps the total number of in-flight backend operations across all diff --git a/objectstore-service/src/backend/extensions.rs b/objectstore-service/src/backend/extensions.rs index 53bd0459..6ae28160 100644 --- a/objectstore-service/src/backend/extensions.rs +++ b/objectstore-service/src/backend/extensions.rs @@ -10,8 +10,7 @@ use reqwest::{Response, header}; use serde::Deserialize; use tracing::Instrument; -use crate::error::{BackendDetail, Error, Result}; -use crate::stream; +use crate::error::{BackendDetail, BackendError, Error, Result}; /// Extension trait that sends a request inside a tracing span. pub trait SendTraced { @@ -81,7 +80,7 @@ struct XmlApiError { /// Use [`check_error`](Self::check_error) instead of /// [`error_for_status`](reqwest::Response::error_for_status) to avoid losing the response body on /// 4xx/5xx errors. The method parses the structured error body (JSON or XML) and returns an -/// [`Error::BackendResponse`] with the extracted error code and message. +/// [`Error::Backend`] with the extracted error code and message. /// /// Implemented for both [`reqwest::Response`] and `Result` so it can be /// chained directly. @@ -94,7 +93,7 @@ pub trait ResponseExt { /// [`reqwest::Response::error_for_status`]. /// /// When called on `Result`, transport errors are - /// wrapped as [`Error::Reqwest`] with the same context string. + /// wrapped as [`Error::Backend`] with the same context string. async fn check_error(self, context: &'static str) -> Result; /// Drains the response body of a response we are otherwise done with. @@ -130,11 +129,7 @@ impl ResponseExt for Response { return Err(Error::reqwest(context, e)); }; - Err(Error::BackendResponse { - context, - status, - detail, - }) + Err(BackendError::from_status(context, status, detail).into()) } async fn drain_body(mut self) { @@ -146,10 +141,7 @@ impl ResponseExt for Result { async fn check_error(self, context: &'static str) -> Result { match self { Ok(resp) => resp.check_error(context).await, - Err(e) => Err(match stream::unpack_client_error(&e) { - Some(ce) => Error::Client(ce), - None => Error::reqwest(context, e), - }), + Err(e) => Err(Error::reqwest(context, e)), } } diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 728a9da5..20b0563d 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -189,14 +189,29 @@ impl GcsObject { .metadata .remove(&GcsMetaKey::Expiration) .map(|s| s.parse()) - .transpose()? + .transpose() + .map_err(|cause| { + Error::metadata( + "GCS: failed to parse expiration policy from object metadata", + cause, + ) + })? .unwrap_or_default(); let origin = self.metadata.remove(&GcsMetaKey::Origin); let filename = self.metadata.remove(&GcsMetaKey::Filename); let content_type = self.content_type; - let compression = self.content_encoding.map(|s| s.parse()).transpose()?; + let compression = self + .content_encoding + .map(|s| s.parse()) + .transpose() + .map_err(|cause| { + Error::metadata( + "GCS: failed to parse compression from object metadata", + cause, + ) + })?; let size = self .size .map(|size| size.parse()) @@ -373,17 +388,16 @@ fn insert_gcs_meta_header( Ok(()) } -/// Returns `true` if the error is a transient reqwest failure worth retrying. +/// Returns `true` if the error is a transient backend failure worth retrying. fn error_is_retryable(error: &Error) -> bool { - match error { - Error::Reqwest { cause, .. } => { - cause.is_timeout() - || cause.is_connect() - || cause.is_request() - || cause.status().is_some_and(status_is_retryable) - } - Error::BackendResponse { status, .. } => status_is_retryable(*status), - _ => false, + let Error::Backend(e) = error else { + return false; + }; + match e.status() { + Some(status) => status_is_retryable(status), + None => e + .cause() + .is_some_and(|c| c.is_timeout() || c.is_connect() || c.is_request()), } } @@ -638,10 +652,8 @@ impl Backend for GcsBackend { // NB: Ensure the order of these fields and that a content-type is attached to them. Both // are required by the GCS API. - let metadata_json = serde_json::to_string(&gcs_metadata).map_err(|cause| Error::Serde { - context: "failed to serialize metadata for GCS upload".to_string(), - cause, - })?; + let metadata_json = serde_json::to_string(&gcs_metadata) + .map_err(|cause| Error::serde("GCS: failed to serialize metadata", cause))?; let multipart = multipart::Form::new() .part( diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index af0b233d..b2fa7404 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -21,7 +21,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, Part, PartNumber, UploadId, UploadPartResponse, }; -use crate::stream::{self, ClientStream}; +use crate::stream::ClientStream; /// Configuration for [`LocalFsBackend`]. /// @@ -96,19 +96,12 @@ impl Backend for LocalFsBackend { let mut reader = pin!(StreamReader::new(stream)); let mut writer = BufWriter::new(file); - let metadata_json = serde_json::to_string(metadata).map_err(|cause| Error::Serde { - context: "failed to serialize metadata".to_string(), - cause, - })?; + let metadata_json = serde_json::to_string(metadata) + .map_err(|cause| Error::serde("local_fs: failed to serialize metadata", cause))?; writer.write_all(metadata_json.as_bytes()).await?; writer.write_all(b"\n").await?; - tokio::io::copy(&mut reader, &mut writer) - .await - .map_err(|e| match stream::unpack_client_error(&e) { - Some(ce) => Error::Client(ce), - None => e.into(), - })?; + tokio::io::copy(&mut reader, &mut writer).await?; writer.flush().await?; let file = writer.into_inner(); @@ -136,11 +129,8 @@ impl Backend for LocalFsBackend { let mut metadata_line = String::new(); reader.read_line(&mut metadata_line).await?; let file_len = reader.get_ref().metadata().await?.len(); - let mut metadata: Metadata = - serde_json::from_str(metadata_line.trim_end()).map_err(|cause| Error::Serde { - context: "failed to deserialize metadata".to_string(), - cause, - })?; + let mut metadata: Metadata = serde_json::from_str(metadata_line.trim_end()) + .map_err(|cause| Error::serde("local_fs: failed to deserialize metadata", cause))?; let payload_size = file_len .checked_sub(metadata_line.len() as u64) .ok_or_else(|| Error::generic("local-fs file corrupted: shorter than header"))?; @@ -199,10 +189,8 @@ impl MultipartUploadBackend for LocalFsBackend { tokio::fs::create_dir_all(&dir).await?; let meta_path = dir.join("metadata.json"); - let metadata_json = serde_json::to_string(metadata).map_err(|cause| Error::Serde { - context: "failed to serialize multipart metadata".to_string(), - cause, - })?; + let metadata_json = serde_json::to_string(metadata) + .map_err(|cause| Error::serde("local_fs: failed to serialize metadata", cause))?; tokio::fs::write(meta_path, metadata_json).await?; Ok(upload_id) @@ -229,10 +217,8 @@ impl MultipartUploadBackend for LocalFsBackend { "uploaded_at": SystemTime::now(), "size": content_length, }); - let header_line = serde_json::to_string(&header).map_err(|cause| Error::Serde { - context: "failed to serialize part header".to_string(), - cause, - })?; + let header_line = serde_json::to_string(&header) + .map_err(|cause| Error::serde("local_fs: failed to serialize part header", cause))?; let part_path = dir.join(format!("{part_number}.part")); let file = OpenOptions::new() @@ -247,12 +233,7 @@ impl MultipartUploadBackend for LocalFsBackend { writer.write_all(header_line.as_bytes()).await?; writer.write_all(b"\n").await?; - let _bytes_copied = tokio::io::copy(&mut reader, &mut writer) - .await - .map_err(|e| match stream::unpack_client_error(&e) { - Some(ce) => Error::Client(ce), - None => e.into(), - })?; + let _bytes_copied = tokio::io::copy(&mut reader, &mut writer).await?; // TODO: validate bytes_copied against content_length and return a BadRequest-style // error. Needs a service-layer error variant that maps to HTTP 400 without abusing @@ -300,9 +281,8 @@ impl MultipartUploadBackend for LocalFsBackend { let mut header_line = String::new(); reader.read_line(&mut header_line).await?; let header: serde_json::Value = - serde_json::from_str(header_line.trim_end()).map_err(|cause| Error::Serde { - context: "failed to deserialize part header".to_string(), - cause, + serde_json::from_str(header_line.trim_end()).map_err(|cause| { + Error::serde("local_fs: failed to deserialize part header", cause) })?; parts.push(Part { @@ -359,11 +339,8 @@ impl MultipartUploadBackend for LocalFsBackend { // Read metadata let meta_path = dir.join("metadata.json"); let meta_bytes = tokio::fs::read(&meta_path).await?; - let metadata: Metadata = - serde_json::from_slice(&meta_bytes).map_err(|cause| Error::Serde { - context: "failed to deserialize multipart metadata".to_string(), - cause, - })?; + let metadata: Metadata = serde_json::from_slice(&meta_bytes) + .map_err(|cause| Error::serde("local_fs: failed to deserialize metadata", cause))?; // TODO: validate that parts are in ascending part_number order and reject with // InvalidPartOrder if not (matches S3/GCS behavior). Needs a proper client error variant. @@ -383,9 +360,8 @@ impl MultipartUploadBackend for LocalFsBackend { let mut header_line = String::new(); reader.read_line(&mut header_line).await?; let header: serde_json::Value = - serde_json::from_str(header_line.trim_end()).map_err(|cause| Error::Serde { - context: "failed to deserialize part header".to_string(), - cause, + serde_json::from_str(header_line.trim_end()).map_err(|cause| { + Error::serde("local_fs: failed to deserialize part header", cause) })?; let stored_etag = header["etag"].as_str().unwrap_or(""); @@ -411,10 +387,8 @@ impl MultipartUploadBackend for LocalFsBackend { .await?; let mut writer = BufWriter::new(file); - let metadata_json = serde_json::to_string(&metadata).map_err(|cause| Error::Serde { - context: "failed to serialize metadata".to_string(), - cause, - })?; + let metadata_json = serde_json::to_string(&metadata) + .map_err(|cause| Error::serde("local_fs: failed to serialize metadata", cause))?; writer.write_all(metadata_json.as_bytes()).await?; writer.write_all(b"\n").await?; diff --git a/objectstore-service/src/backend/s3_compatible.rs b/objectstore-service/src/backend/s3_compatible.rs index 12c1f665..cbc8005b 100644 --- a/objectstore-service/src/backend/s3_compatible.rs +++ b/objectstore-service/src/backend/s3_compatible.rs @@ -175,10 +175,7 @@ where let response = builder .send_traced() .await - .map_err(|cause| Error::Reqwest { - context: "S3: failed to send request".to_string(), - cause, - })?; + .map_err(|cause| Error::reqwest("S3: failed to send request", cause))?; if response.status() == StatusCode::NOT_FOUND { objectstore_log::debug!("Object not found"); @@ -205,7 +202,8 @@ where let response = response.check_error("S3: failed to get object").await?; let headers = response.headers(); - let mut metadata = Metadata::from_headers(headers, GCS_CUSTOM_PREFIX)?; + let mut metadata = Metadata::from_headers(headers, GCS_CUSTOM_PREFIX) + .map_err(|cause| Error::metadata("S3: failed to parse object metadata", cause))?; let content_range = if response.status() == StatusCode::PARTIAL_CONTENT { let range = headers @@ -262,7 +260,14 @@ where format!("/{}/{}", self.bucket, id.as_storage_path()), ) .header("x-goog-metadata-directive", "REPLACE") - .headers(metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)?) + .headers( + metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX).map_err(|cause| { + Error::metadata( + "S3: failed to serialize object metadata for TTI update", + cause, + ) + })?, + ) .send_traced() .await .check_error("S3: update expiration time") @@ -312,7 +317,11 @@ impl Backend for S3CompatibleBackend { objectstore_log::debug!("Writing to s3_compatible backend"); self.request(Method::PUT, self.object_url(id)) .await? - .headers(metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)?) + .headers( + metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX).map_err(|cause| { + Error::metadata("S3: failed to serialize object metadata for upload", cause) + })?, + ) .body(Body::wrap_stream(stream)) .send_traced() .await @@ -353,10 +362,7 @@ impl Backend for S3CompatibleBackend { .await? .send_traced() .await - .map_err(|cause| Error::Reqwest { - context: "S3: failed to send delete request".to_string(), - cause, - })?; + .map_err(|cause| Error::reqwest("S3: failed to send delete request", cause))?; // Do not error for objects that do not exist. if response.status() == StatusCode::NOT_FOUND { diff --git a/objectstore-service/src/error.rs b/objectstore-service/src/error.rs index 9643dd43..5f53df88 100644 --- a/objectstore-service/src/error.rs +++ b/objectstore-service/src/error.rs @@ -3,14 +3,49 @@ //! [`Error`] covers I/O, serialization, HTTP, metadata, authentication, //! and backend-specific failures. [`Result`] is the corresponding alias. +// Needed as the `.kind()` methods generated by `derive(ErrorKind)` currently lack documentation. +#![allow(missing_docs)] + use std::any::Any; +use std::borrow::Cow; use std::fmt; use objectstore_log::Level; use reqwest::StatusCode; use thiserror::Error as ThisError; -use crate::stream::ClientError; +use crate::stream::{self, ClientStreamError}; + +/// The category of a service error. +/// +/// These kinds describe the *cause* of a failure, independent of any HTTP +/// semantics. It is up to the API layer to decide how to map each kind onto a +/// status code, and which kinds to group together. +/// +/// Users should rely on the kind for classification and handling, rather than matching on +/// specific variants. +#[derive(Copy, Clone, Debug, Eq, PartialEq)] +pub enum ErrorKind { + /// Error originating from a client-supplied request body stream. + ClientStream, + /// Malformed client input (bad metadata, JSON, upload id, …). + InvalidInput, + /// The requested byte range is not satisfiable for the object's size. + RangeNotSatisfiable, + /// This service instance has reached its own concurrency limit and is + /// shedding load. + AtCapacity, + /// A storage backend rejected the operation with a rate-limit response. + BackendRateLimited, + /// A storage backend timed out. + BackendTimeout, + /// A storage backend is temporarily unavailable. + BackendUnavailable, + /// Operation unsupported by this service instance. + NotImplemented, + /// Internal failure. + Internal, +} /// Structured error detail parsed from a backend HTTP error response. /// @@ -45,67 +80,42 @@ impl fmt::Display for BackendDetail { } /// Error type for service operations. -#[derive(Debug, ThisError)] +#[derive(Debug, ThisError, derive_error_kind::ErrorKind)] +#[error_kind(ErrorKind)] pub enum Error { /// IO errors related to payload streaming or file operations. #[error("i/o error: {0}")] - Io(#[from] std::io::Error), + #[error_kind(ErrorKind, Internal)] + Io(std::io::Error), /// Error originating from a client-supplied input stream. - /// - /// Indicates the client is at fault (e.g. dropped connection mid-upload) and should - /// map to a 4xx response rather than a 5xx. #[error("error reading client stream: {0}")] - Client(#[from] ClientError), + #[error_kind(ErrorKind, ClientStream)] + ClientStream(#[from] ClientStreamError), /// Errors related to de/serialization. - #[error("serde error: {context}")] - Serde { - /// Context describing what was being serialized/deserialized. - context: String, - /// The underlying serde error. - #[source] - cause: serde_json::Error, - }, - - /// All errors stemming from the reqwest client, used in multiple backends to send requests to - /// e.g. GCP APIs. - /// These can be network errors encountered when sending the requests, but can also indicate - /// errors returned by the API itself. - #[error("reqwest error: {context}")] - Reqwest { - /// Context describing the request that failed. - context: String, - /// The underlying reqwest error. - #[source] - cause: reqwest::Error, - }, + #[error(transparent)] + #[error_kind(transparent)] + Serde(#[from] SerdeError), - /// An HTTP error response from a storage backend (e.g., GCS, S3). - /// - /// Unlike [`Reqwest`](Self::Reqwest), which covers transport-level failures, this variant - /// captures application-level error responses where the server returned a 4xx/5xx status code - /// along with a structured error body. - #[error("{context} ({status}). {detail}")] - BackendResponse { - /// Context describing the request that failed. - context: &'static str, - /// The HTTP status code returned by the backend. - status: StatusCode, - /// Parsed error code and message from the response body. - detail: BackendDetail, - }, + /// Errors from storage backend HTTP calls, transport-level or application-level. + #[error(transparent)] + #[error_kind(transparent)] + Backend(#[from] BackendError), /// Errors related to de/serialization and parsing of object metadata. - #[error("metadata error: {0}")] - Metadata(#[from] objectstore_types::metadata::Error), + #[error(transparent)] + #[error_kind(transparent)] + Metadata(#[from] MetadataError), /// Errors encountered when attempting to authenticate with GCP. #[error("GCP authentication error: {0}")] + #[error_kind(ErrorKind, Internal)] GcpAuth(#[from] gcp_auth::Error), /// A spawned service task panicked. #[error("service task failed: {0}")] + #[error_kind(ErrorKind, Internal)] Panic(String), /// A spawned service task was dropped before it could deliver its result. @@ -113,6 +123,7 @@ pub enum Error { /// This is an unexpected condition that can occur when the runtime drops the task for unknown /// reasons. #[error("task dropped")] + #[error_kind(ErrorKind, Internal)] Dropped, /// A redirect tombstone was encountered at a place where it is not supported. @@ -120,10 +131,12 @@ pub enum Error { /// This indicates a caller bug — tombstone-aware reads must go through the /// [`HighVolumeBackend`](crate::backend::common::HighVolumeBackend) methods. #[error("unexpected tombstone")] + #[error_kind(ErrorKind, Internal)] UnexpectedTombstone, /// The requested byte range is not satisfiable for the object's size. #[error("range not satisfiable (object size: {total} bytes)")] + #[error_kind(ErrorKind, RangeNotSatisfiable)] RangeNotSatisfiable { /// Total size of the object in bytes. total: u64, @@ -131,11 +144,13 @@ pub enum Error { /// The service has reached its concurrency limit and cannot accept more operations. #[error("concurrency limit reached")] + #[error_kind(ErrorKind, AtCapacity)] AtCapacity, /// Any other error stemming from one of the storage backends, which might be specific to that /// backend or to a certain operation. #[error("storage backend error: {context}")] + #[error_kind(ErrorKind, Internal)] Generic { /// Context describing the operation that failed. context: String, @@ -146,10 +161,12 @@ pub enum Error { /// The functionality is not implemented by this instance of the service. #[error("not implemented")] + #[error_kind(ErrorKind, NotImplemented)] NotImplemented, - /// Invalid upload ID (e.g. path traversal attempt). + /// Invalid upload ID for a multipart upload. #[error(transparent)] + #[error_kind(ErrorKind, InvalidInput)] InvalidUploadId(#[from] objectstore_types::multipart::InvalidUploadId), } @@ -166,20 +183,60 @@ impl Error { Self::Panic(msg) } - /// Creates an [`Error::Reqwest`] from a reqwest error with context. - pub fn reqwest(context: impl Into, cause: reqwest::Error) -> Self { - Self::Reqwest { + /// Creates an [`Error`] from a reqwest error with context, categorizing it by status code or + /// transport-level failure signal. + pub fn reqwest(context: impl Into>, cause: reqwest::Error) -> Self { + if let Some(client_error) = stream::unpack_client_error(&cause) { + return Self::ClientStream(client_error); + } + + Self::Backend(BackendError::from_reqwest(context, cause)) + } + + /// Creates an [`Error::Serde`] from a serde error with context, categorizing it as an internal + /// error. + pub fn serde(context: impl Into>, cause: serde_json::Error) -> Self { + Self::Serde(SerdeError { + kind: ErrorKind::Internal, context: context.into(), cause, - } + }) } - /// Creates an [`Error::Serde`] from a serde error with context. - pub fn serde(context: impl Into, cause: serde_json::Error) -> Self { - Self::Serde { + /// Creates an [`Error::Serde`] from a serde error with context, categorizing it as a client + /// error. + pub fn serde_client(context: impl Into>, cause: serde_json::Error) -> Self { + Self::Serde(SerdeError { + kind: ErrorKind::InvalidInput, context: context.into(), cause, - } + }) + } + + /// Creates an [`Error::Metadata`] from a metadata error with context, categorizing it as an + /// internal error. + pub fn metadata( + context: impl Into>, + cause: objectstore_types::metadata::Error, + ) -> Self { + Self::Metadata(MetadataError { + kind: ErrorKind::Internal, + context: context.into(), + cause, + }) + } + + /// Creates an [`Error::Metadata`] from a metadata error with context, categorizing it as a + /// client error. + pub fn metadata_client( + context: impl Into>, + cause: objectstore_types::metadata::Error, + ) -> Self { + Self::Metadata(MetadataError { + kind: ErrorKind::InvalidInput, + context: context.into(), + cause, + }) } /// Creates an [`Error::Generic`] with a context string and no cause. @@ -192,27 +249,175 @@ impl Error { /// Returns the appropriate log level for this error. pub fn level(&self) -> Level { - match self { - // Malformed client input at DEBUG level - Self::Client(_) => Level::DEBUG, - Self::Metadata(_) => Level::DEBUG, - Self::RangeNotSatisfiable { .. } => Level::DEBUG, - // Like rate limits, we treat capacity errors as warnings - Self::AtCapacity => Level::WARN, - // All other errors are service or backend failures - Self::Io(_) => Level::ERROR, - Self::Serde { .. } => Level::ERROR, - Self::Reqwest { .. } => Level::ERROR, - Self::BackendResponse { .. } => Level::ERROR, - Self::GcpAuth(_) => Level::ERROR, - Self::Panic(_) => Level::ERROR, - Self::Dropped => Level::ERROR, - Self::UnexpectedTombstone => Level::ERROR, - Self::NotImplemented => Level::ERROR, - Self::InvalidUploadId(_) => Level::DEBUG, - Self::Generic { .. } => Level::ERROR, + match self.kind() { + ErrorKind::ClientStream | ErrorKind::InvalidInput | ErrorKind::RangeNotSatisfiable => { + Level::DEBUG + } + ErrorKind::AtCapacity + | ErrorKind::BackendRateLimited + | ErrorKind::BackendTimeout + | ErrorKind::BackendUnavailable => Level::WARN, + ErrorKind::NotImplemented | ErrorKind::Internal => Level::ERROR, + } + } +} + +impl From for Error { + fn from(cause: std::io::Error) -> Self { + match stream::unpack_client_error(&cause) { + Some(client_error) => Self::ClientStream(client_error), + None => Self::Io(cause), + } + } +} + +/// A backend HTTP or transport-level error. +/// +/// Unifies transport-level failures (no response received, e.g. connection reset or timeout) and +/// HTTP error responses (4xx/5xx), with or without a structured error body, under a single kind +/// classification. +#[derive(Debug)] +pub struct BackendError { + kind: ErrorKind, + context: Cow<'static, str>, + status: Option, + detail: BackendDetail, + cause: Option, +} + +impl BackendError { + /// Creates a [`BackendError`] for an HTTP error response, categorizing it by status code. + pub(crate) fn from_status( + context: impl Into>, + status: StatusCode, + detail: BackendDetail, + ) -> Self { + Self { + kind: kind_for_status(status), + context: context.into(), + status: Some(status), + detail, + cause: None, + } + } + + /// Creates a [`BackendError`] from a reqwest error, categorizing it by status code if the + /// error carries one, or by transport-level failure signal otherwise. + fn from_reqwest(context: impl Into>, cause: reqwest::Error) -> Self { + let status = cause.status(); + let kind = match status { + Some(status) => kind_for_status(status), + None if cause.is_timeout() => ErrorKind::BackendTimeout, + None if cause.is_connect() || cause.is_request() => ErrorKind::BackendUnavailable, + None => ErrorKind::Internal, + }; + + Self { + kind, + context: context.into(), + status, + detail: BackendDetail::none(), + cause: Some(cause), } } + + /// Returns the service-level category for this backend error. + pub fn kind(&self) -> ErrorKind { + self.kind + } + + /// Returns the HTTP status code, if a response was received. + pub fn status(&self) -> Option { + self.status + } + + /// Returns the parsed error code and message from the response body, if any. + pub fn detail(&self) -> &BackendDetail { + &self.detail + } + + /// Returns the underlying reqwest error, if this is a transport-level failure. + pub fn cause(&self) -> Option<&reqwest::Error> { + self.cause.as_ref() + } +} + +impl fmt::Display for BackendError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}", self.context)?; + if let Some(status) = self.status { + write!(f, " ({status})")?; + } + if !self.detail.code.is_empty() || !self.detail.message.is_empty() { + write!(f, ". {}", self.detail)?; + } + Ok(()) + } +} + +impl std::error::Error for BackendError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + self.cause + .as_ref() + .map(|e| e as &(dyn std::error::Error + 'static)) + } +} + +/// Maps an HTTP status code from a storage backend to a service-level [`ErrorKind`]. +fn kind_for_status(status: StatusCode) -> ErrorKind { + match status { + StatusCode::TOO_MANY_REQUESTS => ErrorKind::BackendRateLimited, + StatusCode::REQUEST_TIMEOUT | StatusCode::GATEWAY_TIMEOUT => ErrorKind::BackendTimeout, + StatusCode::INTERNAL_SERVER_ERROR + | StatusCode::BAD_GATEWAY + | StatusCode::SERVICE_UNAVAILABLE => ErrorKind::BackendUnavailable, + status if status.is_client_error() => ErrorKind::InvalidInput, + _ => ErrorKind::Internal, + } +} + +/// Serde error with context and kind. +#[derive(Debug, ThisError)] +#[error("serde error: {context}")] +pub struct SerdeError { + kind: ErrorKind, + context: Cow<'static, str>, + #[source] + cause: serde_json::Error, +} + +impl SerdeError { + /// Returns the service-level category for this serde error. + pub fn kind(&self) -> ErrorKind { + self.kind + } + + /// Returns the underlying serde error. + pub fn cause(&self) -> &serde_json::Error { + &self.cause + } +} + +/// Metadata error with context and kind. +#[derive(Debug, ThisError)] +#[error("metadata error: {context}")] +pub struct MetadataError { + kind: ErrorKind, + context: Cow<'static, str>, + #[source] + cause: objectstore_types::metadata::Error, +} + +impl MetadataError { + /// Returns the service-level category for this metadata error. + pub fn kind(&self) -> ErrorKind { + self.kind + } + + /// Returns the underlying metadata error. + pub fn cause(&self) -> &objectstore_types::metadata::Error { + &self.cause + } } /// Result type for service operations. diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index 74479aff..4f112e1b 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -200,7 +200,9 @@ impl StorageService { metadata: Metadata, stream: ClientStream, ) -> Result { - metadata.validate()?; + metadata + .validate() + .map_err(|cause| Error::metadata("invalid object metadata", cause))?; let id = ObjectId::optional(context, key); let inner = Arc::clone(&self.inner); self.spawn("insert", async move { @@ -254,7 +256,9 @@ impl StorageService { id: ObjectId, metadata: Metadata, ) -> Result { - metadata.validate()?; + metadata + .validate() + .map_err(|cause| Error::metadata("invalid object metadata", cause))?; self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported let inner = self.inner.clone(); self.spawn("initiate_multipart", async move { diff --git a/objectstore-service/src/stream.rs b/objectstore-service/src/stream.rs index efaf9c43..16be5860 100644 --- a/objectstore-service/src/stream.rs +++ b/objectstore-service/src/stream.rs @@ -4,7 +4,7 @@ //! cover the two directions of data flow: //! //! - [`ClientStream`] — incoming data from a client PUT request body. Uses -//! [`ClientError`] as the error type so backends can distinguish a broken +//! [`ClientStreamError`] as the error type so backends can distinguish a broken //! client connection from a backend I/O failure (400 vs 500). //! - [`PayloadStream`] — outgoing data returned from //! [`Backend::get_object`](crate::backend::common::Backend::get_object). @@ -30,10 +30,10 @@ pub type PayloadStream = BoxStream<'static, io::Result>; /// [`unpack_client_error`] to return a 4xx response rather than treating /// it as a 5xx backend failure. #[derive(Clone, Debug)] -pub struct ClientError(Arc); +pub struct ClientStreamError(Arc); -impl ClientError { - /// Creates a new [`ClientError`] wrapping `err`. +impl ClientStreamError { + /// Creates a new [`ClientStreamError`] wrapping `err`. pub fn new(err: E) -> Self where E: Error + Send + Sync + 'static, @@ -42,47 +42,47 @@ impl ClientError { } } -impl fmt::Display for ClientError { +impl fmt::Display for ClientStreamError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { self.0.fmt(f) } } -impl Error for ClientError { +impl Error for ClientStreamError { fn source(&self) -> Option<&(dyn Error + 'static)> { self.0.source() } } /// Required by [`tokio_util::io::StreamReader`] in the local filesystem backend. -impl From for io::Error { - fn from(err: ClientError) -> Self { +impl From for io::Error { + fn from(err: ClientStreamError) -> Self { io::Error::other(err) } } /// Incoming byte stream from a client PUT request body. /// -/// Uses [`ClientError`] as the error type so that a dropped or interrupted +/// Uses [`ClientStreamError`] as the error type so that a dropped or interrupted /// client connection is distinguishable from a backend I/O failure. Backends -/// that detect a [`ClientError`] (via [`unpack_client_error`]) can surface it -/// as [`crate::error::Error::Client`], which the server maps to HTTP 400 rather -/// than 500. +/// that detect a [`ClientStreamError`] (via [`unpack_client_error`]) can surface +/// it as [`crate::error::Error::ClientStream`], which the server maps to HTTP +/// 400 rather than 500. /// /// Use [`single`] to construct a single-chunk `ClientStream` from an owned value. -pub type ClientStream = BoxStream<'static, Result>; +pub type ClientStream = BoxStream<'static, Result>; -/// Walks the source chain of `err` looking for a [`ClientError`]. +/// Walks the source chain of `err` looking for a [`ClientStreamError`]. /// /// At each step, two locations are checked: /// -/// - **Direct**: the error itself is a `ClientError`. +/// - **Direct**: the error itself is a `ClientStreamError`. /// - **Packed in `io::Error`**: the error is an `io::Error` whose custom inner -/// value is a `ClientError`. +/// value is a `ClientStreamError`. /// /// Use this in `put_object` implementations to reclassify body-stream errors -/// as [`crate::error::Error::Client`] instead of an opaque server error. -pub fn unpack_client_error(err: &E) -> Option +/// as [`crate::error::Error::ClientStream`] instead of an opaque server error. +pub fn unpack_client_error(err: &E) -> Option where E: Error + 'static, { @@ -96,7 +96,7 @@ where None => s, }; - if let Some(client_error) = target.downcast_ref::() { + if let Some(client_error) = target.downcast_ref::() { return Some(client_error.clone()); }