diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 07eecc91b..65995b95a 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -38,6 +38,16 @@ pub mod bucket { } pub mod lifecycle { + pub mod legacy_transition_state_reconcile { + pub use crate::bucket::lifecycle::legacy_transition_state_reconcile::{ + LegacyTransitionStateCopyRepresentation, LegacyTransitionStateMetadataAlias, LegacyTransitionStateReconcileError, + LegacyTransitionStateReconcileOutcome, LegacyTransitionStateReconcileReadiness, + LegacyTransitionStateReconcileRequest, LegacyTransitionStateReconcileResponse, + LegacyTransitionStateReconcileSelector, LegacyTransitionStateSetRepresentation, LegacyTransitionStateSource, + LegacyTransitionStateTarget, + }; + } + pub mod bucket_lifecycle_audit { pub use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc; } diff --git a/crates/ecstore/src/bucket/lifecycle/legacy_transition_state_reconcile.rs b/crates/ecstore/src/bucket/lifecycle/legacy_transition_state_reconcile.rs new file mode 100644 index 000000000..e34cefa30 --- /dev/null +++ b/crates/ecstore/src/bucket/lifecycle/legacy_transition_state_reconcile.rs @@ -0,0 +1,1169 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! 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`. + +use std::collections::HashMap; +use std::time::Duration; + +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use uuid::Uuid; + +use crate::bucket::utils::check_bucket_and_object_names; +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, +}; +use crate::services::tier::tier::{TierConfigMgr, 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; + +use rustfs_filemeta::{FileInfo, FileMeta, TRANSITION_COMPLETE, TransitionVersionState}; +use rustfs_utils::http::metadata_compat::{ + SUFFIX_TRANSITION_TIER_DESTINATION_ID, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, + strip_internal_prefix_preserving_case, +}; + +const LIVE_PROBE_TIMEOUT: Duration = Duration::from_secs(10); + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateReconcileSelector { + pub bucket: String, + pub object: String, + pub version_id: String, +} + +impl LegacyTransitionStateReconcileSelector { + fn canonicalize(mut self) -> Result<(Self, Option), LegacyTransitionStateReconcileError> { + check_bucket_and_object_names(&self.bucket, &self.object) + .map_err(|err| LegacyTransitionStateReconcileError::InvalidSelector(err.to_string()))?; + if self.version_id.is_empty() { + return Err(LegacyTransitionStateReconcileError::InvalidSelector( + "versionId is required; use the literal null for an unversioned object".to_string(), + )); + } + if self.version_id == "null" { + return Ok((self, None)); + } + let version_id = Uuid::parse_str(&self.version_id).map_err(|_| { + LegacyTransitionStateReconcileError::InvalidSelector( + "versionId must be a non-nil UUID or the literal null".to_string(), + ) + })?; + if version_id.is_nil() { + return Err(LegacyTransitionStateReconcileError::InvalidSelector( + "a nil UUID is not a valid selector; use the literal null".to_string(), + )); + } + self.version_id = version_id.to_string(); + Ok((self, Some(version_id))) + } +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateMetadataAlias { + pub key: String, + /// Lowercase hexadecimal preserves empty and non-UTF-8 values exactly. + pub value_hex: String, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateSource { + pub bucket_incarnation: String, + /// Absent when the fleet has not supplied a current topology proof. + pub topology_generation: Option, + pub bucket: String, + pub object: String, + pub version_id: String, + pub data_dir: String, + pub modification_time_unix_nanos: i128, + pub size: i64, + pub etag: String, + pub transition_status: String, + pub tier: String, + pub remote_object: String, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateCopyRepresentation { + pub disk_index: usize, + pub metadata_digest: String, + pub state_aliases: Vec, + pub version_aliases: Vec, + pub destination_aliases: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateSetRepresentation { + pub pool_index: usize, + pub set_index: usize, + pub total_copies: usize, + pub available_copies: usize, + pub copies: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateTarget { + pub state: TransitionVersionState, + pub remote_version: Option, + pub destination_id: String, + pub tier_generation: u64, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateReconcileRequest { + pub confirm: bool, + pub selector: LegacyTransitionStateReconcileSelector, + pub source: LegacyTransitionStateSource, + pub original_sets: Vec, + pub target: LegacyTransitionStateTarget, + pub reconciliation_digest: String, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum LegacyTransitionStateReconcileOutcome { + ReadyToMigrate, + Migrated, + RetainedAmbiguous, + Corrupt, + BackendUnavailable, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateReconcileReadiness { + pub fleet_ready: bool, + pub topology_ready: bool, + pub tier_generation_ready: bool, + pub metadata_quorum_ready: bool, + pub post_ready: bool, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransitionStateReconcileResponse { + pub outcome: LegacyTransitionStateReconcileOutcome, + pub reason_code: String, + pub reason: String, + pub retryable: bool, + pub changed: bool, + pub selector: LegacyTransitionStateReconcileSelector, + pub source: Option, + pub original_sets: Vec, + pub target: Option, + pub reconciliation_digest: Option, + pub readiness: LegacyTransitionStateReconcileReadiness, +} + +#[derive(Debug, thiserror::Error)] +pub enum LegacyTransitionStateReconcileError { + #[error("invalid legacy transition state selector: {0}")] + InvalidSelector(String), + #[error("invalid legacy transition state reconciliation request: {0}")] + InvalidRequest(String), + #[error("legacy transition state reconciliation expected tuple is stale: {0}")] + StaleExpectedTuple(String), + #[error("legacy transition state metadata is corrupt: {0}")] + Corrupt(String), + #[error("legacy transition state backend is unavailable: {0}")] + BackendUnavailable(String), + #[error("legacy transition state write fence is unavailable: {0}")] + WriteFenceUnavailable(String), +} + +struct InspectedCopy { + file_info: FileInfo, + representation: LegacyTransitionStateCopyRepresentation, +} + +fn digest_hex(bytes: &[u8]) -> String { + rustfs_utils::crypto::hex(Sha256::digest(bytes)) +} + +fn metadata_aliases(metadata: &HashMap>, suffix: &str) -> Vec { + let mut aliases = metadata + .iter() + .filter(|(key, _)| internal_suffix_matches(key, suffix)) + .map(|(key, value)| LegacyTransitionStateMetadataAlias { + key: key.clone(), + value_hex: rustfs_utils::crypto::hex(value), + }) + .collect::>(); + aliases.sort_by(|left, right| left.key.cmp(&right.key)); + aliases +} + +fn internal_suffix_matches(key: &str, suffix: &str) -> bool { + strip_internal_prefix_preserving_case(key).is_some_and(|candidate| candidate.eq_ignore_ascii_case(suffix)) +} + +fn inspect_copy( + raw: &[u8], + disk_index: usize, + bucket: &str, + object: &str, + version_id: Option, +) -> Result, LegacyTransitionStateReconcileError> { + let metadata = FileMeta::load(raw) + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(format!("xl.meta decode failed: {err}")))?; + let (_, version) = match metadata.find_version(version_id) { + Ok(version) => version, + Err(rustfs_filemeta::Error::FileVersionNotFound) => return Ok(None), + Err(err) => return Err(LegacyTransitionStateReconcileError::Corrupt(err.to_string())), + }; + let object_metadata = version.object.as_ref().ok_or_else(|| { + LegacyTransitionStateReconcileError::Corrupt("the selected local version is not an object version".to_string()) + })?; + validate_remote_version_bytes(&object_metadata.meta_sys)?; + let file_info = version + .into_fileinfo(bucket, object, true) + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?; + file_info + .validate_for_metadata_read() + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?; + Ok(Some(InspectedCopy { + file_info, + representation: LegacyTransitionStateCopyRepresentation { + disk_index, + metadata_digest: digest_hex(raw), + 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), + }, + })) +} + +fn validate_remote_version_bytes(metadata: &HashMap>) -> Result<(), LegacyTransitionStateReconcileError> { + for (key, value) in metadata + .iter() + .filter(|(key, _)| internal_suffix_matches(key, SUFFIX_TRANSITIONED_VERSION_ID)) + { + if value.is_empty() { + continue; + } + if let Ok(version_id) = Uuid::from_slice(value) { + if !version_id.is_nil() { + continue; + } + return Err(LegacyTransitionStateReconcileError::Corrupt(format!( + "legacy remote version alias {key} contains a nil UUID" + ))); + } + let value = std::str::from_utf8(value).map_err(|_| { + LegacyTransitionStateReconcileError::Corrupt(format!( + "legacy remote version alias {key} is not valid UTF-8 or a raw UUID" + )) + })?; + if value.len() > 1024 + || value.chars().any(char::is_control) + || Uuid::parse_str(value).is_ok_and(|version_id| version_id.is_nil()) + { + return Err(LegacyTransitionStateReconcileError::Corrupt(format!( + "legacy remote version alias {key} contains an invalid identifier" + ))); + } + } + Ok(()) +} + +fn source_from_file_info( + selector: &LegacyTransitionStateReconcileSelector, + bucket_incarnation: Uuid, + file_info: &FileInfo, +) -> Result { + let data_dir = file_info.data_dir.ok_or_else(|| { + LegacyTransitionStateReconcileError::Corrupt("the selected transition source is missing its data directory".to_string()) + })?; + let modification_time = file_info.mod_time.ok_or_else(|| { + LegacyTransitionStateReconcileError::Corrupt( + "the selected transition source is missing its modification time".to_string(), + ) + })?; + let etag = crate::object_api::object_api_utils::get_raw_etag(&file_info.metadata); + if etag.is_empty() { + return Err(LegacyTransitionStateReconcileError::Corrupt( + "the selected transition source is missing its ETag".to_string(), + )); + } + Ok(LegacyTransitionStateSource { + bucket_incarnation: bucket_incarnation.to_string(), + topology_generation: None, + bucket: selector.bucket.clone(), + object: selector.object.clone(), + version_id: selector.version_id.clone(), + data_dir: data_dir.to_string(), + modification_time_unix_nanos: modification_time.unix_timestamp_nanos(), + size: file_info.size, + etag, + transition_status: file_info.transition_status.clone(), + tier: file_info.transition_tier.clone(), + remote_object: file_info.transitioned_objname.clone(), + }) +} + +fn response_digest( + source: &LegacyTransitionStateSource, + sets: &[LegacyTransitionStateSetRepresentation], + target: &LegacyTransitionStateTarget, +) -> Result { + let encoded = serde_json::to_vec(&(source, sets, target)) + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?; + Ok(digest_hex(&encoded)) +} + +impl ECStore { + pub async fn inspect_legacy_transition_state( + &self, + selector: LegacyTransitionStateReconcileSelector, + ) -> Result { + let (canonical_selector, _) = selector.clone().canonicalize()?; + // Keep the multi-pool inspection off the request stack and erase its + // concrete future so admin callers do not repeat its Send proof. + let inspection: futures::future::BoxFuture< + '_, + Result, + > = Box::pin(self.inspect_legacy_transition_state_inner(canonical_selector.clone())); + match inspection.await { + Ok(response) => Ok(response), + Err(err) => Ok(error_response(canonical_selector, err)), + } + } + + async fn inspect_legacy_transition_state_inner( + &self, + selector: LegacyTransitionStateReconcileSelector, + ) -> Result { + let (selector, local_version_id) = selector.canonicalize()?; + let remote_fleet_proof = acquire_remote_version_state_fleet_proof(); + let topology_proof = acquire_cross_pool_fence_fleet_proof(); + // 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_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( + "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; + let mut topology_ready = topology_proof.is_some(); + let mut source = None; + let mut original_sets = Vec::new(); + let mut canonical_file_info: Option = None; + let mut observed_file_infos = Vec::new(); + let mut observed_state_alias_presence = Vec::new(); + + for set in self.all_set_disks() { + let raw_copies = read_legacy_transition_state_metadata_copies(&set, &selector.bucket, &selector.object) + .await + .map_err(|err| match err { + crate::disk::error::Error::FileCorrupt => { + LegacyTransitionStateReconcileError::Corrupt("an xl.meta copy is corrupt".to_string()) + } + _ => LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()), + })?; + let total_copies = raw_copies.len(); + let available_copies = raw_copies.iter().filter(|copy| copy.is_some()).count(); + let mut representations = Vec::new(); + let mut set_file_infos = Vec::new(); + for (disk_index, raw) in raw_copies.into_iter().enumerate() { + let Some(raw) = raw else { continue }; + let Some(inspected) = inspect_copy(&raw, disk_index, &selector.bucket, &selector.object, local_version_id)? + else { + continue; + }; + if let Some(expected) = &canonical_file_info { + if !matches_immutable_transition_source(expected, &inspected.file_info) { + return Err(LegacyTransitionStateReconcileError::Corrupt( + "authoritative copies disagree on the immutable transition source tuple".to_string(), + )); + } + } else { + source = Some(source_from_file_info(&selector, bucket_incarnation, &inspected.file_info)?); + canonical_file_info = Some(inspected.file_info.clone()); + } + observed_state_alias_presence.push(!inspected.representation.state_aliases.is_empty()); + observed_file_infos.push(inspected.file_info.clone()); + set_file_infos.push(inspected.file_info.clone()); + representations.push(inspected.representation); + } + if !representations.is_empty() { + if representations.len() != available_copies { + return Err(LegacyTransitionStateReconcileError::BackendUnavailable( + "the selected version is missing from an existing metadata copy".to_string(), + )); + } + let required = required_reconcile_copy_quorum(&set_file_infos, set.default_write_quorum()); + if representations.len() < required { + return Err(LegacyTransitionStateReconcileError::BackendUnavailable(format!( + "pool {} set {} has only {} selected-version copies; {required} are required", + set.pool_index, + set.set_index, + representations.len() + ))); + } + original_sets.push(LegacyTransitionStateSetRepresentation { + pool_index: set.pool_index, + set_index: set.set_index, + total_copies, + available_copies, + copies: representations, + }); + } + } + if object_guards.iter().any(|guard| guard.is_lock_lost()) || bucket_guard.is_lock_lost() { + return Err(LegacyTransitionStateReconcileError::BackendUnavailable( + "a local metadata snapshot lock was lost before inspection completed".to_string(), + )); + } + drop(object_guards); + drop(bucket_guard); + + let Some(mut source) = source else { + return Err(LegacyTransitionStateReconcileError::Corrupt( + "the selected local object version was not found in any pool or set".to_string(), + )); + }; + source.topology_generation = topology_proof.as_ref().map(cross_pool_fence_topology_generation); + let file_info = canonical_file_info.ok_or_else(|| { + LegacyTransitionStateReconcileError::Corrupt( + "the selected source tuple disappeared while its metadata was being inspected".to_string(), + ) + })?; + if file_info.transition_status != TRANSITION_COMPLETE + || file_info.transition_tier.is_empty() + || file_info.transitioned_objname.is_empty() + { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::Corrupt, + reason_code: "partial_transition_tuple".to_string(), + reason: "the selected object does not contain a complete transition source tuple".to_string(), + retryable: false, + changed: false, + selector, + source: Some(source), + original_sets, + target: None, + reconciliation_digest: None, + readiness: readiness(fleet_ready, topology_ready, false, true, false), + }); + } + let has_explicit_unknown = + observed_file_infos + .iter() + .zip(&observed_state_alias_presence) + .any(|(observed, state_present)| { + observed.transition_version_state == TransitionVersionState::Unknown && *state_present + }); + if has_explicit_unknown { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::Corrupt, + reason_code: "explicit_unknown_state".to_string(), + reason: "an explicit unknown transition state is not legacy absence".to_string(), + retryable: false, + changed: false, + selector, + source: Some(source), + original_sets, + target: None, + reconciliation_digest: None, + readiness: readiness(fleet_ready, topology_ready, false, true, false), + }); + } + + let mut persisted_destination = None; + let mut persisted_remote_version: Option = None; + for observed in &observed_file_infos { + let destination = tier_destination_id_from_metadata(&observed.metadata) + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?; + if let Some(destination) = destination { + if persisted_destination.is_some_and(|expected| expected != destination) { + return Err(LegacyTransitionStateReconcileError::Corrupt( + "authoritative copies contain conflicting tier destination identities".to_string(), + )); + } + persisted_destination = Some(destination); + } + // A converged or partially migrated exact tuple still identifies + // its historical version, even when other remote versions exist. + if let Some(version) = observed.transition_version.as_deref() { + if persisted_remote_version + .as_deref() + .is_some_and(|expected| expected != version) + { + return Err(LegacyTransitionStateReconcileError::Corrupt( + "authoritative copies contain conflicting nonempty remote versions".to_string(), + )); + } + persisted_remote_version = Some(version.to_string()); + } + } + 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, + } + .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( + LIVE_PROBE_TIMEOUT, + lease.probe_legacy_transition_state(&file_info.transitioned_objname, persisted_remote_version.as_deref()), + ) + .await + .map_err(|_| LegacyTransitionStateReconcileError::BackendUnavailable("live tier probe timed out".to_string()))? + .map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?; + let remote_state_fleet_current = remote_fleet_proof + .as_ref() + .is_some_and(remote_version_state_fleet_proof_matches); + topology_ready = topology_proof.as_ref().is_some_and(cross_pool_fence_fleet_proof_matches) && remote_state_fleet_current; + if !lease.is_current_generation() + || self + .bucket_incarnation_id_from_disk(&selector.bucket) + .await + .map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))? + != bucket_incarnation + { + return Err(LegacyTransitionStateReconcileError::BackendUnavailable( + "the tier generation or bucket incarnation changed during the live probe".to_string(), + )); + } + let target = match probe { + LegacyTransitionStateProbe::UnversionedPresent => Some(LegacyTransitionStateTarget { + state: TransitionVersionState::KnownDisabled, + remote_version: None, + destination_id, + tier_generation, + }), + LegacyTransitionStateProbe::SuspendedNullPresent => Some(LegacyTransitionStateTarget { + state: TransitionVersionState::SuspendedNull, + remote_version: Some("null".to_string()), + destination_id, + tier_generation, + }), + LegacyTransitionStateProbe::VersionedPresent(version) => { + if version.is_empty() || version == "null" || Uuid::parse_str(&version).is_ok_and(|id| id.is_nil()) { + return Err(LegacyTransitionStateReconcileError::BackendUnavailable( + "the live tier probe returned an invalid exact version identifier".to_string(), + )); + } + rustfs_s3_client::provider_versions::validate_remote_version_id(&version) + .map_err(|err| LegacyTransitionStateReconcileError::BackendUnavailable(err.to_string()))?; + lease + .validate_remote_version_id(&version) + .map_err(|err| LegacyTransitionStateReconcileError::Corrupt(err.to_string()))?; + Some(LegacyTransitionStateTarget { + state: TransitionVersionState::Exact, + remote_version: Some(version), + destination_id, + tier_generation, + }) + } + LegacyTransitionStateProbe::Missing + | LegacyTransitionStateProbe::Ambiguous + | LegacyTransitionStateProbe::Unsupported => None, + }; + let Some(target) = target else { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::RetainedAmbiguous, + reason_code: "live_probe_ambiguous".to_string(), + reason: "the live backend probe did not prove exactly one remote version model".to_string(), + retryable: true, + changed: false, + selector, + source: Some(source), + original_sets, + target: None, + reconciliation_digest: None, + readiness: readiness(fleet_ready, topology_ready, true, true, false), + }); + }; + + let already_explicit = file_info.transition_version_state == target.state + && observed_file_infos.iter().all(|observed| { + observed.transition_version_state == target.state + && observed.transition_version == target.remote_version + && tier_destination_id_from_metadata(&observed.metadata).ok().flatten() == Some(lease.backend_identity()) + }); + let legacy_remote_versions_match = observed_file_infos + .iter() + .zip(&observed_state_alias_presence) + .filter(|(observed, state_present)| { + observed.transition_version_state == TransitionVersionState::Unknown && !**state_present + }) + .all(|(observed, _)| legacy_remote_version_matches_target(observed, &target)); + if !legacy_remote_versions_match + || persisted_remote_version + .as_deref() + .is_some_and(|persisted| target.remote_version.as_deref() != Some(persisted)) + { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::RetainedAmbiguous, + reason_code: "legacy_remote_version_changed".to_string(), + reason: "the live backend candidate does not match the persisted nonempty legacy remote version".to_string(), + retryable: true, + changed: false, + selector, + source: Some(source), + original_sets, + target: None, + reconciliation_digest: None, + readiness: readiness(fleet_ready, topology_ready, true, true, false), + }); + } + let allowed_retry_subset = + observed_file_infos + .iter() + .zip(&observed_state_alias_presence) + .all(|(observed, state_present)| { + let destination = tier_destination_id_from_metadata(&observed.metadata).ok(); + let explicit_target = observed.transition_version_state == target.state + && observed.transition_version == target.remote_version + && matches!(destination, Some(Some(value)) if value == lease.backend_identity()); + let missing_destination_matches = match destination { + Some(None) => true, + Some(Some(value)) => value == lease.backend_identity(), + None => false, + }; + let missing_state = observed.transition_version_state == TransitionVersionState::Unknown + && !*state_present + && missing_destination_matches; + explicit_target || missing_state + }); + if !allowed_retry_subset { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::Corrupt, + reason_code: "conflicting_retry_subset".to_string(), + reason: "authoritative copies are outside the allowed original-missing plus exact-target retry subset" + .to_string(), + retryable: false, + changed: false, + selector, + source: Some(source), + original_sets, + target: Some(target), + reconciliation_digest: None, + readiness: readiness(fleet_ready, topology_ready, true, true, false), + }); + } + let digest = response_digest(&source, &original_sets, &target)?; + let post_ready = !already_explicit && fleet_ready && topology_ready; + Ok(LegacyTransitionStateReconcileResponse { + outcome: if already_explicit { + LegacyTransitionStateReconcileOutcome::Migrated + } else { + LegacyTransitionStateReconcileOutcome::ReadyToMigrate + }, + reason_code: if already_explicit { + "already_converged".to_string() + } else if post_ready { + "ready_to_migrate".to_string() + } else { + "fleet_write_capability_unavailable".to_string() + }, + reason: if already_explicit { + "all inspected copies already contain the proven explicit transition state".to_string() + } else if post_ready { + "the live backend probe proved a single target tuple".to_string() + } else { + "the target tuple is proven, but the required fleet write capability is not available".to_string() + }, + retryable: !already_explicit, + changed: false, + selector, + source: Some(source), + original_sets, + target: Some(target), + reconciliation_digest: Some(digest), + readiness: readiness(fleet_ready, topology_ready, true, true, post_ready), + }) + } + + pub async fn reconcile_legacy_transition_state( + &self, + request: LegacyTransitionStateReconcileRequest, + ) -> Result { + if !request.confirm { + return Err(LegacyTransitionStateReconcileError::InvalidRequest("confirm must be true".to_string())); + } + let (selector, _) = request.selector.clone().canonicalize()?; + if selector != request.selector { + return Err(LegacyTransitionStateReconcileError::InvalidRequest( + "selector UUID must use its canonical representation".to_string(), + )); + } + let expected_digest = response_digest(&request.source, &request.original_sets, &request.target)?; + if expected_digest != request.reconciliation_digest { + return Ok(LegacyTransitionStateReconcileResponse { + outcome: LegacyTransitionStateReconcileOutcome::Corrupt, + reason_code: "stale_expected_tuple".to_string(), + reason: "reconciliation digest does not match the supplied source, sets, and target".to_string(), + retryable: false, + 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: unavailable_write_readiness(), + }); + } + let current = self.inspect_legacy_transition_state(request.selector.clone()).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, + }); + } + if current.outcome == LegacyTransitionStateReconcileOutcome::Migrated { + return Ok(current); + } + 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, + }) + } +} + +fn legacy_remote_version_matches_target(file_info: &FileInfo, target: &LegacyTransitionStateTarget) -> bool { + match (&file_info.transition_version, &target.remote_version) { + (None, _) => true, + (Some(persisted), Some(proven)) => persisted == proven, + (Some(_), None) => false, + } +} + +fn unavailable_write_readiness() -> LegacyTransitionStateReconcileReadiness { + LegacyTransitionStateReconcileReadiness { + fleet_ready: false, + topology_ready: false, + tier_generation_ready: false, + metadata_quorum_ready: false, + post_ready: false, + } +} + +fn readiness( + fleet_ready: bool, + topology_ready: bool, + tier_generation_ready: bool, + metadata_quorum_ready: bool, + post_ready: bool, +) -> LegacyTransitionStateReconcileReadiness { + LegacyTransitionStateReconcileReadiness { + fleet_ready, + topology_ready, + tier_generation_ready, + metadata_quorum_ready, + post_ready, + } +} + +fn error_response( + selector: LegacyTransitionStateReconcileSelector, + err: LegacyTransitionStateReconcileError, +) -> LegacyTransitionStateReconcileResponse { + let (outcome, reason_code, retryable) = match &err { + LegacyTransitionStateReconcileError::Corrupt(_) | LegacyTransitionStateReconcileError::StaleExpectedTuple(_) => { + (LegacyTransitionStateReconcileOutcome::Corrupt, "corrupt", false) + } + LegacyTransitionStateReconcileError::BackendUnavailable(_) => { + (LegacyTransitionStateReconcileOutcome::BackendUnavailable, "backend_unavailable", true) + } + LegacyTransitionStateReconcileError::WriteFenceUnavailable(_) => { + (LegacyTransitionStateReconcileOutcome::BackendUnavailable, "write_fence_unavailable", true) + } + LegacyTransitionStateReconcileError::InvalidSelector(_) | LegacyTransitionStateReconcileError::InvalidRequest(_) => { + (LegacyTransitionStateReconcileOutcome::Corrupt, "invalid_request", false) + } + }; + LegacyTransitionStateReconcileResponse { + outcome, + reason_code: reason_code.to_string(), + reason: err.to_string(), + retryable, + changed: false, + selector, + source: None, + original_sets: Vec::new(), + target: None, + reconciliation_digest: None, + readiness: unavailable_write_readiness(), + } +} + +fn required_reconcile_copy_quorum(file_infos: &[FileInfo], default_write_quorum: usize) -> usize { + file_infos + .iter() + .map(|info| info.write_quorum(default_write_quorum)) + .max() + .unwrap_or(default_write_quorum) + .max(default_write_quorum) +} + +fn matches_immutable_transition_source(expected: &FileInfo, observed: &FileInfo) -> bool { + expected.version_id == observed.version_id + && expected.data_dir == observed.data_dir + && expected.mod_time == observed.mod_time + && expected.size == observed.size + && crate::object_api::object_api_utils::get_raw_etag(&expected.metadata) + == crate::object_api::object_api_utils::get_raw_etag(&observed.metadata) + && expected.transition_status == observed.transition_status + && expected.transition_tier == observed.transition_tier + && expected.transitioned_objname == observed.transitioned_objname +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn pinned_minio_and_legacy_rustfs_envelopes_are_inspected_without_reencoding() { + for bytes in [ + rustfs_filemeta::test_data::create_minio_small_object_xlmeta().expect("pinned MinIO inline object"), + rustfs_filemeta::test_data::create_minio_large_object_xlmeta().expect("pinned MinIO external data object"), + rustfs_filemeta::test_data::create_minio_versioned_object_xlmeta().expect("pinned MinIO versioned object"), + rustfs_filemeta::test_data::create_issue_2265_legacy_meta_v2_object_xlmeta().expect("pinned legacy RustFS object"), + ] { + let original = bytes.clone(); + let metadata = FileMeta::load(&bytes).expect("load pinned metadata envelope"); + for version in &metadata.versions { + let inspected = inspect_copy(&bytes, 0, "bucket", "object", version.header.version_id); + if version.header.version_type == rustfs_filemeta::VersionType::Delete { + assert!( + matches!(inspected, Err(LegacyTransitionStateReconcileError::Corrupt(_))), + "a pinned delete marker is not a transition-state repair candidate" + ); + continue; + } + let inspected = inspected.expect("inspect pinned version").expect("pinned version is present"); + assert_eq!(inspected.file_info.transition_version_state, TransitionVersionState::Unknown); + assert!(inspected.representation.state_aliases.is_empty()); + assert_eq!(inspected.representation.metadata_digest, digest_hex(&original)); + } + assert_eq!(bytes, original, "foreign metadata must never be round-tripped by inspection"); + } + } + + fn legacy_metadata_fixture(version_bytes: Option<&[u8]>, state: Option<&[u8]>) -> Vec { + let mut metadata = FileMeta::new(); + metadata + .add_version(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 modification time")), + size: 7, + parts: vec![rustfs_filemeta::ObjectPartInfo { + number: 1, + size: 7, + actual_size: 7, + ..Default::default() + }], + erasure: rustfs_filemeta::ErasureInfo { + algorithm: "ReedSolomon".to_string(), + data_blocks: 2, + parity_blocks: 2, + block_size: 1_048_576, + index: 1, + distribution: vec![1, 2, 3, 4], + ..Default::default() + }, + metadata: HashMap::from([("etag".to_string(), "legacy-etag".to_string())]), + ..Default::default() + }) + .expect("encode source object fixture"); + let (index, mut version) = metadata.find_version(Some(Uuid::from_u128(1))).expect("fixture version"); + let object = version.object.as_mut().expect("fixture object metadata"); + for (suffix, value) in [ + ("transition-status", b"complete".as_slice()), + ("transition-tier", b"WARM".as_slice()), + ("transitioned-object", b"archive/object".as_slice()), + ] { + object.meta_sys.insert(format!("x-minio-internal-{suffix}"), value.to_vec()); + } + if let Some(value) = version_bytes { + object + .meta_sys + .insert(format!("x-minio-internal-{SUFFIX_TRANSITIONED_VERSION_ID}"), value.to_vec()); + } + if let Some(value) = state { + object + .meta_sys + .insert(format!("x-minio-internal-{SUFFIX_TRANSITIONED_VERSION_STATE}"), value.to_vec()); + } + metadata.versions[index].header = version.header(); + metadata.versions[index].meta = version.marshal_msg().expect("encode legacy object version"); + metadata.marshal_msg().expect("encode legacy xl.meta") + } + + #[test] + fn raw_minio_and_rustfs_fixtures_preserve_legacy_provenance_without_backfill() { + for raw_version in [ + None, + Some(b"".as_slice()), + Some(b"opaque-version".as_slice()), + Some(Uuid::from_u128(9).as_bytes().as_slice()), + ] { + let bytes = legacy_metadata_fixture(raw_version, None); + let before = bytes.clone(); + let inspected = inspect_copy(&bytes, 0, "bucket", "object", Some(Uuid::from_u128(1))) + .expect("legacy metadata should decode") + .expect("legacy version should exist"); + assert_eq!(inspected.file_info.transition_version_state, TransitionVersionState::Unknown); + assert!(inspected.representation.state_aliases.is_empty()); + assert_eq!(inspected.representation.version_aliases.len(), usize::from(raw_version.is_some())); + assert_eq!(inspected.representation.metadata_digest, digest_hex(&before)); + assert_eq!(bytes, before, "inspection must not rewrite legacy bytes"); + } + } + + #[test] + fn corrupt_legacy_bytes_cannot_become_an_unversioned_candidate() { + for raw_version in [ + b"bad\nversion".as_slice(), + b"\xff\xfe".as_slice(), + Uuid::nil().as_bytes().as_slice(), + ] { + let bytes = legacy_metadata_fixture(Some(raw_version), None); + assert!(matches!( + inspect_copy(&bytes, 0, "bucket", "object", Some(Uuid::from_u128(1))), + Err(LegacyTransitionStateReconcileError::Corrupt(_)) + )); + } + let bytes = legacy_metadata_fixture(Some(b""), Some(b"unknown")); + let inspected = inspect_copy(&bytes, 0, "bucket", "object", Some(Uuid::from_u128(1))) + .expect("explicit unknown is structurally decodable") + .expect("fixture version"); + assert_eq!(inspected.file_info.transition_version_state, TransitionVersionState::Unknown); + assert!( + !inspected.representation.state_aliases.is_empty(), + "explicit unknown must remain distinct from legacy absence" + ); + } + + #[test] + fn selector_requires_explicit_non_nil_local_version_identity() { + for version_id in ["", "00000000-0000-0000-0000-000000000000", "not-a-version"] { + let selector = LegacyTransitionStateReconcileSelector { + bucket: "bucket".to_string(), + object: "object".to_string(), + version_id: version_id.to_string(), + }; + assert!(matches!( + selector.canonicalize(), + Err(LegacyTransitionStateReconcileError::InvalidSelector(_)) + )); + } + assert!( + LegacyTransitionStateReconcileSelector { + bucket: "bucket".to_string(), + object: "object".to_string(), + version_id: "null".to_string(), + } + .canonicalize() + .is_ok() + ); + } + + #[test] + fn digest_binds_source_sets_and_target() { + let source = LegacyTransitionStateSource { + bucket_incarnation: Uuid::from_u128(1).to_string(), + topology_generation: Some("topology-a".to_string()), + bucket: "bucket".to_string(), + object: "object".to_string(), + version_id: "null".to_string(), + data_dir: Uuid::from_u128(2).to_string(), + modification_time_unix_nanos: 1, + size: 1, + etag: "etag".to_string(), + transition_status: TRANSITION_COMPLETE.to_string(), + tier: "WARM".to_string(), + remote_object: "remote/object".to_string(), + }; + let sets = vec![LegacyTransitionStateSetRepresentation { + pool_index: 0, + set_index: 0, + total_copies: 1, + available_copies: 1, + copies: vec![], + }]; + let target = LegacyTransitionStateTarget { + state: TransitionVersionState::KnownDisabled, + remote_version: None, + destination_id: "00".repeat(32), + tier_generation: 1, + }; + let digest = response_digest(&source, &sets, &target).expect("valid source digest"); + let mut changed = source.clone(); + changed.remote_object.push_str("-changed"); + assert_ne!(digest, response_digest(&changed, &sets, &target).expect("changed source digest")); + } + + #[test] + fn nonempty_legacy_remote_version_must_match_live_candidate() { + let mut file_info = FileInfo { + transition_version: Some("remote-version-a".to_string()), + ..Default::default() + }; + let mut target = LegacyTransitionStateTarget { + state: TransitionVersionState::Exact, + remote_version: Some("remote-version-b".to_string()), + destination_id: "00".repeat(32), + tier_generation: 1, + }; + assert!(!legacy_remote_version_matches_target(&file_info, &target)); + target.remote_version = Some("remote-version-a".to_string()); + assert!(legacy_remote_version_matches_target(&file_info, &target)); + file_info.transition_version = None; + assert!(legacy_remote_version_matches_target(&file_info, &target)); + } + + #[test] + fn missing_state_accepts_only_empty_utf8_or_non_nil_raw_uuid_version_values() { + let key = format!("x-minio-internal-{SUFFIX_TRANSITIONED_VERSION_ID}"); + let mut metadata = HashMap::from([(key.clone(), Vec::new())]); + validate_remote_version_bytes(&metadata).expect("empty MinIO value is preserved legacy provenance"); + + metadata.insert(key.clone(), Uuid::from_u128(1).as_bytes().to_vec()); + validate_remote_version_bytes(&metadata).expect("non-nil historical RustFS UUID is valid provenance"); + + metadata.insert(key.clone(), Uuid::nil().as_bytes().to_vec()); + assert!(matches!( + validate_remote_version_bytes(&metadata), + Err(LegacyTransitionStateReconcileError::Corrupt(_)) + )); + + metadata.insert(key, vec![0xff, 0xfe]); + assert!(matches!( + validate_remote_version_bytes(&metadata), + Err(LegacyTransitionStateReconcileError::Corrupt(_)) + )); + } + + #[test] + fn raw_copy_quorum_cannot_be_lowered_by_embedded_erasure_geometry() { + let undersized = FileInfo { + erasure: rustfs_filemeta::ErasureInfo { + data_blocks: 1, + parity_blocks: 0, + ..Default::default() + }, + ..Default::default() + }; + let balanced = FileInfo { + erasure: rustfs_filemeta::ErasureInfo { + data_blocks: 4, + parity_blocks: 4, + ..Default::default() + }, + ..Default::default() + }; + + assert_eq!(required_reconcile_copy_quorum(&[undersized], 5), 5); + assert_eq!(required_reconcile_copy_quorum(&[balanced], 4), 5); + assert_eq!(required_reconcile_copy_quorum(&[], 5), 5); + + let old_pool = FileInfo { + erasure: rustfs_filemeta::ErasureInfo { + data_blocks: 8, + parity_blocks: 4, + ..Default::default() + }, + ..Default::default() + }; + let new_pool = FileInfo { + erasure: rustfs_filemeta::ErasureInfo { + data_blocks: 2, + parity_blocks: 2, + ..Default::default() + }, + ..Default::default() + }; + assert_eq!(required_reconcile_copy_quorum(&[old_pool], 7), 8); + assert_eq!(required_reconcile_copy_quorum(&[new_pool], 3), 3); + } +} diff --git a/crates/ecstore/src/bucket/lifecycle/mod.rs b/crates/ecstore/src/bucket/lifecycle/mod.rs index de64dbbcf..ea56b5f8e 100644 --- a/crates/ecstore/src/bucket/lifecycle/mod.rs +++ b/crates/ecstore/src/bucket/lifecycle/mod.rs @@ -18,6 +18,7 @@ mod config_boundary; pub mod core; mod durable_namespace; pub mod evaluator; +pub mod legacy_transition_state_reconcile; pub mod manual_transition_job; mod metadata_boundary; pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, get_lifecycle_config}; diff --git a/crates/ecstore/src/services/notification_sys.rs b/crates/ecstore/src/services/notification_sys.rs index ee6bc7dd7..ce207bd5c 100644 --- a/crates/ecstore/src/services/notification_sys.rs +++ b/crates/ecstore/src/services/notification_sys.rs @@ -568,6 +568,10 @@ pub(crate) fn tier_delete_journal_topology_generation(proof: &TierDeleteJournalF stable_tier_delete_journal_topology_generation(&proof.token.topology_fingerprint) } +pub(crate) fn cross_pool_fence_topology_generation(proof: &CrossPoolFenceFleetProofToken) -> String { + stable_tier_delete_journal_topology_generation(&proof.0.topology_fingerprint) +} + /// Acquire one non-cloneable authority that must span the complete reconcile /// effect window, including its final strong readback. pub async fn acquire_legacy_transition_state_reconcile_fleet_proof() -> Option { @@ -1120,6 +1124,14 @@ pub(crate) fn install_remote_version_state_fleet_proof_for_test(topology_fingerp RemoteVersionStateFleetProofGuard } +#[cfg(all(test, feature = "test-util"))] +pub(crate) fn install_current_remote_version_state_fleet_proof_for_test() -> RemoteVersionStateFleetProofGuard { + let topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY + .get() + .expect("the test store must bind its fleet topology before installing a writer proof"); + install_remote_version_state_fleet_proof_for_test(topology) +} + #[cfg(all(test, feature = "test-util"))] pub(crate) struct TransitionTransactionCompactionFleetProofGuard; diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index 09c6dee81..96e9b92fc 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -653,6 +653,26 @@ impl MockWarmBackend { #[async_trait] impl WarmBackend for MockWarmBackend { + async fn probe_legacy_metadata( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + use super::warm_backend::LegacyTransitionStateProbe as Probe; + let candidate = match remote_version { + Some(version) if !version.is_empty() => self.probe_transition_version(object, version).await?, + _ => self.probe_transition_candidate(object).await?, + }; + Ok(match candidate { + TransitionCandidateProbe::Missing => Probe::Missing, + TransitionCandidateProbe::UnversionedPresent => Probe::UnversionedPresent, + TransitionCandidateProbe::VersionedPresent(version) if version == "null" => Probe::SuspendedNullPresent, + TransitionCandidateProbe::VersionedPresent(version) => Probe::VersionedPresent(version), + TransitionCandidateProbe::Ambiguous => Probe::Ambiguous, + TransitionCandidateProbe::Unsupported => Probe::Unsupported, + }) + } + fn validate_remote_version_id(&self, remote_version_id: &str) -> Result<(), std::io::Error> { if remote_version_id.is_empty() { return Ok(()); @@ -874,8 +894,9 @@ pub async fn register_mock_tier_backend(handle: &Arc>, tie ..Default::default() }, ); - tier_config_mgr - .install_test_driver(tier_name, Box::new(backend)) + drop(tier_config_mgr); + TierConfigMgr::install_test_driver_in(handle, tier_name, Box::new(backend)) + .await .expect("mock tier driver should install"); } diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 5e67673e9..1d5adf39b 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -2322,6 +2322,14 @@ struct SharedWarmBackendProxy(SharedWarmBackend); #[async_trait::async_trait] impl WarmBackend for SharedWarmBackendProxy { + async fn probe_legacy_metadata( + &self, + object: &str, + remote_version: Option<&str>, + ) -> io::Result { + self.0.probe_legacy_metadata(object, remote_version).await + } + async fn validate(&self) -> io::Result<()> { self.0.validate().await } @@ -2490,6 +2498,27 @@ impl TierOperationLease { self.inner.driver.probe_transition_version(object, remote_version_id).await } + pub(crate) async fn probe_legacy_transition_state( + &self, + object: &str, + remote_version: Option<&str>, + ) -> io::Result { + let Some(reconciler) = self + .inner + .reconciler + .get_or_try_init(|| async { + crate::services::tier::warm_backend::new_transition_candidate_reconciler(&self.inner.tier_config) + .await + .map(|reconciler| reconciler.map(Arc::from)) + }) + .await + .map_err(|err| io::Error::other(err.message))? + else { + return self.inner.driver.probe_legacy_metadata(object, remote_version).await; + }; + reconciler.probe_legacy_transition_state(object, remote_version).await + } + pub(crate) fn is_current_generation(&self) -> bool { lock_unpoisoned(&self.runtime) .generations @@ -6013,6 +6042,19 @@ impl TierConfigMgr { Ok(()) } + #[cfg(any(test, feature = "test-util"))] + pub(crate) async fn install_test_driver_in( + handle: &Arc>, + tier_name: &str, + driver: WarmBackendImpl, + ) -> std::result::Result<(), AdminError> { + let mut manager = handle.write().await; + // Register the generation runtime before installing the mock so its + // explicit lack of a network reconciler survives the first lease. + tier_driver_runtime(handle, &manager); + manager.install_test_driver(tier_name, driver) + } + #[cfg(any(test, feature = "test-util"))] pub(crate) fn install_test_driver( &mut self, diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index fe4b5e118..8e7a4c8f4 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -89,6 +89,18 @@ pub enum TransitionCandidateProbe { Unsupported, } +/// Live evidence for repairing legacy metadata. Ordinary candidate GETs do not +/// establish the bucket's versioning model and cannot supply this authority. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum LegacyTransitionStateProbe { + Missing, + UnversionedPresent, + SuspendedNullPresent, + VersionedPresent(String), + Ambiguous, + Unsupported, +} + #[derive(Clone, Copy)] pub(crate) struct TransitionCandidateIdentity { pub transaction_id: uuid::Uuid, @@ -97,6 +109,14 @@ pub(crate) struct TransitionCandidateIdentity { #[async_trait::async_trait] pub(crate) trait TransitionCandidateReconciler { + async fn probe_legacy_transition_state( + &self, + _object: &str, + _remote_version: Option<&str>, + ) -> Result { + Ok(LegacyTransitionStateProbe::Unsupported) + } + async fn probe_transition_candidate_for( &self, object: &str, @@ -106,6 +126,14 @@ pub(crate) trait TransitionCandidateReconciler { #[async_trait::async_trait] pub trait WarmBackend { + async fn probe_legacy_metadata( + &self, + _object: &str, + _remote_version: Option<&str>, + ) -> Result { + Ok(LegacyTransitionStateProbe::Unsupported) + } + async fn validate(&self) -> Result<(), std::io::Error> { Ok(()) } @@ -448,6 +476,18 @@ impl MeteredWarmBackend { #[async_trait::async_trait] impl WarmBackend for MeteredWarmBackend { + async fn probe_legacy_metadata( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + let result = self.inner.probe_legacy_metadata(object, remote_version).await; + if matches!(result, Ok(LegacyTransitionStateProbe::Unsupported)) { + return result; + } + Self::record(TierRequestOperation::Probe, result) + } + /// Delegated without a counter: only one backend issues a remote request /// here, and every other one takes the trait default, so a `validate` /// counter would mostly record requests that never happened. @@ -524,6 +564,18 @@ struct MeteredTransitionCandidateReconciler { #[async_trait::async_trait] impl TransitionCandidateReconciler for MeteredTransitionCandidateReconciler { + async fn probe_legacy_transition_state( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + let result = self.inner.probe_legacy_transition_state(object, remote_version).await; + if matches!(result, Ok(LegacyTransitionStateProbe::Unsupported)) { + return result; + } + MeteredWarmBackend::record(TierRequestOperation::Probe, result) + } + async fn probe_transition_candidate_for( &self, object: &str, diff --git a/crates/ecstore/src/services/tier/warm_backend_minio.rs b/crates/ecstore/src/services/tier/warm_backend_minio.rs index 6baef149b..34657ef9a 100644 --- a/crates/ecstore/src/services/tier/warm_backend_minio.rs +++ b/crates/ecstore/src/services/tier/warm_backend_minio.rs @@ -101,6 +101,19 @@ impl WarmBackend for WarmBackendMinIO { #[async_trait::async_trait] impl crate::services::tier::warm_backend::TransitionCandidateReconciler for WarmBackendMinIO { + async fn probe_legacy_transition_state( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_legacy_transition_state( + &self.0, + object, + remote_version, + ) + .await + } + async fn probe_transition_candidate_for( &self, object: &str, diff --git a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs index 0bb19bcdf..9c1934a6d 100644 --- a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs +++ b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs @@ -146,6 +146,19 @@ impl WarmBackend for WarmBackendRustFS { #[async_trait::async_trait] impl crate::services::tier::warm_backend::TransitionCandidateReconciler for WarmBackendRustFS { + async fn probe_legacy_transition_state( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_legacy_transition_state( + &self.0, + object, + remote_version, + ) + .await + } + async fn probe_transition_candidate_for( &self, object: &str, diff --git a/crates/ecstore/src/services/tier/warm_backend_s3.rs b/crates/ecstore/src/services/tier/warm_backend_s3.rs index 8c1d6ef9b..4e790b560 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -376,7 +376,6 @@ struct TransitionCandidateVersions { } impl TransitionCandidateVersions { - #[cfg(test)] fn extend(&mut self, remote_object: &str, versions: &ListVersionsResult) { for version in versions.versions.iter().filter(|version| version.key == remote_object) { if self.version_id.is_some() { @@ -512,6 +511,20 @@ mod tests { } async fn candidate_probe_fixture() -> Option<(WarmBackendS3, tokio::task::JoinHandle>)> { + scripted_probe_fixture([ + "HTTP/1.1 206 Partial Content\r\nContent-Length: 1\r\nx-amz-version-id: opaque-version\r\nConnection: close\r\n\r\nx", + "HTTP/1.1 206 Partial Content\r\nContent-Length: 1\r\nConnection: close\r\n\r\nx", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 66\r\nConnection: close\r\n\r\nNoSuchObjectmissing", + "HTTP/1.1 403 Forbidden\r\nContent-Type: application/xml\r\nContent-Length: 65\r\nConnection: close\r\n\r\nAccessDenieddenied", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", + "HTTP/1.1 416 Range Not Satisfiable\r\nContent-Type: application/xml\r\nContent-Length: 72\r\nConnection: close\r\n\r\nInvalidRangeempty version", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 67\r\nConnection: close\r\n\r\nNoSuchVersionmissing", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", + ].into_iter().map(str::to_owned).collect()).await + } + + async fn scripted_probe_fixture(responses: Vec) -> Option<(WarmBackendS3, tokio::task::JoinHandle>)> { let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await { Ok(listener) => listener, Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, @@ -522,17 +535,6 @@ mod tests { .expect("listener local address should be available") .to_string(); let fixture = tokio::spawn(async move { - let responses = [ - "HTTP/1.1 206 Partial Content\r\nContent-Length: 1\r\nx-amz-version-id: opaque-version\r\nConnection: close\r\n\r\nx", - "HTTP/1.1 206 Partial Content\r\nContent-Length: 1\r\nConnection: close\r\n\r\nx", - "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", - "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 66\r\nConnection: close\r\n\r\nNoSuchObjectmissing", - "HTTP/1.1 403 Forbidden\r\nContent-Type: application/xml\r\nContent-Length: 65\r\nConnection: close\r\n\r\nAccessDenieddenied", - "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", - "HTTP/1.1 416 Range Not Satisfiable\r\nContent-Type: application/xml\r\nContent-Length: 72\r\nConnection: close\r\n\r\nInvalidRangeempty version", - "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 67\r\nConnection: close\r\n\r\nNoSuchVersionmissing", - "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", - ]; let mut requests = Vec::new(); for response in responses { let (mut stream, _) = listener.accept().await.expect("fixture should accept candidate GET"); @@ -673,6 +675,117 @@ mod tests { assert!(requests[8].to_ascii_lowercase().contains("?versionid=historical-version")); } + fn legacy_probe_xml_response(body: &str) -> String { + format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/xml\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", + body.len() + ) + } + + fn legacy_probe_versioning_response(status: &str) -> String { + let state = if status.is_empty() { + String::new() + } else { + format!("{status}") + }; + legacy_probe_xml_response(&format!( + "{state}" + )) + } + + fn legacy_probe_versions_response(versions: &[&str]) -> String { + let versions = versions.iter().map(|version| format!( + "archive/object{version}true2026-09-01T00:00:00Z\"legacy-etag\"7STANDARD" + )).collect::(); + legacy_probe_xml_response(&format!( + "bucketarchive/object1000false{versions}" + )) + } + + #[tokio::test] + async fn legacy_transition_state_probe_verifies_disabled_suspended_and_enabled_responses() { + use super::super::warm_backend::LegacyTransitionStateProbe as Probe; + for (initial, confirmed, version, expected) in [ + ("", "", "null", Probe::UnversionedPresent), + ("Suspended", "Suspended", "null", Probe::SuspendedNullPresent), + ("Enabled", "Enabled", "version-a", Probe::VersionedPresent("version-a".to_string())), + ("Enabled", "Enabled", "null", Probe::SuspendedNullPresent), + ("", "", "unexpected-version", Probe::Ambiguous), + ("Suspended", "Enabled", "null", Probe::Ambiguous), + ] { + let responses = vec![ + legacy_probe_versioning_response(initial), + legacy_probe_versions_response(&[version]), + legacy_probe_versioning_response(confirmed), + ]; + let (backend, fixture) = scripted_probe_fixture(responses) + .await + .expect("legacy probe loopback fixture"); + let result = + tokio::time::timeout(Duration::from_secs(10), backend.probe_legacy_transition_state("archive/object", None)) + .await + .expect("legacy probe must finish") + .expect("legacy probe should decode provider XML"); + assert_eq!(result, expected, "initial={initial} confirmed={confirmed} version={version}"); + let requests = fixture.await.expect("legacy probe fixture should finish"); + assert_eq!(requests.len(), 3); + assert!(requests.iter().all(|request| request.starts_with("GET "))); + assert!(requests[0].lines().next().expect("request line").contains("versioning")); + assert!(requests[1].lines().next().expect("request line").contains("versions")); + assert!(requests[2].lines().next().expect("request line").contains("versioning")); + } + } + + #[tokio::test] + async fn legacy_transition_state_probe_preserves_historical_exact_version() { + use super::super::warm_backend::LegacyTransitionStateProbe as Probe; + let responses = vec![ + legacy_probe_versioning_response("Enabled"), + "HTTP/1.1 206 Partial Content\r\nContent-Length: 1\r\nx-amz-version-id: historical-version\r\nConnection: close\r\n\r\nx".to_string(), + legacy_probe_versioning_response("Enabled"), + ]; + let (backend, fixture) = scripted_probe_fixture(responses).await.expect("exact legacy probe fixture"); + assert_eq!( + backend + .probe_legacy_transition_state("archive/object", Some("historical-version")) + .await + .expect("exact version proof"), + Probe::VersionedPresent("historical-version".to_string()) + ); + let requests = fixture.await.expect("exact probe fixture should finish"); + assert!( + requests[1] + .lines() + .next() + .expect("request line") + .contains("versionId=historical-version") + ); + assert!(requests[1].to_ascii_lowercase().contains("range: bytes=0-0")); + assert!(requests.iter().all(|request| request.starts_with("GET "))); + } + + #[tokio::test] + async fn legacy_transition_state_probe_retains_multiple_candidates() { + use super::super::warm_backend::LegacyTransitionStateProbe as Probe; + let responses = vec![ + legacy_probe_versioning_response("Enabled"), + legacy_probe_versions_response(&["version-a", "version-b"]), + ]; + let (backend, fixture) = scripted_probe_fixture(responses) + .await + .expect("ambiguous legacy probe fixture"); + assert_eq!( + backend + .probe_legacy_transition_state("archive/object", None) + .await + .expect("ambiguous proof"), + Probe::Ambiguous + ); + let requests = fixture.await.expect("ambiguous fixture should finish"); + assert_eq!(requests.len(), 2); + assert!(requests.iter().all(|request| request.starts_with("GET "))); + } + fn list_versions(versions: &[(&str, &str)], delete_markers: &[(&str, &str)], is_truncated: bool) -> ListVersionsResult { ListVersionsResult { versions: versions @@ -921,6 +1034,56 @@ impl WarmBackend for WarmBackendS3 { #[async_trait::async_trait] impl TransitionCandidateReconciler for WarmBackendS3 { + async fn probe_legacy_transition_state( + &self, + object: &str, + remote_version: Option<&str>, + ) -> Result { + use super::warm_backend::LegacyTransitionStateProbe as Probe; + + let initial_versioning = self.remote_bucket_versioning().await?; + let candidate = if let Some(version) = remote_version.filter(|version| !version.is_empty()) { + validate_remote_version_id(version)?; + match self.probe_transition_version(object, version).await? { + TransitionCandidateProbe::VersionedPresent(actual) if actual == version => Some(actual), + TransitionCandidateProbe::Missing => return Ok(Probe::Missing), + _ => return Ok(Probe::Ambiguous), + } + } else { + let remote_object = self.get_dest(object); + let mut opts = ListObjectsOptions::default(); + opts.set("prefix", &remote_object); + opts.set("max-keys", "1000"); + let mut key_marker = String::new(); + let mut version_marker = String::new(); + let mut candidates = TransitionCandidateVersions::default(); + let mut complete = false; + // This is one synchronous record inspection, not an unbounded + // remote history scan. The caller also bounds the whole probe. + for _ in 0..128 { + let page = self + .client + .list_object_versions_query(&self.bucket, &opts, &key_marker, &version_marker, "") + .await?; + candidates.extend(&remote_object, &page); + if candidates.ambiguous { + return Ok(Probe::Ambiguous); + } + if !page.is_truncated { + complete = true; + break; + } + advance_version_markers(&mut key_marker, &mut version_marker, &page)?; + } + if !complete { + return Ok(Probe::Ambiguous); + } + candidates.version_id + }; + let confirmed_versioning = self.remote_bucket_versioning().await?; + classify_legacy_transition_state(candidate.as_deref(), initial_versioning, confirmed_versioning) + } + async fn probe_transition_candidate_for( &self, object: &str, @@ -931,3 +1094,32 @@ impl TransitionCandidateReconciler for WarmBackendS3 { .await } } + +fn classify_legacy_transition_state( + candidate: Option<&str>, + initial: RemoteBucketVersioning, + confirmed: RemoteBucketVersioning, +) -> Result { + use super::warm_backend::LegacyTransitionStateProbe as Probe; + if initial != confirmed { + return Ok(Probe::Ambiguous); + } + let Some(version) = candidate else { + return Ok(Probe::Missing); + }; + if !version.is_empty() { + validate_remote_version_id(version)?; + if uuid::Uuid::parse_str(version).is_ok_and(|id| id.is_nil()) { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "legacy tier probe returned a nil version identifier", + )); + } + } + Ok(match (confirmed, version) { + (RemoteBucketVersioning::Disabled, "" | "null") => Probe::UnversionedPresent, + (RemoteBucketVersioning::Disabled, _) | (_, "") => Probe::Ambiguous, + (_, "null") => Probe::SuspendedNullPresent, + (_, version) => Probe::VersionedPresent(version.to_string()), + }) +} diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index d2b63795d..ed7840923 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -3888,6 +3888,31 @@ pub struct SetDisks { >, } +/// Read every physical copy before selecting a version quorum. A minority +/// legacy record is still evidence and must not disappear behind a majority +/// not-found result. Only an explicit file/volume absence produces `None`; +/// an unreadable disk cannot prove that no conflicting copy exists. +pub(crate) async fn read_legacy_transition_state_metadata_copies( + set: &SetDisks, + bucket: &str, + object: &str, +) -> std::result::Result>>, DiskError> { + let disk_object = rustfs_utils::path::encode_dir_object(object); + let disks = set.get_disks_internal().await; + if disks.is_empty() { + return Err(DiskError::DiskNotFound); + } + + let (copies, errs) = SetDisks::read_all_raw_file_info(&disks, bucket, disk_object.as_str(), false).await; + for err in errs.into_iter().flatten() { + if !matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) { + return Err(err); + } + } + + Ok(copies.into_iter().map(|copy| copy.map(|copy| copy.buf)).collect()) +} + // DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones // each endpoint's canonical Arc, so an exact Arc set identifies the lock domain. pub(crate) fn same_distributed_lock_domain(left: &[Arc], right: &[Arc]) -> bool { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index c8d44beb0..07c97418a 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -12803,6 +12803,227 @@ mod tests { body } + #[cfg(feature = "test-util")] + #[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); + } + + #[cfg(feature = "test-util")] + async fn legacy_transition_state_inspection_and_apply_case() { + use crate::bucket::lifecycle::legacy_transition_state_reconcile::{ + LegacyTransitionStateReconcileOutcome as Outcome, LegacyTransitionStateReconcileRequest, + LegacyTransitionStateReconcileSelector, + }; + for (remote_version, expected_state) in [ + ("", rustfs_filemeta::TransitionVersionState::KnownDisabled), + ("null", rustfs_filemeta::TransitionVersionState::SuspendedNull), + ("opaque-version", rustfs_filemeta::TransitionVersionState::Exact), + ] { + let temp_dir = tempfile::tempdir().expect("legacy reconcile store directory"); + let (ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "legacy-state-reconcile-inspect", &[4])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let tier_name = "LEGACY-RECONCILE"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + backend.set_put_remote_version(Some(remote_version.to_string())).await; + let bucket = "legacy-state-reconcile-bucket"; + let object = "archive.bin"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("create legacy fixture bucket"); + let mut reader = PutObjReader::from_vec(b"legacy reconcile body".repeat(1024)); + let source = store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("write source"); + { + // Create the fixture under the existing remote-version writer + // gate. This does not authorize legacy metadata reconciliation. + let _proof = crate::services::notification_sys::install_current_remote_version_state_fleet_proof_for_test(); + temp_env::async_with_vars( + [ + (rustfs_config::ENV_TIER_REMOTE_VERSION_STATE_WRITE, Some("true")), + (rustfs_config::ENV_TIER_REMOTE_VERSION_STATE_FLEET_CONFIRMED, Some("true")), + ], + store.transition_object( + bucket, + object, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: source.etag.clone().expect("source ETag"), + ..Default::default() + }, + mod_time: source.mod_time, + ..Default::default() + }, + ), + ) + .await + .expect("transition source"); + } + assert!( + crate::services::notification_sys::acquire_legacy_transition_state_reconcile_fleet_proof() + .await + .is_none(), + "fixture setup must not grant the missing reconciliation write capability" + ); + let selector = LegacyTransitionStateReconcileSelector { + bucket: bucket.to_string(), + object: object.to_string(), + version_id: "null".to_string(), + }; + if expected_state == rustfs_filemeta::TransitionVersionState::Exact { + backend + .set_transition_candidate_probe_override(Some( + crate::services::tier::warm_backend::TransitionCandidateProbe::Ambiguous, + )) + .await; + } + let converged = store + .inspect_legacy_transition_state(selector.clone()) + .await + .expect("inspect an already explicit transition"); + assert_eq!(converged.outcome, Outcome::Migrated, "{converged:?}"); + assert!(!converged.changed); + backend.set_transition_candidate_probe_override(None).await; + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, remote_version.is_empty()).await; + let paths = (0..4) + .map(|disk| { + temp_dir + .path() + .join(format!("pool0/set0/disk{disk}/{bucket}/{object}/{STORAGE_FORMAT_FILE}")) + }) + .collect::>(); + let mut original = Vec::new(); + for path in &paths { + original.push(tokio::fs::read(path).await.expect("original xl.meta")); + } + backend.clear_op_log().await; + let inspection = store.inspect_legacy_transition_state(selector.clone()); + assert!( + std::mem::size_of_val(&inspection) <= 4 * 1024, + "admin inspection future must remain stack-bounded" + ); + let inspected = inspection.await.expect("inspect legacy state"); + assert_eq!(inspected.outcome, Outcome::ReadyToMigrate, "{inspected:?}"); + assert!(!inspected.readiness.post_ready, "current fleet cannot authorize conditional writes"); + let target = inspected.target.expect("live probe should establish one model"); + assert_eq!(target.state, expected_state); + let request = LegacyTransitionStateReconcileRequest { + confirm: true, + selector, + source: inspected.source.expect("immutable source"), + original_sets: inspected.original_sets, + target, + reconciliation_digest: inspected.reconciliation_digest.expect("expected tuple digest"), + }; + let mut tampered = request.clone(); + tampered.source.remote_object.push_str("-other"); + let probes_before = backend.op_log().await.len(); + let rejected = store + .reconcile_legacy_transition_state(tampered) + .await + .expect("reject tampered tuple"); + assert_eq!(rejected.outcome, Outcome::Corrupt); + assert_eq!(backend.op_log().await.len(), probes_before, "invalid digest must not probe the backend"); + let applied = store + .reconcile_legacy_transition_state(request) + .await + .expect("apply must report unavailable write authority"); + assert_eq!(applied.outcome, Outcome::BackendUnavailable, "{applied:?}"); + assert_eq!(applied.reason_code, "write_fence_unavailable"); + assert!(!applied.changed); + for (path, expected) in paths.iter().zip(&original) { + assert_eq!(tokio::fs::read(path).await.expect("xl.meta after inspection"), *expected); + } + assert_eq!(backend.remove_count().await, 0); + assert!( + backend + .op_log() + .await + .iter() + .all(|operation| matches!(operation, MockWarmOp::Probe { .. })) + ); + + backend.set_unreachable(true).await; + let unavailable = store + .inspect_legacy_transition_state(LegacyTransitionStateReconcileSelector { + bucket: bucket.to_string(), + object: object.to_string(), + version_id: "null".to_string(), + }) + .await + .expect("unreachable tier is a diagnostic outcome"); + assert_eq!(unavailable.outcome, Outcome::BackendUnavailable); + assert!(!unavailable.changed); + backend.set_unreachable(false).await; + for candidate in ["", "00000000-0000-0000-0000-000000000000", "bad\nversion"] { + backend + .set_transition_candidate_probe_override(Some( + crate::services::tier::warm_backend::TransitionCandidateProbe::VersionedPresent(candidate.to_string()), + )) + .await; + let invalid_proof = store + .inspect_legacy_transition_state(LegacyTransitionStateReconcileSelector { + bucket: bucket.to_string(), + object: object.to_string(), + version_id: "null".to_string(), + }) + .await + .expect("invalid backend proof is a diagnostic outcome"); + assert_eq!(invalid_proof.outcome, Outcome::BackendUnavailable, "{invalid_proof:?}"); + assert!(invalid_proof.target.is_none()); + } + backend.set_transition_candidate_probe_override(None).await; + + backend.clear_op_log().await; + for path in &paths[1..] { + tokio::fs::remove_file(path) + .await + .expect("hide majority metadata copies in fixture"); + } + let minority = store + .inspect_legacy_transition_state(LegacyTransitionStateReconcileSelector { + bucket: bucket.to_string(), + object: object.to_string(), + version_id: "null".to_string(), + }) + .await + .expect("inspect minority legacy record"); + assert_eq!( + minority.outcome, + Outcome::BackendUnavailable, + "a minority owner must remain visible: {minority:?}" + ); + assert!( + backend.op_log().await.is_empty(), + "unproven metadata quorum cannot initiate a remote probe" + ); + for (path, bytes) in paths.iter().zip(&original) { + tokio::fs::write(path, bytes).await.expect("restore fixture copies"); + } + tokio::fs::write(&paths[0], b"corrupt-xl-meta") + .await + .expect("inject corrupt copy"); + let corrupt = store + .inspect_legacy_transition_state(LegacyTransitionStateReconcileSelector { + bucket: bucket.to_string(), + object: object.to_string(), + version_id: "null".to_string(), + }) + .await + .expect("inspect corrupt legacy record"); + assert_eq!(corrupt.outcome, Outcome::Corrupt, "{corrupt:?}"); + assert!(backend.op_log().await.is_empty(), "corruption must fail before backend I/O"); + } + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] diff --git a/docs/architecture/ilm-tiering-persistence-contracts.md b/docs/architecture/ilm-tiering-persistence-contracts.md index 991dc6e8e..44a161505 100644 --- a/docs/architecture/ilm-tiering-persistence-contracts.md +++ b/docs/architecture/ilm-tiering-persistence-contracts.md @@ -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`. The current GET and free-version cleanup paths reject that state rather than interpreting an empty remote version as unversioned. There is no admin route that repairs this field in `xl.meta`. 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 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 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. diff --git a/docs/operations/tier-ilm-debugging.md b/docs/operations/tier-ilm-debugging.md index ceeff86cc..3cc99bc4e 100644 --- a/docs/operations/tier-ilm-debugging.md +++ b/docs/operations/tier-ilm-debugging.md @@ -262,9 +262,11 @@ The full schema, lease, mixed-version, retry, privacy, and metric requirements a ## Reconcile legacy transition-version metadata -This section describes an **approved target that is not implemented yet**. The current server has no admin route that backfills a missing `transitioned-version-state` in `xl.meta`. Do not use the transaction reconcile route above for this purpose: that route owns an upload transaction candidate and may delete it, while legacy metadata reconciliation is non-destructive and may update only the exact local metadata version. +The single-record inspection routes below are implemented. GET audits every physical metadata copy, verifies the legacy tuple, and performs a bounded live version-model probe. POST validates confirmation and the complete expected tuple, repeats inspection, and returns `backend-unavailable` with `reason_code: write_fence_unavailable`, `changed: false`, and `post_ready: false` when migration is needed. Persistent backfill remains disabled until conditional per-generation `xl.meta` writes and the dedicated fleet capability are available. An already explicit, converged record can return `migrated` without changing bytes. -The approved interface is synchronous and accepts exactly one bucket/object/local-version tuple: +Do not use the transaction reconcile route above for this purpose: that route owns an upload transaction candidate and may delete it, while legacy metadata inspection never mutates the remote tier or local metadata. + +The interface is synchronous and accepts exactly one bucket/object/local-version tuple: ```text GET /rustfs/admin/v3/ilm/transition/state/reconcile?bucket=&object=&versionId= @@ -273,7 +275,7 @@ POST /rustfs/admin/v3/ilm/transition/state/reconcile?bucket=&object=) -> std::i format!("{ADMIN_PREFIX}/v3/ilm/recovery/exports/{{export_id}}").as_str(), AdminOperation(&IlmRecoveryExportDownloadHandler {}), )?; + r.insert( + Method::GET, + format!("{ADMIN_PREFIX}/v3/ilm/transition/state/reconcile").as_str(), + AdminOperation(&LegacyTransitionStateReconcileInspectHandler {}), + )?; + r.insert( + Method::POST, + format!("{ADMIN_PREFIX}/v3/ilm/transition/state/reconcile").as_str(), + AdminOperation(&LegacyTransitionStateReconcileApplyHandler {}), + )?; Ok(()) } @@ -1001,6 +1012,92 @@ fn validate_recovery_observation_receipt( Ok(receipt.observation) } +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyTransitionStateReconcileQuery { + bucket: Option, + object: Option, + #[serde(rename = "versionId")] + version_id: Option, +} + +fn parse_legacy_transition_state_reconcile_query(query: Option<&str>) -> S3Result { + let query: LegacyTransitionStateReconcileQuery = serde_urlencoded::from_bytes(query.unwrap_or_default().as_bytes()) + .map_err(|_| s3_error!(InvalidArgument, "invalid legacy transition-state reconcile query"))?; + let bucket = query + .bucket + .filter(|bucket| !bucket.is_empty()) + .ok_or_else(|| s3_error!(InvalidRequest, "bucket is required"))?; + if is_reserved_or_invalid_bucket(&bucket, false) { + return Err(s3_error!(InvalidBucketName, "invalid bucket name")); + } + + let object = query + .object + .filter(|object| !object.is_empty()) + .ok_or_else(|| s3_error!(InvalidRequest, "object is required"))?; + if !is_valid_object_prefix(&object) || object.contains('\n') || object.contains('\r') { + return Err(s3_error!(InvalidArgument, "invalid object name")); + } + + let version_id = query + .version_id + .filter(|version_id| !version_id.is_empty()) + .ok_or_else(|| s3_error!(InvalidRequest, "versionId is required"))?; + let version_id = if version_id == "null" { + version_id + } else { + let parsed = Uuid::parse_str(&version_id).map_err(|_| s3_error!(InvalidArgument, "invalid local versionId"))?; + if parsed.is_nil() { + return Err(s3_error!(InvalidArgument, "invalid local versionId")); + } + parsed.to_string() + }; + + Ok(LegacyTransitionStateReconcileSelector { + bucket, + object, + version_id, + }) +} + +fn validate_legacy_transition_state_reconcile_request( + query_selector: &LegacyTransitionStateReconcileSelector, + confirm: bool, + request_selector: &LegacyTransitionStateReconcileSelector, +) -> S3Result<()> { + if !confirm { + return Err(s3_error!( + InvalidRequest, + "legacy transition-state reconciliation requires confirm=true; use GET to inspect without changes" + )); + } + if request_selector != query_selector { + return Err(s3_error!(InvalidRequest, "request selector must exactly match the query selector")); + } + Ok(()) +} + +fn map_legacy_transition_state_reconcile_error(err: LegacyTransitionStateReconcileError) -> S3Error { + match err { + LegacyTransitionStateReconcileError::InvalidSelector(_) | LegacyTransitionStateReconcileError::InvalidRequest(_) => { + s3_error!(InvalidRequest, "invalid legacy transition-state reconciliation request") + } + LegacyTransitionStateReconcileError::StaleExpectedTuple(_) | LegacyTransitionStateReconcileError::Corrupt(_) => { + s3_error!(OperationAborted, "legacy transition-state reconciliation metadata is stale or corrupt") + } + LegacyTransitionStateReconcileError::WriteFenceUnavailable(_) => { + s3_error!( + OperationAborted, + "legacy transition-state reconciliation could not acquire safe write authority" + ) + } + LegacyTransitionStateReconcileError::BackendUnavailable(_) => { + s3_error!(InternalError, "legacy transition-state reconciliation backend is unavailable") + } + } +} + fn map_transition_operator_error(err: TransitionOperatorError) -> S3Error { match err { TransitionOperatorError::NotFound => s3_error!(NoSuchKey, "transition transaction not found"), @@ -1961,6 +2058,52 @@ impl Operation for TransitionReconcileApplyHandler { } } +pub struct LegacyTransitionStateReconcileInspectHandler {} + +#[async_trait::async_trait] +impl Operation for LegacyTransitionStateReconcileInspectHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + authorize_transition_admin_request(&req, AdminAction::ListTierAction).await?; + let selector = parse_legacy_transition_state_reconcile_query(req.uri.query())?; + let Some(store) = object_store_from_extensions(&req.extensions) else { + return Err(s3_error!(InternalError, "object store is not initialized")); + }; + let response: LegacyTransitionStateReconcileResponse = store + .inspect_legacy_transition_state(selector) + .await + .map_err(map_legacy_transition_state_reconcile_error)?; + json_response(StatusCode::OK, &response) + } +} + +pub struct LegacyTransitionStateReconcileApplyHandler {} + +#[async_trait::async_trait] +impl Operation for LegacyTransitionStateReconcileApplyHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + authorize_transition_admin_request(&req, AdminAction::SetTierAction).await?; + let selector = parse_legacy_transition_state_reconcile_query(req.uri.query())?; + let store = object_store_from_extensions(&req.extensions); + let mut input = req.input; + let body = input + .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) + .await + .map_err(|_| s3_error!(InvalidRequest, "legacy transition-state reconciliation body is too large or unreadable"))?; + let request: LegacyTransitionStateReconcileRequest = serde_json::from_slice(&body) + .map_err(|_| s3_error!(InvalidRequest, "legacy transition-state reconciliation request must be valid JSON"))?; + validate_legacy_transition_state_reconcile_request(&selector, request.confirm, &request.selector)?; + let Some(store) = store else { + return Err(s3_error!(InternalError, "object store is not initialized")); + }; + + let response: LegacyTransitionStateReconcileResponse = store + .reconcile_legacy_transition_state(request) + .await + .map_err(map_legacy_transition_state_reconcile_error)?; + json_response(StatusCode::OK, &response) + } +} + #[cfg(test)] mod tests { use super::*; @@ -2592,12 +2735,87 @@ mod tests { let apply = src .split("impl Operation for TransitionReconcileApplyHandler") .nth(1) - .and_then(|block| block.split("#[cfg(test)]").next()) + .and_then(|block| block.split("pub struct LegacyTransitionStateReconcileInspectHandler").next()) .expect("apply handler block"); assert!(apply.contains("AdminAction::SetTierAction")); assert!(!apply.contains("AdminAction::ListTierAction")); } + #[test] + fn legacy_transition_state_query_requires_one_exact_selector() { + let version_id = Uuid::new_v4(); + let query = format!("bucket=test-bucket&object=logs%2F2026%20report&versionId={version_id}"); + let selector = parse_legacy_transition_state_reconcile_query(Some(&query)).expect("exact version selector should parse"); + assert_eq!(selector.bucket, "test-bucket"); + assert_eq!(selector.object, "logs/2026 report"); + assert_eq!(selector.version_id, version_id.to_string()); + + let unversioned = + parse_legacy_transition_state_reconcile_query(Some("bucket=test-bucket&object=logs%2Fcurrent&versionId=null")) + .expect("explicit null selector should parse"); + assert_eq!(unversioned.version_id, "null"); + + for query in [ + None, + Some("bucket=test-bucket&object=key"), + Some("bucket=test-bucket&object=&versionId=null"), + Some("bucket=test-bucket&object=key&versionId="), + Some("bucket=test-bucket&object=key&versionId=not-a-uuid"), + Some("bucket=test-bucket&object=key&versionId=00000000-0000-0000-0000-000000000000"), + Some("bucket=test-bucket&object=bad%0Akey&versionId=null"), + Some("bucket=test-bucket&object=key&versionId=null&prefix=wide"), + Some("bucket=test-bucket&bucket=other-bucket&object=key&versionId=null"), + ] { + assert!( + parse_legacy_transition_state_reconcile_query(query).is_err(), + "query should fail closed: {query:?}" + ); + } + } + + #[test] + fn legacy_transition_state_apply_requires_confirmation_and_matching_selector() { + let selector = LegacyTransitionStateReconcileSelector { + bucket: "test-bucket".to_string(), + object: "key".to_string(), + version_id: "null".to_string(), + }; + let different = LegacyTransitionStateReconcileSelector { + object: "other-key".to_string(), + ..selector.clone() + }; + + assert!(validate_legacy_transition_state_reconcile_request(&selector, false, &selector).is_err()); + assert!(validate_legacy_transition_state_reconcile_request(&selector, true, &different).is_err()); + validate_legacy_transition_state_reconcile_request(&selector, true, &selector) + .expect("confirmed exact selector should pass handler validation"); + } + + #[test] + fn legacy_transition_state_routes_use_read_and_write_tier_actions() { + let src = include_str!("ilm_transition.rs"); + let inspect = src + .split("impl Operation for LegacyTransitionStateReconcileInspectHandler") + .nth(1) + .and_then(|block| { + block + .split("impl Operation for LegacyTransitionStateReconcileApplyHandler") + .next() + }) + .expect("legacy inspect handler block"); + assert!(inspect.contains("AdminAction::ListTierAction")); + assert!(!inspect.contains("AdminAction::SetTierAction")); + + let apply = src + .split("impl Operation for LegacyTransitionStateReconcileApplyHandler") + .nth(1) + .and_then(|block| block.split("#[cfg(test)]").next()) + .expect("legacy apply handler block"); + assert!(apply.contains("AdminAction::SetTierAction")); + assert!(!apply.contains("AdminAction::ListTierAction")); + assert!(apply.contains("validate_legacy_transition_state_reconcile_request")); + } + #[test] fn manual_transition_query_defaults_to_bounded_run() { let (bucket, options, run_mode) = @@ -3066,6 +3284,24 @@ mod tests { assert_eq!(cancel_err.message(), Some("authentication required")); } + #[tokio::test] + async fn legacy_transition_state_handlers_reject_missing_credentials_before_selector_or_body() { + let path = "/rustfs/admin/v3/ilm/transition/state/reconcile?bucket=test-bucket&object=key&versionId=null"; + let inspect_err = LegacyTransitionStateReconcileInspectHandler {} + .call(credential_less_admin_request(Method::GET, path), Params::new()) + .await + .expect_err("inspect handler must reject unsigned requests"); + assert_eq!(inspect_err.code(), &S3ErrorCode::InvalidRequest); + assert_eq!(inspect_err.message(), Some("authentication required")); + + let apply_err = LegacyTransitionStateReconcileApplyHandler {} + .call(credential_less_admin_request(Method::POST, path), Params::new()) + .await + .expect_err("apply handler must reject unsigned requests"); + assert_eq!(apply_err.code(), &S3ErrorCode::InvalidRequest); + assert_eq!(apply_err.message(), Some("authentication required")); + } + #[test] fn manual_transition_job_handlers_authorize_validate_and_load_store() { let src = include_str!("ilm_transition.rs"); diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index aebf92849..7a6137c21 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -550,6 +550,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ SET_TIER, RouteRiskLevel::High, ), + admin( + HttpMethod::Get, + "/rustfs/admin/v3/ilm/transition/state/reconcile", + LIST_TIER, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Post, + "/rustfs/admin/v3/ilm/transition/state/reconcile", + SET_TIER, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/rustfs/admin/v3/audit/target/list", @@ -2212,12 +2224,16 @@ mod tests { assert_action(HttpMethod::Delete, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SET_TIER); assert_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/reconcile/{transaction_id}", LIST_TIER); assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/reconcile/{transaction_id}", SET_TIER); + assert_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/state/reconcile", LIST_TIER); + assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/state/reconcile", SET_TIER); assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SERVER_INFO); assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/recovery/records", SERVER_INFO); assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SERVER_INFO); assert_not_action(HttpMethod::Delete, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SERVER_INFO); assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/reconcile/{transaction_id}", SET_TIER); assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/reconcile/{transaction_id}", LIST_TIER); + assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/state/reconcile", SET_TIER); + assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/state/reconcile", LIST_TIER); } #[test] diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 5a1e66ced..1673a00a9 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -241,6 +241,8 @@ fn expected_admin_route_matrix() -> Vec { "/v3/ilm/transition/reconcile/{transaction_id}", "/v3/ilm/transition/reconcile/11111111-1111-4111-8111-111111111111", ), + admin_route(Method::GET, "/v3/ilm/transition/state/reconcile"), + admin_route(Method::POST, "/v3/ilm/transition/state/reconcile"), admin_route_sample( Method::DELETE, "/v3/ilm/transition/jobs/{job_id}", @@ -993,6 +995,8 @@ fn test_register_routes_cover_representative_admin_paths() { Method::POST, &admin_path("/v3/ilm/transition/reconcile/11111111-1111-4111-8111-111111111111"), ); + assert_route(&router, Method::GET, &admin_path("/v3/ilm/transition/state/reconcile")); + assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/state/reconcile")); assert_route(&router, Method::GET, &table_catalog_path("/config")); assert_route(&router, Method::PUT, &table_catalog_path("/buckets/analytics")); diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index a9ef50c4e..0ee0f9099 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -231,6 +231,10 @@ pub(crate) mod lifecycle { renew_manual_transition_job_lease_if_owned, request_manual_transition_job_cancel, save_manual_transition_job_record, update_manual_transition_job_record, }; + pub(crate) use crate::storage::storage_api::ecstore_bucket::lifecycle::legacy_transition_state_reconcile::{ + LegacyTransitionStateReconcileError, LegacyTransitionStateReconcileRequest, LegacyTransitionStateReconcileResponse, + LegacyTransitionStateReconcileSelector, + }; pub(crate) type ManualTransitionCancelCheck = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionCancelCheck; pub(crate) type ManualTransitionQueueSnapshot =