Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
181 changes: 114 additions & 67 deletions objectstore-service/src/backend/gcs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -902,35 +902,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?;

Ok(xml.try_into()?)
let body = resp
.bytes()
.await
.map_err(|e| Error::reqwest("GCS: read initiate multipart body", e))?;

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
Comment thread
lcian marked this conversation as resolved.
}

#[tracing::instrument(level = "debug", skip(self, content_md5, body))]
Expand All @@ -943,6 +953,7 @@ impl MultipartUploadBackend for GcsBackend {
content_md5: Option<&str>,
body: ClientStream,
) -> Result<UploadPartResponse> {
// 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()
Expand Down Expand Up @@ -998,26 +1009,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))]
Expand All @@ -1030,16 +1047,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
Comment thread
lcian marked this conversation as resolved.
}

#[tracing::instrument(level = "debug", skip(self, parts))]
Expand All @@ -1053,32 +1086,43 @@ impl MultipartUploadBackend for GcsBackend {
let mut url = self.xml_object_url(id)?;
url.query_pairs_mut().append_pair("uploadId", upload_id);

// Serialize once — the payload is small and identical across retries.
let body = XmlCompleteMultipartUpload::from(parts);
let xml = quick_xml::se::to_string(&body).map_err(|e| Error::Generic {
context: "GCS: failed to serialize complete multipart request".into(),
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))?;
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);
// Semantic complete errors (InvalidPart, etc.) come back as XML bodies under a
// success-ish status in our mapping, and are returned as Ok(Some(error)) — not
// retried, since they are permanent request errors.
let error = quick_xml::de::from_reader::<_, XmlError>(body.as_ref())
.ok()
.map(Into::into);

Ok(error)
Ok(error)
}
})
.await
}
}

Expand Down Expand Up @@ -1730,6 +1774,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(())
}

Expand Down
Loading