Skip to content
Closed
Show file tree
Hide file tree
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
37 changes: 20 additions & 17 deletions objectstore-service/src/backend/bigtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1133,27 +1133,30 @@ impl HighVolumeBackend for BigTableBackend {
/// required by BigTable, the resulting timestamp has millisecond precision, with the last digits at
/// 0.
fn ttl_to_micros(ttl: Duration, from: SystemTime) -> Result<i64> {
let deadline = from.checked_add(ttl).ok_or_else(|| Error::Generic {
context: format!(
let deadline = from.checked_add(ttl).ok_or_else(|| {
Error::generic(format!(
"TTL duration overflow: {} plus {}s cannot be represented as SystemTime",
humantime::format_rfc3339_seconds(from),
ttl.as_secs()
),
cause: None,
))
})?;
let millis = deadline
.duration_since(SystemTime::UNIX_EPOCH)
.map_err(|e| Error::Generic {
context: format!(
"unable to get duration since UNIX_EPOCH for SystemTime {}",
humantime::format_rfc3339_seconds(deadline)
),
cause: Some(Box::new(e)),
.map_err(|e| {
Error::generic_cause(
format!(
"unable to get duration since UNIX_EPOCH for SystemTime {}",
humantime::format_rfc3339_seconds(deadline)
),
e,
)
})?
.as_millis();
(millis * 1000).try_into().map_err(|e| Error::Generic {
context: format!("failed to convert {millis}ms to i64 microseconds"),
cause: Some(Box::new(e)),
(millis * 1000).try_into().map_err(|e| {
Error::generic_cause(
format!("failed to convert {millis}ms to i64 microseconds"),
e,
)
})
}

Expand Down Expand Up @@ -1186,10 +1189,10 @@ where
Ok(res) => return Ok(res),
Err(e) if retry_count >= REQUEST_RETRY_COUNT || !is_retryable(&e) => {
objectstore_metrics::count!("bigtable.failures", action = context);
return Err(Error::Generic {
context: format!("Bigtable: `{context}` failed"),
cause: Some(Box::new(e)),
});
return Err(Error::generic_cause(
format!("Bigtable: `{context}` failed"),
e,
));
}
Err(e) => {
retry_count += 1;
Expand Down
108 changes: 43 additions & 65 deletions objectstore-service/src/backend/gcs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -203,9 +203,8 @@ impl GcsObject {
.size
.map(|size| size.parse())
.transpose()
.map_err(|e| Error::Generic {
context: "GCS: failed to parse size from object metadata".to_string(),
cause: Some(Box::new(e)),
.map_err(|e| {
Error::generic_cause("GCS: failed to parse size from object metadata", e)
})?;
let time_created = self.time_created;

Expand All @@ -215,12 +214,9 @@ impl GcsObject {
if let GcsMetaKey::Custom(custom_key) = key {
custom.insert(custom_key, value);
} else {
return Err(Error::Generic {
context: format!(
"GCS: unexpected built-in metadata key in object metadata: {key}"
),
cause: None,
});
return Err(Error::generic(format!(
"GCS: unexpected built-in metadata key in object metadata: {key}"
)));
}
}

Expand Down Expand Up @@ -314,23 +310,19 @@ fn metadata_to_gcs_headers(metadata: &Metadata) -> Result<header::HeaderMap> {
let formatted = humantime::format_rfc3339_seconds(custom_time);
headers.insert(
HeaderName::from_static("x-goog-custom-time"),
formatted.to_string().parse().map_err(|e| Error::Generic {
context: "GCS: invalid custom-time header value".into(),
cause: Some(Box::new(e)),
})?,
formatted
.to_string()
.parse()
.map_err(|e| Error::generic_cause("GCS: invalid custom-time header value", e))?,
);
}

if let Some(compression) = metadata.compression {
headers.insert(
header::CONTENT_ENCODING,
compression
.to_string()
.parse()
.map_err(|e| Error::Generic {
context: "GCS: invalid content-encoding header value".into(),
cause: Some(Box::new(e)),
})?,
compression.to_string().parse().map_err(|e| {
Error::generic_cause("GCS: invalid content-encoding header value", e)
})?,
);
}

Expand Down Expand Up @@ -364,13 +356,11 @@ fn insert_gcs_meta_header(
) -> Result<()> {
let header_name = format!("x-goog-meta-{key}");
headers.insert(
HeaderName::try_from(&header_name).map_err(|e| Error::Generic {
context: format!("GCS: invalid header name: {header_name}"),
cause: Some(Box::new(e)),
HeaderName::try_from(&header_name).map_err(|e| {
Error::generic_cause(format!("GCS: invalid header name: {header_name}"), e)
})?,
value.parse().map_err(|e| Error::Generic {
context: format!("GCS: invalid header value for {header_name}"),
cause: Some(Box::new(e)),
value.parse().map_err(|e| {
Error::generic_cause(format!("GCS: invalid header value for {header_name}"), e)
})?,
);
Ok(())
Expand Down Expand Up @@ -438,12 +428,11 @@ impl GcsBackend {

let path = id.as_storage_path().to_string();
url.path_segments_mut()
.map_err(|()| Error::Generic {
context: format!(
.map_err(|()| {
Error::generic(format!(
"GCS: invalid endpoint URL, {} cannot be a base",
self.endpoint
),
cause: None,
))
})?
.extend(&["storage", "v1", "b", &self.bucket, "o", &path]);

Expand All @@ -455,12 +444,11 @@ impl GcsBackend {
let mut url = self.endpoint.clone();

url.path_segments_mut()
.map_err(|()| Error::Generic {
context: format!(
.map_err(|()| {
Error::generic(format!(
"GCS: invalid endpoint URL, {} cannot be a base",
self.endpoint
),
cause: None,
))
})?
.extend(&["upload", "storage", "v1", "b", &self.bucket, "o"]);

Expand All @@ -480,12 +468,11 @@ impl GcsBackend {
fn xml_object_url(&self, id: &ObjectId) -> Result<Url> {
let mut url = self.endpoint.clone();
{
let mut segments = url.path_segments_mut().map_err(|()| Error::Generic {
context: format!(
let mut segments = url.path_segments_mut().map_err(|()| {
Error::generic(format!(
"GCS: invalid endpoint URL, {} cannot be a base",
self.endpoint
),
cause: None,
))
})?;
segments.push(&self.bucket);
for part in id.as_storage_path().to_string().split('/') {
Expand Down Expand Up @@ -636,10 +623,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("failed to serialize metadata for GCS upload", cause))?;

let multipart = multipart::Form::new()
.part(
Expand All @@ -652,9 +637,11 @@ impl Backend for GcsBackend {
"media",
multipart::Part::stream(Body::wrap_stream(stream))
.mime_str(&metadata.content_type)
.map_err(|e| Error::Generic {
context: format!("invalid mime type: {}", &metadata.content_type),
cause: Some(Box::new(e)),
.map_err(|e| {
Error::generic_cause(
format!("invalid mime type: {}", &metadata.content_type),
e,
)
})?,
);

Expand Down Expand Up @@ -707,12 +694,9 @@ impl Backend for GcsBackend {
match total {
Some(total) => return Err(Error::RangeNotSatisfiable { total }),
None => {
return Err(Error::Generic {
context: format!(
"GCS: 416 response with invalid Content-Range: {raw:?}"
),
cause: None,
});
return Err(Error::generic(format!(
"GCS: 416 response with invalid Content-Range: {raw:?}"
)));
}
}
}
Expand All @@ -728,9 +712,8 @@ impl Backend for GcsBackend {
.get(header::CONTENT_RANGE)
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse::<ContentRange>().ok())
.ok_or_else(|| Error::Generic {
context: "GCS: 206 response missing valid Content-Range header".to_owned(),
cause: None,
.ok_or_else(|| {
Error::generic("GCS: 206 response missing valid Content-Range header")
})?,
)
} else {
Expand Down Expand Up @@ -920,9 +903,8 @@ impl MultipartUploadBackend for GcsBackend {
.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)),
.map_err(|e| {
Error::generic_cause("GCS: failed to parse initiate multipart response", e)
})?;

Ok(xml.try_into()?)
Expand Down Expand Up @@ -1000,11 +982,8 @@ impl MultipartUploadBackend for GcsBackend {
.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_cause("GCS: failed to parse list parts response", e))?;

Ok(xml.into())
}
Expand Down Expand Up @@ -1041,9 +1020,8 @@ impl MultipartUploadBackend for GcsBackend {
url.query_pairs_mut().append_pair("uploadId", upload_id);

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 xml = quick_xml::se::to_string(&body).map_err(|e| {
Error::generic_cause("GCS: failed to serialize complete multipart request", e)
})?;

let resp = self
Expand Down
52 changes: 16 additions & 36 deletions objectstore-service/src/backend/local_fs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -96,10 +96,8 @@ 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("failed to serialize metadata", cause))?;
writer.write_all(metadata_json.as_bytes()).await?;
writer.write_all(b"\n").await?;

Expand Down Expand Up @@ -136,11 +134,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("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"))?;
Expand Down Expand Up @@ -199,10 +194,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("failed to serialize multipart metadata", cause))?;
tokio::fs::write(meta_path, metadata_json).await?;

Ok(upload_id)
Expand All @@ -229,10 +222,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("failed to serialize part header", cause))?;

let part_path = dir.join(format!("{part_number}.part"));
let file = OpenOptions::new()
Expand Down Expand Up @@ -299,11 +290,8 @@ impl MultipartUploadBackend for LocalFsBackend {
let mut reader = BufReader::new(file);
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,
})?;
let header: serde_json::Value = serde_json::from_str(header_line.trim_end())
.map_err(|cause| Error::serde("failed to deserialize part header", cause))?;

parts.push(Part {
part_number: pn,
Expand Down Expand Up @@ -359,11 +347,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("failed to deserialize multipart 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.
Expand All @@ -382,11 +367,8 @@ impl MultipartUploadBackend for LocalFsBackend {
let mut reader = BufReader::new(file);
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,
})?;
let header: serde_json::Value = serde_json::from_str(header_line.trim_end())
.map_err(|cause| Error::serde("failed to deserialize part header", cause))?;

let stored_etag = header["etag"].as_str().unwrap_or("");
if stored_etag != completed.etag {
Expand All @@ -411,10 +393,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("failed to serialize metadata", cause))?;
writer.write_all(metadata_json.as_bytes()).await?;
writer.write_all(b"\n").await?;

Expand Down
Loading
Loading