mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-09 13:46:05 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0b6a3ab10e | |||
| 3d71a45373 |
@@ -12,14 +12,11 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Exact, single-record inspection for legacy transitioned-version metadata.
|
||||
//!
|
||||
//! Mutation intentionally remains fail closed until the disk boundary exposes
|
||||
//! a conditional metadata-generation write and the fleet advertises the
|
||||
//! reconciliation capability. The existing best-effort metadata writers are
|
||||
//! not safe here: quorum failure may remove an existing `xl.meta`.
|
||||
//! Exact, generation-conditional reconciliation of legacy transition metadata.
|
||||
//! Repairs are monotonic and never use metadata rollback or remote deletion.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
@@ -27,12 +24,16 @@ use sha2::{Digest, Sha256};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::bucket::utils::check_bucket_and_object_names;
|
||||
use crate::disk::{DiskAPI as _, TransitionStateReconcileCondition, UpdateMetadataOpts};
|
||||
use crate::object_api::ObjectOptions;
|
||||
use crate::services::notification_sys::{
|
||||
acquire_cross_pool_fence_fleet_proof, acquire_remote_version_state_fleet_proof, cross_pool_fence_fleet_proof_matches,
|
||||
cross_pool_fence_topology_generation, remote_version_state_fleet_proof_matches,
|
||||
LegacyTransitionStateReconcileFleetProofToken, acquire_cross_pool_fence_fleet_proof,
|
||||
acquire_legacy_transition_state_reconcile_fleet_proof, acquire_remote_version_state_fleet_proof,
|
||||
cross_pool_fence_fleet_proof_matches, cross_pool_fence_topology_generation,
|
||||
legacy_transition_state_reconcile_fleet_proof_current, legacy_transition_state_reconcile_fleet_proof_matches,
|
||||
remote_version_state_fleet_proof_matches,
|
||||
};
|
||||
use crate::services::tier::tier::{TierConfigMgr, tier_destination_id_from_metadata};
|
||||
use crate::services::tier::tier::{TierConfigMgr, TierOperationLease, tier_destination_id_from_metadata};
|
||||
use crate::services::tier::warm_backend::LegacyTransitionStateProbe;
|
||||
use crate::set_disk::read_legacy_transition_state_metadata_copies;
|
||||
use crate::store::ECStore;
|
||||
@@ -111,6 +112,7 @@ pub struct LegacyTransitionStateSource {
|
||||
pub struct LegacyTransitionStateCopyRepresentation {
|
||||
pub disk_index: usize,
|
||||
pub metadata_digest: String,
|
||||
pub unchanged_metadata_digest: String,
|
||||
pub state_aliases: Vec<LegacyTransitionStateMetadataAlias>,
|
||||
pub version_aliases: Vec<LegacyTransitionStateMetadataAlias>,
|
||||
pub destination_aliases: Vec<LegacyTransitionStateMetadataAlias>,
|
||||
@@ -174,6 +176,9 @@ pub struct LegacyTransitionStateReconcileResponse {
|
||||
pub reason: String,
|
||||
pub retryable: bool,
|
||||
pub changed: bool,
|
||||
/// A failed RPC may have committed even when its response/readback is lost.
|
||||
#[serde(default)]
|
||||
pub changes_indeterminate: bool,
|
||||
pub selector: LegacyTransitionStateReconcileSelector,
|
||||
pub source: Option<LegacyTransitionStateSource>,
|
||||
pub original_sets: Vec<LegacyTransitionStateSetRepresentation>,
|
||||
@@ -203,6 +208,63 @@ struct InspectedCopy {
|
||||
representation: LegacyTransitionStateCopyRepresentation,
|
||||
}
|
||||
|
||||
/// Ownership accompanies a local publication into its blocking executor.
|
||||
/// Remote disks reconstruct fleet/backend ownership from current local state;
|
||||
/// serialized options never carry an authority supplied by another process.
|
||||
pub(crate) struct TransitionStateReconcileAuthority {
|
||||
// Declaration order releases physical locks, tier, bucket, then fleet.
|
||||
objects: Vec<crate::store::ObjectLockDiagGuard>,
|
||||
tier: TierOperationLease,
|
||||
bucket: Option<rustfs_lock::NamespaceLockGuard>,
|
||||
fleet: LegacyTransitionStateReconcileFleetProofToken,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for TransitionStateReconcileAuthority {
|
||||
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
formatter
|
||||
.debug_struct("TransitionStateReconcileAuthority")
|
||||
.finish_non_exhaustive()
|
||||
}
|
||||
}
|
||||
|
||||
impl TransitionStateReconcileAuthority {
|
||||
pub(crate) fn is_current(&self) -> bool {
|
||||
legacy_transition_state_reconcile_fleet_proof_current(&self.fleet)
|
||||
&& self.tier.is_current_generation()
|
||||
&& self.bucket.as_ref().is_none_or(|bucket| !bucket.is_lock_lost())
|
||||
&& self.objects.iter().all(|guard| !guard.is_lock_lost())
|
||||
}
|
||||
|
||||
pub(crate) async fn for_disk(condition: &TransitionStateReconcileCondition) -> crate::disk::error::Result<Arc<Self>> {
|
||||
use crate::disk::error::Error as DiskError;
|
||||
if let Some(authority) = &condition.authority {
|
||||
if authority.is_current() {
|
||||
return Ok(Arc::clone(authority));
|
||||
}
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
let fleet = acquire_legacy_transition_state_reconcile_fleet_proof()
|
||||
.await
|
||||
.ok_or(DiskError::OutdatedXLMeta)?;
|
||||
let topology = acquire_cross_pool_fence_fleet_proof().ok_or(DiskError::OutdatedXLMeta)?;
|
||||
if cross_pool_fence_topology_generation(&topology) != condition.topology_generation {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
let tier = TierConfigMgr::acquire_operation_lease(&crate::runtime::sources::global_tier_config_mgr(), &condition.tier)
|
||||
.await
|
||||
.map_err(|_| DiskError::OutdatedXLMeta)?;
|
||||
if rustfs_utils::crypto::hex(tier.backend_identity()) != condition.target.destination_id {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
Ok(Arc::new(Self {
|
||||
fleet,
|
||||
tier,
|
||||
objects: Vec::new(),
|
||||
bucket: None,
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
fn digest_hex(bytes: &[u8]) -> String {
|
||||
rustfs_utils::crypto::hex(Sha256::digest(bytes))
|
||||
}
|
||||
@@ -253,6 +315,11 @@ fn inspect_copy(
|
||||
representation: LegacyTransitionStateCopyRepresentation {
|
||||
disk_index,
|
||||
metadata_digest: digest_hex(raw),
|
||||
unchanged_metadata_digest: digest_hex(
|
||||
&metadata
|
||||
.transition_reconcile_generation(version_id)
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?,
|
||||
),
|
||||
state_aliases: metadata_aliases(&object_metadata.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE),
|
||||
version_aliases: metadata_aliases(&object_metadata.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID),
|
||||
destination_aliases: metadata_aliases(&object_metadata.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID),
|
||||
@@ -349,7 +416,7 @@ impl ECStore {
|
||||
let inspection: futures::future::BoxFuture<
|
||||
'_,
|
||||
Result<LegacyTransitionStateReconcileResponse, LegacyTransitionStateReconcileError>,
|
||||
> = Box::pin(self.inspect_legacy_transition_state_inner(canonical_selector.clone()));
|
||||
> = Box::pin(self.inspect_legacy_transition_state_inner(canonical_selector.clone(), false, None));
|
||||
match inspection.await {
|
||||
Ok(response) => Ok(response),
|
||||
Err(err) => Ok(error_response(canonical_selector, err)),
|
||||
@@ -359,6 +426,8 @@ impl ECStore {
|
||||
async fn inspect_legacy_transition_state_inner(
|
||||
&self,
|
||||
selector: LegacyTransitionStateReconcileSelector,
|
||||
write_locked: bool,
|
||||
bound_lease: Option<&TierOperationLease>,
|
||||
) -> Result<LegacyTransitionStateReconcileResponse, LegacyTransitionStateReconcileError> {
|
||||
let (selector, local_version_id) = selector.canonicalize()?;
|
||||
let remote_fleet_proof = acquire_remote_version_state_fleet_proof();
|
||||
@@ -366,29 +435,34 @@ impl ECStore {
|
||||
// Snapshot lock order: bucket lifecycle READ, then the fixed and
|
||||
// physical object READ domains in pool/set order. Release these locks
|
||||
// before acquiring the tier lease or waiting on remote I/O.
|
||||
let bucket_guard = self
|
||||
.acquire_bucket_lifecycle_read_lock(&selector.bucket)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
let bucket_guard = if write_locked {
|
||||
None
|
||||
} else {
|
||||
Some(
|
||||
self.acquire_bucket_lifecycle_read_lock(&selector.bucket)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?,
|
||||
)
|
||||
};
|
||||
let bucket_incarnation = self
|
||||
.bucket_incarnation_id_from_disk(&selector.bucket)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
let encoded_object = rustfs_utils::path::encode_dir_object(&selector.object);
|
||||
let mut lock_options = ObjectOptions::default();
|
||||
let object_guards = self
|
||||
.acquire_all_physical_object_read_locks(
|
||||
let object_guards = if write_locked {
|
||||
Vec::new()
|
||||
} else {
|
||||
self.acquire_all_physical_object_read_locks(
|
||||
"legacy_transition_state_reconcile_inspect",
|
||||
&selector.bucket,
|
||||
&encoded_object,
|
||||
&mut lock_options,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
// The deployed probe capability predates the reconciliation wire
|
||||
// format and destination-identity preservation guarantee. It may make
|
||||
// GET diagnostics possible, but it cannot authorize POST.
|
||||
let fleet_ready = false;
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?
|
||||
};
|
||||
let fleet_ready = acquire_legacy_transition_state_reconcile_fleet_proof().await.is_some();
|
||||
let mut topology_ready = topology_proof.is_some();
|
||||
let mut source = None;
|
||||
let mut original_sets = Vec::new();
|
||||
@@ -454,7 +528,9 @@ impl ECStore {
|
||||
});
|
||||
}
|
||||
}
|
||||
if object_guards.iter().any(|guard| guard.is_lock_lost()) || bucket_guard.is_lock_lost() {
|
||||
if object_guards.iter().any(|guard| guard.is_lock_lost())
|
||||
|| bucket_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|
||||
{
|
||||
return Err(LegacyTransitionStateReconcileError::BackendUnavailable(
|
||||
"a local metadata snapshot lock was lost before inspection completed".to_string(),
|
||||
));
|
||||
@@ -483,6 +559,7 @@ impl ECStore {
|
||||
reason: "the selected object does not contain a complete transition source tuple".to_string(),
|
||||
retryable: false,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -505,6 +582,7 @@ impl ECStore {
|
||||
reason: "an explicit unknown transition state is not legacy absence".to_string(),
|
||||
retryable: false,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -542,18 +620,32 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
let tier_manager = self.tier_config_mgr();
|
||||
let lease = match persisted_destination {
|
||||
Some(destination) => {
|
||||
TierConfigMgr::acquire_operation_lease_for_backend_identity(
|
||||
&tier_manager,
|
||||
&file_info.transition_tier,
|
||||
destination,
|
||||
)
|
||||
.await
|
||||
}
|
||||
None => TierConfigMgr::acquire_operation_lease(&tier_manager, &file_info.transition_tier).await,
|
||||
let owned_lease = if bound_lease.is_none() {
|
||||
Some(
|
||||
match persisted_destination {
|
||||
Some(destination) => {
|
||||
TierConfigMgr::acquire_operation_lease_for_backend_identity(
|
||||
&tier_manager,
|
||||
&file_info.transition_tier,
|
||||
destination,
|
||||
)
|
||||
.await
|
||||
}
|
||||
None => TierConfigMgr::acquire_operation_lease(&tier_manager, &file_info.transition_tier).await,
|
||||
}
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let lease = bound_lease.or(owned_lease.as_ref()).ok_or_else(|| {
|
||||
LegacyTransitionStateReconcileError::BackendUnavailable("tier generation lease is unavailable".to_string())
|
||||
})?;
|
||||
if persisted_destination.is_some_and(|identity| identity != lease.backend_identity()) {
|
||||
return Err(LegacyTransitionStateReconcileError::Corrupt(
|
||||
"persisted tier destination differs from the leased backend".to_string(),
|
||||
));
|
||||
}
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
let tier_generation = lease.generation();
|
||||
let destination_id = rustfs_utils::crypto::hex(lease.backend_identity());
|
||||
let probe = tokio::time::timeout(
|
||||
@@ -620,6 +712,7 @@ impl ECStore {
|
||||
reason: "the live backend probe did not prove exactly one remote version model".to_string(),
|
||||
retryable: true,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -653,6 +746,7 @@ impl ECStore {
|
||||
reason: "the live backend candidate does not match the persisted nonempty legacy remote version".to_string(),
|
||||
retryable: true,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -688,6 +782,7 @@ impl ECStore {
|
||||
.to_string(),
|
||||
retryable: false,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -720,6 +815,7 @@ impl ECStore {
|
||||
},
|
||||
retryable: !already_explicit,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: Some(source),
|
||||
original_sets,
|
||||
@@ -750,6 +846,7 @@ impl ECStore {
|
||||
reason: "reconciliation digest does not match the supplied source, sets, and target".to_string(),
|
||||
retryable: false,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector: request.selector,
|
||||
source: Some(request.source),
|
||||
original_sets: request.original_sets,
|
||||
@@ -758,54 +855,279 @@ impl ECStore {
|
||||
readiness: unavailable_write_readiness(),
|
||||
});
|
||||
}
|
||||
let current = self.inspect_legacy_transition_state(request.selector.clone()).await?;
|
||||
let mut changed = false;
|
||||
let mut indeterminate = false;
|
||||
let result = Box::pin(self.apply_legacy_transition_state(&request, &mut changed, &mut indeterminate)).await;
|
||||
let mut response = match result {
|
||||
Ok(response) => response,
|
||||
Err(err) => error_response(request.selector.clone(), err),
|
||||
};
|
||||
response.changed = changed;
|
||||
response.changes_indeterminate = indeterminate;
|
||||
Ok(response)
|
||||
}
|
||||
|
||||
async fn apply_legacy_transition_state(
|
||||
&self,
|
||||
request: &LegacyTransitionStateReconcileRequest,
|
||||
changed: &mut bool,
|
||||
indeterminate: &mut bool,
|
||||
) -> Result<LegacyTransitionStateReconcileResponse, LegacyTransitionStateReconcileError> {
|
||||
let unavailable = |message: &str| LegacyTransitionStateReconcileError::WriteFenceUnavailable(message.to_string());
|
||||
let fleet = acquire_legacy_transition_state_reconcile_fleet_proof()
|
||||
.await
|
||||
.ok_or_else(|| unavailable("every metadata writer must support conditional transition reconciliation"))?;
|
||||
// Admission -> bucket lifecycle WRITE -> exact tier generation ->
|
||||
// all physical object WRITE domains -> disk metadata mutation domain.
|
||||
let bucket = self
|
||||
.acquire_bucket_lifecycle_write_lock(&request.selector.bucket)
|
||||
.await
|
||||
.map_err(|err| unavailable(&err.to_string()))?;
|
||||
let tier = TierConfigMgr::acquire_operation_lease(&self.tier_config_mgr(), &request.source.tier)
|
||||
.await
|
||||
.map_err(|err| unavailable(&err.to_string()))?;
|
||||
if tier.generation() != request.target.tier_generation
|
||||
|| rustfs_utils::crypto::hex(tier.backend_identity()) != request.target.destination_id
|
||||
{
|
||||
return Err(LegacyTransitionStateReconcileError::StaleExpectedTuple(
|
||||
"tier generation or destination changed".to_string(),
|
||||
));
|
||||
}
|
||||
let object = rustfs_utils::path::encode_dir_object(&request.selector.object);
|
||||
let objects = self
|
||||
.acquire_all_physical_object_write_locks("legacy_transition_state_reconcile", &request.selector.bucket, &object)
|
||||
.await
|
||||
.map_err(|err| unavailable(&err.to_string()))?;
|
||||
let authority = Arc::new(TransitionStateReconcileAuthority {
|
||||
fleet,
|
||||
tier,
|
||||
objects,
|
||||
bucket: Some(bucket),
|
||||
});
|
||||
let mut current = self
|
||||
.inspect_legacy_transition_state_inner(request.selector.clone(), true, Some(&authority.tier))
|
||||
.await?;
|
||||
if !matches!(
|
||||
current.outcome,
|
||||
LegacyTransitionStateReconcileOutcome::ReadyToMigrate | LegacyTransitionStateReconcileOutcome::Migrated
|
||||
) {
|
||||
return Ok(current);
|
||||
}
|
||||
let current_matches_request = current.source.as_ref() == Some(&request.source)
|
||||
&& current.original_sets == request.original_sets
|
||||
&& current.target.as_ref() == Some(&request.target)
|
||||
&& current.reconciliation_digest.as_deref() == Some(request.reconciliation_digest.as_str());
|
||||
if !current_matches_request {
|
||||
return Ok(LegacyTransitionStateReconcileResponse {
|
||||
outcome: LegacyTransitionStateReconcileOutcome::Corrupt,
|
||||
reason_code: "stale_expected_tuple".to_string(),
|
||||
reason: "the authoritative source, target, or per-set metadata changed after inspection".to_string(),
|
||||
retryable: false,
|
||||
changed: false,
|
||||
selector: request.selector,
|
||||
source: current.source,
|
||||
original_sets: current.original_sets,
|
||||
target: current.target,
|
||||
reconciliation_digest: current.reconciliation_digest,
|
||||
readiness: current.readiness,
|
||||
});
|
||||
validate_reconcile_snapshot(request, current.source.as_ref(), ¤t.original_sets, current.target.as_ref())?;
|
||||
if !current.readiness.fleet_ready || !current.readiness.topology_ready {
|
||||
return Err(unavailable("the current fleet/topology snapshot cannot authorize metadata writes"));
|
||||
}
|
||||
if current.outcome == LegacyTransitionStateReconcileOutcome::Migrated {
|
||||
return Ok(current);
|
||||
let topology = request
|
||||
.source
|
||||
.topology_generation
|
||||
.as_ref()
|
||||
.ok_or_else(|| unavailable("topology proof is absent"))?;
|
||||
let (_, version_id) = request.selector.clone().canonicalize()?;
|
||||
// Visit every observed copy, including an already committed retry
|
||||
// subset. The second pass is a zero-write barrier in the same disk
|
||||
// mutation domain; a delayed first attempt can only be idempotent.
|
||||
for verify_only in [false, true] {
|
||||
for set in self.all_set_disks() {
|
||||
let Some(original_set) = request
|
||||
.original_sets
|
||||
.iter()
|
||||
.find(|item| item.pool_index == set.pool_index && item.set_index == set.set_index)
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
let disks = set.disk_inventory().await;
|
||||
for original in &original_set.copies {
|
||||
if !legacy_transition_state_reconcile_fleet_proof_matches(&authority.fleet).await
|
||||
|| self
|
||||
.bucket_incarnation_id_from_disk(&request.selector.bucket)
|
||||
.await
|
||||
.map_err(|err| unavailable(&err.to_string()))?
|
||||
.to_string()
|
||||
!= request.source.bucket_incarnation
|
||||
|| !authority.is_current()
|
||||
{
|
||||
return Err(unavailable("fleet, tier, bucket, or object fence changed before metadata publication"));
|
||||
}
|
||||
let disk = disks
|
||||
.get(original.disk_index)
|
||||
.and_then(Option::as_ref)
|
||||
.ok_or_else(|| unavailable("an authoritative disk is unavailable"))?;
|
||||
let opts = UpdateMetadataOpts {
|
||||
transition_reconcile: Some(Box::new(TransitionStateReconcileCondition {
|
||||
expected_metadata_digest: original.metadata_digest.clone(),
|
||||
unchanged_metadata_digest: original.unchanged_metadata_digest.clone(),
|
||||
target: rustfs_filemeta::TransitionStateReconcileTarget {
|
||||
state: request.target.state,
|
||||
remote_version: request.target.remote_version.clone(),
|
||||
destination_id: request.target.destination_id.clone(),
|
||||
},
|
||||
tier: request.source.tier.clone(),
|
||||
topology_generation: topology.clone(),
|
||||
verify_only,
|
||||
authority: Some(Arc::clone(&authority)),
|
||||
})),
|
||||
..Default::default()
|
||||
};
|
||||
let already_target = current
|
||||
.original_sets
|
||||
.iter()
|
||||
.find(|item| item.pool_index == set.pool_index && item.set_index == set.set_index)
|
||||
.and_then(|item| item.copies.iter().find(|copy| copy.disk_index == original.disk_index))
|
||||
.is_some_and(|copy| representation_matches_target(copy, &request.target));
|
||||
// Empty metadata makes a server which ignores the new
|
||||
// conditional option reject the legacy update operation.
|
||||
let result = disk
|
||||
.update_metadata(
|
||||
&request.selector.bucket,
|
||||
&object,
|
||||
FileInfo {
|
||||
version_id,
|
||||
..Default::default()
|
||||
},
|
||||
&opts,
|
||||
)
|
||||
.await;
|
||||
set.invalidate_get_object_metadata_cache(&request.selector.bucket, &object)
|
||||
.await;
|
||||
match result {
|
||||
Ok(()) => {
|
||||
*changed |= !verify_only && !already_target;
|
||||
}
|
||||
Err(err) => {
|
||||
*indeterminate |= !verify_only && !already_target;
|
||||
return Err(match err {
|
||||
crate::disk::error::Error::OutdatedXLMeta => {
|
||||
unavailable("a disk generation, rollback, or publication fence changed; inspect or retry")
|
||||
}
|
||||
crate::disk::error::Error::FileCorrupt => LegacyTransitionStateReconcileError::Corrupt(
|
||||
"a disk rejected conflicting transition metadata".to_string(),
|
||||
),
|
||||
_ => LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
let mut write_readiness = current.readiness;
|
||||
write_readiness.fleet_ready = false;
|
||||
write_readiness.post_ready = false;
|
||||
Ok(LegacyTransitionStateReconcileResponse {
|
||||
outcome: LegacyTransitionStateReconcileOutcome::BackendUnavailable,
|
||||
reason_code: "write_fence_unavailable".to_string(),
|
||||
reason:
|
||||
"conditional per-set xl.meta generation writes and a fleet reconciliation-capability fence are not implemented"
|
||||
.to_string(),
|
||||
retryable: true,
|
||||
changed: false,
|
||||
selector: request.selector,
|
||||
source: Some(request.source),
|
||||
original_sets: request.original_sets,
|
||||
target: Some(request.target),
|
||||
reconciliation_digest: Some(request.reconciliation_digest),
|
||||
readiness: write_readiness,
|
||||
})
|
||||
let (source, copies) = self.read_reconciled_transition_snapshot(&request.selector).await?;
|
||||
validate_reconcile_snapshot(request, Some(&source), &copies, Some(&request.target))?;
|
||||
if copies
|
||||
.iter()
|
||||
.flat_map(|set| &set.copies)
|
||||
.any(|copy| !representation_matches_target(copy, &request.target))
|
||||
{
|
||||
return Err(unavailable("strong readback did not prove convergence of every authoritative copy"));
|
||||
}
|
||||
if !legacy_transition_state_reconcile_fleet_proof_matches(&authority.fleet).await || !authority.is_current() {
|
||||
return Err(unavailable("publication authority changed before final readback completed"));
|
||||
}
|
||||
current.outcome = LegacyTransitionStateReconcileOutcome::Migrated;
|
||||
current.reason_code = if *changed { "migrated" } else { "already_converged" }.to_string();
|
||||
current.reason = "every authoritative metadata copy contains the proven state and destination".to_string();
|
||||
current.retryable = false;
|
||||
current.reconciliation_digest = Some(response_digest(&source, &copies, &request.target)?);
|
||||
current.source = Some(source);
|
||||
current.original_sets = copies;
|
||||
current.readiness = readiness(true, true, true, true, false);
|
||||
Ok(current)
|
||||
}
|
||||
|
||||
async fn read_reconciled_transition_snapshot(
|
||||
&self,
|
||||
selector: &LegacyTransitionStateReconcileSelector,
|
||||
) -> Result<(LegacyTransitionStateSource, Vec<LegacyTransitionStateSetRepresentation>), LegacyTransitionStateReconcileError>
|
||||
{
|
||||
let (_, version_id) = selector.clone().canonicalize()?;
|
||||
let incarnation = self
|
||||
.bucket_incarnation_id_from_disk(&selector.bucket)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
let topology = acquire_cross_pool_fence_fleet_proof()
|
||||
.ok_or_else(|| LegacyTransitionStateReconcileError::WriteFenceUnavailable("topology proof expired".to_string()))?;
|
||||
let mut source = None;
|
||||
let mut sets = Vec::new();
|
||||
for set in self.all_set_disks() {
|
||||
let raw = read_legacy_transition_state_metadata_copies(&set, &selector.bucket, &selector.object)
|
||||
.await
|
||||
.map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?;
|
||||
let total_copies = raw.len();
|
||||
let available_copies = raw.iter().flatten().count();
|
||||
let mut copies = Vec::new();
|
||||
for (index, raw) in raw.into_iter().enumerate() {
|
||||
let Some(raw) = raw else { continue };
|
||||
let Some(copy) = inspect_copy(&raw, index, &selector.bucket, &selector.object, version_id)? else {
|
||||
continue;
|
||||
};
|
||||
let mut observed_source = source_from_file_info(selector, incarnation, ©.file_info)?;
|
||||
observed_source.topology_generation = Some(cross_pool_fence_topology_generation(&topology));
|
||||
if source.as_ref().is_some_and(|source| source != &observed_source) {
|
||||
return Err(LegacyTransitionStateReconcileError::StaleExpectedTuple(
|
||||
"immutable source changed during readback".to_string(),
|
||||
));
|
||||
}
|
||||
source = Some(observed_source);
|
||||
copies.push(copy.representation);
|
||||
}
|
||||
if !copies.is_empty() {
|
||||
sets.push(LegacyTransitionStateSetRepresentation {
|
||||
pool_index: set.pool_index,
|
||||
set_index: set.set_index,
|
||||
total_copies,
|
||||
available_copies,
|
||||
copies,
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok((
|
||||
source.ok_or_else(|| LegacyTransitionStateReconcileError::StaleExpectedTuple("source disappeared".to_string()))?,
|
||||
sets,
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
fn representation_matches_target(copy: &LegacyTransitionStateCopyRepresentation, target: &LegacyTransitionStateTarget) -> bool {
|
||||
let aliases_equal = |aliases: &[LegacyTransitionStateMetadataAlias], value: &str| {
|
||||
!aliases.is_empty()
|
||||
&& aliases
|
||||
.iter()
|
||||
.all(|alias| alias.value_hex == rustfs_utils::crypto::hex(value.as_bytes()))
|
||||
};
|
||||
aliases_equal(©.state_aliases, target.state.as_str())
|
||||
&& aliases_equal(©.destination_aliases, &target.destination_id)
|
||||
&& match &target.remote_version {
|
||||
Some(version) => aliases_equal(©.version_aliases, version),
|
||||
None => copy.version_aliases.iter().all(|alias| alias.value_hex.is_empty()),
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_reconcile_snapshot(
|
||||
request: &LegacyTransitionStateReconcileRequest,
|
||||
source: Option<&LegacyTransitionStateSource>,
|
||||
sets: &[LegacyTransitionStateSetRepresentation],
|
||||
target: Option<&LegacyTransitionStateTarget>,
|
||||
) -> Result<(), LegacyTransitionStateReconcileError> {
|
||||
let matches = source == Some(&request.source)
|
||||
&& target == Some(&request.target)
|
||||
&& !sets.is_empty()
|
||||
&& sets.len() == request.original_sets.len()
|
||||
&& sets.iter().zip(&request.original_sets).all(|(current, original)| {
|
||||
current.pool_index == original.pool_index
|
||||
&& current.set_index == original.set_index
|
||||
&& current.total_copies == original.total_copies
|
||||
&& current.available_copies == original.available_copies
|
||||
&& current.copies.len() == original.copies.len()
|
||||
&& current.copies.iter().zip(&original.copies).all(|(current, original)| {
|
||||
current.disk_index == original.disk_index
|
||||
&& current.unchanged_metadata_digest == original.unchanged_metadata_digest
|
||||
&& (current == original || representation_matches_target(current, &request.target))
|
||||
})
|
||||
});
|
||||
if !matches {
|
||||
return Err(LegacyTransitionStateReconcileError::StaleExpectedTuple(
|
||||
"source, ownership, target, or unrelated metadata differs from the inspected generation".to_string(),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn legacy_remote_version_matches_target(file_info: &FileInfo, target: &LegacyTransitionStateTarget) -> bool {
|
||||
@@ -866,6 +1188,7 @@ fn error_response(
|
||||
reason: err.to_string(),
|
||||
retryable,
|
||||
changed: false,
|
||||
changes_indeterminate: false,
|
||||
selector,
|
||||
source: None,
|
||||
original_sets: Vec::new(),
|
||||
|
||||
@@ -219,6 +219,7 @@ fn restore_part_transaction_file(current: &Path, backup: &Path, absent: &Path, r
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn write_metadata_rollback_backup(object_dir: &Path, rollback_dir: Uuid, data: &[u8]) -> Result<()> {
|
||||
write_delete_rollback_file(object_dir, rollback_dir, STORAGE_FORMAT_FILE_BACKUP, data, None).await
|
||||
}
|
||||
@@ -250,6 +251,7 @@ async fn write_delete_rollback_file(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn restore_metadata_backup(
|
||||
object_dir: &Path,
|
||||
xl_path: &Path,
|
||||
@@ -277,15 +279,6 @@ async fn restore_metadata_backup_with_namespace_owner(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn restore_delete_rollback(
|
||||
object_dir: &Path,
|
||||
xl_path: &Path,
|
||||
rollback_dir: Uuid,
|
||||
publication_root: &os::PublicationRoot,
|
||||
) -> Result<()> {
|
||||
restore_delete_rollback_with_namespace_owner(object_dir, xl_path, rollback_dir, publication_root, None).await
|
||||
}
|
||||
|
||||
async fn restore_delete_rollback_with_namespace_owner(
|
||||
object_dir: &Path,
|
||||
xl_path: &Path,
|
||||
@@ -5857,6 +5850,8 @@ impl LocalDisk {
|
||||
check_path_length(file_path.to_string_lossy().as_ref())?;
|
||||
|
||||
let xl_path = path_join(&[file_path.as_path(), Path::new(STORAGE_FORMAT_FILE)]);
|
||||
let namespace_owner: Option<Arc<dyn Send + Sync>> =
|
||||
Some(os::acquire_metadata_mutation_lease(&self.get_object_path(volume, path)?, namespace_owner).await);
|
||||
if opts.old_data_dir.is_some() && opts.undo_write {
|
||||
return self.undo_write(file_path.as_path(), &fi, &opts, namespace_owner).await;
|
||||
}
|
||||
@@ -6714,6 +6709,8 @@ impl LocalDisk {
|
||||
|
||||
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &[FileInfo], opts: &DeleteOptions) -> Result<()> {
|
||||
let volume_dir = self.io_get_bucket_path(volume)?;
|
||||
let object_path = self.get_object_path(volume, path)?;
|
||||
let namespace_owner: Option<Arc<dyn Send + Sync>> = Some(os::acquire_metadata_mutation_lease(&object_path, None).await);
|
||||
let xlpath = self.io_get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?;
|
||||
let object_dir = xlpath
|
||||
.parent()
|
||||
@@ -6723,10 +6720,24 @@ impl LocalDisk {
|
||||
&& opts.undo_write
|
||||
{
|
||||
if opts.undo_delete {
|
||||
return restore_delete_rollback(object_dir, &xlpath, rollback_dir, &self.publication_root).await;
|
||||
return restore_delete_rollback_with_namespace_owner(
|
||||
object_dir,
|
||||
&xlpath,
|
||||
rollback_dir,
|
||||
&self.publication_root,
|
||||
namespace_owner.clone(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
return restore_metadata_backup(object_dir, &xlpath, rollback_dir, &self.publication_root).await;
|
||||
return restore_metadata_backup_with_namespace_owner(
|
||||
object_dir,
|
||||
&xlpath,
|
||||
rollback_dir,
|
||||
&self.publication_root,
|
||||
namespace_owner.clone(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
let (data, _) = match self.read_all_data_with_dmtime(volume, volume_dir.as_path(), &xlpath).await {
|
||||
@@ -6738,7 +6749,14 @@ impl LocalDisk {
|
||||
return Err(DiskError::FileNotFound);
|
||||
};
|
||||
return self
|
||||
.write_missing_delete_marker(volume, path, delete_marker, object_dir, opts.old_data_dir, None)
|
||||
.write_missing_delete_marker(
|
||||
volume,
|
||||
path,
|
||||
delete_marker,
|
||||
object_dir,
|
||||
opts.old_data_dir,
|
||||
namespace_owner.clone(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
@@ -6754,7 +6772,8 @@ impl LocalDisk {
|
||||
let rollback_dir = opts.old_data_dir;
|
||||
let mut reserved_version_delete = false;
|
||||
if let Some(rollback_dir) = rollback_dir {
|
||||
write_metadata_rollback_backup(object_dir, rollback_dir, &data).await?;
|
||||
write_delete_rollback_file(object_dir, rollback_dir, STORAGE_FORMAT_FILE_BACKUP, &data, namespace_owner.clone())
|
||||
.await?;
|
||||
}
|
||||
|
||||
for fi in fis.iter() {
|
||||
@@ -6768,13 +6787,16 @@ impl LocalDisk {
|
||||
|
||||
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_metadata_update",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_metadata_update",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6787,7 +6809,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_metadata_update",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -6804,13 +6826,16 @@ impl LocalDisk {
|
||||
Err(err) => {
|
||||
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_data_path",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_data_path",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6823,7 +6848,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_data_path",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -6836,13 +6861,16 @@ impl LocalDisk {
|
||||
let err: DiskError = to_file_error(err).into();
|
||||
if reserved_version_delete {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_rollback_dir",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_rollback_dir",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6855,23 +6883,29 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_rollback_dir",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
.await);
|
||||
}
|
||||
let reserved = match self.reserve_version_delete(volume, path, dir, rollback_dir).await {
|
||||
let reserved = match self
|
||||
.reserve_version_delete_with_namespace_owner(volume, path, dir, rollback_dir, namespace_owner.clone())
|
||||
.await
|
||||
{
|
||||
Ok(reserved) => reserved,
|
||||
Err(err) => {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_reserve_data",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_reserve_data",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6879,11 +6913,12 @@ impl LocalDisk {
|
||||
reserved_version_delete |= reserved;
|
||||
let rollback_data_path = rollback_path.join(dir.to_string());
|
||||
if !reserved
|
||||
&& let Err(err) = rename_all_ignore_missing_source(
|
||||
&& let Err(err) = os::rename_all_ignore_missing_source_with_owner(
|
||||
&dir_path,
|
||||
&rollback_data_path,
|
||||
&rollback_path,
|
||||
&self.publication_root,
|
||||
namespace_owner.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -6896,7 +6931,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_stage_data",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -6905,13 +6940,16 @@ impl LocalDisk {
|
||||
if should_fail_after_delete_data_staged(path) {
|
||||
if reserved_version_delete {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_test_after_stage",
|
||||
DiskError::Unexpected,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_test_after_stage",
|
||||
error: DiskError::Unexpected,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6924,13 +6962,15 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_test_after_stage",
|
||||
error: DiskError::Unexpected,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
.await);
|
||||
}
|
||||
} else if let Err(err) = self.move_to_trash(&dir_path, true, false).await
|
||||
} else if let Err(err) = self
|
||||
.move_to_trash_with_namespace_owner(&dir_path, true, false, namespace_owner.clone())
|
||||
.await
|
||||
&& !(err == DiskError::FileNotFound || err == DiskError::VolumeNotFound)
|
||||
{
|
||||
return Err(err);
|
||||
@@ -6945,16 +6985,22 @@ impl LocalDisk {
|
||||
|
||||
// Remove xl.meta when no versions remain
|
||||
if fm.versions.is_empty() {
|
||||
if let Err(err) = self.delete_file(&volume_dir, &xlpath, true, false).await {
|
||||
if let Err(err) = self
|
||||
.delete_file_with_namespace_owner(&volume_dir, &xlpath, true, false, namespace_owner.clone())
|
||||
.await
|
||||
{
|
||||
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_commit_delete",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_delete",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -6967,7 +7013,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_delete",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -6975,10 +7021,22 @@ impl LocalDisk {
|
||||
}
|
||||
if reserved_version_delete
|
||||
&& let Some(rollback_dir) = rollback_dir
|
||||
&& let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await
|
||||
&& let Err(err) = self
|
||||
.commit_reserved_version_delete_with_namespace_owner(volume, path, rollback_dir, namespace_owner.clone())
|
||||
.await
|
||||
{
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err)
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_intent",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
if should_fail_after_delete_commit(self.root.as_path(), path) {
|
||||
@@ -6995,13 +7053,16 @@ impl LocalDisk {
|
||||
let err: DiskError = err.into();
|
||||
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
"delete_versions_metadata_encode",
|
||||
err,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_metadata_encode",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
@@ -7014,7 +7075,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_metadata_encode",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -7023,12 +7084,28 @@ impl LocalDisk {
|
||||
};
|
||||
|
||||
if let Err(err) = self
|
||||
.write_all_meta(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str(), &buf, true)
|
||||
.write_all_meta_with_namespace_owner(
|
||||
volume,
|
||||
format!("{path}/{STORAGE_FORMAT_FILE}").as_str(),
|
||||
&buf,
|
||||
true,
|
||||
namespace_owner.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_write", err)
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_write",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
return Err(restore_delete_rollback_after_error(
|
||||
@@ -7040,7 +7117,7 @@ impl LocalDisk {
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_write",
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
&self.publication_root,
|
||||
)
|
||||
@@ -7049,10 +7126,22 @@ impl LocalDisk {
|
||||
|
||||
if reserved_version_delete
|
||||
&& let Some(rollback_dir) = rollback_dir
|
||||
&& let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await
|
||||
&& let Err(err) = self
|
||||
.commit_reserved_version_delete_with_namespace_owner(volume, path, rollback_dir, namespace_owner.clone())
|
||||
.await
|
||||
{
|
||||
return Err(self
|
||||
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err)
|
||||
.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
path,
|
||||
DeleteRollbackFailure {
|
||||
stage: "delete_versions_commit_intent",
|
||||
error: err,
|
||||
namespace_owner: namespace_owner.clone(),
|
||||
},
|
||||
)
|
||||
.await);
|
||||
}
|
||||
|
||||
@@ -7063,6 +7152,107 @@ impl LocalDisk {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn reconcile_transition_state_metadata(
|
||||
&self,
|
||||
volume: &str,
|
||||
object: &str,
|
||||
version_id: Option<Uuid>,
|
||||
condition: &super::TransitionStateReconcileCondition,
|
||||
namespace_owner: Option<Arc<dyn Send + Sync>>,
|
||||
authority: Arc<crate::bucket::lifecycle::legacy_transition_state_reconcile::TransitionStateReconcileAuthority>,
|
||||
) -> Result<()> {
|
||||
condition.target.validate()?;
|
||||
if !authority.is_current()
|
||||
|| [
|
||||
condition.expected_metadata_digest.as_str(),
|
||||
condition.unchanged_metadata_digest.as_str(),
|
||||
]
|
||||
.iter()
|
||||
.any(|digest| {
|
||||
digest.len() != 64
|
||||
|| !digest
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
|
||||
})
|
||||
{
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
let metadata_path = format!("{object}/{STORAGE_FORMAT_FILE}");
|
||||
let original = self.read_all(volume, &metadata_path).await?;
|
||||
let original_digest = rustfs_utils::crypto::hex(<sha2::Sha256 as sha2::Digest>::digest(&original));
|
||||
let mut metadata = FileMeta::load(&original)?;
|
||||
let (_, selected) = metadata.find_version(version_id)?;
|
||||
if selected.into_fileinfo(volume, object, true)?.transition_tier != condition.tier {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
let generation = metadata.transition_reconcile_generation(version_id)?;
|
||||
if rustfs_utils::crypto::hex(<sha2::Sha256 as sha2::Digest>::digest(&generation)) != condition.unchanged_metadata_digest {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
let changed = metadata.reconcile_transition_state(version_id, &condition.target)?;
|
||||
if !changed {
|
||||
return Ok(());
|
||||
}
|
||||
if condition.verify_only || original_digest != condition.expected_metadata_digest {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
// fsync_dir_std is a no-op outside Unix, so those platforms cannot
|
||||
// yet prove this repair's durable publication requirement.
|
||||
if !cfg!(unix) || !effective_durability(volume).syncs_commit_metadata() {
|
||||
return Err(DiskError::other(
|
||||
"transition reconciliation requires Unix directory sync and enabled bucket metadata durability",
|
||||
));
|
||||
}
|
||||
// An outstanding rollback can still restore an older whole xl.meta.
|
||||
// Its preparation and execution share this mutation domain; refuse
|
||||
// repair until that transaction has settled and removed its backup.
|
||||
let mut entries = fs::read_dir(self.io_get_object_path(volume, object)?)
|
||||
.await
|
||||
.map_err(to_file_error)?;
|
||||
let mut remaining = 4096usize;
|
||||
while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? {
|
||||
remaining = remaining.checked_sub(1).ok_or(DiskError::OutdatedXLMeta)?;
|
||||
if Uuid::parse_str(&entry.file_name().to_string_lossy()).is_err() {
|
||||
continue;
|
||||
}
|
||||
for marker in [STORAGE_FORMAT_FILE_BACKUP, DELETE_MARKER_ROLLBACK_FILE] {
|
||||
if fs::try_exists(entry.path().join(marker)).await.map_err(to_file_error)? {
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
}
|
||||
}
|
||||
let replacement = metadata.marshal_msg()?;
|
||||
let tmp_volume = self.io_get_bucket_path(RUSTFS_META_TMP_BUCKET)?;
|
||||
let tmp_file = self.io_get_object_path(RUSTFS_META_TMP_BUCKET, &Uuid::new_v4().to_string())?;
|
||||
// Admission above requires metadata durability. Keep rename and its
|
||||
// directory sync in the same owned executor even after cancellation.
|
||||
self.write_all_internal(&tmp_file, InternalBuf::Ref(&replacement), SyncMode::FileOnly, &tmp_volume)
|
||||
.await?;
|
||||
if crash_inject::should_crash_at(CrashPoint::MetaWriteAfterTmpBeforeRename, &metadata_path) {
|
||||
return Err(DiskError::Unexpected);
|
||||
}
|
||||
os::rename_reconciled_metadata(
|
||||
tmp_file,
|
||||
self.io_get_object_path(volume, &metadata_path)?,
|
||||
self.io_get_bucket_path(volume)?,
|
||||
self.publication_root.clone(),
|
||||
namespace_owner.clone(),
|
||||
authority,
|
||||
)
|
||||
.await?;
|
||||
// Keep the same mutation lease through strong readback. Response loss
|
||||
// leaves a monotonic subset for the coordinator's exact-copy retry.
|
||||
let committed = self.read_all(volume, &metadata_path).await?;
|
||||
let mut committed = FileMeta::load(&committed)?;
|
||||
if committed.reconcile_transition_state(version_id, &condition.target)?
|
||||
|| committed.transition_reconcile_generation(version_id)? != generation
|
||||
{
|
||||
return Err(DiskError::OutdatedXLMeta);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn write_all_meta(&self, volume: &str, path: &str, buf: &[u8], sync: bool) -> Result<()> {
|
||||
self.write_all_meta_with_namespace_owner(volume, path, buf, sync, None).await
|
||||
}
|
||||
@@ -8321,6 +8511,7 @@ impl LocalDisk {
|
||||
Ok(Arc::new(QuotaMutationFenceClaim { state }))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn reserve_version_delete(&self, volume: &str, object: &str, data_dir: Uuid, rollback_dir: Uuid) -> Result<bool> {
|
||||
self.reserve_version_delete_with_namespace_owner(volume, object, data_dir, rollback_dir, None)
|
||||
.await
|
||||
@@ -8373,6 +8564,7 @@ impl LocalDisk {
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn commit_reserved_version_delete(&self, volume: &str, object: &str, rollback_dir: Uuid) -> Result<()> {
|
||||
self.commit_reserved_version_delete_with_namespace_owner(volume, object, rollback_dir, None)
|
||||
.await
|
||||
@@ -8464,29 +8656,6 @@ impl LocalDisk {
|
||||
first_err.map_or(Ok(found), Err)
|
||||
}
|
||||
|
||||
async fn abort_reserved_version_delete(
|
||||
&self,
|
||||
object_dir: &Path,
|
||||
rollback_dir: Uuid,
|
||||
volume: &str,
|
||||
object: &str,
|
||||
stage: &'static str,
|
||||
err: DiskError,
|
||||
) -> DiskError {
|
||||
self.abort_reserved_version_delete_with_failure(
|
||||
object_dir,
|
||||
rollback_dir,
|
||||
volume,
|
||||
object,
|
||||
DeleteRollbackFailure {
|
||||
stage,
|
||||
error: err,
|
||||
namespace_owner: None,
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn abort_reserved_version_delete_with_failure(
|
||||
&self,
|
||||
object_dir: &Path,
|
||||
@@ -10031,6 +10200,23 @@ impl DiskAPI for LocalDisk {
|
||||
|
||||
#[tracing::instrument(level = "trace", skip_all)]
|
||||
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> {
|
||||
let object_path = self.get_object_path(volume, path)?;
|
||||
if let Some(condition) = &opts.transition_reconcile {
|
||||
if !fi.metadata.is_empty() || opts.no_persistence || opts.replace_user_metadata {
|
||||
return Err(DiskError::FileCorrupt);
|
||||
}
|
||||
let authority =
|
||||
crate::bucket::lifecycle::legacy_transition_state_reconcile::TransitionStateReconcileAuthority::for_disk(
|
||||
condition,
|
||||
)
|
||||
.await?;
|
||||
let owner: Option<Arc<dyn Send + Sync>> = Some(authority.clone());
|
||||
let owner: Option<Arc<dyn Send + Sync>> = Some(os::acquire_metadata_mutation_lease(&object_path, owner).await);
|
||||
return self
|
||||
.reconcile_transition_state_metadata(volume, path, fi.version_id, condition, owner, authority)
|
||||
.await;
|
||||
}
|
||||
let namespace_owner: Option<Arc<dyn Send + Sync>> = Some(os::acquire_metadata_mutation_lease(&object_path, None).await);
|
||||
if !fi.metadata.is_empty() {
|
||||
let file_path = self.io_get_object_path(volume, path)?;
|
||||
|
||||
@@ -10058,7 +10244,13 @@ impl DiskAPI for LocalDisk {
|
||||
let wbuf = xl_meta.marshal_msg()?;
|
||||
|
||||
return self
|
||||
.write_all_meta(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str(), &wbuf, !opts.no_persistence)
|
||||
.write_all_meta_with_namespace_owner(
|
||||
volume,
|
||||
format!("{path}/{STORAGE_FORMAT_FILE}").as_str(),
|
||||
&wbuf,
|
||||
!opts.no_persistence,
|
||||
namespace_owner,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
@@ -10066,7 +10258,10 @@ impl DiskAPI for LocalDisk {
|
||||
}
|
||||
|
||||
async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
|
||||
self.write_metadata_with_namespace_owner(volume, path, fi, None).await
|
||||
let object_path = self.get_object_path(volume, path)?;
|
||||
let namespace_owner: Option<Arc<dyn Send + Sync>> = Some(os::acquire_metadata_mutation_lease(&object_path, None).await);
|
||||
self.write_metadata_with_namespace_owner(volume, path, fi, namespace_owner)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tracing::instrument(level = "trace", skip_all)]
|
||||
|
||||
@@ -279,11 +279,14 @@ impl LocalDisk {
|
||||
Some(token) => Some(self.claim_quota_mutation_fence(dst_volume, dst_path, token).await?),
|
||||
None => None,
|
||||
};
|
||||
// Quota admission -> metadata RMW -> namespace/volume publication.
|
||||
let metadata_lease =
|
||||
os::acquire_metadata_mutation_lease(&self.get_object_path(dst_volume, dst_path)?, state.namespace_owner.take()).await;
|
||||
let mutation_lease = os::acquire_rename_data_mutation_lease_with_owner(
|
||||
&self.root,
|
||||
dst_volume,
|
||||
&destination_object_path,
|
||||
state.namespace_owner.take(),
|
||||
Some(metadata_lease),
|
||||
)
|
||||
.await;
|
||||
if let Some(claim) = quota_fence_claim {
|
||||
|
||||
@@ -1277,6 +1277,25 @@ pub struct CheckPartsResp {
|
||||
pub struct UpdateMetadataOpts {
|
||||
pub no_persistence: bool,
|
||||
pub replace_user_metadata: bool,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub transition_reconcile: Option<Box<TransitionStateReconcileCondition>>,
|
||||
}
|
||||
|
||||
/// An exact-copy precondition for the single-version tier repair protocol.
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct TransitionStateReconcileCondition {
|
||||
pub expected_metadata_digest: String,
|
||||
pub unchanged_metadata_digest: String,
|
||||
pub target: rustfs_filemeta::TransitionStateReconcileTarget,
|
||||
pub tier: String,
|
||||
pub topology_generation: String,
|
||||
pub verify_only: bool,
|
||||
/// Local ownership is never accepted from the wire. A remote disk acquires
|
||||
/// its own fleet and backend leases before entering the mutation domain.
|
||||
#[serde(skip)]
|
||||
pub(crate) authority:
|
||||
Option<Arc<crate::bucket::lifecycle::legacy_transition_state_reconcile::TransitionStateReconcileAuthority>>,
|
||||
}
|
||||
|
||||
pub struct DiskLocation {
|
||||
|
||||
@@ -1449,6 +1449,38 @@ fn disk_namespace_mutation_lock(path: &Path) -> Arc<NamespaceMutationLock> {
|
||||
lock
|
||||
}
|
||||
|
||||
static DISK_METADATA_MUTATION_LOCKS: LazyLock<Mutex<NamespaceMutationLockRegistry>> =
|
||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
|
||||
/// Serializes the complete xl.meta read/modify/commit transaction. This domain
|
||||
/// precedes namespace/volume publication locks, whose narrower syscall leases
|
||||
/// may retain it after cancellation of the async caller.
|
||||
pub(crate) struct MetadataMutationLease {
|
||||
_guard: OwnedMutexGuard<()>,
|
||||
_owner: Option<Arc<dyn Send + Sync>>,
|
||||
}
|
||||
|
||||
pub(crate) async fn acquire_metadata_mutation_lease(
|
||||
object: &Path,
|
||||
owner: Option<Arc<dyn Send + Sync>>,
|
||||
) -> Arc<MetadataMutationLease> {
|
||||
let lock = {
|
||||
let mut locks = DISK_METADATA_MUTATION_LOCKS.lock();
|
||||
locks.retain(|_, lock| lock.strong_count() > 0);
|
||||
if let Some(lock) = locks.get(object).and_then(Weak::upgrade) {
|
||||
lock
|
||||
} else {
|
||||
let lock = Arc::new(AsyncMutex::new(()));
|
||||
locks.insert(object.to_path_buf(), Arc::downgrade(&lock));
|
||||
lock
|
||||
}
|
||||
};
|
||||
Arc::new(MetadataMutationLease {
|
||||
_guard: lock.lock_owned().await,
|
||||
_owner: owner,
|
||||
})
|
||||
}
|
||||
|
||||
/// Keeps a namespace transaction serialized even when its async waiter is
|
||||
/// cancelled while a blocking filesystem call is still running.
|
||||
pub(crate) struct NamespaceMutationLease {
|
||||
@@ -2112,6 +2144,41 @@ pub(crate) async fn rename_all_with_lease(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Publish a conditional repair and sync its directory in one owned executor.
|
||||
/// Cancellation cannot release its metadata/fleet/tier leases between rename
|
||||
/// and fsync. The last authority check runs after destination preparation.
|
||||
pub(crate) async fn rename_reconciled_metadata(
|
||||
source: PathBuf,
|
||||
destination: PathBuf,
|
||||
base_dir: PathBuf,
|
||||
publication_root: PublicationRoot,
|
||||
owner: Option<Arc<dyn Send + Sync>>,
|
||||
authority: Arc<crate::bucket::lifecycle::legacy_transition_state_reconcile::TransitionStateReconcileAuthority>,
|
||||
) -> Result<()> {
|
||||
let lease = acquire_namespace_mutation_lease_with_owner(&destination, owner).await;
|
||||
run_blocking_namespace_operation(lease, move || {
|
||||
let preparation = prepare_rename_with_retry(&source, &destination, &base_dir, &publication_root)?;
|
||||
#[cfg(all(any(test, feature = "test-util"), not(windows)))]
|
||||
prepared_publication_test_hooks::run(prepared_publication_test_hooks::Stage::Rename, &destination);
|
||||
if !authority.is_current() {
|
||||
return Err(io::Error::new(io::ErrorKind::WouldBlock, "transition reconciliation authority expired"));
|
||||
}
|
||||
rename_prepared(&source, &destination, &preparation)?;
|
||||
if let Some(parent) = destination.parent() {
|
||||
fsync_dir_std(parent)?;
|
||||
}
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.map_err(|err| {
|
||||
if err.kind() == io::ErrorKind::WouldBlock {
|
||||
DiskError::OutdatedXLMeta
|
||||
} else {
|
||||
to_file_error(err).into()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tracing::instrument(level = "debug", skip_all)]
|
||||
pub(crate) async fn rename_all_with_commit_guard(
|
||||
|
||||
@@ -67,12 +67,9 @@ const DECOMMISSION_TARGET_FENCE_POLICY_SUPPORTED_VERSION: u32 = 4;
|
||||
// Keep this synchronized with the version served by node_service. Including
|
||||
// the local member in the minimum prevents an older coordinator from
|
||||
// self-authorizing a policy implemented only by newer remote peers.
|
||||
const LOCAL_CROSS_POOL_FENCE_POLICY_SUPPORTED_VERSION: u32 = 4;
|
||||
/// Version 5 is reserved for a fleet whose every metadata writer preserves
|
||||
/// explicit transition version state and destination identity, and implements
|
||||
/// conditional per-generation `xl.meta` writes with strong readback. The node
|
||||
/// service must not advertise this version until the conditional writer from
|
||||
/// rustfs/backlog#684 is available.
|
||||
const LOCAL_CROSS_POOL_FENCE_POLICY_SUPPORTED_VERSION: u32 = 5;
|
||||
/// Version 5 preserves explicit transition state/destination bindings and
|
||||
/// supports exact-generation metadata repair with strong all-copy readback.
|
||||
const LEGACY_TRANSITION_STATE_RECONCILE_POLICY_SUPPORTED_VERSION: u32 = 5;
|
||||
|
||||
fn resolve_admin_peer_probe_timeout_secs(configured: Option<u64>) -> u64 {
|
||||
@@ -607,6 +604,10 @@ fn acquire_legacy_transition_state_reconcile_fleet_proof_from(
|
||||
}
|
||||
|
||||
async fn observe_legacy_transition_state_reconcile_fleet(expected_topology: &str) -> Option<BTreeMap<String, Uuid>> {
|
||||
#[cfg(all(test, feature = "test-util"))]
|
||||
if let Ok(observation) = LEGACY_RECONCILE_TEST_OBSERVATION.try_with(Clone::clone) {
|
||||
return Some(observation);
|
||||
}
|
||||
let notification_sys = get_global_notification_sys()?;
|
||||
let (peer_epochs, minimum_version) = timeout(
|
||||
REMOTE_VERSION_STATE_PROBE_TIMEOUT,
|
||||
@@ -619,6 +620,34 @@ async fn observe_legacy_transition_state_reconcile_fleet(expected_topology: &str
|
||||
reconcile_result.ok()
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "test-util"))]
|
||||
tokio::task_local! {
|
||||
static LEGACY_RECONCILE_TEST_OBSERVATION: BTreeMap<String, Uuid>;
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "test-util"))]
|
||||
pub(crate) async fn with_legacy_transition_state_fleet_proof_for_test<F: std::future::Future>(future: F) -> F::Output {
|
||||
struct Revoke;
|
||||
impl Drop for Revoke {
|
||||
fn drop(&mut self) {
|
||||
revoke_fleet_capability_proof(legacy_transition_state_reconcile_fleet_proof_slot());
|
||||
}
|
||||
}
|
||||
let topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get().expect("test store topology");
|
||||
assert!(
|
||||
publish_fleet_capability_probe_result(
|
||||
legacy_transition_state_reconcile_fleet_proof_slot(),
|
||||
topology,
|
||||
Ok(BTreeMap::new()),
|
||||
Instant::now(),
|
||||
)
|
||||
.is_none()
|
||||
);
|
||||
let _revoke = Revoke;
|
||||
let _remote_version = install_current_remote_version_state_fleet_proof_for_test();
|
||||
LEGACY_RECONCILE_TEST_OBSERVATION.scope(BTreeMap::new(), future).await
|
||||
}
|
||||
|
||||
/// Revalidate the exact fleet generation captured by a reconcile token with a
|
||||
/// fresh synchronous observation. Callers must await this before each
|
||||
/// conditional metadata write and after the final strong readback.
|
||||
@@ -637,6 +666,18 @@ pub async fn legacy_transition_state_reconcile_fleet_proof_matches(
|
||||
.await
|
||||
}
|
||||
|
||||
/// Final local check in the disk publication executor. The corresponding
|
||||
/// counted permit remains owned until the filesystem operation has drained.
|
||||
pub(crate) fn legacy_transition_state_reconcile_fleet_proof_current(
|
||||
proof: &LegacyTransitionStateReconcileFleetProofToken,
|
||||
) -> bool {
|
||||
let Some(topology) = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get() else { return false };
|
||||
let state = legacy_transition_state_reconcile_fleet_proof_slot()
|
||||
.read()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
legacy_transition_state_reconcile_fleet_proof_matches_at(&state, proof, topology, Instant::now())
|
||||
}
|
||||
|
||||
pub async fn acquire_ilm_recovery_export_fleet_proof() -> Option<IlmRecoveryExportFleetProofToken> {
|
||||
let expected_topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get()?;
|
||||
let proof = {
|
||||
@@ -3926,15 +3967,11 @@ mod tests {
|
||||
assert!(decommission_v3.is_err(), "v3 members do not understand the per-target decommission fence");
|
||||
assert!(reconcile_v3.is_err());
|
||||
|
||||
let (generic_v4, journal_v4, decommission_v4, reconcile_v4) =
|
||||
cross_pool_fence_policy_results(peers.clone(), LOCAL_CROSS_POOL_FENCE_POLICY_SUPPORTED_VERSION);
|
||||
let (generic_v4, journal_v4, decommission_v4, reconcile_v4) = cross_pool_fence_policy_results(peers.clone(), 4);
|
||||
assert!(generic_v4.is_ok());
|
||||
assert!(journal_v4.is_ok());
|
||||
assert!(decommission_v4.is_ok(), "an all-v4 fleet may create sticky per-target reservations");
|
||||
assert!(
|
||||
reconcile_v4.is_err(),
|
||||
"the current local policy lacks the conditional xl.meta writer required by reconcile"
|
||||
);
|
||||
assert!(reconcile_v4.is_err(), "v4 does not support conditional transition metadata writes");
|
||||
|
||||
let (generic_v5, journal_v5, decommission_v5, reconcile_v5) = cross_pool_fence_policy_results(peers, 5);
|
||||
assert!(generic_v5.is_ok());
|
||||
@@ -4612,7 +4649,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn legacy_transition_state_reconcile_single_node_stays_closed_before_local_cas_support() {
|
||||
async fn legacy_transition_state_reconcile_single_node_advertises_conditional_writer() {
|
||||
let notification_sys = NotificationSys {
|
||||
peer_clients: Vec::new(),
|
||||
all_peer_clients: vec![None],
|
||||
@@ -4628,8 +4665,8 @@ mod tests {
|
||||
assert_eq!(minimum_version, LOCAL_CROSS_POOL_FENCE_POLICY_SUPPORTED_VERSION);
|
||||
let (_, _, _, reconcile_result) = cross_pool_fence_policy_results(peers, minimum_version);
|
||||
assert!(
|
||||
reconcile_result.is_err(),
|
||||
"the current node must not self-authorize reconcile before the conditional writer lands"
|
||||
reconcile_result.is_ok(),
|
||||
"the current node implements the conditional writer and preserves repaired bindings"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -51,12 +51,12 @@ use super::super::{
|
||||
can_try_inline_data_shards_direct, capacity_scope_from_disks, codec_streaming_rollout_applies, coding,
|
||||
collect_inline_data_shard_fileinfos_by_index_or_reason, current_dirty_generation, debug, disk,
|
||||
file_info_is_valid_for_metadata, get_metadata_slowtail_fault_request, info, inline_erasure_shard_file_offset,
|
||||
inline_erasure_shard_size, is_err_object_not_found, is_err_version_not_found, is_get_metadata_data_read_early_stop_enabled,
|
||||
is_get_metadata_early_stop_bounded_fanout_enabled, is_get_metadata_early_stop_enabled,
|
||||
is_get_metadata_non_inline_data_read_early_stop_enabled, is_object_dangling, is_version_early_stop_enabled,
|
||||
issue3031_diag_enabled, join_all, join_errs, log_multipart_write_quorum_failure, merge_file_meta_versions,
|
||||
object_fits_single_block, path_join_buf, record_global_dirty_scope, reduce_read_quorum_errs, reduce_write_quorum_errs,
|
||||
send_heal_request_with_admission, should_prevent_write, to_object_err, try_read_inline_data_shards_direct, warn,
|
||||
inline_erasure_shard_size, is_get_metadata_data_read_early_stop_enabled, is_get_metadata_early_stop_bounded_fanout_enabled,
|
||||
is_get_metadata_early_stop_enabled, is_get_metadata_non_inline_data_read_early_stop_enabled, is_object_dangling,
|
||||
is_version_early_stop_enabled, issue3031_diag_enabled, join_all, join_errs, log_multipart_write_quorum_failure,
|
||||
merge_file_meta_versions, object_fits_single_block, path_join_buf, record_global_dirty_scope, reduce_read_quorum_errs,
|
||||
reduce_write_quorum_errs, 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;
|
||||
@@ -3086,21 +3086,61 @@ impl SetDisks {
|
||||
let read_quorum = disks.len().div_ceil(2).max(1);
|
||||
let (raw_fileinfos, errs) = Self::read_all_raw_file_info(&disks, bucket, disk_object.as_str(), false).await;
|
||||
|
||||
if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) {
|
||||
let object_err = to_object_err(err.into(), vec![bucket, object]);
|
||||
if is_err_object_not_found(&object_err) || is_err_version_not_found(&object_err) {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(object_err);
|
||||
if let Some(err) = errs
|
||||
.iter()
|
||||
.flatten()
|
||||
.find(|err| !matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound))
|
||||
{
|
||||
return Err(to_object_err(err.clone().into(), vec![bucket, object]));
|
||||
}
|
||||
// A minority live owner must not disappear behind majority absence.
|
||||
// Only explicit absence on every readable disk proves no ownership.
|
||||
if raw_fileinfos.iter().all(Option::is_none) {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let mut shallow_versions = Vec::with_capacity(raw_fileinfos.len());
|
||||
type TransitionCopy = (FileInfo, Option<crate::services::tier::tier::TierDestinationId>);
|
||||
let mut transition_copies: std::collections::HashMap<Option<Uuid>, Vec<TransitionCopy>> =
|
||||
std::collections::HashMap::new();
|
||||
let decode_error = |err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}"));
|
||||
for raw_fileinfo in raw_fileinfos.into_iter().flatten() {
|
||||
let meta = FileMeta::load(&raw_fileinfo.buf)
|
||||
.map_err(|err| Error::other(format!("exact object metadata decode failed for {bucket}/{object}: {err}")))?;
|
||||
let versions = meta.get_all_file_info_versions(bucket, object, true).map_err(decode_error)?;
|
||||
for version in versions.versions.into_iter().chain(versions.free_versions) {
|
||||
if version.transition_status != rustfs_filemeta::TRANSITION_COMPLETE {
|
||||
continue;
|
||||
}
|
||||
let destination =
|
||||
crate::services::tier::tier::tier_destination_id_from_metadata(&version.metadata).map_err(Error::other)?;
|
||||
transition_copies
|
||||
.entry(version.version_id.filter(|id| !id.is_nil()))
|
||||
.or_default()
|
||||
.push((version, destination));
|
||||
}
|
||||
shallow_versions.push(meta.versions);
|
||||
}
|
||||
|
||||
// Exact cleanup/recovery reads must not select a repaired majority
|
||||
// while another physical copy still carries legacy absence. Missing
|
||||
// copies permit deletion retries; an unreadable disk proves nothing.
|
||||
for copies in transition_copies
|
||||
.values()
|
||||
.filter(|copies| copies.iter().any(|(_, destination)| destination.is_some()))
|
||||
{
|
||||
let (first, destination) = &copies[0];
|
||||
if copies.iter().any(|(copy, identity)| {
|
||||
identity != destination
|
||||
|| copy.transition_version_state != first.transition_version_state
|
||||
|| copy.transition_version != first.transition_version
|
||||
|| copy.transition_tier != first.transition_tier
|
||||
|| copy.transitioned_objname != first.transitioned_objname
|
||||
}) {
|
||||
return Err(Error::other("exact transition metadata has not converged across physical copies"));
|
||||
}
|
||||
}
|
||||
|
||||
if shallow_versions.len() < read_quorum {
|
||||
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
|
||||
}
|
||||
@@ -3117,7 +3157,7 @@ impl SetDisks {
|
||||
..Default::default()
|
||||
}
|
||||
.get_all_file_info_versions(bucket, object, true)
|
||||
.map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}")))?;
|
||||
.map_err(decode_error)?;
|
||||
|
||||
for file_info in file_info_versions
|
||||
.versions
|
||||
@@ -5815,17 +5855,17 @@ impl SetDisks {
|
||||
}
|
||||
|
||||
if let Some(disk) = disks[i].as_ref() {
|
||||
let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]);
|
||||
// A failed version-only copy owns only its new version.
|
||||
// Removing the whole xl.meta would also erase existing
|
||||
// versions and any concurrently reconciled tier binding.
|
||||
let mut rollback = FileInfo {
|
||||
version_id: files[i].version_id,
|
||||
..Default::default()
|
||||
};
|
||||
rollback.set_skip_tier_free_version();
|
||||
revert_futures.push(async move {
|
||||
if let Err(err) = disk
|
||||
.delete(
|
||||
bucket,
|
||||
&path,
|
||||
DeleteOptions {
|
||||
recursive: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.delete_version(bucket, prefix, rollback, false, DeleteOptions::default())
|
||||
.await
|
||||
{
|
||||
warn!("write meta revert err {:?}", err);
|
||||
@@ -12092,7 +12132,9 @@ mod tests {
|
||||
let bucket = "write-unique-bucket";
|
||||
let object = "object";
|
||||
let (_dir, disk) = read_multiple_test_disk(bucket, &[]).await;
|
||||
let files = vec![metadata_test_fileinfo(object), metadata_test_fileinfo(object)];
|
||||
let mut fi = metadata_test_fileinfo(object);
|
||||
fi.mod_time = Some(OffsetDateTime::now_utc());
|
||||
let files = vec![fi.clone(), fi];
|
||||
|
||||
let result = SetDisks::write_unique_file_info(&[Some(disk.clone()), None], bucket, bucket, object, &files, 2).await;
|
||||
|
||||
@@ -12106,6 +12148,53 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn write_unique_file_info_rollback_preserves_existing_reconciled_version() {
|
||||
let bucket = "write-unique-existing";
|
||||
let object = "object";
|
||||
let (_dir, disk) = read_multiple_test_disk(bucket, &[]).await;
|
||||
let mut original = metadata_test_fileinfo(object);
|
||||
original.version_id = Some(Uuid::from_u128(1));
|
||||
original.data_dir = Some(Uuid::from_u128(3));
|
||||
original.mod_time = Some(OffsetDateTime::now_utc());
|
||||
original.transition_status = rustfs_filemeta::TRANSITION_COMPLETE.to_string();
|
||||
original.transition_tier = "WARM".to_string();
|
||||
original.transitioned_objname = "remote-original".to_string();
|
||||
original.transition_version_state = rustfs_filemeta::TransitionVersionState::KnownDisabled;
|
||||
rustfs_utils::http::insert_str(
|
||||
&mut original.metadata,
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
"ab".repeat(32),
|
||||
);
|
||||
disk.write_metadata(bucket, bucket, object, original.clone())
|
||||
.await
|
||||
.expect("existing reconciled source");
|
||||
let raw = disk
|
||||
.read_all(bucket, &format!("{object}/{STORAGE_FORMAT_FILE}"))
|
||||
.await
|
||||
.expect("original metadata");
|
||||
let before = FileMeta::load(&raw).unwrap().find_version(original.version_id).unwrap().1;
|
||||
let mut added = original.clone();
|
||||
added.version_id = Some(Uuid::from_u128(2));
|
||||
let result =
|
||||
SetDisks::write_unique_file_info(&[Some(disk.clone()), None], bucket, bucket, object, &[added.clone(), added], 2)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
let raw = disk
|
||||
.read_all(bucket, &format!("{object}/{STORAGE_FORMAT_FILE}"))
|
||||
.await
|
||||
.expect("preserved xl.meta");
|
||||
let after = FileMeta::load(&raw).expect("preserved metadata");
|
||||
assert_eq!(
|
||||
after
|
||||
.find_version(original.version_id)
|
||||
.expect("original source survives rollback")
|
||||
.1,
|
||||
before
|
||||
);
|
||||
assert!(after.find_version(Some(Uuid::from_u128(2))).is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn update_object_meta_handles_empty_metadata_and_missing_quorum() {
|
||||
let set = io_primitives_test_set(vec![None, None], 1).await;
|
||||
|
||||
@@ -3904,7 +3904,9 @@ pub(crate) async fn read_legacy_transition_state_metadata_copies(
|
||||
return Err(DiskError::DiskNotFound);
|
||||
}
|
||||
|
||||
let (copies, errs) = SetDisks::read_all_raw_file_info(&disks, bucket, disk_object.as_str(), false).await;
|
||||
// Include inline bytes in the generation: a conditional repair preserves
|
||||
// the entire xl.meta, including payloads belonging to other versions.
|
||||
let (copies, errs) = SetDisks::read_all_raw_file_info(&disks, bucket, disk_object.as_str(), true).await;
|
||||
for err in errs.into_iter().flatten() {
|
||||
if !matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) {
|
||||
return Err(err);
|
||||
@@ -4262,7 +4264,7 @@ impl SetDisks {
|
||||
self.get_object_metadata_cache_generations[generation.index].load(Ordering::Acquire) == generation.value
|
||||
}
|
||||
|
||||
async fn invalidate_get_object_metadata_cache(&self, bucket: &str, object: &str) {
|
||||
pub(crate) async fn invalidate_get_object_metadata_cache(&self, bucket: &str, object: &str) {
|
||||
let hash = self.get_object_metadata_cache_hash(bucket, object);
|
||||
let hash_bytes = hash.to_le_bytes();
|
||||
let index = usize::from(u16::from_le_bytes([hash_bytes[0], hash_bytes[1]]) % GET_OBJECT_METADATA_CACHE_FENCE_SHARDS);
|
||||
|
||||
@@ -12807,11 +12807,25 @@ mod tests {
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn legacy_transition_state_inspection_and_apply_keep_all_disk_copies_unchanged() {
|
||||
run_large_stack_async_test("legacy-state-reconcile-inspection", legacy_transition_state_inspection_and_apply_case);
|
||||
run_large_stack_async_test("legacy-state-reconcile-inspection", || {
|
||||
legacy_transition_state_inspection_and_apply_case(false)
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
async fn legacy_transition_state_inspection_and_apply_case() {
|
||||
#[cfg(not(windows))]
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn legacy_transition_state_backfill_retries_partial_commits_and_preserves_other_bytes() {
|
||||
run_large_stack_async_test("legacy-state-reconcile-backfill", || {
|
||||
legacy_transition_state_inspection_and_apply_case(true)
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
async fn legacy_transition_state_inspection_and_apply_case(write_enabled: bool) {
|
||||
#[cfg(windows)]
|
||||
assert!(!write_enabled, "Windows supports inspection but cannot prove repair directory durability");
|
||||
use crate::bucket::lifecycle::legacy_transition_state_reconcile::{
|
||||
LegacyTransitionStateReconcileOutcome as Outcome, LegacyTransitionStateReconcileRequest,
|
||||
LegacyTransitionStateReconcileSelector,
|
||||
@@ -12923,6 +12937,187 @@ mod tests {
|
||||
target,
|
||||
reconciliation_digest: inspected.reconciliation_digest.expect("expected tuple digest"),
|
||||
};
|
||||
#[cfg(not(windows))]
|
||||
if write_enabled {
|
||||
crate::services::notification_sys::with_legacy_transition_state_fleet_proof_for_test(async {
|
||||
crate::disk::local::bucket_durability::set(bucket, Some(crate::disk::local::DurabilityMode::None));
|
||||
let unsynced = store.reconcile_legacy_transition_state(request.clone()).await;
|
||||
crate::disk::local::bucket_durability::set(bucket, None);
|
||||
let unsynced = unsynced.expect("repair without metadata durability");
|
||||
assert_eq!(unsynced.outcome, Outcome::BackendUnavailable);
|
||||
assert!(!unsynced.changed);
|
||||
let rollback = paths[0].parent().expect("object directory").join(Uuid::new_v4().to_string());
|
||||
tokio::fs::create_dir(&rollback).await.expect("pending rollback directory");
|
||||
tokio::fs::write(rollback.join(crate::disk::STORAGE_FORMAT_FILE_BACKUP), &original[0])
|
||||
.await
|
||||
.expect("pending old metadata backup");
|
||||
let unsettled = store.reconcile_legacy_transition_state(request.clone()).await;
|
||||
tokio::fs::remove_dir_all(&rollback).await.expect("settle fixture rollback");
|
||||
let unsettled = unsettled.expect("repair must wait for rollback");
|
||||
assert_eq!(unsettled.outcome, Outcome::BackendUnavailable);
|
||||
assert!(!unsettled.changed);
|
||||
for (path, bytes) in paths.iter().zip(&original) {
|
||||
assert_eq!(tokio::fs::read(path).await.expect("blocked repair leaves original bytes"), *bytes);
|
||||
}
|
||||
// The first disk commits; the second stops after staging.
|
||||
// This models an interrupted cross-disk effect without rollback.
|
||||
let disks = store.all_set_disks()[0].disk_inventory().await;
|
||||
let first_disk = disks[0].as_ref().expect("first physical disk");
|
||||
let publication_path = first_disk
|
||||
.get_object_path_for_io_if_local(bucket, &format!("{object}/{STORAGE_FORMAT_FILE}"))
|
||||
.expect("local disk")
|
||||
.expect("publication path");
|
||||
let crash_key = format!("{object}/{STORAGE_FORMAT_FILE}");
|
||||
let _hook = crate::disk::os::prepared_publication_test_hooks::install_at(
|
||||
crate::disk::os::prepared_publication_test_hooks::Stage::Rename,
|
||||
&publication_path,
|
||||
move || {
|
||||
crate::crash_inject::arm(crate::crash_inject::CrashPoint::MetaWriteAfterTmpBeforeRename, &crash_key);
|
||||
},
|
||||
);
|
||||
let partial = store
|
||||
.reconcile_legacy_transition_state(request.clone())
|
||||
.await
|
||||
.expect("partial repair response");
|
||||
assert_eq!(partial.outcome, Outcome::BackendUnavailable, "{partial:?}");
|
||||
assert!(partial.changed, "first copy was committed: {partial:?}");
|
||||
assert!(partial.changes_indeterminate);
|
||||
assert_ne!(tokio::fs::read(&paths[0]).await.expect("first committed copy"), original[0]);
|
||||
for (path, bytes) in paths[1..].iter().zip(&original[1..]) {
|
||||
assert_eq!(tokio::fs::read(path).await.expect("uncommitted copy"), *bytes);
|
||||
}
|
||||
assert!(
|
||||
store.all_set_disks()[0]
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.is_err(),
|
||||
"cleanup cannot select a partial repair subset"
|
||||
);
|
||||
let repaired = store
|
||||
.reconcile_legacy_transition_state(request.clone())
|
||||
.await
|
||||
.expect("retry original snapshot");
|
||||
assert_eq!(repaired.outcome, Outcome::Migrated, "{repaired:?}");
|
||||
assert!(repaired.changed);
|
||||
assert!(!repaired.changes_indeterminate);
|
||||
let mut committed = Vec::new();
|
||||
for (path, original) in paths.iter().zip(&original) {
|
||||
let raw = tokio::fs::read(path).await.expect("repaired copy");
|
||||
let metadata = FileMeta::load(&raw).expect("decode repaired copy");
|
||||
let previous = FileMeta::load(original).expect("decode original copy");
|
||||
assert_eq!(
|
||||
metadata.transition_reconcile_generation(None).unwrap(),
|
||||
previous.transition_reconcile_generation(None).unwrap()
|
||||
);
|
||||
let (_, version) = metadata.find_version(None).expect("selected version");
|
||||
let info = version.into_fileinfo(bucket, object, true).expect("repaired FileInfo");
|
||||
assert_eq!(info.transition_version_state, expected_state);
|
||||
assert_eq!(info.transition_version, request.target.remote_version);
|
||||
committed.push(raw);
|
||||
}
|
||||
assert!(
|
||||
store.all_set_disks()[0]
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect("converged cleanup snapshot")
|
||||
.is_some()
|
||||
);
|
||||
let replay = store
|
||||
.reconcile_legacy_transition_state(request.clone())
|
||||
.await
|
||||
.expect("idempotent original request replay");
|
||||
assert_eq!(replay.outcome, Outcome::Migrated, "{replay:?}");
|
||||
assert!(!replay.changed);
|
||||
for (path, expected) in paths.iter().zip(&committed) {
|
||||
assert_eq!(
|
||||
tokio::fs::read(path).await.expect("replayed copy"),
|
||||
*expected,
|
||||
"idempotence preserves raw encoding"
|
||||
);
|
||||
}
|
||||
for (path, bytes) in paths.iter().zip(&original) {
|
||||
tokio::fs::write(path, bytes)
|
||||
.await
|
||||
.expect("reset independent cancellation fixture");
|
||||
}
|
||||
let (entered_tx, entered) = tokio::sync::oneshot::channel();
|
||||
let (release, released) = std::sync::mpsc::channel::<()>();
|
||||
let _pause = crate::disk::os::prepared_publication_test_hooks::install_at(
|
||||
crate::disk::os::prepared_publication_test_hooks::Stage::Rename,
|
||||
&publication_path,
|
||||
move || {
|
||||
let _ = entered_tx.send(());
|
||||
let _ = released.recv();
|
||||
},
|
||||
);
|
||||
let mut repair = Box::pin(store.reconcile_legacy_transition_state(request.clone()));
|
||||
tokio::select! {
|
||||
result = &mut repair => panic!("repair completed before publication pause: {result:?}"),
|
||||
result = entered => result.expect("publication executor entered"),
|
||||
}
|
||||
let update_options = crate::disk::UpdateMetadataOpts::default();
|
||||
let mut update = Box::pin(first_disk.update_metadata(
|
||||
bucket,
|
||||
object,
|
||||
FileInfo {
|
||||
metadata: HashMap::from([("x-amz-meta-concurrent".to_string(), "kept".to_string())]),
|
||||
..Default::default()
|
||||
},
|
||||
&update_options,
|
||||
));
|
||||
assert!(
|
||||
tokio::time::timeout(std::time::Duration::from_millis(25), update.as_mut())
|
||||
.await
|
||||
.is_err(),
|
||||
"another metadata RMW must wait for publication"
|
||||
);
|
||||
drop(repair);
|
||||
assert!(
|
||||
tokio::time::timeout(std::time::Duration::from_millis(25), update.as_mut())
|
||||
.await
|
||||
.is_err(),
|
||||
"cancelling the coordinator must not release an in-flight disk mutation"
|
||||
);
|
||||
release.send(()).expect("resume owned publication");
|
||||
update.await.expect("serialized metadata update");
|
||||
let raw = tokio::fs::read(&paths[0])
|
||||
.await
|
||||
.expect("cancelled repair and later metadata update");
|
||||
let (_, version) = FileMeta::load(&raw)
|
||||
.expect("metadata after cancellation")
|
||||
.find_version(None)
|
||||
.expect("selected version");
|
||||
let info = version
|
||||
.into_fileinfo(bucket, object, true)
|
||||
.expect("metadata after serialized update");
|
||||
assert_eq!(info.transition_version_state, expected_state);
|
||||
assert_eq!(info.metadata.get("x-amz-meta-concurrent").map(String::as_str), Some("kept"));
|
||||
let stale = store
|
||||
.reconcile_legacy_transition_state(request.clone())
|
||||
.await
|
||||
.expect("stale original request");
|
||||
assert_eq!(
|
||||
stale.outcome,
|
||||
Outcome::Corrupt,
|
||||
"unrelated metadata change invalidates the original generation: {stale:?}"
|
||||
);
|
||||
assert!(!stale.changed);
|
||||
assert_eq!(tokio::fs::read(&paths[0]).await.expect("stale write leaves bytes unchanged"), raw);
|
||||
assert_eq!(backend.remove_count().await, 0);
|
||||
// A pinned version probe uses GET to verify that exact
|
||||
// candidate; every backend operation still targets it.
|
||||
let operations = backend.op_log().await;
|
||||
assert!(
|
||||
operations.iter().all(|operation| match operation {
|
||||
MockWarmOp::Probe { object } | MockWarmOp::Get { object } => object == &request.source.remote_object,
|
||||
_ => false,
|
||||
}),
|
||||
"unexpected backend effects: {operations:?}"
|
||||
);
|
||||
})
|
||||
.await;
|
||||
continue;
|
||||
}
|
||||
let mut tampered = request.clone();
|
||||
tampered.source.remote_object.push_str("-other");
|
||||
let probes_before = backend.op_log().await.len();
|
||||
|
||||
@@ -50,6 +50,9 @@ use tracing::{error, warn};
|
||||
use uuid::Uuid;
|
||||
use xxhash_rust::xxh64;
|
||||
|
||||
mod transition_reconcile;
|
||||
pub use transition_reconcile::TransitionStateReconcileTarget;
|
||||
|
||||
// XL header specifies the format
|
||||
pub static XL_FILE_HEADER: [u8; 4] = *b"XL2 ";
|
||||
// pub static XL_FILE_VERSION_CURRENT: [u8; 4] = [0; 4];
|
||||
@@ -391,6 +394,16 @@ impl FileMeta {
|
||||
|
||||
if ver_vid == fi_vid {
|
||||
let mut ver = FileMetaVersion::try_from(version.meta.as_slice())?;
|
||||
let previous = ver
|
||||
.object
|
||||
.as_ref()
|
||||
.is_some_and(|object| {
|
||||
rustfs_utils::http::contains_key_bytes(
|
||||
&object.meta_sys,
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
)
|
||||
})
|
||||
.then(|| ver.clone());
|
||||
|
||||
if let Some(ref mut obj) = ver.object {
|
||||
if replace_user_metadata {
|
||||
@@ -447,6 +460,9 @@ impl FileMeta {
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(previous) = previous {
|
||||
transition_reconcile::preserve_reconciled_transition(&previous, &mut ver)?;
|
||||
}
|
||||
// Update
|
||||
version.header = ver.header();
|
||||
version.meta = ver.marshal_msg()?;
|
||||
@@ -492,7 +508,7 @@ impl FileMeta {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn add_version_filemata(&mut self, version: FileMetaVersion) -> Result<()> {
|
||||
pub fn add_version_filemata(&mut self, mut version: FileMetaVersion) -> Result<()> {
|
||||
if !version.valid() {
|
||||
return Err(Error::other("file meta version invalid"));
|
||||
}
|
||||
@@ -512,6 +528,7 @@ impl FileMeta {
|
||||
if existing.free_version() != version.free_version() {
|
||||
return Err(Error::other("cannot replace a free version with a non-free version"));
|
||||
}
|
||||
transition_reconcile::preserve_reconciled_transition(&existing, &mut version)?;
|
||||
return self.set_idx(fidx, version);
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,384 @@
|
||||
// Copyright 2026 RustFS Team
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
use super::{FileMeta, FileMetaVersion};
|
||||
use crate::{Error, Result, TRANSITION_COMPLETE, TransitionVersionState};
|
||||
use rustfs_utils::http::metadata_compat::{
|
||||
SUFFIX_TRANSITION_TIER_DESTINATION_ID, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_bytes,
|
||||
get_consistent_bytes, insert_bytes, remove_bytes,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use uuid::Uuid;
|
||||
|
||||
const RECONCILE_SUFFIXES: [&str; 3] = [
|
||||
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
];
|
||||
|
||||
/// The only fields a legacy transition repair is allowed to persist.
|
||||
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct TransitionStateReconcileTarget {
|
||||
pub state: TransitionVersionState,
|
||||
pub remote_version: Option<String>,
|
||||
pub destination_id: String,
|
||||
}
|
||||
|
||||
impl TransitionStateReconcileTarget {
|
||||
pub fn validate(&self) -> Result<()> {
|
||||
let valid_version = match self.state {
|
||||
TransitionVersionState::KnownDisabled => self.remote_version.is_none(),
|
||||
TransitionVersionState::SuspendedNull => self.remote_version.as_deref() == Some("null"),
|
||||
TransitionVersionState::Exact => self.remote_version.as_deref().is_some_and(|value| {
|
||||
!value.is_empty()
|
||||
&& value.len() <= 1024
|
||||
&& value != "null"
|
||||
&& !value.chars().any(char::is_control)
|
||||
&& !Uuid::parse_str(value).is_ok_and(|id| id.is_nil())
|
||||
}),
|
||||
TransitionVersionState::Unknown => false,
|
||||
};
|
||||
if !valid_version
|
||||
|| self.destination_id.len() != 64
|
||||
|| !self
|
||||
.destination_id
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl FileMeta {
|
||||
/// Canonical identity of every version and inline byte, excluding only the
|
||||
/// three repairable suffixes on the selected version. It survives a repair
|
||||
/// and encoding-order changes, while detecting unrelated metadata changes.
|
||||
pub fn transition_reconcile_generation(&self, version_id: Option<Uuid>) -> Result<Vec<u8>> {
|
||||
if self
|
||||
.versions
|
||||
.iter()
|
||||
.filter(|version| version.header.version_id.unwrap_or_default() == version_id.unwrap_or_default())
|
||||
.count()
|
||||
!= 1
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
let (selected, _) = self.find_version(version_id)?;
|
||||
let mut versions = Vec::with_capacity(self.versions.len());
|
||||
for index in 0..self.versions.len() {
|
||||
let mut version = self.get_idx(index)?;
|
||||
if index == selected {
|
||||
let object = version.object.as_mut().ok_or(Error::FileCorrupt)?;
|
||||
for suffix in RECONCILE_SUFFIXES {
|
||||
remove_bytes(&mut object.meta_sys, suffix);
|
||||
}
|
||||
}
|
||||
versions.push(version);
|
||||
}
|
||||
versions.sort_by_key(|version| version.get_version_id().unwrap_or_default());
|
||||
let mut value = serde_json::to_value((&versions, &self.data)).map_err(|_| Error::FileCorrupt)?;
|
||||
value.sort_all_objects();
|
||||
serde_json::to_vec(&value).map_err(|_| Error::FileCorrupt)
|
||||
}
|
||||
|
||||
/// Returns false for an already converged record. Callers must serialize
|
||||
/// the read/check/commit and compare the observed metadata generation.
|
||||
pub fn reconcile_transition_state(
|
||||
&mut self,
|
||||
version_id: Option<Uuid>,
|
||||
target: &TransitionStateReconcileTarget,
|
||||
) -> Result<bool> {
|
||||
target.validate()?;
|
||||
let (index, mut version) = self.find_version(version_id)?;
|
||||
let info = version.into_fileinfo("", "", true)?;
|
||||
info.validate_for_metadata_read()?;
|
||||
if info.transition_status != TRANSITION_COMPLETE
|
||||
|| info.transition_tier.is_empty()
|
||||
|| info.transitioned_objname.is_empty()
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
let object = version.object.as_mut().ok_or(Error::FileCorrupt)?;
|
||||
let destination = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID);
|
||||
if contains_key_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID)
|
||||
&& destination != Some(target.destination_id.as_bytes())
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
if contains_key_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE) {
|
||||
if info.transition_version_state == target.state
|
||||
&& info.transition_version == target.remote_version
|
||||
&& destination == Some(target.destination_id.as_bytes())
|
||||
{
|
||||
return Ok(false);
|
||||
}
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
if info.transition_version_state != TransitionVersionState::Unknown
|
||||
|| info
|
||||
.transition_version
|
||||
.as_deref()
|
||||
.filter(|value| !value.is_empty())
|
||||
.is_some_and(|value| Some(value) != target.remote_version.as_deref())
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
insert_bytes(
|
||||
&mut object.meta_sys,
|
||||
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
target.state.as_str().as_bytes().to_vec(),
|
||||
);
|
||||
remove_bytes(&mut object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID);
|
||||
if let Some(version) = &target.remote_version {
|
||||
insert_bytes(&mut object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, version.as_bytes().to_vec());
|
||||
}
|
||||
insert_bytes(
|
||||
&mut object.meta_sys,
|
||||
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
target.destination_id.as_bytes().to_vec(),
|
||||
);
|
||||
version.into_fileinfo("", "", true)?.validate_for_metadata_read()?;
|
||||
self.set_idx(index, version)?;
|
||||
Ok(true)
|
||||
}
|
||||
}
|
||||
|
||||
/// A stale healer or metadata writer may carry the original absent fields.
|
||||
/// Preserve a proven binding for the same immutable transition, or reject an
|
||||
/// attempted change of meaning. A new payload/version or a delete is separate.
|
||||
pub(super) fn preserve_reconciled_transition(previous: &FileMetaVersion, next: &mut FileMetaVersion) -> Result<()> {
|
||||
let (Some(previous_object), Some(next_object)) = (&previous.object, &mut next.object) else {
|
||||
return Ok(());
|
||||
};
|
||||
if !contains_key_bytes(&previous_object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID)
|
||||
|| !contains_key_bytes(&previous_object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE)
|
||||
|| previous_object.data_dir != next_object.data_dir
|
||||
{
|
||||
return Ok(());
|
||||
}
|
||||
let previous_info = previous.into_fileinfo("", "", true)?;
|
||||
if previous_info.transition_version_state == TransitionVersionState::Unknown
|
||||
|| previous_info.transition_status != TRANSITION_COMPLETE
|
||||
{
|
||||
return Ok(());
|
||||
}
|
||||
let next_info = next.into_fileinfo("", "", true)?;
|
||||
if previous_info.transition_tier != next_info.transition_tier
|
||||
|| previous_info.transitioned_objname != next_info.transitioned_objname
|
||||
|| previous_info.transition_status != next_info.transition_status
|
||||
|| previous_info.size != next_info.size
|
||||
|| previous_info.metadata.get("etag") != next_info.metadata.get("etag")
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
let destination =
|
||||
get_consistent_bytes(&previous_object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID).ok_or(Error::FileCorrupt)?;
|
||||
previous_info.validate_for_metadata_read()?;
|
||||
TransitionStateReconcileTarget {
|
||||
state: previous_info.transition_version_state,
|
||||
remote_version: previous_info.transition_version.clone(),
|
||||
destination_id: std::str::from_utf8(destination).map_err(|_| Error::FileCorrupt)?.to_string(),
|
||||
}
|
||||
.validate()?;
|
||||
let next_object = next.object.as_mut().ok_or(Error::FileCorrupt)?;
|
||||
if next_info
|
||||
.transition_version
|
||||
.as_ref()
|
||||
.is_some_and(|version| Some(version) != previous_info.transition_version.as_ref())
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
if contains_key_bytes(&next_object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE) {
|
||||
if previous_info.transition_version_state != next_info.transition_version_state
|
||||
|| previous_info.transition_version != next_info.transition_version
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
} else if next_info.transition_version_state != TransitionVersionState::Unknown {
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
if contains_key_bytes(&next_object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID)
|
||||
&& get_consistent_bytes(&next_object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID) != Some(destination)
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
for suffix in RECONCILE_SUFFIXES {
|
||||
remove_bytes(&mut next_object.meta_sys, suffix);
|
||||
if let Some(value) = get_consistent_bytes(&previous_object.meta_sys, suffix) {
|
||||
insert_bytes(&mut next_object.meta_sys, suffix, value.to_vec());
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::{ErasureInfo, FileInfo, ObjectPartInfo};
|
||||
|
||||
fn legacy() -> (FileMeta, FileInfo) {
|
||||
let info = FileInfo {
|
||||
version_id: Some(Uuid::from_u128(1)),
|
||||
data_dir: Some(Uuid::from_u128(2)),
|
||||
mod_time: Some(time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixture time")),
|
||||
size: 7,
|
||||
parts: vec![ObjectPartInfo {
|
||||
number: 1,
|
||||
size: 7,
|
||||
actual_size: 7,
|
||||
..Default::default()
|
||||
}],
|
||||
erasure: ErasureInfo {
|
||||
algorithm: "ReedSolomon".to_string(),
|
||||
data_blocks: 2,
|
||||
parity_blocks: 2,
|
||||
block_size: 1024 * 1024,
|
||||
index: 1,
|
||||
distribution: vec![1, 2, 3, 4],
|
||||
..Default::default()
|
||||
},
|
||||
transition_status: TRANSITION_COMPLETE.to_string(),
|
||||
transition_tier: "WARM".to_string(),
|
||||
transitioned_objname: "remote-object".to_string(),
|
||||
metadata: std::collections::HashMap::from([("etag".to_string(), "source-etag".to_string())]),
|
||||
data: Some(bytes::Bytes::from_static(b"payload")),
|
||||
..Default::default()
|
||||
};
|
||||
let mut metadata = FileMeta::new();
|
||||
metadata.add_version(info.clone()).expect("legacy fixture");
|
||||
(metadata, info)
|
||||
}
|
||||
|
||||
fn target(state: TransitionVersionState) -> TransitionStateReconcileTarget {
|
||||
TransitionStateReconcileTarget {
|
||||
state,
|
||||
remote_version: match state {
|
||||
TransitionVersionState::Exact => Some("opaque-version".to_string()),
|
||||
TransitionVersionState::SuspendedNull => Some("null".to_string()),
|
||||
_ => None,
|
||||
},
|
||||
destination_id: "ab".repeat(32),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transition_reconcile_preserves_payload_and_generation_and_is_idempotent() {
|
||||
for state in [
|
||||
TransitionVersionState::KnownDisabled,
|
||||
TransitionVersionState::SuspendedNull,
|
||||
TransitionVersionState::Exact,
|
||||
] {
|
||||
let (mut metadata, info) = legacy();
|
||||
let mut other = info.clone();
|
||||
other.version_id = Some(Uuid::from_u128(3));
|
||||
other.data_dir = Some(Uuid::from_u128(4));
|
||||
other.transition_status.clear();
|
||||
other.transition_tier.clear();
|
||||
other.transitioned_objname.clear();
|
||||
other.data = Some(bytes::Bytes::from_static(b"other!!"));
|
||||
metadata.add_version(other.clone()).expect("unrelated inline version");
|
||||
let other_before = metadata.find_version(other.version_id).expect("unrelated version").1;
|
||||
let original_data = metadata.data.clone();
|
||||
let generation = metadata
|
||||
.transition_reconcile_generation(info.version_id)
|
||||
.expect("initial generation");
|
||||
let target = target(state);
|
||||
assert!(metadata.reconcile_transition_state(info.version_id, &target).expect("repair"));
|
||||
let bytes = metadata.marshal_msg().expect("encode repair");
|
||||
let mut reloaded = FileMeta::load(&bytes).expect("reload repair");
|
||||
assert_eq!(reloaded.data, original_data);
|
||||
assert_eq!(
|
||||
reloaded
|
||||
.find_version(other.version_id)
|
||||
.expect("preserved unrelated version")
|
||||
.1,
|
||||
other_before
|
||||
);
|
||||
assert_eq!(
|
||||
reloaded
|
||||
.transition_reconcile_generation(info.version_id)
|
||||
.expect("repaired generation"),
|
||||
generation
|
||||
);
|
||||
assert!(!reloaded.reconcile_transition_state(info.version_id, &target).expect("retry"));
|
||||
let (_, version) = reloaded.find_version(info.version_id).expect("selected version");
|
||||
let repaired = version.into_fileinfo("", "", true).expect("decode explicit state");
|
||||
assert_eq!(repaired.transition_version_state, state);
|
||||
assert_eq!(repaired.transition_version, target.remote_version);
|
||||
assert_eq!(repaired.parts, info.parts);
|
||||
for prefix in [
|
||||
rustfs_utils::http::RUSTFS_INTERNAL_PREFIX,
|
||||
rustfs_utils::http::MINIO_INTERNAL_PREFIX,
|
||||
] {
|
||||
assert_eq!(
|
||||
version
|
||||
.object
|
||||
.as_ref()
|
||||
.expect("object")
|
||||
.meta_sys
|
||||
.get(&format!("{prefix}{SUFFIX_TRANSITIONED_VERSION_STATE}")),
|
||||
Some(&state.as_str().as_bytes().to_vec())
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transition_reconcile_generation_detects_unrelated_metadata_and_inline_changes() {
|
||||
let (mut metadata, info) = legacy();
|
||||
let original = metadata.transition_reconcile_generation(info.version_id).expect("generation");
|
||||
let mut updated = info.clone();
|
||||
updated.metadata.insert("user-tag".to_string(), "changed".to_string());
|
||||
metadata.update_object_version(updated).expect("update unrelated field");
|
||||
assert_ne!(metadata.transition_reconcile_generation(info.version_id).expect("generation"), original);
|
||||
let (mut metadata, mut info) = legacy();
|
||||
info.data = Some(bytes::Bytes::from_static(b"changed"));
|
||||
metadata.add_version(info.clone()).expect("change inline bytes");
|
||||
assert_ne!(metadata.transition_reconcile_generation(info.version_id).expect("generation"), original);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transition_reconcile_binding_survives_stale_heal_and_metadata_writes() {
|
||||
let (mut metadata, mut stale) = legacy();
|
||||
let target = target(TransitionVersionState::Exact);
|
||||
metadata
|
||||
.reconcile_transition_state(stale.version_id, &target)
|
||||
.expect("repair");
|
||||
metadata
|
||||
.add_version(stale.clone())
|
||||
.expect("stale heal must preserve the binding");
|
||||
stale.metadata.insert("user-tag".to_string(), "updated".to_string());
|
||||
metadata
|
||||
.update_object_version(stale.clone())
|
||||
.expect("ordinary metadata update");
|
||||
assert!(
|
||||
!metadata
|
||||
.reconcile_transition_state(stale.version_id, &target)
|
||||
.expect("binding remains exact")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transition_reconcile_rejects_explicit_unknown_and_conflicting_binding() {
|
||||
let (mut metadata, mut info) = legacy();
|
||||
rustfs_utils::http::insert_str(&mut info.metadata, SUFFIX_TRANSITIONED_VERSION_STATE, "unknown".to_string());
|
||||
metadata.add_version(info.clone()).expect("explicit unknown fixture");
|
||||
assert!(
|
||||
metadata
|
||||
.reconcile_transition_state(info.version_id, &target(TransitionVersionState::Exact))
|
||||
.is_err()
|
||||
);
|
||||
let (mut metadata, mut stale) = legacy();
|
||||
metadata
|
||||
.reconcile_transition_state(stale.version_id, &target(TransitionVersionState::Exact))
|
||||
.expect("repair");
|
||||
stale.transition_version_state = TransitionVersionState::KnownDisabled;
|
||||
assert!(
|
||||
metadata.add_version(stale).is_err(),
|
||||
"an explicit state cannot be replaced with a different model"
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -522,7 +522,7 @@ The following single-record protocol approves how historical objects without Rus
|
||||
|
||||
### Current
|
||||
|
||||
An absent `transitioned-version-state` key decodes as `TransitionVersionState::Unknown`. Non-destructive compatibility reads distinguish legacy absence from explicit `Unknown`; empty-version reads require a bounded probe. Legacy free-version cleanup requires a persisted exact remote version and does not infer unversioned deletion from an empty field. The single-record Admin routes below now inspect physical copies and live remote-version evidence. They do not yet persist repaired state: POST reports `write_fence_unavailable` until conditional per-generation metadata writes and the dedicated fleet capability are available. The existing transition-transaction reconcile route operates on expired `UploadOutcomeUnknown` transaction records and can exact-delete their canonical candidates; it is a separate protocol and must not be reused for metadata reconciliation.
|
||||
An absent `transitioned-version-state` key decodes as `TransitionVersionState::Unknown`. Non-destructive compatibility reads distinguish legacy absence from explicit `Unknown`; empty-version reads require a bounded probe. Legacy free-version cleanup requires a persisted exact remote version and does not infer unversioned deletion from an empty field. The single-record Admin routes below inspect physical copies and live remote-version evidence. POST persists the proven three-field repair through exact-copy conditional metadata writes when every fleet member advertises policy version 5. An unsupported or unknown member produces `write_fence_unavailable`. The existing transition-transaction reconcile route operates on expired `UploadOutcomeUnknown` transaction records and can exact-delete their canonical candidates; it is a separate protocol and must not be reused for metadata reconciliation.
|
||||
|
||||
An explicitly persisted `unknown`, a malformed state, conflicting RustFS/MinIO compatibility keys, an invalid or nil version identifier, and a partial transition tuple are not legacy absence. They remain invalid or ambiguous and fail closed.
|
||||
|
||||
@@ -560,13 +560,13 @@ The POST may write only the derived `transitioned-version-state`, its correspond
|
||||
|
||||
### Outcome contract
|
||||
|
||||
POST returns exactly one of the following outcomes and whether it changed bytes. GET uses the same diagnostic names for non-applicable cases, returns `ready-to-migrate` when a missing state is provable, and returns `migrated` only when strong readback shows the record was already explicit and converged:
|
||||
POST returns exactly one of the following outcomes. `changed` records confirmed writes; `changes_indeterminate` reports a failed write whose completion could not be established. After an indeterminate result, retry the same request or inspect again instead of assuming that no bytes changed. GET uses the same diagnostic names for non-applicable cases, returns `ready-to-migrate` when a missing state is provable, and returns `migrated` only when strong readback shows the record was already explicit and converged:
|
||||
|
||||
| Outcome | Meaning and permitted effect |
|
||||
|---|---|
|
||||
| `migrated` | All authoritative copies already contain, or were monotonically advanced to, the same proven state and destination identity. Only this outcome makes the record eligible for later ordinary read/delete semantics. |
|
||||
| `retained-ambiguous` | The tuple is structurally legacy-compatible, but the live probe is missing, multiple, changing, unsupported, or otherwise cannot prove exactly one state. No metadata or remote object is changed. |
|
||||
| `corrupt` | Explicit `Unknown`, malformed/contradictory compatibility keys, nil/invalid identifiers, partial transition metadata, or authoritative copies outside the one allowed `{original missing representation, exact proven target}` retry subset were observed. No backend probe is required after corruption is established, and nothing is changed. |
|
||||
| `corrupt` | Explicit `Unknown`, malformed/contradictory compatibility keys, nil/invalid identifiers, partial transition metadata, or authoritative copies outside the one allowed `{original missing representation, exact proven target}` retry subset were observed. No backend probe is required after corruption is established. The writer stops further changes and retains any subset committed before a later conflict was discovered. |
|
||||
| `backend-unavailable` | The bound tier generation/destination cannot be acquired, the bounded probe fails, or a metadata quorum/strong readback needed to complete the operation is unavailable. Any already-persisted monotonic subset is retained for an idempotent retry; it is never rolled back. |
|
||||
|
||||
HTTP failure detail may distinguish a stale expected tuple, lost fence, timeout, or unavailable quorum, but it must preserve one of these machine-readable outcomes. Logs and audit events include request identity, object identity, tier, generations, outcome, and whether bytes changed; they never include credentials or raw credential-derived configuration.
|
||||
@@ -588,11 +588,15 @@ The approved POST executes the following order. A step that cannot be proven sto
|
||||
|
||||
Cross-pool and cross-set partial success is monotonic. The only legal repair edge is `missing state -> one proven {state, remote version, destination identity}`. A retry may accept an already-written subset only when every known copy equals the newly proven target, every remaining copy equals its original missing-state representation captured by the reconciliation digest, and all immutable source fields still match; it then fills only the missing copies. This exact target-plus-original subset is neither stale nor corrupt. The retry never clears a known state, rewrites it to another state, changes destination identity, or rolls a successful set back to missing/`Unknown`. Any other divergent value produces `corrupt`; an unavailable set/readback produces `backend-unavailable`, and destructive cleanup remains blocked until a later strong all-pool read proves complete convergence.
|
||||
|
||||
Each copy carries a raw SHA-256 `metadata_digest` and an `unchanged_metadata_digest` over canonical decoded metadata, all other versions, and all inline bytes, excluding only the three repairable suffixes on the selected version. An initial write requires the raw digest. An already-target retry requires the unchanged-content digest and the exact state/version/destination tuple. GET therefore reads the complete `xl.meta`, including inline payloads belonging to other versions; the reply contains digests, not payload bytes.
|
||||
|
||||
Disk mutation serialization covers metadata updates, version writes, data publication, deletion, and rollback. Reconciliation refuses an outstanding rollback backup and bounds the object-directory backup check to 4,096 entries; it fails closed if that bound is exceeded. A stale rewrite of the same payload preserves the proven tuple or fails closed on a conflicting meaning. Each receiving disk acquires its own fleet/backend ownership, validates the conditional generation, and retains ownership through an atomic rename and mandatory directory fsync. A repair requires a Unix directory-sync implementation and metadata durability to be enabled for the bucket; unsupported platforms and disabled metadata sync block new repair writes instead of weakening the durability contract or overriding the configured policy. A zero-write validation pass enters the same disk mutation domain before final all-copy readback. Exact cleanup/recovery reads reject divergent transition bindings and cannot hide a minority owner behind majority absence.
|
||||
|
||||
GET takes the same fleet/topology snapshot and authoritative all-pool read but no write locks that imply mutation authority. Because GET is advisory, POST always repeats every fence, read, and live proof rather than promoting the GET result.
|
||||
|
||||
### Mixed-version and future batch work
|
||||
|
||||
The writer gate requires every node that can serve, rewrite, heal, decommission, or recover the affected `xl.meta` to preserve the explicit state and destination identity. A rolling fleet with an unknown/unsupported node is inspect-only. Downgrade is blocked while reconciled records could be rewritten by readers that erase or misinterpret those fields. Cross-pool movement must either copy the proven tuple unchanged or block reconciliation; a first-match lookup is never sufficient.
|
||||
The writer gate requires every node that can serve, rewrite, heal, decommission, or recover the affected `xl.meta` to preserve the explicit state and destination identity. A rolling fleet with an unknown/unsupported node is inspect-only. A capability downgrade blocks further reconciliation admission. Deployment controls must also prevent running an older metadata writer against repaired records; the capability probe cannot prevent an operator from replacing the executable. Cross-pool movement must either copy the proven tuple unchanged or block reconciliation; a first-match lookup is never sufficient.
|
||||
|
||||
Explicit `Unknown`, corruption, and ambiguity remain fail closed for reads that cannot prove non-destructive semantics and for every destructive path. A migrated record becomes ordinary explicit metadata, but reconciliation itself never transfers remote DELETE ownership.
|
||||
|
||||
|
||||
@@ -182,10 +182,10 @@ fn remove_heal_control_replay(
|
||||
static HEAL_CONTROL_REPLAY_CACHE: OnceLock<tokio::sync::Mutex<HashMap<String, Arc<HealControlReplayEntry>>>> = OnceLock::new();
|
||||
static NODE_CAPABILITY_SERVER_EPOCH: LazyLock<Uuid> = LazyLock::new(Uuid::new_v4);
|
||||
// v3 additionally promises the v6 tier-delete dispatch-manifest policy; v4
|
||||
// promises the sticky per-target decommission capacity fence. The
|
||||
// existing periodic topology probe carries both capabilities so normal object
|
||||
// operations do not add another peer RPC.
|
||||
const CROSS_POOL_FENCE_SUPPORTED_VERSION: u32 = 4;
|
||||
// promises the sticky per-target decommission capacity fence; v5 supports
|
||||
// conditional transition-state repair and preserves its destination binding.
|
||||
// Normal object operations reuse the periodic topology capability probe.
|
||||
const CROSS_POOL_FENCE_SUPPORTED_VERSION: u32 = 5;
|
||||
|
||||
fn encode_heal_capability_response(
|
||||
topology_member: &str,
|
||||
|
||||
Reference in New Issue
Block a user