mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 20:46:11 +00:00
chore: sync scanner publication with latest main
This commit is contained in:
@@ -196,15 +196,16 @@ pub mod bucket {
|
||||
pub use crate::bucket::metadata_sys::ConfigWriteLockProbe;
|
||||
pub use crate::bucket::metadata_sys::{
|
||||
BucketMetadataMutationGuard, BucketMetadataSys, ObjectLockConfigState, acquire_bucket_metadata_transaction_lock,
|
||||
acquire_bucket_metadata_transaction_lock_for_incarnation, capture_bucket_metadata_incarnation, delete,
|
||||
delete_if_incarnation, delete_under_transaction_lock, get, get_accelerate_config, get_bucket_policy,
|
||||
get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, get_cors_config, get_durability_config,
|
||||
get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config,
|
||||
get_object_lock_config, get_object_lock_config_state, get_on_demand_migration_config, get_public_access_block_config,
|
||||
get_quota_config, get_replication_config, get_request_payment_config, get_sse_config, get_tagging_config,
|
||||
get_versioning_config, get_website_config, init_bucket_metadata_sys, list_bucket_targets, reload_bucket_metadata,
|
||||
remove_bucket_metadata, set_bucket_metadata, update, update_bucket_targets_under_transaction_lock,
|
||||
update_config_with, update_if_incarnation, update_quota_if_incarnation, update_under_transaction_lock,
|
||||
acquire_bucket_metadata_transaction_lock_for_incarnation, acquire_scanner_bucket_incarnation_fence,
|
||||
capture_bucket_metadata_incarnation, delete, delete_if_incarnation, delete_under_transaction_lock, get,
|
||||
get_accelerate_config, get_bucket_policy, get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk,
|
||||
get_cors_config, get_durability_config, get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config,
|
||||
get_notification_config, get_object_lock_config, get_object_lock_config_state, get_on_demand_migration_config,
|
||||
get_public_access_block_config, get_quota_config, get_replication_config, get_request_payment_config, get_sse_config,
|
||||
get_tagging_config, get_versioning_config, get_website_config, init_bucket_metadata_sys, list_bucket_targets,
|
||||
reload_bucket_metadata, remove_bucket_metadata, set_bucket_metadata, update,
|
||||
update_bucket_targets_under_transaction_lock, update_config_with, update_if_incarnation, update_quota_if_incarnation,
|
||||
update_under_transaction_lock,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -655,6 +655,12 @@ pub struct BucketMetadataMutationGuard {
|
||||
}
|
||||
|
||||
impl BucketMetadataMutationGuard {
|
||||
/// Returns the storage-verified identity while both incarnation fences remain valid.
|
||||
pub fn checked_bucket_incarnation(&self) -> Result<(&str, Uuid)> {
|
||||
self.ensure_valid(&self.bucket)?;
|
||||
Ok((&self.bucket, self.incarnation_id))
|
||||
}
|
||||
|
||||
fn ensure_valid(&self, bucket: &str) -> Result<()> {
|
||||
if self.bucket != bucket {
|
||||
return Err(Error::other("bucket metadata mutation guard does not match bucket"));
|
||||
@@ -674,6 +680,29 @@ async fn acquire_config_write_guard_for_incarnation(
|
||||
sys: Arc<RwLock<BucketMetadataSys>>,
|
||||
bucket: &str,
|
||||
expected_incarnation_id: Option<Uuid>,
|
||||
) -> Result<BucketMetadataMutationGuard> {
|
||||
acquire_config_write_guard_with_migration(sys, bucket, expected_incarnation_id, true).await
|
||||
}
|
||||
|
||||
/// Scanner probes must not create an incarnation to make a capability available.
|
||||
pub async fn acquire_scanner_bucket_incarnation_fence(
|
||||
bucket: &str,
|
||||
expected_incarnation_id: Uuid,
|
||||
expected_owner_id: Uuid,
|
||||
) -> Result<BucketMetadataMutationGuard> {
|
||||
super::utils::check_valid_bucket_name(bucket)?;
|
||||
let sys = get_bucket_metadata_sys()?;
|
||||
if expected_owner_id.is_nil() || sys.read().await.api.id != expected_owner_id || expected_incarnation_id.is_nil() {
|
||||
return Err(Error::other("scanner bucket incarnation owner does not match"));
|
||||
}
|
||||
acquire_config_write_guard_with_migration(sys, bucket, Some(expected_incarnation_id), false).await
|
||||
}
|
||||
|
||||
async fn acquire_config_write_guard_with_migration(
|
||||
sys: Arc<RwLock<BucketMetadataSys>>,
|
||||
bucket: &str,
|
||||
expected_incarnation_id: Option<Uuid>,
|
||||
migrate: bool,
|
||||
) -> Result<BucketMetadataMutationGuard> {
|
||||
let metadata_sys = sys.read().await.clone();
|
||||
let lifecycle_guard = metadata_sys.api.acquire_bucket_lifecycle_read_lock(bucket).await?;
|
||||
@@ -681,13 +710,15 @@ async fn acquire_config_write_guard_for_incarnation(
|
||||
// Legacy buckets are migrated while the lifecycle fence prevents a
|
||||
// same-name replacement. The second read under the write transaction is
|
||||
// the CAS source of truth for the actual rewrite.
|
||||
await_bucket_namespace_operation(
|
||||
Some(&lifecycle_guard),
|
||||
bucket,
|
||||
"bucket config incarnation migration",
|
||||
metadata_sys.get_bucket_incarnation_id(bucket),
|
||||
)
|
||||
.await?;
|
||||
if migrate {
|
||||
await_bucket_namespace_operation(
|
||||
Some(&lifecycle_guard),
|
||||
bucket,
|
||||
"bucket config incarnation migration",
|
||||
metadata_sys.get_bucket_incarnation_id(bucket),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
let transaction_guard = await_bucket_namespace_operation(
|
||||
Some(&lifecycle_guard),
|
||||
bucket,
|
||||
@@ -3176,6 +3207,82 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scoped_dirty_usage_incarnation_probe_does_not_migrate_legacy_metadata() {
|
||||
let (dirs, store) = isolated_store_over_temp_disks().await;
|
||||
let sys = Arc::new(RwLock::new(BucketMetadataSys::new(store.clone())));
|
||||
let bucket = "scoped-ack-legacy";
|
||||
for dir in &dirs {
|
||||
std::fs::create_dir_all(dir.path().join(bucket)).expect("create legacy bucket");
|
||||
}
|
||||
let mut metadata = BucketMetadata::new(bucket);
|
||||
metadata.bucket_incarnation_id = Uuid::nil();
|
||||
sys.read()
|
||||
.await
|
||||
.persist_and_set(metadata)
|
||||
.await
|
||||
.expect("persist legacy metadata");
|
||||
assert!(
|
||||
acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(Uuid::new_v4()), false)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
assert!(load_bucket_incarnation(store, bucket).await.expect("read sidecar").is_none());
|
||||
assert!(
|
||||
sys.read()
|
||||
.await
|
||||
.get_config_from_disk(bucket)
|
||||
.await
|
||||
.expect("read metadata")
|
||||
.bucket_incarnation_id
|
||||
.is_nil()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
#[serial]
|
||||
async fn scoped_dirty_usage_incarnation_rejects_deleted_and_recreated_bucket() {
|
||||
let (_dirs, store) = isolated_store_over_temp_disks().await;
|
||||
init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let sys = bucket_metadata_sys_of(&store.ctx).expect("metadata owner");
|
||||
let bucket = "scoped-ack-recreated";
|
||||
store
|
||||
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create bucket");
|
||||
let old = store.bucket_incarnation_id_from_disk(bucket).await.expect("old incarnation");
|
||||
let guard = acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false)
|
||||
.await
|
||||
.expect("trusted incarnation fence");
|
||||
assert_eq!(guard.checked_bucket_incarnation().expect("valid fences"), (bucket, old));
|
||||
drop(guard);
|
||||
store
|
||||
.delete_bucket(bucket, &DeleteBucketOptions::default())
|
||||
.await
|
||||
.expect("delete bucket");
|
||||
assert!(
|
||||
acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
store
|
||||
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("recreate bucket");
|
||||
let new = store.bucket_incarnation_id_from_disk(bucket).await.expect("new incarnation");
|
||||
assert_ne!(old, new);
|
||||
assert!(
|
||||
acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
assert!(
|
||||
acquire_config_write_guard_with_migration(sys, bucket, Some(new), false)
|
||||
.await
|
||||
.is_ok()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn old_node_metadata_rewrite_cannot_replace_bucket_incarnation_sidecar() {
|
||||
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
|
||||
|
||||
@@ -30,6 +30,7 @@ use rustfs_protos::{
|
||||
ChannelClass, create_new_channel, get_channel_for_class,
|
||||
proto_gen::node_service::{
|
||||
heal_control_service_client::HealControlServiceClient, node_service_client::NodeServiceClient,
|
||||
scanner_control_service_client::ScannerControlServiceClient,
|
||||
tier_mutation_control_service_client::TierMutationControlServiceClient,
|
||||
},
|
||||
};
|
||||
@@ -60,6 +61,24 @@ pub async fn node_service_time_out_client(
|
||||
node_service_time_out_client_for_class(addr, interceptor, ChannelClass::Control).await
|
||||
}
|
||||
|
||||
pub(crate) async fn scanner_control_time_out_client(
|
||||
addr: &str,
|
||||
interceptor: TonicInterceptor,
|
||||
) -> crate::error::Result<ScannerControlServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
|
||||
let interceptor = interceptor.with_rpc_audience(addr)?;
|
||||
let channel = match runtime_sources::cached_node_channel(addr).await {
|
||||
Some(channel) => channel,
|
||||
None => create_new_channel(addr)
|
||||
.await
|
||||
.map_err(|err| crate::error::Error::other(err.to_string()))?,
|
||||
};
|
||||
let channel = ReplayScopeChannel::new(channel, interceptor.replay_scope_audience());
|
||||
let limit = rustfs_protos::scoped_dirty_usage::SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize;
|
||||
Ok(ScannerControlServiceClient::with_interceptor(channel, interceptor)
|
||||
.max_decoding_message_size(limit)
|
||||
.max_encoding_message_size(limit))
|
||||
}
|
||||
|
||||
pub async fn heal_control_time_out_client(
|
||||
addr: &str,
|
||||
interceptor: TonicInterceptor,
|
||||
|
||||
@@ -2050,6 +2050,53 @@ impl PeerRestClient {
|
||||
.await
|
||||
}
|
||||
|
||||
/// Probe only: scoped ACK production requires a durable per-bucket proof.
|
||||
pub async fn scanner_scoped_dirty_usage_capability(
|
||||
&self,
|
||||
owner_id: String,
|
||||
instance_id: String,
|
||||
entries: Vec<rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageEntry>,
|
||||
) -> Result<bool> {
|
||||
use rustfs_protos::scoped_dirty_usage::*;
|
||||
let payload = rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageAckRequest {
|
||||
challenge: Uuid::new_v4().as_bytes().to_vec().into(),
|
||||
protocol_version: SCOPED_DIRTY_USAGE_PROTOCOL_VERSION,
|
||||
owner_id,
|
||||
instance_id,
|
||||
scope: SCOPED_DIRTY_USAGE_BUCKET_SCOPE,
|
||||
probe_only: true,
|
||||
entries,
|
||||
};
|
||||
let canonical = canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?;
|
||||
self.finalize_result(
|
||||
async {
|
||||
let mut client = super::client::scanner_control_time_out_client(
|
||||
&self.grid_host,
|
||||
TonicInterceptor::Signature(gen_tonic_signature_interceptor()),
|
||||
)
|
||||
.await?;
|
||||
let mut request = Request::new(payload.clone());
|
||||
set_tonic_canonical_body_digest(&mut request, &canonical)?;
|
||||
let response = client.scanner_scoped_dirty_usage_ack(request).await?.into_inner();
|
||||
let body = canonical_scoped_dirty_usage_response(&canonical, &response)
|
||||
.map_err(|_| Error::other("scoped dirty usage capability response is too large"))?;
|
||||
verify_tonic_rpc_response_proof(&body, response.response_proof.as_ref())?;
|
||||
if response.protocol_version != SCOPED_DIRTY_USAGE_PROTOCOL_VERSION
|
||||
|| response.owner_id != payload.owner_id
|
||||
|| response.instance_id != payload.instance_id
|
||||
|| response.max_entries != SCOPED_DIRTY_USAGE_MAX_ENTRIES
|
||||
|| response.max_request_bytes != SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES
|
||||
|| response.cleared != 0
|
||||
{
|
||||
return Err(Error::other("scoped dirty usage capability response does not match request"));
|
||||
}
|
||||
Ok(response.supported)
|
||||
}
|
||||
.await,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn acknowledge_scanner_dirty_usage(&self, instance_id: String, generation: u64) -> Result<ScannerPeerActivity> {
|
||||
let result = self
|
||||
.scanner_activity_request_with_protocol(instance_id.clone(), generation, SCANNER_ACTIVITY_PROTOCOL_VERSION)
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -59,27 +59,34 @@ use super::super::{
|
||||
send_heal_request_with_admission, should_prevent_write, to_object_err, try_read_inline_data_shards_direct, warn,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(in crate::set_disk) use super::metadata_quorum::MetadataEarlyStopDecision;
|
||||
pub(in crate::set_disk) use super::metadata_quorum::{
|
||||
MetadataQuorumAccumulator, is_metadata_fanout_ignored_error, metadata_early_stop_candidate_matches,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
||||
#[cfg(test)]
|
||||
use crate::diagnostics::get::GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH;
|
||||
#[cfg(test)]
|
||||
use crate::diagnostics::get::GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD;
|
||||
use crate::diagnostics::get::{
|
||||
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_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_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_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_INSUFFICIENT_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, 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_DIRECT_MEMORY,
|
||||
GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, 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, get_stage_timer_if_enabled,
|
||||
record_get_stage_duration_if_enabled,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use crate::diagnostics::get::{
|
||||
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_DIRECT_MEMORY, GET_OBJECT_PATH_INTERNAL_META,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX, 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, get_stage_timer_if_enabled, record_get_stage_duration_if_enabled,
|
||||
GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use crate::disk::CHECK_PART_FILE_NOT_FOUND;
|
||||
@@ -690,324 +697,6 @@ impl MetadataFanoutDiagnostics {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(in crate::set_disk) struct MetadataEarlyStopDecision {
|
||||
pub(in crate::set_disk) reason: &'static str,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub(in crate::set_disk) struct MetadataQuorumAccumulator {
|
||||
pub(in crate::set_disk) total_disks: usize,
|
||||
pub(in crate::set_disk) default_parity_count: usize,
|
||||
pub(in crate::set_disk) allow_early_stop: bool,
|
||||
pub(in crate::set_disk) valid_responses: usize,
|
||||
pub(in crate::set_disk) not_found_responses: usize,
|
||||
pub(in crate::set_disk) version_not_found_responses: usize,
|
||||
pub(in crate::set_disk) ignored_errors: usize,
|
||||
pub(in crate::set_disk) hard_errors: usize,
|
||||
pub(in crate::set_disk) candidate: Option<FileInfo>,
|
||||
pub(in crate::set_disk) candidate_votes: usize,
|
||||
// Bitset of shard indexes whose metadata matches the candidate. Erasure
|
||||
// layouts are capped at 16 shards, so this stays allocation-free on the
|
||||
// GET metadata hot path.
|
||||
candidate_shard_mask: u16,
|
||||
pub(in crate::set_disk) conflicting_metadata: bool,
|
||||
pub(in crate::set_disk) delete_marker_seen: bool,
|
||||
pub(in crate::set_disk) delete_marker_candidates: Vec<(FileInfo, usize)>,
|
||||
pub(in crate::set_disk) delete_marker_votes: usize,
|
||||
pub(in crate::set_disk) requested_version_id: String,
|
||||
pub(in crate::set_disk) matching_version_votes: usize,
|
||||
}
|
||||
|
||||
impl MetadataQuorumAccumulator {
|
||||
pub(in crate::set_disk) fn new(total_disks: usize, default_parity_count: usize, allow_early_stop: bool) -> Self {
|
||||
Self {
|
||||
total_disks,
|
||||
default_parity_count,
|
||||
allow_early_stop,
|
||||
valid_responses: 0,
|
||||
not_found_responses: 0,
|
||||
version_not_found_responses: 0,
|
||||
ignored_errors: 0,
|
||||
hard_errors: 0,
|
||||
candidate: None,
|
||||
candidate_votes: 0,
|
||||
candidate_shard_mask: 0,
|
||||
conflicting_metadata: false,
|
||||
delete_marker_seen: false,
|
||||
delete_marker_candidates: Vec::new(),
|
||||
delete_marker_votes: 0,
|
||||
requested_version_id: String::new(),
|
||||
matching_version_votes: 0,
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn with_requested_version_id(mut self, version_id: &str) -> Self {
|
||||
self.requested_version_id = version_id.to_string();
|
||||
self
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_file_info(&mut self, file_info: &FileInfo) {
|
||||
self.observe_file_info_with_index(None, file_info);
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_file_info_at(&mut self, disk_index: usize, file_info: &FileInfo) {
|
||||
self.observe_file_info_with_index(Some(disk_index), file_info);
|
||||
}
|
||||
|
||||
fn observe_file_info_with_index(&mut self, disk_index: Option<usize>, file_info: &FileInfo) {
|
||||
if !file_info_is_valid_for_metadata(file_info) {
|
||||
self.hard_errors = self.hard_errors.saturating_add(1);
|
||||
return;
|
||||
}
|
||||
|
||||
self.valid_responses = self.valid_responses.saturating_add(1);
|
||||
|
||||
// Track version match for versioned requests
|
||||
if !self.requested_version_id.is_empty()
|
||||
&& let Some(ref vid) = file_info.version_id
|
||||
&& vid.to_string() == self.requested_version_id
|
||||
{
|
||||
self.matching_version_votes = self.matching_version_votes.saturating_add(1);
|
||||
}
|
||||
|
||||
if file_info.is_canonical_delete_marker() {
|
||||
self.delete_marker_seen = true;
|
||||
if let Some((_, votes)) = self
|
||||
.delete_marker_candidates
|
||||
.iter_mut()
|
||||
.find(|(candidate, _)| metadata_early_stop_candidate_matches(candidate, file_info))
|
||||
{
|
||||
*votes = votes.saturating_add(1);
|
||||
} else {
|
||||
self.delete_marker_candidates.push((file_info.clone(), 1));
|
||||
}
|
||||
self.delete_marker_votes = self
|
||||
.delete_marker_candidates
|
||||
.iter()
|
||||
.map(|(_, votes)| *votes)
|
||||
.max()
|
||||
.unwrap_or_default();
|
||||
self.conflicting_metadata |= self.delete_marker_candidates.len() > 1;
|
||||
return;
|
||||
}
|
||||
|
||||
match &self.candidate {
|
||||
Some(candidate) if metadata_early_stop_candidate_matches(candidate, file_info) => {
|
||||
self.candidate_votes = self.candidate_votes.saturating_add(1);
|
||||
if let Some(disk_index) = disk_index
|
||||
&& let Some(bit) = Self::candidate_shard_bit(candidate, file_info, disk_index)
|
||||
{
|
||||
self.candidate_shard_mask |= bit;
|
||||
}
|
||||
}
|
||||
Some(_) => {
|
||||
self.conflicting_metadata = true;
|
||||
}
|
||||
None => {
|
||||
self.candidate = Some(file_info.clone());
|
||||
self.candidate_votes = 1;
|
||||
if let Some(disk_index) = disk_index
|
||||
&& let Some(bit) = Self::candidate_shard_bit(file_info, file_info, disk_index)
|
||||
{
|
||||
self.candidate_shard_mask |= bit;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn candidate_shard_bit(candidate: &FileInfo, file_info: &FileInfo, disk_index: usize) -> Option<u16> {
|
||||
let &erasure_index = candidate.erasure.distribution.get(disk_index)?;
|
||||
if erasure_index == 0 || erasure_index > u16::BITS as usize || file_info.erasure.index != erasure_index {
|
||||
return None;
|
||||
}
|
||||
Some(1u16 << (erasure_index - 1))
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_has_read_reserve(&self) -> bool {
|
||||
self.candidate_read_reserve_target()
|
||||
.is_some_and(|required| self.candidate_shard_mask.count_ones() as usize >= required)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_read_reserve_target(&self) -> Option<usize> {
|
||||
let candidate = self.candidate.as_ref()?;
|
||||
Some(
|
||||
candidate
|
||||
.erasure
|
||||
.data_blocks
|
||||
.saturating_add(usize::from(candidate.erasure.parity_blocks > 0)),
|
||||
)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_error(&mut self, err: &DiskError) {
|
||||
match err {
|
||||
DiskError::FileNotFound | DiskError::VolumeNotFound => {
|
||||
self.not_found_responses = self.not_found_responses.saturating_add(1);
|
||||
}
|
||||
DiskError::FileVersionNotFound => {
|
||||
self.version_not_found_responses = self.version_not_found_responses.saturating_add(1);
|
||||
}
|
||||
_ if is_metadata_fanout_ignored_error(err) => {
|
||||
self.ignored_errors = self.ignored_errors.saturating_add(1);
|
||||
}
|
||||
_ => {
|
||||
self.hard_errors = self.hard_errors.saturating_add(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
|
||||
if !self.allow_early_stop {
|
||||
return None;
|
||||
}
|
||||
if self.delete_marker_votes >= self.default_write_quorum() {
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
|
||||
});
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return None;
|
||||
}
|
||||
if self
|
||||
.candidate
|
||||
.as_ref()
|
||||
.and_then(|candidate| self.candidate_latest_quorum(candidate))
|
||||
.is_some_and(|latest_quorum| self.candidate_votes >= latest_quorum)
|
||||
{
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
|
||||
});
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
/// Check if a versioned request can early-stop because the requested
|
||||
/// version_id has reached quorum across disks.
|
||||
pub(in crate::set_disk) fn version_early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
|
||||
if !self.allow_early_stop {
|
||||
return None;
|
||||
}
|
||||
if self.requested_version_id.is_empty() {
|
||||
return None;
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return None;
|
||||
}
|
||||
if self.matching_version_votes >= self.read_quorum_for_version() {
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
|
||||
});
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn can_still_reach_early_stop_with_pending(&self, pending: usize) -> bool {
|
||||
if !self.allow_early_stop {
|
||||
return false;
|
||||
}
|
||||
if self.delete_marker_votes.saturating_add(pending) >= self.default_write_quorum() {
|
||||
return true;
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if !self.requested_version_id.is_empty()
|
||||
&& self.matching_version_votes.saturating_add(pending) >= self.read_quorum_for_version()
|
||||
{
|
||||
return true;
|
||||
}
|
||||
match &self.candidate {
|
||||
Some(candidate) => self
|
||||
.candidate_latest_quorum(candidate)
|
||||
.is_some_and(|latest_quorum| self.candidate_votes.saturating_add(pending) >= latest_quorum),
|
||||
None => pending >= self.default_write_quorum(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Compute the read quorum threshold for version-aware early-stop.
|
||||
/// Uses `total_disks / 2` (like `missing_response_quorum`) when
|
||||
/// `default_parity_count` is set, otherwise requires all disks.
|
||||
pub(in crate::set_disk) fn read_quorum_for_version(&self) -> usize {
|
||||
self.missing_response_quorum()
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn final_miss_reason(&self) -> &'static str {
|
||||
if !self.allow_early_stop {
|
||||
return GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST;
|
||||
}
|
||||
if self.conflicting_metadata {
|
||||
return GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA;
|
||||
}
|
||||
if self.delete_marker_seen {
|
||||
return GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER;
|
||||
}
|
||||
let missing_response_quorum = self.missing_response_quorum();
|
||||
if self.version_not_found_responses >= missing_response_quorum {
|
||||
return GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND;
|
||||
}
|
||||
if self.not_found_responses >= missing_response_quorum {
|
||||
return GET_METADATA_EARLY_STOP_REASON_NOT_FOUND;
|
||||
}
|
||||
if self.hard_errors > 0 {
|
||||
return GET_METADATA_EARLY_STOP_REASON_ERROR;
|
||||
}
|
||||
if self.ignored_errors > 0 {
|
||||
return GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM;
|
||||
}
|
||||
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_latest_quorum(&self, candidate: &FileInfo) -> Option<usize> {
|
||||
if self.default_parity_count == 0 {
|
||||
return Some(self.total_disks);
|
||||
}
|
||||
if candidate.is_canonical_delete_marker() || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
|
||||
return None;
|
||||
}
|
||||
let data_blocks = candidate.erasure.data_blocks;
|
||||
Some(if data_blocks == candidate.erasure.parity_blocks {
|
||||
data_blocks.saturating_add(1)
|
||||
} else {
|
||||
data_blocks
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn default_write_quorum(&self) -> usize {
|
||||
if self.default_parity_count == 0 || self.default_parity_count >= self.total_disks {
|
||||
return self.total_disks;
|
||||
}
|
||||
let data_blocks = self.total_disks.saturating_sub(self.default_parity_count);
|
||||
if data_blocks == self.default_parity_count {
|
||||
data_blocks.saturating_add(1)
|
||||
} else {
|
||||
data_blocks
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn missing_response_quorum(&self) -> usize {
|
||||
if self.default_parity_count == 0 || self.default_parity_count >= self.total_disks {
|
||||
self.total_disks
|
||||
} else {
|
||||
self.total_disks / 2
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub(in crate::set_disk) enum MetadataCacheLookup {
|
||||
Hit(Arc<GetObjectMetadataCacheEntry>),
|
||||
@@ -1015,39 +704,6 @@ pub(in crate::set_disk) enum MetadataCacheLookup {
|
||||
RejectedInsufficientQuorum,
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo, right: &FileInfo) -> bool {
|
||||
left.volume == right.volume
|
||||
&& left.name == right.name
|
||||
&& left.version_id == right.version_id
|
||||
&& left.is_latest == right.is_latest
|
||||
&& left.deleted == right.deleted
|
||||
&& left.mark_deleted == right.mark_deleted
|
||||
&& left.transition_status == right.transition_status
|
||||
&& left.transitioned_objname == right.transitioned_objname
|
||||
&& left.transition_tier == right.transition_tier
|
||||
&& left.transition_version_id == right.transition_version_id
|
||||
&& left.transition_version == right.transition_version
|
||||
&& left.transition_version_state == right.transition_version_state
|
||||
&& left.expire_restored == right.expire_restored
|
||||
&& left.size == right.size
|
||||
&& left.mod_time == right.mod_time
|
||||
&& left.mode == right.mode
|
||||
&& left.written_by_version == right.written_by_version
|
||||
&& left.metadata == right.metadata
|
||||
&& left.replication_state_internal == right.replication_state_internal
|
||||
&& left.parts == right.parts
|
||||
&& left.checksum == right.checksum
|
||||
&& left.versioned == right.versioned
|
||||
&& left.num_versions == right.num_versions
|
||||
&& left.successor_mod_time == right.successor_mod_time
|
||||
&& left.data_dir == right.data_dir
|
||||
&& left.erasure.algorithm == right.erasure.algorithm
|
||||
&& left.erasure.data_blocks == right.erasure.data_blocks
|
||||
&& left.erasure.parity_blocks == right.erasure.parity_blocks
|
||||
&& left.erasure.block_size == right.erasure.block_size
|
||||
&& left.erasure.distribution == right.erasure.distribution
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
@@ -1247,10 +903,6 @@ pub(in crate::set_disk) fn classify_metadata_response_error(err: &DiskError) ->
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn is_metadata_fanout_ignored_error(err: &DiskError) -> bool {
|
||||
OBJECT_OP_IGNORED_ERRS.iter().any(|ignored| ignored == err)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn is_confirmed_missing_part_error(err: Option<&str>) -> bool {
|
||||
let Some(err) = err else {
|
||||
return false;
|
||||
|
||||
@@ -0,0 +1,385 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Pure metadata quorum and early-stop decisions for `SetDisks` reads.
|
||||
//!
|
||||
//! Disk scheduling, coalescing, cancellation, and late shard materialization
|
||||
//! remain with their existing owners; this module only classifies observations.
|
||||
|
||||
use crate::diagnostics::get::{
|
||||
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,
|
||||
};
|
||||
use crate::disk::error::DiskError;
|
||||
use crate::disk::error_reduce::OBJECT_OP_IGNORED_ERRS;
|
||||
use crate::set_disk::file_info_is_valid_for_metadata;
|
||||
use rustfs_filemeta::FileInfo;
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(in crate::set_disk) struct MetadataEarlyStopDecision {
|
||||
pub(in crate::set_disk) reason: &'static str,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub(in crate::set_disk) struct MetadataQuorumAccumulator {
|
||||
pub(in crate::set_disk) total_disks: usize,
|
||||
pub(in crate::set_disk) default_parity_count: usize,
|
||||
pub(in crate::set_disk) allow_early_stop: bool,
|
||||
pub(in crate::set_disk) valid_responses: usize,
|
||||
pub(in crate::set_disk) not_found_responses: usize,
|
||||
pub(in crate::set_disk) version_not_found_responses: usize,
|
||||
pub(in crate::set_disk) ignored_errors: usize,
|
||||
pub(in crate::set_disk) hard_errors: usize,
|
||||
pub(in crate::set_disk) candidate: Option<FileInfo>,
|
||||
pub(in crate::set_disk) candidate_votes: usize,
|
||||
// Bitset of shard indexes whose metadata matches the candidate. Erasure
|
||||
// layouts are capped at 16 shards, so this stays allocation-free on the
|
||||
// GET metadata hot path.
|
||||
candidate_shard_mask: u16,
|
||||
pub(in crate::set_disk) conflicting_metadata: bool,
|
||||
pub(in crate::set_disk) delete_marker_seen: bool,
|
||||
pub(in crate::set_disk) delete_marker_candidates: Vec<(FileInfo, usize)>,
|
||||
pub(in crate::set_disk) delete_marker_votes: usize,
|
||||
pub(in crate::set_disk) requested_version_id: String,
|
||||
pub(in crate::set_disk) matching_version_votes: usize,
|
||||
}
|
||||
|
||||
impl MetadataQuorumAccumulator {
|
||||
pub(in crate::set_disk) fn new(total_disks: usize, default_parity_count: usize, allow_early_stop: bool) -> Self {
|
||||
Self {
|
||||
total_disks,
|
||||
default_parity_count,
|
||||
allow_early_stop,
|
||||
valid_responses: 0,
|
||||
not_found_responses: 0,
|
||||
version_not_found_responses: 0,
|
||||
ignored_errors: 0,
|
||||
hard_errors: 0,
|
||||
candidate: None,
|
||||
candidate_votes: 0,
|
||||
candidate_shard_mask: 0,
|
||||
conflicting_metadata: false,
|
||||
delete_marker_seen: false,
|
||||
delete_marker_candidates: Vec::new(),
|
||||
delete_marker_votes: 0,
|
||||
requested_version_id: String::new(),
|
||||
matching_version_votes: 0,
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn with_requested_version_id(mut self, version_id: &str) -> Self {
|
||||
self.requested_version_id = version_id.to_string();
|
||||
self
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_file_info(&mut self, file_info: &FileInfo) {
|
||||
self.observe_file_info_with_index(None, file_info);
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_file_info_at(&mut self, disk_index: usize, file_info: &FileInfo) {
|
||||
self.observe_file_info_with_index(Some(disk_index), file_info);
|
||||
}
|
||||
|
||||
fn observe_file_info_with_index(&mut self, disk_index: Option<usize>, file_info: &FileInfo) {
|
||||
if !file_info_is_valid_for_metadata(file_info) {
|
||||
self.hard_errors = self.hard_errors.saturating_add(1);
|
||||
return;
|
||||
}
|
||||
|
||||
self.valid_responses = self.valid_responses.saturating_add(1);
|
||||
|
||||
// Track version match for versioned requests
|
||||
if !self.requested_version_id.is_empty()
|
||||
&& let Some(ref vid) = file_info.version_id
|
||||
&& vid.to_string() == self.requested_version_id
|
||||
{
|
||||
self.matching_version_votes = self.matching_version_votes.saturating_add(1);
|
||||
}
|
||||
|
||||
if file_info.is_canonical_delete_marker() {
|
||||
self.delete_marker_seen = true;
|
||||
if let Some((_, votes)) = self
|
||||
.delete_marker_candidates
|
||||
.iter_mut()
|
||||
.find(|(candidate, _)| metadata_early_stop_candidate_matches(candidate, file_info))
|
||||
{
|
||||
*votes = votes.saturating_add(1);
|
||||
} else {
|
||||
self.delete_marker_candidates.push((file_info.clone(), 1));
|
||||
}
|
||||
self.delete_marker_votes = self
|
||||
.delete_marker_candidates
|
||||
.iter()
|
||||
.map(|(_, votes)| *votes)
|
||||
.max()
|
||||
.unwrap_or_default();
|
||||
self.conflicting_metadata |= self.delete_marker_candidates.len() > 1;
|
||||
return;
|
||||
}
|
||||
|
||||
match &self.candidate {
|
||||
Some(candidate) if metadata_early_stop_candidate_matches(candidate, file_info) => {
|
||||
self.candidate_votes = self.candidate_votes.saturating_add(1);
|
||||
if let Some(disk_index) = disk_index
|
||||
&& let Some(bit) = Self::candidate_shard_bit(candidate, file_info, disk_index)
|
||||
{
|
||||
self.candidate_shard_mask |= bit;
|
||||
}
|
||||
}
|
||||
Some(_) => {
|
||||
self.conflicting_metadata = true;
|
||||
}
|
||||
None => {
|
||||
self.candidate = Some(file_info.clone());
|
||||
self.candidate_votes = 1;
|
||||
if let Some(disk_index) = disk_index
|
||||
&& let Some(bit) = Self::candidate_shard_bit(file_info, file_info, disk_index)
|
||||
{
|
||||
self.candidate_shard_mask |= bit;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn candidate_shard_bit(candidate: &FileInfo, file_info: &FileInfo, disk_index: usize) -> Option<u16> {
|
||||
let &erasure_index = candidate.erasure.distribution.get(disk_index)?;
|
||||
if erasure_index == 0 || erasure_index > u16::BITS as usize || file_info.erasure.index != erasure_index {
|
||||
return None;
|
||||
}
|
||||
Some(1u16 << (erasure_index - 1))
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_has_read_reserve(&self) -> bool {
|
||||
self.candidate_read_reserve_target()
|
||||
.is_some_and(|required| self.candidate_shard_mask.count_ones() as usize >= required)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_read_reserve_target(&self) -> Option<usize> {
|
||||
let candidate = self.candidate.as_ref()?;
|
||||
Some(
|
||||
candidate
|
||||
.erasure
|
||||
.data_blocks
|
||||
.saturating_add(usize::from(candidate.erasure.parity_blocks > 0)),
|
||||
)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_error(&mut self, err: &DiskError) {
|
||||
match err {
|
||||
DiskError::FileNotFound | DiskError::VolumeNotFound => {
|
||||
self.not_found_responses = self.not_found_responses.saturating_add(1);
|
||||
}
|
||||
DiskError::FileVersionNotFound => {
|
||||
self.version_not_found_responses = self.version_not_found_responses.saturating_add(1);
|
||||
}
|
||||
_ if is_metadata_fanout_ignored_error(err) => {
|
||||
self.ignored_errors = self.ignored_errors.saturating_add(1);
|
||||
}
|
||||
_ => {
|
||||
self.hard_errors = self.hard_errors.saturating_add(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
|
||||
if !self.allow_early_stop {
|
||||
return None;
|
||||
}
|
||||
if self.delete_marker_votes >= self.default_write_quorum() {
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
|
||||
});
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return None;
|
||||
}
|
||||
if self
|
||||
.candidate
|
||||
.as_ref()
|
||||
.and_then(|candidate| self.candidate_latest_quorum(candidate))
|
||||
.is_some_and(|latest_quorum| self.candidate_votes >= latest_quorum)
|
||||
{
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
|
||||
});
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
/// Check if a versioned request can early-stop because the requested
|
||||
/// version_id has reached quorum across disks.
|
||||
pub(in crate::set_disk) fn version_early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
|
||||
if !self.allow_early_stop {
|
||||
return None;
|
||||
}
|
||||
if self.requested_version_id.is_empty() {
|
||||
return None;
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return None;
|
||||
}
|
||||
if self.matching_version_votes >= self.read_quorum_for_version() {
|
||||
return Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
|
||||
});
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn can_still_reach_early_stop_with_pending(&self, pending: usize) -> bool {
|
||||
if !self.allow_early_stop {
|
||||
return false;
|
||||
}
|
||||
if self.delete_marker_votes.saturating_add(pending) >= self.default_write_quorum() {
|
||||
return true;
|
||||
}
|
||||
if self.conflicting_metadata
|
||||
|| self.delete_marker_seen
|
||||
|| self.not_found_responses > 0
|
||||
|| self.version_not_found_responses > 0
|
||||
|| self.hard_errors > 0
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if !self.requested_version_id.is_empty()
|
||||
&& self.matching_version_votes.saturating_add(pending) >= self.read_quorum_for_version()
|
||||
{
|
||||
return true;
|
||||
}
|
||||
match &self.candidate {
|
||||
Some(candidate) => self
|
||||
.candidate_latest_quorum(candidate)
|
||||
.is_some_and(|latest_quorum| self.candidate_votes.saturating_add(pending) >= latest_quorum),
|
||||
None => pending >= self.default_write_quorum(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Compute the read quorum threshold for version-aware early-stop.
|
||||
/// Uses `total_disks / 2` (like `missing_response_quorum`) when
|
||||
/// `default_parity_count` is set, otherwise requires all disks.
|
||||
pub(in crate::set_disk) fn read_quorum_for_version(&self) -> usize {
|
||||
self.missing_response_quorum()
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn final_miss_reason(&self) -> &'static str {
|
||||
if !self.allow_early_stop {
|
||||
return GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST;
|
||||
}
|
||||
if self.conflicting_metadata {
|
||||
return GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA;
|
||||
}
|
||||
if self.delete_marker_seen {
|
||||
return GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER;
|
||||
}
|
||||
let missing_response_quorum = self.missing_response_quorum();
|
||||
if self.version_not_found_responses >= missing_response_quorum {
|
||||
return GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND;
|
||||
}
|
||||
if self.not_found_responses >= missing_response_quorum {
|
||||
return GET_METADATA_EARLY_STOP_REASON_NOT_FOUND;
|
||||
}
|
||||
if self.hard_errors > 0 {
|
||||
return GET_METADATA_EARLY_STOP_REASON_ERROR;
|
||||
}
|
||||
if self.ignored_errors > 0 {
|
||||
return GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM;
|
||||
}
|
||||
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn candidate_latest_quorum(&self, candidate: &FileInfo) -> Option<usize> {
|
||||
if self.default_parity_count == 0 {
|
||||
return Some(self.total_disks);
|
||||
}
|
||||
if candidate.is_canonical_delete_marker() || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
|
||||
return None;
|
||||
}
|
||||
let data_blocks = candidate.erasure.data_blocks;
|
||||
Some(if data_blocks == candidate.erasure.parity_blocks {
|
||||
data_blocks.saturating_add(1)
|
||||
} else {
|
||||
data_blocks
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn default_write_quorum(&self) -> usize {
|
||||
if self.default_parity_count == 0 || self.default_parity_count >= self.total_disks {
|
||||
return self.total_disks;
|
||||
}
|
||||
let data_blocks = self.total_disks.saturating_sub(self.default_parity_count);
|
||||
if data_blocks == self.default_parity_count {
|
||||
data_blocks.saturating_add(1)
|
||||
} else {
|
||||
data_blocks
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn missing_response_quorum(&self) -> usize {
|
||||
if self.default_parity_count == 0 || self.default_parity_count >= self.total_disks {
|
||||
self.total_disks
|
||||
} else {
|
||||
self.total_disks / 2
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo, right: &FileInfo) -> bool {
|
||||
left.volume == right.volume
|
||||
&& left.name == right.name
|
||||
&& left.version_id == right.version_id
|
||||
&& left.is_latest == right.is_latest
|
||||
&& left.deleted == right.deleted
|
||||
&& left.mark_deleted == right.mark_deleted
|
||||
&& left.transition_status == right.transition_status
|
||||
&& left.transitioned_objname == right.transitioned_objname
|
||||
&& left.transition_tier == right.transition_tier
|
||||
&& left.transition_version_id == right.transition_version_id
|
||||
&& left.transition_version == right.transition_version
|
||||
&& left.transition_version_state == right.transition_version_state
|
||||
&& left.expire_restored == right.expire_restored
|
||||
&& left.size == right.size
|
||||
&& left.mod_time == right.mod_time
|
||||
&& left.mode == right.mode
|
||||
&& left.written_by_version == right.written_by_version
|
||||
&& left.metadata == right.metadata
|
||||
&& left.replication_state_internal == right.replication_state_internal
|
||||
&& left.parts == right.parts
|
||||
&& left.checksum == right.checksum
|
||||
&& left.versioned == right.versioned
|
||||
&& left.num_versions == right.num_versions
|
||||
&& left.successor_mod_time == right.successor_mod_time
|
||||
&& left.data_dir == right.data_dir
|
||||
&& left.erasure.algorithm == right.erasure.algorithm
|
||||
&& left.erasure.data_blocks == right.erasure.data_blocks
|
||||
&& left.erasure.parity_blocks == right.erasure.parity_blocks
|
||||
&& left.erasure.block_size == right.erasure.block_size
|
||||
&& left.erasure.distribution == right.erasure.distribution
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn is_metadata_fanout_ignored_error(err: &DiskError) -> bool {
|
||||
OBJECT_OP_IGNORED_ERRS.iter().any(|ignored| ignored == err)
|
||||
}
|
||||
@@ -18,3 +18,4 @@
|
||||
//! duplicating read/write/erasure logic.
|
||||
|
||||
pub(crate) mod io_primitives;
|
||||
mod metadata_quorum;
|
||||
|
||||
Reference in New Issue
Block a user