mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-25 21:46:50 +00:00
perf(get): reduce request entry allocations (#6029)
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -1585,7 +1585,7 @@ impl ECStore {
|
|||||||
) -> Result<GetObjectReader> {
|
) -> Result<GetObjectReader> {
|
||||||
check_get_obj_args(bucket, object)?;
|
check_get_obj_args(bucket, object)?;
|
||||||
|
|
||||||
let object = encode_dir_object(object);
|
let object = rustfs_utils::path::encode_dir_object_ref(object);
|
||||||
let mut opts = opts.clone();
|
let mut opts = opts.clone();
|
||||||
let read_lock_guard = self
|
let read_lock_guard = self
|
||||||
.acquire_object_read_lock_if_needed("get_object", bucket, &object, &mut opts)
|
.acquire_object_read_lock_if_needed("get_object", bucket, &object, &mut opts)
|
||||||
@@ -1593,14 +1593,14 @@ impl ECStore {
|
|||||||
|
|
||||||
let reader = if self.single_pool() {
|
let reader = if self.single_pool() {
|
||||||
self.pools[0]
|
self.pools[0]
|
||||||
.get_object_reader(bucket, object.as_str(), range, h, &opts)
|
.get_object_reader(bucket, object.as_ref(), range, h, &opts)
|
||||||
.await?
|
.await?
|
||||||
} else {
|
} else {
|
||||||
let (_, idx) = self
|
let (_, idx) = self
|
||||||
.get_latest_accessible_object_info_with_idx(bucket, &object, &opts)
|
.get_latest_accessible_object_info_with_idx(bucket, &object, &opts)
|
||||||
.await?;
|
.await?;
|
||||||
self.pools[idx]
|
self.pools[idx]
|
||||||
.get_object_reader(bucket, object.as_str(), range, h, &opts)
|
.get_object_reader(bucket, object.as_ref(), range, h, &opts)
|
||||||
.await?
|
.await?
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -12,6 +12,7 @@
|
|||||||
// See the License for the specific language governing permissions and
|
// See the License for the specific language governing permissions and
|
||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
|
use std::borrow::Cow;
|
||||||
use std::path::Component;
|
use std::path::Component;
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
@@ -44,14 +45,19 @@ pub fn has_suffix(s: &str, suffix: &str) -> bool {
|
|||||||
/// If the object name ends with a slash, it is considered a directory object.
|
/// If the object name ends with a slash, it is considered a directory object.
|
||||||
/// The trailing slash is removed and `GLOBAL_DIR_SUFFIX` is appended.
|
/// The trailing slash is removed and `GLOBAL_DIR_SUFFIX` is appended.
|
||||||
/// If it does not end with a slash, the name is returned as is.
|
/// If it does not end with a slash, the name is returned as is.
|
||||||
pub fn encode_dir_object(object: &str) -> String {
|
pub fn encode_dir_object_ref(object: &str) -> Cow<'_, str> {
|
||||||
if has_suffix(object, SLASH_SEPARATOR) {
|
if has_suffix(object, SLASH_SEPARATOR) {
|
||||||
format!("{}{}", object.trim_end_matches(SLASH_SEPARATOR), GLOBAL_DIR_SUFFIX)
|
Cow::Owned(format!("{}{}", object.trim_end_matches(SLASH_SEPARATOR), GLOBAL_DIR_SUFFIX))
|
||||||
} else {
|
} else {
|
||||||
object.to_string()
|
Cow::Borrowed(object)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Owned compatibility wrapper for callers that retain or mutate the encoded name.
|
||||||
|
pub fn encode_dir_object(object: &str) -> String {
|
||||||
|
encode_dir_object_ref(object).into_owned()
|
||||||
|
}
|
||||||
|
|
||||||
/// Checks if the given object name represents a directory object.
|
/// Checks if the given object name represents a directory object.
|
||||||
///
|
///
|
||||||
/// Returns true if the object name ends with `GLOBAL_DIR_SUFFIX`.
|
/// Returns true if the object name ends with `GLOBAL_DIR_SUFFIX`.
|
||||||
@@ -602,6 +608,17 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
use proptest::prelude::*;
|
use proptest::prelude::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn encode_dir_object_ref_borrows_objects_and_encodes_directories() {
|
||||||
|
let object = "prefix/object";
|
||||||
|
let encoded = encode_dir_object_ref(object);
|
||||||
|
assert!(matches!(encoded, Cow::Borrowed(value) if value == object));
|
||||||
|
|
||||||
|
let encoded = encode_dir_object_ref("prefix/directory/");
|
||||||
|
assert!(matches!(encoded, Cow::Owned(ref value) if value == "prefix/directory__XLDIR__"));
|
||||||
|
assert_eq!(encode_dir_object("prefix/directory/"), encoded);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_trim_etag() {
|
fn test_trim_etag() {
|
||||||
// Test with quoted ETag
|
// Test with quoted ETag
|
||||||
|
|||||||
@@ -2044,6 +2044,18 @@ struct GetObjectResumeContext {
|
|||||||
identity: GetObjectResumeIdentity,
|
identity: GetObjectResumeIdentity,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn get_object_store_headers(request_headers: &HeaderMap) -> HeaderMap {
|
||||||
|
let mut headers = HeaderMap::new();
|
||||||
|
for name in [SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER] {
|
||||||
|
if let Some(value) = request_headers.get(name) {
|
||||||
|
let mut value = value.clone();
|
||||||
|
value.set_sensitive(true);
|
||||||
|
headers.insert(name, value);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
headers
|
||||||
|
}
|
||||||
|
|
||||||
impl GetObjectResumeContext {
|
impl GetObjectResumeContext {
|
||||||
#[allow(clippy::too_many_arguments)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
fn new(
|
fn new(
|
||||||
@@ -2061,17 +2073,9 @@ impl GetObjectResumeContext {
|
|||||||
{
|
{
|
||||||
opts.version_id = Some(version_id.to_string());
|
opts.version_id = Some(version_id.to_string());
|
||||||
}
|
}
|
||||||
let mut ssec_headers = HeaderMap::new();
|
// Store spans record their header argument at debug level. Retain only
|
||||||
for name in [SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER] {
|
// the SSE-C inputs needed to reopen the reader and keep them redacted.
|
||||||
if let Some(value) = request_headers.get(name) {
|
let ssec_headers = get_object_store_headers(request_headers);
|
||||||
// The store's instrumented spans record the header argument at
|
|
||||||
// debug level; mark the replayed values sensitive so the SSE-C
|
|
||||||
// key is redacted there on every resume attempt.
|
|
||||||
let mut value = value.clone();
|
|
||||||
value.set_sensitive(true);
|
|
||||||
ssec_headers.insert(name, value);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Self {
|
Self {
|
||||||
store,
|
store,
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
@@ -4455,6 +4459,7 @@ impl DefaultObjectUsecase {
|
|||||||
) -> S3Result<GetObjectPreparedRead> {
|
) -> S3Result<GetObjectPreparedRead> {
|
||||||
let read_start = std::time::Instant::now();
|
let read_start = std::time::Instant::now();
|
||||||
let read_stage_start = rustfs_io_metrics::get_stage_metrics_enabled().then_some(read_start);
|
let read_stage_start = rustfs_io_metrics::get_stage_metrics_enabled().then_some(read_start);
|
||||||
|
let store_headers = get_object_store_headers(&req.headers);
|
||||||
let cache_adapter = self.object_data_cache();
|
let cache_adapter = self.object_data_cache();
|
||||||
if cache_adapter.is_disabled() || !cache_adapter.materialize_fill_enabled() {
|
if cache_adapter.is_disabled() || !cache_adapter.materialize_fill_enabled() {
|
||||||
let io_planning = Self::acquire_get_object_io_planning(
|
let io_planning = Self::acquire_get_object_io_planning(
|
||||||
@@ -4469,7 +4474,7 @@ impl DefaultObjectUsecase {
|
|||||||
.await?;
|
.await?;
|
||||||
let reader = track_object_read_setup(
|
let reader = track_object_read_setup(
|
||||||
object_traffic_health.as_deref(),
|
object_traffic_health.as_deref(),
|
||||||
store.get_object_reader(bucket, key, rs.clone(), req.headers.clone(), opts),
|
store.get_object_reader(bucket, key, rs.clone(), store_headers.clone(), opts),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(map_get_object_reader_error)?;
|
.map_err(map_get_object_reader_error)?;
|
||||||
@@ -4596,7 +4601,7 @@ impl DefaultObjectUsecase {
|
|||||||
drop(metadata_admission.take());
|
drop(metadata_admission.take());
|
||||||
let outcome = coordinate_cold_fill(&coordinator, cache_key, waiter_deadline, Some(proposed_producer_deadline), {
|
let outcome = coordinate_cold_fill(&coordinator, cache_key, waiter_deadline, Some(proposed_producer_deadline), {
|
||||||
let adapter = &cache_adapter;
|
let adapter = &cache_adapter;
|
||||||
let headers = &req.headers;
|
let headers = &store_headers;
|
||||||
let store = &store;
|
let store = &store;
|
||||||
let range = &rs;
|
let range = &rs;
|
||||||
let object_traffic_health = &object_traffic_health;
|
let object_traffic_health = &object_traffic_health;
|
||||||
@@ -4760,7 +4765,7 @@ impl DefaultObjectUsecase {
|
|||||||
.ok_or_else(|| s3_error!(InternalError, "prepared metadata admission is unavailable"))?;
|
.ok_or_else(|| s3_error!(InternalError, "prepared metadata admission is unavailable"))?;
|
||||||
let reader = track_object_read_setup(
|
let reader = track_object_read_setup(
|
||||||
object_traffic_health.as_deref(),
|
object_traffic_health.as_deref(),
|
||||||
prepared.with_headers(req.headers.clone()).into_reader(),
|
prepared.with_headers(store_headers.clone()).into_reader(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(map_get_object_reader_error)?;
|
.map_err(map_get_object_reader_error)?;
|
||||||
@@ -4785,14 +4790,14 @@ impl DefaultObjectUsecase {
|
|||||||
.map_err(map_get_object_reader_error)?;
|
.map_err(map_get_object_reader_error)?;
|
||||||
track_object_read_setup(
|
track_object_read_setup(
|
||||||
object_traffic_health.as_deref(),
|
object_traffic_health.as_deref(),
|
||||||
prepared.with_headers(req.headers.clone()).into_reader(),
|
prepared.with_headers(store_headers.clone()).into_reader(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(map_get_object_reader_error)?
|
.map_err(map_get_object_reader_error)?
|
||||||
} else {
|
} else {
|
||||||
track_object_read_setup(
|
track_object_read_setup(
|
||||||
object_traffic_health.as_deref(),
|
object_traffic_health.as_deref(),
|
||||||
store.get_object_reader(bucket, key, rs.clone(), req.headers.clone(), opts),
|
store.get_object_reader(bucket, key, rs.clone(), store_headers, opts),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(map_get_object_reader_error)?
|
.map_err(map_get_object_reader_error)?
|
||||||
@@ -13722,6 +13727,11 @@ mod tests {
|
|||||||
request_headers.insert(SSEC_KEY_MD5_HEADER, HeaderValue::from_static("bWQ1"));
|
request_headers.insert(SSEC_KEY_MD5_HEADER, HeaderValue::from_static("bWQ1"));
|
||||||
request_headers.insert(http::header::AUTHORIZATION, HeaderValue::from_static("AWS4-HMAC-SHA256 Credential=test"));
|
request_headers.insert(http::header::AUTHORIZATION, HeaderValue::from_static("AWS4-HMAC-SHA256 Credential=test"));
|
||||||
request_headers.insert("x-amz-security-token", HeaderValue::from_static("session-token"));
|
request_headers.insert("x-amz-security-token", HeaderValue::from_static("session-token"));
|
||||||
|
let store_headers = get_object_store_headers(&request_headers);
|
||||||
|
assert_eq!(store_headers.len(), 3, "only store-consumed SSE-C headers are forwarded");
|
||||||
|
assert!(store_headers.values().all(HeaderValue::is_sensitive));
|
||||||
|
assert!(store_headers.get(http::header::AUTHORIZATION).is_none());
|
||||||
|
assert!(store_headers.get("x-amz-security-token").is_none());
|
||||||
let plain_info = ObjectInfo {
|
let plain_info = ObjectInfo {
|
||||||
size: 11,
|
size: 11,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
|
|||||||
Reference in New Issue
Block a user