From e374cc0ea90e70e0b5d18deb33717e85ccf78aa6 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 24 Jul 2026 16:40:54 +0200 Subject: [PATCH 1/4] fix(service): Retry GCS multipart initiate/list/abort Wrap initiate, list, and abort in the existing GCS request retry helper so transient 5xx/transport failures are absorbed internally. Leave put-part and complete unretried for now (streamed body / separate PR). Make abort idempotent on 404 after a lost successful response. --- objectstore-service/docs/architecture.md | 5 + objectstore-service/src/backend/gcs.rs | 141 +++++++++++++++-------- 2 files changed, 96 insertions(+), 50 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 0b745d46..8c09ebe9 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -198,6 +198,11 @@ delegate to the [`MultipartUploadBackend`](backend::common::MultipartUploadBacke trait, accessed via [`Backend::as_multipart_upload_backend`](backend::common::Backend::as_multipart_upload_backend). Multipart operations share the same concurrency limiter as regular operations. +On the GCS backend, initiate/list/abort use the same request-level retry helper +as JSON API calls for transient transport/`5xx` failures. Part upload is not +retried (the body is a one-shot stream). Abort treats HTTP 404 as success so a +retry after a lost successful abort response remains idempotent. + ## Streaming Concurrency The [`streaming`](streaming) module provides [`StreamExecutor`](streaming::StreamExecutor) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 876c254a..7c7db175 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -886,6 +886,11 @@ impl From for crate::multipart::CompleteMultipartError { /// XXX: Any change that affects this implementation should be manually tested against real GCS. /// That's because the fork of [storage-testbench](https://github.com/googleapis/storage-testbench) /// that we test against has an incomplete implementation of the XML multipart API that likely doesn't match GCS's behavior in many cases. +/// +/// Request-level retries use the same internal retry helper as JSON API calls for +/// initiate/list/abort. Part upload is intentionally not retried: the body is streamed and not +/// buffered. Abort treats HTTP 404 as success (idempotent if the session is already gone, e.g. +/// after a successful abort whose response was lost). #[async_trait::async_trait] impl MultipartUploadBackend for GcsBackend { #[tracing::instrument(level = "debug", fields(?id), skip_all)] @@ -898,35 +903,45 @@ impl MultipartUploadBackend for GcsBackend { let mut url = self.xml_object_url(id)?; url.set_query(Some("uploads")); - let mut builder = self - .request(Method::POST, url) - .await? - .header(header::CONTENT_TYPE, metadata.content_type.as_ref()) - .header(header::CONTENT_LENGTH, "0"); - + let content_type = metadata.content_type.clone(); let meta_headers = metadata_to_gcs_headers(metadata)?; - for (name, value) in &meta_headers { - builder = builder.header(name, value); - } - let resp = builder - .send_traced() - .await - .check_error("GCS: initiate multipart upload") - .await?; + self.with_retry("initiate_multipart", || { + let url = url.clone(); + let content_type = content_type.clone(); + let meta_headers = meta_headers.clone(); + async move { + let mut builder = self + .request(Method::POST, url) + .await? + .header(header::CONTENT_TYPE, content_type.as_ref()) + .header(header::CONTENT_LENGTH, "0"); - let body = resp - .bytes() - .await - .map_err(|e| Error::reqwest("GCS: read initiate multipart body", e))?; + for (name, value) in &meta_headers { + builder = builder.header(name, value); + } - let xml: XmlInitiateMultipartUploadResponse = quick_xml::de::from_reader(body.as_ref()) - .map_err(|e| Error::Generic { - context: "GCS: failed to parse initiate multipart response".to_owned(), - cause: Some(Box::new(e)), - })?; + let resp = builder + .send_traced() + .await + .check_error("GCS: initiate multipart upload") + .await?; + + let body = resp + .bytes() + .await + .map_err(|e| Error::reqwest("GCS: read initiate multipart body", e))?; - Ok(xml.try_into()?) + let xml: XmlInitiateMultipartUploadResponse = + quick_xml::de::from_reader(body.as_ref()).map_err(|e| Error::Generic { + context: "GCS: failed to parse initiate multipart response".to_owned(), + cause: Some(Box::new(e)), + })?; + + xml.try_into() + } + }) + .await } #[tracing::instrument(level = "debug", skip(self, content_md5, body))] @@ -939,6 +954,7 @@ impl MultipartUploadBackend for GcsBackend { content_md5: Option<&str>, body: ClientStream, ) -> Result { + // Not retried: the part body is a one-shot stream and is not buffered. objectstore_log::debug!("Uploading part to GCS backend"); let mut url = self.xml_object_url(id)?; url.query_pairs_mut() @@ -994,26 +1010,32 @@ impl MultipartUploadBackend for GcsBackend { } } - let resp = self - .request(Method::GET, url) - .await? - .send_traced() - .await - .check_error("GCS: list parts") - .await?; + self.with_retry("list_parts", || { + let url = url.clone(); + async move { + let resp = self + .request(Method::GET, url) + .await? + .send_traced() + .await + .check_error("GCS: list parts") + .await?; - let body = resp - .bytes() - .await - .map_err(|e| Error::reqwest("GCS: read list parts body", e))?; + let body = resp + .bytes() + .await + .map_err(|e| Error::reqwest("GCS: read list parts body", e))?; - let xml: XmlListPartsResponse = - quick_xml::de::from_reader(body.as_ref()).map_err(|e| Error::Generic { - context: "GCS: failed to parse list parts response".to_owned(), - cause: Some(Box::new(e)), - })?; + let xml: XmlListPartsResponse = + quick_xml::de::from_reader(body.as_ref()).map_err(|e| Error::Generic { + context: "GCS: failed to parse list parts response".to_owned(), + cause: Some(Box::new(e)), + })?; - Ok(xml.into()) + Ok(xml.into()) + } + }) + .await } #[tracing::instrument(level = "debug", skip(self))] @@ -1026,16 +1048,32 @@ impl MultipartUploadBackend for GcsBackend { let mut url = self.xml_object_url(id)?; url.query_pairs_mut().append_pair("uploadId", upload_id); - self.request(Method::DELETE, url) - .await? - .send_traced() - .await - .check_error("GCS: abort multipart upload") - .await? - .drain_body() - .await; + self.with_retry("abort_multipart", || { + let url = url.clone(); + async move { + let resp = self + .request(Method::DELETE, url) + .await? + .send_traced() + .await + .map_err(|e| Error::reqwest("GCS: abort multipart upload", e))?; - Ok(()) + // Idempotent: absent upload means the postcondition is already satisfied + // (including retry after a successful abort whose response was lost). + if resp.status() == StatusCode::NOT_FOUND { + resp.drain_body().await; + return Ok(()); + } + + resp.check_error("GCS: abort multipart upload") + .await? + .drain_body() + .await; + + Ok(()) + } + }) + .await } #[tracing::instrument(level = "debug", skip(self, parts))] @@ -1761,6 +1799,9 @@ mod tests { let result = backend.get_object(&id, None).await?; assert!(result.is_none(), "object should not exist after abort"); + // A second abort should still succeed (idempotent 404 handling). + backend.abort_multipart(&id, &upload_id).await?; + Ok(()) } From ff2e2204bc48daefada3494352d1cdb99b8eb401 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 24 Jul 2026 16:43:25 +0200 Subject: [PATCH 2/4] improve --- objectstore-service/docs/architecture.md | 5 ----- objectstore-service/src/backend/gcs.rs | 6 ------ 2 files changed, 11 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 8c09ebe9..0b745d46 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -198,11 +198,6 @@ delegate to the [`MultipartUploadBackend`](backend::common::MultipartUploadBacke trait, accessed via [`Backend::as_multipart_upload_backend`](backend::common::Backend::as_multipart_upload_backend). Multipart operations share the same concurrency limiter as regular operations. -On the GCS backend, initiate/list/abort use the same request-level retry helper -as JSON API calls for transient transport/`5xx` failures. Part upload is not -retried (the body is a one-shot stream). Abort treats HTTP 404 as success so a -retry after a lost successful abort response remains idempotent. - ## Streaming Concurrency The [`streaming`](streaming) module provides [`StreamExecutor`](streaming::StreamExecutor) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 7c7db175..c65f9c49 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -886,11 +886,6 @@ impl From for crate::multipart::CompleteMultipartError { /// XXX: Any change that affects this implementation should be manually tested against real GCS. /// That's because the fork of [storage-testbench](https://github.com/googleapis/storage-testbench) /// that we test against has an incomplete implementation of the XML multipart API that likely doesn't match GCS's behavior in many cases. -/// -/// Request-level retries use the same internal retry helper as JSON API calls for -/// initiate/list/abort. Part upload is intentionally not retried: the body is streamed and not -/// buffered. Abort treats HTTP 404 as success (idempotent if the session is already gone, e.g. -/// after a successful abort whose response was lost). #[async_trait::async_trait] impl MultipartUploadBackend for GcsBackend { #[tracing::instrument(level = "debug", fields(?id), skip_all)] @@ -954,7 +949,6 @@ impl MultipartUploadBackend for GcsBackend { content_md5: Option<&str>, body: ClientStream, ) -> Result { - // Not retried: the part body is a one-shot stream and is not buffered. objectstore_log::debug!("Uploading part to GCS backend"); let mut url = self.xml_object_url(id)?; url.query_pairs_mut() From 88aa47d3c76614509b186f4b77043603a8db55bc Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 24 Jul 2026 16:59:38 +0200 Subject: [PATCH 3/4] improve --- objectstore-service/src/backend/gcs.rs | 28 ++++++++++++++------------ 1 file changed, 15 insertions(+), 13 deletions(-) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index c65f9c49..1ac96859 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -898,25 +898,27 @@ impl MultipartUploadBackend for GcsBackend { let mut url = self.xml_object_url(id)?; url.set_query(Some("uploads")); - let content_type = metadata.content_type.clone(); - let meta_headers = metadata_to_gcs_headers(metadata)?; + let mut headers = metadata_to_gcs_headers(metadata)?; + headers.insert( + header::CONTENT_TYPE, + metadata.content_type.parse().map_err(|e| Error::Generic { + context: "GCS: invalid content-type header value".into(), + cause: Some(Box::new(e)), + })?, + ); + headers.insert( + header::CONTENT_LENGTH, + header::HeaderValue::from_static("0"), + ); self.with_retry("initiate_multipart", || { let url = url.clone(); - let content_type = content_type.clone(); - let meta_headers = meta_headers.clone(); + let headers = headers.clone(); async move { - let mut builder = self + let resp = self .request(Method::POST, url) .await? - .header(header::CONTENT_TYPE, content_type.as_ref()) - .header(header::CONTENT_LENGTH, "0"); - - for (name, value) in &meta_headers { - builder = builder.header(name, value); - } - - let resp = builder + .headers(headers) .send_traced() .await .check_error("GCS: initiate multipart upload") From 0b6fca4c311f7861c2df2d757916870e5db0459c Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 27 Jul 2026 13:42:03 +0200 Subject: [PATCH 4/4] dont need to handle 404; retry complete too --- objectstore-service/src/backend/gcs.rs | 54 +++++++++++++++----------- 1 file changed, 31 insertions(+), 23 deletions(-) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 1ac96859..3bc48b41 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -1054,12 +1054,9 @@ impl MultipartUploadBackend for GcsBackend { .await .map_err(|e| Error::reqwest("GCS: abort multipart upload", e))?; - // Idempotent: absent upload means the postcondition is already satisfied - // (including retry after a successful abort whose response was lost). - if resp.status() == StatusCode::NOT_FOUND { - resp.drain_body().await; - return Ok(()); - } + // XXX: real S3 would return 404 here if the upload has been recently completed and we + // would have to handle it. It turns out GCS returns 204 instead, so we don't need to + // handle that case. resp.check_error("GCS: abort multipart upload") .await? @@ -1089,26 +1086,37 @@ impl MultipartUploadBackend for GcsBackend { cause: Some(Box::new(e)), })?; - let resp = self - .request(Method::POST, url) - .await? - .header(header::CONTENT_TYPE, "application/xml") - .body(xml) - .send_traced() - .await - .check_error("GCS: complete multipart upload") - .await?; + self.with_retry("complete_multipart", || { + let url = url.clone(); + let xml = xml.clone(); + async move { + let resp = self + .request(Method::POST, url) + .await? + .header(header::CONTENT_TYPE, "application/xml") + .body(xml) + .send_traced() + .await + .check_error("GCS: complete multipart upload") + .await?; - let body = resp - .bytes() - .await - .map_err(|e| Error::reqwest("GCS: read complete multipart body", e))?; + // XXX: real S3 would return 404 here if the upload has been recently completed and we + // would have to handle it. It turns out GCS returns 200 instead, so we don't need to + // handle that case. - let error = quick_xml::de::from_reader::<_, XmlError>(body.as_ref()) - .ok() - .map(Into::into); + let body = resp + .bytes() + .await + .map_err(|e| Error::reqwest("GCS: read complete multipart body", e))?; + + let error = quick_xml::de::from_reader::<_, XmlError>(body.as_ref()) + .ok() + .map(Into::into); - Ok(error) + Ok(error) + } + }) + .await } }