mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-15 17:43:13 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7f23a1ba91 | |||
| 1619c4be60 |
@@ -190,6 +190,17 @@ pub(crate) const GET_METADATA_CACHE_REASON_VERSION_SUSPENDED: &str = "version_su
|
||||
pub(crate) const GET_METADATA_CACHE_REASON_VERSIONED: &str = "versioned";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA: &str = "conflicting_metadata";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER: &str = "delete_marker";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY: &str = "data_read_inline_body_verify";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED: &str = "data_read_inline_deleted";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY: &str = "data_read_inline_geometry";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH: &str = "data_read_inline_identity_mismatch";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD: &str = "data_read_inline_missing_payload";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD: &str = "data_read_inline_missing_shard";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE: &str = "data_read_inline_not_inline";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE: &str = "data_read_inline_part_shape";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE: &str = "data_read_inline_remote";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE: &str = "data_read_inline_size";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED: &str = "data_read_inline_transformed";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_ERROR: &str = "error";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM: &str = "insufficient_quorum";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_NOT_FOUND: &str = "not_found";
|
||||
@@ -551,6 +562,32 @@ mod tests {
|
||||
assert_eq!(GET_METADATA_CACHE_REASON_VERSIONED, "versioned");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, "conflicting_metadata");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER, "delete_marker");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY,
|
||||
"data_read_inline_body_verify"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, "data_read_inline_deleted");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY, "data_read_inline_geometry");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
"data_read_inline_identity_mismatch"
|
||||
);
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
"data_read_inline_missing_payload"
|
||||
);
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD,
|
||||
"data_read_inline_missing_shard"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE, "data_read_inline_not_inline");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, "data_read_inline_part_shape");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE, "data_read_inline_remote");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, "data_read_inline_size");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED,
|
||||
"data_read_inline_transformed"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_ERROR, "error");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, "insufficient_quorum");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_NOT_FOUND, "not_found");
|
||||
|
||||
@@ -32,15 +32,22 @@ use crate::diagnostics::get::{
|
||||
GET_METADATA_CACHE_REASON_NOT_READ_DATA, GET_METADATA_CACHE_REASON_PART_NUMBER,
|
||||
GET_METADATA_CACHE_REASON_RAW_DATA_MOVEMENT_READ, GET_METADATA_CACHE_REASON_USABLE, GET_METADATA_CACHE_REASON_VERSION_ID,
|
||||
GET_METADATA_CACHE_REASON_VERSION_SUSPENDED, GET_METADATA_CACHE_REASON_VERSIONED,
|
||||
GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
|
||||
GET_METADATA_EARLY_STOP_REASON_ERROR, GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_NOT_FOUND, GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST,
|
||||
GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM, GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND, GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND,
|
||||
GET_METADATA_RESPONSE_ERROR, GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT,
|
||||
GET_METADATA_RESPONSE_VALID, GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_OBJECT_PATH_DIRECT_MEMORY, GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK,
|
||||
GET_STAGE_DECODE, GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP,
|
||||
GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED,
|
||||
GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER, GET_METADATA_EARLY_STOP_REASON_ERROR,
|
||||
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, GET_METADATA_EARLY_STOP_REASON_NOT_FOUND,
|
||||
GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM, GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND,
|
||||
GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND, GET_METADATA_RESPONSE_ERROR,
|
||||
GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT, GET_METADATA_RESPONSE_VALID,
|
||||
GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING, GET_OBJECT_PATH_DIRECT_MEMORY,
|
||||
GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK, GET_STAGE_DECODE,
|
||||
GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP,
|
||||
GET_STAGE_READER_SETUP_DROP_PENDING, GET_STAGE_READER_SETUP_SCHEDULE, GET_STAGE_READER_SETUP_WAIT_QUORUM,
|
||||
GET_STAGE_READER_TASK_BITROT_READER_INIT, GET_STAGE_READER_TASK_FILE_OPEN, GET_STAGE_READER_TASK_READER_CONSTRUCTION,
|
||||
GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure,
|
||||
@@ -652,36 +659,48 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo
|
||||
&& left.erasure.distribution == right.erasure.distribution
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) async fn data_read_early_stop_inline_body_verified(
|
||||
pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
candidate: &FileInfo,
|
||||
parts_metadata: &[FileInfo],
|
||||
disks: &[Option<DiskStore>],
|
||||
) -> bool {
|
||||
if !candidate.inline_data()
|
||||
|| candidate.is_compressed()
|
||||
) -> Option<&'static str> {
|
||||
if !candidate.inline_data() {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE);
|
||||
}
|
||||
if candidate.is_compressed()
|
||||
|| candidate
|
||||
.metadata
|
||||
.keys()
|
||||
.any(|key| rustfs_utils::http::is_object_encryption_marker(key))
|
||||
|| candidate.is_remote()
|
||||
|| candidate.deleted
|
||||
|| candidate.size <= 0
|
||||
|| candidate.parts.len() != 1
|
||||
|| !candidate.has_valid_erasure_geometry()
|
||||
{
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED);
|
||||
}
|
||||
if candidate.is_remote() {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE);
|
||||
}
|
||||
if candidate.deleted {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED);
|
||||
}
|
||||
if candidate.size <= 0 {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
}
|
||||
if candidate.parts.len() != 1 {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
}
|
||||
if !candidate.has_valid_erasure_geometry() {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
}
|
||||
|
||||
let Ok(object_size) = usize::try_from(candidate.size) else {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
};
|
||||
if candidate.parts.first().is_none_or(|part| part.size != object_size) {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
}
|
||||
if !can_try_inline_data_shards_direct(object_size, candidate.erasure.block_size) {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
}
|
||||
|
||||
let Ok(erasure) = coding::Erasure::try_new_with_options(
|
||||
@@ -690,18 +709,18 @@ pub(in crate::set_disk) async fn data_read_early_stop_inline_body_verified(
|
||||
candidate.erasure.block_size,
|
||||
candidate.uses_legacy_checksum,
|
||||
) else {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
};
|
||||
let Some(data_files) =
|
||||
collect_inline_data_shard_fileinfos_by_index(parts_metadata, candidate, erasure.data_shards, |index| {
|
||||
let data_files =
|
||||
match collect_inline_data_shard_fileinfos_by_index_or_reason(parts_metadata, candidate, erasure.data_shards, |index| {
|
||||
disks.get(index).is_some_and(Option::is_some)
|
||||
})
|
||||
else {
|
||||
return false;
|
||||
};
|
||||
}) {
|
||||
Ok(data_files) => data_files,
|
||||
Err(reason) => return Some(reason),
|
||||
};
|
||||
|
||||
let Some(part) = candidate.parts.first() else {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
};
|
||||
let checksum_info = candidate.erasure.get_checksum_info(part.number);
|
||||
let checksum_algo = if candidate.uses_legacy_checksum && checksum_info.algorithm == HashAlgorithm::HighwayHash256S {
|
||||
@@ -721,12 +740,13 @@ pub(in crate::set_disk) async fn data_read_early_stop_inline_body_verified(
|
||||
let Ok(mut readers) =
|
||||
build_inline_bitrot_readers_from_refs(&data_files, bucket, object, read_length, shard_size, &checksum_algo, false).await
|
||||
else {
|
||||
return false;
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY);
|
||||
};
|
||||
|
||||
try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size)
|
||||
.await
|
||||
.is_some_and(|body| body.len() == object_size)
|
||||
match try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size).await {
|
||||
Some(body) if body.len() == object_size => None,
|
||||
_ => Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn classify_metadata_response_error(err: &DiskError) -> &'static str {
|
||||
@@ -2469,6 +2489,7 @@ impl SetDisks {
|
||||
let mut next_fanout_index = 0usize;
|
||||
let mut scheduled_count = 0usize;
|
||||
let mut force_full_wait = false;
|
||||
let mut final_miss_reason_override = None;
|
||||
let spawn_read_version =
|
||||
|join_set: &mut JoinSet<(usize, disk::error::Result<FileInfo>, Duration)>, index: usize, disk: Option<DiskStore>| {
|
||||
let task_opts = opts;
|
||||
@@ -2541,17 +2562,29 @@ impl SetDisks {
|
||||
.or_else(|| accumulator.version_early_stop_decision())
|
||||
{
|
||||
let should_return_early = if read_data {
|
||||
let allow_data_read_early_stop = match accumulator.candidate.as_ref() {
|
||||
Some(candidate) => {
|
||||
data_read_early_stop_inline_body_verified(bucket.as_ref(), object.as_ref(), candidate, &ress, disks)
|
||||
.await
|
||||
match accumulator.candidate.as_ref() {
|
||||
Some(candidate) => match data_read_early_stop_inline_body_miss_reason(
|
||||
bucket.as_ref(),
|
||||
object.as_ref(),
|
||||
candidate,
|
||||
&ress,
|
||||
disks,
|
||||
)
|
||||
.await
|
||||
{
|
||||
None => true,
|
||||
Some(reason) => {
|
||||
force_full_wait = true;
|
||||
final_miss_reason_override = Some(reason);
|
||||
false
|
||||
}
|
||||
},
|
||||
None => {
|
||||
force_full_wait = true;
|
||||
final_miss_reason_override = Some(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM);
|
||||
false
|
||||
}
|
||||
None => false,
|
||||
};
|
||||
if !allow_data_read_early_stop {
|
||||
force_full_wait = true;
|
||||
}
|
||||
allow_data_read_early_stop
|
||||
} else {
|
||||
true
|
||||
};
|
||||
@@ -2613,7 +2646,12 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, accumulator.final_miss_reason());
|
||||
let accumulator_miss_reason = accumulator.final_miss_reason();
|
||||
let final_miss_reason = match (final_miss_reason_override, accumulator_miss_reason) {
|
||||
(Some(reason), GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM) => reason,
|
||||
_ => accumulator_miss_reason,
|
||||
};
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, final_miss_reason);
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0);
|
||||
rustfs_io_metrics::record_get_object_metadata_fanout_lifecycle(metrics_path, scheduled_count, scheduled_count, 0);
|
||||
let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations);
|
||||
@@ -6067,11 +6105,126 @@ mod tests {
|
||||
.clone();
|
||||
|
||||
assert!(
|
||||
data_read_early_stop_inline_body_verified(bucket, object, &candidate, &parts_metadata, &disks).await,
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &parts_metadata, &disks)
|
||||
.await
|
||||
.is_none(),
|
||||
"legacy inline metadata must use the legacy bitrot shard sizing and checksum algorithm"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn data_read_early_stop_reports_inline_miss_reasons() {
|
||||
let bucket = "inline-data-get-miss-reason-bucket";
|
||||
let object = "inline-data-get-miss-reason-object";
|
||||
let payload = b"verified inline payload";
|
||||
let (_dirs, disks) = call_counter_local_disks(bucket, 4).await;
|
||||
let files = inline_metadata_fanout_fileinfos_with_mode(bucket, object, payload, false).await;
|
||||
let distribution = files
|
||||
.first()
|
||||
.map(|file| file.erasure.distribution.clone())
|
||||
.expect("fixture should include metadata");
|
||||
let order = bounded_metadata_fanout_order(bucket, object, 4, 2);
|
||||
let mut parts_metadata = vec![FileInfo::default(); 4];
|
||||
for disk_index in order.into_iter().take(3) {
|
||||
let block_index = distribution
|
||||
.get(disk_index)
|
||||
.copied()
|
||||
.expect("fixture distribution should cover every disk");
|
||||
parts_metadata[disk_index] = files
|
||||
.get(block_index.checked_sub(1).expect("erasure block indexes are one-based"))
|
||||
.expect("fixture should include every distributed shard")
|
||||
.clone();
|
||||
}
|
||||
let candidate = parts_metadata
|
||||
.iter()
|
||||
.find(|file| file.name == object)
|
||||
.expect("fixture should include observed metadata")
|
||||
.clone();
|
||||
let data_disk = distribution
|
||||
.iter()
|
||||
.position(|block_index| *block_index == 1)
|
||||
.expect("fixture distribution should include first data shard");
|
||||
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &parts_metadata, &disks).await,
|
||||
None
|
||||
);
|
||||
|
||||
let mut not_inline = candidate.clone();
|
||||
rustfs_utils::http::remove_str(&mut not_inline.metadata, rustfs_utils::http::SUFFIX_INLINE_DATA);
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, ¬_inline, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE)
|
||||
);
|
||||
|
||||
let mut transformed = candidate.clone();
|
||||
rustfs_utils::http::insert_str(&mut transformed.metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &transformed, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED)
|
||||
);
|
||||
|
||||
let mut deleted = candidate.clone();
|
||||
deleted.deleted = true;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &deleted, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED)
|
||||
);
|
||||
|
||||
let mut zero_size = candidate.clone();
|
||||
zero_size.size = 0;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &zero_size, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE)
|
||||
);
|
||||
|
||||
let mut multipart = candidate.clone();
|
||||
multipart.parts.push(multipart.parts[0].clone());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &multipart, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE)
|
||||
);
|
||||
|
||||
let mut invalid_geometry = candidate.clone();
|
||||
invalid_geometry.erasure.data_blocks = 0;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &invalid_geometry, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY)
|
||||
);
|
||||
|
||||
let mut missing_shard = parts_metadata.clone();
|
||||
missing_shard[data_disk] = FileInfo::default();
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &missing_shard, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD)
|
||||
);
|
||||
|
||||
let mut missing_payload = parts_metadata.clone();
|
||||
missing_payload[data_disk].data = None;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &missing_payload, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD)
|
||||
);
|
||||
|
||||
let mut identity_mismatch = parts_metadata.clone();
|
||||
identity_mismatch[data_disk].version_id = Some(Uuid::new_v4());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &identity_mismatch, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH)
|
||||
);
|
||||
|
||||
let mut corrupt = parts_metadata.clone();
|
||||
if let Some(data) = corrupt[data_disk].data.as_mut() {
|
||||
let mut corrupt_data = data.to_vec();
|
||||
corrupt_data[0] ^= 0x01;
|
||||
*data = Bytes::from(corrupt_data);
|
||||
}
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &corrupt, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial]
|
||||
fn metadata_fanout_lifecycle_records_real_early_stop_abort() {
|
||||
@@ -6161,7 +6314,7 @@ mod tests {
|
||||
&[
|
||||
("path", GET_OBJECT_PATH_INTERNAL_META),
|
||||
("decision", "miss"),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
],
|
||||
),
|
||||
1,
|
||||
@@ -6173,7 +6326,7 @@ mod tests {
|
||||
&[
|
||||
("path", GET_OBJECT_PATH_LEGACY_DUPLEX),
|
||||
("decision", "miss"),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
],
|
||||
),
|
||||
0,
|
||||
|
||||
@@ -59,7 +59,10 @@ use crate::client::{object_api_utils::get_raw_etag, transition_api::ReaderImpl};
|
||||
use crate::cluster::rpc::heal_bucket_local_on_disks;
|
||||
use crate::data_usage::record_compression_total_memory;
|
||||
use crate::diagnostics::get::{
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, GET_OBJECT_PATH_BODY_CACHE, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, GET_OBJECT_PATH_BODY_CACHE, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_OBJECT_PATH_CODEC_STREAMING_LEGACY_ENGINE, GET_OBJECT_PATH_CODEC_STREAMING_RUSTFS_ENGINE, GET_OBJECT_PATH_DIRECT_MEMORY,
|
||||
GET_OBJECT_PATH_EMPTY, GET_OBJECT_PATH_INLINE_DIRECT, GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
GET_OBJECT_PATH_REMOTE_TRANSITION, GET_OBJECT_PATH_SET_DISK, GET_STAGE_DECODE, GET_STAGE_EMIT, GET_STAGE_INLINE_PREPARE,
|
||||
@@ -3866,8 +3869,17 @@ fn collect_inline_data_shard_fileinfos_by_index<'a>(
|
||||
parts_metadata: &'a [FileInfo],
|
||||
fi: &FileInfo,
|
||||
data_shards: usize,
|
||||
mut disk_is_online: impl FnMut(usize) -> bool,
|
||||
disk_is_online: impl FnMut(usize) -> bool,
|
||||
) -> Option<Vec<&'a FileInfo>> {
|
||||
collect_inline_data_shard_fileinfos_by_index_or_reason(parts_metadata, fi, data_shards, disk_is_online).ok()
|
||||
}
|
||||
|
||||
fn collect_inline_data_shard_fileinfos_by_index_or_reason<'a>(
|
||||
parts_metadata: &'a [FileInfo],
|
||||
fi: &FileInfo,
|
||||
data_shards: usize,
|
||||
mut disk_is_online: impl FnMut(usize) -> bool,
|
||||
) -> std::result::Result<Vec<&'a FileInfo>, &'static str> {
|
||||
let distribution = &fi.erasure.distribution;
|
||||
let mut data_files = vec![None; data_shards];
|
||||
|
||||
@@ -3875,27 +3887,35 @@ fn collect_inline_data_shard_fileinfos_by_index<'a>(
|
||||
if !disk_is_online(disk_index) {
|
||||
continue;
|
||||
}
|
||||
let block_index = *distribution.get(disk_index)?;
|
||||
let Some(&block_index) = distribution.get(disk_index) else {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
};
|
||||
if block_index == 0 || block_index > data_shards {
|
||||
continue;
|
||||
}
|
||||
if file_info.name.is_empty() {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD);
|
||||
}
|
||||
if file_info.erasure.index != block_index {
|
||||
continue;
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH);
|
||||
}
|
||||
if !file_info.has_valid_erasure_geometry() {
|
||||
continue;
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
}
|
||||
if !core::io_primitives::metadata_early_stop_candidate_matches(file_info, fi) {
|
||||
continue;
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH);
|
||||
}
|
||||
if file_info.data.as_ref().is_none_or(|data| data.is_empty()) {
|
||||
continue;
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD);
|
||||
}
|
||||
|
||||
data_files[block_index - 1] = Some(file_info);
|
||||
}
|
||||
|
||||
data_files.into_iter().collect()
|
||||
data_files
|
||||
.into_iter()
|
||||
.collect::<Option<Vec<_>>>()
|
||||
.ok_or(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD)
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
|
||||
@@ -135,10 +135,7 @@ impl FileMeta {
|
||||
let i = buf.len() as u64;
|
||||
|
||||
// check version, buf = buf[8..]
|
||||
let (buf, _, _) = Self::check_xl2_v1(buf).map_err(|e| {
|
||||
error!("failed to check XL2 v1 format: {}", e);
|
||||
e
|
||||
})?;
|
||||
let (buf, _, _) = Self::check_xl2_v1(buf)?;
|
||||
|
||||
if buf.len() < 5 {
|
||||
error!(
|
||||
|
||||
@@ -102,7 +102,7 @@ bytes.workspace = true
|
||||
hex-simd.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
tracing-subscriber = { workspace = true, features = ["env-filter", "time"] }
|
||||
tracing-subscriber = { workspace = true, features = ["json", "env-filter", "time"] }
|
||||
serial_test = { workspace = true }
|
||||
temp-env = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
|
||||
@@ -65,6 +65,7 @@ const LOG_SUBSYSTEM_FOLDER: &str = "folder";
|
||||
const LOG_SUBSYSTEM_LIFECYCLE: &str = "lifecycle";
|
||||
const LOG_SUBSYSTEM_HEAL: &str = "heal";
|
||||
const EVENT_SCANNER_FOLDER_STATE: &str = "scanner_folder_state";
|
||||
const EVENT_SCANNER_METADATA_CORRUPT: &str = "scanner_metadata_corrupt";
|
||||
const EVENT_SCANNER_LIFECYCLE_ACTION: &str = "scanner_lifecycle_action";
|
||||
const EVENT_SCANNER_HEAL_ADMISSION: &str = "scanner_heal_admission";
|
||||
const EVENT_SCANNER_ALERT_STATE: &str = "scanner_alert_state";
|
||||
@@ -2154,17 +2155,34 @@ impl FolderScanner {
|
||||
self.record_failed(&item.path);
|
||||
|
||||
if should_log_failed_object(into.failed_objects) {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "get_size_failed",
|
||||
error = %e,
|
||||
"Scanner folder failed to get object size"
|
||||
);
|
||||
if let GetSizeFailureAction::HealMetadata { object } = &failure_action {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_METADATA_CORRUPT,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
drive = %self.local_disk.path().display(),
|
||||
bucket = %item.bucket,
|
||||
object = %object,
|
||||
metadata_path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "metadata_corrupt",
|
||||
error = %e,
|
||||
"Scanner detected corrupt object metadata"
|
||||
);
|
||||
} else {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "get_size_failed",
|
||||
error = %e,
|
||||
"Scanner folder failed to get object size"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3054,12 +3072,59 @@ mod tests {
|
||||
use crate::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, TierStats, new_disk, storageclass};
|
||||
use rustfs_filemeta::{FileInfo, FileMeta};
|
||||
use serial_test::serial;
|
||||
use std::io::Write;
|
||||
#[cfg(unix)]
|
||||
use std::os::unix::fs::{PermissionsExt, symlink};
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use temp_env::{with_var, with_var_unset};
|
||||
use tracing_subscriber::fmt::MakeWriter;
|
||||
use uuid::Uuid;
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct CapturedLogs {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
struct CapturedLogWriter {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl CapturedLogs {
|
||||
fn contents(&self) -> String {
|
||||
let buffer = self
|
||||
.buffer
|
||||
.lock()
|
||||
.expect("captured logs mutex should not be poisoned")
|
||||
.clone();
|
||||
String::from_utf8(buffer).expect("captured logs should be valid UTF-8")
|
||||
}
|
||||
}
|
||||
|
||||
impl Write for CapturedLogWriter {
|
||||
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
self.buffer
|
||||
.lock()
|
||||
.expect("captured logs mutex should not be poisoned")
|
||||
.extend_from_slice(buf);
|
||||
Ok(buf.len())
|
||||
}
|
||||
|
||||
fn flush(&mut self) -> std::io::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a> MakeWriter<'a> for CapturedLogs {
|
||||
type Writer = CapturedLogWriter;
|
||||
|
||||
fn make_writer(&'a self) -> Self::Writer {
|
||||
CapturedLogWriter {
|
||||
buffer: Arc::clone(&self.buffer),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_size_summary_application_saturates_usage_counters() {
|
||||
let target = "arn:minio:replication::target".to_string();
|
||||
@@ -4542,9 +4607,19 @@ mod tests {
|
||||
assert!(budget.entries_visited() >= 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
#[serial]
|
||||
async fn test_scan_folder_corrupt_xl_meta_stops_erasure_data_dir_descent() {
|
||||
let logs = CapturedLogs::default();
|
||||
let subscriber = tracing_subscriber::fmt()
|
||||
.json()
|
||||
.with_max_level(tracing::Level::ERROR)
|
||||
.with_writer(logs.clone())
|
||||
.with_ansi(false)
|
||||
.without_time()
|
||||
.finish();
|
||||
let _subscriber_guard = tracing::subscriber::set_default(subscriber);
|
||||
|
||||
let (mut scanner, temp_dir) = build_test_scanner().await;
|
||||
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone());
|
||||
|
||||
@@ -4596,6 +4671,30 @@ mod tests {
|
||||
assert!(!budget.budget_elapsed());
|
||||
assert_eq!(budget.reason(), None);
|
||||
|
||||
let captured = logs.contents();
|
||||
assert!(
|
||||
!captured.contains("failed to check XL2 v1 format"),
|
||||
"the context-free filemeta parser error must not be emitted"
|
||||
);
|
||||
let events = captured
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<serde_json::Value>(line).expect("captured scanner log should be valid JSON"))
|
||||
.filter(|line| line["fields"]["event"] == EVENT_SCANNER_METADATA_CORRUPT)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
events.len(),
|
||||
1,
|
||||
"one corrupt metadata observation must emit one scanner-owned diagnostic event"
|
||||
);
|
||||
let fields = &events[0]["fields"];
|
||||
assert_eq!(fields["component"], LOG_COMPONENT_SCANNER);
|
||||
assert_eq!(fields["subsystem"], LOG_SUBSYSTEM_FOLDER);
|
||||
assert_eq!(fields["drive"], temp_dir.to_string_lossy().as_ref());
|
||||
assert_eq!(fields["bucket"], "bucket");
|
||||
assert_eq!(fields["object"], "object");
|
||||
assert_eq!(fields["metadata_path"], metadata_path.to_string_lossy().as_ref());
|
||||
assert_eq!(fields["state"], "metadata_corrupt");
|
||||
|
||||
let retry_budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
crate::scanner_budget::ScannerCycleBudgetConfig {
|
||||
|
||||
@@ -3849,17 +3849,6 @@ impl ScannerIODisk for Disk {
|
||||
let fivs = match meta.get_file_info_versions(item.bucket.as_str(), item.object_path().as_str(), false) {
|
||||
Ok(versions) => versions,
|
||||
Err(e) => {
|
||||
error!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_DISK_BUCKET_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
bucket = %item.bucket,
|
||||
object = %item.object_path(),
|
||||
state = "file_info_versions_failed",
|
||||
error = %e,
|
||||
"Scanner disk bucket failed to resolve file info versions"
|
||||
);
|
||||
return Err(scanner_metadata_corrupt_error(
|
||||
format!("failed to resolve file info versions: {e}"),
|
||||
&item.bucket,
|
||||
|
||||
Reference in New Issue
Block a user